debate-runner.ts
4,276 bytes
| 1 | /** |
|---|---|
| 2 | * The thin adapter between an HTTP request and the framework-agnostic |
| 3 | * orchestrator: it selects the LLM client (mock vs OpenRouter), creates the |
| 4 | * debate row, and streams `DebateEvent`s over SSE while persisting each durable |
| 5 | * stage. |
| 6 | * |
| 7 | * Crucially, the orchestrator is NOT tied to the client connection - if the |
| 8 | * browser disconnects, `runDebate` keeps going server-side and keeps writing |
| 9 | * StageResults, so a reconnect can replay current state (graceful resume). |
| 10 | */ |
| 11 | import { participantsFor } from '@/core/config'; |
| 12 | import type { DebateEvent } from '@/core/events'; |
| 13 | import { MockLlmClient } from '@/core/mock-client'; |
| 14 | import { runDebate, type Logger } from '@/core/orchestrator'; |
| 15 | import { PROMPT_VERSION } from '@/core/prompts'; |
| 16 | import type { DebateConfig } from '@/core/types'; |
| 17 | import { createDebate, finalizeDebate, persistEvent } from '@/db/repositories'; |
| 18 | import { resolveApiKey } from './byok'; |
| 19 | import { registerDebate, unregisterDebate } from './debate-registry'; |
| 20 | import { env } from './env'; |
| 21 | import { createOpenRouterClient } from './openrouter'; |
| 22 | import { getPricing } from './model-cache'; |
| 23 | import { encodeEvent, encodeHeartbeat, encodeNamed, SSE_HEADERS } from './sse'; |
| 24 | |
| 25 | export class KeyResolutionError extends Error {} |
| 26 | |
| 27 | const serverLogger: Logger = { |
| 28 | info: (m, meta) => console.info(`[debate] ${m}`, meta ?? ''), |
| 29 | warn: (m, meta) => console.warn(`[debate] ${m}`, meta ?? ''), |
| 30 | error: (m, meta) => console.error(`[debate] ${m}`, meta ?? ''), |
| 31 | }; |
| 32 | |
| 33 | export interface StartDebateArgs { |
| 34 | config: DebateConfig; |
| 35 | userId: string | null; |
| 36 | } |
| 37 | |
| 38 | export async function startDebateStream( |
| 39 | args: StartDebateArgs, |
| 40 | ): Promise<{ response: Response; debateId: string }> { |
| 41 | const key = await resolveApiKey(); |
| 42 | if (!key.ok) throw new KeyResolutionError(key.reason); |
| 43 | |
| 44 | const participants = participantsFor(args.config.models); |
| 45 | const debateId = await createDebate({ |
| 46 | userId: args.userId, |
| 47 | config: args.config, |
| 48 | participants, |
| 49 | promptVersion: PROMPT_VERSION, |
| 50 | }); |
| 51 | |
| 52 | const llm = env.MOCK_LLM |
| 53 | ? new MockLlmClient() |
| 54 | : createOpenRouterClient({ apiKey: key.apiKey, pricing: await getPricing() }); |
| 55 | |
| 56 | // Abort controller for explicit server-side cancellation (not client |
| 57 | // disconnect) plus a gavel signal for graceful early conclusion. |
| 58 | const abort = new AbortController(); |
| 59 | const gavel = new AbortController(); |
| 60 | registerDebate(debateId, { abort, gavel }); |
| 61 | |
| 62 | let closed = false; |
| 63 | let heartbeat: ReturnType<typeof setInterval> | undefined; |
| 64 | |
| 65 | const stream = new ReadableStream<Uint8Array>({ |
| 66 | start(controller) { |
| 67 | const safeEnqueue = (chunk: Uint8Array) => { |
| 68 | if (closed) return; |
| 69 | try { |
| 70 | controller.enqueue(chunk); |
| 71 | } catch { |
| 72 | closed = true; |
| 73 | } |
| 74 | }; |
| 75 | |
| 76 | safeEnqueue(encodeNamed('ready', { debateId })); |
| 77 | heartbeat = setInterval(() => safeEnqueue(encodeHeartbeat()), 15_000); |
| 78 | |
| 79 | const emit = async (event: DebateEvent) => { |
| 80 | safeEnqueue(encodeEvent(event)); |
| 81 | try { |
| 82 | await persistEvent(debateId, event); |
| 83 | } catch (err) { |
| 84 | serverLogger.error('persist_failed', { type: event.type, err: String(err) }); |
| 85 | } |
| 86 | }; |
| 87 | |
| 88 | void (async () => { |
| 89 | try { |
| 90 | const result = await runDebate( |
| 91 | args.config, |
| 92 | { llm, emit, logger: serverLogger }, |
| 93 | { debateId, signal: abort.signal, gavelSignal: gavel.signal }, |
| 94 | ); |
| 95 | await finalizeDebate(debateId, result); |
| 96 | } catch (err) { |
| 97 | serverLogger.error('run_failed', { debateId, err: String(err) }); |
| 98 | safeEnqueue(encodeNamed('error', { message: err instanceof Error ? err.message : 'Debate failed' })); |
| 99 | } finally { |
| 100 | unregisterDebate(debateId); |
| 101 | if (heartbeat) clearInterval(heartbeat); |
| 102 | if (!closed) { |
| 103 | try { |
| 104 | controller.close(); |
| 105 | } catch { |
| 106 | /* already closed */ |
| 107 | } |
| 108 | } |
| 109 | } |
| 110 | })(); |
| 111 | }, |
| 112 | cancel() { |
| 113 | // Client went away. Stop writing to the socket but let the debate finish |
| 114 | // and persist server-side - do NOT abort the run. |
| 115 | closed = true; |
| 116 | if (heartbeat) clearInterval(heartbeat); |
| 117 | }, |
| 118 | }); |
| 119 | |
| 120 | return { response: new Response(stream, { headers: SSE_HEADERS }), debateId }; |
| 121 | } |
| 122 | |