server: transient entry flow (button -> ticket -> signed entry -> open)
Closes the long-dangling thread from device-input-flow. On an access device's rising input edge: print the ticket (failover), then sign vehicle_entry, then pulseOpen, then cache the session projection. Two invariants enforced: - signed BEFORE open (an open with no signed event is the fraud signal); - HOLD on print failure — no ticket means a transient can't pay on exit, so sign an anomaly and do NOT open, and do NOT write a vehicle_entry for a car that never got in. Subscribes the same input bus as the device-telemetry writer (independent: telemetry always records; entry acts only on an access device's on-edge, debounced). Verified end to end against stubs: success path signs+opens+caches and verifyChain ok; printer-down path emits only an anomaly with no open and no entry; release edge ignored.
This commit is contained in:
@@ -0,0 +1,184 @@
|
|||||||
|
import { randomUUID } from "node:crypto";
|
||||||
|
import { and, eq, laneDevices, sessions, type Db } from "@parking/db";
|
||||||
|
import {
|
||||||
|
NoPrinterAvailableError,
|
||||||
|
printWithFailover,
|
||||||
|
registry,
|
||||||
|
type AccessControlDevice,
|
||||||
|
type PrinterDevice,
|
||||||
|
type PrinterInstance,
|
||||||
|
type TicketData,
|
||||||
|
} from "@parking/devices";
|
||||||
|
import type { FastifyBaseLogger } from "fastify";
|
||||||
|
import type { DeviceInputEvent } from "./device-events.js";
|
||||||
|
import type { EventLog } from "./event-log.js";
|
||||||
|
import type { LaneMap } from "./lane-map.js";
|
||||||
|
|
||||||
|
// The transient ENTRY flow: a button press → print a ticket → sign a vehicle_entry
|
||||||
|
// → open the barrier. This is the step the device layer left dangling
|
||||||
|
// (wiki/concepts/device-input-flow.md "the entry flow itself is the next build").
|
||||||
|
//
|
||||||
|
// Two invariants from the threat model + safety analysis:
|
||||||
|
// 1. SIGNED BEFORE OPEN — the vehicle_entry is appended to the signed ledger
|
||||||
|
// BEFORE pulseOpen fires; an open with no matching signed event is the fraud
|
||||||
|
// signal (wiki/concepts/append-only-event-chain.md).
|
||||||
|
// 2. HOLD ON PRINT FAILURE — a transient with no ticket can't pay on exit, so if
|
||||||
|
// all printers are down we do NOT open. We sign an `anomaly` (attempt, ticket
|
||||||
|
// unprinted) and leave the barrier closed; the operator handles the held car.
|
||||||
|
// Crucially, NO vehicle_entry is written in that case — we never record an
|
||||||
|
// "entered" event for a car that didn't get in (decision 2026-06-15).
|
||||||
|
//
|
||||||
|
// Ordering, therefore: print → (ok) sign vehicle_entry → pulseOpen → cache session.
|
||||||
|
// (fail) sign anomaly, stop.
|
||||||
|
|
||||||
|
/** Map a 1-based entry input to the relay/door it opens. Default: same channel. */
|
||||||
|
function doorForInput(input: number): number {
|
||||||
|
return input;
|
||||||
|
}
|
||||||
|
|
||||||
|
export class EntryFlow {
|
||||||
|
readonly #db: Db;
|
||||||
|
readonly #log: EventLog;
|
||||||
|
readonly #laneMap: LaneMap;
|
||||||
|
readonly #logger: FastifyBaseLogger;
|
||||||
|
/** Guard against double-fire from the same physical press (on edge only). */
|
||||||
|
readonly #inFlight = new Set<string>();
|
||||||
|
|
||||||
|
constructor(db: Db, log: EventLog, laneMap: LaneMap, logger: FastifyBaseLogger) {
|
||||||
|
this.#db = db;
|
||||||
|
this.#log = log;
|
||||||
|
this.#laneMap = laneMap;
|
||||||
|
this.#logger = logger;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Handle a device input edge. Acts only on the rising ("on") edge of an entry
|
||||||
|
* button in a lane that has an access (barrier) device. */
|
||||||
|
async onInput(e: DeviceInputEvent): Promise<void> {
|
||||||
|
if (e.edge !== "on") return; // release edge is just telemetry
|
||||||
|
|
||||||
|
const lane = this.#laneMap.laneFor(e.deviceId);
|
||||||
|
if (lane == null) return; // unmapped device — telemetry already recorded, no entry
|
||||||
|
|
||||||
|
// Only treat this as an entry trigger if the firing device IS the lane's
|
||||||
|
// access controller (a reader/printer input edge isn't an entry button).
|
||||||
|
const access = await this.#loadAccess(lane, e.deviceId);
|
||||||
|
if (!access) return;
|
||||||
|
|
||||||
|
const key = `${e.deviceId}:${e.input}`;
|
||||||
|
if (this.#inFlight.has(key)) return; // ignore re-fire while one is processing
|
||||||
|
this.#inFlight.add(key);
|
||||||
|
try {
|
||||||
|
await this.#runEntry(lane, e.input, access);
|
||||||
|
} catch (err) {
|
||||||
|
this.#logger.error(`entry-flow failed (lane ${lane}): ${(err as Error).message}`);
|
||||||
|
} finally {
|
||||||
|
this.#inFlight.delete(key);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async #runEntry(lane: number, input: number, access: AccessControlDevice): Promise<void> {
|
||||||
|
const ticketId = newTicketId();
|
||||||
|
const issuedAt = new Date().toISOString();
|
||||||
|
const printers = await this.#loadPrinters(lane);
|
||||||
|
|
||||||
|
// 1. PRINT FIRST. The ticket is the transient's session key — no ticket, no entry.
|
||||||
|
const ticket: TicketData = { ticketId, lane, issuedAt };
|
||||||
|
try {
|
||||||
|
const printedBy = await printWithFailover(printers, "entry-dispenser", (d: PrinterDevice) =>
|
||||||
|
d.printTicket(ticket),
|
||||||
|
);
|
||||||
|
this.#logger.info(`entry ticket ${ticketId} printed on ${printedBy} (lane ${lane})`);
|
||||||
|
} catch (err) {
|
||||||
|
// HOLD: do not open, do not record a vehicle_entry. Sign an anomaly so the
|
||||||
|
// failed attempt is in the tamper-evident record for the operator.
|
||||||
|
const reason =
|
||||||
|
err instanceof NoPrinterAvailableError ? err.message : (err as Error).message;
|
||||||
|
await this.#log.append({
|
||||||
|
type: "anomaly",
|
||||||
|
lane,
|
||||||
|
identity: ticketId,
|
||||||
|
payload: { reason: `entry held — ticket not printed: ${reason}`, ticketPrinted: false },
|
||||||
|
});
|
||||||
|
this.#logger.warn(`entry HELD on lane ${lane}: ${reason} (barrier NOT opened)`);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
// 2. SIGN the vehicle_entry — BEFORE the relay fires (the core invariant).
|
||||||
|
await this.#log.append({
|
||||||
|
type: "vehicle_entry",
|
||||||
|
lane,
|
||||||
|
direction: "entry",
|
||||||
|
source: "ticket",
|
||||||
|
identity: ticketId,
|
||||||
|
payload: { sessionRef: ticketId, ticketPrinted: true },
|
||||||
|
occurredAt: issuedAt,
|
||||||
|
});
|
||||||
|
|
||||||
|
// 3. OPEN the barrier (intent only; the barrier owns the close).
|
||||||
|
await access.pulseOpen(doorForInput(input));
|
||||||
|
|
||||||
|
// 4. Update the session projection cache (rebuildable from the ledger; this is
|
||||||
|
// just a fast read-model, never the source of truth).
|
||||||
|
try {
|
||||||
|
this.#db
|
||||||
|
.insert(sessions)
|
||||||
|
.values({ id: ticketId, lane, identity: ticketId, source: "ticket", enteredAt: issuedAt, state: "open" })
|
||||||
|
.run();
|
||||||
|
} catch (err) {
|
||||||
|
// Cache miss is non-fatal — the ledger is authoritative and the projection
|
||||||
|
// can be rebuilt. Log it; don't fail the (already-open) entry.
|
||||||
|
this.#logger.error(`session-cache insert failed for ${ticketId}: ${(err as Error).message}`);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** The lane's access device, but only if it's the one that fired (the entry
|
||||||
|
* button). Returns a live adapter or null. */
|
||||||
|
async #loadAccess(lane: number, deviceId: string): Promise<AccessControlDevice | null> {
|
||||||
|
const row = await this.#db
|
||||||
|
.select()
|
||||||
|
.from(laneDevices)
|
||||||
|
.where(and(eq(laneDevices.id, deviceId), eq(laneDevices.category, "access")))
|
||||||
|
.get();
|
||||||
|
if (!row || !row.enabled || row.lane !== lane) return null;
|
||||||
|
const driver = registry.get(row.driverId);
|
||||||
|
if (!driver) return null;
|
||||||
|
try {
|
||||||
|
return driver.create(row.config as never) as AccessControlDevice;
|
||||||
|
} catch {
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Build live printer instances for a lane (for failover selection). */
|
||||||
|
async #loadPrinters(lane: number): Promise<PrinterInstance[]> {
|
||||||
|
const rows = await this.#db
|
||||||
|
.select()
|
||||||
|
.from(laneDevices)
|
||||||
|
.where(and(eq(laneDevices.category, "printer"), eq(laneDevices.lane, lane)))
|
||||||
|
.all();
|
||||||
|
const out: PrinterInstance[] = [];
|
||||||
|
for (const row of rows) {
|
||||||
|
if (!row.enabled) continue;
|
||||||
|
const driver = registry.get(row.driverId);
|
||||||
|
if (!driver) continue;
|
||||||
|
const cfg = row.config as Record<string, unknown>;
|
||||||
|
const role = cfg.role === "booth-receipt" ? "booth-receipt" : "entry-dispenser";
|
||||||
|
try {
|
||||||
|
out.push({
|
||||||
|
id: row.id,
|
||||||
|
role,
|
||||||
|
failoverRank: typeof cfg.failoverRank === "number" ? cfg.failoverRank : 0,
|
||||||
|
device: driver.create(cfg as never) as PrinterDevice,
|
||||||
|
});
|
||||||
|
} catch {
|
||||||
|
// skip a printer whose config won't build
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Opaque, unguessable transient ticket id (wiki/concepts/ticket-encoding.md). */
|
||||||
|
function newTicketId(): string {
|
||||||
|
return `T-${randomUUID()}`;
|
||||||
|
}
|
||||||
@@ -5,6 +5,7 @@ import { randomUUID } from "node:crypto";
|
|||||||
import { createDb, deviceEvents as deviceEventsTable, type Db } from "@parking/db";
|
import { createDb, deviceEvents as deviceEventsTable, type Db } from "@parking/db";
|
||||||
import { TOKEN_COOKIE, requireJwtSecret } from "./auth.js";
|
import { TOKEN_COOKIE, requireJwtSecret } from "./auth.js";
|
||||||
import { deviceEvents } from "./device-events.js";
|
import { deviceEvents } from "./device-events.js";
|
||||||
|
import { EntryFlow } from "./entry-flow.js";
|
||||||
import { EventLog } from "./event-log.js";
|
import { EventLog } from "./event-log.js";
|
||||||
import { LaneMap } from "./lane-map.js";
|
import { LaneMap } from "./lane-map.js";
|
||||||
import { PrinterMonitor } from "./printer-monitor.js";
|
import { PrinterMonitor } from "./printer-monitor.js";
|
||||||
@@ -79,6 +80,17 @@ export async function buildServer(opts: BuildOptions = {}): Promise<FastifyInsta
|
|||||||
// See wiki/decisions/event-streams-split.md.
|
// See wiki/decisions/event-streams-split.md.
|
||||||
const eventLog = new EventLog(db, buildSigner(app.log));
|
const eventLog = new EventLog(db, buildSigner(app.log));
|
||||||
await eventRoutes(app, db, eventLog);
|
await eventRoutes(app, db, eventLog);
|
||||||
|
|
||||||
|
// Entry flow: a button press → print ticket → signed vehicle_entry → pulseOpen.
|
||||||
|
// Subscribes to the SAME input bus as the telemetry writer below; the two are
|
||||||
|
// independent (telemetry always records; the entry flow acts only on an access
|
||||||
|
// device's rising edge). See wiki/concepts/device-input-flow.md + parking-session.md.
|
||||||
|
const entryFlow = new EntryFlow(db, eventLog, laneMap, app.log);
|
||||||
|
const unsubscribeEntry = deviceEvents.onInput((e) => {
|
||||||
|
void entryFlow.onInput(e);
|
||||||
|
});
|
||||||
|
app.addHook("onClose", async () => unsubscribeEntry());
|
||||||
|
|
||||||
const unsubscribeInput = deviceEvents.onInput((e) => {
|
const unsubscribeInput = deviceEvents.onInput((e) => {
|
||||||
// Resolve which lane the device belongs to. -1 marks "device fired but isn't
|
// Resolve which lane the device belongs to. -1 marks "device fired but isn't
|
||||||
// mapped to a lane" (assigned without a lane, or a stale id) — still recorded
|
// mapped to a lane" (assigned without a lane, or a stale id) — still recorded
|
||||||
|
|||||||
Reference in New Issue
Block a user