mirror of
https://github.com/Tria-plc/edr-platform.git
synced 2026-08-28 04:20:55 +00:00
Merge branch 'dev' of github.com:Tria-plc/edr-platform into freight_feature/usermanagement
This commit is contained in:
@@ -19,6 +19,7 @@ import { CargoTypesService } from '../rule-engine/services/cargo-types.service';
|
||||
import { DropdownSettingsService } from '../dropdown-settings/dropdown-settings.service';
|
||||
import { FilesService } from '../files/files.service';
|
||||
import { SignaturesService } from '../signatures/signatures.service';
|
||||
import { OtpService } from '../otp/otp.service';
|
||||
import { ContractPricingService } from './contract-pricing.service';
|
||||
import { ClearanceMilestoneService } from './clearance-milestone.service';
|
||||
import { ContractsRepository } from './contracts.repository';
|
||||
@@ -63,6 +64,7 @@ export class ContractTransitionService {
|
||||
private readonly renderer: ContractRendererService,
|
||||
private readonly pdfService: ContractPdfService,
|
||||
private readonly minioService: MinioService,
|
||||
private readonly otpService: OtpService,
|
||||
) {}
|
||||
|
||||
/** Customer submits the contract for approval → SUBMITTED; freeze unit rates. */
|
||||
@@ -520,6 +522,12 @@ export class ContractTransitionService {
|
||||
if (existing) {
|
||||
throw new BadRequestException('Customer has already signed this contract');
|
||||
}
|
||||
// Sudo-mode gate: a fresh, single-use OTP (SMS'd to the customer's phone)
|
||||
// must be verified before the signature is applied.
|
||||
if (!dto.otpPhone || !dto.otp) {
|
||||
throw new BadRequestException('OTP verification is required to sign the contract');
|
||||
}
|
||||
await this.otpService.verifyOtpForAction(dto.otpPhone, dto.otp);
|
||||
await this.applySignature(contract, dto, options);
|
||||
await this.contractsRepository.update(contractId, {
|
||||
status: 'SIGNED_CUSTOMER',
|
||||
|
||||
@@ -11,6 +11,7 @@ import { RuleEngineModule } from '../rule-engine/rule-engine.module';
|
||||
import { FileUploadSettingsModule } from '../file-upload-settings/file-upload-settings.module';
|
||||
import { DropdownSettingsModule } from '../dropdown-settings/dropdown-settings.module';
|
||||
import { SignaturesModule } from '../signatures/signatures.module';
|
||||
import { OtpModule } from '../otp/otp.module';
|
||||
import { BookingsModule } from '../bookings/bookings.module';
|
||||
import { TrainSchedulingModule } from '../train-scheduling/train-scheduling.module';
|
||||
|
||||
@@ -73,6 +74,7 @@ import { ContractDocumentViewModelBuilder } from '../../contracts/contract-docum
|
||||
FilesModule,
|
||||
MinioModule,
|
||||
SignaturesModule,
|
||||
OtpModule,
|
||||
CompaniesModule,
|
||||
// BookingsModule provides BookingsRepository/BookingPricingService used by the
|
||||
// contract PDF builders (they read a Booking today — see docs/new-doc.md §3.3).
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger';
|
||||
import { IsIn, IsOptional, IsString, MinLength } from 'class-validator';
|
||||
import { IsIn, IsOptional, IsString, Matches, MinLength } from 'class-validator';
|
||||
|
||||
export class SignContractDto {
|
||||
@ApiProperty({ enum: ['CUSTOMER', 'STAFF', 'DIRECTOR', 'CEO'] })
|
||||
@@ -26,4 +26,19 @@ export class SignContractDto {
|
||||
@IsOptional()
|
||||
@IsString()
|
||||
consentText?: string;
|
||||
|
||||
// Sudo-mode OTP challenge. Required when role=CUSTOMER: a fresh 6-digit code
|
||||
// SMS'd to the signer's phone, verified server-side before the signature is
|
||||
// applied. `otpPhone` is the number the code was sent to (the signed-in
|
||||
// customer's registered phone).
|
||||
@ApiPropertyOptional({ description: '6-digit OTP; required when role=CUSTOMER' })
|
||||
@IsOptional()
|
||||
@IsString()
|
||||
@Matches(/^\d{6}$/, { message: 'otp must be 6 digits' })
|
||||
otp?: string;
|
||||
|
||||
@ApiPropertyOptional({ description: 'Phone the OTP was sent to; required when role=CUSTOMER' })
|
||||
@IsOptional()
|
||||
@IsString()
|
||||
otpPhone?: string;
|
||||
}
|
||||
|
||||
@@ -14,13 +14,17 @@ import { FleetManage, FleetView } from '../../common/booking-guards';
|
||||
import { DriversService } from './drivers.service';
|
||||
import { CreateDriverDto } from './dto/create-driver.dto';
|
||||
import { UpdateDriverDto } from './dto/update-driver.dto';
|
||||
import { FleetHistoryService } from '../fleet-history/fleet-history.service';
|
||||
|
||||
@ApiTags('drivers')
|
||||
@ApiBearerAuth()
|
||||
@Controller('drivers')
|
||||
@FleetView()
|
||||
export class DriversController {
|
||||
constructor(private readonly driversService: DriversService) {}
|
||||
constructor(
|
||||
private readonly driversService: DriversService,
|
||||
private readonly fleetHistory: FleetHistoryService,
|
||||
) {}
|
||||
|
||||
@Post()
|
||||
@FleetManage()
|
||||
@@ -55,6 +59,12 @@ export class DriversController {
|
||||
return this.driversService.findById(id);
|
||||
}
|
||||
|
||||
@Get(':id/history')
|
||||
@ApiOperation({ summary: 'Get driver assignment & activity history' })
|
||||
history(@Param('id', ParseUUIDPipe) id: string) {
|
||||
return this.fleetHistory.getDriverHistory(id);
|
||||
}
|
||||
|
||||
@Patch(':id')
|
||||
@FleetManage()
|
||||
@ApiOperation({ summary: 'Update a driver' })
|
||||
|
||||
@@ -1,18 +1,27 @@
|
||||
import { Injectable, NotFoundException, ConflictException } from '@nestjs/common';
|
||||
import { Injectable, NotFoundException, ConflictException, BadRequestException } from '@nestjs/common';
|
||||
import { InjectRepository } from '@nestjs/typeorm';
|
||||
import { Repository } from 'typeorm';
|
||||
import { CreateDriverDto } from './dto/create-driver.dto';
|
||||
import { UpdateDriverDto } from './dto/update-driver.dto';
|
||||
import { Driver, DriverStatus } from './entities/driver.entity';
|
||||
import { FleetHistoryService } from '../fleet-history/fleet-history.service';
|
||||
import { FleetEventType } from '../fleet-history/entities/fleet-event.entity';
|
||||
|
||||
@Injectable()
|
||||
export class DriversService {
|
||||
constructor(
|
||||
@InjectRepository(Driver)
|
||||
private readonly driverRepo: Repository<Driver>,
|
||||
private readonly history: FleetHistoryService,
|
||||
) {}
|
||||
|
||||
async create(dto: CreateDriverDto): Promise<Driver> {
|
||||
if (dto.faydaVerified !== true) {
|
||||
throw new BadRequestException(
|
||||
'Driver identity must be verified with Fayda before saving',
|
||||
);
|
||||
}
|
||||
|
||||
const existing = await this.driverRepo.findOne({
|
||||
where: [
|
||||
{ licenseNumber: dto.licenseNumber },
|
||||
@@ -33,8 +42,28 @@ export class DriversService {
|
||||
}
|
||||
}
|
||||
|
||||
if (dto.faydaSub) {
|
||||
const dupe = await this.driverRepo.findOne({
|
||||
where: { faydaSub: dto.faydaSub },
|
||||
});
|
||||
if (dupe) {
|
||||
throw new ConflictException(
|
||||
'A driver is already registered for this Fayda identity',
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
const driver = this.driverRepo.create(dto);
|
||||
return this.driverRepo.save(driver);
|
||||
const saved = await this.driverRepo.save(driver);
|
||||
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.DRIVER_REGISTERED,
|
||||
driverId: saved.id,
|
||||
label: `${saved.firstName ?? ''} ${saved.lastName ?? ''}`.trim() || null,
|
||||
toValue: saved.status ?? null,
|
||||
});
|
||||
|
||||
return saved;
|
||||
}
|
||||
|
||||
async findAll(query: {
|
||||
@@ -106,7 +135,25 @@ export class DriversService {
|
||||
}
|
||||
}
|
||||
|
||||
if (dto.faydaSub && dto.faydaSub !== driver.faydaSub) {
|
||||
const dupe = await this.driverRepo.findOne({
|
||||
where: { faydaSub: dto.faydaSub },
|
||||
});
|
||||
if (dupe) {
|
||||
throw new ConflictException(
|
||||
'A driver is already registered for this Fayda identity',
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Object.assign(driver, dto);
|
||||
|
||||
if (driver.faydaVerified !== true) {
|
||||
throw new BadRequestException(
|
||||
'Driver identity must be verified with Fayda before saving',
|
||||
);
|
||||
}
|
||||
|
||||
return this.driverRepo.save(driver);
|
||||
}
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { IsString, IsEmail, IsDateString, IsEnum, IsOptional, IsArray } from 'class-validator';
|
||||
import { DriverStatus } from '../entities/driver.entity';
|
||||
import { IsString, IsEmail, IsDateString, IsEnum, IsOptional, IsArray, IsBoolean } from 'class-validator';
|
||||
import { DriverStatus, DriverGender } from '../entities/driver.entity';
|
||||
|
||||
export class CreateDriverDto {
|
||||
@IsString()
|
||||
@@ -20,6 +20,10 @@ export class CreateDriverDto {
|
||||
@IsDateString()
|
||||
dateOfBirth!: string;
|
||||
|
||||
@IsOptional()
|
||||
@IsEnum(DriverGender)
|
||||
gender?: DriverGender;
|
||||
|
||||
@IsDateString()
|
||||
licenseExpiryDate!: string;
|
||||
|
||||
@@ -42,4 +46,12 @@ export class CreateDriverDto {
|
||||
@IsOptional()
|
||||
@IsString()
|
||||
notes?: string;
|
||||
|
||||
@IsOptional()
|
||||
@IsBoolean()
|
||||
faydaVerified?: boolean;
|
||||
|
||||
@IsOptional()
|
||||
@IsString()
|
||||
faydaSub?: string;
|
||||
}
|
||||
|
||||
@@ -8,6 +8,12 @@ export enum DriverStatus {
|
||||
ON_LEAVE = 'ON_LEAVE',
|
||||
}
|
||||
|
||||
export enum DriverGender {
|
||||
MALE = 'MALE',
|
||||
FEMALE = 'FEMALE',
|
||||
OTHER = 'OTHER',
|
||||
}
|
||||
|
||||
@Entity({ name: 'drivers', schema: 'freight' })
|
||||
export class Driver extends BaseEntity {
|
||||
@Column({ name: 'license_number', unique: true, nullable: true })
|
||||
@@ -28,6 +34,9 @@ export class Driver extends BaseEntity {
|
||||
@Column({ name: 'date_of_birth', type: 'date', nullable: true })
|
||||
dateOfBirth?: Date;
|
||||
|
||||
@Column({ type: 'varchar', nullable: true })
|
||||
gender?: DriverGender | null;
|
||||
|
||||
@Column({ name: 'license_expiry_date', type: 'date', nullable: true })
|
||||
licenseExpiryDate?: Date;
|
||||
|
||||
@@ -51,4 +60,12 @@ export class Driver extends BaseEntity {
|
||||
|
||||
@Column({ type: 'numeric', precision: 3, scale: 2, nullable: true })
|
||||
rating?: number | null;
|
||||
|
||||
@Column({ name: 'fayda_verified', type: 'boolean', default: false, nullable: true })
|
||||
faydaVerified?: boolean;
|
||||
|
||||
/** Fayda OIDC subject the identity was verified against. Unique — one driver
|
||||
* record per verified Fayda identity (NULLs allowed for legacy/unverified). */
|
||||
@Column({ name: 'fayda_sub', type: 'varchar', unique: true, nullable: true })
|
||||
faydaSub?: string | null;
|
||||
}
|
||||
|
||||
@@ -71,7 +71,7 @@ export class FirstMileInvoiceService {
|
||||
type: 'DELIVERY_FEE',
|
||||
companyId: fm.booking!.companyId,
|
||||
companyProfileId: fm.booking!.companyProfileId || '',
|
||||
currency: 'ETB',
|
||||
currency: fm.booking!.paymentCurrency || 'ETB',
|
||||
lines: [
|
||||
{
|
||||
chargeType: 'DELIVERY',
|
||||
|
||||
@@ -92,6 +92,7 @@ export class FirstMileController {
|
||||
const record = await this.firstMileService.update(id, dto);
|
||||
// Auto-generate invoice if distance or payment was updated
|
||||
const booking = await this.bookingsService.findById(record.bookingId);
|
||||
const currency = booking.paymentCurrency || "ETB";
|
||||
if (dto.exactKm !== undefined || dto.exactKm != record.exactKm || dto.remainingPayment !== undefined || dto.remainingPayment !== record.remainingPayment) {
|
||||
await this.billingService.generateInvoice({
|
||||
source: Freight.InvoiceSource.FirstMile,
|
||||
@@ -99,7 +100,7 @@ export class FirstMileController {
|
||||
type: "FIRST_MILE",
|
||||
companyId: booking.companyId,
|
||||
companyProfileId: booking.companyProfileId,
|
||||
currency: "ETB",
|
||||
currency,
|
||||
|
||||
lines: [
|
||||
{
|
||||
@@ -108,7 +109,7 @@ export class FirstMileController {
|
||||
quantity: 1,
|
||||
unitRate: record.remainingPayment,
|
||||
amount: record.remainingPayment,
|
||||
currency: "ETB",
|
||||
currency,
|
||||
},
|
||||
],
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { Injectable, Logger, NotFoundException } from '@nestjs/common';
|
||||
import { FindOptionsWhere, In } from 'typeorm';
|
||||
import { BadRequestException, Injectable, Logger, NotFoundException } from '@nestjs/common';
|
||||
import { FindOptionsWhere, In, IsNull, Not } from 'typeorm';
|
||||
import { InjectDataSource } from '@nestjs/typeorm';
|
||||
import { DataSource } from 'typeorm';
|
||||
import { VehicleAvailability } from '../vehicles/entities/vehicle.entity';
|
||||
@@ -14,6 +14,8 @@ import { FirstMileContainerAllocation } from "./entities/first-mile-container-al
|
||||
import { FirstMileRepository } from "./first-mile.repository";
|
||||
import { OnEvent } from "@nestjs/event-emitter";
|
||||
import { InvoiceEventPayload } from "../billing/billing.service";
|
||||
import { FleetHistoryService } from "../fleet-history/fleet-history.service";
|
||||
import { FleetEventType } from "../fleet-history/entities/fleet-event.entity";
|
||||
|
||||
type FirstMileListFilter = {
|
||||
status?: FirstMileStatus;
|
||||
@@ -43,8 +45,54 @@ export class FirstMileService {
|
||||
private readonly vehiclesService: VehiclesService,
|
||||
private readonly driversService: DriversService,
|
||||
private readonly smsClient: SmsClientService,
|
||||
private readonly history: FleetHistoryService,
|
||||
) { }
|
||||
|
||||
/** Resolve a vehicle's driver + human labels, for stamping mile events onto
|
||||
* the driver's timeline and naming the vehicle. Best-effort — never throws. */
|
||||
private async vehicleInfo(
|
||||
vehicleId?: string | null,
|
||||
): Promise<{ driverId: string | null; plate: string | null; driverName: string | null }> {
|
||||
if (!vehicleId) return { driverId: null, plate: null, driverName: null };
|
||||
try {
|
||||
const v = await this.vehiclesService.findById(vehicleId);
|
||||
return {
|
||||
driverId: v.assignedDriverId ?? null,
|
||||
plate: v.plateNumber ?? v.code ?? null,
|
||||
driverName: v.assignedDriverName ?? null,
|
||||
};
|
||||
} catch {
|
||||
return { driverId: null, plate: null, driverName: null };
|
||||
}
|
||||
}
|
||||
|
||||
/** A leg counts as having a vehicle if it has a direct assignment or at least
|
||||
* one container allocation carrying a vehicle. Gates the IN_TRANSIT move. */
|
||||
private async hasAssignedVehicle(
|
||||
recordId: string,
|
||||
directVehicleId?: string | null,
|
||||
): Promise<boolean> {
|
||||
if (directVehicleId) return true;
|
||||
const count = await this.dataSource.manager.count(FirstMileContainerAllocation, {
|
||||
where: { firstMileId: recordId, vehicleId: Not(IsNull()) },
|
||||
});
|
||||
return count > 0;
|
||||
}
|
||||
|
||||
/** Human booking reference for a first-mile record, for the history timeline. */
|
||||
private async resolveBookingRef(record: FirstMile): Promise<string | null> {
|
||||
const loaded = (record as FirstMile & { booking?: { reference?: string } })
|
||||
.booking?.reference;
|
||||
if (loaded) return loaded;
|
||||
if (!record.bookingId) return null;
|
||||
try {
|
||||
const b = await this.bookingsRepository.findById(record.bookingId);
|
||||
return (b as { reference?: string } | null)?.reference ?? null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Look up a booking by its human-readable reference and confirm it has been
|
||||
* paid before any first-mile work proceeds. Throws if the reference is
|
||||
@@ -210,6 +258,20 @@ export class FirstMileService {
|
||||
|
||||
if (dto.vehicleId) {
|
||||
await this.vehiclesService.setAvailability(dto.vehicleId, VehicleAvailability.BUSY);
|
||||
const info = await this.vehicleInfo(dto.vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_VEHICLE_ASSIGNED,
|
||||
vehicleId: dto.vehicleId,
|
||||
firstMileId: record.id,
|
||||
driverId: info.driverId,
|
||||
label: record.status,
|
||||
metadata: {
|
||||
mile: 'FIRST',
|
||||
bookingRef: await this.resolveBookingRef(record),
|
||||
vehiclePlate: info.plate,
|
||||
driverName: info.driverName,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
return record;
|
||||
@@ -248,6 +310,18 @@ export class FirstMileService {
|
||||
async update(id: string, dto: UpdateFirstMileDto): Promise<FirstMile> {
|
||||
const existing = await this.findById(id);
|
||||
|
||||
// A leg can only go IN_TRANSIT once a vehicle is assigned (allowing a vehicle
|
||||
// assigned in this same request).
|
||||
if (dto.status === 'IN_TRANSIT' && existing.status !== 'IN_TRANSIT') {
|
||||
const vehicleId =
|
||||
dto.vehicleId !== undefined ? dto.vehicleId : existing.vehicleId;
|
||||
if (!(await this.hasAssignedVehicle(id, vehicleId))) {
|
||||
throw new BadRequestException(
|
||||
'Assign a vehicle before marking this first-mile leg in transit',
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
const dtoAny = dto as any;
|
||||
const updated = await this.firstMileRepository.update(id, {
|
||||
...(dto.bookingId !== undefined ? { bookingId: dto.bookingId } : {}),
|
||||
@@ -278,6 +352,39 @@ export class FirstMileService {
|
||||
if (existing.vehicleId) {
|
||||
await this.vehiclesService.releaseIfUnused([existing.vehicleId]);
|
||||
}
|
||||
// Audit the mile↔vehicle (re)assignment on both vehicle and driver lines.
|
||||
const bookingRef = await this.resolveBookingRef(existing);
|
||||
if (existing.vehicleId) {
|
||||
const info = await this.vehicleInfo(existing.vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_VEHICLE_RELEASED,
|
||||
vehicleId: existing.vehicleId,
|
||||
firstMileId: id,
|
||||
driverId: info.driverId,
|
||||
metadata: {
|
||||
mile: 'FIRST',
|
||||
bookingRef,
|
||||
vehiclePlate: info.plate,
|
||||
driverName: info.driverName,
|
||||
},
|
||||
});
|
||||
}
|
||||
if (dto.vehicleId) {
|
||||
const info = await this.vehicleInfo(dto.vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_VEHICLE_ASSIGNED,
|
||||
vehicleId: dto.vehicleId,
|
||||
firstMileId: id,
|
||||
driverId: info.driverId,
|
||||
label: updated.status,
|
||||
metadata: {
|
||||
mile: 'FIRST',
|
||||
bookingRef,
|
||||
vehiclePlate: info.plate,
|
||||
driverName: info.driverName,
|
||||
},
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
// Notify assigned driver on every explicit vehicle assignment or reassignment
|
||||
@@ -285,6 +392,25 @@ export class FirstMileService {
|
||||
void this.notifyDriverAssignment(dto.vehicleId, existing);
|
||||
}
|
||||
|
||||
if (dto.status !== undefined && dto.status !== existing.status) {
|
||||
const vehicleId = updated.vehicleId ?? existing.vehicleId ?? null;
|
||||
const info = await this.vehicleInfo(vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_STATUS_CHANGED,
|
||||
firstMileId: id,
|
||||
vehicleId,
|
||||
driverId: info.driverId,
|
||||
fromValue: existing.status,
|
||||
toValue: dto.status,
|
||||
metadata: {
|
||||
mile: 'FIRST',
|
||||
bookingRef: await this.resolveBookingRef(existing),
|
||||
vehiclePlate: info.plate,
|
||||
driverName: info.driverName,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
// Trip finished — release the vehicles it was holding
|
||||
if (dto.status === 'RECEIVED_TO_PORT' && existing.status !== 'RECEIVED_TO_PORT') {
|
||||
await this.releaseVehicles(updated);
|
||||
@@ -295,12 +421,40 @@ export class FirstMileService {
|
||||
|
||||
async updateStatus(id: string, status: FirstMileStatus): Promise<FirstMile> {
|
||||
const existing = await this.findById(id);
|
||||
|
||||
if (status === 'IN_TRANSIT' && existing.status !== 'IN_TRANSIT') {
|
||||
if (!(await this.hasAssignedVehicle(id, existing.vehicleId))) {
|
||||
throw new BadRequestException(
|
||||
'Assign a vehicle before marking this first-mile leg in transit',
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
const updated = await this.firstMileRepository.update(id, { status });
|
||||
|
||||
if (!updated) {
|
||||
throw new NotFoundException(`First-mile record ${id} not found`);
|
||||
}
|
||||
|
||||
if (status !== existing.status) {
|
||||
const vehicleId = updated.vehicleId ?? existing.vehicleId ?? null;
|
||||
const info = await this.vehicleInfo(vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_STATUS_CHANGED,
|
||||
firstMileId: id,
|
||||
vehicleId,
|
||||
driverId: info.driverId,
|
||||
fromValue: existing.status,
|
||||
toValue: status,
|
||||
metadata: {
|
||||
mile: 'FIRST',
|
||||
bookingRef: await this.resolveBookingRef(existing),
|
||||
vehiclePlate: info.plate,
|
||||
driverName: info.driverName,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
if (status === 'RECEIVED_TO_PORT' && existing.status !== 'RECEIVED_TO_PORT') {
|
||||
await this.releaseVehicles(updated);
|
||||
}
|
||||
@@ -430,6 +584,34 @@ export class FirstMileService {
|
||||
previousVehicleIds.filter((id) => !vehicleIds.has(id)),
|
||||
);
|
||||
|
||||
// History: one event per vehicle actually added or removed by this
|
||||
// multi-car (re)allocation, so reassignments show on every timeline.
|
||||
const prevSet = new Set(previousVehicleIds);
|
||||
const bookingRef = await this.resolveBookingRef(firstMile);
|
||||
for (const vehicleId of vehicleIds) {
|
||||
if (prevSet.has(vehicleId)) continue;
|
||||
const info = await this.vehicleInfo(vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_VEHICLE_ASSIGNED,
|
||||
vehicleId,
|
||||
firstMileId,
|
||||
driverId: info.driverId,
|
||||
label: firstMile.status,
|
||||
metadata: { mile: 'FIRST', bookingRef, vehiclePlate: info.plate, driverName: info.driverName },
|
||||
});
|
||||
}
|
||||
for (const vehicleId of previousVehicleIds) {
|
||||
if (vehicleIds.has(vehicleId)) continue;
|
||||
const info = await this.vehicleInfo(vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_VEHICLE_RELEASED,
|
||||
vehicleId,
|
||||
firstMileId,
|
||||
driverId: info.driverId,
|
||||
metadata: { mile: 'FIRST', bookingRef, vehiclePlate: info.plate, driverName: info.driverName },
|
||||
});
|
||||
}
|
||||
|
||||
return {
|
||||
success: true,
|
||||
allocated: allocations.length,
|
||||
|
||||
@@ -0,0 +1,55 @@
|
||||
import { Entity, Column, Index } from 'typeorm';
|
||||
import { BaseEntity } from '@edr/api-common';
|
||||
|
||||
/**
|
||||
* Append-only audit log for fleet activity. One row per transition. Queried by
|
||||
* `vehicleId` (vehicle timeline) or `driverId` (driver timeline); an event may
|
||||
* carry both so a driver↔vehicle assignment or a mile assignment shows on both.
|
||||
* `createdAt` (from BaseEntity) is the event time.
|
||||
*/
|
||||
export enum FleetEventType {
|
||||
DRIVER_REGISTERED = 'DRIVER_REGISTERED',
|
||||
VEHICLE_REGISTERED = 'VEHICLE_REGISTERED',
|
||||
DRIVER_ASSIGNED = 'DRIVER_ASSIGNED',
|
||||
DRIVER_UNASSIGNED = 'DRIVER_UNASSIGNED',
|
||||
VEHICLE_STATUS_CHANGED = 'VEHICLE_STATUS_CHANGED',
|
||||
VEHICLE_AVAILABILITY_CHANGED = 'VEHICLE_AVAILABILITY_CHANGED',
|
||||
MILE_VEHICLE_ASSIGNED = 'MILE_VEHICLE_ASSIGNED',
|
||||
MILE_VEHICLE_RELEASED = 'MILE_VEHICLE_RELEASED',
|
||||
MILE_STATUS_CHANGED = 'MILE_STATUS_CHANGED',
|
||||
}
|
||||
|
||||
@Entity({ name: 'fleet_events', schema: 'freight' })
|
||||
export class FleetEvent extends BaseEntity {
|
||||
@Column({ name: 'event_type', type: 'varchar' })
|
||||
eventType!: FleetEventType;
|
||||
|
||||
@Index()
|
||||
@Column({ name: 'vehicle_id', type: 'uuid', nullable: true })
|
||||
vehicleId?: string | null;
|
||||
|
||||
@Index()
|
||||
@Column({ name: 'driver_id', type: 'uuid', nullable: true })
|
||||
driverId?: string | null;
|
||||
|
||||
@Column({ name: 'first_mile_id', type: 'uuid', nullable: true })
|
||||
firstMileId?: string | null;
|
||||
|
||||
@Column({ name: 'last_mile_id', type: 'uuid', nullable: true })
|
||||
lastMileId?: string | null;
|
||||
|
||||
/** Previous value for a transition (e.g. old status/availability). */
|
||||
@Column({ name: 'from_value', type: 'varchar', nullable: true })
|
||||
fromValue?: string | null;
|
||||
|
||||
/** New value for a transition (e.g. new status/availability). */
|
||||
@Column({ name: 'to_value', type: 'varchar', nullable: true })
|
||||
toValue?: string | null;
|
||||
|
||||
/** Human-readable summary token (driver name, plate, booking ref, mile). */
|
||||
@Column({ name: 'label', type: 'varchar', nullable: true })
|
||||
label?: string | null;
|
||||
|
||||
@Column({ name: 'metadata', type: 'jsonb', nullable: true })
|
||||
metadata?: Record<string, unknown> | null;
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
import { Global, Module } from '@nestjs/common';
|
||||
import { TypeOrmModule } from '@nestjs/typeorm';
|
||||
import { FleetEvent } from './entities/fleet-event.entity';
|
||||
import { FleetHistoryService } from './fleet-history.service';
|
||||
|
||||
/**
|
||||
* Global so any fleet-touching service (vehicles, drivers, first/last-mile) can
|
||||
* inject FleetHistoryService to append audit events without each module having
|
||||
* to import this one.
|
||||
*/
|
||||
@Global()
|
||||
@Module({
|
||||
imports: [TypeOrmModule.forFeature([FleetEvent])],
|
||||
providers: [FleetHistoryService],
|
||||
exports: [FleetHistoryService],
|
||||
})
|
||||
export class FleetHistoryModule {}
|
||||
@@ -0,0 +1,54 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
import { InjectRepository } from '@nestjs/typeorm';
|
||||
import { Repository } from 'typeorm';
|
||||
import { FleetEvent, FleetEventType } from './entities/fleet-event.entity';
|
||||
|
||||
export interface FleetEventInput {
|
||||
eventType: FleetEventType;
|
||||
vehicleId?: string | null;
|
||||
driverId?: string | null;
|
||||
firstMileId?: string | null;
|
||||
lastMileId?: string | null;
|
||||
fromValue?: string | null;
|
||||
toValue?: string | null;
|
||||
label?: string | null;
|
||||
metadata?: Record<string, unknown> | null;
|
||||
}
|
||||
|
||||
@Injectable()
|
||||
export class FleetHistoryService {
|
||||
private readonly logger = new Logger(FleetHistoryService.name);
|
||||
|
||||
constructor(
|
||||
@InjectRepository(FleetEvent)
|
||||
private readonly eventRepo: Repository<FleetEvent>,
|
||||
) {}
|
||||
|
||||
/**
|
||||
* Append an audit event. Best-effort: recording history must never break the
|
||||
* business operation that triggered it, so failures are logged and swallowed.
|
||||
*/
|
||||
async record(input: FleetEventInput): Promise<void> {
|
||||
try {
|
||||
await this.eventRepo.save(this.eventRepo.create(input));
|
||||
} catch (err) {
|
||||
this.logger.error(
|
||||
`Failed to record fleet event ${input.eventType}: ${String(err)}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
getVehicleHistory(vehicleId: string): Promise<FleetEvent[]> {
|
||||
return this.eventRepo.find({
|
||||
where: { vehicleId },
|
||||
order: { createdAt: 'DESC' },
|
||||
});
|
||||
}
|
||||
|
||||
getDriverHistory(driverId: string): Promise<FleetEvent[]> {
|
||||
return this.eventRepo.find({
|
||||
where: { driverId },
|
||||
order: { createdAt: 'DESC' },
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
import { BaseEntity } from '@edr/api-common';
|
||||
import { Column, Entity, Index, JoinColumn, ManyToOne, Unique } from 'typeorm';
|
||||
|
||||
import { Vehicle } from '../../vehicles/entities/vehicle.entity';
|
||||
import { LastMile } from './last-mile.entity';
|
||||
|
||||
/**
|
||||
* One row per vehicle assigned to a last-mile delivery. A delivery can be
|
||||
* served by several vehicles at once (multi-truck bookings); the legacy
|
||||
* `last_mile.vehicle_id` column keeps pointing at the first assignment for
|
||||
* backward compatibility.
|
||||
*/
|
||||
@Entity({ name: 'last_mile_vehicle_assignments', schema: 'freight' })
|
||||
@Unique(['lastMileId', 'vehicleId'])
|
||||
@Index(['vehicleId'])
|
||||
export class LastMileVehicleAssignment extends BaseEntity {
|
||||
@Column({ name: 'last_mile_id', type: 'uuid' })
|
||||
lastMileId!: string;
|
||||
|
||||
@ManyToOne(() => LastMile, (lm) => lm.vehicleAssignments, { nullable: false, onDelete: 'CASCADE' })
|
||||
@JoinColumn({ name: 'last_mile_id' })
|
||||
lastMile?: LastMile;
|
||||
|
||||
@Column({ name: 'vehicle_id', type: 'uuid' })
|
||||
vehicleId!: string;
|
||||
|
||||
@ManyToOne(() => Vehicle, { nullable: false, eager: false })
|
||||
@JoinColumn({ name: 'vehicle_id' })
|
||||
vehicle?: Vehicle;
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import { Column, Entity, Index, JoinColumn, ManyToOne, OneToMany } from 'typeorm
|
||||
import { Booking } from '../../bookings/entities/booking.entity';
|
||||
import { Vehicle } from '../../vehicles/entities/vehicle.entity';
|
||||
import { LastMileContainerAllocation } from './last-mile-container-allocation.entity';
|
||||
import { LastMileVehicleAssignment } from './last-mile-vehicle-assignment.entity';
|
||||
|
||||
export const LAST_MILE_STATUSES = [
|
||||
'PAYMENT_PENDING',
|
||||
@@ -57,4 +58,7 @@ export class LastMile extends BaseEntity {
|
||||
|
||||
@OneToMany(() => LastMileContainerAllocation, (ca) => ca.lastMile)
|
||||
containerAllocations?: LastMileContainerAllocation[];
|
||||
|
||||
@OneToMany(() => LastMileVehicleAssignment, (va) => va.lastMile)
|
||||
vehicleAssignments?: LastMileVehicleAssignment[];
|
||||
}
|
||||
|
||||
@@ -61,7 +61,7 @@ export class LastMileInvoiceService {
|
||||
type: 'DELIVERY_FEE',
|
||||
companyId: lm.booking!.companyId,
|
||||
companyProfileId: lm.booking!.companyProfileId || '',
|
||||
currency: 'ETB',
|
||||
currency: lm.booking!.paymentCurrency || 'ETB',
|
||||
lines: [
|
||||
{
|
||||
chargeType: 'DELIVERY',
|
||||
|
||||
@@ -86,6 +86,7 @@ export class LastMileController {
|
||||
const record = await this.lastMileService.update(id, dto);
|
||||
// Auto-generate invoice if distance or payment was updated
|
||||
const booking = await this.bookingsService.findById(record.bookingId);
|
||||
const currency = booking.paymentCurrency || "ETB";
|
||||
if (dto.exactKm !== undefined || dto.exactKm != record.exactKm || dto.remainingPayment !== undefined || dto.remainingPayment !== record.remainingPayment) {
|
||||
await this.billingService.generateInvoice({
|
||||
source: Freight.InvoiceSource.LastMile,
|
||||
@@ -93,8 +94,8 @@ export class LastMileController {
|
||||
type: "LAST_MILE",
|
||||
companyId: booking.companyId,
|
||||
companyProfileId: booking.companyProfileId,
|
||||
currency: "ETB",
|
||||
|
||||
currency,
|
||||
|
||||
lines: [
|
||||
{
|
||||
chargeType: "LAST_MILE",
|
||||
@@ -102,7 +103,7 @@ export class LastMileController {
|
||||
quantity: 1,
|
||||
unitRate: record.remainingPayment,
|
||||
amount: record.remainingPayment,
|
||||
currency: "ETB",
|
||||
currency,
|
||||
},
|
||||
],
|
||||
|
||||
|
||||
@@ -1,10 +1,11 @@
|
||||
import { Injectable, Logger, NotFoundException } from '@nestjs/common';
|
||||
import { DataSource, FindOptionsWhere } from 'typeorm';
|
||||
import { BadRequestException, Injectable, Logger, NotFoundException } from '@nestjs/common';
|
||||
import { DataSource, FindOptionsWhere, In, IsNull, Not } from 'typeorm';
|
||||
|
||||
import { BookingsRepository } from '../bookings/bookings.repository';
|
||||
import { DriversService } from '../drivers/drivers.service';
|
||||
import { SmsClientService } from '../notifications/sms-client.service';
|
||||
import { VehiclesService } from '../vehicles/vehicles.service';
|
||||
import { VehicleAvailability } from '../vehicles/entities/vehicle.entity';
|
||||
import { CreateLastMileDto } from './dto/create-last-mile.dto';
|
||||
import { UpdateLastMileDto } from './dto/update-last-mile.dto';
|
||||
import { LastMile, LastMileStatus } from './entities/last-mile.entity';
|
||||
@@ -12,6 +13,8 @@ import { LastMileContainerAllocation } from './entities/last-mile-container-allo
|
||||
import { LastMileRepository } from './last-mile.repository';
|
||||
import { InvoiceEventPayload } from '../billing/billing.service';
|
||||
import { OnEvent } from '@nestjs/event-emitter';
|
||||
import { FleetHistoryService } from '../fleet-history/fleet-history.service';
|
||||
import { FleetEventType } from '../fleet-history/entities/fleet-event.entity';
|
||||
|
||||
type LastMileListFilter = {
|
||||
status?: LastMileStatus;
|
||||
@@ -41,8 +44,57 @@ export class LastMileService {
|
||||
private readonly driversService: DriversService,
|
||||
private readonly smsClient: SmsClientService,
|
||||
private readonly dataSource: DataSource,
|
||||
private readonly history: FleetHistoryService,
|
||||
) {}
|
||||
|
||||
/** Resolve a vehicle's driver + human labels, for stamping mile events onto
|
||||
* the driver's timeline and naming the vehicle. Best-effort — never throws. */
|
||||
private async vehicleInfo(
|
||||
vehicleId?: string | null,
|
||||
): Promise<{ driverId: string | null; plate: string | null; driverName: string | null }> {
|
||||
if (!vehicleId) return { driverId: null, plate: null, driverName: null };
|
||||
try {
|
||||
const v = await this.vehiclesService.findById(vehicleId);
|
||||
return {
|
||||
driverId: v.assignedDriverId ?? null,
|
||||
plate: v.plateNumber ?? v.code ?? null,
|
||||
driverName: v.assignedDriverName ?? null,
|
||||
};
|
||||
} catch {
|
||||
return { driverId: null, plate: null, driverName: null };
|
||||
}
|
||||
}
|
||||
|
||||
/** A leg counts as having a vehicle if it has a direct assignment or at least
|
||||
* one container allocation carrying a vehicle. Gates the IN_TRANSIT move. */
|
||||
private async hasAssignedVehicle(
|
||||
recordId: string,
|
||||
directVehicleId?: string | null,
|
||||
): Promise<boolean> {
|
||||
if (directVehicleId) return true;
|
||||
const count = await this.dataSource.manager.count(LastMileContainerAllocation, {
|
||||
where: { lastMileId: recordId, vehicleId: Not(IsNull()) },
|
||||
});
|
||||
return count > 0;
|
||||
}
|
||||
|
||||
/** Human booking reference for a last-mile record, for the history timeline.
|
||||
* Uses the already-loaded relation when present, else looks it up. */
|
||||
private async resolveBookingRef(
|
||||
record: LastMile,
|
||||
): Promise<string | null> {
|
||||
const loaded = (record as LastMile & { booking?: { reference?: string } })
|
||||
.booking?.reference;
|
||||
if (loaded) return loaded;
|
||||
if (!record.bookingId) return null;
|
||||
try {
|
||||
const booking = await this.bookingsRepository.findById(record.bookingId);
|
||||
return (booking as { reference?: string } | null)?.reference ?? null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
async acceptBooking(bookingReference: string): Promise<LastMile | null> {
|
||||
const booking = await this.bookingsRepository.findByReference(bookingReference);
|
||||
|
||||
@@ -132,7 +184,7 @@ export class LastMileService {
|
||||
}
|
||||
|
||||
async create(dto: CreateLastMileDto): Promise<LastMile> {
|
||||
return this.lastMileRepository.create({
|
||||
const record = await this.lastMileRepository.create({
|
||||
bookingId: dto.bookingId,
|
||||
status: dto.status ?? 'READY_TO_TRANSIT',
|
||||
advancedPayment: dto.advancedPayment ?? 0,
|
||||
@@ -142,6 +194,26 @@ export class LastMileService {
|
||||
vehicleId: dto.vehicleId ?? null,
|
||||
paid: (dto as any).paid ?? false,
|
||||
});
|
||||
|
||||
if (dto.vehicleId) {
|
||||
await this.vehiclesService.setAvailability(dto.vehicleId, VehicleAvailability.BUSY);
|
||||
const info = await this.vehicleInfo(dto.vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_VEHICLE_ASSIGNED,
|
||||
vehicleId: dto.vehicleId,
|
||||
lastMileId: record.id,
|
||||
driverId: info.driverId,
|
||||
label: record.status,
|
||||
metadata: {
|
||||
mile: 'LAST',
|
||||
bookingRef: await this.resolveBookingRef(record),
|
||||
vehiclePlate: info.plate,
|
||||
driverName: info.driverName,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
return record;
|
||||
}
|
||||
|
||||
@OnEvent("lastmile.invoice.paid")
|
||||
@@ -159,6 +231,18 @@ export class LastMileService {
|
||||
async update(id: string, dto: UpdateLastMileDto): Promise<LastMile> {
|
||||
const existing = await this.findById(id);
|
||||
|
||||
// A leg can only go IN_TRANSIT once a vehicle is assigned (allowing a vehicle
|
||||
// assigned in this same request).
|
||||
if (dto.status === 'IN_TRANSIT' && existing.status !== 'IN_TRANSIT') {
|
||||
const vehicleId =
|
||||
dto.vehicleId !== undefined ? dto.vehicleId : existing.vehicleId;
|
||||
if (!(await this.hasAssignedVehicle(id, vehicleId))) {
|
||||
throw new BadRequestException(
|
||||
'Assign a vehicle before marking this last-mile leg in transit',
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
const dtoAny = dto as any;
|
||||
const updated = await this.lastMileRepository.update(id, {
|
||||
...(dto.bookingId !== undefined ? { bookingId: dto.bookingId } : {}),
|
||||
@@ -180,9 +264,95 @@ export class LastMileService {
|
||||
void this.notifyDriverAssignment(dto.vehicleId, existing);
|
||||
}
|
||||
|
||||
// Audit the mile↔vehicle (re)assignment on both vehicle and driver lines.
|
||||
const bookingRef = await this.resolveBookingRef(existing);
|
||||
if (dto.vehicleId !== undefined && dto.vehicleId !== existing.vehicleId) {
|
||||
// Keep vehicle availability in sync: new vehicle goes BUSY, replaced one
|
||||
// is freed if no other active trip still holds it.
|
||||
if (dto.vehicleId) {
|
||||
await this.vehiclesService.setAvailability(dto.vehicleId, VehicleAvailability.BUSY);
|
||||
}
|
||||
if (existing.vehicleId) {
|
||||
await this.vehiclesService.releaseIfUnused([existing.vehicleId]);
|
||||
}
|
||||
if (existing.vehicleId) {
|
||||
const info = await this.vehicleInfo(existing.vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_VEHICLE_RELEASED,
|
||||
vehicleId: existing.vehicleId,
|
||||
lastMileId: id,
|
||||
driverId: info.driverId,
|
||||
metadata: {
|
||||
mile: 'LAST',
|
||||
bookingRef,
|
||||
vehiclePlate: info.plate,
|
||||
driverName: info.driverName,
|
||||
},
|
||||
});
|
||||
}
|
||||
if (dto.vehicleId) {
|
||||
const info = await this.vehicleInfo(dto.vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_VEHICLE_ASSIGNED,
|
||||
vehicleId: dto.vehicleId,
|
||||
lastMileId: id,
|
||||
driverId: info.driverId,
|
||||
label: updated.status,
|
||||
metadata: {
|
||||
mile: 'LAST',
|
||||
bookingRef,
|
||||
vehiclePlate: info.plate,
|
||||
driverName: info.driverName,
|
||||
},
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
if (dto.status !== undefined && dto.status !== existing.status) {
|
||||
const vehicleId = updated.vehicleId ?? existing.vehicleId ?? null;
|
||||
const info = await this.vehicleInfo(vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_STATUS_CHANGED,
|
||||
lastMileId: id,
|
||||
vehicleId,
|
||||
driverId: info.driverId,
|
||||
fromValue: existing.status,
|
||||
toValue: dto.status,
|
||||
metadata: {
|
||||
mile: 'LAST',
|
||||
bookingRef,
|
||||
vehiclePlate: info.plate,
|
||||
driverName: info.driverName,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
// Delivery finished — free the vehicles this trip was holding.
|
||||
if (dto.status === 'DELIVERED' && existing.status !== 'DELIVERED') {
|
||||
await this.releaseVehicles(updated);
|
||||
}
|
||||
|
||||
return updated;
|
||||
}
|
||||
|
||||
/**
|
||||
* Free every vehicle held by this record (direct assignment + container
|
||||
* allocations), unless still in use by another active trip.
|
||||
*/
|
||||
private async releaseVehicles(record: LastMile): Promise<void> {
|
||||
const recordAllocations = await this.dataSource.manager.find(
|
||||
LastMileContainerAllocation,
|
||||
{ where: { lastMileId: record.id } },
|
||||
);
|
||||
const vehicleIds = recordAllocations
|
||||
.map((a) => a.vehicleId)
|
||||
.filter((id): id is string => Boolean(id));
|
||||
if (record.vehicleId) {
|
||||
vehicleIds.push(record.vehicleId);
|
||||
}
|
||||
await this.vehiclesService.releaseIfUnused(vehicleIds);
|
||||
}
|
||||
|
||||
private async notifyDriverAssignment(vehicleId: string, record: LastMile): Promise<void> {
|
||||
try {
|
||||
const vehicle = await this.vehiclesService.findById(vehicleId);
|
||||
@@ -236,6 +406,21 @@ export class LastMileService {
|
||||
throw new NotFoundException(`Last-mile record ${lastMileId} not found`);
|
||||
}
|
||||
|
||||
// Capture the vehicles currently on these containers so a reallocation can
|
||||
// be diffed into assigned/released history events below.
|
||||
const previousAllocations = await this.dataSource.manager.find(
|
||||
LastMileContainerAllocation,
|
||||
{
|
||||
where: {
|
||||
lastMileId,
|
||||
containerId: In(allocations.map((a) => a.containerId)),
|
||||
},
|
||||
},
|
||||
);
|
||||
const previousVehicleIds = previousAllocations
|
||||
.map((a) => a.vehicleId)
|
||||
.filter((id): id is string => Boolean(id));
|
||||
|
||||
await this.dataSource.transaction(async (manager) => {
|
||||
for (const allocation of allocations) {
|
||||
await manager.delete(LastMileContainerAllocation, {
|
||||
@@ -252,6 +437,46 @@ export class LastMileService {
|
||||
}
|
||||
});
|
||||
|
||||
// Keep vehicle availability in sync: newly-allocated cars go BUSY, cars no
|
||||
// longer on any of these containers are freed if unused elsewhere.
|
||||
const vehicleIds = new Set(allocations.map((a) => a.vehicleId));
|
||||
await Promise.all(
|
||||
[...vehicleIds].map((id) =>
|
||||
this.vehiclesService.setAvailability(id, VehicleAvailability.BUSY),
|
||||
),
|
||||
);
|
||||
await this.vehiclesService.releaseIfUnused(
|
||||
previousVehicleIds.filter((id) => !vehicleIds.has(id)),
|
||||
);
|
||||
|
||||
// History: one event per vehicle actually added or removed by this
|
||||
// multi-car (re)allocation, so reassignments show on every timeline.
|
||||
const prevSet = new Set(previousVehicleIds);
|
||||
const bookingRef = await this.resolveBookingRef(lastMile);
|
||||
for (const vehicleId of vehicleIds) {
|
||||
if (prevSet.has(vehicleId)) continue;
|
||||
const info = await this.vehicleInfo(vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_VEHICLE_ASSIGNED,
|
||||
vehicleId,
|
||||
lastMileId,
|
||||
driverId: info.driverId,
|
||||
label: lastMile.status,
|
||||
metadata: { mile: 'LAST', bookingRef, vehiclePlate: info.plate, driverName: info.driverName },
|
||||
});
|
||||
}
|
||||
for (const vehicleId of previousVehicleIds) {
|
||||
if (vehicleIds.has(vehicleId)) continue;
|
||||
const info = await this.vehicleInfo(vehicleId);
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.MILE_VEHICLE_RELEASED,
|
||||
vehicleId,
|
||||
lastMileId,
|
||||
driverId: info.driverId,
|
||||
metadata: { mile: 'LAST', bookingRef, vehiclePlate: info.plate, driverName: info.driverName },
|
||||
});
|
||||
}
|
||||
|
||||
return {
|
||||
success: true,
|
||||
allocated: allocations.length,
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
import { ApiProperty, ApiPropertyOptional } from "@nestjs/swagger";
|
||||
import { IsEmail, IsNotEmpty, IsOptional, IsString } from "class-validator";
|
||||
|
||||
export class SendEmailDto {
|
||||
@ApiProperty({
|
||||
description: "Recipient email address",
|
||||
example: "customer@example.com",
|
||||
})
|
||||
@IsEmail()
|
||||
@IsNotEmpty()
|
||||
to!: string;
|
||||
|
||||
@ApiProperty({
|
||||
description: "Email subject",
|
||||
example: "Your EDR Freight verification code",
|
||||
})
|
||||
@IsString()
|
||||
@IsNotEmpty()
|
||||
subject!: string;
|
||||
|
||||
@ApiPropertyOptional()
|
||||
@IsOptional()
|
||||
@IsString()
|
||||
text?: string;
|
||||
|
||||
@ApiPropertyOptional()
|
||||
@IsOptional()
|
||||
@IsString()
|
||||
html?: string;
|
||||
}
|
||||
@@ -0,0 +1,51 @@
|
||||
import {
|
||||
Inject,
|
||||
Injectable,
|
||||
Logger,
|
||||
OnApplicationBootstrap,
|
||||
} from "@nestjs/common";
|
||||
import { ClientProxy } from "@nestjs/microservices";
|
||||
import { SendEmailDto } from "./dtos/email.dto";
|
||||
|
||||
@Injectable()
|
||||
export class EmailClientService implements OnApplicationBootstrap {
|
||||
private readonly logger = new Logger(EmailClientService.name);
|
||||
|
||||
constructor(
|
||||
@Inject("EMAIL_SERVICE")
|
||||
private readonly emailClient: ClientProxy,
|
||||
) {}
|
||||
|
||||
private readonly enabled = process.env.RABBITMQ_ENABLED !== "false";
|
||||
|
||||
async onApplicationBootstrap() {
|
||||
if (!this.enabled) return;
|
||||
this.emailClient
|
||||
.connect()
|
||||
.then(() => this.logger.log("connected to Email service"))
|
||||
.catch((err) => {
|
||||
console.error("Error happened at Email service", err);
|
||||
});
|
||||
}
|
||||
|
||||
async sendEmail(dto: SendEmailDto): Promise<{ queued: boolean }> {
|
||||
if (!this.enabled) {
|
||||
this.logger.warn(`RABBITMQ disabled — skipped EMAIL to=${dto.to}`);
|
||||
return { queued: false };
|
||||
}
|
||||
this.emailClient.emit("send-email", {
|
||||
to: dto.to,
|
||||
subject: dto.subject,
|
||||
text: dto.text,
|
||||
html: dto.html,
|
||||
appKey: "IFHCRS-LICENSE-MANAGEMENT",
|
||||
});
|
||||
// Fire-and-forget enqueue: confirms hand-off to RabbitMQ, NOT delivery.
|
||||
this.logger.log(
|
||||
`EMAIL queued to RabbitMQ [${process.env.EMAIL_QUEUE ?? "email_queue"}] pattern='send-email'`,
|
||||
);
|
||||
// Recipient + content are PII — debug only.
|
||||
this.logger.debug(`EMAIL payload to=${dto.to} subject="${dto.subject}"`);
|
||||
return { queued: true };
|
||||
}
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import { ClientsModule, Transport } from "@nestjs/microservices";
|
||||
|
||||
import { NotificationsService } from "./notifications.service";
|
||||
import { SmsClientService } from "./sms-client.service";
|
||||
import { EmailClientService } from "./email-client.service";
|
||||
import { EmailNotificationStrategy } from "./strategies/notification.email.strategy";
|
||||
import { SmsNotificationStrategy } from "./strategies/notification.sms.strategy";
|
||||
|
||||
@@ -20,10 +21,25 @@ import { SmsNotificationStrategy } from "./strategies/notification.sms.strategy"
|
||||
queueOptions: { durable: true },
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "EMAIL_SERVICE",
|
||||
transport: Transport.RMQ,
|
||||
options: {
|
||||
urls: [process.env.RABBITMQ_URL as string],
|
||||
queue: process.env.EMAIL_QUEUE ?? "email_queue",
|
||||
queueOptions: { durable: true },
|
||||
},
|
||||
},
|
||||
]),
|
||||
],
|
||||
controllers: [],
|
||||
providers: [EmailNotificationStrategy, SmsNotificationStrategy, NotificationsService, SmsClientService],
|
||||
exports: [NotificationsService, SmsClientService],
|
||||
providers: [
|
||||
EmailNotificationStrategy,
|
||||
SmsNotificationStrategy,
|
||||
NotificationsService,
|
||||
SmsClientService,
|
||||
EmailClientService,
|
||||
],
|
||||
exports: [NotificationsService, SmsClientService, EmailClientService],
|
||||
})
|
||||
export class NotificationsModule {}
|
||||
|
||||
@@ -1,15 +1,24 @@
|
||||
// otp.controller.ts
|
||||
|
||||
import {
|
||||
BadRequestException,
|
||||
Body,
|
||||
Controller,
|
||||
Post,
|
||||
} from "@nestjs/common";
|
||||
|
||||
|
||||
import { OtpService } from "./otp.service";
|
||||
import { OtpService, OtpTarget } from "./otp.service";
|
||||
import { Public } from "@edr/api-common";
|
||||
|
||||
// Exactly one of phone/email must be present per request — the channel the
|
||||
// code is sent through / checked against.
|
||||
function toTarget(phone?: string, email?: string): OtpTarget {
|
||||
if (email) return { email };
|
||||
if (phone) return { phone };
|
||||
throw new BadRequestException("phone or email is required");
|
||||
}
|
||||
|
||||
@Controller("otp")
|
||||
@Public()
|
||||
export class OtpController {
|
||||
@@ -24,9 +33,12 @@ export class OtpController {
|
||||
@Post("send")
|
||||
async sendOtp(
|
||||
@Body("phone")
|
||||
phone: string
|
||||
phone?: string,
|
||||
|
||||
@Body("email")
|
||||
email?: string
|
||||
) {
|
||||
return this.otpService.sendOtp(phone);
|
||||
return this.otpService.sendOtp(toTarget(phone, email));
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -36,13 +48,16 @@ export class OtpController {
|
||||
@Post("verify")
|
||||
async verifyOtp(
|
||||
@Body("phone")
|
||||
phone: string,
|
||||
phone: string | undefined,
|
||||
|
||||
@Body("email")
|
||||
email: string | undefined,
|
||||
|
||||
@Body("otp")
|
||||
otp: string
|
||||
) {
|
||||
return this.otpService.verifyOtp(
|
||||
phone,
|
||||
toTarget(phone, email),
|
||||
otp
|
||||
);
|
||||
}
|
||||
|
||||
@@ -10,10 +10,19 @@ import { BaseEntity } from "@edr/api-common";
|
||||
name: "otp_verifications",
|
||||
})
|
||||
export class OtpVerification extends BaseEntity{
|
||||
// Exactly one of phone/email is set per row — the channel the code was sent
|
||||
// through.
|
||||
@Column({
|
||||
unique: true,
|
||||
nullable: true,
|
||||
})
|
||||
phone!: string;
|
||||
phone?: string;
|
||||
|
||||
@Column({
|
||||
unique: true,
|
||||
nullable: true,
|
||||
})
|
||||
email?: string;
|
||||
|
||||
@Column()
|
||||
otp!: string;
|
||||
|
||||
@@ -31,6 +31,7 @@ import { NotificationsModule } from "../notifications/notifications.module";
|
||||
|
||||
exports: [
|
||||
OtpRepository,
|
||||
OtpService,
|
||||
],
|
||||
})
|
||||
export class OtpModule {}
|
||||
@@ -31,17 +31,44 @@ export class OtpRepository {
|
||||
});
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Find By Email
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
async findByEmail(
|
||||
email: string
|
||||
) {
|
||||
return this.repository.findOne({
|
||||
where: {
|
||||
email,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Find By Target (either channel)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
async findByTarget(
|
||||
target: { phone?: string; email?: string }
|
||||
) {
|
||||
return target.email
|
||||
? this.findByEmail(target.email)
|
||||
: this.findByPhone(target.phone!);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Create OTP
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
async createOtp(
|
||||
phone: string,
|
||||
target: { phone?: string; email?: string },
|
||||
otp: string
|
||||
) {
|
||||
const entity =
|
||||
this.repository.create({
|
||||
phone,
|
||||
phone: target.phone,
|
||||
email: target.email,
|
||||
otp,
|
||||
verified: false,
|
||||
});
|
||||
@@ -70,10 +97,10 @@ export class OtpRepository {
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Verify Phone
|
||||
// Mark Verified
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
async verifyPhone(
|
||||
async markVerified(
|
||||
otpVerification: OtpVerification
|
||||
) {
|
||||
otpVerification.verified =
|
||||
@@ -83,4 +110,18 @@ export class OtpRepository {
|
||||
otpVerification
|
||||
);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Delete OTP (single-use consume)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
// Hard delete so the unique `phone` row is freed and a fresh code can be
|
||||
// requested for the same number on the next action.
|
||||
async deleteOtp(
|
||||
otpVerification: OtpVerification
|
||||
) {
|
||||
return this.repository.remove(
|
||||
otpVerification
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -1,80 +1,80 @@
|
||||
// otp.service.ts
|
||||
|
||||
import {
|
||||
BadRequestException,
|
||||
Injectable,
|
||||
} from "@nestjs/common";
|
||||
import { BadRequestException, Injectable, Logger } from "@nestjs/common";
|
||||
|
||||
import { OtpRepository } from "./otp.repository";
|
||||
|
||||
import { SmsClientService } from "../notifications/sms-client.service";
|
||||
import { EmailClientService } from "../notifications/email-client.service";
|
||||
|
||||
// Exactly one of phone/email is set — enforced by the controller before it
|
||||
// reaches here.
|
||||
export type OtpTarget = { phone?: string; email?: string };
|
||||
|
||||
@Injectable()
|
||||
export class OtpService {
|
||||
logger = new Logger(OtpService.name);
|
||||
constructor(
|
||||
private readonly otpRepository: OtpRepository,
|
||||
private readonly smsClient: SmsClientService
|
||||
) {}
|
||||
private readonly smsClient: SmsClientService,
|
||||
private readonly emailClient: EmailClientService,
|
||||
) { }
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Generate OTP
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
generateOtp(): string {
|
||||
return Math.floor(
|
||||
100000 + Math.random() * 900000
|
||||
).toString();
|
||||
return Math.floor(100000 + Math.random() * 900000).toString();
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Send OTP
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
async sendOtp(phone: string) {
|
||||
async sendOtp(target: OtpTarget) {
|
||||
try {
|
||||
// The verification code is generated server-side — never supplied by the
|
||||
// caller — so the OTP stays a secret known only to the server and the
|
||||
// recipient of the SMS.
|
||||
// recipient of the SMS/email.
|
||||
const otp = this.generateOtp();
|
||||
|
||||
// find existing phone
|
||||
const existingPhone =
|
||||
await this.otpRepository.findByPhone(
|
||||
phone
|
||||
);
|
||||
// find existing row for this channel
|
||||
const existing = await this.otpRepository.findByTarget(target);
|
||||
|
||||
// update existing otp
|
||||
if (existingPhone) {
|
||||
await this.otpRepository.updateOtp(
|
||||
existingPhone,
|
||||
otp
|
||||
);
|
||||
if (existing) {
|
||||
await this.otpRepository.updateOtp(existing, otp);
|
||||
} else {
|
||||
// create new otp
|
||||
await this.otpRepository.createOtp(
|
||||
phone,
|
||||
otp
|
||||
);
|
||||
await this.otpRepository.createOtp(target, otp);
|
||||
}
|
||||
|
||||
// send sms (queued to RabbitMQ via the shared SMS service)
|
||||
await this.smsClient.sendSms({
|
||||
to: phone,
|
||||
message: `Your verification code is ${otp}`,
|
||||
});
|
||||
if (target.email) {
|
||||
// send email (queued to RabbitMQ via the shared Email service)
|
||||
await this.emailClient.sendEmail({
|
||||
to: target.email,
|
||||
subject: "Your EDR Freight verification code",
|
||||
text: `Your verification code is ${otp}`,
|
||||
});
|
||||
} else {
|
||||
// send sms (queued to RabbitMQ via the shared SMS service)
|
||||
await this.smsClient.sendSms({
|
||||
to: target.phone as string,
|
||||
message: `Your verification code is ${otp}`,
|
||||
});
|
||||
}
|
||||
|
||||
this.logger.log(`OTP send for ${target.email ?? target.phone}: ${otp}`);
|
||||
return {
|
||||
success: true,
|
||||
|
||||
message:
|
||||
"OTP sent successfully",
|
||||
message: "OTP sent successfully",
|
||||
};
|
||||
} catch (error) {
|
||||
console.log(error);
|
||||
|
||||
throw new BadRequestException(
|
||||
"Failed to send OTP"
|
||||
);
|
||||
throw new BadRequestException("Failed to send OTP");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -82,40 +82,70 @@ export class OtpService {
|
||||
// Verify OTP
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
async verifyOtp(
|
||||
phone: string,
|
||||
otp: string
|
||||
) {
|
||||
// find phone
|
||||
const otpData =
|
||||
await this.otpRepository.findByPhone(
|
||||
phone
|
||||
);
|
||||
async verifyOtp(target: OtpTarget, otp: string) {
|
||||
// find the channel's row
|
||||
const otpData = await this.otpRepository.findByTarget(target);
|
||||
|
||||
// phone not found
|
||||
// not found
|
||||
if (!otpData) {
|
||||
throw new BadRequestException(
|
||||
"Phone number not found"
|
||||
target.email ? "Email address not found" : "Phone number not found",
|
||||
);
|
||||
}
|
||||
|
||||
// invalid otp
|
||||
if (otpData.otp !== otp) {
|
||||
throw new BadRequestException(
|
||||
"Invalid OTP"
|
||||
);
|
||||
throw new BadRequestException("Invalid OTP");
|
||||
}
|
||||
|
||||
// verify phone
|
||||
await this.otpRepository.verifyPhone(
|
||||
otpData
|
||||
);
|
||||
// mark verified
|
||||
await this.otpRepository.markVerified(otpData);
|
||||
|
||||
return {
|
||||
success: true,
|
||||
|
||||
message:
|
||||
"Phone verified successfully",
|
||||
message: target.email
|
||||
? "Email verified successfully"
|
||||
: "Phone verified successfully",
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Verify OTP for a sensitive action (sudo mode)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
// Fresh, single-use challenge gating a sensitive action (e.g. applying a
|
||||
// contract signature). Unlike verifyOtp above — which marks a phone verified
|
||||
// and leaves the code in place — this enforces a short TTL and consumes the
|
||||
// code on success so it can never be replayed.
|
||||
private readonly ACTION_OTP_TTL_MS = 5 * 60 * 1000;
|
||||
|
||||
async verifyOtpForAction(phone: string, otp: string) {
|
||||
const otpData = await this.otpRepository.findByPhone(phone);
|
||||
|
||||
if (!otpData) {
|
||||
throw new BadRequestException(
|
||||
"No verification code was requested for this phone",
|
||||
);
|
||||
}
|
||||
|
||||
const ageMs = Date.now() - new Date(otpData.updatedAt).getTime();
|
||||
|
||||
if (ageMs > this.ACTION_OTP_TTL_MS) {
|
||||
await this.otpRepository.deleteOtp(otpData);
|
||||
|
||||
throw new BadRequestException(
|
||||
"Verification code has expired. Request a new one.",
|
||||
);
|
||||
}
|
||||
|
||||
if (otpData.otp !== otp) {
|
||||
throw new BadRequestException("Invalid verification code");
|
||||
}
|
||||
|
||||
// single-use: consume on success
|
||||
await this.otpRepository.deleteOtp(otpData);
|
||||
|
||||
return { success: true };
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,9 +1,10 @@
|
||||
import {
|
||||
Body,
|
||||
Controller,
|
||||
HttpCode,
|
||||
HttpStatus,
|
||||
Post,
|
||||
Body,
|
||||
Controller,
|
||||
HttpCode,
|
||||
HttpStatus,
|
||||
Logger,
|
||||
Post,
|
||||
} from "@nestjs/common";
|
||||
import { ApiOperation, ApiTags } from "@nestjs/swagger";
|
||||
import { Public } from "@edr/api-common";
|
||||
@@ -22,14 +23,17 @@ import { PaymentService } from "./payment.service";
|
||||
@Public()
|
||||
@Controller("internal/payments")
|
||||
export class InternalPaymentController {
|
||||
constructor(private readonly paymentService: PaymentService) { }
|
||||
private readonly logger = new Logger(InternalPaymentController.name);
|
||||
constructor(private readonly paymentService: PaymentService) { }
|
||||
|
||||
@Post("mark-paid")
|
||||
@HttpCode(HttpStatus.OK)
|
||||
@ApiOperation({
|
||||
summary: "Apply a payment.succeeded / payment.failed event from the payment service (idempotent)",
|
||||
})
|
||||
async markPaid(@Body() event: PaymentEventDto): Promise<MarkPaidResponseDto> {
|
||||
return this.paymentService.handlePaymentEvent(event);
|
||||
}
|
||||
@Post("mark-paid")
|
||||
@HttpCode(HttpStatus.OK)
|
||||
@ApiOperation({
|
||||
summary:
|
||||
"Apply a payment.succeeded / payment.failed event from the payment service (idempotent)",
|
||||
})
|
||||
async markPaid(@Body() event: PaymentEventDto): Promise<MarkPaidResponseDto> {
|
||||
this.logger.log(`Marking payment ${event} as PAID`);
|
||||
return this.paymentService.handlePaymentEvent(event);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,9 +1,15 @@
|
||||
import { BaseEntity } from '@edr/api-common';
|
||||
import { LoadingStatus } from '@edr/types';
|
||||
import { Entity, Index, JoinColumn, ManyToOne, Column } from 'typeorm';
|
||||
|
||||
import { Booking } from '../../bookings/entities/booking.entity';
|
||||
import { TrainSchedule } from './train-schedule.entity';
|
||||
|
||||
export const TRAIN_SCHEDULE_BOOKING_LOADING_STATUSES = [
|
||||
LoadingStatus.Unloaded,
|
||||
LoadingStatus.Loaded,
|
||||
] as const;
|
||||
|
||||
@Entity({ schema: 'freight', name: 'train_schedule_bookings' })
|
||||
@Index(['trainScheduleId', 'bookingId'], { unique: true })
|
||||
@Index(['bookingId'], { unique: true })
|
||||
@@ -23,4 +29,7 @@ export class TrainScheduleBooking extends BaseEntity {
|
||||
@ManyToOne(() => Booking)
|
||||
@JoinColumn({ name: 'booking_id' })
|
||||
booking?: Booking;
|
||||
|
||||
@Column({ name: 'loading_status', type: 'varchar', length: 20, default: 'UNLOADED' })
|
||||
loadingStatus!: string;
|
||||
}
|
||||
|
||||
@@ -47,4 +47,27 @@ export class TrainScheduleBookingsRepository extends BaseRepository<TrainSchedul
|
||||
select: { id: true, bookingId: true, trainScheduleId: true },
|
||||
});
|
||||
}
|
||||
|
||||
findByScheduleId(
|
||||
trainScheduleId: string,
|
||||
manager?: EntityManager,
|
||||
): Promise<TrainScheduleBooking[]> {
|
||||
return this.repo(manager).find({
|
||||
where: { trainScheduleId },
|
||||
select: { id: true, bookingId: true, trainScheduleId: true, loadingStatus: true },
|
||||
});
|
||||
}
|
||||
|
||||
async updateLoadingStatusMany(
|
||||
trainScheduleId: string,
|
||||
bookingIds: string[],
|
||||
loadingStatus: string,
|
||||
manager?: EntityManager,
|
||||
): Promise<void> {
|
||||
if (!bookingIds.length) return;
|
||||
await this.repo(manager).update(
|
||||
{ trainScheduleId, bookingId: In(bookingIds) },
|
||||
{ loadingStatus },
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,15 @@
|
||||
import { LoadingStatus } from '@edr/types';
|
||||
import { ApiProperty } from '@nestjs/swagger';
|
||||
import { ArrayMinSize, IsArray, IsEnum, IsUUID } from 'class-validator';
|
||||
|
||||
export class UpdateImportLoadingStatusDto {
|
||||
@ApiProperty({ type: [String] })
|
||||
@IsArray()
|
||||
@ArrayMinSize(1)
|
||||
@IsUUID('4', { each: true })
|
||||
bookingIds!: string[];
|
||||
|
||||
@ApiProperty({ enum: LoadingStatus })
|
||||
@IsEnum(LoadingStatus)
|
||||
loadingStatus!: LoadingStatus;
|
||||
}
|
||||
@@ -28,6 +28,7 @@ import { GetEligibleBulkBookingsDto } from "./dto/get-eligible-bulk-bookings.dto
|
||||
import { GetEligibleContainerBookingsDto } from "./dto/get-eligible-container-bookings.dto";
|
||||
import { PinWagonsDto } from "./dto/pin-wagons.dto";
|
||||
import { UpdateContainerItemDto } from "./dto/update-container-item.dto";
|
||||
import { UpdateImportLoadingStatusDto } from "./dto/update-import-loading-status.dto";
|
||||
import { PreviewBulkTrainScheduleDto } from "./dto/preview-bulk-train-schedule.dto";
|
||||
import { PreviewContainerTrainScheduleDto } from "./dto/preview-container-train-schedule.dto";
|
||||
import { PreviewTrainScheduleDto } from "./dto/preview-train-schedule.dto";
|
||||
@@ -336,6 +337,28 @@ export class TrainSchedulingController {
|
||||
return this.trainSchedulingService.getCompositionRemovals(id);
|
||||
}
|
||||
|
||||
@Get("schedules/:id/import-loading-bookings")
|
||||
@TrainSchedulingView()
|
||||
@ApiOperation({
|
||||
summary: "List import bookings eligible for loading confirmation on this schedule",
|
||||
})
|
||||
getImportLoadingBookings(@Param("id", ParseUUIDPipe) id: string) {
|
||||
return this.trainSchedulingService.getImportLoadingBookings(id);
|
||||
}
|
||||
|
||||
@Patch("schedules/:id/import-loading-status")
|
||||
@TrainSchedulingManage()
|
||||
@ApiOperation({
|
||||
summary:
|
||||
"Mark import bookings loaded/unloaded on this schedule (tracking only, does not affect dispatch)",
|
||||
})
|
||||
updateImportLoadingStatus(
|
||||
@Param("id", ParseUUIDPipe) id: string,
|
||||
@Body() dto: UpdateImportLoadingStatusDto,
|
||||
) {
|
||||
return this.trainSchedulingService.updateImportLoadingStatus(id, dto);
|
||||
}
|
||||
|
||||
@Post("schedules/:id/pin-wagons")
|
||||
@TrainSchedulingManage()
|
||||
@ApiOperation({ summary: "Pin physical wagons to train set slots" })
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import {
|
||||
AllocationLoadType,
|
||||
LoadingStatus,
|
||||
SchedulingStatus,
|
||||
TrainCheckpointKind,
|
||||
TrainScheduleStatus as TrainScheduleStatusEnum,
|
||||
@@ -46,6 +47,7 @@ import { GetEligibleBulkBookingsDto } from './dto/get-eligible-bulk-bookings.dto
|
||||
import { GetEligibleContainerBookingsDto } from './dto/get-eligible-container-bookings.dto';
|
||||
import { PinWagonsDto } from './dto/pin-wagons.dto';
|
||||
import { UpdateContainerItemDto } from './dto/update-container-item.dto';
|
||||
import { UpdateImportLoadingStatusDto } from './dto/update-import-loading-status.dto';
|
||||
import { PreviewBulkTrainScheduleDto } from './dto/preview-bulk-train-schedule.dto';
|
||||
import { PreviewContainerTrainScheduleDto } from './dto/preview-container-train-schedule.dto';
|
||||
import { PreviewTrainScheduleDto } from './dto/preview-train-schedule.dto';
|
||||
@@ -866,6 +868,87 @@ export class TrainSchedulingService {
|
||||
);
|
||||
}
|
||||
|
||||
async getImportLoadingBookings(scheduleId: string) {
|
||||
const schedule = await this.trainSchedulesRepository.findById(scheduleId);
|
||||
if (!schedule) {
|
||||
throw new NotFoundException(`Train schedule ${scheduleId} not found`);
|
||||
}
|
||||
|
||||
const [scheduleBookings, allocations] = await Promise.all([
|
||||
this.trainScheduleBookingsRepository.findByScheduleId(scheduleId),
|
||||
this.wagonBookingAllocationsRepository.findByScheduleId(scheduleId),
|
||||
]);
|
||||
if (!scheduleBookings.length) {
|
||||
return { count: 0, items: [] };
|
||||
}
|
||||
|
||||
const allocatedBookingIds = new Set(allocations.map((a) => a.bookingId));
|
||||
const statusByBookingId = new Map(
|
||||
scheduleBookings.map((sb) => [sb.bookingId, sb.loadingStatus]),
|
||||
);
|
||||
const candidateIds = scheduleBookings
|
||||
.map((sb) => sb.bookingId)
|
||||
.filter((id) => allocatedBookingIds.has(id));
|
||||
if (!candidateIds.length) {
|
||||
return { count: 0, items: [] };
|
||||
}
|
||||
|
||||
const bookings = await this.bookingsRepository.findByIdsForScheduling(candidateIds);
|
||||
const items = bookings
|
||||
.filter((b) => b.tradeDirection === 'IMPORT' && b.paymentStatus === 'PAID')
|
||||
.map((b) => ({
|
||||
id: b.id,
|
||||
reference: b.reference ?? null,
|
||||
customer: b.company?.name ?? null,
|
||||
weightTons: b.cargoTotalWeightVgm,
|
||||
loadingStatus: statusByBookingId.get(b.id) ?? LoadingStatus.Unloaded,
|
||||
}));
|
||||
return { count: items.length, items };
|
||||
}
|
||||
|
||||
async updateImportLoadingStatus(scheduleId: string, dto: UpdateImportLoadingStatusDto) {
|
||||
const schedule = await this.trainSchedulesRepository.findById(scheduleId);
|
||||
if (!schedule) {
|
||||
throw new NotFoundException(`Train schedule ${scheduleId} not found`);
|
||||
}
|
||||
|
||||
const [scheduleBookings, allocations, bookings] = await Promise.all([
|
||||
this.trainScheduleBookingsRepository.findByScheduleId(scheduleId),
|
||||
this.wagonBookingAllocationsRepository.findByScheduleId(scheduleId),
|
||||
this.bookingsRepository.findByIdsForScheduling(dto.bookingIds),
|
||||
]);
|
||||
|
||||
const scheduledIds = new Set(scheduleBookings.map((sb) => sb.bookingId));
|
||||
const allocatedIds = new Set(allocations.map((a) => a.bookingId));
|
||||
const bookingById = new Map(bookings.map((b) => [b.id, b]));
|
||||
|
||||
const invalid: string[] = [];
|
||||
for (const id of dto.bookingIds) {
|
||||
const booking = bookingById.get(id);
|
||||
if (
|
||||
!scheduledIds.has(id) ||
|
||||
!allocatedIds.has(id) ||
|
||||
!booking ||
|
||||
booking.tradeDirection !== 'IMPORT' ||
|
||||
booking.paymentStatus !== 'PAID'
|
||||
) {
|
||||
invalid.push(id);
|
||||
}
|
||||
}
|
||||
if (invalid.length) {
|
||||
throw new BadRequestException(
|
||||
`Not eligible for import loading confirmation on this schedule: ${invalid.join(', ')}`,
|
||||
);
|
||||
}
|
||||
|
||||
await this.trainScheduleBookingsRepository.updateLoadingStatusMany(
|
||||
scheduleId,
|
||||
dto.bookingIds,
|
||||
dto.loadingStatus,
|
||||
);
|
||||
return this.getImportLoadingBookings(scheduleId);
|
||||
}
|
||||
|
||||
async pinWagons(scheduleId: string, dto: PinWagonsDto) {
|
||||
const schedule = await this.trainSchedulesRepository.findByIdWithFullGraph(scheduleId);
|
||||
if (!schedule) {
|
||||
|
||||
@@ -14,13 +14,17 @@ import { FleetManage, FleetView } from '../../common/booking-guards';
|
||||
import { VehiclesService } from './vehicles.service';
|
||||
import { CreateVehicleDto } from './dto/create-vehicle.dto';
|
||||
import { UpdateVehicleDto } from './dto/update-vehicle.dto';
|
||||
import { FleetHistoryService } from '../fleet-history/fleet-history.service';
|
||||
|
||||
@ApiTags('vehicles')
|
||||
@ApiBearerAuth()
|
||||
@Controller('vehicles')
|
||||
@FleetView()
|
||||
export class VehiclesController {
|
||||
constructor(private readonly vehiclesService: VehiclesService) {}
|
||||
constructor(
|
||||
private readonly vehiclesService: VehiclesService,
|
||||
private readonly fleetHistory: FleetHistoryService,
|
||||
) {}
|
||||
|
||||
@Post()
|
||||
@FleetManage()
|
||||
@@ -57,6 +61,12 @@ export class VehiclesController {
|
||||
return this.vehiclesService.findById(id);
|
||||
}
|
||||
|
||||
@Get(':id/history')
|
||||
@ApiOperation({ summary: 'Get vehicle assignment, status & mile history' })
|
||||
history(@Param('id', ParseUUIDPipe) id: string) {
|
||||
return this.fleetHistory.getVehicleHistory(id);
|
||||
}
|
||||
|
||||
@Patch(':id')
|
||||
@FleetManage()
|
||||
@ApiOperation({ summary: 'Update a vehicle' })
|
||||
|
||||
@@ -9,12 +9,15 @@ import { FirstMileContainerAllocation } from '../first-mile/entities/first-mile-
|
||||
import { LastMile, LastMileStatus } from '../last-mile/entities/last-mile.entity';
|
||||
import { LastMileContainerAllocation } from '../last-mile/entities/last-mile-container-allocation.entity';
|
||||
import { BookingContainerAllocation } from '../bookings/entities/booking-container-allocation.entity';
|
||||
import { FleetHistoryService } from '../fleet-history/fleet-history.service';
|
||||
import { FleetEventType } from '../fleet-history/entities/fleet-event.entity';
|
||||
|
||||
@Injectable()
|
||||
export class VehiclesService {
|
||||
constructor(
|
||||
@InjectRepository(Vehicle)
|
||||
private readonly vehicleRepo: Repository<Vehicle>,
|
||||
private readonly history: FleetHistoryService,
|
||||
) {}
|
||||
|
||||
async create(dto: CreateVehicleDto): Promise<Vehicle> {
|
||||
@@ -34,7 +37,28 @@ export class VehiclesService {
|
||||
registrationNumber,
|
||||
});
|
||||
|
||||
return this.vehicleRepo.save(vehicle);
|
||||
const saved = await this.vehicleRepo.save(vehicle);
|
||||
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.VEHICLE_REGISTERED,
|
||||
vehicleId: saved.id,
|
||||
label: saved.plateNumber ?? saved.code ?? null,
|
||||
toValue: saved.availability ?? null,
|
||||
});
|
||||
if (saved.assignedDriverId) {
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.DRIVER_ASSIGNED,
|
||||
vehicleId: saved.id,
|
||||
driverId: saved.assignedDriverId,
|
||||
label: saved.assignedDriverName ?? null,
|
||||
metadata: {
|
||||
vehiclePlate: saved.plateNumber ?? saved.code ?? null,
|
||||
driverName: saved.assignedDriverName ?? null,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
return saved;
|
||||
}
|
||||
|
||||
async findAll(query: {
|
||||
@@ -97,12 +121,76 @@ export class VehiclesService {
|
||||
}
|
||||
}
|
||||
|
||||
const prev = {
|
||||
assignedDriverId: vehicle.assignedDriverId,
|
||||
assignedDriverName: vehicle.assignedDriverName,
|
||||
status: vehicle.status,
|
||||
availability: vehicle.availability,
|
||||
};
|
||||
|
||||
Object.assign(vehicle, dto);
|
||||
return this.vehicleRepo.save(vehicle);
|
||||
const saved = await this.vehicleRepo.save(vehicle);
|
||||
|
||||
// Driver (re)assignment — emit an unassign for the old driver and/or an
|
||||
// assign for the new one so both drivers' timelines and the vehicle's line up.
|
||||
if (
|
||||
dto.assignedDriverId !== undefined &&
|
||||
dto.assignedDriverId !== prev.assignedDriverId
|
||||
) {
|
||||
const vehiclePlate = saved.plateNumber ?? saved.code ?? null;
|
||||
if (prev.assignedDriverId) {
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.DRIVER_UNASSIGNED,
|
||||
vehicleId: id,
|
||||
driverId: prev.assignedDriverId,
|
||||
label: prev.assignedDriverName ?? null,
|
||||
metadata: { vehiclePlate, driverName: prev.assignedDriverName ?? null },
|
||||
});
|
||||
}
|
||||
if (saved.assignedDriverId) {
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.DRIVER_ASSIGNED,
|
||||
vehicleId: id,
|
||||
driverId: saved.assignedDriverId,
|
||||
label: saved.assignedDriverName ?? null,
|
||||
metadata: { vehiclePlate, driverName: saved.assignedDriverName ?? null },
|
||||
});
|
||||
}
|
||||
}
|
||||
if (dto.status !== undefined && dto.status !== prev.status) {
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.VEHICLE_STATUS_CHANGED,
|
||||
vehicleId: id,
|
||||
fromValue: prev.status ?? null,
|
||||
toValue: saved.status ?? null,
|
||||
});
|
||||
}
|
||||
if (dto.availability !== undefined && dto.availability !== prev.availability) {
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.VEHICLE_AVAILABILITY_CHANGED,
|
||||
vehicleId: id,
|
||||
fromValue: prev.availability ?? null,
|
||||
toValue: saved.availability ?? null,
|
||||
});
|
||||
}
|
||||
|
||||
return saved;
|
||||
}
|
||||
|
||||
async setAvailability(id: string, availability: VehicleAvailability): Promise<void> {
|
||||
// Read the current value so the audit event records an accurate from→to and
|
||||
// we skip logging no-op writes (setAvailability is called in release loops).
|
||||
const vehicle = await this.vehicleRepo.findOne({ where: { id } });
|
||||
const previous = vehicle?.availability;
|
||||
await this.vehicleRepo.update(id, { availability });
|
||||
if (previous !== availability) {
|
||||
await this.history.record({
|
||||
eventType: FleetEventType.VEHICLE_AVAILABILITY_CHANGED,
|
||||
vehicleId: id,
|
||||
fromValue: previous ?? null,
|
||||
toValue: availability,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
import { Column, Entity, Index } from 'typeorm';
|
||||
import { BaseEntity } from '@edr/api-common';
|
||||
|
||||
/**
|
||||
* One row per started Fayda verification. Mirrors the passenger-api Prisma
|
||||
* model `FaydaVerificationSession`, but stored in the freight schema via
|
||||
* TypeORM. `state` is the single-use CSRF token that links the eSignet
|
||||
* redirect back to this session.
|
||||
*/
|
||||
@Entity({ name: 'fayda_verification_sessions', schema: 'freight' })
|
||||
@Index(['expiresAt'])
|
||||
@Index(['iamUserId'])
|
||||
export class FaydaVerificationSession extends BaseEntity {
|
||||
@Column({ name: 'state', unique: true })
|
||||
state!: string;
|
||||
|
||||
@Column({ name: 'code_verifier' })
|
||||
codeVerifier!: string;
|
||||
|
||||
/** VERIFY | LOGIN */
|
||||
@Column({ name: 'purpose', default: 'VERIFY' })
|
||||
purpose!: string;
|
||||
|
||||
/** WEB | MOBILE — recorded for audit */
|
||||
@Column({ name: 'platform', default: 'WEB' })
|
||||
platform!: string;
|
||||
|
||||
@Column({ name: 'save_to_account', type: 'boolean', default: false })
|
||||
saveToAccount!: boolean;
|
||||
|
||||
/** PENDING | COMPLETED | FAILED */
|
||||
@Column({ name: 'status', default: 'PENDING' })
|
||||
status!: string;
|
||||
|
||||
@Column({ name: 'error_code', type: 'varchar', nullable: true })
|
||||
errorCode?: string | null;
|
||||
|
||||
@Column({ name: 'error_description', type: 'text', nullable: true })
|
||||
errorDescription?: string | null;
|
||||
|
||||
@Column({ name: 'iam_user_id', type: 'uuid', nullable: true })
|
||||
iamUserId?: string | null;
|
||||
|
||||
@Column({ name: 'expires_at', type: 'timestamptz' })
|
||||
expiresAt!: Date;
|
||||
|
||||
@Column({ name: 'completed_at', type: 'timestamptz', nullable: true })
|
||||
completedAt?: Date | null;
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
import { Controller, Get, Query } from '@nestjs/common';
|
||||
import { ApiOkResponse, ApiOperation, ApiTags } from '@nestjs/swagger';
|
||||
import { IsPublic } from '@tria-plc/api-common/modules/auth/decorators/public.decorator';
|
||||
import { VerifaydaCallbackDto } from './verifayda.dto';
|
||||
|
||||
/**
|
||||
* Plain acknowledgement endpoint for the Fayda redirect_uri when it points at
|
||||
* the API instead of the web app (e.g. MOBILE clients or connectivity checks).
|
||||
* Registered at /callback (excluded from the global /api prefix in main.ts).
|
||||
* It does NOT consume the verification session — the client must still call
|
||||
* GET /api/fayda/verification/complete with the echoed code+state.
|
||||
*/
|
||||
@ApiTags('Fayda Verification')
|
||||
@Controller('callback')
|
||||
export class FaydaCallbackController {
|
||||
@Get()
|
||||
@IsPublic()
|
||||
@ApiOperation({ summary: 'Acknowledge a Fayda redirect (returns OK, echoes code/state)' })
|
||||
@ApiOkResponse({
|
||||
schema: { example: { status: 'ok', code: '...', state: '...' } },
|
||||
})
|
||||
ok(@Query() query: VerifaydaCallbackDto) {
|
||||
return {
|
||||
status: 'ok',
|
||||
...(query.code ? { code: query.code } : {}),
|
||||
...(query.state ? { state: query.state } : {}),
|
||||
...(query.error ? { error: query.error } : {}),
|
||||
...(query.error_description ? { error_description: query.error_description } : {}),
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
// return res.redirect(url.toString());
|
||||
@@ -0,0 +1,30 @@
|
||||
import { CanActivate, ExecutionContext, Injectable } from '@nestjs/common';
|
||||
import { Reflector } from '@nestjs/core';
|
||||
import { InjectDataSource } from '@nestjs/typeorm';
|
||||
import { JwtGuard as IamJwtGuard } from '@tria-plc/api-common/modules/auth/services/jwt.guard';
|
||||
import { DataSource } from 'typeorm';
|
||||
|
||||
/**
|
||||
* Like the IAM JwtGuard, but never rejects the request.
|
||||
*
|
||||
* When a valid IAM bearer token is present, `request.user` is populated with
|
||||
* the package `TCurrentUser`. Missing or invalid tokens continue as guests.
|
||||
*/
|
||||
@Injectable()
|
||||
export class OptionalJwtGuard extends IamJwtGuard implements CanActivate {
|
||||
constructor(
|
||||
reflector: Reflector,
|
||||
@InjectDataSource() dataSource: DataSource,
|
||||
) {
|
||||
super(reflector, dataSource);
|
||||
}
|
||||
|
||||
async canActivate(context: ExecutionContext): Promise<boolean> {
|
||||
try {
|
||||
await super.canActivate(context);
|
||||
} catch {
|
||||
context.switchToHttp().getRequest().user = undefined;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,71 @@
|
||||
import { exportJWK, generateKeyPair, importJWK, jwtVerify, type JWK } from 'jose';
|
||||
import { generateClientAssertion } from './client-assertion.util';
|
||||
|
||||
describe('generateClientAssertion', () => {
|
||||
let privateJwk: JWK;
|
||||
let publicJwk: JWK;
|
||||
|
||||
beforeAll(async () => {
|
||||
const kp = await generateKeyPair('RS256', { extractable: true });
|
||||
privateJwk = await exportJWK(kp.privateKey);
|
||||
publicJwk = await exportJWK(kp.publicKey);
|
||||
});
|
||||
|
||||
it('produces a JWT verifiable with the matching public key', async () => {
|
||||
const jwt = await generateClientAssertion({
|
||||
clientId: 'edr-passenger-test',
|
||||
audience: 'https://esignet.example.com/token',
|
||||
privateJwk,
|
||||
});
|
||||
|
||||
const verifier = await importJWK(publicJwk, 'RS256');
|
||||
const { payload, protectedHeader } = await jwtVerify(jwt, verifier, {
|
||||
issuer: 'edr-passenger-test',
|
||||
subject: 'edr-passenger-test',
|
||||
audience: 'https://esignet.example.com/token',
|
||||
});
|
||||
|
||||
expect(protectedHeader.alg).toBe('RS256');
|
||||
expect(protectedHeader.typ).toBe('JWT');
|
||||
expect(payload.iss).toBe('edr-passenger-test');
|
||||
expect(payload.sub).toBe('edr-passenger-test');
|
||||
expect(payload.aud).toBe('https://esignet.example.com/token');
|
||||
expect(typeof payload.iat).toBe('number');
|
||||
expect(typeof payload.exp).toBe('number');
|
||||
});
|
||||
|
||||
it('defaults exp to 120 seconds after iat', async () => {
|
||||
const jwt = await generateClientAssertion({
|
||||
clientId: 'c',
|
||||
audience: 'https://a/token',
|
||||
privateJwk,
|
||||
});
|
||||
const verifier = await importJWK(publicJwk, 'RS256');
|
||||
const { payload } = await jwtVerify(jwt, verifier);
|
||||
expect(payload.exp! - payload.iat!).toBe(120);
|
||||
});
|
||||
|
||||
it('honors a custom expiresIn', async () => {
|
||||
const jwt = await generateClientAssertion({
|
||||
clientId: 'c',
|
||||
audience: 'https://a/token',
|
||||
privateJwk,
|
||||
expiresIn: '5m',
|
||||
});
|
||||
const verifier = await importJWK(publicJwk, 'RS256');
|
||||
const { payload } = await jwtVerify(jwt, verifier);
|
||||
expect(payload.exp! - payload.iat!).toBe(300);
|
||||
});
|
||||
|
||||
it('fails verification against a wrong audience', async () => {
|
||||
const jwt = await generateClientAssertion({
|
||||
clientId: 'c',
|
||||
audience: 'https://a/token',
|
||||
privateJwk,
|
||||
});
|
||||
const verifier = await importJWK(publicJwk, 'RS256');
|
||||
await expect(
|
||||
jwtVerify(jwt, verifier, { audience: 'https://other/token' }),
|
||||
).rejects.toThrow();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,22 @@
|
||||
import { SignJWT, importJWK, type JWK } from 'jose';
|
||||
|
||||
export interface GenerateClientAssertionInput {
|
||||
clientId: string;
|
||||
audience: string;
|
||||
privateJwk: JWK;
|
||||
expiresIn?: string;
|
||||
}
|
||||
|
||||
export async function generateClientAssertion(
|
||||
input: GenerateClientAssertionInput,
|
||||
): Promise<string> {
|
||||
const privateKey = await importJWK(input.privateJwk, 'RS256');
|
||||
return new SignJWT({})
|
||||
.setProtectedHeader({ alg: 'RS256', typ: 'JWT' })
|
||||
.setIssuer(input.clientId)
|
||||
.setSubject(input.clientId)
|
||||
.setAudience(input.audience)
|
||||
.setIssuedAt()
|
||||
.setExpirationTime(input.expiresIn ?? '2m')
|
||||
.sign(privateKey);
|
||||
}
|
||||
@@ -0,0 +1,65 @@
|
||||
import { createHash } from 'crypto';
|
||||
import {
|
||||
base64Url,
|
||||
generateCodeChallenge,
|
||||
generateCodeVerifier,
|
||||
generateState,
|
||||
} from './pkce.util';
|
||||
|
||||
describe('pkce.util', () => {
|
||||
describe('base64Url', () => {
|
||||
it('strips padding and replaces + and / with - and _', () => {
|
||||
const input = Buffer.from([0xfb, 0xff, 0xbf, 0xfe]);
|
||||
const out = base64Url(input);
|
||||
expect(out).not.toMatch(/[+/=]/);
|
||||
});
|
||||
});
|
||||
|
||||
describe('generateCodeVerifier', () => {
|
||||
it('returns a base64url-safe string', () => {
|
||||
expect(generateCodeVerifier()).toMatch(/^[A-Za-z0-9_-]+$/);
|
||||
});
|
||||
|
||||
it('produces unique values across calls', () => {
|
||||
const a = generateCodeVerifier();
|
||||
const b = generateCodeVerifier();
|
||||
expect(a).not.toEqual(b);
|
||||
});
|
||||
|
||||
it('produces at least 43 characters (RFC 7636 minimum)', () => {
|
||||
expect(generateCodeVerifier().length).toBeGreaterThanOrEqual(43);
|
||||
});
|
||||
});
|
||||
|
||||
describe('generateCodeChallenge', () => {
|
||||
it('equals base64url(sha256(verifier))', () => {
|
||||
const verifier = 'fixed-test-verifier';
|
||||
const expected = createHash('sha256')
|
||||
.update(verifier)
|
||||
.digest('base64')
|
||||
.replace(/\+/g, '-')
|
||||
.replace(/\//g, '_')
|
||||
.replace(/=/g, '');
|
||||
expect(generateCodeChallenge(verifier)).toBe(expected);
|
||||
});
|
||||
|
||||
it('is deterministic for the same verifier', () => {
|
||||
const verifier = generateCodeVerifier();
|
||||
expect(generateCodeChallenge(verifier)).toBe(generateCodeChallenge(verifier));
|
||||
});
|
||||
|
||||
it('differs for different verifiers', () => {
|
||||
expect(generateCodeChallenge('a')).not.toBe(generateCodeChallenge('b'));
|
||||
});
|
||||
});
|
||||
|
||||
describe('generateState', () => {
|
||||
it('returns a base64url-safe string', () => {
|
||||
expect(generateState()).toMatch(/^[A-Za-z0-9_-]+$/);
|
||||
});
|
||||
|
||||
it('produces unique values across calls', () => {
|
||||
expect(generateState()).not.toEqual(generateState());
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,21 @@
|
||||
import { createHash, randomBytes } from 'crypto';
|
||||
|
||||
export function base64Url(buffer: Buffer): string {
|
||||
return buffer
|
||||
.toString('base64')
|
||||
.replace(/\+/g, '-')
|
||||
.replace(/\//g, '_')
|
||||
.replace(/=/g, '');
|
||||
}
|
||||
|
||||
export function generateCodeVerifier(): string {
|
||||
return base64Url(randomBytes(64));
|
||||
}
|
||||
|
||||
export function generateCodeChallenge(codeVerifier: string): string {
|
||||
return base64Url(createHash('sha256').update(codeVerifier).digest());
|
||||
}
|
||||
|
||||
export function generateState(): string {
|
||||
return base64Url(randomBytes(32));
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
import {
|
||||
Body,
|
||||
Controller,
|
||||
Get,
|
||||
HttpCode,
|
||||
HttpStatus,
|
||||
Post,
|
||||
Query,
|
||||
Req,
|
||||
UseGuards,
|
||||
} from '@nestjs/common';
|
||||
import {
|
||||
ApiBearerAuth,
|
||||
ApiOkResponse,
|
||||
ApiOperation,
|
||||
ApiTags,
|
||||
} from '@nestjs/swagger';
|
||||
import type { TCurrentUser } from '@tria-plc/api-common/modules/auth/types/current-user.type';
|
||||
import { IsPublic } from '@tria-plc/api-common/modules/auth/decorators/public.decorator';
|
||||
import { JwtGuard } from '@tria-plc/api-common/modules/auth/services/jwt.guard';
|
||||
import { OptionalJwtGuard } from './optional-jwt.guard';
|
||||
import {
|
||||
CompleteVerificationResultDto,
|
||||
StartVerificationDto,
|
||||
VerifaydaCallbackDto,
|
||||
VerificationStatusDto,
|
||||
} from './verifayda.dto';
|
||||
import { VerifaydaService } from './verifayda.service';
|
||||
|
||||
/** Minimal slices of the Express req we touch (avoids a hard dependency on
|
||||
* `@types/express`, which isn't resolved in this package). */
|
||||
interface RequestWithOptionalUser {
|
||||
user?: TCurrentUser;
|
||||
}
|
||||
interface RequestWithUser {
|
||||
user: TCurrentUser;
|
||||
}
|
||||
|
||||
@ApiTags('Fayda Verification')
|
||||
@Controller('fayda/verification')
|
||||
export class VerifaydaController {
|
||||
constructor(private readonly service: VerifaydaService) {}
|
||||
|
||||
@Post('start')
|
||||
@IsPublic()
|
||||
@HttpCode(HttpStatus.OK)
|
||||
@UseGuards(OptionalJwtGuard)
|
||||
@ApiBearerAuth('JWT-auth')
|
||||
@ApiOperation({
|
||||
summary: 'Start a VeriFayda 2.0 verification session',
|
||||
description: `Creates a verification session and returns the eSignet authorize URL the frontend should send the user to.
|
||||
|
||||
- Works for **logged-in users** and **guests**. If a valid bearer token is present, the verification is tied to that user.
|
||||
- **VERIFY** (default): the user proves their identity and \`/complete\` returns the verified attributes (name, email, phone, dob, gender).
|
||||
- **LOGIN**: \`/complete\` resolves/creates the user and returns a JWT.
|
||||
- The returned \`authorizationUrl\` already carries the PKCE \`code_challenge\`, CSRF \`state\`, requested \`claims\`, and \`code_challenge_method=S256\`. The frontend simply navigates to it (full page or popup).`,
|
||||
})
|
||||
@ApiOkResponse({
|
||||
description: 'Authorize URL the frontend should redirect the user to.',
|
||||
schema: {
|
||||
example: {
|
||||
authorizationUrl:
|
||||
'https://esignet.example.com/authorize?client_id=...&state=...&code_challenge=...',
|
||||
},
|
||||
},
|
||||
})
|
||||
async start(
|
||||
@Body() dto: StartVerificationDto,
|
||||
@Req() req: RequestWithOptionalUser,
|
||||
): Promise<{ authorizationUrl: string }> {
|
||||
const authorizationUrl = await this.service.startVerification({
|
||||
purpose: dto.purpose ?? 'VERIFY',
|
||||
platform: dto.platform ?? 'WEB',
|
||||
userId: req.user?.id,
|
||||
wantsPasswordSetup: dto.wantsPasswordSetup ?? false,
|
||||
});
|
||||
return { authorizationUrl };
|
||||
}
|
||||
|
||||
@Get('complete')
|
||||
@IsPublic()
|
||||
@ApiOperation({
|
||||
summary: 'Complete a verification (Fayda redirect / client callback lands here)',
|
||||
description: `This is the registered Fayda \`redirect_uri\`. Fayda redirects the browser here with \`?code&state\``,
|
||||
})
|
||||
@ApiOkResponse({ type: CompleteVerificationResultDto })
|
||||
async complete(
|
||||
@Query() dto: VerifaydaCallbackDto,
|
||||
): Promise<CompleteVerificationResultDto> {
|
||||
return this.service.completeVerification(dto);
|
||||
}
|
||||
|
||||
@Get('status')
|
||||
@UseGuards(JwtGuard)
|
||||
@ApiBearerAuth('JWT-auth')
|
||||
@ApiOperation({
|
||||
summary: "Get the current user's Fayda verification status",
|
||||
description:
|
||||
'Returns whether the authenticated user has linked a verified Fayda identity to their account, when, and the name on file.',
|
||||
})
|
||||
@ApiOkResponse({ type: VerificationStatusDto })
|
||||
async status(
|
||||
@Req() req: RequestWithUser,
|
||||
): Promise<VerificationStatusDto> {
|
||||
return this.service.getVerificationStatus(req.user.id);
|
||||
}
|
||||
}
|
||||
105
apps/edr-freight-api/src/modules/verifayda/verifayda.dto.ts
Normal file
105
apps/edr-freight-api/src/modules/verifayda/verifayda.dto.ts
Normal file
@@ -0,0 +1,105 @@
|
||||
import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger';
|
||||
import { IsIn, IsOptional, IsString } from 'class-validator';
|
||||
|
||||
export class StartVerificationDto {
|
||||
@ApiPropertyOptional({
|
||||
enum: ['LOGIN', 'VERIFY'],
|
||||
default: 'VERIFY',
|
||||
description:
|
||||
'Reason for verification. VERIFY returns the verified identity attributes; LOGIN resolves/creates a user and returns a JWT.',
|
||||
})
|
||||
@IsOptional()
|
||||
@IsIn(['LOGIN', 'VERIFY'])
|
||||
purpose?: 'LOGIN' | 'VERIFY';
|
||||
|
||||
@ApiPropertyOptional({
|
||||
enum: ['WEB', 'MOBILE'],
|
||||
default: 'WEB',
|
||||
description:
|
||||
'Client platform. Selects which OAuth redirect_uri is sent to eSignet: WEB uses FAYDA_WEB_REDIRECT_URI, MOBILE uses FAYDA_REDIRECT_URI. Both land on the same /complete endpoint with identical handling.',
|
||||
})
|
||||
@IsOptional()
|
||||
@IsIn(['WEB', 'MOBILE'])
|
||||
platform?: 'WEB' | 'MOBILE';
|
||||
|
||||
@ApiPropertyOptional({
|
||||
type: Boolean,
|
||||
default: false,
|
||||
description:
|
||||
'Set to true when the user opts in to full account registration (checkbox). ' +
|
||||
'When true, the /complete response includes a short-lived token and promptPasswordSetup=true ' +
|
||||
'so the frontend can immediately prompt for a password via POST /v1/auth/set-fayda-password.',
|
||||
})
|
||||
@IsOptional()
|
||||
wantsPasswordSetup?: boolean;
|
||||
}
|
||||
|
||||
export class CompleteVerificationResultDto {
|
||||
@ApiProperty({ enum: ['LOGIN', 'VERIFY'] })
|
||||
purpose!: 'LOGIN' | 'VERIFY';
|
||||
|
||||
@ApiProperty() verified!: boolean;
|
||||
|
||||
@ApiPropertyOptional({ description: 'JWT. LOGIN: session token for the authenticated user. VERIFY: short-lived token for calling /v1/auth/set-fayda-password.' })
|
||||
token?: string;
|
||||
|
||||
@ApiPropertyOptional()
|
||||
refreshToken?: string;
|
||||
|
||||
@ApiPropertyOptional({
|
||||
description: 'Authenticated user summary (LOGIN flow only; same shape as /auth/login).',
|
||||
})
|
||||
user?: {
|
||||
id: string;
|
||||
email: string;
|
||||
role: string;
|
||||
passengerId?: string;
|
||||
agentId?: string;
|
||||
};
|
||||
|
||||
@ApiPropertyOptional({ description: 'Verified full name from Fayda (VERIFY flow).' })
|
||||
fullName?: string;
|
||||
|
||||
@ApiPropertyOptional({ description: 'Verified email from Fayda (VERIFY flow).' })
|
||||
email?: string;
|
||||
|
||||
@ApiPropertyOptional({ description: 'Verified phone number from Fayda (VERIFY flow).' })
|
||||
phoneNumber?: string;
|
||||
|
||||
@ApiPropertyOptional({
|
||||
description: 'Verified date of birth from Fayda, ISO yyyy-MM-dd (VERIFY flow).',
|
||||
})
|
||||
birthdate?: string;
|
||||
|
||||
@ApiPropertyOptional({ description: 'Verified gender from Fayda (VERIFY flow).' })
|
||||
gender?: string;
|
||||
|
||||
@ApiPropertyOptional({ description: 'Whether the verified identity was saved to IAM. False if the IAM write failed.' })
|
||||
userDataSaved?: boolean;
|
||||
|
||||
@ApiPropertyOptional({ description: 'IAM user ID of the verified identity (VERIFY flow).' })
|
||||
iamUserId?: string;
|
||||
|
||||
@ApiPropertyOptional({ description: 'True when the IAM account has not yet set a password (VERIFY flow).' })
|
||||
requiresPassword?: boolean;
|
||||
|
||||
@ApiPropertyOptional({
|
||||
description:
|
||||
'True when the user opted in to immediate password setup (wantsPasswordSetup=true at start) ' +
|
||||
'AND they have not yet set a password. Frontend should navigate to the set-password screen.',
|
||||
})
|
||||
promptPasswordSetup?: boolean;
|
||||
}
|
||||
|
||||
export class VerifaydaCallbackDto {
|
||||
@ApiPropertyOptional() @IsOptional() @IsString() code?: string;
|
||||
@ApiPropertyOptional() @IsOptional() @IsString() state?: string;
|
||||
@ApiPropertyOptional() @IsOptional() @IsString() error?: string;
|
||||
@ApiPropertyOptional() @IsOptional() @IsString() error_description?: string;
|
||||
}
|
||||
|
||||
export class VerificationStatusDto {
|
||||
@ApiProperty() verified!: boolean;
|
||||
@ApiPropertyOptional() verifiedAt?: Date;
|
||||
@ApiPropertyOptional() fullName?: string;
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
import { BadGatewayException, ConflictException } from '@nestjs/common';
|
||||
|
||||
export class FaydaTokenExchangeException extends BadGatewayException {
|
||||
constructor(message = 'Fayda token exchange failed') {
|
||||
super({ code: 'FAYDA_TOKEN_EXCHANGE_FAILED', message });
|
||||
}
|
||||
}
|
||||
|
||||
export class FaydaUserInfoException extends BadGatewayException {
|
||||
constructor(message = 'Fayda userinfo fetch failed') {
|
||||
super({ code: 'FAYDA_USERINFO_FAILED', message });
|
||||
}
|
||||
}
|
||||
|
||||
export class FaydaIdentityConflictException extends ConflictException {
|
||||
constructor(message = 'This Fayda identity is already linked to another account') {
|
||||
super({ code: 'FAYDA_IDENTITY_CONFLICT', message });
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,14 @@
|
||||
import { Module } from '@nestjs/common';
|
||||
import { TypeOrmModule } from '@nestjs/typeorm';
|
||||
import { VerifaydaController } from './verifayda.controller';
|
||||
import { FaydaCallbackController } from './fayda-callback.controller';
|
||||
import { VerifaydaService } from './verifayda.service';
|
||||
import { FaydaVerificationSession } from './entities/fayda-verification-session.entity';
|
||||
|
||||
@Module({
|
||||
imports: [TypeOrmModule.forFeature([FaydaVerificationSession])],
|
||||
controllers: [VerifaydaController, FaydaCallbackController],
|
||||
providers: [VerifaydaService],
|
||||
exports: [VerifaydaService],
|
||||
})
|
||||
export class VerifaydaModule {}
|
||||
597
apps/edr-freight-api/src/modules/verifayda/verifayda.service.ts
Normal file
597
apps/edr-freight-api/src/modules/verifayda/verifayda.service.ts
Normal file
@@ -0,0 +1,597 @@
|
||||
import {
|
||||
BadRequestException,
|
||||
Injectable,
|
||||
Logger,
|
||||
ServiceUnavailableException,
|
||||
UnauthorizedException,
|
||||
} from '@nestjs/common';
|
||||
import { ConfigService } from '@nestjs/config';
|
||||
import { InjectDataSource, InjectRepository } from '@nestjs/typeorm';
|
||||
import { DataSource, Repository } from 'typeorm';
|
||||
import { generateToken, generateRefreshToken } from '@tria-plc/api-common/utils/token';
|
||||
import { FaydaConfig, FaydaPlatform } from '../../config/fayda.config';
|
||||
import { FaydaVerificationSession } from './entities/fayda-verification-session.entity';
|
||||
import {
|
||||
generateCodeChallenge,
|
||||
generateCodeVerifier,
|
||||
generateState,
|
||||
} from './utils/pkce.util';
|
||||
import { generateClientAssertion } from './utils/client-assertion.util';
|
||||
import { VerifaydaCallbackDto, VerificationStatusDto } from './verifayda.dto';
|
||||
import {
|
||||
FaydaTokenExchangeException,
|
||||
FaydaUserInfoException,
|
||||
} from './verifayda.errors';
|
||||
import {
|
||||
FaydaTokenResponse,
|
||||
FaydaUserInfo,
|
||||
NormalizedFaydaUserInfo,
|
||||
VerifaydaPurpose,
|
||||
} from './verifayda.types';
|
||||
|
||||
export interface StartVerificationInput {
|
||||
purpose: VerifaydaPurpose;
|
||||
platform?: FaydaPlatform;
|
||||
userId?: string; // iamUserId of the authenticated user, if any
|
||||
wantsPasswordSetup?: boolean;
|
||||
}
|
||||
|
||||
export interface FaydaUserSummary {
|
||||
id: string;
|
||||
email: string;
|
||||
role: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Result of completing a verification. `verified` is always true on success.
|
||||
* LOGIN additionally returns a JWT + user; VERIFY returns the verified identity
|
||||
* attributes (name, email, phone, dob, gender) for the caller to consume.
|
||||
*/
|
||||
export interface CompleteVerificationResult {
|
||||
purpose: VerifaydaPurpose;
|
||||
verified: boolean;
|
||||
token?: string;
|
||||
refreshToken?: string;
|
||||
requiresPassword?: boolean;
|
||||
promptPasswordSetup?: boolean;
|
||||
iamUserId?: string;
|
||||
user?: FaydaUserSummary;
|
||||
fullName?: string;
|
||||
email?: string;
|
||||
phoneNumber?: string;
|
||||
birthdate?: string;
|
||||
gender?: string;
|
||||
userDataSaved?: boolean;
|
||||
}
|
||||
|
||||
@Injectable()
|
||||
export class VerifaydaService {
|
||||
private readonly logger = new Logger(VerifaydaService.name);
|
||||
|
||||
private readonly faydaConfig: FaydaConfig;
|
||||
|
||||
constructor(
|
||||
private readonly config: ConfigService,
|
||||
@InjectRepository(FaydaVerificationSession)
|
||||
private readonly sessionRepo: Repository<FaydaVerificationSession>,
|
||||
@InjectDataSource() private readonly dataSource: DataSource,
|
||||
) {
|
||||
const fayda = this.config.get<FaydaConfig>('fayda');
|
||||
if (!fayda) {
|
||||
throw new Error('Fayda config namespace not registered');
|
||||
}
|
||||
this.faydaConfig = fayda;
|
||||
}
|
||||
|
||||
// ==========================================================================
|
||||
// OIDC flow
|
||||
// ==========================================================================
|
||||
|
||||
async startVerification(input: StartVerificationInput): Promise<string> {
|
||||
if (!this.faydaConfig.enabled) {
|
||||
throw new ServiceUnavailableException({
|
||||
code: 'FAYDA_DISABLED',
|
||||
message: 'Fayda integration is not enabled',
|
||||
});
|
||||
}
|
||||
|
||||
const state = generateState();
|
||||
const codeVerifier = generateCodeVerifier();
|
||||
const codeChallenge = generateCodeChallenge(codeVerifier);
|
||||
const expiresAt = new Date(
|
||||
Date.now() + this.faydaConfig.sessionTtlMinutes * 60_000,
|
||||
);
|
||||
|
||||
await this.sessionRepo.save(
|
||||
this.sessionRepo.create({
|
||||
state,
|
||||
codeVerifier,
|
||||
purpose: input.purpose,
|
||||
platform: input.platform ?? 'WEB',
|
||||
saveToAccount: input.wantsPasswordSetup ?? false,
|
||||
iamUserId: input.userId ?? null,
|
||||
expiresAt,
|
||||
}),
|
||||
);
|
||||
|
||||
this.logger.log(
|
||||
`Fayda verification started: purpose=${input.purpose} platform=${input.platform ?? 'WEB'} userId=${input.userId ?? 'none'}`,
|
||||
);
|
||||
|
||||
return this.buildAuthorizationUrl({
|
||||
state,
|
||||
codeChallenge,
|
||||
redirectUri: this.redirectUriForPlatform(input.platform ?? 'WEB'),
|
||||
});
|
||||
}
|
||||
|
||||
/** WEB clients use `webRedirectUri`; MOBILE uses the base `redirectUri`. */
|
||||
private redirectUriForPlatform(platform?: FaydaPlatform): string {
|
||||
return platform === 'MOBILE'
|
||||
? this.faydaConfig.redirectUri
|
||||
: this.faydaConfig.webRedirectUri;
|
||||
}
|
||||
|
||||
async completeVerification(
|
||||
query: VerifaydaCallbackDto,
|
||||
): Promise<CompleteVerificationResult> {
|
||||
if (query.error) {
|
||||
this.logger.warn(`Fayda callback returned error: ${query.error}`);
|
||||
if (query.state) {
|
||||
await this.markSessionFailed(
|
||||
query.state,
|
||||
query.error,
|
||||
query.error_description,
|
||||
);
|
||||
}
|
||||
throw new BadRequestException({
|
||||
code: 'FAYDA_AUTH_ERROR',
|
||||
message: query.error,
|
||||
description: query.error_description,
|
||||
});
|
||||
}
|
||||
|
||||
if (!query.code || !query.state) {
|
||||
throw new BadRequestException({
|
||||
code: 'FAYDA_MISSING_PARAMETERS',
|
||||
message: 'code and state are required',
|
||||
});
|
||||
}
|
||||
|
||||
const session = await this.sessionRepo.findOne({
|
||||
where: { state: query.state },
|
||||
});
|
||||
if (!session || session.status !== 'PENDING') {
|
||||
this.logger.warn('Fayda complete with unknown or non-pending state');
|
||||
throw new BadRequestException({
|
||||
code: 'FAYDA_INVALID_STATE',
|
||||
message: 'Verification session is invalid or already used',
|
||||
});
|
||||
}
|
||||
if (session.expiresAt.getTime() < Date.now()) {
|
||||
await this.markSessionFailed(query.state, 'session_expired');
|
||||
throw new BadRequestException({
|
||||
code: 'FAYDA_SESSION_EXPIRED',
|
||||
message: 'Verification session has expired; start again',
|
||||
});
|
||||
}
|
||||
|
||||
try {
|
||||
const tokens = await this.exchangeCodeForTokens(
|
||||
query.code,
|
||||
session.codeVerifier,
|
||||
this.redirectUriForPlatform(session.platform as FaydaPlatform),
|
||||
);
|
||||
const userInfo = await this.fetchUserInfo(tokens.access_token);
|
||||
const normalized = this.normalizeUserInfo(userInfo);
|
||||
|
||||
if (!normalized.sub) {
|
||||
throw new FaydaUserInfoException('Fayda userinfo missing required sub');
|
||||
}
|
||||
|
||||
let result: CompleteVerificationResult;
|
||||
if (session.purpose === 'LOGIN') {
|
||||
const { userId } = await this.handleLoginSuccess(normalized);
|
||||
const login = await this.issueLoginToken(userId);
|
||||
result = { purpose: 'LOGIN', verified: true, ...login };
|
||||
} else {
|
||||
// VERIFY — prove identity, save to IAM, return verified attributes + short-lived token.
|
||||
const { iamUserId, userDataSaved } = await this.upsertIamUser(normalized);
|
||||
|
||||
let sessionToken: { token: string; refreshToken: string; requiresPassword: boolean } | undefined;
|
||||
if (iamUserId) {
|
||||
try {
|
||||
sessionToken = await this.createFaydaSession(iamUserId);
|
||||
} catch (err) {
|
||||
this.logger.warn(`Fayda session creation failed: ${(err as Error).message}`);
|
||||
}
|
||||
}
|
||||
|
||||
result = {
|
||||
purpose: 'VERIFY',
|
||||
verified: true,
|
||||
fullName: normalized.fullName,
|
||||
email: normalized.email,
|
||||
phoneNumber: normalized.phoneNumber,
|
||||
birthdate: normalized.birthdate,
|
||||
gender: normalized.gender,
|
||||
userDataSaved,
|
||||
iamUserId: iamUserId ?? undefined,
|
||||
token: sessionToken?.token,
|
||||
refreshToken: sessionToken?.refreshToken,
|
||||
requiresPassword: sessionToken?.requiresPassword,
|
||||
promptPasswordSetup: session.saveToAccount && (sessionToken?.requiresPassword ?? false),
|
||||
};
|
||||
}
|
||||
|
||||
await this.sessionRepo.update(session.id, {
|
||||
status: 'COMPLETED',
|
||||
completedAt: new Date(),
|
||||
codeVerifier: '',
|
||||
});
|
||||
|
||||
this.logger.log(
|
||||
`Fayda verification completed: purpose=${session.purpose} platform=${session.platform}`,
|
||||
);
|
||||
return result;
|
||||
} catch (err) {
|
||||
const reason = this.classifyFailureReason(err);
|
||||
this.logger.error(
|
||||
`Fayda verification failed: reason=${reason} message=${(err as Error).message}`,
|
||||
);
|
||||
await this.markSessionFailed(
|
||||
query.state,
|
||||
reason,
|
||||
(err as Error).message,
|
||||
);
|
||||
throw err;
|
||||
}
|
||||
}
|
||||
|
||||
private async issueLoginToken(
|
||||
_userId: string,
|
||||
): Promise<{ token: string; user: FaydaUserSummary }> {
|
||||
throw new UnauthorizedException({
|
||||
code: 'FAYDA_LOGIN_MIGRATED_TO_IAM',
|
||||
message: 'Fayda login tokens are issued by the IAM package auth endpoints.',
|
||||
});
|
||||
}
|
||||
|
||||
async getVerificationStatus(iamUserId: string): Promise<VerificationStatusDto> {
|
||||
const rows = await this.dataSource.query<{ verified_by: string | null; updated_at: Date | null; name: { en: string; am: string } | null }[]>(
|
||||
`SELECT verified_by, updated_at, name FROM iam.users WHERE id = $1 LIMIT 1`,
|
||||
[iamUserId],
|
||||
);
|
||||
const iam = rows[0] ?? null;
|
||||
const faydaVerified = iam?.verified_by === 'fayda';
|
||||
const faydaVerifiedAt = faydaVerified && iam?.updated_at ? new Date(iam.updated_at) : undefined;
|
||||
const fullName = iam?.name?.en ?? iam?.name?.am ?? undefined;
|
||||
return { verified: faydaVerified, verifiedAt: faydaVerifiedAt, fullName };
|
||||
}
|
||||
|
||||
// ==========================================================================
|
||||
// OIDC internals
|
||||
// ==========================================================================
|
||||
|
||||
private buildAuthorizationUrl(args: {
|
||||
state: string;
|
||||
codeChallenge: string;
|
||||
redirectUri: string;
|
||||
}): string {
|
||||
const params = new URLSearchParams({
|
||||
client_id: this.faydaConfig.clientId,
|
||||
response_type: 'code',
|
||||
redirect_uri: args.redirectUri,
|
||||
scope: this.faydaConfig.scope,
|
||||
state: args.state,
|
||||
code_challenge: args.codeChallenge,
|
||||
code_challenge_method: 'S256',
|
||||
acr_values: this.faydaConfig.acrValues,
|
||||
claims_locales: this.faydaConfig.claimsLocales,
|
||||
});
|
||||
|
||||
// Every claim is marked essential so eSignet shows them locked/pre-checked
|
||||
// on the consent screen — the user cannot toggle any off; they either
|
||||
// consent to all of them or the whole flow is cancelled (?error=...).
|
||||
const claims = {
|
||||
userinfo: {
|
||||
name: { essential: true },
|
||||
phone_number: { essential: true },
|
||||
email: { essential: true },
|
||||
birthdate: { essential: true },
|
||||
gender: { essential: true },
|
||||
address: { essential: true },
|
||||
nationality: { essential: true },
|
||||
picture: { essential: true },
|
||||
},
|
||||
id_token: {},
|
||||
};
|
||||
params.set('claims', JSON.stringify(claims));
|
||||
|
||||
return `${this.faydaConfig.authorizationEndpoint}?${params.toString()}`;
|
||||
}
|
||||
|
||||
private async exchangeCodeForTokens(
|
||||
code: string,
|
||||
codeVerifier: string,
|
||||
redirectUri: string,
|
||||
): Promise<FaydaTokenResponse> {
|
||||
const clientAssertion = await generateClientAssertion({
|
||||
clientId: this.faydaConfig.clientId,
|
||||
audience: this.faydaConfig.tokenEndpoint,
|
||||
privateJwk: this.faydaConfig.privateJwk,
|
||||
});
|
||||
|
||||
const body = new URLSearchParams({
|
||||
grant_type: 'authorization_code',
|
||||
code,
|
||||
redirect_uri: redirectUri,
|
||||
client_id: this.faydaConfig.clientId,
|
||||
client_assertion_type:
|
||||
'urn:ietf:params:oauth:client-assertion-type:jwt-bearer',
|
||||
client_assertion: clientAssertion,
|
||||
code_verifier: codeVerifier,
|
||||
});
|
||||
|
||||
const response = await fetch(this.faydaConfig.tokenEndpoint, {
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': 'application/x-www-form-urlencoded' },
|
||||
body,
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
let detail = '';
|
||||
try {
|
||||
detail = await response.text();
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
throw new FaydaTokenExchangeException(
|
||||
`Fayda token endpoint returned ${response.status}${detail ? `: ${detail}` : ''}`,
|
||||
);
|
||||
}
|
||||
|
||||
return (await response.json()) as FaydaTokenResponse;
|
||||
}
|
||||
|
||||
private async fetchUserInfo(accessToken: string): Promise<FaydaUserInfo> {
|
||||
const response = await fetch(this.faydaConfig.userInfoEndpoint, {
|
||||
method: 'GET',
|
||||
headers: { Authorization: `Bearer ${accessToken}` },
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
throw new FaydaUserInfoException(
|
||||
`Fayda userinfo endpoint returned ${response.status}`,
|
||||
);
|
||||
}
|
||||
|
||||
const contentType = response.headers.get('content-type') ?? '';
|
||||
const raw = await response.text();
|
||||
|
||||
if (contentType.includes('application/json')) {
|
||||
return JSON.parse(raw) as FaydaUserInfo;
|
||||
}
|
||||
|
||||
// Signed JWT response — decode payload (signature verification = production TODO)
|
||||
if (raw.split('.').length === 3) {
|
||||
const payloadB64 = raw.split('.')[1];
|
||||
const normalizedB64 = payloadB64.replace(/-/g, '+').replace(/_/g, '/');
|
||||
const json = Buffer.from(normalizedB64, 'base64').toString('utf8');
|
||||
return JSON.parse(json) as FaydaUserInfo;
|
||||
}
|
||||
|
||||
throw new FaydaUserInfoException(
|
||||
'Unsupported Fayda userinfo response format',
|
||||
);
|
||||
}
|
||||
|
||||
private normalizeUserInfo(raw: FaydaUserInfo): NormalizedFaydaUserInfo {
|
||||
const nameEn = raw['name#en'] as string | undefined;
|
||||
const nameAm = raw['name#am'] as string | undefined;
|
||||
const genderEn = raw['gender#en'] as string | undefined;
|
||||
const genderAm = raw['gender#am'] as string | undefined;
|
||||
const addressEn = raw['address#en'] as string | undefined;
|
||||
const addressAm = raw['address#am'] as string | undefined;
|
||||
const rawPhone = (raw.phone_number ?? raw['phone_number#en'] ?? raw['phone_number#am'] ?? raw.phone) as string | undefined;
|
||||
|
||||
return {
|
||||
sub: raw.sub,
|
||||
fullName: (raw.name as string | undefined) ?? nameEn ?? nameAm,
|
||||
phoneNumber: rawPhone ? this.standardizePhoneNumber(rawPhone) : undefined,
|
||||
rawPhoneNumber: rawPhone,
|
||||
email: raw.email as string | undefined,
|
||||
gender: genderEn ?? genderAm ?? (raw.gender as string | undefined),
|
||||
birthdate: raw.birthdate as string | undefined,
|
||||
picture: raw.picture as string | undefined,
|
||||
nameEn,
|
||||
nameAm,
|
||||
genderEn,
|
||||
genderAm,
|
||||
addressEn,
|
||||
addressAm,
|
||||
};
|
||||
}
|
||||
|
||||
private standardizePhoneNumber(phone: string): string {
|
||||
const digits = phone.replace(/\D/g, '');
|
||||
if (digits.startsWith('251')) return `+${digits}`;
|
||||
if (digits.startsWith('0')) return `+251${digits.slice(1)}`;
|
||||
return `+${digits}`;
|
||||
}
|
||||
|
||||
// LOGIN via Fayda is handled entirely by the IAM package's own OIDC flow.
|
||||
// This method is kept as a stub so completeVerification() still compiles;
|
||||
// it throws immediately without touching the database.
|
||||
private async handleLoginSuccess(
|
||||
_normalized: NormalizedFaydaUserInfo,
|
||||
): Promise<{ userId: string }> {
|
||||
throw new UnauthorizedException({
|
||||
code: 'FAYDA_LOGIN_MIGRATED_TO_IAM',
|
||||
message: 'Fayda login tokens are issued by the IAM package at /v1/auth/fayda endpoints.',
|
||||
});
|
||||
}
|
||||
|
||||
private async upsertIamUser(
|
||||
normalized: NormalizedFaydaUserInfo,
|
||||
): Promise<{ iamUserId: string | null; userDataSaved: boolean }> {
|
||||
try {
|
||||
const iamMetadata = {
|
||||
sub: normalized.sub,
|
||||
address: { am: normalized.addressAm ?? '', en: normalized.addressEn ?? '' },
|
||||
email: normalized.email ?? '',
|
||||
gender: { am: normalized.genderAm ?? '', en: normalized.genderEn ?? '' },
|
||||
name: { am: normalized.nameAm ?? '', en: normalized.nameEn ?? '' },
|
||||
phoneNumber: normalized.rawPhoneNumber ?? '',
|
||||
};
|
||||
|
||||
// Step 1 — already linked to this Fayda sub; ensure verified_by is set
|
||||
const bySub = await this.dataSource.query<{ id: string }[]>(
|
||||
`SELECT id FROM iam.users WHERE metadata->>'sub' = $1 LIMIT 1`,
|
||||
[normalized.sub],
|
||||
);
|
||||
if (bySub.length > 0) {
|
||||
await this.dataSource.query(
|
||||
`UPDATE iam.users SET verified_by = 'fayda', updated_at = NOW() WHERE id = $1`,
|
||||
[bySub[0].id],
|
||||
);
|
||||
return { iamUserId: bySub[0].id, userDataSaved: true };
|
||||
}
|
||||
|
||||
// Step 2 — existing user by phone or email, not yet Fayda-verified
|
||||
const conditions: string[] = [];
|
||||
const params: unknown[] = [];
|
||||
if (normalized.phoneNumber) {
|
||||
params.push(normalized.phoneNumber);
|
||||
conditions.push(`phone_number = $${params.length}`);
|
||||
}
|
||||
if (normalized.email) {
|
||||
params.push(normalized.email);
|
||||
conditions.push(`email = $${params.length}`);
|
||||
}
|
||||
if (conditions.length > 0) {
|
||||
const byContact = await this.dataSource.query<{ id: string }[]>(
|
||||
`SELECT id FROM iam.users WHERE ${conditions.join(' OR ')} LIMIT 1`,
|
||||
params,
|
||||
);
|
||||
if (byContact.length > 0) {
|
||||
const existingId = byContact[0].id;
|
||||
await this.dataSource.query(
|
||||
`UPDATE iam.users
|
||||
SET metadata = COALESCE(metadata, '{}'::jsonb) || $1::jsonb,
|
||||
verified_by = 'fayda',
|
||||
updated_at = NOW()
|
||||
WHERE id = $2`,
|
||||
[JSON.stringify(iamMetadata), existingId],
|
||||
);
|
||||
return { iamUserId: existingId, userDataSaved: true };
|
||||
}
|
||||
}
|
||||
|
||||
// Step 3 — new user
|
||||
const name = { am: normalized.nameAm ?? '', en: normalized.nameEn ?? '' };
|
||||
const username = normalized.phoneNumber ?? normalized.email ?? normalized.sub;
|
||||
const inserted = await this.dataSource.query<{ id: string }[]>(
|
||||
`INSERT INTO iam.users (
|
||||
id, name, username, email, phone_number, metadata,
|
||||
user_type, status, is_active, has_set_password,
|
||||
is_phone_number_verified, verified_by,
|
||||
created_at, updated_at
|
||||
) VALUES (
|
||||
gen_random_uuid(), $1::jsonb, $2, $3, $4, $5::jsonb,
|
||||
'individual', 'submitted', true, false,
|
||||
false, 'fayda',
|
||||
NOW(), NOW()
|
||||
) RETURNING id`,
|
||||
[
|
||||
JSON.stringify(name),
|
||||
username,
|
||||
normalized.email ?? null,
|
||||
normalized.phoneNumber ?? null,
|
||||
JSON.stringify(iamMetadata),
|
||||
],
|
||||
);
|
||||
return { iamUserId: inserted[0].id, userDataSaved: true };
|
||||
} catch (err) {
|
||||
this.logger.error(`Fayda IAM upsert failed: ${(err as Error).message}`);
|
||||
return { iamUserId: null, userDataSaved: false };
|
||||
}
|
||||
}
|
||||
|
||||
private async createFaydaSession(
|
||||
iamUserId: string,
|
||||
): Promise<{ token: string; refreshToken: string; requiresPassword: boolean }> {
|
||||
const rows = await this.dataSource.query<{
|
||||
id: string;
|
||||
email: string;
|
||||
name: { en: string; am: string } | null;
|
||||
username: string;
|
||||
phone_number: string | null;
|
||||
has_set_password: boolean;
|
||||
status: string;
|
||||
}[]>(
|
||||
`SELECT id, email, name, username, phone_number, has_set_password, status
|
||||
FROM iam.users WHERE id = $1 LIMIT 1`,
|
||||
[iamUserId],
|
||||
);
|
||||
if (!rows.length) throw new Error(`IAM user ${iamUserId} not found`);
|
||||
const u = rows[0];
|
||||
|
||||
const userInfo = {
|
||||
id: u.id,
|
||||
email: u.email ?? '',
|
||||
name: u.name ?? { en: '', am: '' },
|
||||
userType: 'individual',
|
||||
status: u.status,
|
||||
hasSetPassword: u.has_set_password,
|
||||
isPhoneNumberVerified: false,
|
||||
hasFinishedRegistration: false,
|
||||
hasFinishedDMSOnboarding: false,
|
||||
username: u.username,
|
||||
phoneNumber: u.phone_number ?? '',
|
||||
roles: [],
|
||||
permissions: [],
|
||||
employee: [],
|
||||
};
|
||||
|
||||
const sessions = await this.dataSource.query<{ id: string }[]>(
|
||||
`INSERT INTO iam.sessions
|
||||
(id, email, device, "userInfo", expiry_time, refresh_count, status, user_id)
|
||||
VALUES (gen_random_uuid(), $1, 'fayda-verify', $2::jsonb, NOW() + INTERVAL '1 day', 0, 'ACTIVE', $3)
|
||||
ON CONFLICT (user_id, device) DO UPDATE
|
||||
SET status = 'ACTIVE', "userInfo" = EXCLUDED."userInfo",
|
||||
expiry_time = NOW() + INTERVAL '1 day', updated_at = NOW()
|
||||
RETURNING id`,
|
||||
[u.email ?? '', JSON.stringify(userInfo), iamUserId],
|
||||
);
|
||||
|
||||
const sessionId = sessions[0].id;
|
||||
const token = generateToken({ id: sessionId });
|
||||
const refreshToken = generateRefreshToken({ id: sessionId });
|
||||
|
||||
return { token, refreshToken, requiresPassword: !u.has_set_password };
|
||||
}
|
||||
|
||||
private async markSessionFailed(
|
||||
state: string,
|
||||
errorCode: string,
|
||||
errorDescription?: string,
|
||||
): Promise<void> {
|
||||
await this.sessionRepo.update(
|
||||
{ state, status: 'PENDING' },
|
||||
{
|
||||
status: 'FAILED',
|
||||
errorCode,
|
||||
errorDescription: errorDescription ?? null,
|
||||
completedAt: new Date(),
|
||||
codeVerifier: '',
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
private classifyFailureReason(err: unknown): string {
|
||||
if (err instanceof FaydaTokenExchangeException) return 'token_exchange_failed';
|
||||
if (err instanceof FaydaUserInfoException) return 'userinfo_failed';
|
||||
return 'verification_failed';
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
export type VerifaydaPurpose = 'LOGIN' | 'VERIFY';
|
||||
|
||||
export interface FaydaTokenResponse {
|
||||
access_token: string;
|
||||
id_token?: string;
|
||||
token_type: string;
|
||||
expires_in?: number;
|
||||
scope?: string;
|
||||
}
|
||||
|
||||
export interface FaydaUserInfo {
|
||||
sub: string;
|
||||
name?: string;
|
||||
'name#en'?: string;
|
||||
'name#am'?: string;
|
||||
phone_number?: string;
|
||||
'phone_number#en'?: string;
|
||||
'phone_number#am'?: string;
|
||||
phone?: string;
|
||||
email?: string;
|
||||
gender?: string;
|
||||
birthdate?: string;
|
||||
picture?: string;
|
||||
address?: Record<string, unknown>;
|
||||
[key: string]: unknown;
|
||||
}
|
||||
|
||||
export interface NormalizedFaydaUserInfo {
|
||||
sub: string;
|
||||
// Convenience / display fields
|
||||
fullName?: string;
|
||||
phoneNumber?: string; // standardized e.g. +251911234567
|
||||
email?: string;
|
||||
gender?: string;
|
||||
birthdate?: string;
|
||||
picture?: string;
|
||||
// Raw localized fields — preserved for IAM-identical writes
|
||||
nameEn?: string;
|
||||
nameAm?: string;
|
||||
genderEn?: string;
|
||||
genderAm?: string;
|
||||
addressEn?: string;
|
||||
addressAm?: string;
|
||||
rawPhoneNumber?: string; // unstandardized, stored in IAM metadata
|
||||
}
|
||||
@@ -1,9 +1,10 @@
|
||||
import { ApiProperty, ApiPropertyOptional, PartialType } from '@nestjs/swagger';
|
||||
import { IsBoolean, IsInt, IsOptional, IsString } from 'class-validator';
|
||||
import { IsBoolean, IsInt, IsOptional, IsString, Matches } from 'class-validator';
|
||||
|
||||
export class CreateAllocationRuleDto {
|
||||
@ApiProperty()
|
||||
@IsString()
|
||||
@Matches(/^[A-Za-z\s]+$/, { message: 'name may only contain letters and spaces' })
|
||||
name!: string;
|
||||
|
||||
@ApiPropertyOptional({ default: 100 })
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger';
|
||||
import { IsEnum, IsNumber, IsOptional, IsString, IsUUID, MaxLength, Min } from 'class-validator';
|
||||
import { IsEnum, IsNumber, IsOptional, IsString, IsUUID, Matches, MaxLength, Min } from 'class-validator';
|
||||
|
||||
import { WAREHOUSE_TYPES, WarehouseType } from '../entities/warehouse.entity';
|
||||
|
||||
@@ -7,6 +7,7 @@ export class CreateWarehouseDto {
|
||||
@ApiProperty()
|
||||
@IsString()
|
||||
@MaxLength(160)
|
||||
@Matches(/^[A-Za-z\s]+$/, { message: 'name may only contain letters and spaces' })
|
||||
name!: string;
|
||||
|
||||
@ApiProperty()
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { ApiProperty, ApiPropertyOptional, PartialType } from '@nestjs/swagger';
|
||||
import { Type } from 'class-transformer';
|
||||
import { IsArray, IsEnum, IsInt, IsNumber, IsOptional, IsString, IsUUID, Min, ValidateNested } from 'class-validator';
|
||||
import { IsArray, IsEnum, IsInt, IsNumber, IsOptional, IsString, IsUUID, Matches, Min, ValidateNested } from 'class-validator';
|
||||
|
||||
import { FEE_RULE_TYPES, FeeRuleType } from '../entities/warehouse-fee-rule.entity';
|
||||
|
||||
@@ -25,6 +25,7 @@ export class FeeRuleTierDto {
|
||||
export class CreateFeeRuleDto {
|
||||
@ApiProperty()
|
||||
@IsString()
|
||||
@Matches(/^[A-Za-z\s]+$/, { message: 'name may only contain letters and spaces' })
|
||||
name!: string;
|
||||
|
||||
@ApiProperty({ enum: FEE_RULE_TYPES })
|
||||
|
||||
@@ -52,6 +52,11 @@ export class FilterWarehouseInventoryDto {
|
||||
@IsEnum(WAREHOUSE_INVENTORY_STATUSES)
|
||||
status?: WarehouseInventoryStatus;
|
||||
|
||||
@ApiPropertyOptional({ enum: ['IMPORT', 'EXPORT'] })
|
||||
@IsOptional()
|
||||
@IsEnum(['IMPORT', 'EXPORT'])
|
||||
direction?: 'IMPORT' | 'EXPORT';
|
||||
|
||||
@ApiPropertyOptional()
|
||||
@IsOptional()
|
||||
@IsString()
|
||||
|
||||
@@ -430,6 +430,7 @@ export class WarehouseInventoryService {
|
||||
...(filter.status ? { status: filter.status } : {}),
|
||||
...(createdAt ? { createdAt } : {}),
|
||||
...(filter.facilityId ? { warehouse: { stationId: filter.facilityId } } : {}),
|
||||
...(filter.direction ? { booking: { tradeDirection: filter.direction } } : {}),
|
||||
};
|
||||
|
||||
const search = filter.search?.trim();
|
||||
@@ -2032,6 +2033,22 @@ export class WarehouseInventoryService {
|
||||
const isTruckLeaving = dto.grossWeight !== undefined && Boolean(dto.gateOutTime);
|
||||
if (isTruckLeaving) {
|
||||
await this.invoices.assertClearanceAllowed(id);
|
||||
|
||||
if (item.bookingId) {
|
||||
const [truckInfo]: Array<{ customerTruckAssignedAt: string | null }> =
|
||||
await this.dataSource.query(
|
||||
`SELECT customer_truck_assigned_at AS "customerTruckAssignedAt"
|
||||
FROM freight.bookings
|
||||
WHERE id = $1 AND deleted_at IS NULL`,
|
||||
[item.bookingId],
|
||||
);
|
||||
const usesCustomerTruck = Boolean(truckInfo?.customerTruckAssignedAt);
|
||||
if (usesCustomerTruck && !this.extractCustomerDeliveryApproval(item.notes)) {
|
||||
throw new BadRequestException(
|
||||
'Customer must approve delivery (sign the handover) before the exit paper can be generated',
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
const releaseDate = isTruckLeaving
|
||||
? dto.releaseDate ? new Date(dto.releaseDate) : new Date()
|
||||
|
||||
Reference in New Issue
Block a user