Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 24 additions & 2 deletions src/routes/api/spark/run/+server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,25 @@ interface SparkRunBody {
};
}

export function _createOrderedMissionRelayQueue<T>(
deliver: (event: T) => Promise<void>,
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',
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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}).`);
Expand Down Expand Up @@ -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}).`, {
Expand All @@ -391,6 +411,7 @@ export const POST: RequestHandler = async (event) => {
}
}
});
await relayQueue.flush();

return json({
success: true,
Expand All @@ -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 });
}
Expand Down
80 changes: 79 additions & 1 deletion src/routes/api/spark/run/spark-run.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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<void>((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<void>((resolve) => {
releaseFirst = resolve;
});
const secondGate = new Promise<void>((resolve) => {
releaseSecond = resolve;
});
const secondStarted = new Promise<void>((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 }> = [];
Expand Down