import { Injectable, Logger } from '@nestjs/common'; import { Cron } from '@nestjs/schedule'; import { PrismaService } from '../../common/prisma.service'; import { SmsClientService } from '../notifications/sms-client.service'; import { CurrencyService } from '../currency/currency.service'; import { MAX_PAYMENT_HOURS, CUTOFF_MINUTES, computePaymentDeadline } from '../../common/utils/payment-deadline.utils'; // Retention windows const OTP_RETENTION_HOURS = 1; const FAYDA_SESSION_RETENTION_HOURS = 1; const AUDIT_LOG_RETENTION_DAYS = 365; const WEBHOOK_EVENT_RETENTION_DAYS = 90; const GATE_LOG_RETENTION_DAYS = 180; function fmtTime(d: Date): string { return d.toLocaleTimeString('en-GB', { hour: '2-digit', minute: '2-digit', timeZone: 'Africa/Addis_Ababa', }); } @Injectable() export class TasksService { private readonly logger = new Logger(TasksService.name); constructor( private readonly prisma: PrismaService, private readonly sms: SmsClientService, private readonly currencyService: CurrencyService, ) {} // ───────────────────────────────────────────────────────────────────────── // Every 1 min: advance TrainSchedule statuses. // // SCHEDULED → BOARDING when departure ≤ 30 min away (closed to new bookings) // BOARDING → EN_ROUTE at actual departure // EN_ROUTE → ARRIVED at arrival time // ───────────────────────────────────────────────────────────────────────── @Cron('*/1 * * * *') async syncScheduleStatuses() { const now = new Date(); const thirtyMinFromNow = new Date(now.getTime() + CUTOFF_MINUTES * 60 * 1000); // ── Schedule-level transitions (operational display) ─────────────────── const [boarding, departed, arrived] = await Promise.all([ this.prisma.trainSchedule.updateMany({ where: { status: 'SCHEDULED', departureAt: { lte: thirtyMinFromNow } }, data: { status: 'BOARDING' }, }), this.prisma.trainSchedule.updateMany({ where: { status: 'BOARDING', departureAt: { lte: now } }, data: { status: 'EN_ROUTE' }, }), this.prisma.trainSchedule.updateMany({ where: { status: 'EN_ROUTE', arrivalAt: { lte: now } }, data: { status: 'ARRIVED' }, }), ]); // ── Per-stop transitions (segment-level status) ──────────────────────── // OPEN → CHECKIN_CLOSED: each RouteStop carries its own checkinMinutesBefore // override; falls back to the Route-level value when null. // Group by effective cutoff → one updateMany per (effectiveMins, routeId) pair. const routeStops = await this.prisma.routeStop.findMany({ select: { routeId: true, stationId: true, checkinMinutesBefore: true, route: { select: { checkinMinutesBefore: true } }, }, }); // Map: effectiveMins → Map const byMins = new Map>(); for (const stop of routeStops) { const mins = stop.checkinMinutesBefore ?? stop.route.checkinMinutesBefore; if (!byMins.has(mins)) byMins.set(mins, new Map()); const byRoute = byMins.get(mins)!; if (!byRoute.has(stop.routeId)) byRoute.set(stop.routeId, []); byRoute.get(stop.routeId)!.push(stop.stationId); } // Arrival basis: each stop's own estimated arrival time, not its departure. The first // stop of a route has no arrival (nothing to arrive at), so it falls back to its // departure — expressed below as COALESCE(plannedArrivalAt, plannedDepartureAt). let reopenedCount = 0; let checkinClosedCount = 0; for (const [mins, byRoute] of byMins) { const cutoffAt = new Date(now.getTime() + mins * 60 * 1000); for (const [routeId, stationIds] of byRoute) { // Revert first: if the cutoff was reduced, stops that were prematurely closed // should reopen (arrival is still beyond the new cutoff window). const reverted = await this.prisma.tripStopTime.updateMany({ where: { status: 'CHECKIN_CLOSED', OR: [ { plannedArrivalAt: { gt: cutoffAt } }, { AND: [{ plannedArrivalAt: null }, { plannedDepartureAt: { gt: cutoffAt } }] }, ], stationId: { in: stationIds }, schedule: { routeId }, }, data: { status: 'OPEN' }, }); reopenedCount += reverted.count; // Forward: close stops now within the cutoff window. const closed = await this.prisma.tripStopTime.updateMany({ where: { status: 'OPEN', OR: [ { plannedArrivalAt: { lte: cutoffAt } }, { AND: [{ plannedArrivalAt: null }, { plannedDepartureAt: { lte: cutoffAt } }] }, ], stationId: { in: stationIds }, schedule: { routeId }, }, data: { status: 'CHECKIN_CLOSED' }, }); checkinClosedCount += closed.count; } } const boardedStops = await this.prisma.tripStopTime.updateMany({ where: { status: 'CHECKIN_CLOSED', plannedDepartureAt: { lte: now } }, data: { status: 'BOARDED' }, }); if (boarding.count > 0 || departed.count > 0 || arrived.count > 0 || reopenedCount > 0 || checkinClosedCount > 0 || boardedStops.count > 0) { this.logger.log( `Schedule sync: ${boarding.count} → BOARDING, ${departed.count} → EN_ROUTE, ${arrived.count} → ARRIVED | ` + `Stops: ${reopenedCount} → OPEN (reverted), ${checkinClosedCount} → CHECKIN_CLOSED, ${boardedStops.count} → BOARDED`, ); } } // ───────────────────────────────────────────────────────────────────────── // Every 1 min: payment deadline enforcement. // // Reminder — sent once at the midpoint of the booking's payment window: // reminder_at = booking_time + total_window / 2 // // Cancel — when now ≥ payment_deadline // payment_deadline = MIN(booking_time + 2h, departure_time - 30min) // // Examples (departure 10:00, cutoff 9:30): // Booked 8:00 → deadline 9:30, window 1.5h, reminder at 8:45 // Booked 9:00 → deadline 9:30, window 30min, reminder at 9:15 // ───────────────────────────────────────────────────────────────────────── @Cron('*/1 * * * *') async enforcePaymentDeadlines() { const now = new Date(); await Promise.all([ this.sendPaymentReminders(now), this.cancelExpiredPendingBookings(now), this.cancelExpiredPendingPackageBookings(now), ]); } // ── Send reminder at the midpoint of each booking's payment window ──────── private async sendPaymentReminders(now: Date) { // Only look at bookings created within the last 3 h with a future departure. const threeHoursAgo = new Date(now.getTime() - 3 * 60 * 60 * 1000); const bookings = await this.prisma.booking.findMany({ where: { status: 'PENDING_PAYMENT', paymentReminderSentAt: null, createdAt: { gte: threeHoursAgo }, // Do NOT filter by schedule.departureAt here: for multi-stop routes the // passenger's segment may depart well after the schedule's first stop, and // that first-stop time could already be in the past even though B→C is still open. }, include: { schedule: { include: { originStation: { select: { name: true } }, destinationStation: { select: { name: true } }, stopTimes: { select: { stationId: true, plannedArrivalAt: true, plannedDepartureAt: true } }, route: { select: { checkinMinutesBefore: true, stops: { select: { stationId: true, checkinMinutesBefore: true } } } }, }, }, }, }); for (const booking of bookings) { try { const createdAt = booking.createdAt as Date; // Use the booking's origin-segment estimated arrival (falling back to its departure // for the first stop) and that stop's own check-in window (falling back to the route // default), same resolution as holdSeats/search. const originStop = (booking.schedule as any).stopTimes?.find( (s: any) => s.stationId === (booking as any).originStationId, ); const dep = (originStop?.plannedArrivalAt ?? originStop?.plannedDepartureAt ?? booking.schedule.departureAt) as Date; const originRouteStop = (booking.schedule as any).route?.stops?.find( (s: any) => s.stationId === (booking as any).originStationId, ); const checkinMinutes = originRouteStop?.checkinMinutesBefore ?? (booking.schedule as any).route?.checkinMinutesBefore ?? 30; if (dep <= now) continue; // segment has already departed; cancel job handles clean-up const paymentDeadline = computePaymentDeadline(createdAt, dep, checkinMinutes); const totalWindowMs = paymentDeadline.getTime() - createdAt.getTime(); // Skip degenerate windows (< 2 min) — the cancel job will handle these immediately if (totalWindowMs < 2 * 60 * 1000) continue; // Remind once, at the midpoint of the total payment window const reminderAt = new Date(createdAt.getTime() + totalWindowMs / 2); if (now < reminderAt) continue; const origin = booking.schedule.originStation?.name ?? ''; const dest = booking.schedule.destinationStation?.name ?? ''; const remainingMs = Math.max(0, paymentDeadline.getTime() - now.getTime()); const remainingMin = Math.round(remainingMs / 60_000); const message = `EDR: Your booking ${booking.bookingRef} ` + `(${origin} → ${dest}) departs at ${fmtTime(dep)}. ` + `Complete payment within ${remainingMin} minute(s) (by ${fmtTime(paymentDeadline)}) ` + `or your booking will be cancelled.`; if (booking.contactPhone) { await this.sms.sendSms({ to: booking.contactPhone, message }).catch(() => null); } await this.prisma.booking.update({ where: { id: booking.id }, data: { paymentReminderSentAt: now } as any, }); this.logger.log( `Payment reminder sent: ${booking.bookingRef} ` + `(deadline ${fmtTime(paymentDeadline)}, ${remainingMin} min remaining)`, ); } catch (err) { this.logger.error( `Reminder failed for ${booking.bookingRef}: ${err instanceof Error ? err.message : String(err)}`, ); } } } // ── Cancel bookings whose payment deadline has passed ───────────────────── private async cancelExpiredPendingBookings(now: Date) { const twoHoursAgo = new Date(now.getTime() - MAX_PAYMENT_HOURS * 60 * 60 * 1000); // The departure pre-filter below is a query-scoping optimization only — the real // deadline check happens per-row further down. It must be widened to the largest // configured checkinMinutes across all routes/stops, or a booking on a route with a // cutoff bigger than the CUTOFF_MINUTES default would never even be fetched here, // silently never getting auto-cancelled. const [maxRouteCutoff, maxStopCutoff] = await Promise.all([ this.prisma.route.aggregate({ _max: { checkinMinutesBefore: true } }), this.prisma.routeStop.aggregate({ _max: { checkinMinutesBefore: true } }), ]); const effectiveMaxCutoffMinutes = Math.max( CUTOFF_MINUTES, maxRouteCutoff._max.checkinMinutesBefore ?? 0, maxStopCutoff._max.checkinMinutesBefore ?? 0, ); const departureCutoff = new Date(now.getTime() + effectiveMaxCutoffMinutes * 60 * 1000); // payment_deadline = MIN(createdAt + 2h, segment_arrival - checkinMinutes) // Deadline is reached when either branch of the MIN is in the past: // (a) createdAt ≤ now - 2h → 2-hour max window elapsed // (b) departureAt ≤ now + effectiveMaxCutoff → within the widest possible cutoff window const expiredBookings = await this.prisma.booking.findMany({ where: { status: 'PENDING_PAYMENT', OR: [ { createdAt: { lte: twoHoursAgo } }, { schedule: { departureAt: { lte: departureCutoff } } }, ], }, include: { schedule: { include: { originStation: { select: { name: true } }, destinationStation: { select: { name: true } }, stopTimes: { select: { stationId: true, plannedArrivalAt: true, plannedDepartureAt: true } }, route: { select: { checkinMinutesBefore: true, stops: { select: { stationId: true, checkinMinutesBefore: true } } } }, }, }, paymentIntent: { select: { method: true } }, seats: { select: { seatId: true } }, }, }); let cancelledCount = 0; for (const booking of expiredBookings) { try { // Re-verify exact deadline to avoid racing with a concurrent payment confirmation. // Use the booking's origin-segment estimated arrival (falling back to its departure // for the first stop) and that stop's own check-in window, so a B→C booking on an // A→B→C→D schedule gets the correct payment window anchored to B, not A. const createdAt = booking.createdAt as Date; const originStop = (booking.schedule as any).stopTimes?.find( (s: any) => s.stationId === (booking as any).originStationId, ); const dep = (originStop?.plannedArrivalAt ?? originStop?.plannedDepartureAt ?? booking.schedule.departureAt) as Date; const originRouteStop = (booking.schedule as any).route?.stops?.find( (s: any) => s.stationId === (booking as any).originStationId, ); const checkinMinutes = originRouteStop?.checkinMinutesBefore ?? (booking.schedule as any).route?.checkinMinutesBefore ?? CUTOFF_MINUTES; const paymentDeadline = computePaymentDeadline(createdAt, dep, checkinMinutes); if (now < paymentDeadline) continue; // 1a. Release held seats (Journey rows are the occupancy source of truth once paid) await this.prisma.journey.deleteMany({ where: { bookingId: booking.id } as any }); // 1b. Also release the SeatHold(s) covering this booking's seats — SeatsService // extends these to the payment deadline when the booking is created, so without // this they'd otherwise keep the seat locked for up to MAX_PAYMENT_HOURS even // though the booking is now cancelled. Scoped to this booking's own schedule, // since the same physical Seat row is reused across other recurring dates. const seatIds = booking.seats.map((s: any) => s.seatId); if (seatIds.length > 0) { await this.prisma.seatHold.deleteMany({ where: { scheduleId: booking.scheduleId, seatIds: { hasSome: seatIds } }, }); } // 2. Audit record (no refund — payment was never completed) await this.prisma.bookingCancellation.create({ data: { bookingId: booking.id, cancelledBy: 'SYSTEM', reason: 'Payment not completed before deadline', refundAmount: 0, refundMethod: booking.paymentIntent?.method ?? 'NONE', refundStatus: 'NOT_APPLICABLE', }, }).catch(() => null); // 3. Mark cancelled await this.prisma.booking.update({ where: { id: booking.id }, data: { status: 'CANCELLED' }, }); // 4. Notify passenger const origin = booking.schedule.originStation?.name ?? ''; const dest = booking.schedule.destinationStation?.name ?? ''; const message = `EDR: Your booking ${booking.bookingRef} ` + `(${origin} → ${dest}, departs ${fmtTime(dep)}) has been cancelled ` + `because payment was not completed before the deadline (${fmtTime(paymentDeadline)}).`; if (booking.contactPhone) { await this.sms.sendSms({ to: booking.contactPhone, message }).catch(() => null); } this.logger.log(`Auto-cancelled: ${booking.bookingRef} (deadline was ${fmtTime(paymentDeadline)})`); cancelledCount++; } catch (err) { this.logger.error( `Auto-cancel failed for ${booking.bookingRef}: ${err instanceof Error ? err.message : String(err)}`, ); } } if (cancelledCount > 0) { this.logger.log(`Auto-cancelled ${cancelledCount} expired pending booking(s)`); } } // ── Cancel PackageBookings whose payment deadline has passed ────────────── private async cancelExpiredPendingPackageBookings(now: Date) { const twoHoursAgo = new Date(now.getTime() - MAX_PAYMENT_HOURS * 60 * 60 * 1000); const departureCutoff = new Date(now.getTime() + CUTOFF_MINUTES * 60 * 1000); const expiredBookings = await this.prisma.packageBooking.findMany({ where: { status: 'PENDING_PAYMENT', OR: [ { createdAt: { lte: twoHoursAgo } }, { package: { outboundSchedule: { departureAt: { lte: departureCutoff } } } }, ], }, include: { package: { include: { outboundSchedule: { include: { route: { select: { checkinMinutesBefore: true } } }, }, }, }, }, }); let cancelledCount = 0; for (const booking of expiredBookings) { try { const createdAt = booking.createdAt as Date; const dep = (booking.package as any).outboundSchedule.departureAt as Date; const checkinMinutes = (booking.package as any).outboundSchedule.route?.checkinMinutesBefore ?? CUTOFF_MINUTES; const paymentDeadline = computePaymentDeadline(createdAt, dep, checkinMinutes); if (now < paymentDeadline) continue; // Revert the tier's seat counters that were incremented when the booking was created. const seatsReserved = booking.adultCount + Math.max(0, booking.childCount - booking.adultCount); await this.prisma.packagePriceTier.update({ where: { id: booking.priceTierId }, data: { bookedSeats: { decrement: seatsReserved }, availableSeats: { increment: seatsReserved }, }, }); await this.prisma.packageBooking.update({ where: { id: booking.id }, data: { status: 'CANCELLED' }, }); const message = `EDR: Your package booking ${booking.bookingRef} ` + `(departs ${fmtTime(dep)}) has been cancelled ` + `because payment was not completed by ${fmtTime(paymentDeadline)}.`; if (booking.contactPhone) { await this.sms.sendSms({ to: booking.contactPhone, message }).catch(() => null); } this.logger.log(`Auto-cancelled package booking: ${booking.bookingRef} (deadline was ${fmtTime(paymentDeadline)})`); cancelledCount++; } catch (err) { this.logger.error( `Auto-cancel failed for package booking ${(booking as any).bookingRef}: ${err instanceof Error ? err.message : String(err)}`, ); } } if (cancelledCount > 0) { this.logger.log(`Auto-cancelled ${cancelledCount} expired pending package booking(s)`); } } // ───────────────────────────────────────────────────────────────────────── // Daily at 02:00 EAT: purge expired/stale records to enforce data retention. // ───────────────────────────────────────────────────────────────────────── @Cron('0 2 * * *') async purgeExpiredData() { const now = new Date(); const otpCutoff = new Date(now.getTime() - OTP_RETENTION_HOURS * 60 * 60 * 1000); const faydaCutoff = new Date(now.getTime() - FAYDA_SESSION_RETENTION_HOURS * 60 * 60 * 1000); const auditCutoff = new Date(now.getTime() - AUDIT_LOG_RETENTION_DAYS * 24 * 60 * 60 * 1000); const webhookCutoff = new Date(now.getTime() - WEBHOOK_EVENT_RETENTION_DAYS * 24 * 60 * 60 * 1000); const gateCutoff = new Date(now.getTime() - GATE_LOG_RETENTION_DAYS * 24 * 60 * 60 * 1000); const [otps, faydaSessions, auditLogs, webhookEvents, gateLogs] = await Promise.all([ this.prisma.otpCode.deleteMany({ where: { OR: [{ expiresAt: { lte: otpCutoff } }, { verified: true, createdAt: { lte: otpCutoff } }] }, }), this.prisma.faydaVerificationSession.deleteMany({ where: { OR: [{ expiresAt: { lte: faydaCutoff } }, { status: { in: ['COMPLETED', 'FAILED'] }, createdAt: { lte: faydaCutoff } }] }, }), this.prisma.auditLog.deleteMany({ where: { createdAt: { lte: auditCutoff } } }), this.prisma.paymentWebhookEvent.deleteMany({ where: { receivedAt: { lte: webhookCutoff } } }), this.prisma.gateValidationLog.deleteMany({ where: { validatedAt: { lte: gateCutoff } } }), ]); this.logger.log( `Data retention purge: ${otps.count} OTPs, ${faydaSessions.count} Fayda sessions, ` + `${auditLogs.count} audit logs, ${webhookEvents.count} webhook events, ${gateLogs.count} gate logs deleted`, ); } }