/** * Dependency-free, in-memory authoritative realtime sessions. * * The service owns credentials, interest filtering, sequence allocation and * movement validation. HTTP/SSE routes are adapters only. No poses survive a * process restart and usable bearer tokens are never retained in memory. */ import { createHash, randomBytes, timingSafeEqual } from "node:crypto"; import { isEntityPoseSnapshot, isInterestCell, validateMotion, } from "../../../src/realtime/protocol.ts"; import type { EntityPoseSnapshot, InterestCell, JoinGrant, MotionLimits, PoseDelta, RealtimeMemberRole, ResumeGrant, ResumeRequest, ServerRealtimeMessage, } from "../../../src/realtime/types.ts"; const UINT32_MAX = 4_294_967_295; export interface RealtimeServiceOptions { now?: () => number; sessionTtlMs?: number; disconnectGraceMs?: number; maximumSessions?: number; maximumSessionsPerCell?: number; maximumInterestsPerSession?: number; maximumPoseUpdatesPerWindow?: number; rateWindowMs?: number; actorMotionLimits?: MotionLimits; vehicleMotionLimits?: MotionLimits; aircraftMotionLimits?: MotionLimits; } export interface RealtimeJoinInput { requestId: string; actorId: string; subject: string; role: RealtimeMemberRole; interests: readonly InterestCell[]; } export interface RealtimeJoinResult { sessionId: string; token: string; grant: JoinGrant; } export type RealtimeFailureCode = | "invalid" | "unauthorized" | "expired" | "capacity" | "conflict" | "rate-limited" | "sequence" | "interest" | "ownership" | "motion"; export type RealtimeResult = | { ok: true; value: T } | { ok: false; code: RealtimeFailureCode; message: string }; export type RealtimeListener = (message: ServerRealtimeMessage) => void; interface Session { id: string; tokenHash: Buffer; previousTokenHash: Buffer | null; previousTokenExpiresAt: number; subject: string; actorId: string; role: RealtimeMemberRole; interests: readonly InterestCell[]; interestKeys: ReadonlySet; createdAt: number; expiresAt: number; disconnectedAt: number | null; lastClientSequence: number; rateWindowStartedAt: number; rateCount: number; proposedPoses: Map; poses: Map; listeners: Set; } interface ResolvedOptions { now: () => number; sessionTtlMs: number; disconnectGraceMs: number; maximumSessions: number; maximumSessionsPerCell: number; maximumInterestsPerSession: number; maximumPoseUpdatesPerWindow: number; rateWindowMs: number; actorMotionLimits: MotionLimits; vehicleMotionLimits: MotionLimits; aircraftMotionLimits: MotionLimits; } export interface RealtimeService { join(input: RealtimeJoinInput): RealtimeResult; authenticate(sessionId: string, token: string): RealtimeResult; connect(sessionId: string, token: string, listener: RealtimeListener): RealtimeResult; resume(request: ResumeRequest, subject: string): RealtimeResult; connectResume( request: ResumeRequest, subject: string, listener: RealtimeListener, ): RealtimeResult; submitPose( sessionId: string, token: string, snapshot: unknown, ): RealtimeResult; leave(sessionId: string, token: string): RealtimeResult; revokeSubject(subject: string): number; cleanup(): number; sessionCount(): number; } export interface RealtimeSessionView { sessionId: string; actorId: string; role: RealtimeMemberRole; interests: readonly InterestCell[]; expiresAt: number; snapshot: readonly EntityPoseSnapshot[]; } export interface RealtimeConnection { view: RealtimeSessionView; disconnect(): void; } export interface RealtimeResumeConnection extends RealtimeConnection { grant: ResumeGrant; } function finitePositive(value: number | undefined, fallback: number): number { return typeof value === "number" && Number.isFinite(value) && value > 0 ? value : fallback; } function finiteCount(value: number | undefined, fallback: number): number { return Math.max(1, Math.floor(finitePositive(value, fallback))); } function resolveOptions(value: RealtimeServiceOptions): ResolvedOptions { return { now: value.now ?? Date.now, sessionTtlMs: finitePositive(value.sessionTtlMs, 5 * 60_000), disconnectGraceMs: finitePositive(value.disconnectGraceMs, 20_000), maximumSessions: finiteCount(value.maximumSessions, 1_024), maximumSessionsPerCell: finiteCount(value.maximumSessionsPerCell, 96), maximumInterestsPerSession: finiteCount(value.maximumInterestsPerSession, 16), maximumPoseUpdatesPerWindow: finiteCount(value.maximumPoseUpdatesPerWindow, 120), rateWindowMs: finitePositive(value.rateWindowMs, 10_000), actorMotionLimits: value.actorMotionLimits ?? { maximumHorizontalSpeedMps: 24, maximumVerticalSpeedMps: 18, maximumTurnRateDegPerSec: 540, positionSlackM: 1.5, maximumDeltaSeconds: 2, }, vehicleMotionLimits: value.vehicleMotionLimits ?? { maximumHorizontalSpeedMps: 105, maximumVerticalSpeedMps: 45, maximumTurnRateDegPerSec: 240, positionSlackM: 4, maximumDeltaSeconds: 2, }, aircraftMotionLimits: value.aircraftMotionLimits ?? { maximumHorizontalSpeedMps: 125, maximumVerticalSpeedMps: 45, maximumTurnRateDegPerSec: 180, positionSlackM: 6, maximumDeltaSeconds: 2, }, }; } function hashToken(token: string): Buffer { return createHash("sha256").update(token).digest(); } function tokenMatches(expected: Buffer, token: string): boolean { if (typeof token !== "string" || token.length < 16 || token.length > 512) return false; const found = hashToken(token); return found.length === expected.length && timingSafeEqual(found, expected); } function interestKey(cell: InterestCell): string { switch (cell.kind) { case "california-tile": return `tile:${cell.level}:${cell.x}:${cell.y}`; case "city": return `city:${cell.cityId}`; case "office": return `office:${cell.officeId}`; case "floor": return `floor:${cell.officeId}:${cell.floorId}`; case "room": return `room:${cell.officeId}:${cell.floorId}:${cell.roomId}`; } } function cloneCell(cell: InterestCell): InterestCell { return { ...cell }; } function clonePose(pose: EntityPoseSnapshot): EntityPoseSnapshot { return JSON.parse(JSON.stringify(pose)) as EntityPoseSnapshot; } function entityKey(pose: EntityPoseSnapshot): string { if (pose.entity === "actor") return `actor:${pose.actorId}`; if (pose.entity === "vehicle") return `vehicle:${pose.vehicleId}`; return `aircraft:${pose.aircraftId}`; } function poseInterestKey(pose: EntityPoseSnapshot): string | null { if (pose.pose.space === "local") return interestKey(pose.pose.cell); // Geographic coordinates are interpreted through the session's authored // outdoor interest: either one exact California tile or one exact city. return null; } function validIdentity(value: string): boolean { return typeof value === "string" && value.length > 0 && value.length <= 256; } export function createRealtimeService(options: RealtimeServiceOptions = {}): RealtimeService { const settings = resolveOptions(options); const sessions = new Map(); const epoch = randomBytes(16).toString("base64url"); let serverSequence = 0; function nextSequence(): number { serverSequence = serverSequence >= UINT32_MAX ? 0 : serverSequence + 1; return serverSequence; } function geographicInterestKeys(session: Session): string[] { return [...session.interestKeys].filter( (key) => key.startsWith("tile:") || key.startsWith("city:"), ); } function poseBroadcastKeys(session: Session, pose: EntityPoseSnapshot): string[] { const local = poseInterestKey(pose); return local === null ? geographicInterestKeys(session) : [local]; } function broadcast(delta: PoseDelta, keys: readonly string[], excluding?: Session): void { if (keys.length === 0) return; for (const peer of sessions.values()) { if (peer === excluding || !keys.some((key) => peer.interestKeys.has(key))) continue; for (const listener of peer.listeners) listener(delta); } } function removalDelta(session: Session): PoseDelta | null { const removedEntityIds = [...session.poses.keys()]; if (removedEntityIds.length === 0) return null; return { type: "pose-delta", protocolVersion: 1, serverEpoch: epoch, sequence: nextSequence(), timestampMs: settings.now(), updates: [], removedEntityIds, }; } function broadcastRemoval(session: Session, keys = [...session.interestKeys]): void { const delta = removalDelta(session); if (delta) broadcast(delta, keys, session); } function remove(session: Session, announce = true): void { if (announce) broadcastRemoval(session); sessions.delete(session.id); session.listeners.clear(); session.proposedPoses.clear(); session.poses.clear(); session.tokenHash.fill(0); session.previousTokenHash?.fill(0); } function cleanup(): number { const now = settings.now(); let removed = 0; for (const session of sessions.values()) { const graceExpired = session.disconnectedAt !== null && now - session.disconnectedAt > settings.disconnectGraceMs; if (now >= session.expiresAt || graceExpired) { remove(session); removed += 1; } } return removed; } function authenticate( sessionId: string, token: string, allowPrevious = false, ): RealtimeResult { cleanup(); if (!validIdentity(sessionId)) { return { ok: false, code: "unauthorized", message: "Realtime session is invalid." }; } const session = sessions.get(sessionId); const matchesCurrent = session ? tokenMatches(session.tokenHash, token) : false; const matchesPrevious = allowPrevious && session?.previousTokenHash !== null && settings.now() <= (session?.previousTokenExpiresAt ?? 0) && tokenMatches(session?.previousTokenHash ?? Buffer.alloc(0), token); if (!session || (!matchesCurrent && !matchesPrevious)) { return { ok: false, code: "unauthorized", message: "Realtime session is invalid." }; } if (settings.now() >= session.expiresAt) { remove(session); return { ok: false, code: "expired", message: "Realtime session expired." }; } return { ok: true, value: session }; } function view(session: Session): RealtimeSessionView { const snapshot: EntityPoseSnapshot[] = []; for (const peer of sessions.values()) { for (const pose of peer.poses.values()) { const visible = poseBroadcastKeys(peer, pose).some((key) => session.interestKeys.has(key)); if (visible) snapshot.push(clonePose(pose)); } } return { sessionId: session.id, actorId: session.actorId, role: session.role, interests: session.interests.map(cloneCell), expiresAt: session.expiresAt, snapshot, }; } function join(input: RealtimeJoinInput): RealtimeResult { cleanup(); if ( !validIdentity(input.requestId) || !validIdentity(input.actorId) || !validIdentity(input.subject) || (input.role !== "member" && input.role !== "admin") || !Array.isArray(input.interests) || input.interests.length === 0 || input.interests.length > settings.maximumInterestsPerSession || !input.interests.every(isInterestCell) ) return { ok: false, code: "invalid", message: "Realtime join request is invalid." }; if (sessions.size >= settings.maximumSessions) { return { ok: false, code: "capacity", message: "Realtime service is at capacity." }; } if ([...sessions.values()].some((session) => session.actorId === input.actorId)) { return { ok: false, code: "conflict", message: "Realtime actor id is already active." }; } const unique = new Map(input.interests.map((cell) => [interestKey(cell), cloneCell(cell)])); for (const key of unique.keys()) { let count = 0; for (const session of sessions.values()) if (session.interestKeys.has(key)) count += 1; if (count >= settings.maximumSessionsPerCell) { return { ok: false, code: "capacity", message: "Realtime interest cell is at capacity." }; } } const now = settings.now(); const id = randomBytes(18).toString("base64url"); const token = randomBytes(32).toString("base64url"); const interests = [...unique.values()]; const session: Session = { id, tokenHash: hashToken(token), previousTokenHash: null, previousTokenExpiresAt: 0, subject: input.subject, actorId: input.actorId, role: input.role, interests, interestKeys: new Set(unique.keys()), createdAt: now, expiresAt: now + settings.sessionTtlMs, disconnectedAt: null, lastClientSequence: -1, rateWindowStartedAt: now, rateCount: 0, proposedPoses: new Map(), poses: new Map(), listeners: new Set(), }; sessions.set(id, session); const grant: JoinGrant = { type: "join-grant", protocolVersion: 1, requestId: input.requestId, sessionId: id, actorId: input.actorId, role: input.role, serverEpoch: epoch, serverTimeMs: now, nextSequence: serverSequence >= UINT32_MAX ? 0 : serverSequence + 1, resumeToken: token, interests: interests.map(cloneCell), initial: view(session).snapshot, }; return { ok: true, value: { sessionId: id, token, grant } }; } function connect( sessionId: string, token: string, listener: RealtimeListener, ): RealtimeResult { if (typeof listener !== "function") { return { ok: false, code: "invalid", message: "Realtime listener is invalid." }; } const found = authenticate(sessionId, token); if (!found.ok) return found; const session = found.value; session.disconnectedAt = null; session.listeners.add(listener); let connected = true; return { ok: true, value: { view: view(session), disconnect() { if (!connected) return; connected = false; session.listeners.delete(listener); if (session.listeners.size === 0) session.disconnectedAt = settings.now(); }, }, }; } function updateInterests( session: Session, requested: readonly InterestCell[], ): RealtimeResult { if ( !Array.isArray(requested) || requested.length === 0 || requested.length > settings.maximumInterestsPerSession || !requested.every(isInterestCell) ) return { ok: false, code: "invalid", message: "Realtime interests are invalid." }; const unique = new Map(requested.map((cell) => [interestKey(cell), cloneCell(cell)])); for (const key of unique.keys()) { let count = 0; for (const peer of sessions.values()) { if (peer.id !== session.id && peer.interestKeys.has(key)) count += 1; } if (count >= settings.maximumSessionsPerCell) { return { ok: false, code: "capacity", message: "Realtime interest cell is at capacity." }; } } const changed = unique.size !== session.interestKeys.size || [...unique.keys()].some((key) => !session.interestKeys.has(key)); const previousKeys = [...session.interestKeys]; if (changed) broadcastRemoval(session, previousKeys); session.interests = [...unique.values()]; session.interestKeys = new Set(unique.keys()); if (changed) { // A pose from the old cell cannot become the first authoritative point in // a new cell's speed check or leak through its next snapshot. session.poses.clear(); session.proposedPoses.clear(); } return { ok: true, value: session.interests }; } function rotateToken(session: Session): string { const token = randomBytes(32).toString("base64url"); session.previousTokenHash?.fill(0); session.previousTokenHash = session.tokenHash; session.previousTokenExpiresAt = settings.now() + 10_000; session.tokenHash = hashToken(token); return token; } function resume( request: ResumeRequest, subject: string, ): RealtimeResult { // The immediately previous token is accepted only for resume. This small // overlap lets a client retry when the rotated grant was lost in transit, // without letting an old token continue publishing or leaving a session. const found = authenticate(request.sessionId, request.resumeToken, true); if (!found.ok) return found; const session = found.value; if (session.subject !== subject) { return { ok: false, code: "unauthorized", message: "Realtime session is invalid." }; } if (request.serverEpoch !== epoch || request.lastReceivedSequence > serverSequence) { return { ok: false, code: "sequence", message: "Realtime resume point is invalid." }; } const updated = updateInterests(session, request.interests); if (!updated.ok) return updated; const continuous = request.lastReceivedSequence === serverSequence; const token = rotateToken(session); return { ok: true, value: { type: "resume-grant", protocolVersion: 1, requestId: request.requestId, sessionId: session.id, serverEpoch: epoch, serverTimeMs: settings.now(), nextSequence: serverSequence >= UINT32_MAX ? 0 : serverSequence + 1, resumeToken: token, continuous, snapshot: view(session).snapshot, }, }; } function connectResume( request: ResumeRequest, subject: string, listener: RealtimeListener, ): RealtimeResult { if (typeof listener !== "function") { return { ok: false, code: "invalid", message: "Realtime listener is invalid." }; } const resumed = resume(request, subject); if (!resumed.ok) return resumed; const session = sessions.get(request.sessionId); if (!session) return { ok: false, code: "expired", message: "Realtime session expired." }; session.disconnectedAt = null; session.listeners.add(listener); let connected = true; return { ok: true, value: { grant: resumed.value, view: view(session), disconnect() { if (!connected) return; connected = false; session.listeners.delete(listener); if (session.listeners.size === 0) session.disconnectedAt = settings.now(); }, }, }; } function submitPose( sessionId: string, token: string, input: unknown, ): RealtimeResult { const found = authenticate(sessionId, token); if (!found.ok) return found; const session = found.value; const now = settings.now(); if (now - session.rateWindowStartedAt >= settings.rateWindowMs) { session.rateWindowStartedAt = now; session.rateCount = 0; } session.rateCount += 1; if (session.rateCount > settings.maximumPoseUpdatesPerWindow) { return { ok: false, code: "rate-limited", message: "Realtime pose rate exceeded." }; } if (!isEntityPoseSnapshot(input)) { return { ok: false, code: "invalid", message: "Realtime pose is invalid." }; } const snapshot = input; if (snapshot.sequence <= session.lastClientSequence) { return { ok: false, code: "sequence", message: "Realtime pose sequence is stale." }; } if (snapshot.timestampMs > now + 5_000 || snapshot.timestampMs < now - 10_000) { return { ok: false, code: "sequence", message: "Realtime pose timestamp is stale." }; } if ( (snapshot.entity === "actor" && snapshot.actorId !== session.actorId) || (snapshot.entity === "vehicle" && snapshot.driverActorId !== session.actorId) || (snapshot.entity === "aircraft" && snapshot.pilotActorId !== session.actorId) ) return { ok: false, code: "ownership", message: "Realtime entity is not owned by this session." }; if (snapshot.entity === "aircraft" && snapshot.pose.space !== "geographic") { return { ok: false, code: "interest", message: "Aircraft poses require an outdoor geographic interest.", }; } const localKey = poseInterestKey(snapshot); if (localKey !== null && !session.interestKeys.has(localKey)) { return { ok: false, code: "interest", message: "Realtime pose is outside the joined interest." }; } if (localKey === null && geographicInterestKeys(session).length !== 1) { return { ok: false, code: "interest", message: "Geographic pose requires exactly one outdoor interest.", }; } const key = entityKey(snapshot); const previous = session.proposedPoses.get(key); if (previous) { const motionLimits = snapshot.entity === "vehicle" ? settings.vehicleMotionLimits : snapshot.entity === "aircraft" ? settings.aircraftMotionLimits : settings.actorMotionLimits; const motion = validateMotion( previous, snapshot, motionLimits, ); if (!motion.ok) { return { ok: false, code: "motion", message: `Realtime motion rejected: ${motion.reason}.` }; } } session.lastClientSequence = snapshot.sequence; const authoritative = clonePose({ ...snapshot, sequence: nextSequence(), timestampMs: now }); const removedEntityIds = [...session.poses.keys()].filter((existing) => existing !== key); for (const removed of removedEntityIds) { session.poses.delete(removed); session.proposedPoses.delete(removed); } session.proposedPoses.set(key, clonePose(snapshot)); session.poses.set(key, authoritative); const delta: PoseDelta = { type: "pose-delta", protocolVersion: 1, serverEpoch: epoch, sequence: authoritative.sequence, timestampMs: now, updates: [clonePose(authoritative)], removedEntityIds, }; broadcast(delta, poseBroadcastKeys(session, snapshot)); return { ok: true, value: delta }; } function leave(sessionId: string, token: string): RealtimeResult { const found = authenticate(sessionId, token); if (!found.ok) return found; remove(found.value); return { ok: true, value: null }; } function revokeSubject(subject: string): number { let count = 0; for (const session of sessions.values()) { if (session.subject !== subject) continue; const message: ServerRealtimeMessage = { type: "membership-revoked", protocolVersion: 1, serverEpoch: epoch, sequence: nextSequence(), timestampMs: settings.now(), sessionId: session.id, actorId: session.actorId, reason: "membership-revoked", reconnectAllowed: false, }; for (const listener of session.listeners) listener(message); broadcastRemoval(session); remove(session, false); count += 1; } return count; } return { join, authenticate(sessionId, token) { const found = authenticate(sessionId, token); return found.ok ? { ok: true, value: view(found.value) } : found; }, connect, resume, connectResume, submitPose, leave, revokeSubject, cleanup, sessionCount: () => sessions.size, }; }