From a56ae7797d6081ea01f5d154a14446737f99cd12 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Thu, 17 Sep 2026 19:15:49 -0700 Subject: [PATCH] fix(server): keep tool payloads out of completion queues Filter terminal-worker subscriptions before live buffering and historical payload decoding. Preserve matching run updates without retaining unrelated tool output while queue promotion waits on providers or thread locks. --- apps/server/src/orchestration-v2/EventSink.ts | 69 ++++++++++---- .../server/src/orchestration-v2/EventStore.ts | 2 + .../FoundationPersistence.test.ts | 95 +++++++++++++++++++ .../src/orchestration-v2/Orchestrator.ts | 32 ++++--- .../Layers/OrchestrationEventStore.ts | 3 + .../Services/OrchestrationEventStore.ts | 1 + 6 files changed, 169 insertions(+), 33 deletions(-) diff --git a/apps/server/src/orchestration-v2/EventSink.ts b/apps/server/src/orchestration-v2/EventSink.ts index b6339fddc6b3..e9306e144980 100644 --- a/apps/server/src/orchestration-v2/EventSink.ts +++ b/apps/server/src/orchestration-v2/EventSink.ts @@ -159,6 +159,8 @@ export interface EventSinkV2Shape { readonly stream: (input?: { readonly threadId?: ThreadId; readonly afterSequence?: number; + /** Filter before queuing live events so a busy worker retains only the events it handles. */ + readonly eventType?: OrchestrationV2DomainEvent["type"]; /** Bound RPC subscribers; internal workers must not drop their subscription under load. */ readonly bounded?: boolean; }) => Stream.Stream; @@ -196,6 +198,20 @@ const baseLayer: Layer.Layer< const projectionStore = yield* ProjectionStoreV2; const turnItemPositions = yield* TurnItemPositionStoreV2; const liveEvents = yield* PubSub.unbounded(); + const liveEventsByType = new Map< + OrchestrationV2DomainEvent["type"], + PubSub.PubSub + >(); + const publishLiveEvents = (events: ReadonlyArray) => + Effect.gen(function* () { + yield* PubSub.publishAll(liveEvents, events); + for (const [type, pubsub] of liveEventsByType) { + yield* PubSub.publishAll( + pubsub, + events.filter((stored) => stored.event.type === type), + ); + } + }); // A user can answer after terminal normalization reads the pending request. // Recheck inside the write transaction so stale cleanup cannot erase answers. @@ -324,7 +340,7 @@ const baseLayer: Layer.Layer< yield* effectOutbox.notifyAvailable(input.effects.length); } yield* eventStore.publishCommitted(storedEvents); - yield* PubSub.publishAll(liveEvents, storedEvents); + yield* publishLiveEvents(storedEvents); return storedEvents; }); @@ -378,7 +394,7 @@ const baseLayer: Layer.Layer< ); if (result.committed) { yield* eventStore.publishCommitted(result.storedEvents); - yield* PubSub.publishAll(liveEvents, result.storedEvents); + yield* publishLiveEvents(result.storedEvents); } return result; }, @@ -439,7 +455,7 @@ const baseLayer: Layer.Layer< ); if (result.committed) { yield* eventStore.publishCommitted(result.storedEvents); - yield* PubSub.publishAll(liveEvents, result.storedEvents); + yield* publishLiveEvents(result.storedEvents); } return result; }); @@ -517,7 +533,7 @@ const baseLayer: Layer.Layer< } if (result.committed) { yield* eventStore.publishCommitted(result.storedEvents); - yield* PubSub.publishAll(liveEvents, result.storedEvents); + yield* publishLiveEvents(result.storedEvents); } return { receipt: result.receipt, @@ -556,6 +572,7 @@ const baseLayer: Layer.Layer< readonly afterSequence: number; readonly throughSequence: number; readonly threadId?: ThreadId; + readonly eventType?: OrchestrationV2DomainEvent["type"]; }): Stream.Stream => { const pageSize = 256; const loop = (afterSequence: number): Stream.Stream => @@ -565,6 +582,7 @@ const baseLayer: Layer.Layer< afterSequence, throughSequence: input.throughSequence, ...(input.threadId === undefined ? {} : { threadId: input.threadId }), + ...(input.eventType === undefined ? {} : { eventType: input.eventType }), limit: pageSize, }) .pipe( @@ -587,31 +605,44 @@ const baseLayer: Layer.Layer< const stream = (input?: Parameters[0]) => { const afterSequence = input?.afterSequence ?? 0; - const matchesThread = (stored: OrchestrationV2StoredEvent) => - input?.threadId === undefined || stored.event.threadId === input.threadId; + const matches = (stored: OrchestrationV2StoredEvent) => + (input?.threadId === undefined || stored.event.threadId === input.threadId) && + (input?.eventType === undefined || stored.event.type === input.eventType); const replay = (throughSequence: number) => catchUp({ afterSequence, throughSequence, ...(input?.threadId === undefined ? {} : { threadId: input.threadId }), - }); - if (input?.bounded === true) { - return replayAndBufferProjectedLiveEvents({ - subscribe: PubSub.subscribe(liveEvents), - latestSequence: eventStore.latestSequence(), - afterSequence, - filter: matchesThread, - replay, - project: (stored) => ({ ...stored, event: projectDomainEventForWire(stored.event) }), - }); - } + ...(input?.eventType === undefined ? {} : { eventType: input.eventType }), + }).pipe(Stream.filter(matches)); return Stream.unwrap( Effect.gen(function* () { - const subscription = yield* PubSub.subscribe(liveEvents); + let pubsub = liveEvents; + if (input?.eventType !== undefined) { + const existing = liveEventsByType.get(input.eventType); + if (existing !== undefined) { + pubsub = existing; + } else { + const created = yield* PubSub.unbounded(); + pubsub = liveEventsByType.get(input.eventType) ?? created; + liveEventsByType.set(input.eventType, pubsub); + } + } + if (input?.bounded === true) { + return replayAndBufferProjectedLiveEvents({ + subscribe: PubSub.subscribe(pubsub), + latestSequence: eventStore.latestSequence(), + afterSequence, + filter: matches, + replay, + project: (stored) => ({ ...stored, event: projectDomainEventForWire(stored.event) }), + }); + } + const subscription = yield* PubSub.subscribe(pubsub); const highWater = yield* eventStore.latestSequence(); const live = Stream.fromSubscription(subscription).pipe( Stream.filter((stored) => stored.sequence > Math.max(highWater, afterSequence)), - Stream.filter(matchesThread), + Stream.filter(matches), ); return Stream.concat(replay(highWater), live); }), diff --git a/apps/server/src/orchestration-v2/EventStore.ts b/apps/server/src/orchestration-v2/EventStore.ts index a21abfd389c2..34c4c1f80051 100644 --- a/apps/server/src/orchestration-v2/EventStore.ts +++ b/apps/server/src/orchestration-v2/EventStore.ts @@ -56,6 +56,7 @@ export interface EventStoreV2Shape { readonly afterSequence?: number; readonly throughSequence?: number; readonly threadId?: ThreadId; + readonly eventType?: OrchestrationV2DomainEvent["type"]; readonly limit?: number; }) => Stream.Stream; readonly readByCommandId: (input: { @@ -86,6 +87,7 @@ const baseLayer: Layer.Layer = Lay ? {} : { throughSequence: input.throughSequence }), ...(input?.threadId === undefined ? {} : { threadId: input.threadId }), + ...(input?.eventType === undefined ? {} : { eventType: input.eventType }), ...(input?.limit === undefined ? {} : { limit: input.limit }), }) .pipe( diff --git a/apps/server/src/orchestration-v2/FoundationPersistence.test.ts b/apps/server/src/orchestration-v2/FoundationPersistence.test.ts index f32be226ead7..1e0952b4aea0 100644 --- a/apps/server/src/orchestration-v2/FoundationPersistence.test.ts +++ b/apps/server/src/orchestration-v2/FoundationPersistence.test.ts @@ -550,6 +550,101 @@ it.layer(TestLayer)("orchestration V2 foundation persistence", (it) => { ), ); + it.effect("filters worker replay and live queues without losing matching events", () => + Effect.scoped( + Effect.gen(function* () { + const sink = yield* EventSinkV2; + const sql = yield* SqlClient.SqlClient; + const now = yield* DateTime.now; + const thread = makeThread(ThreadId.make("thread:filtered-worker"), now); + const other = makeThread(ThreadId.make("thread:filtered-worker-other"), now); + const run: OrchestrationV2Run = { + id: RunId.make("run:filtered-worker"), + threadId: thread.id, + ordinal: 1, + providerInstanceId, + modelSelection, + providerThreadId: null, + userMessageId: MessageId.make("message:filtered-worker"), + rootNodeId: null, + activeAttemptId: null, + status: "completed", + queuePosition: null, + requestedAt: now, + startedAt: now, + completedAt: now, + checkpointId: null, + contextHandoffId: null, + }; + const runEvent = (id: string, payload: OrchestrationV2Run): OrchestrationV2DomainEvent => ({ + id: EventId.make(id), + type: "run.updated", + threadId: payload.threadId, + runId: payload.id, + occurredAt: now, + payload, + }); + const history = yield* sink.write({ + events: [ + threadCreatedEvent({ id: "event:filtered-worker:created", thread, now }), + threadCreatedEvent({ id: "event:filtered-worker:other", thread: other, now }), + runEvent("event:filtered-worker:history", run), + runEvent("event:filtered-worker:other-history", { + ...run, + id: RunId.make("run:filtered-worker-other"), + threadId: other.id, + }), + ], + }); + // Unrelated payloads must be skipped in SQL, before decoding or + // allocating their bodies, even when a retained row is unreadable. + const original = yield* sql<{ readonly payload_json: string }>` + SELECT payload_json FROM orchestration_events + WHERE sequence = ${history[0]!.sequence} + `; + yield* Effect.acquireRelease( + sql` + UPDATE orchestration_events SET payload_json = 'unreadable unrelated payload' + WHERE sequence = ${history[0]!.sequence} + `, + () => + sql` + UPDATE orchestration_events SET payload_json = ${original[0]!.payload_json} + WHERE sequence = ${history[0]!.sequence} + `.pipe(Effect.orDie), + ); + const pull = yield* Stream.toPull( + sink.stream({ threadId: thread.id, eventType: "run.updated" }), + ); + assert.deepEqual( + (yield* pull).map((stored) => stored.sequence), + [history[2]!.sequence], + ); + // The worker is occupied with the previous batch while the thread + // publishes output. Its live queue must receive just run updates. + yield* sink.write({ + events: Array.from({ length: LIVE_STREAM_MAX_ITEMS + 1 }, (_, index) => ({ + id: EventId.make(`event:filtered-worker:output:${index}`), + type: "thread.metadata-updated" as const, + threadId: thread.id, + occurredAt: now, + payload: { ...thread, title: `Output ${index}` }, + })), + }); + const live = yield* sink.write({ + events: [ + runEvent("event:filtered-worker:live:1", { ...run, status: "interrupted" }), + runEvent("event:filtered-worker:live:2", { ...run, status: "failed" }), + ], + }); + assert.deepEqual( + (yield* pull).map((stored) => stored.sequence), + live.map((stored) => stored.sequence), + ); + }), + ), + ); + it.effect("paginates catch-up beyond the event-store read limit", () => Effect.gen(function* () { const eventSink = yield* EventSinkV2; diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 71675be18b0a..f0450477629a 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -8549,20 +8549,24 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio // below. Replaying the full event table on every server start delays live // queue promotion in proportion to the lifetime size of the database. const terminalEventsAfterSequence = yield* eventSink.latestSequence().pipe(Effect.orDie); - yield* eventSink.stream({ afterSequence: terminalEventsAfterSequence }).pipe( - Stream.filter( - (stored) => - stored.event.type === "run.updated" && - !String(stored.commandId).startsWith("command:runtime-reconcile:") && - (stored.event.payload.status === "completed" || - stored.event.payload.status === "interrupted" || - stored.event.payload.status === "failed" || - stored.event.payload.status === "cancelled" || - stored.event.payload.status === "rolled_back"), - ), - Stream.runForEach(handleTerminalRun), - Effect.forkDetach, - ); + // Queue promotion can wait on a provider or a thread lock. Subscribe to run + // updates before buffering so that wait never retains unrelated tool bodies. + yield* eventSink + .stream({ afterSequence: terminalEventsAfterSequence, eventType: "run.updated" }) + .pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && + !String(stored.commandId).startsWith("command:runtime-reconcile:") && + (stored.event.payload.status === "completed" || + stored.event.payload.status === "interrupted" || + stored.event.payload.status === "failed" || + stored.event.payload.status === "cancelled" || + stored.event.payload.status === "rolled_back"), + ), + Stream.runForEach(handleTerminalRun), + Effect.forkDetach, + ); // The high-water subscription deliberately skips history, so recover the // two terminal side effects from current projections instead: one queued diff --git a/apps/server/src/persistence/Layers/OrchestrationEventStore.ts b/apps/server/src/persistence/Layers/OrchestrationEventStore.ts index b7d9ad71664c..0c427609be0d 100644 --- a/apps/server/src/persistence/Layers/OrchestrationEventStore.ts +++ b/apps/server/src/persistence/Layers/OrchestrationEventStore.ts @@ -372,6 +372,7 @@ const makeEventStore = Effect.gen(function* () { readonly throughSequence?: number; readonly threadId?: ThreadId; readonly commandId?: CommandId; + readonly eventType?: OrchestrationV2DomainEvent["type"]; readonly onlyAgentEvents?: boolean; readonly limit: number; }) => @@ -408,6 +409,7 @@ const makeEventStore = Effect.gen(function* () { AND ${sql.and([ ...(input.threadId === undefined ? [] : [sql`stream_id = ${input.threadId}`]), ...(input.commandId === undefined ? [] : [sql`command_id = ${input.commandId}`]), + ...(input.eventType === undefined ? [] : [sql`event_type = ${input.eventType}`]), ])} ORDER BY sequence ASC LIMIT ${input.limit} @@ -492,6 +494,7 @@ const makeEventStore = Effect.gen(function* () { : { throughSequence: input.throughSequence }), ...(input?.threadId === undefined ? {} : { threadId: input.threadId }), ...(input?.commandId === undefined ? {} : { commandId: input.commandId }), + ...(input?.eventType === undefined ? {} : { eventType: input.eventType }), onlyAgentEvents: true, limit: pageLimit, }).pipe( diff --git a/apps/server/src/persistence/Services/OrchestrationEventStore.ts b/apps/server/src/persistence/Services/OrchestrationEventStore.ts index 826d65e54818..f6c2297badbf 100644 --- a/apps/server/src/persistence/Services/OrchestrationEventStore.ts +++ b/apps/server/src/persistence/Services/OrchestrationEventStore.ts @@ -78,6 +78,7 @@ export interface OrchestrationEventStoreShape { readonly throughSequence?: number; readonly threadId?: ThreadId; readonly commandId?: CommandId; + readonly eventType?: OrchestrationV2DomainEvent["type"]; readonly limit?: number; }) => Stream.Stream;