diff --git a/apps/edr-passenger-api/package.json b/apps/edr-passenger-api/package.json index c048ef1e9..d05d44c3f 100644 --- a/apps/edr-passenger-api/package.json +++ b/apps/edr-passenger-api/package.json @@ -4,11 +4,14 @@ "private": true, "scripts": { "dev": "nest start --watch", + "start:debug": "nest start --debug --watch", + "start:debug:brk": "nest build && node --inspect-brk dist/main.js", "build": "prisma generate && nest build", "start": "node dist/main.js", "start:prod": "node dist/main.js", "lint": "eslint src", "test": "jest", + "test:debug": "node --inspect-brk node_modules/jest/bin/jest.js --runInBand", "test:e2e": "jest --config ./test/jest-e2e.json", "test:e2e:report": "jest --config ./test/jest-e2e.json; open e2e-report/index.html", "test:e2e:all": "bash ../../e2e/run.sh", @@ -18,6 +21,7 @@ "type-check": "tsc --noEmit", "iam:migrate": "node --env-file=.env scripts/run-iam-migrations.cjs", "iam:seed-dev-user": "node --env-file=.env scripts/seed-iam-dev-user.cjs", + "iam:seed-legacy-users": "node --env-file=.env scripts/seed-legacy-role-users.cjs", "prisma:generate": "prisma generate", "prisma:migrate": "prisma migrate deploy", "prisma:migrate:dev": "prisma migrate dev", diff --git a/apps/edr-passenger-api/scripts/repro-double-booking.cjs b/apps/edr-passenger-api/scripts/repro-double-booking.cjs new file mode 100644 index 000000000..d4964045d --- /dev/null +++ b/apps/edr-passenger-api/scripts/repro-double-booking.cjs @@ -0,0 +1,456 @@ +#!/usr/bin/env node +/** + * Double-booking reproduction harness — drives the real HTTP API only. + * + * No direct database access, no row edits: every seat here is claimed the same way a + * passenger's browser claims it (POST /seats/hold, POST /bookings/guest). If a scenario + * reports FAIL, two bookings hold the same seat on the same schedule and the API let it + * happen through its own public endpoints. + * + * node scripts/repro-double-booking.cjs \ + * --schedule --origin --destination [--scenario S3] [--concurrency 8] + * + * The three UUIDs come straight off any booking-detail response (data.schedule.id, + * data.schedule.origin.id, data.schedule.destination.id). Pick a FUTURE departure — holds + * are refused inside the check-in cutoff. + * + * Bookings land in PENDING_PAYMENT and are never paid, so no JourneySegment or ticket is + * produced. Every bookingRef created is printed at the end for cleanup. + */ + +const BASE = (process.env.API_URL || 'http://localhost:3002').replace(/\/$/, ''); + +// ── args ────────────────────────────────────────────────────────────────────── +const argv = process.argv.slice(2); +const arg = (name, fallback) => { + const i = argv.indexOf(`--${name}`); + return i !== -1 && argv[i + 1] ? argv[i + 1] : fallback; +}; +const SCHEDULE_ID = arg('schedule'); +const ORIGIN_ID = arg('origin'); +const DESTINATION_ID = arg('destination'); +const SEAT_CLASS_ID = arg('seat-class-id'); +const ONLY = arg('scenario'); +const CONCURRENCY = Number(arg('concurrency', 8)); + +if (!SCHEDULE_ID || !ORIGIN_ID || !DESTINATION_ID) { + console.error('Usage: --schedule --origin --destination '); + process.exit(1); +} + +const isLocal = /^https?:\/\/(localhost|127\.0\.0\.1|0\.0\.0\.0|\[::1\])(:|\/|$)/.test(BASE); +if (!isLocal && process.env.ALLOW_REMOTE !== '1') { + console.error(`Refusing to run against ${BASE} — this CREATES REAL BOOKINGS.`); + console.error('Set ALLOW_REMOTE=1 only if this is a staging/dev API you are willing to dirty.'); + process.exit(1); +} + +// ── http ────────────────────────────────────────────────────────────────────── +let callCount = 0; +async function api(method, path, body) { + callCount++; + const res = await fetch(`${BASE}${path}`, { + method, + headers: { 'content-type': 'application/json' }, + body: body === undefined ? undefined : JSON.stringify(body), + }); + const text = await res.text(); + let json; + try { + json = JSON.parse(text); + } catch { + json = { raw: text.slice(0, 300) }; + } + // A global interceptor wraps every payload as { success, data, timestamp } — unwrap it + // so callers see the resource itself, not the envelope. + const data = + json && typeof json === 'object' && 'success' in json && 'data' in json ? json.data : json; + return { ok: res.ok, status: res.status, json, data }; +} + +const uuid = () => crypto.randomUUID(); +const msg = (r) => r.json?.message ?? r.json?.error ?? JSON.stringify(r.json).slice(0, 160); + +// ── building blocks ─────────────────────────────────────────────────────────── +function holdBody(seatIds, direction, passengerId) { + return { + scheduleId: SCHEDULE_ID, + originStationId: ORIGIN_ID, + destinationStationId: DESTINATION_ID, + ...(direction ? { journeyDirection: direction } : {}), + passengers: seatIds.map((seatId) => ({ passengerId: passengerId ?? uuid(), seatId })), + }; +} + +/** + * POST /seats/hold — public, exactly what the seat map calls. Pass `passengerId` to make + * two holds look like the same traveller, which is what the OUTBOUND/RETURN exemption + * legitimately covers. + */ +async function hold(seatIds, direction, passengerId) { + const r = await api('POST', '/seats/hold', holdBody(seatIds, direction, passengerId)); + return { ...r, holdId: r.data?.id ?? r.data?.holdId }; +} + +// One traveller may not hold two unpaid bookings on the same train (see +// assertIdentitiesNotAlreadyBooked). Every run therefore needs fresh identities, or the +// leftovers from the previous run reject this one for the wrong reason. +const RUN = Date.now().toString(36).slice(-5).toUpperCase(); +let personCounter = 0; +function passenger(seatId, extra = {}) { + personCounter++; + return { + seatId, + passengerName: `Repro ${RUN} ${personCounter}`, + dateOfBirth: '1990-05-15', + idDocumentType: 'PASSPORT', + passportNumber: `RP${RUN}${String(personCounter).padStart(3, '0')}`, + passportCountry: 'Djibouti', + nationality: 'Other', + phone: `+2519${String(10000000 + personCounter).slice(0, 8)}`, + ...extra, + }; +} + +/** + * POST /bookings/guest — public, no auth. `holdId` and the seat ids in `passengers` + * are sent independently, which is the whole point of scenarios 1 and 2. + */ +async function guestBook({ holdId, seatIds, seatClassId }) { + const r = await api('POST', '/bookings/guest', { + bookingType: 'ONE_WAY', + scheduleId: SCHEDULE_ID, + holdId, + originStationId: ORIGIN_ID, + destinationStationId: DESTINATION_ID, + seatClassId, + skipIdentityVerification: true, + passengers: seatIds.map((s) => passenger(s)), + }); + const ref = r.data?.bookingRef ?? r.data?.booking?.bookingRef; + return { ...r, bookingRef: ref }; +} + +// ── ledger: who ended up on which seat ──────────────────────────────────────── +const claims = []; // { seatId, seatNumber, bookingRef, scenario } +const created = []; // every bookingRef we made, for cleanup +function record(scenario, seat, bookingRef) { + claims.push({ seatId: seat.id, seatNumber: seat.seatNumber, bookingRef, scenario }); + created.push(bookingRef); +} + +// ── seat pool ───────────────────────────────────────────────────────────────── +async function loadSeats() { + const q = `originStationId=${ORIGIN_ID}&destinationStationId=${DESTINATION_ID}`; + const r = await api('GET', `/seats/seatmap/${SCHEDULE_ID}?${q}`); + if (!r.ok) throw new Error(`seatmap failed (${r.status}): ${msg(r)}`); + + const pool = []; + for (const coach of r.data?.coaches ?? []) { + for (const s of coach.seats ?? []) { + if (String(s.status).toUpperCase() !== 'AVAILABLE') continue; + pool.push({ + id: s.id, + seatNumber: s.seatNumber, + coach: coach.coachNumber, + coachClass: coach.seatClass, + }); + } + } + return pool; +} + +async function resolveSeatClassId(seat) { + if (SEAT_CLASS_ID) return SEAT_CLASS_ID; + const r = await api('GET', '/seat-classes'); + const list = Array.isArray(r.data) ? r.data : (r.data?.items ?? []); + const norm = (s) => String(s || '').replace(/\s+/g, ' ').trim().toLowerCase(); + const hit = list.find((c) => norm(c.name) === norm(seat.coachClass)) ?? list[0]; + if (!hit) throw new Error('No seat classes returned — pass --seat-class-id explicitly.'); + return hit.id; +} + +// ── scenarios ───────────────────────────────────────────────────────────────── +// Each returns { verdict: 'FAIL' | 'PASS' | 'INCONCLUSIVE', detail }. +// S* are attacks: FAIL means the double booking reproduced. +// P* are positive controls: FAIL means a legitimate booking got blocked, which is just +// as much a regression — a guard that refuses everything is not a fix. + +const scenarios = { + /** P1 — the ordinary happy path: hold a seat, book that same seat. Must succeed. */ + async P1(pool, seatClassId) { + const seat = pool.shift(); + const h = await hold([seat.id]); + if (!h.holdId) return { verdict: 'FAIL', detail: `could not hold seat ${seat.seatNumber}: ${msg(h)}` }; + + const b = await guestBook({ holdId: h.holdId, seatIds: [seat.id], seatClassId }); + if (!b.bookingRef) { + return { verdict: 'FAIL', detail: `held seat ${seat.seatNumber} but booking was refused (${b.status}): ${msg(b)}` }; + } + record('P1', seat, b.bookingRef); + return { verdict: 'PASS', detail: `seat ${seat.seatNumber} held and booked normally — ${b.bookingRef}` }; + }, + + /** + * P2 — one traveller, both legs of a turnaround round trip on the same seat. This is + * what the OUTBOUND/RETURN exemption is for, so it must keep working after the + * exemption is narrowed to the requester's own holds. + */ + async P2(pool) { + const seat = pool.shift(); + const traveller = uuid(); + const outbound = await hold([seat.id], 'OUTBOUND', traveller); + if (!outbound.holdId) { + return { verdict: 'FAIL', detail: `outbound hold refused on seat ${seat.seatNumber}: ${msg(outbound)}` }; + } + const ret = await hold([seat.id], 'RETURN', traveller); + if (!ret.holdId) { + return { + verdict: 'FAIL', + detail: `same traveller was blocked from holding their own return leg on seat ${seat.seatNumber}: ${msg(ret)}`, + }; + } + return { verdict: 'PASS', detail: `one traveller held seat ${seat.seatNumber} on both legs` }; + }, + + /** + * S1 — Hold laundering. Hold seat A, then guest-book seat B while presenting A's holdId. + * bookings.service.ts runs validateSeatIdsAgainstHold here; guest-booking.service.ts + * does not, so the seat you pay for need never be the seat you held. + */ + async S1(pool, seatClassId) { + const [decoy, target] = [pool.shift(), pool.shift()]; + const h = await hold([decoy.id]); + if (!h.holdId) return { verdict: 'INCONCLUSIVE', detail: `hold on decoy failed: ${msg(h)}` }; + + const b = await guestBook({ holdId: h.holdId, seatIds: [target.id], seatClassId }); + if (!b.bookingRef) return { verdict: 'PASS', detail: `rejected (${b.status}): ${msg(b)}` }; + + record('S1', target, b.bookingRef); + return { + verdict: 'FAIL', + detail: `held seat ${decoy.seatNumber} but booked seat ${target.seatNumber} — ${b.bookingRef}`, + }; + }, + + /** + * S2 — Steal a live hold. Someone else holds the seat; we present an unrelated hold + * and book their seat anyway. This is the production shape: a seat that another + * passenger is mid-checkout on gets sold underneath them. + */ + async S2(pool, seatClassId) { + const [victimSeat, mySeat] = [pool.shift(), pool.shift()]; + const victimHold = await hold([victimSeat.id]); + if (!victimHold.holdId) + return { verdict: 'INCONCLUSIVE', detail: `victim hold failed: ${msg(victimHold)}` }; + + const myHold = await hold([mySeat.id]); + if (!myHold.holdId) + return { verdict: 'INCONCLUSIVE', detail: `attacker hold failed: ${msg(myHold)}` }; + + const b = await guestBook({ holdId: myHold.holdId, seatIds: [victimSeat.id], seatClassId }); + if (!b.bookingRef) return { verdict: 'PASS', detail: `rejected (${b.status}): ${msg(b)}` }; + + record('S2', victimSeat, b.bookingRef); + return { + verdict: 'FAIL', + detail: `seat ${victimSeat.seatNumber} was under a live hold, booked anyway — ${b.bookingRef}`, + }; + }, + + /** + * S3 — Concurrent guest bookings, same seat. N independent "browsers" each hold their + * own throwaway seat, then all fire at the same instant for one shared seat. Closest + * analogue to peak-hour traffic. More than one bookingRef means the seat was sold twice. + */ + async S3(pool, seatClassId) { + const target = pool.shift(); + const decoys = pool.splice(0, CONCURRENCY); + if (decoys.length < CONCURRENCY) + return { verdict: 'INCONCLUSIVE', detail: 'not enough free seats' }; + + const holds = await Promise.all(decoys.map((d) => hold([d.id]))); + const usable = holds.filter((h) => h.holdId); + if (usable.length < 2) + return { verdict: 'INCONCLUSIVE', detail: 'fewer than 2 holds succeeded' }; + + const results = await Promise.all( + usable.map((h) => guestBook({ holdId: h.holdId, seatIds: [target.id], seatClassId })), + ); + const won = results.filter((r) => r.bookingRef); + won.forEach((r) => record('S3', target, r.bookingRef)); + + if (won.length <= 1) { + return { + verdict: 'PASS', + detail: `${won.length}/${usable.length} succeeded on seat ${target.seatNumber}`, + }; + } + return { + verdict: 'FAIL', + detail: `seat ${target.seatNumber} sold ${won.length}x concurrently — ${won.map((r) => r.bookingRef).join(', ')}`, + }; + }, + + /** + * S4 — Concurrent holds, same seat. Tests the hold transaction itself, upstream of any + * booking. If two holds coexist on one seat, every downstream check is already poisoned. + */ + async S4(pool) { + const target = pool.shift(); + const results = await Promise.all(Array.from({ length: CONCURRENCY }, () => hold([target.id]))); + const won = results.filter((r) => r.holdId); + if (won.length <= 1) { + return { + verdict: 'PASS', + detail: `${won.length}/${CONCURRENCY} holds granted on seat ${target.seatNumber}`, + }; + } + return { + verdict: 'FAIL', + detail: `seat ${target.seatNumber} held ${won.length}x simultaneously — holdIds ${won.map((r) => r.holdId).join(', ')}`, + }; + }, + + /** + * S5 — Direction bypass. checkDirectionConflict (journey-direction.utils.ts:13) returns + * false for OUTBOUND vs RETURN so a round trip can reuse a seat across its own legs. + * The rule is per-seat, not per-booking, so two *different* passengers can straddle it. + */ + async S5(pool, seatClassId) { + const target = pool.shift(); + const outbound = await hold([target.id], 'OUTBOUND'); + const ret = await hold([target.id], 'RETURN'); + + if (!outbound.holdId || !ret.holdId) { + return { + verdict: 'PASS', + detail: `second direction refused: ${msg(outbound.holdId ? ret : outbound)}`, + }; + } + + const a = await guestBook({ holdId: outbound.holdId, seatIds: [target.id], seatClassId }); + const b = await guestBook({ holdId: ret.holdId, seatIds: [target.id], seatClassId }); + const won = [a, b].filter((r) => r.bookingRef); + won.forEach((r) => record('S5', target, r.bookingRef)); + + if (won.length <= 1) { + return { + verdict: 'INCONCLUSIVE', + detail: `both directions held seat ${target.seatNumber} at once, but only ${won.length} booking stuck`, + }; + } + return { + verdict: 'FAIL', + detail: `seat ${target.seatNumber} sold as OUTBOUND and RETURN — ${won.map((r) => r.bookingRef).join(', ')}`, + }; + }, + + /** + * S6 — Hold expiry gap. A booking sits in PENDING_PAYMENT; its hold lapses; no + * JourneySegment exists yet because those are only written on payment success + * (payments.service.ts:2163). The seat reads as free to everyone until the first + * booking pays. Needs --wait-hold-expiry, since it must outlive SEAT_HOLD_DURATION_MINUTES. + */ + async S6(pool, seatClassId) { + const waitMin = Number(arg('wait-hold-expiry', '0')); + if (!waitMin) { + return { + verdict: 'SKIPPED', + detail: 'pass --wait-hold-expiry SEAT_HOLD_DURATION_MINUTES>', + }; + } + const target = pool.shift(); + const first = await hold([target.id]); + if (!first.holdId) return { verdict: 'INCONCLUSIVE', detail: `hold failed: ${msg(first)}` }; + + const b1 = await guestBook({ holdId: first.holdId, seatIds: [target.id], seatClassId }); + if (!b1.bookingRef) + return { verdict: 'INCONCLUSIVE', detail: `first booking failed: ${msg(b1)}` }; + record('S6', target, b1.bookingRef); + + console.log(` waiting ${waitMin} min for the hold to lapse (booking stays PENDING_PAYMENT)…`); + await new Promise((r) => setTimeout(r, waitMin * 60 * 1000)); + + const second = await hold([target.id]); + if (!second.holdId) + return { verdict: 'PASS', detail: `re-hold refused after expiry: ${msg(second)}` }; + + const b2 = await guestBook({ holdId: second.holdId, seatIds: [target.id], seatClassId }); + if (!b2.bookingRef) return { verdict: 'PASS', detail: `second booking refused: ${msg(b2)}` }; + + record('S6', target, b2.bookingRef); + return { + verdict: 'FAIL', + detail: `seat ${target.seatNumber} rebooked after hold lapsed — ${b1.bookingRef} and ${b2.bookingRef}`, + }; + }, +}; + +// ── runner ──────────────────────────────────────────────────────────────────── +(async () => { + console.log(`API ${BASE}`); + console.log(`Schedule ${SCHEDULE_ID}`); + console.log(`Route ${ORIGIN_ID} -> ${DESTINATION_ID}\n`); + + const pool = await loadSeats(); + console.log(`${pool.length} seats reported AVAILABLE.\n`); + if (pool.length < CONCURRENCY + 6) { + console.error( + `Need at least ${CONCURRENCY + 6} free seats; lower --concurrency or pick an emptier departure.`, + ); + process.exit(1); + } + + const seatClassId = await resolveSeatClassId(pool[0]); + const names = ONLY ? [ONLY] : Object.keys(scenarios); + const summary = []; + + for (const name of names) { + const fn = scenarios[name]; + if (!fn) { + console.error(`Unknown scenario ${name}`); + continue; + } + process.stdout.write(`${name} … `); + let out; + try { + out = await fn(pool, seatClassId); + } catch (e) { + out = { verdict: 'ERROR', detail: e.message }; + } + console.log(`${out.verdict}\n ${out.detail}`); + summary.push({ name, ...out }); + } + + // Collision ledger — built only from what the API handed back to us. + const bySeat = new Map(); + for (const c of claims) { + if (!bySeat.has(c.seatId)) bySeat.set(c.seatId, []); + bySeat.get(c.seatId).push(c); + } + const collisions = [...bySeat.values()].filter((g) => g.length > 1); + + console.log(`\n─────── ${callCount} API calls, ${created.length} bookings created ───────`); + for (const s of summary) console.log(` ${s.verdict.padEnd(13)} ${s.name}`); + + if (collisions.length) { + console.log('\nDOUBLE-BOOKED SEATS:'); + for (const g of collisions) { + console.log(` seat ${g[0].seatNumber} (${g[0].seatId})`); + for (const c of g) console.log(` ${c.bookingRef} [${c.scenario}]`); + } + } else { + console.log('\nNo seat was claimed by more than one booking.'); + } + + if (created.length) { + console.log('\nCreated (PENDING_PAYMENT, unpaid — cancel these):'); + console.log(` ${created.join(' ')}`); + } + + process.exit(collisions.length ? 2 : 0); +})().catch((e) => { + console.error(e); + process.exit(1); +}); diff --git a/apps/edr-passenger-api/scripts/seed-legacy-role-users.cjs b/apps/edr-passenger-api/scripts/seed-legacy-role-users.cjs new file mode 100644 index 000000000..f174a2f81 --- /dev/null +++ b/apps/edr-passenger-api/scripts/seed-legacy-role-users.cjs @@ -0,0 +1,236 @@ +/** + * Create the three legacy backoffice roles and their users. + * + * Deliberately a standalone CLI, NOT part of the app lifecycle: nothing here runs on + * `onApplicationBootstrap`, so a deploy can decide when accounts appear. Contrast + * `EdrPassengerOrgSeeder` / `PassengerStaffUsersSeeder`, which run at boot behind + * SEED_EDR_PASSENGER_ORG / SEED_PASSENGER_STAFF. + * + * pnpm --filter @edr/passenger-api iam:seed-legacy-users + * pnpm --filter @edr/passenger-api iam:seed-legacy-users -- --dry-run + * + * The roles below carry ONLY the permission keys that existed before the granular + * create/edit/delete change (47 of them). None of the newer narrow keys are granted. + * That is the point: these three model how the app was used before, so running them + * against the new guards proves the granular change stayed backwards-compatible. + * + * Idempotent, and safe to re-run: roles and users are upserted by their natural key, + * and each role's permission links are synced to exactly the set declared here — + * extras are pruned so the file stays the source of truth. + * + * Requires DATABASE_* in the environment (`node --env-file=.env` does this). + * The password comes from SEED_USER_PASSWORD, falling back to DEFAULT_PASSWORD. + * In production one of them MUST be set — there is no built-in default there. + */ +const { DataSource } = require('typeorm'); +const { hashPassword } = require('@tria-plc/api-common/utils/argon'); + +const APP = 'edr_passenger_app'; +const ORG_KEY = 'edr'; +const DRY_RUN = process.argv.includes('--dry-run'); + +const p = (s) => `${APP}:${s}`; + +/** Every `:view` / `:view_all` key that existed before the granular change. */ +const LEGACY_VIEW_SUFFIXES = [ + 'agents:view', 'audit:view', 'bookings:view', 'classes:view', 'coaches:view', + 'currencies:view', 'dashboard:view', 'fraud:view', 'inquiries:view', 'packages:view', + 'passengers:view', 'payment_methods:view', 'payments:view', 'payments:view_all', + 'reports:view', 'routes:view', 'schedules:view', 'seats:view', 'stations:view', + 'tariff_rates:view', 'tickets:view', 'trains:view', +]; + +/** The rest of the pre-granular registry — 25 non-view keys. */ +const LEGACY_WRITE_SUFFIXES = [ + 'admin', + 'agents:manage', 'bookings:cancel', 'bookings:manage', 'bookings:reschedule', + 'classes:manage', 'coaches:manage', 'currencies:manage', 'fraud:manage', + 'inquiries:manage', 'notifications:send', 'packages:manage', 'passengers:manage', + 'payment_methods:manage', 'payments:manage', 'payments:manage_methods', + 'payments:refund', 'routes:manage', 'schedules:manage', 'seats:manage', + 'stations:manage', 'tariff_rates:manage', 'tickets:generate', 'tickets:manage', + 'trains:manage', +]; + +const LEGACY_VIEW = LEGACY_VIEW_SUFFIXES.map(p); +const LEGACY_ALL = [...LEGACY_VIEW_SUFFIXES, ...LEGACY_WRITE_SUFFIXES].map(p); + +/** + * Withheld from the chief. `payment_methods:manage` and the legacy alias + * `payments:manage_methods` both open POST/PATCH /payments/methods, so excluding only + * one would leave the ability intact — both have to go for "no managing payment + * methods" to actually hold. + */ +const CHIEF_EXCLUDED = [ + p('payment_methods:manage'), + p('payments:manage_methods'), + p('tickets:generate'), + p('admin'), +]; + +const ROLES = [ + { + key: 'old_ticketofficer', + name: { en: 'Ticket Officer (legacy)', am: 'የቲኬት ኦፊሰር' }, + email: 'old.ticketofficer@edr.local', + username: 'old_ticketofficer', + permissions: [ + p('passengers:view'), + p('bookings:view'), + p('tickets:view'), + p('tickets:manage'), + p('payments:view'), + p('payments:manage'), + ], + }, + { + key: 'old_passengerchief', + name: { en: 'Passenger Chief (legacy)', am: 'የተሳፋሪ ኃላፊ' }, + email: 'old.passengerchief@edr.local', + username: 'old_passengerchief', + permissions: LEGACY_ALL.filter((k) => !CHIEF_EXCLUDED.includes(k)), + }, + { + key: 'old_passengerdirector', + name: { en: 'Passenger Director (legacy)', am: 'የተሳፋሪ ዳይሬክተር' }, + email: 'old.passengerdirector@edr.local', + username: 'old_passengerdirector', + permissions: LEGACY_VIEW, + }, +]; + +function resolvePassword() { + const pw = (process.env.SEED_USER_PASSWORD || process.env.DEFAULT_PASSWORD || '').trim(); + if (pw) return pw; + if (process.env.NODE_ENV === 'production') { + throw new Error('SEED_USER_PASSWORD (or DEFAULT_PASSWORD) must be set in production'); + } + return '12345678'; +} + +const ds = new DataSource({ + type: 'postgres', + host: process.env.DATABASE_HOST, + port: Number(process.env.DATABASE_PORT || 5432), + database: process.env.DATABASE_NAME, + username: process.env.DATABASE_USER, + password: process.env.DATABASE_PASSWORD, +}); + +(async () => { + const password = resolvePassword(); + await ds.initialize(); + + // Fail before writing anything if a key is not seeded — a typo here would otherwise + // create a role that silently grants less than intended. + const wanted = [...new Set(ROLES.flatMap((r) => r.permissions))]; + const found = await ds.query( + `SELECT key FROM iam.permissions WHERE key = ANY($1::text[])`, + [wanted], + ); + const missing = wanted.filter((k) => !found.some((f) => f.key === k)); + if (missing.length) { + throw new Error( + `these permission keys are not in iam.permissions — run the app once with ` + + `SEED_EDR_PASSENGER_ORG=true first:\n ${missing.join('\n ')}`, + ); + } + + const [org] = await ds.query(`SELECT id FROM iam.organizations WHERE key = $1`, [ORG_KEY]); + if (!org) throw new Error(`missing_organization:${ORG_KEY}`); + + if (DRY_RUN) { + console.log('\n=== DRY RUN — nothing written ==='); + for (const r of ROLES) { + console.log(`\n${r.key} (${r.email}) ${r.permissions.length} permissions`); + for (const k of [...r.permissions].sort()) console.log(' ', k); + } + await ds.destroy(); + return; + } + + const hashed = await hashPassword(password); + + await ds.transaction(async (m) => { + for (const r of ROLES) { + const [role] = await m.query( + `INSERT INTO iam.roles (id, key, name, created_at, updated_at) + VALUES (gen_random_uuid(), $1, $2::jsonb, now(), now()) + ON CONFLICT (key) DO UPDATE SET name = EXCLUDED.name, updated_at = now() + RETURNING id`, + [r.key, JSON.stringify(r.name)], + ); + + // Sync links to exactly this set: add what is missing, drop what is extra. + await m.query( + `INSERT INTO iam.role_permissions (id, role_id, permission_id, created_at, updated_at) + SELECT gen_random_uuid(), $1, p.id, now(), now() + FROM iam.permissions p + WHERE p.key = ANY($2::text[]) + AND NOT EXISTS (SELECT 1 FROM iam.role_permissions rp + WHERE rp.role_id = $1 AND rp.permission_id = p.id)`, + [role.id, r.permissions], + ); + // TypeORM's postgres driver returns `[rows, affectedCount]` for a DELETE ... RETURNING, + // so the rows are at [0] — reading `.length` off the outer array would report 2 every time. + const deleted = await m.query( + `DELETE FROM iam.role_permissions rp + USING iam.permissions p + WHERE rp.permission_id = p.id + AND rp.role_id = $1 + AND NOT (p.key = ANY($2::text[])) + RETURNING rp.id`, + [role.id, r.permissions], + ); + const prunedCount = (Array.isArray(deleted[0]) ? deleted[0] : deleted).length; + + const [user] = await m.query( + `INSERT INTO iam.users (id, name, username, email, user_type, status, is_active, + has_set_password, created_at, updated_at) + VALUES (gen_random_uuid(), $1::jsonb, $2, $3, 'individual', 'accepted', true, true, + now(), now()) + ON CONFLICT (email) DO UPDATE SET updated_at = now() + RETURNING id`, + [JSON.stringify(r.name), r.username, r.email], + ); + + // Never overwrite a password that already exists — re-running must not reset a + // credential someone has since changed. + await m.query( + `INSERT INTO iam.user_credentials (id, user_id, password, is_active, created_at, updated_at) + SELECT gen_random_uuid(), $1, $2, true, now(), now() + WHERE NOT EXISTS (SELECT 1 FROM iam.user_credentials + WHERE user_id = $1 AND is_active = true)`, + [user.id, hashed], + ); + + await m.query( + `INSERT INTO iam.user_roles (id, user_id, role_id, organization_id, created_at, updated_at) + SELECT gen_random_uuid(), $1, $2, $3, now(), now() + WHERE NOT EXISTS (SELECT 1 FROM iam.user_roles + WHERE user_id = $1 AND role_id = $2)`, + [user.id, role.id, org.id], + ); + + await m.query( + `INSERT INTO iam.employees (id, user_id, organization_id, is_current, name, created_at, updated_at) + SELECT gen_random_uuid(), $1, $2, true, $3::jsonb, now(), now() + WHERE NOT EXISTS (SELECT 1 FROM iam.employees + WHERE user_id = $1 AND organization_id = $2 AND is_current = true)`, + [user.id, org.id, JSON.stringify(r.name)], + ); + + console.log( + ` ${r.key.padEnd(24)} ${String(r.permissions.length).padStart(2)} permissions` + + `${prunedCount ? ` (${prunedCount} stale link(s) pruned)` : ''} -> ${r.email}`, + ); + } + }); + + console.log('\nDone. Sign in with the email above and the seeded password.'); + console.log('Re-running is safe; an existing password is never overwritten.\n'); + await ds.destroy(); +})().catch((e) => { + console.error('[seed-legacy-role-users] FAIL:', e.message); + process.exit(1); +}); diff --git a/apps/edr-passenger-api/scripts/stress-booking-concurrency.cjs b/apps/edr-passenger-api/scripts/stress-booking-concurrency.cjs new file mode 100644 index 000000000..28e08b595 --- /dev/null +++ b/apps/edr-passenger-api/scripts/stress-booking-concurrency.cjs @@ -0,0 +1,359 @@ +#!/usr/bin/env node +/** + * Booking concurrency + regression harness. + * + * Companion to repro-double-booking.cjs: that one proves specific attacks are closed, this + * one proves the system stays correct under load and that legitimate bookings are not + * rejected as collateral. + * + * node scripts/stress-booking-concurrency.cjs + * --schedule --origin --destination + * [--mode ladder|edge|all] [--levels 2,5,10,25,50,100] + * [--return-schedule --return-origin --return-destination ] + * + * ladder, per concurrency level N: + * hold-contest N browsers click one seat at once -> exactly 1 hold granted + * booking-storm N submits of one valid hold at once -> exactly 1 booking created + * diff-seats N users book N different seats at once -> all N succeed (no false rejects) + * + * edge: multi-passenger, all-or-nothing partial availability, overlapping multi-seat + * requests, and a contested round trip. + * + * Everything goes through the public HTTP API — no direct database writes. Every request + * body is built up front and released together from one Promise.all, and each line reports + * an factor (summed latency / wall time) so the concurrency is provable rather + * than assumed. + * + * Bookings land in PENDING_PAYMENT and are never paid. Every ref created is printed at the + * end; cancel them afterwards. Pick a FUTURE departure with plenty of free seats. + */ +const BASE = (process.env.API_URL || 'http://localhost:3002').replace(/\/$/, ''); +const argv = process.argv.slice(2); +const arg = (n, d) => { const i = argv.indexOf(`--${n}`); return i !== -1 && argv[i + 1] ? argv[i + 1] : d; }; + +const SCHED = arg('schedule'); +const ORIGIN = arg('origin'); +const DEST = arg('destination'); +const RET_SCHED = arg('return-schedule'); +const RET_ORIGIN = arg('return-origin'); +const RET_DEST = arg('return-destination'); +const MODE = arg('mode', 'ladder'); +const LEVELS = arg('levels', '2,5,10,25,50,100').split(',').map(Number); + +// Node's built-in fetch (undici) defaults to an unlimited per-origin connection pool, so +// Promise.all really does put every request on the wire at once. Proven per run by +// comparing wall time against the summed per-request latency (see in each line). + +const RUN = Date.now().toString(36).slice(-5).toUpperCase(); +let pc = 0; +const uuid = () => crypto.randomUUID(); + +async function api(method, path, body) { + const t0 = Date.now(); + try { + const res = await fetch(`${BASE}${path}`, { + method, + headers: { 'content-type': 'application/json' }, + body: body === undefined ? undefined : JSON.stringify(body), + }); + const text = await res.text(); + let json; try { json = JSON.parse(text); } catch { json = { raw: text.slice(0, 200) }; } + const data = json && typeof json === 'object' && 'success' in json && 'data' in json ? json.data : json; + return { ok: res.ok, status: res.status, json, data, ms: Date.now() - t0 }; + } catch (e) { + return { ok: false, status: 0, json: { message: `NETWORK: ${e.message}` }, data: null, ms: Date.now() - t0 }; + } +} +const msg = (r) => String(r.json?.message ?? r.json?.error ?? JSON.stringify(r.json ?? {})).slice(0, 150); + +function person(seatId) { + pc++; + const tag = `${RUN}${String(pc).padStart(4, '0')}`; + return { + seatId, + passengerName: `Stress ${tag}`, + dateOfBirth: '1990-05-15', + idDocumentType: 'PASSPORT', + passportNumber: `ST${tag}`, + passportCountry: 'Djibouti', + nationality: 'Other', + phone: `+2519${String(40000000 + pc).slice(0, 8)}`, + }; +} + +const holdBody = (seatIds, sched = SCHED, o = ORIGIN, d = DEST, dir, pid) => ({ + scheduleId: sched, originStationId: o, destinationStationId: d, + ...(dir ? { journeyDirection: dir } : {}), + passengers: seatIds.map((seatId) => ({ passengerId: pid ?? uuid(), seatId })), +}); +async function hold(seatIds, sched, o, d, dir, pid) { + const r = await api('POST', '/seats/hold', holdBody(seatIds, sched, o, d, dir, pid)); + return { ...r, holdId: r.data?.holdId ?? r.data?.id }; +} + +function oneWayBody(holdId, seatIds, seatClassId) { + return { + bookingType: 'ONE_WAY', scheduleId: SCHED, holdId, + originStationId: ORIGIN, destinationStationId: DEST, + seatClassId, skipIdentityVerification: true, + passengers: seatIds.map((s) => person(s)), + }; +} +function roundTripBody(holdId, returnHoldId, outSeats, retSeats, seatClassId) { + return { + bookingType: 'ROUND_TRIP', scheduleId: SCHED, holdId, + originStationId: ORIGIN, destinationStationId: DEST, + returnScheduleId: RET_SCHED, returnHoldId, + returnOriginStationId: RET_ORIGIN, returnDestinationStationId: RET_DEST, + seatClassId, returnSeatClassId: seatClassId, skipIdentityVerification: true, + passengers: outSeats.map((s, i) => ({ ...person(s), returnSeatId: retSeats[i] })), + }; +} +async function book(body) { + const r = await api('POST', '/bookings/guest', body); + return { ...r, bookingRef: r.data?.bookingRef }; +} + +async function seatPool(sched = SCHED, o = ORIGIN, d = DEST) { + const r = await api('GET', `/seats/seatmap/${sched}?originStationId=${o}&destinationStationId=${d}`); + if (!r.ok) throw new Error(`seatmap ${r.status}: ${msg(r)}`); + const out = []; + for (const c of r.data?.coaches ?? []) + for (const s of c.seats ?? []) + if (String(s.status).toUpperCase() === 'AVAILABLE') + out.push({ id: s.id, seatNumber: s.seatNumber, coach: c.coachNumber, coachClass: c.seatClass }); + return out; +} +async function seatClassId(seat) { + const r = await api('GET', '/seat-classes'); + const list = Array.isArray(r.data) ? r.data : (r.data?.items ?? []); + const norm = (s) => String(s || '').replace(/\s+/g, ' ').trim().toLowerCase(); + return (list.find((c) => norm(c.name) === norm(seat.coachClass)) ?? list[0]).id; +} + +/** Fire everything at once and report the arrival spread so "concurrent" is provable. */ +async function fireAll(tasks) { + const t0 = Date.now(); + const starts = []; + const wrapped = tasks.map((fn) => (async () => { starts.push(Date.now() - t0); return fn(); })()); + const results = await Promise.all(wrapped); + const wallMs = Date.now() - t0; + const sumMs = results.reduce((a, r) => a + (r.ms || 0), 0); + // overlap >> 1 proves the requests were genuinely in flight together rather than queued. + return { results, wallMs, releaseSpreadMs: Math.max(...starts) - Math.min(...starts), overlap: (sumMs / Math.max(wallMs, 1)).toFixed(1) }; +} + +function tally(results) { + const t = {}; + for (const r of results) { + const k = r.bookingRef ? 'created' : r.holdId ? 'held' : `http_${r.status}`; + t[k] = (t[k] || 0) + 1; + } + return t; +} +const errSample = (results) => + [...new Set(results.filter((r) => !r.bookingRef && !r.holdId).map((r) => `${r.status}: ${msg(r)}`))].slice(0, 3); + +const ALL = { created: [], failures: [] }; +function note(refs) { ALL.created.push(...refs); } + +// ───────────────────────────────────────────────────────────────────────────── +/** N users all click the SAME seat: N simultaneous holds. Exactly 1 must win. */ +async function holdContest(pool, n) { + const seat = pool.shift(); + const { results, wallMs, releaseSpreadMs, overlap } = await fireAll( + Array.from({ length: n }, () => () => hold([seat.id])), + ); + const won = results.filter((r) => r.holdId); + return { + name: `hold-contest n=${n}`, + pass: won.length === 1, + detail: `seat ${seat.seatNumber}: ${won.length} hold(s) granted ${JSON.stringify(tally(results))} ` + + `wall=${wallMs}ms release-spread=${releaseSpreadMs}ms overlap=x${overlap}`, + errs: errSample(results), + winner: won[0], + seat, + }; +} + +/** The double-submit / retry storm: N simultaneous bookings on ONE valid hold. */ +async function bookingStorm(seat, holdId, scid, n) { + const bodies = Array.from({ length: n }, () => oneWayBody(holdId, [seat.id], scid)); + const { results, wallMs, releaseSpreadMs, overlap } = await fireAll(bodies.map((b) => () => book(b))); + const won = results.filter((r) => r.bookingRef); + note(won.map((r) => r.bookingRef)); + return { + name: `booking-storm n=${n}`, + pass: won.length === 1, + detail: `seat ${seat.seatNumber}: ${won.length} booking(s) ${JSON.stringify(tally(results))} ` + + `wall=${wallMs}ms release-spread=${releaseSpreadMs}ms refs=${won.map((r) => r.bookingRef).join(',')}`, + errs: errSample(results), + }; +} + +/** N users, N DIFFERENT seats, all at once. All N must succeed — no false rejections. */ +async function differentSeats(pool, scid, n) { + const seats = pool.splice(0, n); + if (seats.length < n) return { name: `diff-seats n=${n}`, pass: false, detail: 'not enough free seats', errs: [] }; + const holds = await Promise.all(seats.map((s) => hold([s.id]))); + const usable = holds.map((h, i) => ({ h, s: seats[i] })).filter((x) => x.h.holdId); + const bodies = usable.map((x) => oneWayBody(x.h.holdId, [x.s.id], scid)); + const { results, wallMs, releaseSpreadMs, overlap } = await fireAll(bodies.map((b) => () => book(b))); + const won = results.filter((r) => r.bookingRef); + note(won.map((r) => r.bookingRef)); + return { + name: `diff-seats n=${n}`, + pass: won.length === usable.length && usable.length === n, + detail: `${won.length}/${usable.length} distinct-seat bookings succeeded (holds granted ${usable.length}/${n}) ` + + `${JSON.stringify(tally(results))} wall=${wallMs}ms release-spread=${releaseSpreadMs}ms overlap=x${overlap}`, + errs: errSample(results), + }; +} + +/** Overlapping multi-seat requests over a small shared pool — deadlock + partial-write probe. */ +async function overlapping(pool, scid, n, poolSize = 4) { + const shared = pool.splice(0, poolSize); + const holds = await Promise.all(shared.map((s) => hold([s.id]))); + const byId = new Map(shared.map((s, i) => [s.id, holds[i]])); + const granted = shared.filter((s) => byId.get(s.id).holdId); + if (granted.length < 2) return { name: `overlap n=${n}`, pass: false, detail: 'too few holds', errs: [] }; + + // Each request asks for 2 seats from the shared pool, in varying order. + const bodies = []; + for (let i = 0; i < n; i++) { + const a = granted[i % granted.length]; + const b = granted[(i + 1 + (i % 2)) % granted.length]; + if (a.id === b.id) continue; + // Present the hold of the FIRST seat; second seat will be rejected as not-in-hold. + bodies.push({ body: oneWayBody(byId.get(a.id).holdId, [a.id, b.id], scid), seats: [a, b] }); + } + const { results, wallMs } = await fireAll(bodies.map((x) => () => book(x.body))); + const won = results.filter((r) => r.bookingRef); + note(won.map((r) => r.bookingRef)); + return { + name: `overlap n=${bodies.length}`, + pass: true, // correctness judged by the DB duplicate check, not the count + detail: `${won.length} booking(s) from ${bodies.length} overlapping 2-seat requests over ${granted.length} seats ` + + `${JSON.stringify(tally(results))} wall=${wallMs}ms`, + errs: errSample(results), + }; +} + +/** Multi-passenger booking where one seat is contested — must be all-or-nothing. */ +async function partialAvailability(pool, scid) { + const [a, b, c] = pool.splice(0, 3); + const hAll = await hold([a.id, b.id, c.id]); + if (!hAll.holdId) return { name: 'partial-availability', pass: false, detail: `3-seat hold failed: ${msg(hAll)}`, errs: [] }; + + // Book seat B alone first (legitimately, from the same hold). + const first = await book(oneWayBody(hAll.holdId, [b.id], scid)); + if (!first.bookingRef) return { name: 'partial-availability', pass: false, detail: `setup booking failed: ${msg(first)}`, errs: [] }; + note([first.bookingRef]); + + // Now try to book all three. B is taken -> the whole request must fail, leaving no row. + const second = await book(oneWayBody(hAll.holdId, [a.id, b.id, c.id], scid)); + if (second.bookingRef) { + note([second.bookingRef]); + return { name: 'partial-availability', pass: false, detail: `3-seat booking succeeded despite seat ${b.seatNumber} being taken — ${second.bookingRef}`, errs: [] }; + } + return { + name: 'partial-availability', + pass: true, + detail: `seat ${b.seatNumber} taken -> 3-seat request rejected whole (${second.status}: ${msg(second)}); seats ${a.seatNumber},${c.seatNumber} must remain free`, + errs: [], + freeSeats: [a, c], + takenSeat: b, + }; +} + +/** N concurrent ROUND_TRIP bookings contesting one outbound+return seat pair. */ +async function roundTripContest(n, scid) { + if (!RET_SCHED) return { name: `round-trip n=${n}`, pass: true, detail: 'skipped (no --return-schedule)', errs: [] }; + const outPool = await seatPool(SCHED, ORIGIN, DEST); + const retPool = await seatPool(RET_SCHED, RET_ORIGIN, RET_DEST); + const outSeat = outPool[0]; + const retSeat = retPool.find((s) => s.id !== outSeat.id) ?? retPool[0]; + + const traveller = uuid(); + const hOut = await hold([outSeat.id], SCHED, ORIGIN, DEST, 'OUTBOUND', traveller); + const hRet = await hold([retSeat.id], RET_SCHED, RET_ORIGIN, RET_DEST, 'RETURN', traveller); + if (!hOut.holdId || !hRet.holdId) { + return { name: `round-trip n=${n}`, pass: false, detail: `hold failed out=${msg(hOut)} ret=${msg(hRet)}`, errs: [] }; + } + const bodies = Array.from({ length: n }, () => + roundTripBody(hOut.holdId, hRet.holdId, [outSeat.id], [retSeat.id], scid)); + const { results, wallMs } = await fireAll(bodies.map((b) => () => book(b))); + const won = results.filter((r) => r.bookingRef); + note(won.map((r) => r.bookingRef)); + return { + name: `round-trip n=${n}`, + pass: won.length === 1, + detail: `out seat ${outSeat.seatNumber} / ret seat ${retSeat.seatNumber}: ${won.length} booking(s) ` + + `${JSON.stringify(tally(results))} wall=${wallMs}ms refs=${won.map((r) => r.bookingRef).join(',')}`, + errs: errSample(results), + seats: { outSeat, retSeat }, + }; +} + +/** A single multi-passenger, multi-seat booking — the ordinary family booking. */ +async function multiPassenger(pool, scid, count = 4) { + const seats = pool.splice(0, count); + const h = await hold(seats.map((s) => s.id)); + if (!h.holdId) return { name: `multi-passenger x${count}`, pass: false, detail: `hold failed: ${msg(h)}`, errs: [] }; + const r = await book(oneWayBody(h.holdId, seats.map((s) => s.id), scid)); + if (r.bookingRef) note([r.bookingRef]); + return { + name: `multi-passenger x${count}`, + pass: !!r.bookingRef, + detail: r.bookingRef + ? `${count} passengers on seats ${seats.map((s) => s.seatNumber).join(',')} -> ${r.bookingRef}` + : `rejected (${r.status}): ${msg(r)}`, + errs: [], + seats, + ref: r.bookingRef, + }; +} + +// ───────────────────────────────────────────────────────────────────────────── +(async () => { + console.log(`API ${BASE}`); + console.log(`Schedule ${SCHED}`); + if (RET_SCHED) console.log(`Return ${RET_SCHED}`); + const pool = await seatPool(); + const scid = await seatClassId(pool[0]); + console.log(`${pool.length} AVAILABLE seats, seatClass ${scid}\n`); + + const out = []; + const report = (r) => { + out.push(r); + console.log(`${r.pass ? 'PASS' : 'FAIL'} ${r.name}`); + console.log(` ${r.detail}`); + if (r.errs?.length) r.errs.forEach((e) => console.log(` · ${e}`)); + console.log(''); + return r; + }; + + if (MODE === 'ladder' || MODE === 'all') { + for (const n of LEVELS) { + console.log(`──────── concurrency ${n} ────────`); + const hc = report(await holdContest(pool, n)); + if (hc.winner) report(await bookingStorm(hc.seat, hc.winner.holdId, scid, n)); + report(await differentSeats(pool, scid, n)); + } + } + if (MODE === 'edge' || MODE === 'all') { + console.log(`──────── edge cases ────────`); + report(await multiPassenger(pool, scid, 4)); + report(await partialAvailability(pool, scid)); + report(await overlapping(pool, scid, 12)); + report(await roundTripContest(10, scid)); + } + + const fails = out.filter((r) => !r.pass); + console.log(`════════ ${out.length} checks, ${fails.length} failed, ${ALL.created.length} bookings created ════════`); + fails.forEach((f) => console.log(` FAIL ${f.name}: ${f.detail}`)); + if (ALL.created.length) { + console.log(`\nrefs: ${ALL.created.join(' ')}`); + } + process.exit(fails.length ? 1 : 0); +})().catch((e) => { console.error(e); process.exit(2); }); diff --git a/apps/edr-passenger-api/src/modules/bookings/booking-identity.util.ts b/apps/edr-passenger-api/src/modules/bookings/booking-identity.util.ts index cd5385765..236d019a6 100644 --- a/apps/edr-passenger-api/src/modules/bookings/booking-identity.util.ts +++ b/apps/edr-passenger-api/src/modules/bookings/booking-identity.util.ts @@ -8,7 +8,7 @@ import { PrismaService } from '../../common/prisma.service'; * the same person can book that train again. PENDING_PAYMENT counts — otherwise the whole check * is bypassable by simply never finishing the first payment. */ -const ACTIVE_BOOKING_STATUSES: BookingStatus[] = [ +export const ACTIVE_BOOKING_STATUSES: BookingStatus[] = [ BookingStatus.DRAFT, BookingStatus.PENDING_PAYMENT, BookingStatus.CONFIRMED, diff --git a/apps/edr-passenger-api/src/modules/bookings/bookings.service.ts b/apps/edr-passenger-api/src/modules/bookings/bookings.service.ts index 004aca024..12ea1f9d8 100644 --- a/apps/edr-passenger-api/src/modules/bookings/bookings.service.ts +++ b/apps/edr-passenger-api/src/modules/bookings/bookings.service.ts @@ -2,7 +2,7 @@ import { Injectable, NotFoundException, BadRequestException, Logger } from '@nes import { InjectDataSource } from '@nestjs/typeorm'; import { DataSource } from 'typeorm'; import { PrismaService } from '../../common/prisma.service'; -import { SeatsService } from '../seats/seats.service'; +import { SeatsService, SeatClaim } from '../seats/seats.service'; import { TicketsService } from '../tickets/tickets.service'; import { EventEmitter2 } from '@nestjs/event-emitter'; import { CreateBookingDto } from './bookings.dto'; @@ -873,15 +873,10 @@ export class BookingsService { return this.createOneWayBooking(dto); } - private validateSeatIdsAgainstHold(holdId: string, holdSeatIds: string[], requestedSeatIds: string[]) { - for (const seatId of requestedSeatIds) { - if (!holdSeatIds.includes(seatId)) { - throw new BadRequestException( - `Seat ${seatId} is not part of hold ${holdId}. Use seat IDs returned from POST /seats/hold.`, - ); - } - } - } + // validateSeatIdsAgainstHold lived here. Every caller now goes through + // SeatsService.assertSeatsClaimable, which performs the same seat-in-hold check alongside + // the hold-is-live, hold-is-for-this-schedule and no-conflict checks — and, at write time, + // re-runs all of them inside the seat lock. /** Resolves contactEmail/contactPhone for an IAM-authenticated passenger booking. */ private async resolveIamContact(passengerId?: string): Promise<{ contactEmail: string | null; contactPhone: string | null }> { @@ -896,11 +891,18 @@ export class BookingsService { } private async createOneWayBooking(dto: CreateBookingDto) { - const hold = await this.prisma.seatHold.findUnique({ where: { id: dto.holdId } }); - if (!hold || hold.expiresAt < new Date()) throw new BadRequestException('Seat hold expired'); - - const requestedSeatIds = (dto.passengers as any[]).map(p => p.seatId); - this.validateSeatIdsAgainstHold(dto.holdId, hold.seatIds, requestedSeatIds); + // Hold live, on this departure, covering these seats, and nobody else holding or booked + // on them. Advisory here — re-run under the seat lock at the write below. + const requestedSeatIds = (dto.passengers as any[]).map(p => p.seatId).filter((id: any): id is string => !!id); + const claims: SeatClaim[] = [{ + holdId: dto.holdId, + seatIds: requestedSeatIds, + scheduleId: dto.scheduleId, + originStationId: dto.originStationId, + destinationStationId: dto.destinationStationId, + journeyDirection: JourneyDirection.ONE_WAY, + }]; + await this.seatsService.assertSeatsClaimableAll(claims); const schedule = await this.prisma.trainSchedule.findUnique({ where: { id: dto.scheduleId }, @@ -912,14 +914,6 @@ export class BookingsService { const destStop = schedule.stopTimes.find(s => s.stationId === dto.destinationStationId); if (!originStop || !destStop) throw new NotFoundException('Origin or destination not found'); - await this.seatsService.assertNoRouteSeatConflict({ - scheduleId: dto.scheduleId, - seatIds: requestedSeatIds, - originStationId: dto.originStationId, - destinationStationId: dto.destinationStationId, - journeyDirection: JourneyDirection.ONE_WAY, - }); - const [passengersData, iamContact] = await Promise.all([ this.processPassengers(dto.passengers as any[]), this.resolveIamContact(dto.passengerId), @@ -1029,7 +1023,9 @@ export class BookingsService { // still pass; the tolerance absorbs FX-conversion rounding. this.assertTotalNotUnderAuthoritative(resolvedTotalMinor, fareCalculation.totalMinor, 'createOneWayBooking'); - const booking = await this.prisma.booking.create({ + // Locked and re-validated inside the lock — see SeatsService.claimSeatsAndWrite. All the + // slow work (fare engine, FX, identity) is already done, so this transaction stays short. + const booking = await this.seatsService.claimSeatsAndWrite(claims, async (tx) => tx.booking.create({ data: { bookingRef: generateRef(), passengerId: dto.passengerId, @@ -1068,7 +1064,7 @@ export class BookingsService { } }, include: { seats: { include: { seat: true } }, schedule: { include: { originStation: true, destinationStation: true, train: true } } } - }); + })); await this.seatsService.confirmSeats(passengersData.map(p => p.seatId)); if (dto.packageId && dto.priceTierId) { @@ -1087,18 +1083,31 @@ export class BookingsService { throw new BadRequestException('Return trip details required for round-trip booking'); } - const [outboundHold, returnHold] = await Promise.all([ - this.prisma.seatHold.findUnique({ where: { id: dto.holdId } }), - this.prisma.seatHold.findUnique({ where: { id: dto.returnHoldId } }) - ]); - - if (!outboundHold || outboundHold.expiresAt < new Date()) throw new BadRequestException('Outbound seat hold expired'); - if (!returnHold || returnHold.expiresAt < new Date()) throw new BadRequestException('Return seat hold expired'); - const holdObSeatIds = (dto.passengers as any[]).map((p: any) => p.seatId ?? p.outboundSeatId).filter(Boolean); const holdRetSeatIds = (dto.passengers as any[]).map((p: any) => p.returnSeatId).filter(Boolean); - if (holdObSeatIds.length) this.validateSeatIdsAgainstHold(dto.holdId, outboundHold.seatIds, holdObSeatIds); - if (holdRetSeatIds.length) this.validateSeatIdsAgainstHold(dto.returnHoldId!, returnHold.seatIds, holdRetSeatIds); + + // Advisory pass; re-run under the seat lock at the write. + const claims: SeatClaim[] = [ + { + holdId: dto.holdId, + seatIds: holdObSeatIds, + scheduleId: dto.scheduleId, + originStationId: dto.originStationId, + destinationStationId: dto.destinationStationId, + journeyDirection: JourneyDirection.OUTBOUND, + legLabel: 'Outbound', + }, + { + holdId: dto.returnHoldId!, + seatIds: holdRetSeatIds, + scheduleId: dto.returnScheduleId!, + originStationId: dto.returnOriginStationId, + destinationStationId: dto.returnDestinationStationId, + journeyDirection: JourneyDirection.RETURN, + legLabel: 'Return', + }, + ]; + await this.seatsService.assertSeatsClaimableAll(claims); const [outboundSchedule, returnSchedule] = await Promise.all([ this.prisma.trainSchedule.findUnique({ @@ -1122,21 +1131,6 @@ export class BookingsService { throw new NotFoundException('Origin or destination stops not found'); } - await this.seatsService.assertNoRouteSeatConflict({ - scheduleId: dto.scheduleId, - seatIds: holdObSeatIds, - originStationId: dto.originStationId, - destinationStationId: dto.destinationStationId, - journeyDirection: JourneyDirection.OUTBOUND, - }); - await this.seatsService.assertNoRouteSeatConflict({ - scheduleId: dto.returnScheduleId, - seatIds: holdRetSeatIds, - originStationId: dto.returnOriginStationId, - destinationStationId: dto.returnDestinationStationId, - journeyDirection: JourneyDirection.RETURN, - }); - const [passengersData, iamContact] = await Promise.all([ this.processRoundTripPassengers(dto.passengers as any[]), this.resolveIamContact(dto.passengerId), @@ -1252,7 +1246,8 @@ export class BookingsService { // C-1 guard: never charge less than the server-recomputed authoritative round-trip fare. this.assertTotalNotUnderAuthoritative(totalMinor, authoritativeTotalMinor, 'createRoundTripBooking'); - const booking = await this.prisma.booking.create({ + // Locked across both legs' seats and re-validated inside the lock. + const booking = await this.seatsService.claimSeatsAndWrite(claims, async (tx) => tx.booking.create({ data: { bookingRef: generateRef(), passengerId: dto.passengerId, @@ -1314,7 +1309,7 @@ export class BookingsService { }, } as any, include: { seats: { include: { seat: true } }, schedule: { include: { originStation: true, destinationStation: true, train: true } } } - }); + })); const outboundSeatIds = passengersData.map(p => p.outboundSeatId); const returnSeatIds = passengersData.map(p => p.returnSeatId); @@ -1355,17 +1350,31 @@ export class BookingsService { throw new BadRequestException('leg2ScheduleId, leg2HoldId, transitStationId and leg2DestinationStationId are required for TRANSIT bookings'); } - const [leg1Hold, leg2Hold] = await Promise.all([ - this.prisma.seatHold.findUnique({ where: { id: dto.holdId } }), - this.prisma.seatHold.findUnique({ where: { id: dto.leg2HoldId } }), - ]); - if (!leg1Hold || leg1Hold.expiresAt < new Date()) throw new BadRequestException('Leg-1 seat hold expired'); - if (!leg2Hold || leg2Hold.expiresAt < new Date()) throw new BadRequestException('Leg-2 seat hold expired'); + const leg1SeatIds = (dto.passengers as any[]).map(p => p.seatId).filter(Boolean); + const leg2SeatIds = (dto.passengers as any[]).map(p => p.leg2SeatId ?? p.seatId).filter(Boolean); - const leg1SeatIds = (dto.passengers as any[]).map(p => p.seatId); - const leg2SeatIds = (dto.passengers as any[]).map(p => p.leg2SeatId ?? p.seatId); - this.validateSeatIdsAgainstHold(dto.holdId, leg1Hold.seatIds, leg1SeatIds); - this.validateSeatIdsAgainstHold(dto.leg2HoldId!, leg2Hold.seatIds, leg2SeatIds); + // Advisory pass; re-run under the seat lock at the write. + const claims: SeatClaim[] = [ + { + holdId: dto.holdId, + seatIds: leg1SeatIds, + scheduleId: dto.scheduleId, + originStationId: dto.originStationId, + destinationStationId: dto.transitStationId, + journeyDirection: JourneyDirection.ONE_WAY, + legLabel: 'Leg-1', + }, + { + holdId: dto.leg2HoldId!, + seatIds: leg2SeatIds, + scheduleId: dto.leg2ScheduleId!, + originStationId: dto.transitStationId, + destinationStationId: dto.leg2DestinationStationId, + journeyDirection: JourneyDirection.ONE_WAY, + legLabel: 'Leg-2', + }, + ]; + await this.seatsService.assertSeatsClaimableAll(claims); const [leg1Schedule, leg2Schedule] = await Promise.all([ this.prisma.trainSchedule.findUnique({ @@ -1387,21 +1396,6 @@ export class BookingsService { if (!leg1OriginStop || !leg1DestStop) throw new NotFoundException('Leg-1 origin or transit station not found on schedule'); if (!leg2OriginStop || !leg2DestStop) throw new NotFoundException('Transit or leg-2 destination station not found on leg-2 schedule'); - await this.seatsService.assertNoRouteSeatConflict({ - scheduleId: dto.scheduleId, - seatIds: leg1SeatIds, - originStationId: dto.originStationId, - destinationStationId: dto.transitStationId, - journeyDirection: JourneyDirection.ONE_WAY, - }); - await this.seatsService.assertNoRouteSeatConflict({ - scheduleId: dto.leg2ScheduleId, - seatIds: leg2SeatIds, - originStationId: dto.transitStationId, - destinationStationId: dto.leg2DestinationStationId, - journeyDirection: JourneyDirection.ONE_WAY, - }); - const [passengersData, iamContact] = await Promise.all([ this.processPassengers(dto.passengers as any[]), this.resolveIamContact(dto.passengerId), @@ -1467,7 +1461,8 @@ export class BookingsService { }); // Single booking — leg-1 seats at leg=1, leg-2 seats at leg=2 - const booking = await this.prisma.booking.create({ + // Locked across both legs' seats and re-validated inside the lock. + const booking = await this.seatsService.claimSeatsAndWrite(claims, async (tx) => tx.booking.create({ data: { bookingRef: generateRef(), passengerId: dto.passengerId, @@ -1526,7 +1521,7 @@ export class BookingsService { }, } as any, include: { seats: { include: { seat: true } }, schedule: { include: { originStation: true, destinationStation: true, train: true } } }, - }); + })); await Promise.all([ this.seatsService.confirmSeats(passengersData.map(p => p.seatId)), @@ -1561,23 +1556,49 @@ export class BookingsService { ); } - // Validate all 4 holds - const [obL1Hold, obL2Hold, retL1Hold, retL2Hold] = await Promise.all([ - this.prisma.seatHold.findUnique({ where: { id: dto.holdId } }), - this.prisma.seatHold.findUnique({ where: { id: dto.leg2HoldId } }), - this.prisma.seatHold.findUnique({ where: { id: dto.returnHoldId } }), - this.prisma.seatHold.findUnique({ where: { id: dto.returnLeg2HoldId } }), - ]); - const now = new Date(); - if (!obL1Hold || obL1Hold.expiresAt < now) throw new BadRequestException('Outbound leg-1 seat hold expired'); - if (!obL2Hold || obL2Hold.expiresAt < now) throw new BadRequestException('Outbound leg-2 seat hold expired'); - if (!retL1Hold || retL1Hold.expiresAt < now) throw new BadRequestException('Return leg-1 seat hold expired'); - if (!retL2Hold || retL2Hold.expiresAt < now) throw new BadRequestException('Return leg-2 seat hold expired'); - - this.validateSeatIdsAgainstHold(dto.holdId, obL1Hold.seatIds, (dto.passengers as any[]).map(p => p.seatId)); - this.validateSeatIdsAgainstHold(dto.leg2HoldId!, obL2Hold.seatIds, (dto.passengers as any[]).map(p => p.leg2SeatId ?? p.seatId)); - this.validateSeatIdsAgainstHold(dto.returnHoldId!, retL1Hold.seatIds, (dto.passengers as any[]).map(p => p.returnSeatId)); - this.validateSeatIdsAgainstHold(dto.returnLeg2HoldId!, retL2Hold.seatIds, (dto.passengers as any[]).map(p => p.returnLeg2SeatId ?? p.returnSeatId)); + // All 4 holds: live, on their own departure, covering their leg's seats, unconflicted. + // Advisory pass; re-run under the seat lock at the write. + const seatsOf = (pick: (p: any) => string | undefined) => + (dto.passengers as any[]).map(pick).filter((id): id is string => !!id); + const claims: SeatClaim[] = [ + { + holdId: dto.holdId, + seatIds: seatsOf(p => p.seatId), + scheduleId: dto.scheduleId, + originStationId: dto.originStationId, + destinationStationId: dto.transitStationId, + journeyDirection: JourneyDirection.OUTBOUND, + legLabel: 'Outbound leg-1', + }, + { + holdId: dto.leg2HoldId!, + seatIds: seatsOf(p => p.leg2SeatId ?? p.seatId), + scheduleId: dto.leg2ScheduleId!, + originStationId: dto.transitStationId, + destinationStationId: dto.leg2DestinationStationId, + journeyDirection: JourneyDirection.OUTBOUND, + legLabel: 'Outbound leg-2', + }, + { + holdId: dto.returnHoldId!, + seatIds: seatsOf(p => p.returnSeatId), + scheduleId: dto.returnScheduleId!, + originStationId: dto.returnOriginStationId, + destinationStationId: dto.returnTransitStationId, + journeyDirection: JourneyDirection.RETURN, + legLabel: 'Return leg-1', + }, + { + holdId: dto.returnLeg2HoldId!, + seatIds: seatsOf(p => p.returnLeg2SeatId ?? p.returnSeatId), + scheduleId: dto.returnLeg2ScheduleId!, + originStationId: dto.returnTransitStationId, + destinationStationId: dto.returnLeg2DestinationStationId, + journeyDirection: JourneyDirection.RETURN, + legLabel: 'Return leg-2', + }, + ]; + await this.seatsService.assertSeatsClaimableAll(claims); // Load all 4 schedules const [obL1Sched, obL2Sched, retL1Sched, retL2Sched] = await Promise.all([ @@ -1604,35 +1625,6 @@ export class BookingsService { if (!retL1Origin || !retL1Dest) throw new NotFoundException('Return leg-1: origin or transit station not found'); if (!retL2Origin || !retL2Dest) throw new NotFoundException('Return leg-2: transit or destination not found'); - await this.seatsService.assertNoRouteSeatConflict({ - scheduleId: dto.scheduleId, - seatIds: (dto.passengers as any[]).map(p => p.seatId), - originStationId: dto.originStationId, - destinationStationId: dto.transitStationId, - journeyDirection: JourneyDirection.OUTBOUND, - }); - await this.seatsService.assertNoRouteSeatConflict({ - scheduleId: dto.leg2ScheduleId, - seatIds: (dto.passengers as any[]).map(p => p.leg2SeatId ?? p.seatId), - originStationId: dto.transitStationId, - destinationStationId: dto.leg2DestinationStationId, - journeyDirection: JourneyDirection.OUTBOUND, - }); - await this.seatsService.assertNoRouteSeatConflict({ - scheduleId: dto.returnScheduleId, - seatIds: (dto.passengers as any[]).map(p => p.returnSeatId), - originStationId: dto.returnOriginStationId, - destinationStationId: dto.returnTransitStationId, - journeyDirection: JourneyDirection.RETURN, - }); - await this.seatsService.assertNoRouteSeatConflict({ - scheduleId: dto.returnLeg2ScheduleId, - seatIds: (dto.passengers as any[]).map(p => p.returnLeg2SeatId ?? p.returnSeatId), - originStationId: dto.returnTransitStationId, - destinationStationId: dto.returnLeg2DestinationStationId, - journeyDirection: JourneyDirection.RETURN, - }); - const [passengersData, iamContact] = await Promise.all([ this.processRoundTripPassengers(dto.passengers as any[]), this.resolveIamContact(dto.passengerId), @@ -1718,7 +1710,8 @@ export class BookingsService { displayCurrency, }); - const booking = await this.prisma.booking.create({ + // Locked across all four legs' seats and re-validated inside the lock. + const booking = await this.seatsService.claimSeatsAndWrite(claims, async (tx) => tx.booking.create({ data: { bookingRef: generateRef(), passengerId: dto.passengerId, @@ -1759,7 +1752,7 @@ export class BookingsService { }, } as any, include: { seats: { include: { seat: true } }, schedule: { include: { originStation: true, destinationStation: true, train: true } } }, - }); + })); await Promise.all([ this.seatsService.confirmSeats(passengersData.map(p => p.outboundSeatId)), diff --git a/apps/edr-passenger-api/src/modules/bookings/guest-booking.service.ts b/apps/edr-passenger-api/src/modules/bookings/guest-booking.service.ts index cf5d92b82..b33c2d8fe 100644 --- a/apps/edr-passenger-api/src/modules/bookings/guest-booking.service.ts +++ b/apps/edr-passenger-api/src/modules/bookings/guest-booking.service.ts @@ -2,7 +2,7 @@ import { Injectable, BadRequestException, NotFoundException, Logger } from '@nes import { InjectDataSource } from '@nestjs/typeorm'; import { DataSource } from 'typeorm'; import { PrismaService } from '../../common/prisma.service'; -import { SeatsService } from '../seats/seats.service'; +import { SeatsService, SeatClaim } from '../seats/seats.service'; import { VerifaydaService } from '../verifayda/verifayda.service'; import { CurrencyService } from '../currency/currency.service'; import { PassengerAuthService } from '../auth/passenger-auth.service'; @@ -146,11 +146,20 @@ export class GuestBookingService { private async createGuestOneWayBooking(dto: CreateGuestBookingDto, req?: any) { const authUserId: string | null = req?.user?.id ?? null; - // Validate hold - const hold = await this.prisma.seatHold.findUnique({ where: { id: dto.holdId } }); - if (!hold || hold.expiresAt < new Date()) { - throw new BadRequestException('Seat hold expired or not found'); - } + const claimedSeatIds = dto.passengers.map((p) => p.seatId).filter((id): id is string => !!id); + // Fail fast, before any fare or identity work: hold is live, belongs to this departure, + // covers the seats asked for, and nobody else holds or has booked them. This is a + // courtesy check for a clean early error — the authoritative one runs under the seat + // lock at the write below, because anything checked out here can change before we write. + const claims: SeatClaim[] = [{ + holdId: dto.holdId, + seatIds: claimedSeatIds, + scheduleId: dto.scheduleId, + originStationId: dto.originStationId, + destinationStationId: dto.destinationStationId, + journeyDirection: JourneyDirection.ONE_WAY, + }]; + await this.seatsService.assertSeatsClaimableAll(claims); // Get schedule const schedule = await this.prisma.trainSchedule.findUnique({ @@ -371,8 +380,15 @@ export class GuestBookingService { } } - // Create booking - const booking = await this.prisma.booking.create({ + // Create booking. + // + // Locked, and re-validated inside the lock. Everything slow — fare engine, FX, Verifayda, + // guest-passenger resolution — is already done above, so this transaction is two reads and + // a write. Re-checking here is the whole point: the courtesy check at the top of the method + // ran hundreds of milliseconds ago, and between then and now another request holding the + // same hold (a double-submit, a retry after a timeout) could have booked these seats. + const booking = await this.seatsService.claimSeatsAndWrite(claims, async (tx) => { + return tx.booking.create({ data: { bookingRef: generateRef(), passengerId: guestPassengerId, @@ -418,6 +434,7 @@ export class GuestBookingService { seats: { include: { seat: { include: { coach: true } } } }, schedule: { include: { originStation: true, destinationStation: true, train: true } }, }, + }); }); // Save passenger details as traveler profiles — guest bookings only. @@ -702,19 +719,35 @@ export class GuestBookingService { throw new BadRequestException('returnScheduleId, returnHoldId, returnOriginStationId and returnDestinationStationId are required for ROUND_TRIP'); } - // Validate both holds - const [outboundHold, returnHold] = await Promise.all([ - this.prisma.seatHold.findUnique({ where: { id: dto.holdId } }), - this.prisma.seatHold.findUnique({ where: { id: dto.returnHoldId } }), - ]); - if (!outboundHold || outboundHold.expiresAt < new Date()) throw new BadRequestException('Outbound seat hold expired or not found'); - if (!returnHold || returnHold.expiresAt < new Date()) throw new BadRequestException('Return seat hold expired or not found'); - // Validate passengers have returnSeatId for (const p of dto.passengers) { if (!p.returnSeatId) throw new BadRequestException(`returnSeatId is required for each passenger in a ROUND_TRIP booking (missing for ${p.passengerName})`); } + // Each leg's hold must be live, belong to that leg's departure, cover that leg's seats, + // and clear the conflict check. Advisory here; re-run under the seat lock at the write. + const claims: SeatClaim[] = [ + { + holdId: dto.holdId, + seatIds: dto.passengers.map((p) => p.seatId).filter((id): id is string => !!id), + scheduleId: dto.scheduleId, + originStationId: dto.originStationId, + destinationStationId: dto.destinationStationId, + journeyDirection: JourneyDirection.OUTBOUND, + legLabel: 'Outbound', + }, + { + holdId: dto.returnHoldId, + seatIds: dto.passengers.map((p) => p.returnSeatId).filter((id): id is string => !!id), + scheduleId: dto.returnScheduleId, + originStationId: dto.returnOriginStationId, + destinationStationId: dto.returnDestinationStationId, + journeyDirection: JourneyDirection.RETURN, + legLabel: 'Return', + }, + ]; + await this.seatsService.assertSeatsClaimableAll(claims); + // Load both schedules const [outboundSchedule, returnSchedule] = await Promise.all([ this.prisma.trainSchedule.findUnique({ @@ -921,7 +954,10 @@ export class GuestBookingService { const outboundSeatIds = dto.passengers.map(p => p.seatId).filter((id): id is string => !!id); const returnSeatIds = dto.passengers.map(p => p.returnSeatId!); - const booking = await this.prisma.booking.create({ + // Locked across BOTH legs' seats, with each leg's claim re-checked inside the lock — + // so a round trip is all-or-nothing: it never commits with the outbound seat secured + // and the return seat sold out from under it. + const booking = await this.seatsService.claimSeatsAndWrite(claims, async (tx) => tx.booking.create({ data: { bookingRef: generateRef(), passengerId: guestPassengerId, @@ -987,7 +1023,7 @@ export class GuestBookingService { seats: { include: { seat: { include: { coach: true } } } }, schedule: { include: { originStation: true, destinationStation: true, train: true } }, }, - }); + })); // Traveler profiles: guest bookings only (authenticated passengers already have one). if (!authUserId) await this.createTravelerProfiles(guestPassengerId, passengersData); @@ -1032,17 +1068,32 @@ export class GuestBookingService { throw new BadRequestException('leg2ScheduleId, leg2HoldId, transitStationId and leg2DestinationStationId are required for TRANSIT bookings'); } - const [leg1Hold, leg2Hold] = await Promise.all([ - this.prisma.seatHold.findUnique({ where: { id: dto.holdId } }), - this.prisma.seatHold.findUnique({ where: { id: dto.leg2HoldId } }), - ]); - if (!leg1Hold || leg1Hold.expiresAt < new Date()) throw new BadRequestException('Leg-1 seat hold expired or not found'); - if (!leg2Hold || leg2Hold.expiresAt < new Date()) throw new BadRequestException('Leg-2 seat hold expired or not found'); - for (const p of dto.passengers) { if (!p.leg2SeatId) throw new BadRequestException(`leg2SeatId is required for each passenger in a TRANSIT booking (missing for ${p.passengerName})`); } + const claims: SeatClaim[] = [ + { + holdId: dto.holdId, + seatIds: dto.passengers.map((p) => p.seatId).filter((id): id is string => !!id), + scheduleId: dto.scheduleId, + originStationId: dto.originStationId, + destinationStationId: dto.transitStationId, + journeyDirection: JourneyDirection.ONE_WAY, + legLabel: 'Leg-1', + }, + { + holdId: dto.leg2HoldId, + seatIds: dto.passengers.map((p) => p.leg2SeatId).filter((id): id is string => !!id), + scheduleId: dto.leg2ScheduleId, + originStationId: dto.transitStationId, + destinationStationId: dto.leg2DestinationStationId, + journeyDirection: JourneyDirection.ONE_WAY, + legLabel: 'Leg-2', + }, + ]; + await this.seatsService.assertSeatsClaimableAll(claims); + const [leg1Schedule, leg2Schedule] = await Promise.all([ this.prisma.trainSchedule.findUnique({ where: { id: dto.scheduleId }, @@ -1140,8 +1191,9 @@ export class GuestBookingService { const { guestPassengerId, iamUserId, createdAccount } = await this.resolveGuestPassenger(dto, passengersData[0], req); const contact = await this.resolveActorContact(req, passengersData[0]); - // Single booking — leg-1 seats at leg=1, leg-2 seats at leg=2 - const booking = await this.prisma.booking.create({ + // Single booking — leg-1 seats at leg=1, leg-2 seats at leg=2. + // Locked across both legs' seats and re-checked inside the lock. + const booking = await this.seatsService.claimSeatsAndWrite(claims, async (tx) => tx.booking.create({ data: { bookingRef: generateRef(), passengerId: guestPassengerId, @@ -1204,7 +1256,7 @@ export class GuestBookingService { seats: { include: { seat: { include: { coach: true } } } }, schedule: { include: { originStation: true, destinationStation: true, train: true } }, }, - }); + })); // Traveler profiles: guest bookings only (authenticated passengers already have one). if (!authUserId) await this.createTravelerProfiles(guestPassengerId, passengersData); @@ -1247,17 +1299,48 @@ export class GuestBookingService { if (!p.returnLeg2SeatId) throw new BadRequestException(`returnLeg2SeatId required for ${p.passengerName}`); } - const now = new Date(); - const [obL1Hold, obL2Hold, retL1Hold, retL2Hold] = await Promise.all([ - this.prisma.seatHold.findUnique({ where: { id: dto.holdId } }), - this.prisma.seatHold.findUnique({ where: { id: dto.leg2HoldId } }), - this.prisma.seatHold.findUnique({ where: { id: dto.returnHoldId } }), - this.prisma.seatHold.findUnique({ where: { id: dto.returnLeg2HoldId } }), - ]); - if (!obL1Hold || obL1Hold.expiresAt < now) throw new BadRequestException('Outbound leg-1 hold expired'); - if (!obL2Hold || obL2Hold.expiresAt < now) throw new BadRequestException('Outbound leg-2 hold expired'); - if (!retL1Hold || retL1Hold.expiresAt < now) throw new BadRequestException('Return leg-1 hold expired'); - if (!retL2Hold || retL2Hold.expiresAt < now) throw new BadRequestException('Return leg-2 hold expired'); + const seatsOf = (pick: (p: (typeof dto.passengers)[number]) => string | undefined) => + dto.passengers.map(pick).filter((id): id is string => !!id); + + const claims: SeatClaim[] = [ + { + holdId: dto.holdId, + seatIds: seatsOf((p) => p.seatId), + scheduleId: dto.scheduleId, + originStationId: dto.originStationId, + destinationStationId: dto.transitStationId, + journeyDirection: JourneyDirection.OUTBOUND, + legLabel: 'Outbound leg-1', + }, + { + holdId: dto.leg2HoldId, + seatIds: seatsOf((p) => p.leg2SeatId), + scheduleId: dto.leg2ScheduleId, + originStationId: dto.transitStationId, + destinationStationId: dto.leg2DestinationStationId, + journeyDirection: JourneyDirection.OUTBOUND, + legLabel: 'Outbound leg-2', + }, + { + holdId: dto.returnHoldId, + seatIds: seatsOf((p) => p.returnSeatId), + scheduleId: dto.returnScheduleId, + originStationId: dto.returnOriginStationId, + destinationStationId: dto.returnTransitStationId, + journeyDirection: JourneyDirection.RETURN, + legLabel: 'Return leg-1', + }, + { + holdId: dto.returnLeg2HoldId, + seatIds: seatsOf((p) => p.returnLeg2SeatId), + scheduleId: dto.returnLeg2ScheduleId, + originStationId: dto.returnTransitStationId, + destinationStationId: dto.returnLeg2DestinationStationId, + journeyDirection: JourneyDirection.RETURN, + legLabel: 'Return leg-2', + }, + ]; + await this.seatsService.assertSeatsClaimableAll(claims); const [obL1Sched, obL2Sched, retL1Sched, retL2Sched] = await Promise.all([ this.prisma.trainSchedule.findUnique({ where: { id: dto.scheduleId }, include: { originStation: true, destinationStation: true, stopTimes: { include: { station: true }, orderBy: { sequence: 'asc' } }, route: { include: { stops: true } } } }), @@ -1371,7 +1454,8 @@ export class GuestBookingService { displayCurrency, }); - const booking = await this.prisma.booking.create({ + // Locked across all four legs' seats and re-checked inside the lock. + const booking = await this.seatsService.claimSeatsAndWrite(claims, async (tx) => tx.booking.create({ data: { bookingRef: generateRef(), passengerId: guestPassengerId, @@ -1410,7 +1494,7 @@ export class GuestBookingService { seats: { include: { seat: { include: { coach: true } } } }, schedule: { include: { originStation: true, destinationStation: true, train: true } }, }, - }); + })); // Traveler profiles: guest bookings only (authenticated passengers already have one). if (!authUserId) await this.createTravelerProfiles(guestPassengerId, passengersData); diff --git a/apps/edr-passenger-api/src/modules/seats/seats.service.double-booking.spec.ts b/apps/edr-passenger-api/src/modules/seats/seats.service.double-booking.spec.ts new file mode 100644 index 000000000..3d03e8dd4 --- /dev/null +++ b/apps/edr-passenger-api/src/modules/seats/seats.service.double-booking.spec.ts @@ -0,0 +1,217 @@ +import { Test, TestingModule } from '@nestjs/testing'; +import { SeatsService } from './seats.service'; +import { PrismaService } from '../../common/prisma.service'; +import { SegmentsService } from '../segments/segments.service'; +import { SystemConfigService } from '../system-config/system-config.service'; +import { AuditService } from '../../common/audit.service'; +import { SmsClientService } from '../notifications/sms-client.service'; + +/** + * Regression guard for the double-booking fix. + * + * Every case below was reproduced against the running API before the fix and must stay + * closed. These are unit tests: they prove the guards exist, are wired in the right order, + * and read through the transaction client. They do NOT prove concurrency safety — only a + * real database can, and that proof lives in scripts/stress-booking-concurrency.cjs. + */ +describe('SeatsService — double-booking guards', () => { + let service: SeatsService; + + const seatHold = { findUnique: jest.fn() }; + const bookingSeat = { findMany: jest.fn() }; + const tripStopTime = { findMany: jest.fn() }; + const seat = { findMany: jest.fn() }; + const queryRaw = jest.fn(); + // SET LOCAL lock_timeout — see SeatsService.lockSeatsForUpdate. + const executeRawUnsafe = jest.fn(); + + const tx: any = { seatHold, bookingSeat, tripStopTime, seat, $queryRaw: queryRaw, $executeRawUnsafe: executeRawUnsafe }; + const mockPrisma: any = { + seatHold, bookingSeat, tripStopTime, seat, $queryRaw: queryRaw, $executeRawUnsafe: executeRawUnsafe, + $transaction: jest.fn((fn: any) => fn(tx)), + }; + const mockSegments = { getSeatAvailabilityMap: jest.fn() }; + + const LIVE_HOLD = { + id: 'hold-1', + scheduleId: 'sched-1', + seatIds: ['seat-1', 'seat-2'], + passengerId: 'pax-1', + expiresAt: new Date(Date.now() + 60_000), + }; + const baseArgs = { + holdId: 'hold-1', + seatIds: ['seat-1'], + scheduleId: 'sched-1', + originStationId: 'a', + destinationStationId: 'c', + }; + + /** A BookingSeat row owned by someone else, spanning a -> c on this schedule. */ + const ownedByOther = [{ + seatId: 'seat-1', + seat: { seatNumber: '7' }, + booking: { + bookingRef: 'AAA111', scheduleId: 'sched-1', + originStationId: 'a', destinationStationId: 'c', + returnScheduleId: null, returnOriginStationId: null, returnDestinationStationId: null, + leg2ScheduleId: null, leg2OriginStationId: null, leg2DestinationStationId: null, + returnLeg2ScheduleId: null, returnLeg2OriginStationId: null, returnLeg2DestStationId: null, + }, + }]; + + beforeEach(async () => { + const module: TestingModule = await Test.createTestingModule({ + providers: [ + SeatsService, + { provide: PrismaService, useValue: mockPrisma }, + { provide: SegmentsService, useValue: mockSegments }, + { provide: SystemConfigService, useValue: { getNumber: jest.fn().mockResolvedValue(5) } }, + { provide: AuditService, useValue: { log: jest.fn() } }, + { provide: SmsClientService, useValue: { send: jest.fn() } }, + ], + }).compile(); + service = module.get(SeatsService); + jest.clearAllMocks(); + + seatHold.findUnique.mockResolvedValue(LIVE_HOLD); + bookingSeat.findMany.mockResolvedValue([]); + tripStopTime.findMany.mockResolvedValue([ + { stationId: 'a', sequence: 1 }, + { stationId: 'b', sequence: 2 }, + { stationId: 'c', sequence: 3 }, + ]); + seat.findMany.mockResolvedValue([{ id: 'seat-1', seatNumber: '7' }]); + mockSegments.getSeatAvailabilityMap.mockResolvedValue(new Map()); + queryRaw.mockResolvedValue([]); + executeRawUnsafe.mockResolvedValue(0); + }); + + describe('assertSeatsClaimable', () => { + it('accepts a seat the presented hold actually covers', async () => { + await expect(service.assertSeatsClaimable(baseArgs)).resolves.toBeUndefined(); + }); + + it('rejects a seat that is not part of the hold (hold laundering)', async () => { + await expect(service.assertSeatsClaimable({ ...baseArgs, seatIds: ['seat-99'] })) + .rejects.toThrow(/not part of hold/i); + }); + + it('rejects a hold taken on a different departure — one Seat row is reused across schedules', async () => { + await expect(service.assertSeatsClaimable({ ...baseArgs, scheduleId: 'sched-OTHER' })) + .rejects.toThrow(/different departure/i); + }); + + it('rejects an expired hold', async () => { + seatHold.findUnique.mockResolvedValue({ ...LIVE_HOLD, expiresAt: new Date(Date.now() - 1) }); + await expect(service.assertSeatsClaimable(baseArgs)).rejects.toThrow(/expired or not found/i); + }); + + it('rejects a hold that does not exist', async () => { + seatHold.findUnique.mockResolvedValue(null); + await expect(service.assertSeatsClaimable(baseArgs)).rejects.toThrow(/expired or not found/i); + }); + }); + + describe('active-booking conflict — enabled for booking writes only', () => { + it('rejects a seat an active booking already owns on an overlapping leg', async () => { + bookingSeat.findMany.mockResolvedValue(ownedByOther); + await expect(service.assertSeatsClaimable({ ...baseArgs, alsoRejectActiveBookings: true })) + .rejects.toThrow(/just booked by someone else/i); + }); + + it('allows a seat whose existing booking covers a non-overlapping stretch', async () => { + // existing leg b->c is [2,3); requested a->b is [1,2) — they only touch, so no conflict + bookingSeat.findMany.mockResolvedValue([{ + ...ownedByOther[0], + booking: { ...ownedByOther[0].booking, originStationId: 'b', destinationStationId: 'c' }, + }]); + await expect(service.assertSeatsClaimable({ + ...baseArgs, destinationStationId: 'b', alsoRejectActiveBookings: true, + })).resolves.toBeUndefined(); + }); + + it('does NOT run for plain holds — reschedule/upgrade legitimately re-hold their own seat', async () => { + bookingSeat.findMany.mockResolvedValue(ownedByOther); + await expect(service.assertSeatsClaimable(baseArgs)).resolves.toBeUndefined(); + expect(bookingSeat.findMany).not.toHaveBeenCalled(); + }); + + it('blocks conservatively when the existing booking leg cannot be resolved', async () => { + bookingSeat.findMany.mockResolvedValue([{ + ...ownedByOther[0], + booking: { ...ownedByOther[0].booking, originStationId: null, destinationStationId: null }, + }]); + await expect(service.assertSeatsClaimable({ ...baseArgs, alsoRejectActiveBookings: true })) + .rejects.toThrow(/just booked by someone else/i); + }); + + it('reads BookingSeat through the transaction client, not the pool', async () => { + await service.claimSeatsAndWrite([baseArgs], async () => ({ id: 'b1' }) as any); + expect(bookingSeat.findMany).toHaveBeenCalled(); + // same mock object is shared by tx and pool here, so assert the availability map — + // whose last positional argument is the client — received the transaction client. + const call = mockSegments.getSeatAvailabilityMap.mock.calls[0]; + expect(call[call.length - 1]).toBe(tx); + }); + }); + + describe('claimSeatsAndWrite', () => { + it('locks the seats FOR UPDATE, then re-checks, then writes — in that order', async () => { + const order: string[] = []; + queryRaw.mockImplementation(() => { order.push('lock'); return Promise.resolve([]); }); + seatHold.findUnique.mockImplementation(() => { order.push('check'); return Promise.resolve(LIVE_HOLD); }); + const write = jest.fn(async () => { order.push('write'); return { id: 'b1' } as any; }); + + await service.claimSeatsAndWrite([baseArgs], write); + + expect(order).toEqual(['lock', 'check', 'write']); + expect(mockPrisma.$transaction).toHaveBeenCalledTimes(1); + }); + + it('never runs the write when a leg fails its claim check', async () => { + seatHold.findUnique.mockResolvedValue({ ...LIVE_HOLD, expiresAt: new Date(Date.now() - 1) }); + const write = jest.fn(); + await expect(service.claimSeatsAndWrite([baseArgs], write as any)).rejects.toThrow(); + expect(write).not.toHaveBeenCalled(); + }); + + it('takes one lock covering every leg of a multi-leg booking, in one transaction', async () => { + const write = jest.fn(async () => ({ id: 'b1' }) as any); + await service.claimSeatsAndWrite( + [baseArgs, { ...baseArgs, seatIds: ['seat-2', 'seat-1'], legLabel: 'Return' }], + write, + ); + expect(mockPrisma.$transaction).toHaveBeenCalledTimes(1); + expect(queryRaw).toHaveBeenCalledTimes(1); + expect(write).toHaveBeenCalledTimes(1); + }); + + it('skips the lock statement when there are no seats to lock (free-child-only booking)', async () => { + const write = jest.fn(async () => ({ id: 'b1' }) as any); + await service.claimSeatsAndWrite([{ ...baseArgs, seatIds: [] }], write); + expect(queryRaw).not.toHaveBeenCalled(); + expect(write).toHaveBeenCalledTimes(1); + }); + }); + + describe('withSeatsLocked', () => { + it('de-duplicates and sorts seat ids so two overlapping requests cannot deadlock', async () => { + await service.withSeatsLocked(['s-b', 's-a', 's-b'], async () => null); + const sql = queryRaw.mock.calls[0][0]; + // Prisma.sql carries the interpolated values in `values`; order proves the sort. + expect(sql.values).toEqual(['s-a', 's-b']); + }); + + it('bounds the lock wait so a deep queue returns 409, not a transaction timeout', async () => { + await service.withSeatsLocked(['s-a'], async () => null); + expect(executeRawUnsafe).toHaveBeenCalledWith(expect.stringMatching(/SET LOCAL lock_timeout/i)); + }); + + it('reports a berth being claimed elsewhere as a conflict, not a 500', async () => { + queryRaw.mockRejectedValue(Object.assign(new Error('raw query failed'), { meta: { code: '55P03' } })); + await expect(service.withSeatsLocked(['s-a'], async () => null)) + .rejects.toThrow(/being booked by someone else/i); + }); + }); +}); diff --git a/apps/edr-passenger-api/src/modules/seats/seats.service.ts b/apps/edr-passenger-api/src/modules/seats/seats.service.ts index e54c085a1..c0c4420b1 100644 --- a/apps/edr-passenger-api/src/modules/seats/seats.service.ts +++ b/apps/edr-passenger-api/src/modules/seats/seats.service.ts @@ -1,16 +1,101 @@ import { Injectable, ConflictException, NotFoundException, BadRequestException, Logger } from '@nestjs/common'; import { randomUUID } from 'crypto'; +import { Prisma } from '@prisma/client'; import { PrismaService } from '../../common/prisma.service'; import { BlockSeatDto, HoldSeatsDto, JourneyDirection, SeatBlockReasonCategory } from './seats.dto'; import { ActingUser } from '../../common/acting-user'; import { Cron, CronExpression } from '@nestjs/schedule'; -import { SegmentsService } from '../segments/segments.service'; +import { SegmentsService, PrismaClientLike } from '../segments/segments.service'; import { SystemConfigService, CONFIG_KEYS } from '../system-config/system-config.service'; import { AuditService } from '../../common/audit.service'; import { AUDIT_ACTIONS, AUDIT_ENTITIES } from '../../common/audit.actions'; import { SmsClientService } from '../notifications/sms-client.service'; import { computePaymentDeadline } from '../../common/utils/payment-deadline.utils'; import { checkDirectionConflict } from '../../common/utils/journey-direction.utils'; +// Same definition the one-ticket-per-identity check uses: the states that still hold a +// traveller's place. CANCELLED / REFUNDED / NO_SHOW must free the berth immediately. +import { ACTIVE_BOOKING_STATUSES } from '../bookings/booking-identity.util'; + +/** + * How long a request will queue for a berth's row lock before giving up and reporting the + * seat as taken. Must stay comfortably under Prisma's interactive-transaction timeout + * (5000 ms by default) so this is what fires, not the transaction ceiling — the latter + * surfaces as an opaque 500 instead of an actionable 409. + */ +const SEAT_LOCK_TIMEOUT_MS = 3000; + +/** Postgres 55P03 lock_not_available, as surfaced through Prisma's raw-query wrapper. */ +function isLockTimeout(err: unknown): boolean { + const meta = (err as any)?.meta; + if (meta?.code === '55P03') return true; + const text = `${(err as any)?.message ?? ''} ${meta?.message ?? ''}`; + return /55P03|lock timeout|canceling statement due to lock timeout/i.test(text); +} + +/** + * One leg's claim on some seats: the hold being presented, the seats it must cover, and the + * stretch of that schedule they are wanted for. A one-way booking has one; a round-trip + * transit booking has four, each with its own hold. + */ +export interface SeatClaim { + holdId: string; + seatIds: string[]; + scheduleId: string; + originStationId?: string; + destinationStationId?: string; + journeyDirection?: JourneyDirection; + /** Label used in error messages when a booking presents several holds (one per leg). */ + legLabel?: string; +} + +/** + * The stretch of track a booking occupies on one schedule. + * + * Resolved by matching the SCHEDULE rather than the leg number. `BookingSeat.leg` is + * numbered per booking type (1=outbound, 2=return/leg-2, 3-4=return transit legs), so + * reading it correctly means re-deriving the same branching PaymentsService.createJourneySegments + * does. Matching on scheduleId asks the question we actually care about — "what stretch of + * THIS departure does this booking hold?" — and stays correct if leg numbering ever changes. + * + * Leg boundaries mirror createJourneySegments exactly: on a transit booking the outbound leg + * ends at the transit station (leg2OriginStationId), not at the booking's final destination. + * + * A booking that somehow lands on one schedule twice yields the union of both stretches, + * which is the conservative reading. + */ +function resolveBookingLegRange( + booking: { + scheduleId: string; + originStationId: string | null; + destinationStationId: string | null; + returnScheduleId: string | null; + returnOriginStationId: string | null; + returnDestinationStationId: string | null; + leg2ScheduleId: string | null; + leg2OriginStationId: string | null; + leg2DestinationStationId: string | null; + returnLeg2ScheduleId: string | null; + returnLeg2OriginStationId: string | null; + returnLeg2DestStationId: string | null; + }, + scheduleId: string, +): { originStationId: string | null; destinationStationId: string | null } | null { + const b = booking; + const legs = [ + // Outbound. On a transit booking this leg stops at the transit station. + { scheduleId: b.scheduleId, from: b.originStationId, to: b.leg2OriginStationId ?? b.destinationStationId }, + { scheduleId: b.leg2ScheduleId, from: b.leg2OriginStationId, to: b.leg2DestinationStationId }, + // Return. Likewise stops at the return transit station when there is one. + { scheduleId: b.returnScheduleId, from: b.returnOriginStationId, to: b.returnLeg2OriginStationId ?? b.returnDestinationStationId }, + { scheduleId: b.returnLeg2ScheduleId, from: b.returnLeg2OriginStationId, to: b.returnLeg2DestStationId }, + ].filter((l) => l.scheduleId === scheduleId && l.from && l.to); + + if (!legs.length) return null; + if (legs.length === 1) { + return { originStationId: legs[0].from, destinationStationId: legs[0].to }; + } + return { originStationId: legs[0].from, destinationStationId: legs[legs.length - 1].to }; +} @Injectable() export class SeatsService { @@ -206,6 +291,239 @@ export class SeatsService { return legacyMap[col?.toUpperCase()] ?? null; } + /** + * The single gate every booking path must pass before it may write a BookingSeat. + * + * It answers both halves of "may this request claim these seats": that the hold it + * presents actually covers them, and that nobody else has them for this leg. Those two + * checks used to live only in BookingsService, so POST /bookings/guest — a public, + * unauthenticated endpoint — validated nothing beyond "some hold exists and hasn't + * expired" and would happily connect any seat id the caller sent, held or not, free or + * not. Keep both booking services routed through here; a path that skips it can sell + * the same berth twice. + */ + async assertSeatsClaimable(args: SeatClaim & { + /** + * Read through this client. Callers inside `withSeatsLocked` pass their `tx` so the + * check runs on the connection that owns the row locks; everyone else gets the pool. + */ + db?: PrismaClientLike; + /** See assertNoRouteSeatConflict — booking writes turn this on, holds do not. */ + alsoRejectActiveBookings?: boolean; + }): Promise { + const { holdId, seatIds, legLabel } = args; + const db = args.db ?? this.prisma; + const prefix = legLabel ? `${legLabel} ` : ''; + + const hold = await db.seatHold.findUnique({ where: { id: holdId } }); + if (!hold || hold.expiresAt < new Date()) { + throw new BadRequestException(`${prefix}seat hold expired or not found`.trim()); + } + + // A hold is scoped to one departure, but the same physical Seat row is reused across + // every schedule its coach is assigned to. Without this check a hold taken on a quiet + // departure redeems that seat on a busy one — and the resulting booking is then + // "protected" by a hold on the wrong schedule, i.e. not protected at all. + if (hold.scheduleId !== args.scheduleId) { + throw new BadRequestException( + `${prefix}hold ${holdId} belongs to a different departure. ` + + `Hold the seats on this schedule before booking them.`.trim(), + ); + } + + const heldSeatIds = new Set(hold.seatIds); + const stray = seatIds.filter((id) => id && !heldSeatIds.has(id)); + if (stray.length) { + throw new BadRequestException( + `${prefix}seat(s) ${stray.join(', ')} are not part of hold ${holdId}. ` + + `Use seat IDs returned from POST /seats/hold.`.trim(), + ); + } + + await this.assertNoRouteSeatConflict({ + scheduleId: args.scheduleId, + seatIds, + originStationId: args.originStationId, + destinationStationId: args.destinationStationId, + journeyDirection: args.journeyDirection, + requestingPassengerId: hold.passengerId, + db, + alsoRejectActiveBookings: args.alsoRejectActiveBookings, + legLabel, + }); + } + + /** + * Runs `work` inside a transaction that holds an exclusive row lock on every seat it + * touches — the write-side twin of the lock `holdSeats` takes. + * + * Validating seat availability and then writing the BookingSeat rows in two separate + * statements is a read-then-write with nothing between them: N concurrent requests all + * read "free" and all write, which is exactly how one berth was sold six times over. + * Wrapping both in this makes the second request block on the lock, then re-read and see + * the first request's committed booking. + * + * `ORDER BY id` matches holdSeats, so a hold and a booking contending for the same two + * seats take them in the same sequence and queue instead of deadlocking. + * + * The transaction must stay short: do fare calculation, FX conversion and identity + * verification BEFORE calling this. Slow I/O in here holds seat locks open for the whole + * round trip. + */ + /** + * Takes the exclusive row locks for `seatIds` inside the caller's transaction. + * + * Bounded on purpose. FOR UPDATE serialises everyone contending for a berth, so under a + * burst the Nth contender waits for the N-1 transactions ahead of it. Left unbounded that + * wait runs past Prisma's 5s interactive-transaction ceiling and the request dies with + * P2028 — a 500 reading "Internal server error" for what is really "someone else got the + * seat". lock_timeout fires first and turns that into the honest answer. + * + * Waiting at all is still right: a fast hand-off (the common case) succeeds normally. Only + * a queue deep enough to mean the berth is genuinely being taken gets the 409. + */ + private async lockSeatsForUpdate(tx: Prisma.TransactionClient, seatIds: string[]): Promise { + if (!seatIds.length) return; + // SET LOCAL — scoped to this transaction, reset on commit/rollback. Must stay below the + // interactive-transaction timeout so it is the one that fires. + await tx.$executeRawUnsafe(`SET LOCAL lock_timeout = '${SEAT_LOCK_TIMEOUT_MS}ms'`); + try { + await tx.$queryRaw( + Prisma.sql`SELECT id FROM passenger."Seat" WHERE id IN (${Prisma.join(seatIds)}) ORDER BY id FOR UPDATE`, + ); + } catch (err) { + // 55P03 lock_not_available: another transaction is mid-claim on one of these berths. + if (isLockTimeout(err)) { + throw new ConflictException( + 'Those seats are being booked by someone else right now. Please pick another seat.', + ); + } + throw err; + } + } + + async withSeatsLocked( + seatIds: string[], + work: (tx: Prisma.TransactionClient) => Promise, + ): Promise { + // Sorted here as well as in the SQL. ORDER BY is the intent, but Postgres is only + // guaranteed to apply it before locking when the plan produces rows in that order — a + // bitmap heap scan can lock first and sort after. Handing the IN-list in id order costs + // nothing and makes the ordering independent of the planner. + const distinct = [...new Set(seatIds.filter(Boolean))].sort(); + return this.prisma.$transaction(async (tx) => { + await this.lockSeatsForUpdate(tx, distinct); + return work(tx); + }); + } + + /** + * The two-line contract every booking-creation path uses. + * + * const booking = await seatsService.claimSeatsAndWrite(claims, (tx) => tx.booking.create(…)); + * + * Locks every seat across every leg, re-runs the full claim gate for each leg inside that + * lock (including the active-booking check), then runs `write` — all in one transaction, so + * a concurrent request either blocks and then sees this booking, or is seen by it. + * + * Call `assertSeatsClaimableAll` first, before the fare and identity work, so an obviously + * doomed request fails early instead of paying for a Verifayda round trip it will discard. + * That earlier pass is advisory only; this one decides. + */ + async claimSeatsAndWrite( + claims: SeatClaim[], + write: (tx: Prisma.TransactionClient) => Promise, + ): Promise { + const allSeatIds = claims.flatMap((c) => c.seatIds); + return this.withSeatsLocked(allSeatIds, async (tx) => { + for (const claim of claims) { + await this.assertSeatsClaimable({ ...claim, db: tx, alsoRejectActiveBookings: true }); + } + return write(tx); + }); + } + + /** Advisory pre-flight for the same claims, outside any lock. See claimSeatsAndWrite. */ + async assertSeatsClaimableAll(claims: SeatClaim[]): Promise { + for (const claim of claims) { + await this.assertSeatsClaimable(claim); + } + } + + /** + * Rejects seats that an active booking already owns on an overlapping stretch of this + * schedule. + * + * This is deliberately NOT part of assertNoRouteSeatConflict. That check is also run by + * holdSeats, and some flows (the backoffice reservation-issue path, reschedule, upgrade) + * legitimately re-hold a seat for a booking that already owns it — folding this in there + * would break them. Here it guards only the moment of writing a NEW booking, where no + * legitimate caller can already hold the seat. + * + * It exists because getSeatAvailabilityMap cannot see a pending booking: it reads + * SeatHolds and JourneySegments, and JourneySegments are only written on payment success + * (PaymentsService.createJourneySegments). Until then a PENDING_PAYMENT booking is + * represented by nothing but its hold — so any request whose own hold is skipped as + * "its own" saw a free seat and booked straight over it. + */ + async assertNoPendingBookingConflict( + db: PrismaClientLike, + args: { scheduleId: string; seatIds: string[]; reqFrom: number; reqTo: number; legLabel?: string }, + ): Promise { + const { scheduleId, seatIds, reqFrom, reqTo, legLabel } = args; + if (!seatIds.length) return; + + const rows = await db.bookingSeat.findMany({ + where: { + scheduleId, + seatId: { in: seatIds }, + booking: { status: { in: ACTIVE_BOOKING_STATUSES } }, + }, + select: { + seatId: true, + seat: { select: { seatNumber: true } }, + booking: { + select: { + bookingRef: true, scheduleId: true, + originStationId: true, destinationStationId: true, + returnScheduleId: true, returnOriginStationId: true, returnDestinationStationId: true, + leg2ScheduleId: true, leg2OriginStationId: true, leg2DestinationStationId: true, + returnLeg2ScheduleId: true, returnLeg2OriginStationId: true, returnLeg2DestStationId: true, + }, + }, + }, + }); + if (!rows.length) return; + + const stopTimes = await db.tripStopTime.findMany({ + where: { scheduleId }, + select: { stationId: true, sequence: true }, + }); + const seqOf = (id?: string | null) => + id ? stopTimes.find(s => s.stationId === id)?.sequence : undefined; + + const conflicts = new Map(); + for (const row of rows) { + const leg = resolveBookingLegRange(row.booking as any, scheduleId); + const from = seqOf(leg?.originStationId); + const to = seqOf(leg?.destinationStationId); + // Conservative: a row whose leg we cannot resolve blocks. Selling a berth twice is + // far worse than refusing one booking we could not prove safe. + const overlaps = from === undefined || to === undefined || (from < reqTo && reqFrom < to); + if (overlaps) { + conflicts.set(row.seatId, row.seat?.seatNumber ?? row.seatId); + } + } + + if (conflicts.size > 0) { + const prefix = legLabel ? `${legLabel} ` : ''; + throw new ConflictException( + `${prefix}seat(s) ${[...conflicts.values()].join(', ')} were just booked by someone else. ` + + `Pick another seat.`.trim(), + ); + } + } + // Delegates the actual "is this seat held/booked for this leg" determination to // SegmentsService.getSeatAvailabilityMap — the same canonical check search results // (availabilityByClass) use — so the seatmap and search results can never disagree @@ -217,11 +535,28 @@ export class SeatsService { originStationId?: string; destinationStationId?: string; journeyDirection?: JourneyDirection; + /** + * Whoever is asking. OUTBOUND/RETURN on one schedule only stop conflicting when both + * legs belong to this same passenger — see getSeatAvailabilityMap. + */ + requestingPassengerId?: string; + /** Read through this client — a locking caller passes its own `tx`. */ + db?: PrismaClientLike; + /** + * Also reject seats an active booking already owns. Off by default: holdSeats and the + * reschedule/upgrade/reservation flows legitimately re-hold a seat whose booking already + * exists. Only the moment of writing a NEW booking turns this on — see + * assertNoPendingBookingConflict. + */ + alsoRejectActiveBookings?: boolean; + /** Label used in error messages when a booking presents several holds (one per leg). */ + legLabel?: string; }): Promise { const { scheduleId, seatIds, originStationId, destinationStationId, journeyDirection = JourneyDirection.ONE_WAY } = args; if (!seatIds.length) return; + const db = args.db ?? this.prisma; - const stopTimes = await this.prisma.tripStopTime.findMany({ + const stopTimes = await db.tripStopTime.findMany({ where: { scheduleId }, select: { stationId: true, sequence: true }, }); @@ -245,11 +580,19 @@ export class SeatsService { reqFrom, reqTo, journeyDirection, + args.requestingPassengerId, + db, ); + if (args.alsoRejectActiveBookings) { + await this.assertNoPendingBookingConflict(db, { + scheduleId, seatIds, reqFrom, reqTo, legLabel: args.legLabel, + }); + } + if (availability.size === 0) return; - const seats = await this.prisma.seat.findMany({ + const seats = await db.seat.findMany({ where: { id: { in: seatIds } }, select: { id: true, seatNumber: true }, }); @@ -413,6 +756,16 @@ export class SeatsService { } const hold = await this.prisma.$transaction(async (tx) => { + // Everything below is read-then-write — check no hold/segment covers these seats, + // then insert a hold — with no unique constraint behind it. Without a lock, N + // simultaneous requests all read "free" and all insert, which is how one berth ends + // up held (and then sold) several times over. Locking the Seat rows serializes those + // requests; ORDER BY id keeps two overlapping multi-seat requests taking the rows in + // the same sequence, so they queue instead of deadlocking. Held only for this + // transaction, so cross-schedule contention on a shared coach is negligible. + // Sorted so holds and booking writes take the same rows in the same sequence. + await this.lockSeatsForUpdate(tx, [...new Set(seatIds)].sort()); + const seats = await tx.seat.findMany({ where: { id: { in: seatIds } }, select: { id: true, status: true, seatNumber: true }, @@ -483,6 +836,16 @@ export class SeatsService { originStationId: dto.originStationId, destinationStationId: dto.destinationStationId, journeyDirection: currentDirection, + // Matches how SeatHold.passengerId is written below, so this passenger's own + // outbound hold doesn't block them from holding their return leg. + requestingPassengerId: dto.passengers[0]?.passengerId, + // MUST be tx, not the pool. This runs while the transaction already owns a pooled + // connection; reading through `this.prisma` would check out a SECOND one. With N + // concurrent holds that is 2N connections against a pool of N_max, so past roughly + // half the pool every transaction ends up waiting for a connection only another + // transaction can release — a pool deadlock that fails every request with a 500, + // not just the losers. Reproduced at 25 concurrent holds before this was passed. + db: tx, }); const activeHolds = await tx.seatHold.findMany({ @@ -511,7 +874,14 @@ export class SeatsService { const legsOverlap = legUnknown || (holdFrom < reqTo && reqFrom < holdTo); if (!legsOverlap) continue; - const directionsConflict = checkDirectionConflict(currentDirection, holdDirection); + // OUTBOUND and RETURN coexist on one schedule only for a single traveller holding + // both legs of their own turnaround trip. Between two different travellers that + // exemption is just a double-sold berth, so it applies only to the requester's + // own holds. + const requestPassengerIds = new Set(dto.passengers.map((p) => p.passengerId)); + const isOwnHold = passengerIds.some((id) => requestPassengerIds.has(id)); + + const directionsConflict = !isOwnHold || checkDirectionConflict(currentDirection, holdDirection); if (!directionsConflict) continue; for (const { passengerId, seatId } of dto.passengers) { @@ -1364,12 +1734,15 @@ export class SeatsService { }); const occupiedIds = new Set(journeySegments.map(js => js.seatId!)); - // Group BookingSeat rows by seatId::leg to find candidate duplicates, - // then filter to only those whose booking segments actually overlap. + // Group BookingSeat rows by seatId to find candidate duplicates, then filter to + // only those whose booking segments actually overlap. type BS = (typeof bookingSeats)[number]; const groups = new Map(); + // Grouped by seat alone. `leg` is a per-booking notion — one booking's leg 1 and + // another's leg 2 are the same physical berth on this departure — so keying on it + // hid every round-trip-vs-one-way collision from this report. for (const bs of bookingSeats) { - const key = `${bs.seatId}::${bs.leg}`; + const key = bs.seatId; if (!groups.has(key)) groups.set(key, []); groups.get(key)!.push(bs); } @@ -1463,16 +1836,20 @@ export class SeatsService { } if (overlapping.length <= 1) continue; - const [seatId] = key.split('::'); + const seatId = key; const seat = coach.seats.find(s => s.id === seatId); duplicates.push({ seatId, seatNumber: seat?.seatNumber ?? seatId, + // A group can now span legs (one booking's outbound against another's return), + // so the authoritative leg is the per-booking one below; this stays for + // backwards compatibility with existing callers. leg: overlapping[0].leg, bookings: overlapping.map(bs => ({ bookingSeatId: bs.id, bookingId: bs.booking.id, bookingRef: bs.booking.bookingRef, + leg: bs.leg, passengerName: bs.passengerName, contactPhone: bs.booking.contactPhone, createdAt: bs.booking.createdAt, diff --git a/apps/edr-passenger-api/src/modules/segments/segments.service.ts b/apps/edr-passenger-api/src/modules/segments/segments.service.ts index aaeb51e03..953fd705d 100644 --- a/apps/edr-passenger-api/src/modules/segments/segments.service.ts +++ b/apps/edr-passenger-api/src/modules/segments/segments.service.ts @@ -1,7 +1,15 @@ import { Injectable, BadRequestException } from '@nestjs/common'; +import { Prisma } from '@prisma/client'; import { PrismaService } from '../../common/prisma.service'; import { JourneyDirection } from '../seats/seats.dto'; -import { checkDirectionConflict } from '../../common/utils/journey-direction.utils'; + +/** + * Either the pooled client or a transaction-scoped one. A caller that has taken row locks + * must pass its own `tx`, or its "is this seat still free" read runs on a second connection + * outside the lock's transaction — burning a pool slot per in-flight booking and reading a + * snapshot the lock does not actually govern. + */ +export type PrismaClientLike = PrismaService | Prisma.TransactionClient; export interface Segment { fromStationId: string; @@ -73,9 +81,10 @@ export class SegmentsService { * P3: A(1) → D(4) reqFrom=1, reqTo=4 * Check P3 vs P2: 1 < 4 AND 2 < 4 → true AND true → CONFLICT ✓ * - * journeyDirection lets a round-trip's OUTBOUND and RETURN holds coexist on the - * same schedule without blocking each other (see checkDirectionConflict) — omit it - * for one-way contexts, where it defaults to ONE_WAY (conflicts with anything). + * A round-trip's OUTBOUND and RETURN holds coexist on the same schedule because they + * belong to the same traveller — pass requestingPassengerId so that traveller's own + * holds are skipped. Direction alone never grants that exemption: between two different + * travellers an OUTBOUND and a RETURN hold on one berth is a double sale. * * Sources checked: * 1. Active SeatHolds — leg + direction decoded from createdBy JSON @@ -93,7 +102,24 @@ export class SegmentsService { stopTimesForSeqLookup: ReadonlyArray<{ stationId: string; sequence: number }>, reqFrom: number, reqTo: number, + /** + * Retained for call-site compatibility (this is a positional signature) and for the + * seatmap/search callers that still describe their leg. Hold conflicts no longer turn + * on it — see the own-hold rule below. + */ journeyDirection: JourneyDirection = JourneyDirection.ONE_WAY, + /** + * The passenger asking. Their own active holds are skipped, which is what lets one + * traveller keep a berth across both legs of a turnaround round trip. Omit it in + * display contexts (seatmap, search) so every hold shows as taken. + */ + requestingPassengerId?: string, + /** + * Transaction client to read through. Defaults to the pooled client, which is right for + * every display caller; a caller holding `FOR UPDATE` locks must pass its own `tx` so the + * read happens on the locked connection. + */ + db: PrismaClientLike = this.prisma, ): Promise> { const result = new Map(); if (seatIds.length === 0) return result; @@ -105,11 +131,11 @@ export class SegmentsService { const now = new Date(); const [allHolds, bookedLegs] = await Promise.all([ - this.prisma.seatHold.findMany({ + db.seatHold.findMany({ where: { scheduleId, expiresAt: { gt: now } }, - select: { seatIds: true, createdBy: true }, + select: { seatIds: true, createdBy: true, passengerId: true }, }), - this.prisma.journeySegment.findMany({ + db.journeySegment.findMany({ where: { scheduleId, seatId: { in: seatIds }, @@ -123,22 +149,29 @@ export class SegmentsService { for (const hold of allHolds) { let holdFrom: number | undefined; let holdTo: number | undefined; - let holdDirection = JourneyDirection.ONE_WAY; try { if (hold.createdBy) { const meta = JSON.parse(hold.createdBy as string); holdFrom = seqOf(meta.originStationId); holdTo = seqOf(meta.destinationStationId); - holdDirection = meta.journeyDirection || JourneyDirection.ONE_WAY; } } catch { /* ignore */ } + // A hold reserves the seat *for* whoever placed it, so it must never be read as an + // obstacle to that same person's own booking — including the other leg of their own + // round trip, which is what the OUTBOUND/RETURN exemption used to cover. Anyone + // else's hold blocks on leg overlap alone: direction is irrelevant between two + // different travellers, and treating OUTBOUND and RETURN as compatible there is + // exactly how one berth got sold to two people. + const isOwnHold = + requestingPassengerId != null && hold.passengerId === requestingPassengerId; + if (isOwnHold) continue; + for (const sid of hold.seatIds) { if (!seatIdSet.has(sid)) continue; // Conservative block if leg can't be resolved; otherwise check overlap. const legsOverlap = holdFrom === undefined || holdTo === undefined || (holdFrom < reqTo && reqFrom < holdTo); if (!legsOverlap) continue; - if (!checkDirectionConflict(journeyDirection, holdDirection)) continue; result.set(sid, 'HELD'); } }