diff --git a/src/routes/api/spark/run/+server.ts b/src/routes/api/spark/run/+server.ts index e9a11c23..dbf1e07b 100644 --- a/src/routes/api/spark/run/+server.ts +++ b/src/routes/api/spark/run/+server.ts @@ -53,6 +53,25 @@ interface SparkRunBody { }; } +export function _createOrderedMissionRelayQueue( + deliver: (event: T) => Promise, + onError: () => void = () => console.warn('[SparkRun] Mission relay delivery failed.') +) { + let tail = Promise.resolve(); + return { + enqueue(event: T) { + tail = tail.then(() => deliver(event)).catch(() => onError()); + }, + async flush() { + for (;;) { + const observedTail = tail; + await observedTail; + if (tail === observedTail) return; + } + } + }; +} + export const GET: RequestHandler = async (event) => { const unauthorized = requireControlAuth(event, { surface: 'SparkRunHealth', @@ -218,6 +237,7 @@ export const POST: RequestHandler = async (event) => { windowMs: 60_000 }); if (rateLimited) return rateLimited; + const relayQueue = _createOrderedMissionRelayQueue(relayMissionControlEvent); try { const body = (await event.request.json().catch(() => ({}))) as SparkRunBody; @@ -309,7 +329,7 @@ export const POST: RequestHandler = async (event) => { } }; eventBridge.emit(bridgeEvent); - void relayMissionControlEvent(bridgeEvent); + relayQueue.enqueue(bridgeEvent); }; emitMissionEvent('mission_created', `Mission created (${mission.id}).`); @@ -381,7 +401,7 @@ export const POST: RequestHandler = async (event) => { } }; eventBridge.emit(relayEvent); - void relayMissionControlEvent(relayEvent); + relayQueue.enqueue(relayEvent); if (bridgeEvent.type === 'dispatch_started' && !missionStartedEmitted) { missionStartedEmitted = true; emitMissionEvent('mission_started', `Mission started (${mission.id}).`, { @@ -391,6 +411,7 @@ export const POST: RequestHandler = async (event) => { } } }); + await relayQueue.flush(); return json({ success: true, @@ -406,6 +427,7 @@ export const POST: RequestHandler = async (event) => { audit: capability }); } catch (error) { + await relayQueue.flush(); if (error instanceof HarnessAuthorityError) { return json({ success: false, error: error.message, code: error.code, authority: error.verdict }, { status: error.status }); } diff --git a/src/routes/api/spark/run/spark-run.integration.test.ts b/src/routes/api/spark/run/spark-run.integration.test.ts index 3194e05d..5ad7ff6c 100644 --- a/src/routes/api/spark/run/spark-run.integration.test.ts +++ b/src/routes/api/spark/run/spark-run.integration.test.ts @@ -23,7 +23,7 @@ vi.mock('$lib/server/provider-runtime', () => ({ } })); -import { GET, POST } from './+server'; +import { _createOrderedMissionRelayQueue, GET, POST } from './+server'; import { providerRuntime } from '$lib/server/provider-runtime'; import { eventBridge } from '$lib/services/event-bridge'; import { getMissionControlPersistPath, getMissionControlRelaySnapshot } from '$lib/server/mission-control-relay'; @@ -348,6 +348,84 @@ describe('/api/spark/run integration', () => { } }); + it('serializes relay delivery and waits for a slow first event before flush resolves', async () => { + let releaseFirst!: () => void; + const firstGate = new Promise((resolve) => { + releaseFirst = resolve; + }); + const started: string[] = []; + const completed: string[] = []; + const queue = _createOrderedMissionRelayQueue(async (event: { type: string }) => { + started.push(event.type); + if (event.type === 'mission_created') await firstGate; + completed.push(event.type); + }); + + queue.enqueue({ type: 'mission_created' }); + queue.enqueue({ type: 'dispatch_started' }); + queue.enqueue({ type: 'mission_started' }); + queue.enqueue({ type: 'task_failed' }); + queue.enqueue({ type: 'mission_failed' }); + let flushSettled = false; + const flush = queue.flush().then(() => { + flushSettled = true; + }); + await Promise.resolve(); + + expect(started).toEqual(['mission_created']); + expect(completed).toEqual([]); + expect(flushSettled).toBe(false); + + releaseFirst(); + await flush; + expect(started).toEqual([ + 'mission_created', 'dispatch_started', 'mission_started', 'task_failed', 'mission_failed' + ]); + expect(completed).toEqual(started); + }); + + it('keeps flush open for relay events appended while an earlier delivery is draining', async () => { + let releaseFirst!: () => void; + let releaseSecond!: () => void; + let markSecondStarted!: () => void; + const firstGate = new Promise((resolve) => { + releaseFirst = resolve; + }); + const secondGate = new Promise((resolve) => { + releaseSecond = resolve; + }); + const secondStarted = new Promise((resolve) => { + markSecondStarted = resolve; + }); + const completed: string[] = []; + const queue = _createOrderedMissionRelayQueue(async (event: { type: string }) => { + if (event.type === 'mission_created') await firstGate; + if (event.type === 'mission_failed') { + markSecondStarted(); + await secondGate; + } + completed.push(event.type); + }); + + queue.enqueue({ type: 'mission_created' }); + let flushSettled = false; + const flush = queue.flush().then(() => { + flushSettled = true; + }); + await Promise.resolve(); + queue.enqueue({ type: 'mission_failed' }); + releaseFirst(); + await secondStarted; + + expect(completed).toEqual(['mission_created']); + expect(flushSettled).toBe(false); + + releaseSecond(); + await flush; + expect(completed).toEqual(['mission_created', 'mission_failed']); + expect(flushSettled).toBe(true); + }); + it('derives one mission_started when provider runtime repeats dispatch_started', async () => { const dispatch = vi.mocked(providerRuntime.dispatch); const emitted: Array<{ type?: string; missionId?: string }> = [];