Files
edr-platform/apps/edr-freight-api/src/modules/train-scheduling/booking-batch.service.ts
Marshal b5a97d344a train
2026-07-14 13:10:00 +00:00

3249 lines
126 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

import {
BadRequestException,
ConflictException,
forwardRef,
Inject,
Injectable,
Logger,
NotFoundException,
OnModuleInit,
Optional,
} from '@nestjs/common';
import { InjectDataSource } from '@nestjs/typeorm';
import { SchedulerRegistry } from '@nestjs/schedule';
import {
Between,
DataSource,
FindOptionsWhere,
ILike,
In,
LessThanOrEqual,
MoreThanOrEqual,
} from 'typeorm';
import { Booking } from '../bookings/entities/booking.entity';
import { BookingsRepository } from '../bookings/bookings.repository';
import { BookingPricingService } from '../bookings/booking-pricing.service';
import { Locomotive } from '../locomotives/entities/locomotive.entity';
import { formatRouteLabel } from '../routes/entities/route.entity';
import { RouteMilestone } from '../routes/entities/route-milestone.entity';
import { TrainSchedule } from '../train-schedules/entities/train-schedule.entity';
import { TrainScheduleBooking } from '../train-schedules/entities/train-schedule-booking.entity';
import { TrainSchedulesRepository } from '../train-schedules/train-schedules.repository';
import { TrainScheduleBookingsRepository } from '../train-schedules/train-schedule-bookings.repository';
import { BookingNotifierService } from './booking-notifier.service';
import { TrainSchedulingService } from './train-scheduling.service';
import { eatDay, groupBookingsIntoBoardWindows } from './batch-window.util';
import {
BATCH_BOARD_STATUSES,
BatchBoardQueryDto,
} from './dto/batch-board-query.dto';
import {
Freight,
PaginatedResponse,
TrainScheduleStatus as TrainScheduleStatusEnum,
} from "@edr/types";
import { BillingService } from "../billing/billing.service";
import {
buildPaginationMeta,
normalizePagination,
} from '../../common/utils/pagination.util';
import {
DEFAULT_BULK_WAGON_CAPACITY_TONS,
DEFAULT_BULK_WAGON_LENGTH_METERS,
DEFAULT_BULK_WAGON_TARE_TONS,
DEFAULT_CONTAINER_WAGON_CAPACITY_TONS,
DEFAULT_CONTAINER_WAGON_LENGTH_METERS,
DEFAULT_CONTAINER_WAGON_TARE_TONS,
DEFAULT_WAGONS_PER_BOOKING,
} from "./booking-batch.constants";
import {
WagonTypeDimensions,
bookingGrossWeightTons,
deriveTrainCapacityFromLocomotive,
sizePartialOfferWagons,
trainHardCaps,
wagonTypeDimensionsFromEntity,
} from './train-capacity.util';
import { WagonType } from '../wagon-types/entities/wagon-type.entity';
import { ClearanceMilestoneService } from '../contracts/clearance-milestone.service';
import { BookingSplitService } from './booking-split.service';
import { BookingWindowGateway } from './booking-window.gateway';
import {
MAX_TEU_SLOTS_PER_WAGON,
containerWagonsForLines,
} from './wagon-plan.util';
import {
Capacity,
CorridorBudget,
CorridorLeg,
OverageTolerance,
stopYardsFor,
} from './corridor-capacity.util';
export type { Capacity } from './corridor-capacity.util';
/**
* A train's fill limits: the base caps the corridor budget spends from, plus
* the locomotive overage tolerance spendable only on whole-booking admission.
*/
type TrainLimits = { base: Capacity; tolerance: OverageTolerance };
/**
* Result of the export whole-booking single-train space check. `scheduleId`
* is the earliest fillable train that carries the whole booking, or null when
* none can — then `bestAvailable` reports the largest single-train leftover
* in the booking's own units and `fullMessage` is the customer-facing copy.
*/
export interface ExportSpaceReport {
scheduleId: string | null;
trainsForDay: boolean;
corridorMatched: boolean;
need: Capacity;
bestAvailable: { wagons: number; cargoTons: number } | null;
fullMessage: string | null;
}
/** A day-level pool key: all trains on this route departing on this EAT day. */
interface RouteDayGroup {
originYardId: string;
destinationYardId: string;
/** EAT calendar day, `yyyy-MM-dd`. */
day: string;
}
/** One wagon type's footprint: its length on the train, the tare it adds to the
* locomotive's gross load, and the payload it carries. */
type PerWagonDims = { lengthMeters: number; tareWeightTons: number; capacityTons: number };
/**
* Wagon dimensions used to size a booking's capacity draw. `byWagonTypeId` holds
* every wagon type so a booking is measured on the type its cargo/container type
* actually rides (the same FK resolution allocation uses); `container`/`bulk` are
* representative fallbacks for bookings whose type has no wagon type configured.
*/
type WagonDims = {
container: PerWagonDims;
bulk: PerWagonDims;
byWagonTypeId: Map<string, PerWagonDims>;
};
export type BatchBoardBookingState =
| "ALLOCATED"
| "SELECTED_FOR_BATCH"
| "READY"
| "WAITING"
| "PENDING_CONTRACT"
| "EXPIRED";
export interface BatchBoardBooking {
id: string;
reference: string;
company: string;
isGovernment: boolean;
wagons: number;
weightTons: number;
lengthMeters: number;
paymentDeadline: string | null;
state: BatchBoardBookingState;
/** Rule-engine priority score used to rank the batch (higher = boards first). */
priorityScore: number;
/** CONTAINER | BULK — for the priority-tracking visuals. */
freightType: string | null;
}
export type BookingAllocationStatus =
| "NOT_ATTEMPTED"
| "ASSIGNED"
| "DEFERRED"
| "FAILED";
export interface BatchBoardBookingDetail extends BatchBoardBooking {
fullyExecutedAt: string | null;
selectedForBatchAt: string | null;
allocationStatus: BookingAllocationStatus;
allocationIssue: string | null;
/** Set when this booking shares a wagon with a consolidation partner. */
consolidationPartnerId: string | null;
consolidationPartnerRef: string | null;
}
export interface BatchWindowGroup {
key: string;
label: string;
/** EAT calendar day as ISO `YYYY-MM-DD` (empty for the pending-contract bucket). */
date: string;
/** Human label for the day, e.g. `Thu, 05 Jun` (empty for pending-contract). */
dateLabel: string;
start: string;
end: string;
counts: {
allocated: number;
selectedForBatch: number;
ready: number;
waiting: number;
expired: number;
pendingContract: number;
};
bookings: BatchBoardBookingDetail[];
}
export interface BatchBoardScheduleDetail {
scheduleId: string;
/** Human-facing schedule reference (S-YYYY-NNNNN). */
scheduleReference: string | null;
trainNumber: string | null;
routeName: string | null;
origin: string | null;
destination: string | null;
scheduleDate: string | null;
status: string;
bookingWindowStatus: string;
direction: string | null;
windowPhase: string | null;
windowOpensAt: string | null;
windowClosesAt: string | null;
docReviewEndsAt: string | null;
paymentPhaseEndsAt: string | null;
bookingCycleNo: number;
/** Built train (Train Builder) behind this departure, when scheduled by train. */
train: BatchBoardSchedule["train"];
locomotive: BatchBoardSchedule["locomotive"];
capacity: BatchBoardSchedule["capacity"];
counts: BatchBoardSchedule["counts"];
windows: BatchWindowGroup[];
pendingContract: BatchWindowGroup;
allocationViolations: string[];
}
export interface BatchBoardSchedule {
scheduleId: string;
/** Human-facing schedule reference (S-YYYY-NNNNN). */
scheduleReference: string | null;
trainNumber: string | null;
routeName: string | null;
origin: string | null;
destination: string | null;
scheduleDate: string | null;
createdAt: string | null;
status: string;
bookingWindowStatus: string;
direction: string | null;
windowPhase: string | null;
windowOpensAt: string | null;
windowClosesAt: string | null;
docReviewEndsAt: string | null;
paymentPhaseEndsAt: string | null;
bookingCycleNo: number;
/** Built train (Train Builder) behind this departure, when scheduled by train. */
train: {
id: string;
code: string;
trainName: string | null;
} | null;
locomotive: {
code: string;
name: string | null;
maxPullWeightTons: number;
maxTrainLengthMeters: number;
} | null;
capacity: {
/** Wagons on bookings already linked to the train (ALLOCATED only). */
allocatedWagons: number;
/** Train length used by allocated bookings (from wagon-type dimensions). */
allocatedLengthMeters: number;
maxLengthMeters: number | null;
/** Weight committed on the train (allocated + selected-for-batch). */
usedWeightTons: number;
maxWeightTons: number | null;
/** Wagon-slot cap for the train (locomotive/wagon-type derived). */
maxWagons: number | null;
};
counts: {
allocated: number;
selectedForBatch: number;
ready: number;
waiting: number;
pendingContract: number;
expired: number;
};
bookings: BatchBoardBooking[];
}
/** Paginated batch-board list in the shared `{items, meta}` envelope — the API
* response wrapper already uses `data`, and the frontend's unwrap() strips one
* `data` level. */
export type BatchBoardListResponse = PaginatedResponse<BatchBoardSchedule>;
/**
* Demand-batching engine: every 3h (EAT) it ranks each OPEN schedule's ready pool
* by priority, greedily fills the train to capacity (skipping bookings that don't fit),
* reserves a 1h pay window for commercial customers (government allocated unpaid,
* preempting lower-priority commercial if needed), then settles each batch 1h later —
* allocating those who paid and expiring those who didn't, topping up from the waiting list.
* Capacity is bounded on three axes at once: wagon count (`schedule.maxWagons`), the
* locomotive's max pull weight, and its max train length (also capped by global rules).
*/
@Injectable()
export class BookingBatchService implements OnModuleInit {
private readonly logger = new Logger(BookingBatchService.name);
/**
* Serialises settle/top-up per schedule. The PAYMENT phase transition and the
* tick's overdue backstop both call settleDueReservations for the same schedule
* in the same second; without this they interleave and the top-up runs against a
* schedule whose phase has already been concluded.
*/
private readonly scheduleLocks = new Map<string, Promise<void>>();
constructor(
@InjectDataSource() private readonly dataSource: DataSource,
private readonly bookingsRepository: BookingsRepository,
private readonly trainSchedulesRepository: TrainSchedulesRepository,
private readonly trainScheduleBookingsRepository: TrainScheduleBookingsRepository,
private readonly notifier: BookingNotifierService,
private readonly scheduler: SchedulerRegistry,
private readonly trainSchedulingService: TrainSchedulingService,
private readonly billing: BillingService,
private readonly bookingWindowGateway: BookingWindowGateway,
@Inject(forwardRef(() => BookingPricingService))
private readonly pricingService: BookingPricingService,
@Optional() private readonly milestoneService?: ClearanceMilestoneService,
@Optional() private readonly splitService?: BookingSplitService,
) {}
/** On boot, reconcile OPEN route-days and re-arm settle timers. */
async onModuleInit(): Promise<void> {
const groups = await this.openRouteDayGroups();
for (const group of groups) {
try {
await this.processRouteDay(group);
} catch (err) {
this.logger.warn(
`Boot reconcile failed for ${this.groupLabel(group)}: ${(err as Error).message}`,
);
}
}
const reserved = await this.dataSource
.getRepository(Booking)
.createQueryBuilder("b")
.select("DISTINCT b.train_schedule_id", "scheduleId")
.where(`b.status IN ('SELECTED_FOR_BATCH', 'AWAITING_PAYMENT')`)
.andWhere("b.train_schedule_id IS NOT NULL")
.getRawMany<{ scheduleId: string }>();
for (const { scheduleId } of reserved) this.armSettle(scheduleId);
}
/**
* Fire-and-forget batch pipeline for the (route, day) a schedule belongs to
* (contract sign, payment). Day-level pooling distributes across all of that
* day's trains, so a single schedule id maps to its whole route-day group.
*/
enqueueScheduleProcessing(scheduleId: string): void {
void this.processRouteDayForSchedule(scheduleId).catch((err) =>
this.logger.error(
`processRouteDay for schedule ${scheduleId} failed: ${(err as Error).message}`,
),
);
}
/**
* Fire-and-forget batch pipeline for a (route, day) directly — used when a
* booking enters the pool without a target train yet (e.g. after the
* operations team accepts an operation request). The booking is already
* FULLY_EXECUTED with its scheduled_date set, so the day-level fill will pick
* it up; this just runs that fill immediately instead of waiting for the cron.
*/
enqueueRouteDayProcessing(
originYardId: string,
destinationYardId: string,
day: string,
): void {
void this.processRouteDay({ originYardId, destinationYardId, day }).catch(
(err) =>
this.logger.error(
`processRouteDay for ${originYardId}${destinationYardId} on ${day} failed: ${(err as Error).message}`,
),
);
}
/** Resolve a schedule's (route, day) group and run the day-level pipeline. */
private async processRouteDayForSchedule(scheduleId: string): Promise<void> {
const schedule = await this.trainSchedulesRepository.findById(scheduleId);
if (!schedule?.scheduledDepartureDate) return;
await this.processRouteDay({
originYardId: schedule.originStationId,
destinationYardId: schedule.destinationStationId,
day: eatDay(schedule.scheduledDepartureDate),
});
}
/**
* Day-level pipeline: distribute the (route, day) pool across all its trains,
* then settle / reconcile / assign wagons per schedule (those steps stay
* schedule-scoped — only the fill is day-level).
*/
async processRouteDay(group: RouteDayGroup): Promise<void> {
this.logger.log(
`[BATCH] processRouteDay START ${group.originYardId}->${group.destinationYardId} ${group.day}`,
);
const scheduleIds = await this.fillRouteDay(
group.originYardId,
group.destinationYardId,
group.day,
);
for (const scheduleId of scheduleIds) {
await this.settleDueReservations(scheduleId);
await this.reconcilePaidUnlinked(scheduleId);
await this.trainSchedulingService.tryAutoWagonAllocation(scheduleId);
}
}
/** Fill pool, settle due reservations, link orphaned PAID, then assign wagons. */
async processSchedule(scheduleId: string): Promise<void> {
await this.fillSchedule(scheduleId);
await this.settleDueReservations(scheduleId);
await this.reconcilePaidUnlinked(scheduleId);
await this.trainSchedulingService.tryAutoWagonAllocation(scheduleId);
}
/**
* Distinct (origin, destination, EAT day) groups across LEGACY OPEN schedules —
* schedules with a `windowPhase` are driven exclusively by the window engine
* (BookingWindowService), never by the periodic legacy fill.
*/
private async openRouteDayGroups(): Promise<RouteDayGroup[]> {
const open = (
await this.trainSchedulesRepository.findAll({
where: [
{ bookingWindowStatus: "OPEN", status: TrainScheduleStatusEnum.Draft },
{ bookingWindowStatus: "OPEN", status: TrainScheduleStatusEnum.Scheduled },
],
})
).filter((s) => s.windowPhase == null);
const groups = new Map<string, RouteDayGroup>();
for (const s of open) {
if (!s.scheduledDepartureDate) continue;
const day = eatDay(s.scheduledDepartureDate);
const key = `${s.originStationId}|${s.destinationStationId}|${day}`;
if (!groups.has(key)) {
groups.set(key, {
originYardId: s.originStationId,
destinationYardId: s.destinationStationId,
day,
});
}
}
return [...groups.values()];
}
private groupLabel(group: RouteDayGroup): string {
return `${group.originYardId}${group.destinationYardId} on ${group.day}`;
}
/**
* Idempotent: link a paid batch booking to its schedule and assign wagons.
* Handles SELECTED_FOR_BATCH, PAID-without-link, and PAID-already-linked cases.
*/
async ensurePaidBookingAllocated(bookingId: string): Promise<void> {
const booking = await this.dataSource.getRepository(Booking).findOne({
where: { id: bookingId },
relations: { company: true },
});
if (!booking) return;
if (!booking.trainScheduleId) {
// A paid booking with no train is money taken and nothing boarding —
// scream so staff pin it to a schedule manually (batch board / assign).
if (booking.paymentStatus === "PAID" || booking.status === "PAID") {
this.logger.error(
`PAID booking ${booking.reference ?? bookingId} has no train_schedule_id — ` +
`its reservation was likely expired before the payment landed. ` +
`Assign it to a schedule manually from the batch board.`,
);
}
return;
}
const isBatchPaid =
booking.status === "SELECTED_FOR_BATCH" ||
booking.status === "AWAITING_PAYMENT" ||
booking.status === "PAID" ||
booking.paymentStatus === "PAID";
if (!isBatchPaid) return;
if (
booking.status === "SELECTED_FOR_BATCH" ||
booking.status === "AWAITING_PAYMENT"
) {
await this.dataSource
.getRepository(Booking)
.update(bookingId, { paymentStatus: "PAID", status: "PAID" });
} else if (booking.paymentStatus !== "PAID") {
await this.dataSource
.getRepository(Booking)
.update(bookingId, { paymentStatus: "PAID" });
}
// Paying inside the window accepts an open partial offer — reduce the booking
// to the offered part before it boards (remainder returns to the contract cap).
if (this.splitService) {
await this.splitService.applySplit(bookingId);
}
const linked =
await this.trainScheduleBookingsRepository.existsForBooking(bookingId);
if (!linked) {
await this.allocate(booking.trainScheduleId, booking, "paid");
this.logger.log(
`Linked PAID booking ${booking.reference ?? bookingId} to schedule ${booking.trainScheduleId}`,
);
} else {
// Already linked at booking time (export FCFS: the customer books a
// specific train, so allocate() ran up front). allocate() is where the
// payment-settled tracking milestones are written, so on this branch we
// record them here — otherwise a paid, already-linked booking leaves
// FREIGHT_PAYMENT_SETTLED stuck PENDING and the clearance step never ticks.
void this.completeTrackingMilestones(bookingId, [
"WAGON_REQUESTED",
"FREIGHT_PAYMENT_PENDING",
"FREIGHT_PAYMENT_SETTLED",
]);
void this.markWagonAllocatedMilestone(bookingId);
}
const schedule = await this.trainSchedulesRepository.findByIdWithFullGraph(
booking.trainScheduleId,
);
if (schedule && (await this.isTrainFull(schedule))) {
await this.setWindow(booking.trainScheduleId, "FULL");
}
const result = await this.trainSchedulingService.tryAutoWagonAllocation(
booking.trainScheduleId,
);
if (result.assignedBookingIds.length) {
this.logger.log(
`Wagon allocation for ${booking.reference ?? bookingId}: ${result.assignedBookingIds.length} assigned`,
);
}
if (
result.issues.some(
(i) => i.bookingId === bookingId && i.status !== "ASSIGNED",
)
) {
const issue = result.issues.find((i) => i.bookingId === bookingId);
this.logger.warn(
`Wagon allocation issue for ${booking.reference ?? bookingId}: ${issue?.issue ?? issue?.status}`,
);
}
this.notifyBoardChanged(booking.trainScheduleId, "booking_paid_allocated");
}
/** Customer paid — delegate to ensurePaidBookingAllocated. */
async confirmPaidAndAllocate(bookingId: string): Promise<void> {
await this.ensurePaidBookingAllocated(bookingId);
}
/** Open partial-capacity offer summary for booking detail payloads (null when none). */
async getOpenOfferSummary(bookingId: string): Promise<{
offeredWagons: number;
totalWagons: number;
offeredAmount: number;
paymentDeadline: Date;
} | null> {
if (!this.splitService) return null;
const offer = await this.splitService.findOpenOffer(bookingId);
if (!offer) return null;
return {
offeredWagons: offer.offeredWagons,
totalWagons: offer.totalWagons,
offeredAmount: Number(offer.offeredAmount),
paymentDeadline: offer.paymentDeadline,
};
}
// ---- export FCFS -----------------------------------------------------------
/**
* Whole-booking single-train space report for an EXPORT booking. Export
* bookings never split — the entire booking must ride ONE train, so the
* report scans every fillable export train on the booking's corridor/day
* (earliest first) for one whose remaining budget fits the whole need. When
* none fits, `bestAvailable` carries the largest single-train leftover
* converted into the booking's own units (base caps, no overage tolerance)
* so the customer can be told exactly how much he COULD book on that day.
*/
async exportSpaceReport(
booking: Booking,
need?: Capacity,
): Promise<ExportSpaceReport> {
if (!booking.scheduledDate) {
throw new BadRequestException('Booking has no scheduled date');
}
const day = eatDay(new Date(booking.scheduledDate));
// Corridor-aware: any train whose route carries the booking's origin
// strictly before its destination qualifies — a Dire→Djibouti booking may
// ride an Addis→…→Djibouti train. The leg check below (legOf) enforces the
// stop order, so we fetch the day's open trains without endpoint filters.
const corridor = await this.trainSchedulesRepository.findAll({
where: [
{ status: TrainScheduleStatusEnum.Draft },
{ status: TrainScheduleStatusEnum.Scheduled },
],
});
const candidates = corridor
.filter(
(s) =>
s.scheduledDepartureDate != null &&
eatDay(s.scheduledDepartureDate) === day &&
this.isFillable(s),
)
.sort(
(a, b) =>
a.scheduledDepartureDate.getTime() - b.scheduledDepartureDate.getTime(),
);
const wagonDims = await this.loadWagonDims();
const required = need ?? this.needFor(booking, wagonDims);
const dims = this.dimsFor(booking, wagonDims);
const report: ExportSpaceReport = {
scheduleId: null,
trainsForDay: candidates.length > 0,
corridorMatched: false,
need: required,
bestAvailable: null,
fullMessage: null,
};
for (const candidate of candidates) {
const schedule = await this.trainSchedulesRepository.findByIdWithFullGraph(
candidate.id,
);
const locomotive = schedule?.trainSet?.locomotive;
if (!schedule || !locomotive) continue;
const limits = await this.capacityLimits(locomotive);
const budget = await this.remainingBudget(schedule, limits, wagonDims);
const leg = budget.legOf(booking.originYardId, booking.destinationYardId);
if (!leg) continue; // this train's route doesn't carry the booking's leg
report.corridorMatched = true;
if (budget.fits(required, leg)) {
// Earliest fitting train wins — no need to keep sizing leftovers.
report.scheduleId = schedule.id;
return report;
}
const available = this.bookableWithin(budget.remainingFor(leg), dims);
if (
!report.bestAvailable ||
available.cargoTons > report.bestAvailable.cargoTons ||
(available.cargoTons === report.bestAvailable.cargoTons &&
available.wagons > report.bestAvailable.wagons)
) {
report.bestAvailable = available;
}
}
report.fullMessage = this.exportFullMessage(booking, report);
return report;
}
/**
* Largest booking (in the requester's own wagon-type units) that a single
* train's leftover base capacity could still admit: bounded by free wagon
* slots, free train length, and the locomotive's remaining pull weight
* (gross — each wagon's tare eats into it before any cargo does).
*/
private bookableWithin(
remaining: Capacity,
dims: PerWagonDims,
): { wagons: number; cargoTons: number } {
const byLength =
dims.lengthMeters > 0
? Math.floor(Math.max(0, remaining.lengthMeters) / dims.lengthMeters)
: Math.floor(Math.max(0, remaining.wagons));
const maxWagons = Math.max(
0,
Math.min(Math.floor(Math.max(0, remaining.wagons)), byLength),
);
let bestTons = 0;
let usableWagons = 0;
for (let w = 1; w <= maxWagons; w++) {
if (w * dims.tareWeightTons > remaining.weightTons) break;
usableWagons = w;
const tons = Math.min(
w * dims.capacityTons,
remaining.weightTons - w * dims.tareWeightTons,
);
if (tons > bestTons) bestTons = tons;
}
return {
wagons: usableWagons,
cargoTons: Math.max(0, Math.floor(bestTons * 1000) / 1000),
};
}
/** Customer-facing "train is full" copy carrying the bookable leftover. */
private exportFullMessage(booking: Booking, report: ExportSpaceReport): string {
if (!report.trainsForDay || !report.corridorMatched) {
return 'No export train is accepting bookings for this day';
}
const best = report.bestAvailable;
const base =
'Not enough train space — an export booking must ride a single train whole, ' +
'and no open train on this day can carry it. ';
if (!best || best.wagons <= 0) {
return base + 'No capacity is left on this day — pick another shipment day.';
}
if (booking.freightType === 'BULK') {
return (
base +
`The largest remaining space is about ${best.cargoTons} tons ` +
`(${best.wagons} wagon${best.wagons === 1 ? '' : 's'}) — book up to that amount or pick another day.`
);
}
return (
base +
`The largest remaining space is ${best.wagons} wagon${best.wagons === 1 ? '' : 's'} ` +
`(up to ${best.wagons * 2} × 20ft or ${best.wagons} × 40ft, weight permitting) — ` +
'reduce the booking or pick another day.'
);
}
/**
* Export is first-come-first-serve: no window cycle, no priority, no batch.
* Pick the earliest open export train on the booking's corridor/day that still
* fits the booking. Throws ConflictException when every train is full — the
* staff accept fails and no more export bookings are taken.
*/
async pickExportSchedule(booking: Booking, need?: Capacity): Promise<string> {
const report = await this.exportSpaceReport(booking, need);
if (report.scheduleId) return report.scheduleId;
throw new ConflictException(
report.fullMessage ?? 'Train is full — no export capacity left for this day',
);
}
/**
* Accept an export booking into the FCFS flow. Solo bookings reserve immediately.
* A consolidated booking reserves as a pair only once BOTH partners are ready
* (FULLY_EXECUTED): the second partner's accept triggers the pair reservation
* against the combined shared-wagon need; the first partner's accept just waits.
* Throws ConflictException (before this booking is persisted-ready) when there is
* no export capacity for the day, so staff accept fails.
*/
async acceptExportBooking(booking: Booking): Promise<void> {
const partnerId = booking.consolidationPartnerId ?? null;
if (!partnerId) {
const scheduleId = await this.pickExportSchedule(booking);
await this.reserveOnExport([booking], scheduleId);
return;
}
const partner = await this.dataSource
.getRepository(Booking)
.findOne({
where: { id: partnerId },
relations: {
company: true,
bookingContainers: { containerType: true },
cargoType: true,
},
});
// Partner not yet accepted → this booking is now FULLY_EXECUTED and simply
// waits; the partner's later accept will reserve the pair.
if (!partner || partner.status !== 'FULLY_EXECUTED') {
return;
}
const wagonDims = await this.loadWagonDims();
const need = this.combinedNeed(booking, partner, wagonDims);
const scheduleId = await this.pickExportSchedule(booking, need);
await this.reserveOnExport([booking, partner], scheduleId);
}
/** Reserve one or two (consolidated) export bookings on a train and open pay windows. */
private async reserveOnExport(
bookings: Booking[],
scheduleId: string,
): Promise<void> {
for (const b of bookings) await this.reserve(b, scheduleId);
this.armSettle(scheduleId);
const schedule =
await this.trainSchedulesRepository.findByIdWithFullGraph(scheduleId);
if (schedule && (await this.isTrainFull(schedule))) {
await this.setWindow(scheduleId, 'FULL');
}
this.notifyBoardChanged(scheduleId, 'export_booking_accepted');
}
/** Link PAID bookings that have no train_schedule_bookings row (cron backstop). */
async reconcilePaidUnlinked(scheduleId: string): Promise<void> {
const unlinked =
await this.bookingsRepository.findPaidUnlinkedForSchedule(scheduleId);
for (const booking of unlinked) {
await this.allocate(scheduleId, booking, "paid");
this.logger.log(
`Reconciled PAID booking ${booking.reference ?? booking.id} → schedule ${scheduleId}`,
);
}
if (unlinked.length > 0) {
this.notifyBoardChanged(scheduleId, "paid_reconciled");
}
}
// ---- legacy fill entry point ----------------------------------------------
/**
* Legacy periodic fill for schedules without a window phase (DOMESTIC and
* pre-migration trains). Invoked by BookingWindowService's tick — the old
* standalone cron was replaced by the window engine.
*/
async runBatchFill(): Promise<void> {
const groups = await this.openRouteDayGroups();
this.logger.log(`Batch fill: ${groups.length} OPEN route-day group(s).`);
for (const group of groups) {
try {
await this.processRouteDay(group);
} catch (err) {
this.logger.error(
`Batch fill failed for ${this.groupLabel(group)}: ${(err as Error).message}`,
);
}
}
}
// ---- monitoring board -----------------------------------------------------
/**
* Read model for the batch monitoring page: every import schedule — including
* dispatched, arrived and cancelled history — with its locomotive, capacity
* usage and its bookings grouped by lifecycle state (allocated / awaiting
* payment / paid-waiting / pending contract / expired). Paginated and
* filterable; per-schedule booking summaries are only computed for the
* requested page.
*/
async getBatchBoard(
query: BatchBoardQueryDto = {},
): Promise<BatchBoardListResponse> {
// Board cards are heavy (per-schedule booking summaries), so the default
// page is smaller than the toolkit-wide 20.
const { page, pageSize, skip, take } = normalizePagination(query, {
defaultPageSize: 12,
});
// Status filter: any subset of the lifecycle. Omitted = all statuses, so
// arrived / cancelled / dispatched schedules stay visible as history.
const allowedStatuses = new Set<string>(BATCH_BOARD_STATUSES);
const statuses = (query.statuses ?? "")
.split(",")
.map((v) => v.trim().toUpperCase())
.filter((v) => allowedStatuses.has(v));
const dateRange = (from?: string, to?: string) => {
const f = from ? new Date(from) : null;
const t = to ? new Date(to) : null;
if (f && t) return Between(f, t);
if (f) return MoreThanOrEqual(f);
if (t) return LessThanOrEqual(t);
return undefined;
};
// Batch board is IMPORT-only: export is FCFS with no batch/priority calc,
// and domestic/legacy schedules run the legacy fill, not the window batch.
const base: FindOptionsWhere<TrainSchedule> = { direction: "IMPORT" };
if (statuses.length) base.status = In(statuses) as never;
if (query.bookingWindowStatus) {
base.bookingWindowStatus = query.bookingWindowStatus;
}
const departure = dateRange(query.departureFrom, query.departureTo);
if (departure) base.scheduledDepartureDate = departure as never;
const created = dateRange(query.createdFrom, query.createdTo);
if (created) base.createdAt = created as never;
// Search fans out across every human-recognizable label. Each OR variant
// repeats the base filters so the search never widens them.
const term = query.search?.trim();
let where: FindOptionsWhere<TrainSchedule> | FindOptionsWhere<TrainSchedule>[] =
base;
if (term) {
const like = ILike(`%${term}%`);
where = [
{ ...base, trainNumber: like as never },
{ ...base, originStation: { label: like } },
{ ...base, destinationStation: { label: like } },
{ ...base, route: { originYard: { label: like } } },
{ ...base, route: { destinationYard: { label: like } } },
{ ...base, trainSet: { locomotive: { code: like } } },
] as FindOptionsWhere<TrainSchedule>[];
}
const sortBy = query.sortBy ?? "createdAt";
const sortOrder = query.sortOrder ?? "DESC";
const [schedules, total] = await this.trainSchedulesRepository.findAndCount({
where,
relations: {
trainSet: { locomotive: true, train: true },
originStation: true,
destinationStation: true,
// Yards supply the route's display name for `routeName` below;
// milestones (with yards) give it the full corridor path.
route: { originYard: true, destinationYard: true, milestones: { yard: true } },
},
order: { [sortBy]: sortOrder } as never,
skip,
take,
});
const wagonDims = await this.loadWagonDims();
const linkRepo = this.dataSource.getRepository(TrainScheduleBooking);
const board: BatchBoardSchedule[] = [];
for (const s of schedules) {
const links = await linkRepo.find({ where: { trainScheduleId: s.id } });
const linkedIds = new Set(links.map((l) => l.bookingId));
const bookings = await this.bookingsRepository.findAllBySchedule(s.id);
const items: BatchBoardBooking[] = bookings.map((b) => {
const need = this.needFor(b, wagonDims);
return {
id: b.id,
reference: b.reference ?? b.id.slice(0, 8),
company: b.isGovernment
? (b.governmentInstitution ?? "Government")
: (b.company?.name ?? "—"),
isGovernment: Boolean(b.isGovernment),
wagons: need.wagons,
weightTons: need.weightTons,
lengthMeters: need.lengthMeters,
paymentDeadline: b.paymentDeadline
? b.paymentDeadline.toISOString()
: null,
state: this.boardState(b, linkedIds.has(b.id)),
priorityScore: Number(b.priorityScore ?? 0),
freightType: b.freightType ?? null,
};
});
board.push(this.buildScheduleSummary(s, items));
}
return { items: board, meta: buildPaginationMeta(total, page, pageSize) };
}
/** Schedule-level batch board with EAT 3h windows grouped by fullyExecutedAt. */
async getBatchBoardDetail(
scheduleId: string,
): Promise<BatchBoardScheduleDetail> {
const s =
await this.trainSchedulesRepository.findByIdWithFullGraph(scheduleId);
if (!s)
throw new NotFoundException(`Train schedule ${scheduleId} not found`);
// Arrived / cancelled schedules stay viewable — the board is also the
// historical record of what each train carried.
// Batch board is IMPORT-only (export is FCFS, no batch/priority calc).
if (s.direction !== "IMPORT") {
throw new BadRequestException(
"The batch board only covers import schedules",
);
}
const wagonDims = await this.loadWagonDims();
const linkRepo = this.dataSource.getRepository(TrainScheduleBooking);
const links = await linkRepo.find({ where: { trainScheduleId: s.id } });
const linkedIds = new Set(links.map((l) => l.bookingId));
const bookings = await this.bookingsRepository.findAllBySchedule(s.id);
let allocationPreview: Awaited<
ReturnType<TrainSchedulingService["previewAllocationForSchedule"]>
>;
try {
allocationPreview =
await this.trainSchedulingService.previewAllocationForSchedule(s.id);
} catch {
allocationPreview = {
assignedBookingIds: [],
deferred: [],
issues: [],
violations: [],
};
}
const allocationByBooking = new Map(
allocationPreview.issues.map((i) => [i.bookingId, i]),
);
// Resolve consolidation-partner references for the shared-wagon badge. Most
// partners are on this same schedule; look up any that aren't in one query.
const refById = new Map(
bookings.map((b) => [b.id, b.reference ?? b.id.slice(0, 8)]),
);
const missingPartnerIds = [
...new Set(
bookings
.map((b) => b.consolidationPartnerId)
.filter((id): id is string => Boolean(id) && !refById.has(id!)),
),
];
if (missingPartnerIds.length) {
const partners = await this.dataSource
.getRepository(Booking)
.find({ where: { id: In(missingPartnerIds) } });
for (const p of partners) {
refById.set(p.id, p.reference ?? p.id.slice(0, 8));
}
}
const items: BatchBoardBookingDetail[] = bookings.map((b) => {
const need = this.needFor(b, wagonDims);
const alloc = allocationByBooking.get(b.id);
return {
id: b.id,
reference: b.reference ?? b.id.slice(0, 8),
company: b.isGovernment
? (b.governmentInstitution ?? "Government")
: (b.company?.name ?? "—"),
isGovernment: Boolean(b.isGovernment),
wagons: need.wagons,
weightTons: need.weightTons,
lengthMeters: need.lengthMeters,
paymentDeadline: b.paymentDeadline
? b.paymentDeadline.toISOString()
: null,
state: this.boardState(b, linkedIds.has(b.id)),
priorityScore: Number(b.priorityScore ?? 0),
freightType: b.freightType ?? null,
fullyExecutedAt: b.fullyExecutedAt
? b.fullyExecutedAt.toISOString()
: null,
selectedForBatchAt: b.selectedForBatchAt
? b.selectedForBatchAt.toISOString()
: null,
allocationStatus: alloc?.status ?? "NOT_ATTEMPTED",
allocationIssue: alloc?.issue ?? null,
consolidationPartnerId: b.consolidationPartnerId ?? null,
consolidationPartnerRef: b.consolidationPartnerId
? (refById.get(b.consolidationPartnerId) ?? null)
: null,
};
});
const loco = s.trainSet?.locomotive ?? null;
// Display windows are the REAL booking-window cycles this schedule was FROZEN
// with at creation (import: opens at its stored window time, lasts its rule's
// duration, reopens per its rule's delay; export: single FCFS lead window) —
// NOT the live global config. A later global-rules edit only re-derives
// not-yet-open schedules (restampPendingWindows), so an already-open schedule
// must keep drawing from its own snapshot, anchored on its stored open time.
// Legacy rows with no snapshot fall back to the live config.
const liveCfg = await this.trainSchedulingService.getWindowConfig();
const num = (v: unknown, fallback: number) => {
const n = v == null ? NaN : Number(v);
return Number.isFinite(n) ? n : fallback;
};
const windowCfg = {
windowOpenHour: num(s.ruleWindowOpenHour, liveCfg.windowOpenHour),
windowCloseHour: num(s.ruleWindowCloseHour, liveCfg.windowCloseHour),
windowDurationHours: num(
s.ruleWindowDurationHours,
liveCfg.windowDurationHours,
),
reopenDelayMinutes: num(s.ruleReopenDelayMinutes, liveCfg.reopenDelayMinutes),
importWindowLeadDays: num(
s.ruleImportWindowLeadDays,
liveCfg.importWindowLeadDays,
),
exportBookingLeadHours: num(
s.ruleExportBookingLeadHours,
liveCfg.exportBookingLeadHours,
),
};
const departureDate = s.scheduledDepartureDate ?? new Date();
const windowBuckets = groupBookingsIntoBoardWindows(
items,
(item) => (item.fullyExecutedAt ? new Date(item.fullyExecutedAt) : null),
s.direction ?? null,
departureDate,
windowCfg,
undefined,
s.windowOpensAt ?? null,
);
const emptyCounts = () => ({
allocated: 0,
selectedForBatch: 0,
ready: 0,
waiting: 0,
expired: 0,
pendingContract: 0,
});
const countFor = (bookingsInWindow: BatchBoardBookingDetail[]) => {
const counts = emptyCounts();
for (const b of bookingsInWindow) {
if (b.state === "ALLOCATED") counts.allocated += 1;
else if (b.state === "SELECTED_FOR_BATCH") counts.selectedForBatch += 1;
else if (b.state === "READY") counts.ready += 1;
else if (b.state === "WAITING") counts.waiting += 1;
else if (b.state === "EXPIRED") counts.expired += 1;
else counts.pendingContract += 1;
}
return counts;
};
const windows: BatchWindowGroup[] = [];
for (const [key, bucket] of windowBuckets) {
if (key === "pending-contract" || !bucket.window) continue;
const w = bucket.window;
windows.push({
key: w.key,
label: w.label,
date: w.date,
dateLabel: w.dateLabel,
start: w.start.toISOString(),
end: w.end.toISOString(),
counts: countFor(bucket.items),
bookings: bucket.items,
});
}
windows.sort(
(a, b) => new Date(a.start).getTime() - new Date(b.start).getTime(),
);
const pendingBookings = windowBuckets.get("pending-contract")?.items ?? [];
return {
scheduleId: s.id,
scheduleReference: s.reference ?? null,
trainNumber: s.trainNumber ?? null,
routeName: s.route ? formatRouteLabel(s.route) : null,
origin: s.originStation?.label ?? s.originStation?.code ?? null,
destination:
s.destinationStation?.label ?? s.destinationStation?.code ?? null,
scheduleDate: s.scheduledDepartureDate
? s.scheduledDepartureDate.toISOString()
: null,
status: s.status,
bookingWindowStatus: s.bookingWindowStatus,
direction: s.direction ?? null,
windowPhase: s.windowPhase ?? null,
windowOpensAt: s.windowOpensAt ? s.windowOpensAt.toISOString() : null,
windowClosesAt: s.windowClosesAt ? s.windowClosesAt.toISOString() : null,
docReviewEndsAt: s.docReviewEndsAt ? s.docReviewEndsAt.toISOString() : null,
paymentPhaseEndsAt: s.paymentPhaseEndsAt
? s.paymentPhaseEndsAt.toISOString()
: null,
bookingCycleNo: s.bookingCycleNo ?? 0,
train: s.trainSet?.train
? {
id: s.trainSet.train.id,
code: s.trainSet.train.code,
trainName: s.trainSet.train.trainName ?? null,
}
: null,
locomotive: loco
? {
code: loco.code,
name: loco.name ?? null,
maxPullWeightTons: Number(loco.maxPullWeightTons),
maxTrainLengthMeters: Number(loco.maxTrainLengthMeters),
}
: null,
capacity: this.computeBoardCapacity(items, loco, s.maxWagons ?? null),
counts: {
allocated: items.filter((i) => i.state === "ALLOCATED").length,
selectedForBatch: items.filter((i) => i.state === "SELECTED_FOR_BATCH")
.length,
ready: items.filter((i) => i.state === "READY").length,
waiting: items.filter((i) => i.state === "WAITING").length,
pendingContract: items.filter((i) => i.state === "PENDING_CONTRACT")
.length,
expired: items.filter((i) => i.state === "EXPIRED").length,
},
windows,
pendingContract: {
key: "pending-contract",
label: "Pending contract",
date: "",
dateLabel: "",
start: "",
end: "",
counts: countFor(pendingBookings),
bookings: pendingBookings,
},
allocationViolations: allocationPreview.violations,
};
}
/** Run wagon-level allocation for all eligible linked bookings on a schedule. */
async runWagonAllocation(scheduleId: string) {
const result =
await this.trainSchedulingService.tryAutoWagonAllocation(scheduleId);
if (result.assignedBookingIds.length > 0) {
this.notifyBoardChanged(scheduleId, "wagon_allocation_run");
}
return result;
}
/**
* Board capacity figures. `usedWeightTons` is GROSS (each item's weight already
* includes the tare of the wagons it occupies), so the ceiling it is measured
* against must be the same one the fill loop spends from: the locomotive's own
* limits widened by its overage tolerance (global rule caps do not apply, same
* as {@link capacityLimits}). Reading the raw `loco.maxPullWeightTons` here
* showed staff a ceiling the batch engine did not use.
*/
private computeBoardCapacity(
items: Array<{
state: BatchBoardBookingState;
wagons: number;
weightTons: number;
lengthMeters: number;
}>,
loco: Locomotive | null,
maxWagons: number | null,
): BatchBoardSchedule["capacity"] {
const allocated = items.filter((i) => i.state === "ALLOCATED");
// Every booking still targeting this train holds gross weight — including
// PAID ones waiting for wagon allocation (WAITING) and post-dispatch
// catch-all states. Counting only ALLOCATED + SELECTED_FOR_BATCH zeroed the
// board's weight the moment customers paid. Only EXPIRED released its hold.
const committed = items.filter((i) => i.state !== "EXPIRED");
const caps = loco
? trainHardCaps({
maxPullWeightTons: Number(loco.maxPullWeightTons),
maxTrainLengthMeters: Number(loco.maxTrainLengthMeters),
overageToleranceTons: Number(loco.overageToleranceTons) || 0,
overageToleranceMeters: Number(loco.overageToleranceMeters) || 0,
})
: null;
const round2 = (value: number) => Math.round(value * 100) / 100;
return {
allocatedWagons: allocated.reduce((sum, i) => sum + i.wagons, 0),
allocatedLengthMeters: round2(
allocated.reduce((sum, i) => sum + i.lengthMeters, 0),
),
maxLengthMeters: caps ? caps.maxLengthMeters : null,
usedWeightTons: round2(committed.reduce((sum, i) => sum + i.weightTons, 0)),
maxWeightTons: caps ? caps.maxWeightTons : null,
maxWagons: maxWagons ?? null,
};
}
private buildScheduleSummary(
s: TrainSchedule,
items: BatchBoardBooking[],
): BatchBoardSchedule {
const loco = s.trainSet?.locomotive ?? null;
return {
scheduleId: s.id,
scheduleReference: s.reference ?? null,
trainNumber: s.trainNumber ?? null,
routeName: s.route ? formatRouteLabel(s.route) : null,
origin: s.originStation?.label ?? s.originStation?.code ?? null,
destination:
s.destinationStation?.label ?? s.destinationStation?.code ?? null,
scheduleDate: s.scheduledDepartureDate
? s.scheduledDepartureDate.toISOString()
: null,
createdAt: s.createdAt ? s.createdAt.toISOString() : null,
status: s.status,
bookingWindowStatus: s.bookingWindowStatus,
direction: s.direction ?? null,
windowPhase: s.windowPhase ?? null,
windowOpensAt: s.windowOpensAt ? s.windowOpensAt.toISOString() : null,
windowClosesAt: s.windowClosesAt ? s.windowClosesAt.toISOString() : null,
docReviewEndsAt: s.docReviewEndsAt ? s.docReviewEndsAt.toISOString() : null,
paymentPhaseEndsAt: s.paymentPhaseEndsAt
? s.paymentPhaseEndsAt.toISOString()
: null,
bookingCycleNo: s.bookingCycleNo ?? 0,
train: s.trainSet?.train
? {
id: s.trainSet.train.id,
code: s.trainSet.train.code,
trainName: s.trainSet.train.trainName ?? null,
}
: null,
locomotive: loco
? {
code: loco.code,
name: loco.name ?? null,
maxPullWeightTons: Number(loco.maxPullWeightTons),
maxTrainLengthMeters: Number(loco.maxTrainLengthMeters),
}
: null,
capacity: this.computeBoardCapacity(items, loco, s.maxWagons ?? null),
counts: {
allocated: items.filter((i) => i.state === "ALLOCATED").length,
selectedForBatch: items.filter((i) => i.state === "SELECTED_FOR_BATCH")
.length,
ready: items.filter((i) => i.state === "READY").length,
waiting: items.filter((i) => i.state === "WAITING").length,
pendingContract: items.filter((i) => i.state === "PENDING_CONTRACT")
.length,
expired: items.filter((i) => i.state === "EXPIRED").length,
},
bookings: items.slice(0, 3),
};
}
private boardState(
booking: Booking,
linked: boolean,
): BatchBoardBookingState {
if (linked) return "ALLOCATED";
if (
booking.status === "SELECTED_FOR_BATCH" ||
booking.status === "AWAITING_PAYMENT"
) {
return "SELECTED_FOR_BATCH";
}
if (booking.status === "EXPIRED") return "EXPIRED";
if (booking.status === "FULLY_EXECUTED" && booking.fullyExecutedAt)
return "READY";
if (booking.status === "PAID") return "WAITING";
return "PENDING_CONTRACT";
}
// ---- core fill ------------------------------------------------------------
/**
* Whether the batch engine may reserve/allocate onto this schedule right now.
* Legacy (no window phase): the customer-facing OPEN gate doubles as the fill gate.
* Import window cycle: the engine fills while the customer window is CLOSED —
* during DOC_REVIEW (early staff trigger) and PAYMENT (batch run + top-ups).
* Export: FCFS while the booking window is open.
*/
isFillable(schedule: TrainSchedule): boolean {
if (schedule.bookingWindowStatus === "FULL") return false;
if (!schedule.windowPhase) return schedule.bookingWindowStatus === "OPEN";
if (schedule.direction === "EXPORT") {
return schedule.windowPhase === "OPEN" && schedule.bookingWindowStatus === "OPEN";
}
return schedule.windowPhase === "DOC_REVIEW" || schedule.windowPhase === "PAYMENT";
}
/**
* Fill one schedule from its priority-ordered pool until full. Returns the
* number of commercial units it RESERVED this pass (0 for government-only or
* no-fit passes) so a top-up caller can extend the payment phase only when a
* fresh pay window actually opened.
*/
async fillSchedule(scheduleId: string): Promise<number> {
const schedule =
await this.trainSchedulesRepository.findByIdWithFullGraph(scheduleId);
if (!schedule || !this.isFillable(schedule)) return 0;
const locomotive = schedule.trainSet?.locomotive;
if (!schedule.trainSetId || !locomotive) {
this.logger.warn(
`Schedule ${scheduleId} has no locomotive/train set — skipped.`,
);
return 0;
}
const wagonDims = await this.loadWagonDims();
const limits = await this.capacityLimits(locomotive);
await this.syncScheduleMaxWagons(schedule, locomotive);
const budget = await this.remainingBudget(schedule, limits, wagonDims);
const minPerWagon = this.minPerWagonNeed(wagonDims);
if (budget.isExhausted(minPerWagon)) {
await this.setWindow(scheduleId, "FULL");
return 0;
}
const pool = await this.bookingsRepository.findBatchPool(scheduleId);
// Same bulk re-score as fillRouteDayInternal — the legacy per-schedule fill
// must rank bulk bookings by their wagon-derived priority too.
await this.recomputeBulkPriorities(pool, wagonDims);
this.resortPoolByPriority(pool);
const units = this.groupConsolidatedPool(pool);
let armed = false;
let preempted = false;
let reservedThisPass = 0;
let commercialReserved = 0;
// Batch fill trace: caps + pool at entry. Kept on debug level — invaluable when
// reservations trickle instead of landing in one pass (a reserve() throwing
// mid-loop, e.g. schema drift, or a mis-synced capacity cap).
this.logger.debug(
`[fillSchedule ${scheduleId}] limits=${JSON.stringify(limits)} ` +
`maxWagons=${schedule.maxWagons} remaining=${JSON.stringify(budget.maxRemaining())} ` +
`poolSize=${pool.length} units=${units.length}`,
);
for (const unit of units) {
const { primary: booking, partner } = unit;
const isPair = partner != null;
const need = isPair
? this.combinedNeed(booking, partner, wagonDims)
: this.needFor(booking, wagonDims);
const isGov = booking.isGovernment || (partner?.isGovernment ?? false);
// Consolidated partners always share one corridor, so the primary's leg
// stands for the pair.
const leg = budget.legForYards(booking.originYardId, booking.destinationYardId);
// Per-unit fit trace: which axis (wagons/weight/length) admits or rejects.
this.logger.debug(
`[fillSchedule ${scheduleId}] unit ${booking.reference}: need=${JSON.stringify(need)} ` +
`roomOnLeg=${JSON.stringify(budget.remainingFor(leg))} fits=${budget.fits(need, leg)}`,
);
if (!budget.fits(need, leg)) {
if (isGov) {
const freed = await this.preemptForGovernment(
scheduleId,
need,
leg,
budget,
wagonDims,
);
preempted = true;
if (!freed) continue; // still doesn't fit even after preempt
} else {
// Doesn't fit whole. A split-eligible import booking is offered the part
// that fits in the remaining room (top-up path splits the boundary
// booking, mirroring fillRouteDay); otherwise skip and try the next.
const cand: { id: string; budget: CorridorBudget; armed: boolean } = {
id: scheduleId,
budget,
armed,
};
if (await this.maybeOfferPartial(booking, isPair, [cand], need)) {
armed = cand.armed;
continue;
}
continue; // skip a unit that exceeds weight/length/wagons, try the next
}
}
// Isolate each unit so a throw in reserve/allocate (e.g. billing hiccup)
// can't abort the whole top-up pass and leave the rest to trickle in one
// per tick. Log + skip the failing unit, keep going.
try {
if (isGov) {
await this.allocate(scheduleId, booking, "gov");
if (partner) await this.allocate(scheduleId, partner, "gov");
} else {
await this.reserve(booking, scheduleId);
if (partner) await this.reserve(partner, scheduleId);
armed = true;
commercialReserved += 1;
}
budget.subtract(need, leg);
reservedThisPass += 1;
} catch (err) {
this.logger.error(
`[fillSchedule ${scheduleId}] reserve/allocate FAILED for ${booking.reference} ` +
`— skipping this unit, continuing: ${(err as Error).message}`,
);
continue;
}
if (budget.maxRemaining().wagons <= 0) break; // every leg exhausted — nothing more can board
}
this.logger.log(
`[fillSchedule ${scheduleId}] reserved ${reservedThisPass}/${units.length} unit(s) this pass`,
);
if (budget.isExhausted(minPerWagon)) await this.setWindow(scheduleId, "FULL");
if (armed) this.armSettle(scheduleId);
// One push per fill pass (never per booking) — only when rows changed.
if (reservedThisPass > 0 || armed || preempted) {
this.notifyBoardChanged(scheduleId, "batch_fill");
}
void this.triggerWagonAllocation(scheduleId);
return commercialReserved;
}
/**
* Distribute one (route, day) pool across ALL of that day's OPEN trains, by
* priority, filling each train (earliest departure first) until it's full and
* spilling overflow to the next. Government bookings that fit no train preempt
* lower-priority commercial; bookings that fit no train at all stay pending and
* trigger a staff `unplaced` warning. Returns the schedule ids that were touched
* (or that had remaining pool work) so the caller can settle them per-schedule.
*/
async fillRouteDay(
originYardId: string,
destinationYardId: string,
day: string,
): Promise<string[]> {
const { scheduleIds } = await this.fillRouteDayInternal(
originYardId,
destinationYardId,
day,
);
return scheduleIds;
}
/**
* Route-day top-up for a single schedule: re-run the DAY pool over the whole
* corridor the schedule belongs to, and report how many commercial units got a
* fresh pay window.
*
* `fillSchedule` cannot do this job. Its pool (`findBatchPool`) is keyed on
* `booking.train_schedule_id = :scheduleId`, but under day-level pooling a
* booking that has not been reserved yet has a NULL `train_schedule_id` — it is
* only pinned by `reserve()`. So the schedule-scoped top-up returned zero rows
* and the waiting list never boarded after an expiry freed capacity; bookings
* trickled in one per window cycle instead.
*/
private async topUpFill(scheduleId: string): Promise<number> {
const schedule = await this.trainSchedulesRepository.findById(scheduleId);
if (!schedule?.scheduledDepartureDate) return 0;
const { commercialReserved } = await this.fillRouteDayInternal(
schedule.originStationId,
schedule.destinationStationId,
eatDay(schedule.scheduledDepartureDate),
);
return commercialReserved;
}
private async fillRouteDayInternal(
originYardId: string,
destinationYardId: string,
day: string,
): Promise<{ scheduleIds: string[]; commercialReserved: number }> {
// The day's fillable schedules on this exact corridor, earliest first. Fillable
// covers legacy OPEN trains and window-cycle trains in DOC_REVIEW/PAYMENT —
// the batch must run while the customer window is closed.
const corridor = await this.trainSchedulesRepository.findAll({
where: [
{
originStationId: originYardId,
destinationStationId: destinationYardId,
status: TrainScheduleStatusEnum.Draft,
},
{
originStationId: originYardId,
destinationStationId: destinationYardId,
status: TrainScheduleStatusEnum.Scheduled,
},
],
});
const onDay = corridor
.filter(
(s) =>
s.scheduledDepartureDate != null &&
eatDay(s.scheduledDepartureDate) === day,
)
.sort(
(a, b) =>
a.scheduledDepartureDate.getTime() - b.scheduledDepartureDate.getTime(),
);
// A schedule flagged FULL is rejected by isFillable() before its budget is
// ever consulted. Re-derive that flag from live capacity first, so a train
// whose bookings all expired is not skipped forever with an empty consist.
for (const s of onDay) {
if (s.bookingWindowStatus === "FULL") {
await this.refreshWindowStatus(s.id);
const fresh = await this.trainSchedulesRepository.findById(s.id);
if (fresh) s.bookingWindowStatus = fresh.bookingWindowStatus;
}
}
const scheduleIds = onDay.filter((s) => this.isFillable(s)).map((s) => s.id);
if (scheduleIds.length === 0) {
return { scheduleIds: [], commercialReserved: 0 };
}
const wagonDims = await this.loadWagonDims();
// Live per-schedule corridor budget + arm/changed flags, in departure order.
const trains: Array<{
id: string;
budget: CorridorBudget;
armed: boolean;
changed: boolean;
}> = [];
for (const id of scheduleIds) {
const schedule =
await this.trainSchedulesRepository.findByIdWithFullGraph(id);
const locomotive = schedule?.trainSet?.locomotive;
if (!schedule || !schedule.trainSetId || !locomotive) {
this.logger.warn(
`Schedule ${id} has no locomotive/train set — skipped.`,
);
continue;
}
const limits = await this.capacityLimits(locomotive);
await this.syncScheduleMaxWagons(schedule, locomotive);
const budget = await this.remainingBudget(schedule, limits, wagonDims);
trains.push({ id, budget, armed: false, changed: false });
}
if (trains.length === 0) return { scheduleIds, commercialReserved: 0 };
// The day pool covers every booking whose leg lies somewhere on one of the
// day's corridors — full-route AND sub-corridor (e.g. Dire→Djibouti on an
// Addis→Djibouti train). Which train actually takes a booking is decided
// by the per-train legOf check below.
const corridorYards = [...new Set(trains.flatMap((t) => t.budget.stops))];
const pool = await this.bookingsRepository.findBatchPoolByCorridorDay(
corridorYards,
day,
);
// BULK bookings only get their real (wagon-derived) priority score now, at
// batch time — stamp it and re-rank before the fill consumes the pool.
await this.recomputeBulkPriorities(pool, wagonDims);
this.resortPoolByPriority(pool);
// Consolidated partners collapse into one atomic unit (both-or-neither); a
// consolidated booking whose partner isn't ready this cycle is skipped.
const units = this.groupConsolidatedPool(pool);
// Batch fill trace: each train's caps + the day pool size at entry.
this.logger.debug(
`[fillRouteDay ${originYardId}->${destinationYardId} ${day}] ` +
`trains=${trains.map((t) => `${t.id}:${JSON.stringify(t.budget.maxRemaining())}`).join(",")} ` +
`poolSize=${pool.length} units=${units.length}`,
);
let reservedThisPass = 0;
let commercialReserved = 0;
for (const unit of units) {
const { primary: booking, partner } = unit;
const isPair = partner != null;
const need = isPair
? this.combinedNeed(booking, partner, wagonDims)
: this.needFor(booking, wagonDims);
const isGov = booking.isGovernment || (partner?.isGovernment ?? false);
const legOn = (t: { budget: CorridorBudget }): CorridorLeg | null =>
t.budget.legOf(booking.originYardId, booking.destinationYardId);
// First train (earliest departure) whose corridor carries this booking's
// leg and still fits it as-is.
let target = trains.find((t) => {
const leg = legOn(t);
return leg != null && t.budget.fits(need, leg);
});
// Per-unit trace: chosen train + each train's remaining room on this leg.
this.logger.debug(
`[fillRouteDay] unit ${booking.reference}: need=${JSON.stringify(need)} ` +
`targetTrain=${target?.id ?? "none"} ` +
`rooms=${trains
.map((t) => {
const leg = legOn(t);
return leg ? `${t.id}:${JSON.stringify(t.budget.remainingFor(leg))}` : `${t.id}:offleg`;
})
.join(",")}`,
);
if (!target && isGov) {
// Government fits nowhere on its own — try to preempt commercial
// on each corridor-matching train (earliest first) until one frees room.
for (const t of trains) {
const leg = legOn(t);
if (!leg) continue;
const freed = await this.preemptForGovernment(
t.id,
need,
leg,
t.budget,
wagonDims,
);
// Preempt may have displaced (expired) victims even when the need
// still doesn't fit — the board must refresh either way.
t.changed = true;
if (freed) {
target = t;
break;
}
}
}
if (!target) {
// Fits no train whole. A split-eligible booking is offered the largest
// part that fits on the train with the most free wagons on its leg (this
// covers both "fits nowhere" and the boundary case where earlier bookings
// already consumed most of the room). Consolidated pairs / government /
// non-import never split — isSplitEligible guards that. Passing the live
// `trains` entries lets maybeOfferPartial mutate the chosen budget/armed.
const offered = await this.maybeOfferPartial(booking, isPair, trains, need);
if (offered) {
// A partial offer opens a real commercial pay window, same as reserve().
commercialReserved += 1;
reservedThisPass += 1;
continue;
}
// Stays in the pool, retried next batch/window cycle.
this.notifier.unplaced(booking, day);
if (partner) this.notifier.unplaced(partner, day);
continue;
}
// A throw here (e.g. a billing/invoice hiccup inside reserve) must NOT abort
// the whole pass — otherwise only the bookings before the failure get a pay
// window and the rest trickle in one-per-tick on later retries (the
// "selected one at a time / staggered" symptom). Isolate each unit: log +
// skip a failing one, keep reserving the others. The skipped unit stays in
// the pool and is retried next cycle.
try {
if (isGov) {
await this.allocate(target.id, booking, "gov");
if (partner) await this.allocate(target.id, partner, "gov");
} else {
await this.reserve(booking, target.id);
if (partner) await this.reserve(partner, target.id);
target.armed = true;
commercialReserved += 1;
}
target.budget.subtract(need, legOn(target)!);
target.changed = true;
reservedThisPass += 1;
} catch (err) {
this.logger.error(
`[fillRouteDay] reserve/allocate FAILED for ${booking.reference} on ${target.id} ` +
`— skipping this unit, continuing the batch: ${(err as Error).message}`,
);
}
}
this.logger.log(
`[fillRouteDay ${originYardId}->${destinationYardId} ${day}] reserved ${reservedThisPass}/${units.length} unit(s) this pass`,
);
const minPerWagon = this.minPerWagonNeed(wagonDims);
for (const t of trains) {
if (t.budget.isExhausted(minPerWagon)) await this.setWindow(t.id, "FULL");
if (t.armed) this.armSettle(t.id);
// One push per touched train per pass (never per booking). `armed` covers
// commercial reserves + partial offers; `changed` covers gov allocations
// and preemption.
if (t.armed || t.changed) this.notifyBoardChanged(t.id, "batch_fill");
void this.triggerWagonAllocation(t.id);
}
return { scheduleIds: trains.map((t) => t.id), commercialReserved };
}
/**
* A lone commercial IMPORT booking on a GENERAL or ONE_TIME contract may be
* offered a partial (split-on-payment). Consolidated pairs never split (both-or-
* neither shared wagon) and government bookings never split (they preempt).
*/
private isSplitEligible(booking: Booking, isPair: boolean): boolean {
return (
!isPair &&
!booking.isGovernment &&
booking.tradeDirection === "IMPORT" &&
(booking.contractKind === "GENERAL" || booking.contractKind === "ONE_TIME") &&
this.splitService != null
);
}
/**
* Offer the largest fitting part of a booking that does not fit any candidate
* train whole, on the train with the most free wagons on the booking's leg.
* Mutates the chosen candidate's budget + armed flag in place. Returns true when
* an offer was opened (caller should `continue` past this unit), false otherwise.
* Shared by fillRouteDay (multi-train) and fillSchedule (single train). The leg
* is computed per candidate from the booking's yards, so callers pass their live
* train entries and only leg-carrying trains are considered.
*/
private async maybeOfferPartial(
booking: Booking,
isPair: boolean,
candidates: Array<{ id: string; budget: CorridorBudget; armed: boolean }>,
need: Capacity,
): Promise<boolean> {
if (!this.isSplitEligible(booking, isPair)) return false;
const target = candidates
.map((c) => {
const leg = c.budget.legOf(booking.originYardId, booking.destinationYardId);
return leg ? { c, leg, room: c.budget.remainingFor(leg) } : null;
})
.filter((x): x is NonNullable<typeof x> => x != null && x.room.wagons >= 1)
.sort((a, b) => b.room.wagons - a.room.wagons)[0];
if (!target) return false;
const offered = await this.tryPartialOffer(
booking,
target.c.id,
target.room,
need,
);
if (!offered) return false;
target.c.budget.subtract(offered, target.leg);
target.c.armed = true;
return true;
}
/**
* Offer the largest fitting part of an over-capacity booking as a partial
* (split-on-payment). Returns the capacity the offer consumes, or null when no
* meaningful partial fits / an offer is already open.
*/
private async tryPartialOffer(
booking: Booking,
scheduleId: string,
budget: Capacity,
need: Capacity,
): Promise<Capacity | null> {
if (!this.splitService) return null;
// A consolidated booking is already half of a shared wagon — never split it.
if (booking.consolidationPartnerId) return null;
if (await this.splitService.findOpenOffer(booking.id)) return null;
const wagonDims = await this.loadWagonDims();
// The wagon-slot axis alone under-constrains the offer. On a weight- or
// length-limited train (slots to spare, but e.g. only 798T of pull weight
// left) sizing by slots either produced an offer the fits() check below
// rejected, or — when the free slots exceeded the booking's own wagon
// count — sizeOffer refused outright, so a bulk booking on a weight-bound
// train was never offered a split at all. Size across all three axes,
// measured on the booking's REAL wagon type — the same one allocation
// validates against. Bulk splits ride FULL wagons only: the offer never
// part-loads its last wagon.
const perWagon = this.dimsFor(booking, wagonDims);
const partial = sizePartialOfferWagons(budget, need.wagons, perWagon, {
fullWagonsOnly: booking.freightType === "BULK",
});
if (!partial) return null;
const sized = await this.splitService.sizeOffer(
booking,
partial.wagons,
need.wagons,
perWagon.capacityTons,
partial.maxCargoTons,
);
if (!sized) return null;
const offeredNeed: Capacity = {
wagons: sized.offeredWagons,
weightTons: bookingGrossWeightTons(
sized.offeredWeightTons,
sized.offeredWagons,
perWagon.tareWeightTons,
),
lengthMeters: sized.offeredWagons * perWagon.lengthMeters,
};
if (!this.fits(offeredNeed, budget)) return null;
const deadline = new Date(Date.now() + (await this.paymentWindowMs()));
await this.splitService.createOffer(booking, scheduleId, sized, deadline);
// Reserve like a normal batch selection, but the partial invoice + partial
// pay-now notification were already produced by createOffer.
await this.bookingsRepository.update(booking.id, {
trainScheduleId: scheduleId,
status: "SELECTED_FOR_BATCH",
selectedForBatchAt: new Date(),
paymentDeadline: deadline,
} as never);
booking.trainScheduleId = scheduleId;
return offeredNeed;
}
/**
* Settle a schedule's reserved bookings. `expireUnpaidUnknownDeadline` decides
* how to treat a reservation with no deadline (durable path: leave it; timeout
* path: expire it). Consolidated pairs settle atomically: both allocate only
* when both paid; if either partner expires, both expire (a half-paid shared
* wagon must not ship). Returns whether anything changed.
*/
private async settleReserved(
scheduleId: string,
expireUnpaidUnknownDeadline: boolean,
): Promise<boolean> {
const reserved =
await this.bookingsRepository.findReservedForSchedule(scheduleId);
const now = Date.now();
const byId = new Map(reserved.map((b) => [b.id, b]));
const done = new Set<string>();
let anySettled = false;
this.logger.debug(
`[settleReserved ${scheduleId}] ${reserved.length} reserved booking(s) to settle`,
);
const isPaid = (b: Booking) =>
b.paymentStatus === "PAID" || b.status === "PAID";
const isExpired = (b: Booking) =>
b.paymentDeadline
? b.paymentDeadline.getTime() <= now
: expireUnpaidUnknownDeadline;
for (const booking of reserved) {
if (done.has(booking.id)) continue;
const partner = booking.consolidationPartnerId
? (byId.get(booking.consolidationPartnerId) ?? null)
: null;
if (partner) {
done.add(booking.id);
done.add(partner.id);
// Both-or-neither: allocate the shared wagon only when both partners paid;
// if either lapsed, expire both so no half-paid wagon rides.
if (isPaid(booking) && isPaid(partner)) {
await this.allocate(scheduleId, booking, "paid");
await this.allocate(scheduleId, partner, "paid");
anySettled = true;
} else if (isExpired(booking) || isExpired(partner)) {
await this.expire(booking);
await this.expire(partner);
anySettled = true;
}
continue;
}
done.add(booking.id);
if (isPaid(booking)) {
await this.allocate(scheduleId, booking, "paid");
anySettled = true;
} else if (isExpired(booking)) {
await this.expire(booking);
anySettled = true;
}
}
return anySettled;
}
/**
* Durable settle: allocate paid / expire overdue reservations, then top up the
* freed capacity from the waiting list.
*
* Serialised per schedule. Two callers race here every time a payment phase
* ends: `advanceImport`'s PAYMENT branch and the tick's `settleOverdueReservations`
* backstop. Both read the same reserved rows in the same second, so without the
* lock the second caller re-settles rows the first is mid-way through expiring,
* and `concludeCycle` observes capacity that is neither pre- nor post-expiry.
*/
async settleDueReservations(scheduleId: string): Promise<void> {
await this.withScheduleLock(scheduleId, () =>
this.settleAndTopUp(scheduleId, false),
);
}
/**
* Settle, then keep promoting the waiting list until the train can take no more.
* Returns whether anything settled.
*
* One top-up pass is not enough: expiring an N-wagon booking can free room for
* several smaller ones, and reserving those can in turn leave room for the next
* size down. Loop until a pass reserves nothing, so the batch ends with the train
* as full as the pool allows — rather than leaving a booking stranded until the
* next window cycle.
*
* Each round that opens a fresh pay window pushes `paymentPhaseEndsAt` out, so
* `concludeCycle` cannot fire before the promoted customers' deadlines.
*/
private async settleAndTopUp(
scheduleId: string,
expireUnpaidUnknownDeadline: boolean,
): Promise<boolean> {
const anySettled = await this.settleReserved(
scheduleId,
expireUnpaidUnknownDeadline,
);
if (!anySettled) return false;
this.logger.log(
`[BATCH] settle changed state on ${scheduleId} — running top-up fill for the waiting list`,
);
// Bounded: every round either reserves at least one unit (shrinking the pool)
// or breaks. The cap is a backstop against a pathological reserve/expire cycle.
let promoted = 0;
for (let round = 0; round < 10; round += 1) {
const reservedThisRound = await this.topUpFill(scheduleId);
if (reservedThisRound <= 0) break;
promoted += reservedThisRound;
await this.extendPaymentPhaseForTopUp(scheduleId);
}
if (promoted > 0) {
this.logger.log(
`[BATCH] top-up promoted ${promoted} waiting booking(s) onto ${scheduleId} ` +
`— payment phase extended for them`,
);
}
// Emitted here (not in settleDueReservations/settleBatch, which both wrap
// this) so one settle produces one push, after every allocation/expiry/
// top-up extension for this schedule has been persisted.
this.notifyBoardChanged(scheduleId, "reservations_settled");
return true;
}
/**
* Run `fn` with exclusive access to `scheduleId`. Concurrent callers await the
* in-flight run rather than interleaving with it. Single-process only — a second
* API replica would need a row lock on the schedule instead.
*/
private async withScheduleLock<T>(
scheduleId: string,
fn: () => Promise<T>,
): Promise<T> {
const inFlight = this.scheduleLocks.get(scheduleId) ?? Promise.resolve();
// Chain onto the previous holder; swallow its rejection so one failure does
// not poison every later caller's lock.
const run = inFlight.catch(() => undefined).then(fn);
const gate = run.then(
() => undefined,
() => undefined,
);
this.scheduleLocks.set(scheduleId, gate);
try {
return await run;
} finally {
// Last one out clears the slot so the map does not grow without bound.
if (this.scheduleLocks.get(scheduleId) === gate) {
this.scheduleLocks.delete(scheduleId);
}
}
}
// ---- settle (1h after a batch) -------------------------------------------
/** Allocate paid reservations, expire the rest, then top up. */
async settleBatch(scheduleId: string): Promise<void> {
this.removeTimeout(scheduleId);
await this.withScheduleLock(scheduleId, () =>
this.settleAndTopUp(scheduleId, true),
);
void this.triggerWagonAllocation(scheduleId);
}
private triggerWagonAllocation(scheduleId: string): void {
void this.trainSchedulingService
.tryAutoWagonAllocation(scheduleId)
.catch((err) =>
this.logger.warn(
`Auto wagon allocation failed for ${scheduleId}: ${(err as Error).message}`,
),
);
}
/**
* Announce that a schedule's batch-board data changed so open boards refetch.
* Called AFTER the state change is persisted; a push failure only logs — it
* must never break the business transaction that triggered it.
*/
private notifyBoardChanged(scheduleId: string, reason: string): void {
try {
this.bookingWindowGateway.emitBatchChanged(scheduleId, reason);
} catch (err) {
this.logger.warn(
`Batch-board push (${reason}) failed for ${scheduleId}: ${(err as Error).message}`,
);
}
}
// ---- staff override actions ----------------------------------------------
/** Staff "mark paid" override → set PAID and allocate immediately (don't wait for settle). */
async markPaid(bookingId: string): Promise<void> {
const booking = await this.dataSource
.getRepository(Booking)
.findOne({ where: { id: bookingId } });
if (!booking) throw new NotFoundException(`Booking ${bookingId} not found`);
if (!booking.trainScheduleId) {
throw new BadRequestException(
"Booking has no target schedule to allocate to",
);
}
await this.dataSource
.getRepository(Booking)
.update(bookingId, { paymentStatus: "PAID" });
await this.allocate(booking.trainScheduleId, booking, "paid");
const schedule = await this.trainSchedulesRepository.findByIdWithFullGraph(
booking.trainScheduleId,
);
if (schedule && (await this.isTrainFull(schedule))) {
await this.setWindow(booking.trainScheduleId, "FULL");
}
void this.triggerWagonAllocation(booking.trainScheduleId!);
this.notifyBoardChanged(booking.trainScheduleId, "booking_marked_paid");
}
/**
* Re-point a booking to another OPEN same-route schedule (keeps approval/contract + priority).
* Used for EXPIRED or full-schedule bookings — no re-approval.
*/
async moveToSchedule(
bookingId: string,
newScheduleId: string,
): Promise<void> {
const booking = await this.dataSource
.getRepository(Booking)
.findOne({ where: { id: bookingId } });
if (!booking) throw new NotFoundException(`Booking ${bookingId} not found`);
const schedule = await this.dataSource
.getRepository(TrainSchedule)
.findOne({ where: { id: newScheduleId } });
if (!schedule)
throw new NotFoundException(`Train schedule ${newScheduleId} not found`);
if (schedule.bookingWindowStatus !== "OPEN") {
throw new BadRequestException(
"Target schedule is not accepting bookings",
);
}
const stops = await this.stopsForSchedule(schedule);
const fromIdx = stops.indexOf(booking.originYardId);
const toIdx = stops.indexOf(booking.destinationYardId);
if (fromIdx < 0 || toIdx < 0 || fromIdx >= toIdx) {
throw new BadRequestException(
"Target schedule is not on the booking route",
);
}
const sourceScheduleId = booking.trainScheduleId ?? null;
await this.dataSource.transaction(async (manager) => {
if (booking.trainScheduleId) {
await this.trainScheduleBookingsRepository.deleteByScheduleAndBooking(
booking.trainScheduleId,
bookingId,
manager,
);
}
const restoredStatus =
booking.status === "EXPIRED"
? booking.isGovernment
? "APPROVED"
: "FULLY_EXECUTED"
: booking.status;
await manager.getRepository(Booking).update(bookingId, {
trainScheduleId: newScheduleId,
status: restoredStatus,
schedulingStatus: "ELIGIBLE",
paymentDeadline: null,
selectedForBatchAt: null,
} as never);
});
// Both boards changed: the booking left the source train and joined the target.
if (sourceScheduleId && sourceScheduleId !== newScheduleId) {
this.notifyBoardChanged(sourceScheduleId, "booking_moved");
}
this.notifyBoardChanged(newScheduleId, "booking_moved");
}
/** Staff "expire" override → free a reservation now (booking becomes EXPIRED). */
async expireReservation(bookingId: string): Promise<void> {
const booking = await this.dataSource
.getRepository(Booking)
.findOne({ where: { id: bookingId } });
if (!booking) throw new NotFoundException(`Booking ${bookingId} not found`);
// Capture the train before expire() detaches the booking from it — the
// top-up has to run against the schedule whose wagons were just freed.
const freedScheduleId = booking.trainScheduleId;
await this.expire(booking);
if (freedScheduleId) {
const topUpReserved = await this.topUpFill(freedScheduleId);
if (topUpReserved > 0) {
await this.extendPaymentPhaseForTopUp(freedScheduleId);
}
// After the top-up + phase extension so one push carries the final state.
this.notifyBoardChanged(freedScheduleId, "reservation_expired");
}
}
// ---- intercity ride-along API ---------------------------------------------
/**
* Remaining corridor capacity budget (per-edge wagons / weight / length) for
* a schedule, and the per-booking need calculator — exposed for the intercity
* accept flow, which reserves ride-along bookings onto import/export trains
* outside the batch engine. Segment-based: an intercity booking fits whenever
* ITS leg has room, even if the train is full on other legs.
*/
async intercityCapacity(scheduleId: string): Promise<{
budget: CorridorBudget;
needFor: (booking: Booking) => Capacity;
} | null> {
const schedule =
await this.trainSchedulesRepository.findByIdWithFullGraph(scheduleId);
const locomotive = schedule?.trainSet?.locomotive;
if (!schedule || !locomotive) return null;
const wagonDims = await this.loadWagonDims();
const limits = await this.capacityLimits(locomotive);
const budget = await this.remainingBudget(schedule, limits, wagonDims);
return { budget, needFor: (booking) => this.needFor(booking, wagonDims) };
}
/**
* Accept an intercity booking onto the given train. Commercial bookings get
* the same pay-window lifecycle as a batch reservation (deadline, invoice
* due-date sync, pay-now notify, settle on the window tick), so payment →
* allocation needs no special path. Government bookings allocate directly.
*/
async acceptIntercity(booking: Booking, scheduleId: string): Promise<void> {
if (booking.isGovernment) {
await this.dataSource
.getRepository(Booking)
.update(booking.id, { trainScheduleId: scheduleId });
booking.trainScheduleId = scheduleId;
await this.allocate(scheduleId, booking, 'gov');
this.notifyBoardChanged(scheduleId, 'intercity_accepted');
return;
}
await this.reserve(booking, scheduleId);
this.armSettle(scheduleId);
this.notifyBoardChanged(scheduleId, 'intercity_accepted');
}
// ---- mutations ------------------------------------------------------------
/**
* Reserve capacity for a commercial booking on a specific train and open its
* pay window. `scheduleId` is persisted so the settle/allocate lifecycle
* (settleDueReservations, settleBatch, ensurePaidBookingAllocated, markPaid),
* which is all keyed off `booking.trainScheduleId`, can find the train — with
* day-level pooling the booking arrives here with `trainScheduleId` still null,
* so the engine sets it as it picks the train.
*/
private async reserve(booking: Booking, scheduleId: string): Promise<void> {
// Idempotency guard: a booking already reserved (pay window open) or already
// paid on THIS schedule must never be re-reserved — that would fire a second
// `payNow` and reset its deadline, the "asked to pay again after paying"
// symptom. Read fresh state (the in-memory `booking` may be stale from the
// pooled query). Only bookings not yet committed to this train pass through.
const fresh = await this.dataSource
.getRepository(Booking)
.findOne({ where: { id: booking.id } });
if (
fresh &&
fresh.trainScheduleId === scheduleId &&
(fresh.status === "SELECTED_FOR_BATCH" ||
fresh.status === "AWAITING_PAYMENT" ||
fresh.status === "PAID" ||
fresh.paymentStatus === "PAID")
) {
this.logger.debug(
`[BATCH] reserve skipped for ${booking.reference} — already ` +
`${fresh.status}/${fresh.paymentStatus} on schedule ${scheduleId}`,
);
return;
}
const now = new Date();
const deadline = new Date(now.getTime() + (await this.paymentWindowMs()));
await this.bookingsRepository.update(booking.id, {
trainScheduleId: scheduleId,
status: "SELECTED_FOR_BATCH",
selectedForBatchAt: now,
paymentDeadline: deadline,
} as never);
booking.trainScheduleId = scheduleId;
// The invoice was generated at booking creation/approval, before this pay
// window opened — refresh its printed due date to the real deadline.
await this.billing.syncPayableDueDate(
Freight.InvoiceSource.Booking,
booking.id,
deadline,
"PREPAID",
);
await this.notifier.payNow(booking, deadline);
const reservedWagons = this.wagonsFor(booking, await this.loadWagonDims());
this.logger.log(
`[BATCH] RESERVED ${booking.reference} (${reservedWagons}w, ` +
`priority ${booking.priorityScore ?? 0}) on schedule ${scheduleId}` +
`pay by ${deadline.toISOString()}`,
);
// Customer tracking: a wagon slot is reserved and the freight pay window is
// open. Doc-trigger path — silent no-op for bookings without milestone rows.
void this.completeTrackingMilestones(booking.id, [
"WAGON_REQUESTED",
"FREIGHT_PAYMENT_PENDING",
]);
}
/** Allocate a booking to the schedule's train (creates the TrainScheduleBooking link). */
private async allocate(
scheduleId: string,
booking: Booking,
reason: "paid" | "gov",
): Promise<void> {
await this.dataSource.transaction(async (manager) => {
const exists =
await this.trainScheduleBookingsRepository.existsForBooking(
booking.id,
manager,
);
if (!exists) {
await this.trainScheduleBookingsRepository.createMany(
[{ trainScheduleId: scheduleId, bookingId: booking.id }],
manager,
);
}
await manager.getRepository(Booking).update(booking.id, {
status: reason === "paid" ? "PAID" : booking.status,
schedulingStatus: "SCHEDULED",
scheduledAt: new Date(),
paymentDeadline: null,
selectedForBatchAt: null,
} as never);
});
this.logger.log(
`[BATCH] ALLOCATED ${booking.reference} (${reason}) to train on schedule ${scheduleId}`,
);
this.notifier.secured(booking, reason, scheduleId);
void this.triggerWagonAllocation(scheduleId);
void this.markWagonAllocatedMilestone(booking.id);
// Customer tracking: freight payment settled (commercial pay-window path).
// Government allocations don't pay upfront — theirs stay pending.
if (reason === 'paid') {
void this.completeTrackingMilestones(booking.id, [
'WAGON_REQUESTED',
'FREIGHT_PAYMENT_PENDING',
'FREIGHT_PAYMENT_SETTLED',
]);
}
}
private async markWagonAllocatedMilestone(bookingId: string): Promise<void> {
if (!this.milestoneService) return;
try {
await this.milestoneService.completeForBooking(bookingId, 'WAGON_ALLOCATED');
} catch {
// Booking may have no milestone rows (non-contract path).
}
}
/**
* Complete customer-tracking milestones on lifecycle events via the
* doc-trigger path — a silent no-op for bookings without milestone rows
* (non-customs bookings). Never blocks the batch action.
*/
private async completeTrackingMilestones(
bookingId: string,
codes: string[],
): Promise<void> {
if (!this.milestoneService) return;
for (const code of codes) {
try {
await this.milestoneService.completeByDocTrigger({ bookingId }, code);
} catch (err) {
this.logger.warn(
`Milestone ${code} completion failed for booking ${bookingId}: ${(err as Error).message}`,
);
}
}
}
/**
* Expire an unpaid reservation and free its capacity. With day-level pooling we
* also clear `trainScheduleId` so the booking is no longer pinned to the train
* it failed to pay for — it's back in the day pool for staff to act on.
* `reason` picks the customer message: 'payment' (pay window lapsed) or
* 'no-capacity' (no train on the chosen day could take the booking).
*
* PAID GUARD: a booking whose payment has landed is never expired — money was
* taken, so it boards, even when the webhook arrived after the deadline or the
* settle read a stale row. It allocates onto the train it was selected for; if
* the wagon planner then finds no physical wagon, the booking stays linked and
* staff assign wagons manually. Consolidated bookings are exempt from the
* rescue: the shared wagon is both-or-neither, and settleReserved owns that
* pair decision.
*/
private async expire(
booking: Booking,
reason: "payment" | "no-capacity" = "payment",
): Promise<void> {
if (!booking.consolidationPartnerId) {
const fresh = await this.dataSource
.getRepository(Booking)
.findOne({ where: { id: booking.id }, relations: { company: true } });
const paid =
fresh != null &&
(fresh.paymentStatus === "PAID" || fresh.status === "PAID");
const paidScheduleId = fresh?.trainScheduleId ?? booking.trainScheduleId;
if (paid && paidScheduleId) {
this.logger.log(
`[BATCH] expire skipped for ${booking.reference} — payment already ` +
`landed; allocating on schedule ${paidScheduleId} instead`,
);
await this.allocate(paidScheduleId, fresh, "paid");
return;
}
}
const freedScheduleId = booking.trainScheduleId;
await this.bookingsRepository.update(booking.id, {
trainScheduleId: null,
status: "EXPIRED",
schedulingStatus: "ELIGIBLE",
paymentDeadline: null,
selectedForBatchAt: null,
} as never);
booking.trainScheduleId = null;
// The wagons this reservation held are back — a schedule parked at FULL
// because of it must reopen, or it can never be filled again.
if (freedScheduleId) await this.refreshWindowStatus(freedScheduleId);
// An unpaid partial offer dies with the reservation — the booking stays whole.
if (this.splitService) {
await this.splitService.expireOpenOffer(booking.id);
}
// Pay window closed before settlement → expire the booking's open invoice too
// (emits `booking.invoice.expired`). Domain owns the reaction; billing stays
// source-agnostic.
await this.billing.expirePayable(Freight.InvoiceSource.Booking, booking.id, "PREPAID");
if (reason === "no-capacity") {
this.notifier.expiredNoCapacity(booking);
} else {
this.notifier.expired(booking);
}
this.logger.log(
`[BATCH] EXPIRED ${booking.reference}` +
(reason === "no-capacity"
? "no train on its day had capacity left"
: "payment window passed; freed its wagons back to the pool for top-up"),
);
}
/**
* End-of-day sweep: once a schedule's window cycle concludes and NO other
* train on the same route-day can still run a cycle, the waiting pool for
* that day is dead — a FULLY_EXECUTED booking left in it would wait forever.
* Expire every leftover commercial booking and tell the customers to rebook
* another day. Government bookings are never auto-expired (they preempt).
* Returns how many bookings were expired.
*/
async expireLeftoverDayPool(scheduleId: string): Promise<number> {
const schedule = await this.trainSchedulesRepository.findById(scheduleId);
if (!schedule?.scheduledDepartureDate) return 0;
const day = eatDay(schedule.scheduledDepartureDate);
const group: RouteDayGroup = {
originYardId: schedule.originStationId,
destinationYardId: schedule.destinationStationId,
day,
};
// Another train on this route-day that can still take bookings keeps the
// pool alive — when IT concludes, its own sweep runs this check again.
const siblings = await this.trainSchedulesRepository.findAll({
where: [
{
originStationId: group.originYardId,
destinationStationId: group.destinationYardId,
status: TrainScheduleStatusEnum.Draft,
},
{
originStationId: group.originYardId,
destinationStationId: group.destinationYardId,
status: TrainScheduleStatusEnum.Scheduled,
},
],
});
const anotherTrainStillOpen = siblings.some(
(s) =>
s.id !== schedule.id &&
s.scheduledDepartureDate != null &&
eatDay(s.scheduledDepartureDate) === day &&
s.windowPhase !== "DONE" &&
s.bookingWindowStatus !== "FULL",
);
if (anotherTrainStillOpen) return 0;
const corridorYards = await this.corridorYardsForRouteDay(group);
const pool = corridorYards.length
? await this.bookingsRepository.findBatchPoolByCorridorDay(corridorYards, day)
: await this.bookingsRepository.findBatchPoolByRouteDay(
group.originYardId,
group.destinationYardId,
day,
);
const leftovers = pool.filter((b) => !b.isGovernment);
// Capture pinned schedules BEFORE expire() clears trainScheduleId, so each
// touched board gets exactly one push at the end of the sweep.
const touchedScheduleIds = new Set<string>();
if (leftovers.length) touchedScheduleIds.add(scheduleId);
for (const booking of leftovers) {
if (booking.trainScheduleId) touchedScheduleIds.add(booking.trainScheduleId);
await this.expire(booking, "no-capacity");
}
if (leftovers.length) {
this.logger.log(
`[BATCH] ${this.groupLabel(group)}: no train left with capacity — ` +
`expired ${leftovers.length} waiting booking(s)`,
);
}
for (const id of touchedScheduleIds) {
this.notifyBoardChanged(id, "day_pool_expired");
}
return leftovers.length;
}
/**
* Union of stop yards across the day's fillable schedules on this corridor —
* the same pool scope fillRouteDay uses, so full-route AND sub-corridor bookings
* are covered. Empty when no fillable schedule exists for the group.
*/
private async corridorYardsForRouteDay(
group: RouteDayGroup,
): Promise<string[]> {
const corridor = await this.trainSchedulesRepository.findAll({
where: [
{
originStationId: group.originYardId,
destinationStationId: group.destinationYardId,
status: TrainScheduleStatusEnum.Draft,
},
{
originStationId: group.originYardId,
destinationStationId: group.destinationYardId,
status: TrainScheduleStatusEnum.Scheduled,
},
],
});
const yards = new Set<string>();
for (const schedule of corridor) {
if (
schedule.scheduledDepartureDate == null ||
eatDay(schedule.scheduledDepartureDate) !== group.day
) {
continue;
}
for (const yardId of await this.stopsForSchedule(schedule)) {
yards.add(yardId);
}
}
return [...yards];
}
/**
* Sweep bookings on a route-day whose operation request staff did NOT accept by
* the time the window's document-review phase ends. They never reached
* FULLY_EXECUTED, so they never enter the batch — expire them (customer must
* rebook a new window). No reservation and no invoice exists yet at this stage,
* so this is a lighter expiry than `expire()`: just flip status + notify, and
* best-effort close any payable if one was issued early. Government/export are
* excluded by the query.
*/
async expireUnacceptedForRouteDay(group: RouteDayGroup): Promise<void> {
const corridorYards = await this.corridorYardsForRouteDay(group);
if (corridorYards.length === 0) return;
const unaccepted = await this.bookingsRepository.findUnacceptedForRouteDay(
corridorYards,
group.day,
);
if (unaccepted.length > 0) {
this.logger.log(
`[BATCH] doc-review end: expiring ${unaccepted.length} un-accepted booking(s) ` +
`on ${group.originYardId}->${group.destinationYardId} ${group.day}`,
);
}
// Only bookings pinned to a train show on a board — collect their schedules
// and push once per schedule after the sweep (most unaccepted rows are
// unpinned under day-level pooling, so this usually emits nothing).
const touchedScheduleIds = new Set<string>();
for (const booking of unaccepted) {
if (booking.trainScheduleId) touchedScheduleIds.add(booking.trainScheduleId);
await this.bookingsRepository.update(booking.id, {
status: "EXPIRED",
schedulingStatus: "ELIGIBLE",
// Free the shipment day so the customer can rebook a fresh window.
scheduledDate: null,
} as never);
// Close any payable issued before doc-review end (normally none — the invoice
// is created at ops-accept, which by definition has not happened here).
await this.billing
.expirePayable(Freight.InvoiceSource.Booking, booking.id, "PREPAID")
.catch(() => undefined);
this.notifier.expired(booking);
this.logger.log(
`[BATCH] EXPIRED (unaccepted) ${booking.reference}:${booking.id} at doc-review end`,
);
}
for (const id of touchedScheduleIds) {
this.notifyBoardChanged(id, "unaccepted_expired");
}
}
/**
* Free capacity for a government booking by displacing the lowest-priority commercial
* bookings (reserved first, then allocated — including PAID). Displaced → EXPIRED + notified.
* Only victims whose legs overlap the government booking's leg actually free useful
* room, so others are skipped. Mutates `budget`; returns whether the need now fits.
*/
private async preemptForGovernment(
scheduleId: string,
need: Capacity,
leg: CorridorLeg,
budget: CorridorBudget,
wagonDims: WagonDims,
): Promise<boolean> {
if (budget.fits(need, leg)) return true;
const reservedCommercial = (
await this.bookingsRepository.findReservedForSchedule(scheduleId)
).filter((b) => !b.isGovernment);
const allocatedCommercial =
await this.bookingsRepository.findAllocatedCommercialForSchedule(
scheduleId,
);
// lowest priority first; reserved are cheaper to free than allocated
const candidates = [...reservedCommercial, ...allocatedCommercial].sort(
(a, b) => (a.priorityScore ?? 0) - (b.priorityScore ?? 0),
);
for (const victim of candidates) {
if (budget.fits(need, leg)) break;
const victimLeg = budget.legForYards(
victim.originYardId,
victim.destinationYardId,
);
// Displacing a booking on a disjoint leg frees nothing the government
// booking can use — don't kill it for nothing.
const overlaps = victimLeg.fromEdge < leg.toEdge && leg.fromEdge < victimLeg.toEdge;
if (!overlaps) continue;
await this.dataSource.transaction(async (manager) => {
await this.trainScheduleBookingsRepository.deleteByScheduleAndBooking(
scheduleId,
victim.id,
manager,
);
await manager.getRepository(Booking).update(victim.id, {
status: "EXPIRED",
schedulingStatus: "ELIGIBLE",
paymentDeadline: null,
selectedForBatchAt: null,
} as never);
// Displaced → EXPIRED: close its open invoice too, so a dead booking
// can't still be paid (mirrors `expire()`; enlisted in this txn).
await this.billing.expirePayable(
Freight.InvoiceSource.Booking,
victim.id,
"PREPAID",
manager,
);
});
this.notifier.displaced(victim);
budget.add(this.needFor(victim, wagonDims), victimLeg);
// Displacing frees wagons the same way an expiry does — don't leave the
// schedule stuck at FULL.
await this.refreshWindowStatus(scheduleId);
}
return budget.fits(need, leg);
}
// ---- capacity helpers -----------------------------------------------------
/**
* Collapse consolidated partners into single pool entries so the fill treats a
* shared-wagon pair as one atomic unit (both-or-neither). For each pool entry:
* - no `consolidationPartnerId` → passes through as a lone booking.
* - consolidated + partner also in this pool → emitted ONCE (at the position of
* whichever partner ranks first) as a pair; the partner is not emitted again.
* - consolidated + partner NOT in this pool → dropped (can't ship half a wagon;
* it waits for the partner to become ready in a later cycle).
* The pool is already priority-ordered, so emitting the pair at the first-seen
* partner's slot ranks it by the stronger (max-priority) partner automatically.
*/
private groupConsolidatedPool(
pool: Booking[],
): Array<{ primary: Booking; partner: Booking | null }> {
const byId = new Map(pool.map((b) => [b.id, b]));
const emitted = new Set<string>();
const units: Array<{ primary: Booking; partner: Booking | null }> = [];
for (const booking of pool) {
if (emitted.has(booking.id)) continue;
const partnerId = booking.consolidationPartnerId ?? null;
if (!partnerId) {
emitted.add(booking.id);
units.push({ primary: booking, partner: null });
continue;
}
const partner = byId.get(partnerId) ?? null;
if (!partner) {
// Both-or-neither: partner not ready in this pool → skip the pair entirely.
emitted.add(booking.id);
continue;
}
emitted.add(booking.id);
emitted.add(partner.id);
units.push({ primary: booking, partner });
}
return units;
}
/**
* Combined capacity need of a consolidated pair sharing wagons. The whole point of
* consolidation is that the two partial 20ft counts pack onto the SAME wagons, so
* the shared wagon count is ceil((c1+c2)/2) — strictly fewer than summing the two
* independently-rounded-up needs (that is the capacity consolidation saves).
*/
private combinedNeed(
primary: Booking,
partner: Booking,
wagonDims: WagonDims,
): Capacity {
const containers = (b: Booking): number =>
(b.bookingContainers ?? []).reduce((sum, c) => sum + Number(c.quantity ?? 0), 0);
const totalContainers = containers(primary) + containers(partner);
const cargoTons =
Number(primary.cargoTotalWeightVgm ?? 0) + Number(partner.cargoTotalWeightVgm ?? 0);
// Consolidation shares TEU slots, never rated payload: the pair still needs
// enough wagons to carry its combined cargo, so the weight axis bounds the
// shared count exactly as it bounds an individual booking's. A pair shares
// wagons, so the primary's wagon type stands for both partners.
const dims = this.dimsFor(primary, wagonDims);
const capacityTons = dims.capacityTons;
const byWeight =
cargoTons > 0 && capacityTons > 0 ? Math.ceil(cargoTons / capacityTons) : 0;
const byLength =
totalContainers > 0
? Math.ceil(totalContainers / MAX_TEU_SLOTS_PER_WAGON)
: this.wagonsFor(primary, wagonDims) + this.wagonsFor(partner, wagonDims);
const sharedWagons = Math.max(byLength, byWeight);
return {
wagons: sharedWagons,
// Consolidation saves tare as well as slots: the pair rides `sharedWagons`
// wagons, so it is charged `sharedWagons` tares, not one per booking.
weightTons: bookingGrossWeightTons(
cargoTons,
sharedWagons,
dims.tareWeightTons,
),
lengthMeters: sharedWagons * dims.lengthMeters,
};
}
/**
* Stamp real priority scores on the pool's BULK bookings before the batch
* ranks it. Submit-time scoring runs with totalWagons = 0 for bulk (a bulk
* booking has no container lines to carry a wagon count), so every
* wagon-range priority config missed and bulk import bookings entered the
* batch at score 0 — they were never prioritized. Their wagon footprint is
* derivable from tonnage vs. live wagon capacity (wagonsFor), so the score
* is computed here — when doc review closes and the batch runs — and
* persisted so the priority board shows the same ranking. The pool arrives
* SQL-ordered by the old scores; the caller must re-sort after this.
*/
private async recomputeBulkPriorities(
pool: Booking[],
wagonDims: WagonDims,
): Promise<void> {
for (const booking of pool) {
if (booking.freightType !== 'BULK') continue;
try {
const wagons = this.wagonsFor(booking, wagonDims);
const score = await this.pricingService.computeSubmitPriorityScore(
booking,
wagons,
);
if (Number(booking.priorityScore ?? 0) === score) continue;
await this.dataSource
.getRepository(Booking)
.update(booking.id, { priorityScore: score });
booking.priorityScore = score;
} catch (err) {
// A failed recompute keeps the stored score — never blocks the batch.
this.logger.warn(
`Bulk priority recompute failed for ${booking.reference ?? booking.id}: ` +
`${(err as Error).message}`,
);
}
}
}
/** Restore the batch pool ordering (mirrors findBatchPool's ORDER BY) after scores changed. */
private resortPoolByPriority(pool: Booking[]): void {
pool.sort(
(a, b) =>
Number(b.isGovernment) - Number(a.isGovernment) ||
Number(b.priorityScore ?? 0) - Number(a.priorityScore ?? 0) ||
(a.fullyExecutedAt?.getTime() ?? Infinity) -
(b.fullyExecutedAt?.getTime() ?? Infinity) ||
a.createdAt.getTime() - b.createdAt.getTime(),
);
}
/**
* Wagons a booking occupies. Two axes bind independently and the booking needs
* enough wagons to satisfy BOTH, so the count is the larger of:
*
* weight — ceil(cargoTons / wagonType.capacityTons), the rated payload
* length — TEU geometry, two 20ft to a wagon (container bookings only)
*
* The weight axis was missing entirely. A BULK booking carries no container
* lines, so `containerWagonsForLines` returned 0 and every bulk booking
* collapsed to a single wagon no matter its tonnage — a 2590T fertilizer
* booking counted as 1 wagon, and `needFor` then charged 1 tare instead of 37.
* That under-reported the board and let the fill loop overbook the train.
*/
private wagonsFor(booking: Booking, wagonDims: WagonDims): number {
// Stored wagonsRequired is a candidate, never an early return: rows written
// while sumWagonsRequired hardcoded BULK to 1 wagon are still in the DB, and
// trusting them charged one tare for a whole bulk consist (a 700T booking on
// 70T wagons read 700 + 1 tare instead of 700 + 10 tares).
const stored =
booking.wagonsRequired && booking.wagonsRequired > 0
? Math.ceil(booking.wagonsRequired)
: 0;
// TEU-aware: two 20ft share one wagon (wagonsPerUnit = 0.5). The old fallback
// summed raw container QUANTITY, so 20×20ft counted as 20 wagons, not 10.
const byLength = containerWagonsForLines(booking.bookingContainers ?? []);
const capacityTons = this.dimsFor(booking, wagonDims).capacityTons;
const cargoTons = Number(booking.cargoTotalWeightVgm ?? 0);
const byWeight =
cargoTons > 0 && capacityTons > 0 ? Math.ceil(cargoTons / capacityTons) : 0;
return Math.max(DEFAULT_WAGONS_PER_BOOKING, stored, byLength, byWeight);
}
/**
* What one booking consumes along all three capacity axes.
*
* The weight axis is GROSS — cargo plus the tare of every wagon the booking
* occupies — because it is spent against the locomotive's pull limit, which
* governs the whole train and not just its payload. Charging cargo alone let a
* 37-wagon box-wagon train read 2590T when it really weighed 3522T.
*/
private needFor(booking: Booking, wagonDims: WagonDims): Capacity {
const wagons = this.wagonsFor(booking, wagonDims);
const dims = this.dimsFor(booking, wagonDims);
return {
wagons,
weightTons: bookingGrossWeightTons(
Number(booking.cargoTotalWeightVgm ?? 0),
wagons,
dims.tareWeightTons,
),
lengthMeters: wagons * dims.lengthMeters,
};
}
private fits(need: Capacity, budget: Capacity): boolean {
return (
need.wagons <= budget.wagons &&
need.weightTons <= budget.weightTons &&
need.lengthMeters <= budget.lengthMeters
);
}
/**
* Caps for a schedule's train: gross pull weight, train length, and the
* length-derived wagon slot count (never a fixed 53). Bookings spend against
* `base` via {@link needFor}, whose weight axis is gross. The locomotive's
* overage tolerance is returned separately — the corridor budget spends it
* only to admit a booking whole, never to size a split.
*
* Limits come from the LOCOMOTIVE ALONE — the global-rules weight/length
* caps deliberately do not apply here (a mis-set global row once capped
* every train at 14m and no export booking could board).
*/
private async capacityLimits(locomotive: Locomotive): Promise<TrainLimits> {
const wagonTypes = await this.loadWagonTypeDimensions();
const derived = deriveTrainCapacityFromLocomotive(
{
maxPullWeightTons: Number(locomotive.maxPullWeightTons),
maxTrainLengthMeters: Number(locomotive.maxTrainLengthMeters),
overageToleranceTons: Number(locomotive.overageToleranceTons) || 0,
overageToleranceMeters: Number(locomotive.overageToleranceMeters) || 0,
},
wagonTypes,
);
return {
base: {
wagons: derived.maxWagonSlots,
weightTons: derived.baseWeightTons,
lengthMeters: derived.baseLengthMeters,
},
tolerance: {
weightTons: derived.toleranceTons,
lengthMeters: derived.toleranceMeters,
},
};
}
/**
* Keep schedule.max_wagons aligned with the train's boarding limit: the
* locomotive's length-derived slot count. The physical wagons currently in
* the train set do NOT cap this — bookings are admitted on length/weight
* alone and yard staff attach the wagons manually before departure.
*/
private async syncScheduleMaxWagons(
schedule: TrainSchedule,
locomotive: Locomotive,
): Promise<void> {
const limits = await this.capacityLimits(locomotive);
const maxWagons = limits.base.wagons;
if ((schedule.maxWagons ?? 0) !== maxWagons) {
await this.dataSource
.getRepository(TrainSchedule)
.update(schedule.id, { maxWagons });
schedule.maxWagons = maxWagons;
}
}
/**
* Every active wagon type, so the slot count is derived from the shortest wagon
* the fleet can actually marshal rather than from an arbitrary two-code sample.
*/
private async loadWagonTypeDimensions(): Promise<WagonTypeDimensions[]> {
const types = await this.dataSource
.getRepository(WagonType)
.find({ where: { isActive: true } });
if (types.length) return types.map(wagonTypeDimensionsFromEntity);
return [
{
lengthMeters: DEFAULT_CONTAINER_WAGON_LENGTH_METERS,
capacityTons: 70,
tareWeightTons: DEFAULT_CONTAINER_WAGON_TARE_TONS,
},
{
lengthMeters: DEFAULT_BULK_WAGON_LENGTH_METERS,
capacityTons: 60,
tareWeightTons: DEFAULT_BULK_WAGON_TARE_TONS,
},
];
}
/**
* Every wagon type keyed by id (drives per-booking dims via the cargo/container
* type's wagon_type_id FK), plus representative fallbacks per freight type
* (NW5 flat for containers, CW3 gondola for bulk) for bookings whose type has
* no wagon type configured yet.
*/
private async loadWagonDims(): Promise<WagonDims> {
const types = await this.dataSource.getRepository(WagonType).find();
const byCode = new Map(
types.map((t) => [t.code, wagonTypeDimensionsFromEntity(t)]),
);
const byWagonTypeId = new Map(
types.map((t) => [t.id, wagonTypeDimensionsFromEntity(t)]),
);
const nw5 = byCode.get("NW5");
const cw3 = byCode.get("CW3");
// capacityTons divides a bulk booking's cargo, so a 0 or missing rated payload
// must fall back rather than yield an infinite wagon count.
const payload = (value: number | undefined, fallback: number): number =>
value && value > 0 ? value : fallback;
return {
container: {
lengthMeters: nw5?.lengthMeters ?? DEFAULT_CONTAINER_WAGON_LENGTH_METERS,
tareWeightTons: nw5?.tareWeightTons ?? DEFAULT_CONTAINER_WAGON_TARE_TONS,
capacityTons: payload(nw5?.capacityTons, DEFAULT_CONTAINER_WAGON_CAPACITY_TONS),
},
bulk: {
lengthMeters: cw3?.lengthMeters ?? DEFAULT_BULK_WAGON_LENGTH_METERS,
tareWeightTons: cw3?.tareWeightTons ?? DEFAULT_BULK_WAGON_TARE_TONS,
capacityTons: payload(cw3?.capacityTons, DEFAULT_BULK_WAGON_CAPACITY_TONS),
},
byWagonTypeId,
};
}
/**
* Dimensions of the wagon type THIS booking rides: bulk resolves through its
* cargo type's allowed wagon-type list, container through the first container
* line's — the same list resolution the scheduling planner applies when the
* paid booking is allocated. Board/fill math measured on a representative
* wagon while allocation validated the real one let a selected batch flunk
* the post-payment gross-weight check; sharing the resolution closes that
* gap. Uses the first configured type (the fill engine has no train context);
* falls back to the representative dims when the list or relation is absent.
*/
private dimsFor(booking: Booking, wagonDims: WagonDims): PerWagonDims {
const fallback =
booking.freightType === "BULK" ? wagonDims.bulk : wagonDims.container;
const wagonTypeId =
booking.freightType === "BULK"
? booking.cargoType?.wagonTypes?.[0]?.id
: (booking.bookingContainers ?? [])
.flatMap((line) => line.containerType?.wagonTypes ?? [])
.map((wagonType) => wagonType.id)
.find((id): id is string => Boolean(id));
const dims = wagonTypeId ? wagonDims.byWagonTypeId.get(wagonTypeId) : undefined;
if (!dims) return fallback;
return {
...dims,
capacityTons: dims.capacityTons > 0 ? dims.capacityTons : fallback.capacityTons,
};
}
/**
* Ordered stop yards of the schedule's route (origin → milestones →
* destination); the legacy two-stop pseudo-route when milestones are absent.
*/
private async stopsForSchedule(schedule: TrainSchedule): Promise<string[]> {
let milestoneYards: string[] | null = null;
if (schedule.routeId) {
const milestones = await this.dataSource
.getRepository(RouteMilestone)
.find({ where: { routeId: schedule.routeId }, order: { sequenceNo: 'ASC' } });
if (milestones.length >= 2) milestoneYards = milestones.map((m) => m.yardId);
}
return stopYardsFor(
milestoneYards,
schedule.originStationId,
schedule.destinationStationId,
);
}
/**
* Remaining capacity per corridor edge = hard caps minus what allocated +
* reserved bookings already use ON THEIR OWN LEGS. A booking riding only
* Dire→Djibouti leaves the Addis→Dire edges untouched.
*
* The wagon axis is the locomotive's length-derived slot count only — the
* physical wagons currently marshalled in the train set do NOT cap it.
* Bookings are admitted on length/weight capacity and yard staff attach
* the missing wagons manually before wagon assignment.
*/
private async remainingBudget(
schedule: TrainSchedule,
limits: TrainLimits,
wagonDims: WagonDims,
): Promise<CorridorBudget> {
const stops = await this.stopsForSchedule(schedule);
const budget = new CorridorBudget(stops, limits.base, limits.tolerance);
const allocated = (schedule.scheduleBookings ?? [])
.map((sb) => sb.booking)
.filter((b): b is Booking => Boolean(b));
const reserved = await this.bookingsRepository.findReservedForSchedule(
schedule.id,
);
for (const b of [...allocated, ...reserved]) {
budget.subtract(
this.needFor(b, wagonDims),
budget.legForYards(b.originYardId, b.destinationYardId),
);
}
return budget;
}
/**
* Wagon slots still boardable somewhere on the corridor (most-open edge).
* ≤ 0 means no leg can take another booking. Slot axis ONLY — the train-wide
* FULL signal is {@link isTrainFull}, which also closes weight/length-bound
* trains that still show free slots.
*/
private async remainingWagons(schedule: TrainSchedule): Promise<number> {
const wagonDims = await this.loadWagonDims();
const budget = await this.remainingBudget(
schedule,
{
base: {
wagons: schedule.maxWagons ?? 0,
weightTons: Number.POSITIVE_INFINITY,
lengthMeters: Number.POSITIVE_INFINITY,
},
tolerance: { weightTons: 0, lengthMeters: 0 },
},
wagonDims,
);
return budget.maxRemaining().wagons;
}
async setWindow(
scheduleId: string,
status: "OPEN" | "FULL" | "CLOSED",
): Promise<void> {
await this.dataSource
.getRepository(TrainSchedule)
.update(scheduleId, { bookingWindowStatus: status });
// Push the change (open / train full / closed) so portal home and GL cards
// flip in real time — FULL in particular happens outside the window tick
// (batch fill, staff mark-paid) and had no live signal before.
try {
const fresh = await this.trainSchedulesRepository.findById(scheduleId);
if (fresh) this.bookingWindowGateway.emitPhase(fresh);
} catch (err) {
this.logger.warn(
`Booking-window push failed for ${scheduleId}: ${(err as Error).message}`,
);
}
}
/**
* A reservation on this schedule still has time left to pay.
*
* The PAYMENT phase ends a hair BEFORE its own reservations do: `paymentPhaseEndsAt`
* is stamped when the phase starts, then `reserve()` gives each booking
* `now + paymentWindow` a few hundred milliseconds later, one booking at a time. So
* the first settle after the phase deadline finds every reservation still in date,
* expires nothing, reports `anySettled = false`, runs no top-up — and the caller
* concludes the cycle out from under customers who still had time to pay. The next
* tick then expires them with no cycle left to promote the waiting list into.
*
* Callers must not conclude the cycle while this returns true.
*/
async hasLiveReservations(scheduleId: string): Promise<boolean> {
const reserved =
await this.bookingsRepository.findReservedForSchedule(scheduleId);
const now = Date.now();
return reserved.some(
(b) =>
b.paymentStatus !== "PAID" &&
b.status !== "PAID" &&
b.paymentDeadline != null &&
b.paymentDeadline.getTime() > now,
);
}
/**
* FULL on ANY capacity axis: out of wagon slots, or out of pull weight /
* train length for even one more loaded wagon. The old slot-only check let
* a weight-bound train (PW2: weight binds at 37 wagons = 3522.4T of
* 3500+90T, slots bind at 44) cycle its booking window forever instead of
* finalizing — 7 phantom slots kept it "not full" while nothing could board.
*/
async isScheduleFull(scheduleId: string): Promise<boolean> {
const schedule =
await this.trainSchedulesRepository.findByIdWithFullGraph(scheduleId);
if (!schedule) return false;
return this.isTrainFull(schedule);
}
/** See {@link isScheduleFull} — same check for callers that already hold the full graph. */
private async isTrainFull(schedule: TrainSchedule): Promise<boolean> {
if ((await this.remainingWagons(schedule)) <= 0) return true;
const locomotive = schedule.trainSet?.locomotive;
if (!locomotive) return false; // no weight/length limits to bind against
const wagonDims = await this.loadWagonDims();
const limits = await this.capacityLimits(locomotive);
const budget = await this.remainingBudget(schedule, limits, wagonDims);
return budget.isExhausted(this.minPerWagonNeed(wagonDims));
}
/**
* Smallest gross weight / shortest length one more wagon could add: the
* lightest wagon type at its rated payload. Feeds CorridorBudget.isExhausted,
* so FULL is only declared when not even this wagon fits anywhere.
*/
private minPerWagonNeed(wagonDims: WagonDims): {
grossWeightTons: number;
lengthMeters: number;
} {
const all = [
wagonDims.container,
wagonDims.bulk,
...wagonDims.byWagonTypeId.values(),
];
return {
grossWeightTons: Math.min(
...all.map((d) => d.tareWeightTons + d.capacityTons),
),
lengthMeters: Math.min(...all.map((d) => d.lengthMeters)),
};
}
/**
* Re-derive `bookingWindowStatus` from live capacity after wagons were freed
* (a reservation expired, a booking was displaced, a link was removed).
*
* FULL used to be a one-way door: `isFillable()` rejects a FULL schedule before
* it ever looks at the budget, and the only writers of OPEN skip a FULL row. So
* a train that filled once and then lost every booking to expiry stayed FULL
* with all its wagons free — permanently unfillable, cycling PRE_WINDOW→PAYMENT
* forever while `concludeCycle` (which reads real capacity, not the flag) kept
* reopening it. Clearing FULL here is what lets the next batch actually run.
*
* Only the customer-facing OPEN phases may go back to OPEN; a schedule mid
* DOC_REVIEW/PAYMENT drops to CLOSED, which `isFillable()` still admits.
*/
async refreshWindowStatus(scheduleId: string): Promise<void> {
const schedule =
await this.trainSchedulesRepository.findByIdWithFullGraph(scheduleId);
if (!schedule || schedule.bookingWindowStatus !== "FULL") return;
// Symmetric with isScheduleFull: a weight/length-bound FULL is not stale
// just because slots remain — clearing it here would reopen a train
// nothing can board.
if (await this.isTrainFull(schedule)) return;
const customerWindowOpen =
schedule.windowPhase == null || schedule.windowPhase === "OPEN";
await this.setWindow(scheduleId, customerWindowOpen ? "OPEN" : "CLOSED");
this.logger.log(
`[BATCH] ${scheduleId} cleared stale FULL — wagons freed, window is now ` +
`${customerWindowOpen ? "OPEN" : "CLOSED"} and the batch can fill it again`,
);
}
// ---- timer plumbing -------------------------------------------------------
/** Configured customer pay window in ms (global rules, with defaults). */
private async paymentWindowMs(): Promise<number> {
const cfg = await this.trainSchedulingService.getWindowConfig();
return cfg.paymentWindowMinutes * 60_000;
}
private timeoutName(scheduleId: string): string {
return `settle:${scheduleId}`;
}
/**
* In-process accelerator only — the durable settle enforcement is the window
* engine's minute tick calling settleDueReservations off `paymentDeadline`.
*/
private armSettle(scheduleId: string): void {
void this.paymentWindowMs()
.then((delayMs) => {
this.removeTimeout(scheduleId);
const handle = setTimeout(() => {
void this.settleBatch(scheduleId).catch((err) =>
this.logger.error(
`settleBatch ${scheduleId} failed: ${(err as Error).message}`,
),
);
}, delayMs);
this.scheduler.addTimeout(this.timeoutName(scheduleId), handle);
})
.catch((err) =>
this.logger.warn(
`armSettle ${scheduleId} skipped: ${(err as Error).message}`,
),
);
}
/**
* A top-up reservation (settle freed capacity mid-cycle, so the next waiting
* booking got a fresh pay window) sets a NEW paymentDeadline. But the schedule's
* `paymentPhaseEndsAt` — which the window tick watches to end PAYMENT and run
* concludeCycle — was frozen when the phase started. Without this, concludeCycle
* fires before the top-up customer's deadline and expires a booking that still
* had time to pay. Push `paymentPhaseEndsAt` to at least cover a full payment
* window from now, but never past departure. Only while the schedule is still
* in the PAYMENT phase (a reopened cycle manages its own phase).
*/
async extendPaymentPhaseForTopUp(scheduleId: string): Promise<void> {
const schedule = await this.dataSource
.getRepository(TrainSchedule)
.findOne({ where: { id: scheduleId } });
if (!schedule || schedule.windowPhase !== "PAYMENT") return;
const windowMs = await this.paymentWindowMs();
let target = new Date(Date.now() + windowMs);
if (
schedule.scheduledDepartureDate &&
target > schedule.scheduledDepartureDate
) {
target = schedule.scheduledDepartureDate;
}
// Only ever push the deadline OUT, never pull it in.
if (
schedule.paymentPhaseEndsAt &&
schedule.paymentPhaseEndsAt.getTime() >= target.getTime()
) {
return;
}
await this.dataSource
.getRepository(TrainSchedule)
.update(scheduleId, { paymentPhaseEndsAt: target });
this.logger.log(
`[BATCH] extended PAYMENT phase for ${scheduleId} to ${target.toISOString()} ` +
`(top-up reservation opened a fresh pay window)`,
);
}
private removeTimeout(scheduleId: string): void {
const name = this.timeoutName(scheduleId);
try {
if (this.scheduler.doesExist("timeout", name)) {
this.scheduler.deleteTimeout(name);
}
} catch {
// ignore — not armed
}
}
}