From b04271ef3a31725c57c6c4e36ed50542d30a6267 Mon Sep 17 00:00:00 2001 From: badcuban <108198679+badcuban@users.noreply.github.com> Date: Mon, 31 Aug 2026 16:22:03 -0400 Subject: [PATCH] perf(agents): reduce subagent streaming overhead --- .../Layers/ProviderRuntimeIngestion.test.ts | 349 ++++++++++++++++++ .../Layers/ProviderRuntimeIngestion.ts | 274 +++++++++++++- apps/web/src/agentsPanelStore.ts | 44 ++- apps/web/src/components/ChatView.tsx | 50 ++- .../components/chat/AgentsPanel.browser.tsx | 30 +- apps/web/src/components/chat/AgentsPanel.tsx | 87 +++-- .../routes/_chat.$environmentId.$threadId.tsx | 1 - apps/web/src/session-logic.test.ts | 22 ++ apps/web/src/session-logic.ts | 94 ++++- 9 files changed, 867 insertions(+), 84 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index dfe3ecc87..27eb2b727 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -1450,6 +1450,56 @@ describe("ProviderRuntimeIngestion", () => { ), ); + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-subagent-result-before-abort-1"), + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + createdAt: "2026-01-01T00:00:03.500Z", + turnId: asTurnId("turn-abort"), + itemId: asItemId("item-subagent-result-before-abort"), + providerRefs: { + providerThreadId: "agent-live-before-abort", + providerTurnId: "agent-turn-before-abort", + providerItemId: asItemId("item-subagent-result-before-abort"), + }, + payload: { + streamKind: "assistant_text", + delta: "First chunk. ", + }, + }); + await waitForThread(harness.readModel, (thread) => + thread.activities.some((activity) => { + const payload = activity.payload as { + sourceAgentThreadId?: string; + data?: { subagentLiveText?: string }; + }; + return ( + activity.kind === "subagent.result" && + payload.sourceAgentThreadId === "agent-live-before-abort" && + payload.data?.subagentLiveText === "First chunk. " + ); + }), + ); + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-subagent-result-before-abort-2"), + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + createdAt: "2026-01-01T00:00:04.000Z", + turnId: asTurnId("turn-abort"), + itemId: asItemId("item-subagent-result-before-abort"), + providerRefs: { + providerThreadId: "agent-live-before-abort", + providerTurnId: "agent-turn-before-abort", + providerItemId: asItemId("item-subagent-result-before-abort"), + }, + payload: { + streamKind: "assistant_text", + delta: "Second chunk.", + }, + }); + harness.emit({ type: "turn.aborted", eventId: asEventId("evt-turn-aborted"), @@ -1482,6 +1532,29 @@ describe("ProviderRuntimeIngestion", () => { (subagent) => subagent.agentThreadId === "agent-completed-before-abort", )?.status, ).toBe("completed"); + const streamedResult = thread.activities.find((activity) => { + const payload = activity.payload as { sourceAgentThreadId?: string }; + return ( + activity.kind === "subagent.result" && + payload.sourceAgentThreadId === "agent-live-before-abort" + ); + }); + expect(streamedResult?.payload).toMatchObject({ + data: { + subagentLiveText: "First chunk. Second chunk.", + }, + }); + + await Effect.runPromise(Effect.sleep("150 millis")); + await harness.drain(); + const settledThread = (await harness.readModel()).threads.find( + (entry) => entry.id === "thread-1", + ); + expect( + settledThread?.subagents?.find( + (subagent) => subagent.agentThreadId === "agent-live-before-abort", + )?.status, + ).toBe("interrupted"); }); it("settles live subagents when turn.completed reports an interruption", async () => { @@ -2171,6 +2244,282 @@ describe("ProviderRuntimeIngestion", () => { status: "completed", message: "The runtime path is correct.", }); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-child-final-late-delta"), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-parent"), + itemId: asItemId("child-message-final"), + providerRefs: { + providerThreadId: "child-provider-thread", + providerTurnId: "child-turn-1", + providerItemId: asItemId("child-message-final"), + }, + payload: { + streamKind: "assistant_text", + delta: " Late duplicate text.", + }, + }); + await harness.drain(); + + const afterLateDelta = (await harness.readModel()).threads.find( + (entry) => entry.id === "thread-1", + ); + const finalActivity = afterLateDelta?.activities.find((entry) => entry.id === activity?.id); + expect(finalActivity?.payload).toMatchObject({ + status: "completed", + data: { + item: { + agentsStates: { + "child-provider-thread": { + status: "completed", + message: "The runtime path is correct.", + }, + }, + }, + }, + }); + }); + + it("retains a child result when turn and session completion arrive before item completion", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-out-of-order-child-completion"), + threadId: ThreadId.make("thread-1"), + session: { + threadId: ThreadId.make("thread-1"), + status: "running", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: asTurnId("turn-out-of-order-child-completion"), + providerThreadId: "parent-provider-thread", + updatedAt: now, + lastError: null, + }, + createdAt: now, + }), + ); + + for (const [index, delta] of ["Full child ", "result."].entries()) { + harness.emit({ + type: "content.delta", + eventId: asEventId(`evt-out-of-order-child-delta-${index}`), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-out-of-order-child-completion"), + itemId: asItemId("out-of-order-child-message"), + providerRefs: { + providerThreadId: "out-of-order-child-provider-thread", + providerTurnId: "out-of-order-child-turn", + providerItemId: asItemId("out-of-order-child-message"), + }, + payload: { + streamKind: "assistant_text", + delta, + }, + }); + } + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-out-of-order-parent-turn-completed"), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-out-of-order-child-completion"), + payload: { + state: "completed", + }, + }); + harness.emit({ + type: "session.exited", + eventId: asEventId("evt-out-of-order-session-exited"), + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + createdAt: now, + payload: { + reason: "provider exited", + exitKind: "graceful", + }, + }); + harness.emit({ + type: "item.completed", + eventId: asEventId("evt-out-of-order-child-item-completed"), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-out-of-order-child-completion"), + itemId: asItemId("out-of-order-child-message"), + providerRefs: { + providerThreadId: "out-of-order-child-provider-thread", + providerTurnId: "out-of-order-child-turn", + providerItemId: asItemId("out-of-order-child-message"), + }, + payload: { + itemType: "assistant_message", + status: "completed", + data: { + item: { + phase: "final_answer", + }, + }, + }, + }); + await harness.drain(); + + const thread = (await harness.readModel()).threads.find((entry) => entry.id === "thread-1"); + const result = thread?.activities.find((entry) => { + const resultPayload = entry.payload as { sourceAgentThreadId?: string; status?: string }; + return ( + entry.kind === "subagent.result" && + resultPayload.sourceAgentThreadId === "out-of-order-child-provider-thread" && + resultPayload.status === "completed" + ); + }); + expect(result?.payload).toMatchObject({ + data: { + item: { + agentsStates: { + "out-of-order-child-provider-thread": { + status: "completed", + message: "Full child result.", + }, + }, + }, + }, + }); + }); + + it("keeps rapid updates from three child agents independent", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + const childProviderThreadIds = ["child-provider-a", "child-provider-b", "child-provider-c"]; + const readSubagentResultEvents = async () => { + const events = await Effect.runPromise( + Stream.runCollect(harness.engine.readEvents(0)).pipe( + Effect.map((chunk) => Array.from(chunk)), + ), + ); + return events.filter((event) => { + if (event.type !== "thread.activity-appended") { + return false; + } + const activity = event.payload.activity; + return ( + activity.kind === "subagent.result" && + childProviderThreadIds.includes( + (activity.payload as { sourceAgentThreadId?: string }).sourceAgentThreadId ?? "", + ) + ); + }); + }; + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-three-streaming-children"), + threadId: ThreadId.make("thread-1"), + session: { + threadId: ThreadId.make("thread-1"), + status: "running", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: asTurnId("turn-three-streaming-children"), + providerThreadId: "parent-provider-thread", + updatedAt: now, + lastError: null, + }, + createdAt: now, + }), + ); + + for (const [childIndex, childProviderThreadId] of childProviderThreadIds.entries()) { + harness.emit({ + type: "content.delta", + eventId: asEventId(`evt-three-streaming-children-${childIndex}-0`), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-three-streaming-children"), + itemId: asItemId(`child-message-${childIndex}`), + providerRefs: { + providerThreadId: childProviderThreadId, + providerTurnId: `child-turn-${childIndex}`, + providerItemId: asItemId(`child-message-${childIndex}`), + }, + payload: { + streamKind: "assistant_text", + delta: `${childProviderThreadId}:0`, + }, + }); + } + + await waitForThread(harness.readModel, (thread) => { + const sources = new Set( + thread.activities + .filter((activity) => activity.kind === "subagent.result") + .map( + (activity) => + (activity.payload as { sourceAgentThreadId?: string }).sourceAgentThreadId, + ), + ); + return childProviderThreadIds.every((childProviderThreadId) => + sources.has(childProviderThreadId), + ); + }); + expect(await readSubagentResultEvents()).toHaveLength(3); + + for (let updateIndex = 1; updateIndex < 10; updateIndex += 1) { + for (const [childIndex, childProviderThreadId] of childProviderThreadIds.entries()) { + harness.emit({ + type: "content.delta", + eventId: asEventId(`evt-three-streaming-children-${childIndex}-${updateIndex}`), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-three-streaming-children"), + itemId: asItemId(`child-message-${childIndex}`), + providerRefs: { + providerThreadId: childProviderThreadId, + providerTurnId: `child-turn-${childIndex}`, + providerItemId: asItemId(`child-message-${childIndex}`), + }, + payload: { + streamKind: "assistant_text", + delta: `|${updateIndex}`, + }, + }); + } + } + + const thread = await waitForThread(harness.readModel, (entry) => + childProviderThreadIds.every((childProviderThreadId) => + entry.activities.some((activity) => { + const payload = activity.payload as { + sourceAgentThreadId?: string; + data?: { subagentLiveText?: string }; + }; + return ( + activity.kind === "subagent.result" && + payload.sourceAgentThreadId === childProviderThreadId && + payload.data?.subagentLiveText === `${childProviderThreadId}:0|1|2|3|4|5|6|7|8|9` + ); + }), + ), + ); + expect( + thread.activities.filter((activity) => activity.kind === "subagent.result"), + ).toHaveLength(3); + + expect(await readSubagentResultEvents()).toHaveLength(6); }); it("uses the durable child roster when the projected parent provider id is missing", async () => { diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 14d9462f3..ec4f2c161 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -24,6 +24,7 @@ import * as Layer from "effect/Layer"; import * as Metric from "effect/Metric"; import * as Option from "effect/Option"; import * as Path from "effect/Path"; +import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; import * as Stream from "effect/Stream"; import { countUnifiedDiffStats, type DiffLineStats } from "@threadlines/shared/diffStats"; @@ -72,6 +73,7 @@ const BUFFERED_PROPOSED_PLAN_BY_ID_TTL = Duration.minutes(120); const BUFFERED_ACTIVITY_STREAM_BY_KEY_CACHE_CAPACITY = 20_000; const BUFFERED_ACTIVITY_STREAM_BY_KEY_TTL = Duration.minutes(120); const STREAMING_ASSISTANT_DELTA_FLUSH_INTERVAL = Duration.millis(50); +const SUBAGENT_RESULT_ACTIVITY_FLUSH_INTERVAL = Duration.millis(100); const MARKDOWN_FENCE_INDENT_LIMIT = 3; type ContentDeltaStreamKind = Extract< ProviderRuntimeEvent, @@ -118,6 +120,15 @@ interface PendingStreamingAssistantMessage { readonly createdAt: string; } +interface PendingSubagentResultActivity { + readonly event: ProviderRuntimeEvent; + readonly threadId: ThreadId; + readonly parentProviderThreadId?: string | undefined; + readonly childProviderThreadId: string; + readonly body: string; + readonly createdAt: string; +} + interface EmittedActivityStreamCursor { readonly byteCount: number; readonly lineCount: number; @@ -142,6 +153,9 @@ type RuntimeIngestionInput = } | { source: "flush"; + } + | { + source: "subagent-result-flush"; }; function toTurnId(value: TurnId | string | undefined): TurnId | undefined { @@ -1189,9 +1203,28 @@ const make = Effect.gen(function* () { timeToLive: BUFFERED_SUBAGENT_RESULT_TEXT_BY_KEY_TTL, lookup: () => Effect.succeed(""), }); + const emittedSubagentResultActivityKeys = yield* Cache.make({ + capacity: BUFFERED_SUBAGENT_RESULT_TEXT_BY_KEY_CACHE_CAPACITY, + timeToLive: BUFFERED_SUBAGENT_RESULT_TEXT_BY_KEY_TTL, + lookup: () => + Effect.die( + new Error("subagent result activity keys should only be read after they are emitted"), + ), + }); + const settledSubagentResultActivityKeys = yield* Cache.make({ + capacity: BUFFERED_SUBAGENT_RESULT_TEXT_BY_KEY_CACHE_CAPACITY, + timeToLive: BUFFERED_SUBAGENT_RESULT_TEXT_BY_KEY_TTL, + lookup: () => + Effect.die(new Error("settled subagent result keys should only be read after completion")), + }); const pendingStreamingAssistantMessages = yield* Ref.make( new Map(), ); + const pendingSubagentResultActivitiesByKey = yield* Ref.make( + new Map(), + ); + const subagentResultFlushQueued = yield* Ref.make(false); + const subagentResultFlushRequests = yield* Queue.unbounded(); const realtimeAssistantMessageIds = yield* Ref.make(new Map()); const streamingAssistantTableMessageIds = yield* Ref.make(new Set()); @@ -1603,7 +1636,8 @@ const make = Effect.gen(function* () { const dispatchSubagentResultActivity = (input: { event: ProviderRuntimeEvent; - thread: Pick; + threadId: ThreadId; + parentProviderThreadId?: string | undefined; childProviderThreadId: string; body: string; status: "inProgress" | "completed"; @@ -1612,7 +1646,7 @@ const make = Effect.gen(function* () { const activity = subagentResultActivity({ event: input.event, childProviderThreadId: input.childProviderThreadId, - parentProviderThreadId: input.thread.session?.providerThreadId ?? undefined, + parentProviderThreadId: input.parentProviderThreadId, body: input.body, status: input.status, createdAt: input.createdAt, @@ -1627,12 +1661,172 @@ const make = Effect.gen(function* () { input.event, input.status === "completed" ? "subagent-result-complete" : "subagent-result-update", ), - threadId: input.thread.id, + threadId: input.threadId, activity, createdAt: activity.createdAt, }); }; + const queuePendingSubagentResultActivity = (key: string, input: PendingSubagentResultActivity) => + Ref.update(pendingSubagentResultActivitiesByKey, (pending) => { + const next = new Map(pending); + next.set(key, input); + return next; + }); + + const pendingSubagentResultActivities = ( + matches: (key: string, input: PendingSubagentResultActivity) => boolean = () => true, + ) => + Ref.get(pendingSubagentResultActivitiesByKey).pipe( + Effect.map((pending) => [...pending].filter(([key, input]) => matches(key, input))), + ); + + const clearPendingSubagentResultActivity = (key: string) => + Ref.update(pendingSubagentResultActivitiesByKey, (pending) => { + if (!pending.has(key)) { + return pending; + } + const next = new Map(pending); + next.delete(key); + return next; + }); + + const requestSubagentResultFlushIfNeeded = Effect.gen(function* () { + const pending = yield* Ref.get(pendingSubagentResultActivitiesByKey); + if (pending.size === 0) { + return; + } + const shouldRequestFlush = yield* Ref.modify(subagentResultFlushQueued, (queued) => + queued ? ([false, true] as const) : ([true, true] as const), + ); + if (shouldRequestFlush) { + yield* Queue.offer(subagentResultFlushRequests, undefined); + } + }); + + const flushPendingSubagentResultActivities = ( + matches?: (key: string, input: PendingSubagentResultActivity) => boolean, + ) => + Effect.gen(function* () { + const pending = yield* pendingSubagentResultActivities(matches); + yield* Effect.forEach( + pending, + ([key, input]) => + dispatchSubagentResultActivity({ + ...input, + status: "inProgress", + }).pipe(Effect.andThen(clearPendingSubagentResultActivity(key))), + { concurrency: 1 }, + ).pipe(Effect.asVoid); + }); + + const flushTimedSubagentResultActivities = flushPendingSubagentResultActivities().pipe( + Effect.ensuring( + Ref.set(subagentResultFlushQueued, false).pipe( + Effect.andThen(requestSubagentResultFlushIfNeeded), + ), + ), + ); + + const queueSubagentResultActivity = (key: string, input: PendingSubagentResultActivity) => + Effect.gen(function* () { + if (!hasRenderableAssistantText(input.body)) { + return; + } + const alreadyEmitted = yield* Cache.getOption(emittedSubagentResultActivityKeys, key); + if (Option.isNone(alreadyEmitted)) { + yield* dispatchSubagentResultActivity({ + ...input, + status: "inProgress", + }); + yield* Cache.set(emittedSubagentResultActivityKeys, key, true); + return; + } + yield* queuePendingSubagentResultActivity(key, input); + yield* requestSubagentResultFlushIfNeeded; + }); + + /** Turn completion seals every known result item against late text deltas, + * but retains its buffer so an out-of-order item.completed can still publish + * the full terminal body. */ + const sealSubagentResultState = (prefix: string) => + Effect.gen(function* () { + const keys = new Set(); + yield* Ref.update(pendingSubagentResultActivitiesByKey, (pending) => { + const next = new Map(pending); + for (const key of next.keys()) { + if (key.startsWith(prefix)) { + keys.add(key); + next.delete(key); + } + } + return next; + }); + for (const key of yield* Cache.keys(bufferedSubagentResultTextByKey)) { + if (key.startsWith(prefix)) { + keys.add(key); + } + } + for (const key of yield* Cache.keys(emittedSubagentResultActivityKeys)) { + if (key.startsWith(prefix)) { + keys.add(key); + } + } + yield* Effect.forEach( + keys, + (key) => + Cache.set(settledSubagentResultActivityKeys, key, true).pipe( + Effect.andThen(Cache.invalidate(emittedSubagentResultActivityKeys, key)), + ), + { concurrency: 1 }, + ).pipe(Effect.asVoid); + }); + + const clearSettledSubagentResultState = (prefix: string) => + Effect.gen(function* () { + const keys = yield* Cache.keys(settledSubagentResultActivityKeys); + yield* Effect.forEach( + keys, + (key) => + key.startsWith(prefix) + ? Cache.invalidate(settledSubagentResultActivityKeys, key) + : Effect.void, + { concurrency: 1 }, + ).pipe(Effect.asVoid); + }); + + const clearSubagentResultState = (prefix: string) => + Effect.gen(function* () { + yield* Ref.update(pendingSubagentResultActivitiesByKey, (pending) => { + const next = new Map(pending); + for (const key of next.keys()) { + if (key.startsWith(prefix)) { + next.delete(key); + } + } + return next; + }); + + const bufferedKeys = Array.from(yield* Cache.keys(bufferedSubagentResultTextByKey)); + const emittedKeys = Array.from(yield* Cache.keys(emittedSubagentResultActivityKeys)); + yield* Effect.forEach( + bufferedKeys, + (key) => + key.startsWith(prefix) + ? Cache.invalidate(bufferedSubagentResultTextByKey, key) + : Effect.void, + { concurrency: 1 }, + ).pipe(Effect.asVoid); + yield* Effect.forEach( + emittedKeys, + (key) => + key.startsWith(prefix) + ? Cache.invalidate(emittedSubagentResultActivityKeys, key) + : Effect.void, + { concurrency: 1 }, + ).pipe(Effect.asVoid); + }); + const queuePendingStreamingAssistantMessage = (input: PendingStreamingAssistantMessage) => Ref.update(pendingStreamingAssistantMessages, (pending) => { const next = new Map(pending); @@ -2365,11 +2559,20 @@ const make = Effect.gen(function* () { // An agent that does report again re-opens its row through the normal // fold, since later lifecycle activities win over this one. if (event.type === "session.started" || event.type === "session.exited") { + yield* flushPendingSubagentResultActivities( + (_key, pending) => pending.threadId === thread.id, + ); yield* settleLiveSubagents({ commandTag: "subagent-orphan", summary: "Subagent no longer tracked by the provider session", belongsToLifecycle: () => true, }); + if (event.type === "session.started") { + yield* clearSubagentResultState(`${thread.id}:`); + yield* clearSettledSubagentResultState(`${thread.id}:`); + } else { + yield* sealSubagentResultState(`${thread.id}:`); + } } if ( @@ -2458,6 +2661,16 @@ const make = Effect.gen(function* () { event.type === "turn.aborted" || completedTurnState === "interrupted" || completedTurnState === "cancelled"; + if ( + (event.type === "turn.completed" || event.type === "turn.aborted") && + shouldApplyThreadLifecycle && + eventTurnId !== undefined + ) { + yield* flushPendingSubagentResultActivities( + (_key, pending) => + pending.threadId === thread.id && sameId(pending.event.turnId, eventTurnId), + ); + } if (turnWasInterrupted && shouldApplyThreadLifecycle && eventTurnId !== undefined) { yield* settleLiveSubagents({ commandTag: "subagent-turn-aborted", @@ -2705,15 +2918,18 @@ const make = Effect.gen(function* () { const childProviderThreadId = childProviderThreadIdForEvent(event, attributionThread); if (childProviderThreadId) { const bufferKey = subagentResultTextKey({ event, childProviderThreadId }); - const body = yield* appendBufferedSubagentResultText(bufferKey, assistantDelta); - yield* dispatchSubagentResultActivity({ - event, - thread, - childProviderThreadId, - body, - status: "inProgress", - createdAt: now, - }); + const settled = yield* Cache.getOption(settledSubagentResultActivityKeys, bufferKey); + if (Option.isNone(settled)) { + const body = yield* appendBufferedSubagentResultText(bufferKey, assistantDelta); + yield* queueSubagentResultActivity(bufferKey, { + event, + threadId: thread.id, + parentProviderThreadId: thread.session?.providerThreadId ?? undefined, + childProviderThreadId, + body, + createdAt: now, + }); + } } else { const turnId = toTurnId(event.turnId); const assistantMessageId = @@ -2842,6 +3058,8 @@ const make = Effect.gen(function* () { const childProviderThreadId = childProviderThreadIdForEvent(event, attributionThread); if (childProviderThreadId) { const bufferKey = subagentResultTextKey({ event, childProviderThreadId }); + yield* clearPendingSubagentResultActivity(bufferKey); + yield* Cache.invalidate(emittedSubagentResultActivityKeys, bufferKey); const bufferedText = yield* takeBufferedSubagentResultText(bufferKey); const body = bufferedText.length > 0 @@ -2851,7 +3069,8 @@ const make = Effect.gen(function* () { : ""; yield* dispatchSubagentResultActivity({ event, - thread, + threadId: thread.id, + parentProviderThreadId: thread.session?.providerThreadId ?? undefined, childProviderThreadId, body, // Codex completes each assistant-message item, including interim @@ -2862,6 +3081,7 @@ const make = Effect.gen(function* () { assistantMessagePhaseFromEvent(event) === "commentary" ? "inProgress" : "completed", createdAt: now, }); + yield* Cache.set(settledSubagentResultActivityKeys, bufferKey, true); } else { const detailedThread = yield* getLoadedThreadDetail(); const messages = detailedThread?.messages ?? []; @@ -3009,6 +3229,9 @@ const make = Effect.gen(function* () { turnId, updatedAt: now, }); + if (shouldApplyThreadLifecycle) { + yield* sealSubagentResultState(`${thread.id}:${turnId}:`); + } } } @@ -3209,6 +3432,8 @@ const make = Effect.gen(function* () { return processDomainEvent(input.event); case "flush": return flushPendingStreamingAssistantMessages.pipe(Effect.asVoid); + case "subagent-result-flush": + return flushTimedSubagentResultActivities; } }; @@ -3220,9 +3445,13 @@ const make = Effect.gen(function* () { } return Effect.logWarning("provider runtime ingestion failed to process event", { source: input.source, - ...(input.source === "flush" - ? { eventId: "streaming-assistant-flush", eventType: "flush" } - : { eventId: input.event.eventId, eventType: input.event.type }), + ...(input.source === "runtime" || input.source === "domain" + ? { eventId: input.event.eventId, eventType: input.event.type } + : { + eventId: + input.source === "flush" ? "streaming-assistant-flush" : "subagent-result-flush", + eventType: "flush", + }), cause: Cause.pretty(cause), }); }), @@ -3259,11 +3488,22 @@ const make = Effect.gen(function* () { Effect.forever, ), ); + yield* Effect.forkScoped( + Queue.take(subagentResultFlushRequests).pipe( + Effect.andThen(Effect.sleep(SUBAGENT_RESULT_ACTIVITY_FLUSH_INTERVAL)), + Effect.andThen( + worker.enqueue({ + source: "subagent-result-flush", + }), + ), + Effect.forever, + ), + ); }); return { start, - drain: worker.drain, + drain: worker.enqueue({ source: "subagent-result-flush" }).pipe(Effect.andThen(worker.drain)), } satisfies ProviderRuntimeIngestionShape; }); diff --git a/apps/web/src/agentsPanelStore.ts b/apps/web/src/agentsPanelStore.ts index e57cf8a46..2907ed3e9 100644 --- a/apps/web/src/agentsPanelStore.ts +++ b/apps/web/src/agentsPanelStore.ts @@ -30,10 +30,6 @@ export interface AgentsPanelSource { * the live items so the panel and the conversation's receipts resolve the * same set of agents. */ history: ReadonlyArray; - /** Child-owned work remains in this durable log even though the main - * conversation suppresses it. The inspector uses it to fill provider - * transcript gaps. */ - workEntries: ReadonlyArray; /** Provider driver label, e.g. `codex`; drives the trunk hue and run chips. */ providerLabel: string | null; /** True from the moment a turn is dispatched until it settles. Lets the panel @@ -49,8 +45,19 @@ export interface AgentsPanelSource { onStopBackgroundRun: (run: ThreadBackgroundRunItem) => void; } +export interface AgentsPanelActivitySource { + environmentId: EnvironmentId; + threadId: ThreadId; + /** Child-owned work remains in this durable log even though the main + * conversation suppresses it. Only a drilled-in inspector subscribes. */ + workEntries: ReadonlyArray; +} + +const EMPTY_WORK_ENTRIES: ReadonlyArray = []; + interface AgentsPanelStoreState { source: AgentsPanelSource | null; + activitySource: AgentsPanelActivitySource | null; /** * The agent whose transcript the panel is drilled into. It lives here rather * than inside the panel because the drill-in is reachable from outside it: @@ -59,11 +66,13 @@ interface AgentsPanelStoreState { */ selectedAgentId: string | null; publishSource: (source: AgentsPanelSource | null) => void; + publishActivitySource: (source: AgentsPanelActivitySource | null) => void; selectAgent: (agentId: string | null) => void; } export const useAgentsPanelStore = create((set) => ({ source: null, + activitySource: null, selectedAgentId: null, publishSource: (source) => { set((state) => ({ @@ -73,6 +82,9 @@ export const useAgentsPanelStore = create((set) => ({ selectedAgentId: state.source?.threadId === source?.threadId ? state.selectedAgentId : null, })); }, + publishActivitySource: (activitySource) => { + set({ activitySource }); + }, selectAgent: (agentId) => { set({ selectedAgentId: agentId }); }, @@ -86,14 +98,36 @@ export function useSelectedAgentId(): string | null { return useAgentsPanelStore((state) => state.selectedAgentId); } +/** Detailed child activity changes much more often than the tree. Keep it out + * of the route subscription and wake only the inspector that can render it. */ +export function useAgentsPanelWorkEntries(input: { + environmentId: EnvironmentId; + threadId: ThreadId; + enabled?: boolean; +}): ReadonlyArray { + return useAgentsPanelStore((state) => { + if (input.enabled === false) { + return EMPTY_WORK_ENTRIES; + } + const source = state.activitySource; + return source?.environmentId === input.environmentId && source.threadId === input.threadId + ? source.workEntries + : EMPTY_WORK_ENTRIES; + }); +} + export function publishAgentsPanelSource(source: AgentsPanelSource | null): void { useAgentsPanelStore.getState().publishSource(source); } +export function publishAgentsPanelActivitySource(source: AgentsPanelActivitySource | null): void { + useAgentsPanelStore.getState().publishActivitySource(source); +} + export function selectAgentsPanelAgent(agentId: string | null): void { useAgentsPanelStore.getState().selectAgent(agentId); } export function resetAgentsPanelSourceForTests(): void { - useAgentsPanelStore.setState({ source: null, selectedAgentId: null }); + useAgentsPanelStore.setState({ source: null, activitySource: null, selectedAgentId: null }); } diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index a0335f1e4..c81ec606b 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -80,10 +80,7 @@ import { deriveActiveStatusLabel, deriveActiveWorkStartedAt, deriveActivePlanState, - deriveSubagentProgressState, - deriveSubagentLiveEntries, - deriveSubagentResultEntries, - deriveThreadSubagentHistory, + deriveSubagentActivityState, findSidebarProposedPlan, findLatestProposedPlan, deriveWorkLogEntries, @@ -139,7 +136,11 @@ import { draftRightPanelStateKey, useChatHeaderBottomVarRef, } from "../rightPanelLayout"; -import { publishAgentsPanelSource, selectAgentsPanelAgent } from "../agentsPanelStore"; +import { + publishAgentsPanelActivitySource, + publishAgentsPanelSource, + selectAgentsPanelAgent, +} from "../agentsPanelStore"; import { summarizeLiveAgents } from "./chat/agentsPanel.logic"; import { buildTemporaryWorktreeBranchName } from "@threadlines/shared/git"; import { BranchToolbar } from "./BranchToolbar"; @@ -2001,9 +2002,9 @@ export default function ChatView(props: ChatViewProps) { : null, [activePlan, taskProgressBadge, taskProgressLabel, taskProgressProposedPlan], ); - const subagentProgress = useMemo( + const subagentActivityState = useMemo( () => - deriveSubagentProgressState({ + deriveSubagentActivityState({ activities: threadActivities, subagents: activeThread?.subagents ?? [], latestTurnId: activeLatestTurn?.turnId ?? null, @@ -2011,13 +2012,11 @@ export default function ChatView(props: ChatViewProps) { }), [activeLatestTurn?.turnId, activeThread?.subagents, latestTurnSettled, threadActivities], ); + const subagentProgress = subagentActivityState.progress; // The turn-scoped progress above empties when the turn settles. The panel's // history section and the conversation's receipts both need the thread's whole // roster, which the same activities answer without the turn filter. - const subagentHistory = useMemo( - () => deriveThreadSubagentHistory(threadActivities, activeThread?.subagents), - [activeThread?.subagents, threadActivities], - ); + const subagentHistory = subagentActivityState.history; const showPlanFollowUpPrompt = pendingUserInputs.length === 0 && interactionMode === "plan" && @@ -2344,15 +2343,15 @@ export default function ChatView(props: ChatViewProps) { timelineMessages, activeThread?.proposedPlans ?? [], workLogEntries, - deriveSubagentResultEntries(threadActivities, activeThread?.subagents), + subagentActivityState.resultEntries, forkContextEntries, - deriveSubagentLiveEntries(threadActivities, activeThread?.subagents), + subagentActivityState.liveEntries, ), [ activeThread?.proposedPlans, - activeThread?.subagents, forkContextEntries, - threadActivities, + subagentActivityState.liveEntries, + subagentActivityState.resultEntries, timelineMessages, workLogEntries, ], @@ -3397,7 +3396,6 @@ export default function ChatView(props: ChatViewProps) { subagents: subagentProgress?.items ?? EMPTY_SUBAGENT_ITEMS, subagentRuns: promotedSubagentRuns, history: subagentHistory, - workEntries: workLogEntries, providerLabel: activeProviderDriver, turnInFlight: activeTurnInProgress, hydrated: threadDetailHydrated, @@ -3415,9 +3413,25 @@ export default function ChatView(props: ChatViewProps) { subagentHistory, subagentProgress?.items, threadDetailHydrated, - workLogEntries, ]); - useEffect(() => () => publishAgentsPanelSource(null), []); + useEffect(() => { + if (!activeThreadId) { + publishAgentsPanelActivitySource(null); + return; + } + publishAgentsPanelActivitySource({ + environmentId, + threadId: activeThreadId, + workEntries: workLogEntries, + }); + }, [activeThreadId, environmentId, workLogEntries]); + useEffect( + () => () => { + publishAgentsPanelSource(null); + publishAgentsPanelActivitySource(null); + }, + [], + ); const confirmPendingTerminalKill = useCallback(() => { if (!pendingTerminalKill) return; diff --git a/apps/web/src/components/chat/AgentsPanel.browser.tsx b/apps/web/src/components/chat/AgentsPanel.browser.tsx index d991615ca..f212289aa 100644 --- a/apps/web/src/components/chat/AgentsPanel.browser.tsx +++ b/apps/web/src/components/chat/AgentsPanel.browser.tsx @@ -10,7 +10,10 @@ import type { ThreadSubagentHistoryEntry, WorkLogEntry, } from "../../session-logic"; -import { resetAgentsPanelSourceForTests } from "../../agentsPanelStore"; +import { + publishAgentsPanelActivitySource, + resetAgentsPanelSourceForTests, +} from "../../agentsPanelStore"; import { AgentsPanel } from "./AgentsPanel"; import { ChatRightPanel } from "../ChatRightPanel"; import { buildRightPanelLauncherStates } from "./rightPanelLauncherState"; @@ -866,19 +869,42 @@ describe("AgentsPanel", () => { toolCallId: "call-2", }, ]; + publishAgentsPanelActivitySource({ + environmentId: ENVIRONMENT_ID, + threadId: ThreadId.make("another-thread"), + workEntries: workEntries.slice(0, 1), + }); const mounted = await renderPanel({ subagents: [buildSubagent({ label: "Router sweep" })], - workEntries, }); try { await page.getByRole("button", { name: "Open Router sweep transcript" }).click(); await expect.element(page.getByText("I am checking the route.")).toBeVisible(); + expect( + document.querySelector("[data-subagent-transcript-activity-toggle='true']"), + ).toBeNull(); + + publishAgentsPanelActivitySource({ + environmentId: ENVIRONMENT_ID, + threadId: THREAD_ID, + workEntries: workEntries.slice(0, 1), + }); const receipt = await vi.waitUntil(() => document.querySelector("[data-subagent-transcript-activity-toggle='true']"), ); expect(receipt.textContent).toContain("Activity"); + expect(receipt.textContent).toContain("1 action"); + + publishAgentsPanelActivitySource({ + environmentId: ENVIRONMENT_ID, + threadId: THREAD_ID, + workEntries, + }); + await vi.waitFor(() => { + expect(receipt.textContent).toContain("2 actions"); + }); expect(receipt.textContent).toContain("2 actions"); expect(receipt.textContent).toContain("Running command"); expect(receipt.textContent).toContain("vp run typecheck"); diff --git a/apps/web/src/components/chat/AgentsPanel.tsx b/apps/web/src/components/chat/AgentsPanel.tsx index f6aa78713..def16f1ef 100644 --- a/apps/web/src/components/chat/AgentsPanel.tsx +++ b/apps/web/src/components/chat/AgentsPanel.tsx @@ -8,7 +8,11 @@ import type { WorkLogEntry, } from "../../session-logic"; import { cn } from "~/lib/utils"; -import { selectAgentsPanelAgent, useSelectedAgentId } from "../../agentsPanelStore"; +import { + selectAgentsPanelAgent, + useAgentsPanelWorkEntries, + useSelectedAgentId, +} from "../../agentsPanelStore"; import type { Icon } from "../Icons"; import { Button } from "../ui/button"; import { LiveNode, SectionLabel } from "../ui/threadline"; @@ -63,6 +67,60 @@ function providerAcceptsSubagentInput(providerLabel: string | null | undefined): return providerLabel?.trim().toLowerCase().includes("codex") ?? false; } +function ConnectedSubagentInspector({ + environmentId, + threadId, + item, + workEntries, + providerLabel, + threadCwd, + onMessageThroughParent, + onClose, +}: { + environmentId: EnvironmentId; + threadId: ThreadId; + item: SubagentProgressItem; + workEntries: ReadonlyArray | undefined; + providerLabel: string | null | undefined; + threadCwd: string | null | undefined; + onMessageThroughParent: ((agentName: string) => void) | undefined; + onClose: () => void; +}) { + const publishedWorkEntries = useAgentsPanelWorkEntries({ + environmentId, + threadId, + enabled: workEntries === undefined, + }); + const selectedWorkEntries = useMemo(() => { + const agentThreadId = item.agentThreadId; + if (!agentThreadId) { + return EMPTY_WORK_ENTRIES; + } + return (workEntries ?? publishedWorkEntries).filter( + (entry) => + entry.sourceAgentThreadId === agentThreadId || + // Background tasks the agent started in its own conversation: for + // Claude the owner spawn call id is the agent's thread id. + entry.ownerAgentToolUseId === agentThreadId, + ); + }, [item.agentThreadId, publishedWorkEntries, workEntries]); + + return ( + + ); +} + /** The trunk takes the provider's own hue so the panel reads as that * provider's work; anything else falls back to the hairline colour. */ function trunkColor(providerLabel: string | null | undefined): string { @@ -303,7 +361,7 @@ export const AgentsPanel = memo(function AgentsPanel({ subagents, subagentRuns, history, - workEntries = EMPTY_WORK_ENTRIES, + workEntries, providerLabel, turnInFlight = false, threadCwd, @@ -326,21 +384,6 @@ export const AgentsPanel = memo(function AgentsPanel({ const providerGlyph = useMemo(() => providerIconForDriverLabel(providerLabel), [providerLabel]); const anyRunning = hasRunningAgentActivity({ subagents }); const selectedSubagent = findAgentsPanelSubagent(view, selectedAgentId); - const selectedSubagentThreadId = selectedSubagent?.agentThreadId ?? null; - const selectedSubagentWorkEntries = useMemo( - () => - selectedSubagentThreadId - ? workEntries.filter( - (entry) => - entry.sourceAgentThreadId === selectedSubagentThreadId || - // Background tasks the agent started in its own conversation: - // for Claude agents the owner spawn call id is the agent's - // thread id, so they belong to the same drill-in view. - entry.ownerAgentToolUseId === selectedSubagentThreadId, - ) - : [], - [selectedSubagentThreadId, workEntries], - ); const handleSelect = useCallback((branch: AgentBranch) => { if (branch.item.agentThreadId) { @@ -362,16 +405,14 @@ export const AgentsPanel = memo(function AgentsPanel({ }, []); const inspector = selectedSubagent ? ( - diff --git a/apps/web/src/routes/_chat.$environmentId.$threadId.tsx b/apps/web/src/routes/_chat.$environmentId.$threadId.tsx index cd8669151..9a045f6dc 100644 --- a/apps/web/src/routes/_chat.$environmentId.$threadId.tsx +++ b/apps/web/src/routes/_chat.$environmentId.$threadId.tsx @@ -500,7 +500,6 @@ function ChatThreadRouteView() { subagents={agentsSource?.subagents ?? EMPTY_SUBAGENTS} subagentRuns={agentsSource?.subagentRuns} history={agentsSource?.history ?? EMPTY_SUBAGENT_HISTORY} - workEntries={agentsSource?.workEntries} providerLabel={agentsSource?.providerLabel} turnInFlight={agentsSource?.turnInFlight ?? false} threadCwd={agentsSource?.threadCwd} diff --git a/apps/web/src/session-logic.test.ts b/apps/web/src/session-logic.test.ts index 836eebc4d..2801a37b2 100644 --- a/apps/web/src/session-logic.test.ts +++ b/apps/web/src/session-logic.test.ts @@ -17,6 +17,7 @@ import { derivePendingApprovals, derivePendingUserInputs, isBlockingUserInput, + deriveSubagentActivityState, deriveSubagentLiveEntries, deriveSubagentProgressState, deriveSubagentResultEntries, @@ -328,6 +329,27 @@ describe("deriveThreadSubagentHistory", () => { expect(history[0]?.resultBody).toBe("50 .tsx files."); }); + it("derives every active-chat subagent view from the shared activity state", () => { + const activities = codexSpawnActivities(); + const latestTurnId = TurnId.make("turn-1"); + const state = deriveSubagentActivityState({ + activities, + latestTurnId, + latestTurnSettled: false, + }); + + expect(state).toEqual({ + progress: deriveSubagentProgressState({ + activities, + latestTurnId, + latestTurnSettled: false, + }), + history: deriveThreadSubagentHistory(activities), + resultEntries: deriveSubagentResultEntries(activities), + liveEntries: deriveSubagentLiveEntries(activities), + }); + }); + it("keeps durable child identity and explicit spawn settings after lifecycle rows age out", () => { const durable: OrchestrationSubagent = { id: "codex-child-1", diff --git a/apps/web/src/session-logic.ts b/apps/web/src/session-logic.ts index 2602f5a36..acc1536aa 100644 --- a/apps/web/src/session-logic.ts +++ b/apps/web/src/session-logic.ts @@ -855,13 +855,29 @@ export function deriveSubagentProgressState(input: { }): SubagentProgressState | null { const records = collectSubagentActivityRecords(input.activities, { subagents: input.subagents, - latestTurnId: input.latestTurnId ?? null, }); + return deriveSubagentProgressStateFromRecords(records, input); +} + +function deriveSubagentProgressStateFromRecords( + records: ReadonlyArray, + input: { + latestTurnId?: TurnId | null | undefined; + latestTurnSettled?: boolean | undefined; + }, +): SubagentProgressState | null { + const latestTurnId = input.latestTurnId ?? null; + const scopedRecords = records.filter( + (record) => + latestTurnId === null || + record.turnId === latestTurnId || + isActiveSubagentStatus(record.status), + ); // Finished agents remain useful while their parent turn is still running: // they explain a shrinking active count and make the completed badge/state // reachable. Once the turn settles, successful agents clear with the rest // of the transient activity UI while failed and stopped work remains visible. - const visibleRecords = records.filter( + const visibleRecords = scopedRecords.filter( (record) => record.status !== "completed" || input.latestTurnSettled === false, ); const items = visibleRecords.map(toSubagentProgressItem); @@ -945,7 +961,15 @@ export function deriveThreadSubagentHistory( activities: ReadonlyArray, subagents: ReadonlyArray = [], ): ThreadSubagentHistoryEntry[] { - return collectSubagentActivityRecords(activities, { subagents }).map((record) => ({ + return deriveThreadSubagentHistoryFromRecords( + collectSubagentActivityRecords(activities, { subagents }), + ); +} + +function deriveThreadSubagentHistoryFromRecords( + records: ReadonlyArray, +): ThreadSubagentHistoryEntry[] { + return records.map((record) => ({ item: toSubagentProgressItem(record), resultBody: record.resultBody, })); @@ -955,7 +979,15 @@ export function deriveSubagentResultEntries( activities: ReadonlyArray, subagents: ReadonlyArray = [], ): SubagentResultEntry[] { - return collectSubagentActivityRecords(activities, { subagents }) + return deriveSubagentResultEntriesFromRecords( + collectSubagentActivityRecords(activities, { subagents }), + ); +} + +function deriveSubagentResultEntriesFromRecords( + records: ReadonlyArray, +): SubagentResultEntry[] { + return records .filter( ( record, @@ -993,7 +1025,15 @@ export function deriveSubagentLiveEntries( activities: ReadonlyArray, subagents: ReadonlyArray = [], ): SubagentLiveEntry[] { - return collectSubagentActivityRecords(activities, { subagents }) + return deriveSubagentLiveEntriesFromRecords( + collectSubagentActivityRecords(activities, { subagents }), + ); +} + +function deriveSubagentLiveEntriesFromRecords( + records: ReadonlyArray, +): SubagentLiveEntry[] { + return records .filter( ( record, @@ -1023,6 +1063,33 @@ export function deriveSubagentLiveEntries( .toSorted((left, right) => left.createdAt.localeCompare(right.createdAt)); } +export interface SubagentActivityState { + readonly progress: SubagentProgressState | null; + readonly history: ThreadSubagentHistoryEntry[]; + readonly resultEntries: SubagentResultEntry[]; + readonly liveEntries: SubagentLiveEntry[]; +} + +/** Derives every subagent view from one ordered activity fold. The active chat + * consumes all four results together, so sorting and scanning the same + * activity history separately only adds work to every streamed update. */ +export function deriveSubagentActivityState(input: { + activities: ReadonlyArray; + subagents?: ReadonlyArray; + latestTurnId?: TurnId | null | undefined; + latestTurnSettled?: boolean | undefined; +}): SubagentActivityState { + const records = collectSubagentActivityRecords(input.activities, { + subagents: input.subagents, + }); + return { + progress: deriveSubagentProgressStateFromRecords(records, input), + history: deriveThreadSubagentHistoryFromRecords(records), + resultEntries: deriveSubagentResultEntriesFromRecords(records), + liveEntries: deriveSubagentLiveEntriesFromRecords(records), + }; +} + function asForkContextPayload(payload: unknown): ThreadForkContextPayload | null { const record = asRecord(payload); if (!record) { @@ -1198,13 +1265,11 @@ function collectTurnModelSelections( function collectSubagentActivityRecords( activities: ReadonlyArray, options: { - latestTurnId?: TurnId | null | undefined; subagents?: ReadonlyArray | undefined; }, ): InternalSubagentRecord[] { const byAgentId = new Map(); const pendingSpawnKeysByCallId = new Map(); - const latestTurnId = options.latestTurnId ?? null; const sortedActivities = [...activities].toSorted(compareActivitiesByOrder); // The server owns an uncapped roster so identity and spawn settings survive @@ -1490,17 +1555,10 @@ function collectSubagentActivityRecords( } } - return [...byAgentId.values()] - .filter( - (record) => - latestTurnId === null || - record.turnId === latestTurnId || - isActiveSubagentStatus(record.status), - ) - .toSorted((left, right) => { - const createdAtComparison = left.createdAt.localeCompare(right.createdAt); - return createdAtComparison === 0 ? left.id.localeCompare(right.id) : createdAtComparison; - }); + return [...byAgentId.values()].toSorted((left, right) => { + const createdAtComparison = left.createdAt.localeCompare(right.createdAt); + return createdAtComparison === 0 ? left.id.localeCompare(right.id) : createdAtComparison; + }); } function collectSubagentActivityTelemetry(