From 62ce5d69e93f4f6432df29e6681d55f2b2439940 Mon Sep 17 00:00:00 2001 From: tarik02 Date: Sat, 4 Jul 2026 12:25:28 +0300 Subject: [PATCH 1/3] revert fork cursor fixes --- .../Layers/ProviderRuntimeIngestion.test.ts | 400 ++---------------- .../Layers/ProviderRuntimeIngestion.ts | 116 +---- .../provider/acp/AcpSessionRuntime.test.ts | 84 ---- .../src/provider/acp/AcpSessionRuntime.ts | 13 - .../src/provider/acp/CursorAcpSupport.ts | 2 - 5 files changed, 42 insertions(+), 573 deletions(-) delete mode 100644 apps/server/src/provider/acp/AcpSessionRuntime.test.ts diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index f908c2c2292..001ba388949 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -61,10 +61,6 @@ const asEventId = (value: string): EventId => EventId.make(value); const asMessageId = (value: string): MessageId => MessageId.make(value); const asThreadId = (value: string): ThreadId => ThreadId.make(value); const asTurnId = (value: string): TurnId => TurnId.make(value); -const turnScopedAssistantMessageId = (turnId: string, itemId: string, segment?: number): string => - segment === undefined - ? `assistant:turn:${turnId}:${itemId}` - : `assistant:turn:${turnId}:${itemId}:segment:${segment}`; type LegacyProviderRuntimeEvent = { readonly type: string; @@ -713,12 +709,11 @@ describe("ProviderRuntimeIngestion", () => { const thread = await waitForThread(harness.readModel, (entry) => entry.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === turnScopedAssistantMessageId("turn-2", "item-1") && !message.streaming, + message.id === "assistant:item-1" && !message.streaming, ), ); const message = thread.messages.find( - (entry: ProviderRuntimeTestMessage) => - entry.id === turnScopedAssistantMessageId("turn-2", "item-1"), + (entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-1", ); expect(message?.text).toBe("hello world"); expect(message?.streaming).toBe(false); @@ -746,13 +741,11 @@ describe("ProviderRuntimeIngestion", () => { const thread = await waitForThread(harness.readModel, (entry) => entry.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === turnScopedAssistantMessageId("turn-no-delta", "item-no-delta") && - !message.streaming, + message.id === "assistant:item-no-delta" && !message.streaming, ), ); const message = thread.messages.find( - (entry: ProviderRuntimeTestMessage) => - entry.id === turnScopedAssistantMessageId("turn-no-delta", "item-no-delta"), + (entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-no-delta", ); expect(message?.text).toBe("assistant-only final text"); expect(message?.streaming).toBe(false); @@ -1622,8 +1615,7 @@ describe("ProviderRuntimeIngestion", () => { const midThread = midReadModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); expect( midThread?.messages.some( - (message: ProviderRuntimeTestMessage) => - message.id === turnScopedAssistantMessageId("turn-buffered", "item-buffered"), + (message: ProviderRuntimeTestMessage) => message.id === "assistant:item-buffered", ), ).toBe(false); @@ -1644,13 +1636,11 @@ describe("ProviderRuntimeIngestion", () => { const thread = await waitForThread(harness.readModel, (entry) => entry.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === turnScopedAssistantMessageId("turn-buffered", "item-buffered") && - !message.streaming, + message.id === "assistant:item-buffered" && !message.streaming, ), ); const message = thread.messages.find( - (entry: ProviderRuntimeTestMessage) => - entry.id === turnScopedAssistantMessageId("turn-buffered", "item-buffered"), + (entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-buffered", ); expect(message?.text).toBe("buffer me"); expect(message?.streaming).toBe(false); @@ -1705,19 +1695,13 @@ describe("ProviderRuntimeIngestion", () => { const thread = await waitForThread(harness.readModel, (entry) => entry.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId( - "turn-buffered-request-flush", - "item-buffered-request-flush", - ) && + message.id === "assistant:item-buffered-request-flush" && !message.streaming && message.text === "visible before approval", ), ); const message = thread.messages.find( - (entry: ProviderRuntimeTestMessage) => - entry.id === - turnScopedAssistantMessageId("turn-buffered-request-flush", "item-buffered-request-flush"), + (entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-buffered-request-flush", ); expect(message?.streaming).toBe(false); }); @@ -1777,22 +1761,14 @@ describe("ProviderRuntimeIngestion", () => { const thread = await waitForThread(harness.readModel, (entry) => entry.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId( - "turn-buffered-user-input-flush", - "item-buffered-user-input-flush", - ) && + message.id === "assistant:item-buffered-user-input-flush" && !message.streaming && message.text === "visible before user input", ), ); const message = thread.messages.find( (entry: ProviderRuntimeTestMessage) => - entry.id === - turnScopedAssistantMessageId( - "turn-buffered-user-input-flush", - "item-buffered-user-input-flush", - ), + entry.id === "assistant:item-buffered-user-input-flush", ); expect(message?.streaming).toBe(false); }); @@ -1852,11 +1828,7 @@ describe("ProviderRuntimeIngestion", () => { expect( thread.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId( - "turn-buffered-whitespace-request", - "item-buffered-whitespace-request", - ), + message.id === "assistant:item-buffered-whitespace-request", ), ).toBe(false); }); @@ -1913,11 +1885,7 @@ describe("ProviderRuntimeIngestion", () => { await waitForThread(harness.readModel, (entry) => entry.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId( - "turn-buffered-request-append", - "item-buffered-request-append", - ) && + message.id === "assistant:item-buffered-request-append" && !message.streaming && message.text === "first half", ), @@ -1953,32 +1921,17 @@ describe("ProviderRuntimeIngestion", () => { const thread = await waitForThread(harness.readModel, (entry) => entry.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId( - "turn-buffered-request-append", - "item-buffered-request-append", - 1, - ) && + message.id === "assistant:item-buffered-request-append:segment:1" && !message.streaming && message.text === " second half", ), ); const firstMessage = thread.messages.find( - (entry: ProviderRuntimeTestMessage) => - entry.id === - turnScopedAssistantMessageId( - "turn-buffered-request-append", - "item-buffered-request-append", - ), + (entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-buffered-request-append", ); const resumedMessage = thread.messages.find( (entry: ProviderRuntimeTestMessage) => - entry.id === - turnScopedAssistantMessageId( - "turn-buffered-request-append", - "item-buffered-request-append", - 1, - ), + entry.id === "assistant:item-buffered-request-append:segment:1", ); expect(firstMessage?.text).toBe("first half"); expect(firstMessage?.streaming).toBe(false); @@ -1993,12 +1946,7 @@ describe("ProviderRuntimeIngestion", () => { const assistantEvents = events.filter( (event): event is Extract<(typeof events)[number], { type: "thread.message-sent" }> => event.type === "thread.message-sent" && - event.payload.messageId.startsWith( - turnScopedAssistantMessageId( - "turn-buffered-request-append", - "item-buffered-request-append", - ).slice(0, "assistant:turn:turn-buffered-request-append:".length), - ), + event.payload.messageId.startsWith("assistant:item-buffered-request-append"), ); expect(assistantEvents).toHaveLength(4); expect(assistantEvents[0]?.payload.streaming).toBe(true); @@ -2006,20 +1954,12 @@ describe("ProviderRuntimeIngestion", () => { expect(assistantEvents[1]?.payload.streaming).toBe(false); expect(assistantEvents[1]?.payload.text).toBe(""); expect(assistantEvents[2]?.payload.messageId).toBe( - turnScopedAssistantMessageId( - "turn-buffered-request-append", - "item-buffered-request-append", - 1, - ), + "assistant:item-buffered-request-append:segment:1", ); expect(assistantEvents[2]?.payload.streaming).toBe(true); expect(assistantEvents[2]?.payload.text).toBe(" second half"); expect(assistantEvents[3]?.payload.messageId).toBe( - turnScopedAssistantMessageId( - "turn-buffered-request-append", - "item-buffered-request-append", - 1, - ), + "assistant:item-buffered-request-append:segment:1", ); expect(assistantEvents[3]?.payload.streaming).toBe(false); expect(assistantEvents[3]?.payload.text).toBe(""); @@ -2077,11 +2017,7 @@ describe("ProviderRuntimeIngestion", () => { await waitForThread(harness.readModel, (entry) => entry.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId( - "turn-streaming-request-segment", - "item-streaming-request-segment", - ) && + message.id === "assistant:item-streaming-request-segment" && !message.streaming && message.text === "before approval", ), @@ -2117,12 +2053,7 @@ describe("ProviderRuntimeIngestion", () => { const thread = await waitForThread(harness.readModel, (entry) => entry.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId( - "turn-streaming-request-segment", - "item-streaming-request-segment", - 1, - ) && + message.id === "assistant:item-streaming-request-segment:segment:1" && !message.streaming && message.text === " after approval", ), @@ -2130,272 +2061,17 @@ describe("ProviderRuntimeIngestion", () => { expect( thread.messages.find( (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId( - "turn-streaming-request-segment", - "item-streaming-request-segment", - ), + message.id === "assistant:item-streaming-request-segment", )?.text, ).toBe("before approval"); expect( thread.messages.find( (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId( - "turn-streaming-request-segment", - "item-streaming-request-segment", - 1, - ), + message.id === "assistant:item-streaming-request-segment:segment:1", )?.text, ).toBe(" after approval"); }); - it("keeps assistant message identities isolated across turns when provider item IDs are reused", async () => { - const harness = await createHarness(); - const now = "2026-01-01T00:00:00.000Z"; - const reusedItemId = asItemId("assistant:cursor-session:segment:4"); - - harness.emit({ - type: "content.delta", - eventId: asEventId("evt-reused-item-turn-1"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-reused-1"), - itemId: reusedItemId, - payload: { - streamKind: "assistant_text", - delta: "first turn response", - }, - }); - harness.emit({ - type: "item.completed", - eventId: asEventId("evt-reused-item-turn-1-complete"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-reused-1"), - itemId: reusedItemId, - payload: { - itemType: "assistant_message", - status: "completed", - }, - }); - - harness.emit({ - type: "content.delta", - eventId: asEventId("evt-reused-item-turn-2"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-reused-2"), - itemId: reusedItemId, - payload: { - streamKind: "assistant_text", - delta: "second turn response", - }, - }); - harness.emit({ - type: "item.completed", - eventId: asEventId("evt-reused-item-turn-2-complete"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-reused-2"), - itemId: reusedItemId, - payload: { - itemType: "assistant_message", - status: "completed", - }, - }); - - const thread = await waitForThread(harness.readModel, (entry) => - entry.messages.some( - (message: ProviderRuntimeTestMessage) => - message.id === turnScopedAssistantMessageId("turn-reused-2", String(reusedItemId)) && - !message.streaming, - ), - ); - - const firstMessage = thread.messages.find( - (message: ProviderRuntimeTestMessage) => - message.id === turnScopedAssistantMessageId("turn-reused-1", String(reusedItemId)), - ); - const secondMessage = thread.messages.find( - (message: ProviderRuntimeTestMessage) => - message.id === turnScopedAssistantMessageId("turn-reused-2", String(reusedItemId)), - ); - - expect(firstMessage?.text).toBe("first turn response"); - expect(secondMessage?.text).toBe("second turn response"); - }); - - it("ignores cursor assistant item replay without a turn id while a new turn is active", async () => { - const harness = await createHarness({ serverSettings: { enableAssistantStreaming: true } }); - const now = "2026-01-01T00:00:00.000Z"; - const replayedItemId = asItemId("assistant:cursor-session:segment:0"); - const secondItemId = asItemId("assistant:cursor-session:segment:1"); - - harness.emit({ - type: "turn.started", - eventId: asEventId("evt-cursor-replay-turn-1-started"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-cursor-replay-1"), - }); - harness.emit({ - type: "content.delta", - eventId: asEventId("evt-cursor-replay-turn-1-delta"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-cursor-replay-1"), - itemId: replayedItemId, - payload: { - streamKind: "assistant_text", - delta: "first turn response", - }, - }); - harness.emit({ - type: "item.completed", - eventId: asEventId("evt-cursor-replay-turn-1-complete"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-cursor-replay-1"), - itemId: replayedItemId, - payload: { - itemType: "assistant_message", - status: "completed", - }, - }); - harness.emit({ - type: "turn.completed", - eventId: asEventId("evt-cursor-replay-turn-1-turn-complete"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-cursor-replay-1"), - status: "completed", - }); - - await waitForThread( - harness.readModel, - (thread) => - thread.session?.status === "ready" && - thread.messages.some( - (message: ProviderRuntimeTestMessage) => message.text === "first turn response", - ), - ); - - harness.emit({ - type: "turn.started", - eventId: asEventId("evt-cursor-replay-turn-2-started"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-cursor-replay-2"), - }); - await waitForThread( - harness.readModel, - (thread) => - thread.session?.status === "running" && - thread.session.activeTurnId === "turn-cursor-replay-2", - ); - - // Cursor can replay the previous assistant segment after an ACP session - // resume without attaching a turn id, then report that same item completed - // under the active turn. Neither event belongs to the new turn. - harness.emit({ - type: "item.started", - eventId: asEventId("evt-cursor-replay-stale-started"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - itemId: replayedItemId, - payload: { - itemType: "assistant_message", - status: "inProgress", - }, - }); - harness.emit({ - type: "content.delta", - eventId: asEventId("evt-cursor-replay-stale-delta"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - itemId: replayedItemId, - payload: { - streamKind: "assistant_text", - delta: "first turn response", - }, - }); - harness.emit({ - type: "item.completed", - eventId: asEventId("evt-cursor-replay-stale-complete"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-cursor-replay-2"), - itemId: replayedItemId, - payload: { - itemType: "assistant_message", - status: "completed", - }, - }); - - harness.emit({ - type: "content.delta", - eventId: asEventId("evt-cursor-replay-turn-2-delta"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-cursor-replay-2"), - itemId: secondItemId, - payload: { - streamKind: "assistant_text", - delta: "second turn response", - }, - }); - harness.emit({ - type: "item.completed", - eventId: asEventId("evt-cursor-replay-turn-2-complete"), - provider: ProviderDriverKind.make("cursor"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-cursor-replay-2"), - itemId: secondItemId, - payload: { - itemType: "assistant_message", - status: "completed", - }, - }); - - const thread = await waitForThread(harness.readModel, (entry) => - entry.messages.some( - (message: ProviderRuntimeTestMessage) => - message.id === turnScopedAssistantMessageId("turn-cursor-replay-2", String(secondItemId)), - ), - ); - expect(thread.messages.map((message: ProviderRuntimeTestMessage) => message.text)).toEqual([ - "first turn response", - "second turn response", - ]); - expect( - thread.messages.some( - (message: ProviderRuntimeTestMessage) => message.id === `assistant:${replayedItemId}`, - ), - ).toBe(false); - expect( - thread.messages.some( - (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId("turn-cursor-replay-2", String(replayedItemId)), - ), - ).toBe(false); - }); - it("streams assistant deltas when thread.turn.start requests streaming mode", async () => { const harness = await createHarness({ serverSettings: { enableAssistantStreaming: true } }); const now = "2026-01-01T00:00:00.000Z"; @@ -2450,15 +2126,13 @@ describe("ProviderRuntimeIngestion", () => { const liveThread = await waitForThread(harness.readModel, (entry) => entry.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId("turn-streaming-mode", "item-streaming-mode") && + message.id === "assistant:item-streaming-mode" && message.streaming && message.text === "hello live", ), ); const liveMessage = liveThread.messages.find( - (entry: ProviderRuntimeTestMessage) => - entry.id === turnScopedAssistantMessageId("turn-streaming-mode", "item-streaming-mode"), + (entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-streaming-mode", ); expect(liveMessage?.streaming).toBe(true); @@ -2480,14 +2154,11 @@ describe("ProviderRuntimeIngestion", () => { const finalThread = await waitForThread(harness.readModel, (entry) => entry.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId("turn-streaming-mode", "item-streaming-mode") && - !message.streaming, + message.id === "assistant:item-streaming-mode" && !message.streaming, ), ); const finalMessage = finalThread.messages.find( - (entry: ProviderRuntimeTestMessage) => - entry.id === turnScopedAssistantMessageId("turn-streaming-mode", "item-streaming-mode"), + (entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-streaming-mode", ); expect(finalMessage?.text).toBe("hello live"); expect(finalMessage?.streaming).toBe(false); @@ -2543,13 +2214,11 @@ describe("ProviderRuntimeIngestion", () => { const thread = await waitForThread(harness.readModel, (entry) => entry.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === turnScopedAssistantMessageId("turn-buffer-spill", "item-buffer-spill") && - !message.streaming, + message.id === "assistant:item-buffer-spill" && !message.streaming, ), ); const message = thread.messages.find( - (entry: ProviderRuntimeTestMessage) => - entry.id === turnScopedAssistantMessageId("turn-buffer-spill", "item-buffer-spill"), + (entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-buffer-spill", ); expect(message?.text.length).toBe(oversizedText.length); expect(message?.text).toBe(oversizedText); @@ -2621,9 +2290,7 @@ describe("ProviderRuntimeIngestion", () => { thread.session?.activeTurnId === null && thread.messages.some( (message: ProviderRuntimeTestMessage) => - message.id === - turnScopedAssistantMessageId("turn-complete-dedup", "item-complete-dedup") && - !message.streaming, + message.id === "assistant:item-complete-dedup" && !message.streaming, ), ); @@ -2637,8 +2304,7 @@ describe("ProviderRuntimeIngestion", () => { return false; } return ( - event.payload.messageId === - turnScopedAssistantMessageId("turn-complete-dedup", "item-complete-dedup") && + event.payload.messageId === "assistant:item-complete-dedup" && event.payload.streaming === false ); }); @@ -2995,9 +2661,7 @@ describe("ProviderRuntimeIngestion", () => { (entry: ProviderRuntimeTestCheckpoint) => entry.turnId === "turn-p1", ); expect(checkpoint?.status).toBe("missing"); - expect(checkpoint?.assistantMessageId).toBe( - turnScopedAssistantMessageId("turn-p1", "item-p1-assistant"), - ); + expect(checkpoint?.assistantMessageId).toBe("assistant:item-p1-assistant"); expect(checkpoint?.checkpointRef).toBe("provider-diff:evt-turn-diff-updated"); }); diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 9ffd67a00fb..bd300ba2ad1 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -53,8 +53,6 @@ type GoalActivityState = Pick; const TURN_MESSAGE_IDS_BY_TURN_CACHE_CAPACITY = 10_000; const TURN_MESSAGE_IDS_BY_TURN_TTL = Duration.minutes(120); -const STALE_REPLAY_ITEM_IDS_BY_TURN_CACHE_CAPACITY = 10_000; -const STALE_REPLAY_ITEM_IDS_BY_TURN_TTL = Duration.minutes(120); const BUFFERED_MESSAGE_TEXT_BY_MESSAGE_ID_CACHE_CAPACITY = 20_000; const BUFFERED_MESSAGE_TEXT_BY_MESSAGE_ID_TTL = Duration.minutes(120); const BUFFERED_PROPOSED_PLAN_BY_ID_CACHE_CAPACITY = 10_000; @@ -243,9 +241,8 @@ function proposedPlanIdFromEvent(event: ProviderRuntimeEvent, threadId: ThreadId return `plan:${threadId}:event:${event.eventId}`; } -function assistantSegmentBaseKeyFromEvent(event: ProviderRuntimeEvent, turnId?: TurnId): string { - const providerKey = String(event.itemId ?? event.turnId ?? event.eventId); - return turnId ? `turn:${turnId}:${providerKey}` : providerKey; +function assistantSegmentBaseKeyFromEvent(event: ProviderRuntimeEvent): string { + return String(event.itemId ?? event.turnId ?? event.eventId); } function assistantSegmentMessageId(baseKey: string, segmentIndex: number): MessageId { @@ -253,41 +250,6 @@ function assistantSegmentMessageId(baseKey: string, segmentIndex: number): Messa segmentIndex === 0 ? `assistant:${baseKey}` : `assistant:${baseKey}:segment:${segmentIndex}`, ); } - -function runtimeAssistantMessageIdFromEvent( - event: ProviderRuntimeEvent, - turnId?: TurnId, -): MessageId { - return assistantSegmentMessageId(assistantSegmentBaseKeyFromEvent(event, turnId), 0); -} - -function isTurnOutputRuntimeEvent(event: ProviderRuntimeEvent): boolean { - switch (event.type) { - case "content.delta": - case "item.started": - case "item.updated": - case "item.completed": - case "request.opened": - case "request.resolved": - case "runtime.warning": - case "task.started": - case "task.progress": - case "task.completed": - case "thread.state.changed": - case "thread.token-usage.updated": - case "tool.denied": - case "turn.completed": - case "turn.diff.updated": - case "turn.plan.updated": - case "turn.proposed.completed": - case "turn.proposed.delta": - case "user-input.requested": - case "user-input.resolved": - return true; - default: - return false; - } -} function buildContextWindowActivityPayload( event: ProviderRuntimeEvent, ): ThreadTokenUsageSnapshot | undefined { @@ -791,12 +753,6 @@ const make = Effect.gen(function* () { ), }); - const staleReplayItemIdsByTurnKey = yield* Cache.make>({ - capacity: STALE_REPLAY_ITEM_IDS_BY_TURN_CACHE_CAPACITY, - timeToLive: STALE_REPLAY_ITEM_IDS_BY_TURN_TTL, - lookup: () => Effect.succeed(new Set()), - }); - const bufferedProposedPlanById = yield* Cache.make({ capacity: BUFFERED_PROPOSED_PLAN_BY_ID_CACHE_CAPACITY, timeToLive: BUFFERED_PROPOSED_PLAN_BY_ID_TTL, @@ -878,37 +834,6 @@ const make = Effect.gen(function* () { const clearAssistantSegmentStateForTurn = (threadId: ThreadId, turnId: TurnId) => Cache.invalidate(assistantSegmentStateByTurnKey, providerTurnKey(threadId, turnId)); - const rememberStaleReplayItemForTurn = (threadId: ThreadId, turnId: TurnId, itemId: string) => - Cache.getOption(staleReplayItemIdsByTurnKey, providerTurnKey(threadId, turnId)).pipe( - Effect.flatMap((existingIds) => - Cache.set( - staleReplayItemIdsByTurnKey, - providerTurnKey(threadId, turnId), - Option.match(existingIds, { - onNone: () => new Set([itemId]), - onSome: (ids) => { - const nextIds = new Set(ids); - nextIds.add(itemId); - return nextIds; - }, - }), - ), - ), - ); - - const hasStaleReplayItemForTurn = (threadId: ThreadId, turnId: TurnId, itemId: string) => - Cache.getOption(staleReplayItemIdsByTurnKey, providerTurnKey(threadId, turnId)).pipe( - Effect.map((existingIds) => - Option.match(existingIds, { - onNone: () => false, - onSome: (ids) => ids.has(itemId), - }), - ), - ); - - const clearStaleReplayItemsForTurn = (threadId: ThreadId, turnId: TurnId) => - Cache.invalidate(staleReplayItemIdsByTurnKey, providerTurnKey(threadId, turnId)); - const getActiveAssistantMessageIdForTurn = (threadId: ThreadId, turnId: TurnId) => getAssistantSegmentStateForTurn(threadId, turnId).pipe( Effect.map((state) => @@ -955,7 +880,7 @@ const make = Effect.gen(function* () { }) => Effect.gen(function* () { if (!input.turnId) { - return runtimeAssistantMessageIdFromEvent(input.event); + return assistantSegmentMessageId(assistantSegmentBaseKeyFromEvent(input.event), 0); } const activeMessageId = yield* getActiveAssistantMessageIdForTurn( @@ -969,7 +894,7 @@ const make = Effect.gen(function* () { return yield* startAssistantSegmentForTurn({ threadId: input.threadId, turnId: input.turnId, - baseKey: assistantSegmentBaseKeyFromEvent(input.event, input.turnId), + baseKey: assistantSegmentBaseKeyFromEvent(input.event), }); }); @@ -1263,7 +1188,6 @@ const make = Effect.gen(function* () { const proposedPlanPrefix = `plan:${threadId}:`; const turnKeys = Array.from(yield* Cache.keys(turnMessageIdsByTurnKey)); const assistantSegmentKeys = Array.from(yield* Cache.keys(assistantSegmentStateByTurnKey)); - const staleReplayKeys = Array.from(yield* Cache.keys(staleReplayItemIdsByTurnKey)); const proposedPlanKeys = Array.from(yield* Cache.keys(bufferedProposedPlanById)); yield* Effect.forEach( turnKeys, @@ -1292,12 +1216,6 @@ const make = Effect.gen(function* () { : Effect.void, { concurrency: 1 }, ).pipe(Effect.asVoid); - yield* Effect.forEach( - staleReplayKeys, - (key) => - key.startsWith(prefix) ? Cache.invalidate(staleReplayItemIdsByTurnKey, key) : Effect.void, - { concurrency: 1 }, - ).pipe(Effect.asVoid); yield* Effect.forEach( proposedPlanKeys, (key) => @@ -1539,23 +1457,6 @@ const make = Effect.gen(function* () { } } - const eventItemId = event.itemId === undefined ? undefined : String(event.itemId); - const staleReplayItemForActiveTurn = - activeTurnId !== null && eventItemId !== undefined - ? yield* hasStaleReplayItemForTurn(thread.id, activeTurnId, eventItemId) - : false; - const shouldSkipRuntimeOutput = - STRICT_PROVIDER_LIFECYCLE_GUARD && - isTurnOutputRuntimeEvent(event) && - (conflictsWithActiveTurn || missingTurnForActiveTurn || staleReplayItemForActiveTurn); - - if (shouldSkipRuntimeOutput) { - if (missingTurnForActiveTurn && activeTurnId !== null && eventItemId !== undefined) { - yield* rememberStaleReplayItemForTurn(thread.id, activeTurnId, eventItemId); - } - return; - } - const assistantDelta = event.type === "content.delta" && event.payload.streamKind === "assistant_text" ? event.payload.delta @@ -1657,7 +1558,9 @@ const make = Effect.gen(function* () { const assistantCompletion = event.type === "item.completed" && event.payload.itemType === "assistant_message" ? { - messageId: runtimeAssistantMessageIdFromEvent(event, toTurnId(event.turnId)), + messageId: MessageId.make( + `assistant:${event.itemId ?? event.turnId ?? event.eventId}`, + ), fallbackText: event.payload.detail, } : undefined; @@ -1759,7 +1662,6 @@ const make = Effect.gen(function* () { ).pipe(Effect.asVoid); yield* clearAssistantMessageIdsForTurn(thread.id, turnId); yield* clearAssistantSegmentStateForTurn(thread.id, turnId); - yield* clearStaleReplayItemsForTurn(thread.id, turnId); yield* finalizeBufferedProposedPlan({ event, @@ -1858,7 +1760,9 @@ const make = Effect.gen(function* () { if (hasCheckpointForTurn(checkpointContext.checkpoints, turnId)) { // Already tracked; no-op. } else { - const assistantMessageId = runtimeAssistantMessageIdFromEvent(event, turnId); + const assistantMessageId = MessageId.make( + `assistant:${event.itemId ?? event.turnId ?? event.eventId}`, + ); yield* orchestrationEngine.dispatch({ type: "thread.turn.diff.complete", commandId: yield* providerCommandId(event, "thread-turn-diff-complete"), diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.test.ts b/apps/server/src/provider/acp/AcpSessionRuntime.test.ts deleted file mode 100644 index 59dc3421983..00000000000 --- a/apps/server/src/provider/acp/AcpSessionRuntime.test.ts +++ /dev/null @@ -1,84 +0,0 @@ -import { assert, it } from "@effect/vitest"; -import * as Effect from "effect/Effect"; -import * as Queue from "effect/Queue"; -import * as Ref from "effect/Ref"; -import type * as EffectAcpSchema from "effect-acp/schema"; - -import { handleSessionUpdateForTest, type AcpSessionRuntimeEvent } from "./AcpSessionRuntime.ts"; -import type { AcpSessionModeState, AcpToolCallState } from "./AcpRuntimeModel.ts"; - -it.effect("suppresses loaded-session replay updates until the first live prompt", () => - Effect.gen(function* () { - const queue = yield* Queue.unbounded(); - const modeStateRef = yield* Ref.make({ - currentModeId: "ask", - availableModes: [ - { id: "ask", name: "Ask" }, - { id: "code", name: "Code" }, - ], - }); - const toolCallsRef = yield* Ref.make(new Map()); - const assistantSegmentRef = yield* Ref.make({ nextSegmentIndex: 0 }); - const suppressSessionUpdatesRef = yield* Ref.make(true); - - const handle = (params: EffectAcpSchema.SessionNotification) => - handleSessionUpdateForTest({ - queue, - modeStateRef, - toolCallsRef, - assistantSegmentRef, - suppressSessionUpdatesRef, - params, - }); - - yield* handle({ - sessionId: "cursor-session", - update: { - sessionUpdate: "current_mode_update", - currentModeId: "code", - }, - }); - yield* handle({ - sessionId: "cursor-session", - update: { - sessionUpdate: "plan", - entries: [{ content: "Old replayed plan", priority: "high", status: "completed" }], - }, - }); - yield* handle({ - sessionId: "cursor-session", - update: { - sessionUpdate: "user_message_chunk", - content: { type: "text", text: "old replayed user prompt" }, - }, - }); - yield* handle({ - sessionId: "cursor-session", - update: { - sessionUpdate: "agent_message_chunk", - content: { type: "text", text: "old replayed assistant text" }, - }, - }); - - assert.equal(yield* Queue.size(queue), 0); - assert.equal((yield* Ref.get(modeStateRef))?.currentModeId, "code"); - - yield* Ref.set(suppressSessionUpdatesRef, false); - yield* handle({ - sessionId: "cursor-session", - update: { - sessionUpdate: "agent_message_chunk", - content: { type: "text", text: "new assistant text" }, - }, - }); - - const started = yield* Queue.take(queue); - const delta = yield* Queue.take(queue); - assert.equal(started._tag, "AssistantItemStarted"); - assert.equal(delta._tag, "ContentDelta"); - if (delta._tag === "ContentDelta") { - assert.equal(delta.text, "new assistant text"); - } - assert.equal(yield* Queue.size(queue), 0); - }), -); diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index 0c108929c0b..ab1f19b09aa 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -70,7 +70,6 @@ export interface AcpSessionRuntimeOptions { }; readonly authMethodId: string; readonly mcpServers?: ReadonlyArray; - readonly suppressSessionUpdatesUntilPrompt?: boolean; readonly requestLogger?: (event: AcpSessionRequestLogEvent) => Effect.Effect; readonly protocolLogging?: { readonly logIncoming?: boolean; @@ -283,9 +282,6 @@ export const make = ( const assistantSegmentRef = yield* Ref.make({ nextSegmentIndex: 0 }); const configOptionsRef = yield* Ref.make(sessionConfigOptionsFromSetup(undefined)); const startStateRef = yield* Ref.make({ _tag: "NotStarted" }); - const suppressSessionUpdatesRef = yield* Ref.make( - options.suppressSessionUpdatesUntilPrompt === true, - ); const promptSerializationSemaphore = yield* Semaphore.make(1); const activePromptFiberRef = yield* Ref.make< Option.Option> @@ -394,7 +390,6 @@ export const make = ( modeStateRef, toolCallsRef, assistantSegmentRef, - suppressSessionUpdatesRef, params: notification, }); }), @@ -724,7 +719,6 @@ export const make = ( sessionId: started.sessionId, ...payload, } satisfies EffectAcpSchema.PromptRequest; - yield* Ref.set(suppressSessionUpdatesRef, false); const cancelledResponse = { stopReason: "cancelled", } satisfies EffectAcpSchema.PromptResponse; @@ -843,14 +837,12 @@ const handleSessionUpdate = ({ modeStateRef, toolCallsRef, assistantSegmentRef, - suppressSessionUpdatesRef, params, }: { readonly queue: Queue.Queue; readonly modeStateRef: Ref.Ref; readonly toolCallsRef: Ref.Ref>; readonly assistantSegmentRef: Ref.Ref; - readonly suppressSessionUpdatesRef: Ref.Ref; readonly params: EffectAcpSchema.SessionNotification; }): Effect.Effect => Effect.gen(function* () { @@ -860,9 +852,6 @@ const handleSessionUpdate = ({ current === undefined ? current : updateModeState(current, parsed.modeId!), ); } - if (yield* Ref.get(suppressSessionUpdatesRef)) { - return; - } for (const event of parsed.events) { if (event._tag === "ToolCallUpdated") { yield* closeActiveAssistantSegment({ @@ -912,8 +901,6 @@ const handleSessionUpdate = ({ } }); -export const handleSessionUpdateForTest = handleSessionUpdate; - function updateModeState(modeState: AcpSessionModeState, nextModeId: string): AcpSessionModeState { const normalized = nextModeId.trim(); if (!normalized) { diff --git a/apps/server/src/provider/acp/CursorAcpSupport.ts b/apps/server/src/provider/acp/CursorAcpSupport.ts index 07f7c41f6c5..169d7c6206d 100644 --- a/apps/server/src/provider/acp/CursorAcpSupport.ts +++ b/apps/server/src/provider/acp/CursorAcpSupport.ts @@ -59,8 +59,6 @@ export const makeCursorAcpRuntime = ( spawn: buildCursorAcpSpawnInput(input.cursorSettings, input.cwd, input.environment), authMethodId: "cursor_login", clientCapabilities: CURSOR_PARAMETERIZED_MODEL_PICKER_CAPABILITIES, - suppressSessionUpdatesUntilPrompt: - input.suppressSessionUpdatesUntilPrompt ?? input.resumeSessionId !== undefined, }).pipe( Layer.provide( Layer.succeed(ChildProcessSpawner.ChildProcessSpawner, input.childProcessSpawner), From 474a1a0c32c963495610b492c39e00f0c9a68996 Mon Sep 17 00:00:00 2001 From: taschaub Date: Thu, 2 Jul 2026 19:40:38 +0200 Subject: [PATCH 2/3] Fix Cursor ACP thread rendering, thought output, and cancel delivery - Send session/cancel as a spec-compliant JSON-RPC notification (no id); the Cursor CLI dropped the malformed message so turns kept running after pausing a thread. - Tag ACP assistant segment item ids with a per-runtime tag so resumed sessions stop appending new output to messages from earlier runs, which pushed assistant text above the latest user message. - Parse agent_thought_chunk into channel-aware segments and surface the accumulated reasoning as expandable thinking rows in the work log. - Drain queued ACP session updates before emitting turn.completed in CursorAdapter (raced against the notification fiber to avoid hanging on mid-turn teardown). Co-authored-by: Cursor --- apps/server/scripts/acp-mock-agent.ts | 29 +++++ .../Layers/ProviderRuntimeIngestion.ts | 35 +++++ .../src/provider/Layers/CursorAdapter.test.ts | 6 +- .../src/provider/Layers/CursorAdapter.ts | 18 +++ .../server/src/provider/Layers/GrokAdapter.ts | 4 + .../provider/acp/AcpCoreRuntimeEvents.test.ts | 44 +++++++ .../src/provider/acp/AcpCoreRuntimeEvents.ts | 18 ++- .../provider/acp/AcpJsonRpcConnection.test.ts | 58 +++++++++ .../src/provider/acp/AcpRuntimeModel.test.ts | 30 +++++ .../src/provider/acp/AcpRuntimeModel.ts | 26 ++++ .../src/provider/acp/AcpSessionRuntime.ts | 122 ++++++++++++++---- apps/web/src/session-logic.ts | 4 +- packages/effect-acp/src/protocol.test.ts | 10 +- packages/effect-acp/src/protocol.ts | 16 ++- 14 files changed, 390 insertions(+), 30 deletions(-) diff --git a/apps/server/scripts/acp-mock-agent.ts b/apps/server/scripts/acp-mock-agent.ts index bc7828dd854..1add889c6b9 100644 --- a/apps/server/scripts/acp-mock-agent.ts +++ b/apps/server/scripts/acp-mock-agent.ts @@ -16,6 +16,7 @@ const exitLogPath = process.env.T3_ACP_EXIT_LOG_PATH; const emitToolCalls = process.env.T3_ACP_EMIT_TOOL_CALLS === "1"; const emitInterleavedAssistantToolCalls = process.env.T3_ACP_EMIT_INTERLEAVED_ASSISTANT_TOOL_CALLS === "1"; +const emitThoughtChunks = process.env.T3_ACP_EMIT_THOUGHT_CHUNKS === "1"; const emitGenericToolPlaceholders = process.env.T3_ACP_EMIT_GENERIC_TOOL_PLACEHOLDERS === "1"; const emitAskQuestion = process.env.T3_ACP_EMIT_ASK_QUESTION === "1"; const emitXAiAskUserQuestion = process.env.T3_ACP_EMIT_XAI_ASK_USER_QUESTION === "1"; @@ -580,6 +581,34 @@ const program = Effect.gen(function* () { return yield* Effect.never; } + if (emitThoughtChunks) { + yield* agent.client.sessionUpdate({ + sessionId: requestedSessionId, + update: { + sessionUpdate: "agent_thought_chunk", + content: { type: "text", text: "thinking about " }, + }, + }); + + yield* agent.client.sessionUpdate({ + sessionId: requestedSessionId, + update: { + sessionUpdate: "agent_thought_chunk", + content: { type: "text", text: "the answer" }, + }, + }); + + yield* agent.client.sessionUpdate({ + sessionId: requestedSessionId, + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "final answer" }, + }, + }); + + return { stopReason: "end_turn" }; + } + if (emitInterleavedAssistantToolCalls) { const toolCallId = "tool-call-1"; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index bd300ba2ad1..9d111e03072 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -214,6 +214,19 @@ function goalUpdatedActivitySummary( } } +/** Full reasoning text is kept expandable in the work log but capped so a + * single thought segment cannot bloat the persisted activity payload. */ +const MAX_REASONING_ACTIVITY_CHARS = 8_000; + +function reasoningSummaryFromText(text: string): string { + const firstLine = text + .split("\n") + .find((line) => line.trim().length > 0) + ?.trim(); + const cleaned = firstLine?.replaceAll("**", "").trim(); + return truncateDetail(cleaned && cleaned.length > 0 ? cleaned : "Thinking", 120); +} + function normalizeProposedPlanMarkdown(planMarkdown: string | undefined): string | undefined { const trimmed = planMarkdown?.trim(); if (!trimmed) { @@ -671,6 +684,28 @@ function runtimeEventToActivities( } case "item.completed": { + // Reasoning segments (e.g. Cursor agent_thought_chunk) become expandable + // "thinking" rows in the work log instead of disappearing entirely. + if (event.payload.itemType === "reasoning") { + const reasoningText = event.payload.detail?.trim(); + if (!reasoningText) { + return []; + } + return [ + { + id: event.eventId, + createdAt: event.createdAt, + tone: "info", + kind: "reasoning", + summary: reasoningSummaryFromText(reasoningText), + payload: { + detail: truncateDetail(reasoningText, MAX_REASONING_ACTIVITY_CHARS), + }, + turnId: toTurnId(event.turnId) ?? null, + ...maybeSequence, + }, + ]; + } if (!isToolLifecycleItemType(event.payload.itemType)) { return []; } diff --git a/apps/server/src/provider/Layers/CursorAdapter.test.ts b/apps/server/src/provider/Layers/CursorAdapter.test.ts index 89e9c56eb8a..125f3ff41bd 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.test.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.test.ts @@ -228,7 +228,9 @@ cursorAdapterTestLayer("CursorAdapterLive", (it) => { assert.isDefined(delta); if (delta?.type === "content.delta") { assert.equal(delta.payload.delta, "hello from mock"); - assert.match(String(delta.itemId), /^assistant:mock-session-1:segment:0$/); + // The middle part is a per-runtime tag that keeps segment ids unique + // across restarts that resume the same ACP session. + assert.match(String(delta.itemId), /^assistant:mock-session-1:[^:]+:segment:0$/); } const assistantCompleted = runtimeEvents.find( @@ -687,7 +689,7 @@ cursorAdapterTestLayer("CursorAdapterLive", (it) => { if (contentDelta?.type === "content.delta") { assert.equal(String(contentDelta.turnId), String(turn.turnId)); assert.equal(contentDelta.payload.delta, "hello from mock"); - assert.equal(String(contentDelta.itemId), "assistant:mock-session-1:segment:0"); + assert.match(String(contentDelta.itemId), /^assistant:mock-session-1:[^:]+:segment:0$/); } }); diff --git a/apps/server/src/provider/Layers/CursorAdapter.ts b/apps/server/src/provider/Layers/CursorAdapter.ts index 319c7d8550b..8a65e67a0de 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.ts @@ -800,6 +800,7 @@ export function makeCursorAdapter( turnId: ctx.activeTurnId, itemId: event.itemId, lifecycle: "item.started", + channel: event.channel, }), ); return; @@ -812,6 +813,8 @@ export function makeCursorAdapter( turnId: ctx.activeTurnId, itemId: event.itemId, lifecycle: "item.completed", + channel: event.channel, + ...(event.text !== undefined ? { detail: event.text } : {}), }), ); return; @@ -862,6 +865,7 @@ export function makeCursorAdapter( threadId: ctx.threadId, turnId: ctx.activeTurnId, ...(event.itemId ? { itemId: event.itemId } : {}), + channel: event.channel, text: event.text, rawPayload: event.rawPayload, }), @@ -1014,6 +1018,20 @@ export function makeCursorAdapter( ), ); + // The prompt RPC can resolve while session/update notifications are + // still queued. Wait until the notification fiber has published all + // of them so content deltas and item completions reach consumers + // before turn.completed (otherwise buffered assistant text may be + // finalized against the wrong turn state). Racing against the + // notification fiber keeps this from hanging when the session is + // torn down mid-turn and nobody can acknowledge the drain barrier. + yield* ctx.notificationFiber + ? Effect.raceFirst( + ctx.acp.drainEvents, + Fiber.await(ctx.notificationFiber).pipe(Effect.asVoid), + ) + : Effect.void; + const turnRecord = ctx.turns.find((turn) => turn.id === turnId); if (turnRecord) { turnRecord.items.push({ prompt: promptParts, result }); diff --git a/apps/server/src/provider/Layers/GrokAdapter.ts b/apps/server/src/provider/Layers/GrokAdapter.ts index 8f1870eb242..7ca54267a0a 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.ts @@ -819,6 +819,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte turnId: notificationTurnId, itemId: event.itemId, lifecycle: "item.started", + channel: event.channel, }), ); return; @@ -831,6 +832,8 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte turnId: notificationTurnId, itemId: event.itemId, lifecycle: "item.completed", + channel: event.channel, + ...(event.text !== undefined ? { detail: event.text } : {}), }), ); return; @@ -864,6 +867,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte threadId: ctx.threadId, turnId: notificationTurnId, ...(event.itemId ? { itemId: event.itemId } : {}), + channel: event.channel, text: event.text, rawPayload: event.rawPayload, }), diff --git a/apps/server/src/provider/acp/AcpCoreRuntimeEvents.test.ts b/apps/server/src/provider/acp/AcpCoreRuntimeEvents.test.ts index 7fe25699bbc..7c7d6219eda 100644 --- a/apps/server/src/provider/acp/AcpCoreRuntimeEvents.test.ts +++ b/apps/server/src/provider/acp/AcpCoreRuntimeEvents.test.ts @@ -152,4 +152,48 @@ describe("AcpCoreRuntimeEvents", () => { }, }); }); + + it("maps thought-channel segments to reasoning items and reasoning_text deltas", () => { + const stamp = { eventId: "event-1" as never, createdAt: "2026-03-27T00:00:00.000Z" }; + const turnId = TurnId.make("turn-1"); + + expect( + makeAcpContentDeltaEvent({ + stamp, + provider: ProviderDriverKind.make("cursor"), + threadId: "thread-1" as never, + turnId, + itemId: "thought:session-1:tag:segment:0", + channel: "thought", + text: "Checking the failing test first.", + rawPayload: { sessionId: "session-1" }, + }), + ).toMatchObject({ + type: "content.delta", + payload: { + streamKind: "reasoning_text", + delta: "Checking the failing test first.", + }, + }); + + expect( + makeAcpAssistantItemEvent({ + stamp, + provider: ProviderDriverKind.make("cursor"), + threadId: "thread-1" as never, + turnId, + itemId: "thought:session-1:tag:segment:0", + lifecycle: "item.completed", + channel: "thought", + detail: "Checking the failing test first.", + }), + ).toMatchObject({ + type: "item.completed", + payload: { + itemType: "reasoning", + status: "completed", + detail: "Checking the failing test first.", + }, + }); + }); }); diff --git a/apps/server/src/provider/acp/AcpCoreRuntimeEvents.ts b/apps/server/src/provider/acp/AcpCoreRuntimeEvents.ts index c93e61dc37b..14a90db85ab 100644 --- a/apps/server/src/provider/acp/AcpCoreRuntimeEvents.ts +++ b/apps/server/src/provider/acp/AcpCoreRuntimeEvents.ts @@ -12,7 +12,12 @@ import { type TurnId, } from "@t3tools/contracts"; -import type { AcpPermissionRequest, AcpPlanUpdate, AcpToolCallState } from "./AcpRuntimeModel.ts"; +import type { + AcpAssistantChannel, + AcpPermissionRequest, + AcpPlanUpdate, + AcpToolCallState, +} from "./AcpRuntimeModel.ts"; type AcpAdapterRawSource = Extract< RuntimeEventRawSource, @@ -198,6 +203,9 @@ export function makeAcpAssistantItemEvent(input: { readonly turnId: TurnId | undefined; readonly itemId: string; readonly lifecycle: "item.started" | "item.completed"; + readonly channel?: AcpAssistantChannel; + /** Full segment text for completed thought segments. */ + readonly detail?: string; }): ProviderRuntimeEvent { return { type: input.lifecycle, @@ -207,8 +215,11 @@ export function makeAcpAssistantItemEvent(input: { turnId: input.turnId, itemId: RuntimeItemId.make(input.itemId), payload: { - itemType: "assistant_message", + itemType: input.channel === "thought" ? "reasoning" : "assistant_message", status: input.lifecycle === "item.completed" ? "completed" : "inProgress", + ...(input.detail !== undefined && input.detail.trim().length > 0 + ? { detail: input.detail } + : {}), }, }; } @@ -219,6 +230,7 @@ export function makeAcpContentDeltaEvent(input: { readonly threadId: ThreadId; readonly turnId: TurnId | undefined; readonly itemId?: string; + readonly channel?: AcpAssistantChannel; readonly text: string; readonly rawPayload: unknown; }): ProviderRuntimeEvent { @@ -230,7 +242,7 @@ export function makeAcpContentDeltaEvent(input: { turnId: input.turnId, ...(input.itemId ? { itemId: RuntimeItemId.make(input.itemId) } : {}), payload: { - streamKind: "assistant_text", + streamKind: input.channel === "thought" ? "reasoning_text" : "assistant_text", delta: input.text, }, raw: { diff --git a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts index 4e9700dab7d..db5dd961494 100644 --- a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts +++ b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts @@ -293,6 +293,64 @@ describe("AcpSessionRuntime", () => { ), ); + it.effect("segments thought chunks separately from assistant text", () => + Effect.gen(function* () { + const runtime = yield* AcpSessionRuntime.AcpSessionRuntime; + yield* runtime.start(); + + const promptResult = yield* runtime.prompt({ + prompt: [{ type: "text", text: "hi" }], + }); + expect(promptResult).toMatchObject({ stopReason: "end_turn" }); + + const notes = Array.from(yield* Stream.runCollect(Stream.take(runtime.getEvents(), 7))); + expect(notes.map((note) => note._tag)).toEqual([ + "AssistantItemStarted", + "ContentDelta", + "ContentDelta", + "AssistantItemCompleted", + "AssistantItemStarted", + "ContentDelta", + "AssistantItemCompleted", + ]); + + const thoughtStarted = notes[0]; + const thoughtCompleted = notes[3]; + const assistantStarted = notes[4]; + if ( + thoughtStarted?._tag === "AssistantItemStarted" && + thoughtCompleted?._tag === "AssistantItemCompleted" && + assistantStarted?._tag === "AssistantItemStarted" + ) { + expect(thoughtStarted.channel).toBe("thought"); + expect(thoughtStarted.itemId).toMatch(/^thought:/); + expect(thoughtCompleted.itemId).toBe(thoughtStarted.itemId); + // The completed thought segment carries the accumulated text. + expect(thoughtCompleted.text).toBe("thinking about the answer"); + expect(assistantStarted.channel).toBe("assistant"); + expect(assistantStarted.itemId).toMatch(/^assistant:/); + expect(assistantStarted.itemId).not.toBe(thoughtStarted.itemId); + } + }).pipe( + Effect.provide( + AcpSessionRuntime.layer({ + spawn: { + command: mockAgentCommand, + args: mockAgentArgs, + env: { + T3_ACP_EMIT_THOUGHT_CHUNKS: "1", + }, + }, + cwd: process.cwd(), + clientInfo: { name: "t3-test", version: "0.0.0" }, + authMethodId: "test", + }), + ), + Effect.scoped, + Effect.provide(NodeServices.layer), + ), + ); + it.effect("suppresses generic placeholder tool updates until completion", () => Effect.gen(function* () { const runtime = yield* AcpSessionRuntime.AcpSessionRuntime; diff --git a/apps/server/src/provider/acp/AcpRuntimeModel.test.ts b/apps/server/src/provider/acp/AcpRuntimeModel.test.ts index 7682c5f5f9c..f358cf6f8cc 100644 --- a/apps/server/src/provider/acp/AcpRuntimeModel.test.ts +++ b/apps/server/src/provider/acp/AcpRuntimeModel.test.ts @@ -321,6 +321,7 @@ describe("AcpRuntimeModel", () => { expect(contentResult.events).toEqual([ { _tag: "ContentDelta", + channel: "assistant", text: "hello from acp", rawPayload: { sessionId: "session-1", @@ -334,6 +335,35 @@ describe("AcpRuntimeModel", () => { }, }, ]); + + const thoughtResult = parseSessionUpdateEvent({ + sessionId: "session-1", + update: { + sessionUpdate: "agent_thought_chunk", + content: { + type: "text", + text: "thinking...", + }, + }, + } satisfies EffectAcpSchema.SessionNotification); + + expect(thoughtResult.events).toEqual([ + { + _tag: "ContentDelta", + channel: "thought", + text: "thinking...", + rawPayload: { + sessionId: "session-1", + update: { + sessionUpdate: "agent_thought_chunk", + content: { + type: "text", + text: "thinking...", + }, + }, + }, + }, + ]); }); it("keeps permission request parsing compatible with loose extension payloads", () => { diff --git a/apps/server/src/provider/acp/AcpRuntimeModel.ts b/apps/server/src/provider/acp/AcpRuntimeModel.ts index e6bfc127e6e..6f25b651e15 100644 --- a/apps/server/src/provider/acp/AcpRuntimeModel.ts +++ b/apps/server/src/provider/acp/AcpRuntimeModel.ts @@ -80,6 +80,12 @@ export interface AcpPermissionRequest { readonly toolCall?: AcpToolCallState; } +/** + * Distinguishes the two ACP text streams: `agent_message_chunk` (final + * assistant prose) and `agent_thought_chunk` (intermediate reasoning). + */ +export type AcpAssistantChannel = "assistant" | "thought"; + export type AcpParsedSessionEvent = | { readonly _tag: "ModeChanged"; @@ -88,10 +94,15 @@ export type AcpParsedSessionEvent = | { readonly _tag: "AssistantItemStarted"; readonly itemId: string; + readonly channel: AcpAssistantChannel; } | { readonly _tag: "AssistantItemCompleted"; readonly itemId: string; + readonly channel: AcpAssistantChannel; + /** Accumulated segment text; populated for thought segments so the + * completed event can carry the full reasoning text downstream. */ + readonly text?: string; } | { readonly _tag: "PlanUpdated"; @@ -106,6 +117,7 @@ export type AcpParsedSessionEvent = | { readonly _tag: "ContentDelta"; readonly itemId?: string; + readonly channel: AcpAssistantChannel; readonly text: string; readonly rawPayload: unknown; }; @@ -568,6 +580,20 @@ export function parseSessionUpdateEvent(params: EffectAcpSchema.SessionNotificat if (upd.content.type === "text" && upd.content.text.length > 0) { events.push({ _tag: "ContentDelta", + channel: "assistant", + text: upd.content.text, + rawPayload: params, + }); + } + break; + } + case "agent_thought_chunk": { + // Reasoning text; agents like the Cursor CLI stream their intermediate + // narration here, so dropping it hides most of the turn's text output. + if (upd.content.type === "text" && upd.content.text.length > 0) { + events.push({ + _tag: "ContentDelta", + channel: "thought", text: upd.content.text, rawPayload: params, }); diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index ab1f19b09aa..dd841a4c979 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -31,6 +31,7 @@ import { sessionUpdateIsReplay, waitForSessionLoadReplayIdle, type SessionLoadGate, + type AcpAssistantChannel, type AcpParsedSessionEvent, type AcpSessionModeState, type AcpToolCallState, @@ -256,16 +257,36 @@ type AcpStartState = } | { readonly _tag: "Started"; readonly result: AcpStartedState }; +interface AcpActiveAssistantSegment { + readonly itemId: string; + readonly channel: AcpAssistantChannel; + /** Accumulated thought text so the segment-completed event can carry it. + * Assistant text is not accumulated here — it streams via ContentDelta. */ + readonly text: string; +} + interface AcpAssistantSegmentState { readonly nextSegmentIndex: number; - readonly activeItemId?: string; + readonly active?: AcpActiveAssistantSegment; } interface EnsureActiveAssistantSegmentResult { readonly itemId: string; - readonly startedEvent?: Extract; + readonly events: ReadonlyArray< + Extract< + AcpParsedSessionEvent, + { readonly _tag: "AssistantItemStarted" | "AssistantItemCompleted" } + > + >; } +/** Keeps thought accumulation bounded for very long reasoning segments. */ +const MAX_THOUGHT_SEGMENT_CHARS = 20_000; + +/** Differentiates runtime instances within one process; combined with the + * startup timestamp it makes segment item ids unique across restarts. */ +let runtimeInstanceCounter = 0; + export const make = ( options: AcpSessionRuntimeOptions, ): Effect.Effect< @@ -280,6 +301,12 @@ export const make = ( const modeStateRef = yield* Ref.make(undefined); const toolCallsRef = yield* Ref.make(new Map()); const assistantSegmentRef = yield* Ref.make({ nextSegmentIndex: 0 }); + // Segment item ids must not repeat across runtimes that resume the same + // ACP session (e.g. after a server restart), otherwise downstream + // consumers derive colliding message ids and new output gets appended to + // messages from a previous run. The tag makes ids runtime-unique. + runtimeInstanceCounter += 1; + const runtimeTag = `${(yield* Clock.currentTimeMillis).toString(36)}-${runtimeInstanceCounter.toString(36)}`; const configOptionsRef = yield* Ref.make(sessionConfigOptionsFromSetup(undefined)); const startStateRef = yield* Ref.make({ _tag: "NotStarted" }); const promptSerializationSemaphore = yield* Semaphore.make(1); @@ -390,6 +417,7 @@ export const make = ( modeStateRef, toolCallsRef, assistantSegmentRef, + runtimeTag, params: notification, }); }), @@ -837,12 +865,14 @@ const handleSessionUpdate = ({ modeStateRef, toolCallsRef, assistantSegmentRef, + runtimeTag, params, }: { readonly queue: Queue.Queue; readonly modeStateRef: Ref.Ref; readonly toolCallsRef: Ref.Ref>; readonly assistantSegmentRef: Ref.Ref; + readonly runtimeTag: string; readonly params: EffectAcpSchema.SessionNotification; }): Effect.Effect => Effect.gen(function* () { @@ -882,7 +912,8 @@ const handleSessionUpdate = ({ if (event._tag === "ContentDelta") { if (event.text.trim().length === 0) { const assistantSegmentState = yield* Ref.get(assistantSegmentRef); - if (!assistantSegmentState.activeItemId) { + // Whitespace-only deltas may not open a segment on their own. + if (assistantSegmentState.active?.channel !== event.channel) { continue; } } @@ -890,7 +921,22 @@ const handleSessionUpdate = ({ queue, assistantSegmentRef, sessionId: params.sessionId, + runtimeTag, + channel: event.channel, }); + if (event.channel === "thought") { + yield* Ref.update(assistantSegmentRef, (current) => + current.active?.itemId === itemId + ? { + ...current, + active: { + ...current.active, + text: appendThoughtSegmentText(current.active.text, event.text), + }, + } + : current, + ); + } yield* Queue.offer(queue, { ...event, itemId, @@ -927,44 +973,79 @@ function shouldEmitToolCallUpdate( return previous === undefined || previous.title !== next.title || previous.detail !== next.detail; } -const assistantItemId = (sessionId: string, segmentIndex: number) => - `assistant:${sessionId}:segment:${segmentIndex}`; +const assistantItemId = (input: { + readonly sessionId: string; + readonly runtimeTag: string; + readonly channel: AcpAssistantChannel; + readonly segmentIndex: number; +}) => + `${input.channel === "thought" ? "thought" : "assistant"}:${input.sessionId}:${input.runtimeTag}:segment:${input.segmentIndex}`; + +function appendThoughtSegmentText(current: string, delta: string): string { + if (current.length >= MAX_THOUGHT_SEGMENT_CHARS) { + return current; + } + return `${current}${delta}`.slice(0, MAX_THOUGHT_SEGMENT_CHARS); +} + +function completedSegmentEvent( + segment: AcpActiveAssistantSegment, +): Extract { + return { + _tag: "AssistantItemCompleted", + itemId: segment.itemId, + channel: segment.channel, + ...(segment.channel === "thought" && segment.text.length > 0 ? { text: segment.text } : {}), + }; +} const ensureActiveAssistantSegment = ({ queue, assistantSegmentRef, sessionId, + runtimeTag, + channel, }: { readonly queue: Queue.Queue; readonly assistantSegmentRef: Ref.Ref; readonly sessionId: string; + readonly runtimeTag: string; + readonly channel: AcpAssistantChannel; }) => Ref.modify( assistantSegmentRef, (current) => { - if (current.activeItemId) { - return [{ itemId: current.activeItemId }, current] as const; + if (current.active?.channel === channel) { + return [{ itemId: current.active.itemId, events: [] }, current] as const; } - const itemId = assistantItemId(sessionId, current.nextSegmentIndex); - return [ + const itemId = assistantItemId({ + sessionId, + runtimeTag, + channel, + segmentIndex: current.nextSegmentIndex, + }); + const events: EnsureActiveAssistantSegmentResult["events"] = [ + // A channel switch (assistant <-> thought) closes the previous segment. + ...(current.active ? [completedSegmentEvent(current.active)] : []), { + _tag: "AssistantItemStarted", itemId, - startedEvent: { - _tag: "AssistantItemStarted", - itemId, - } satisfies Extract, + channel, }, + ]; + return [ + { itemId, events }, { nextSegmentIndex: current.nextSegmentIndex + 1, - activeItemId: itemId, + active: { itemId, channel, text: "" }, } satisfies AcpAssistantSegmentState, ] as const; }, ).pipe( Effect.flatMap((result) => - result.startedEvent - ? Queue.offer(queue, result.startedEvent).pipe(Effect.as(result.itemId)) - : Effect.succeed(result.itemId), + Effect.forEach(result.events, (event) => Queue.offer(queue, event), { + discard: true, + }).pipe(Effect.as(result.itemId)), ), ); @@ -976,14 +1057,11 @@ const closeActiveAssistantSegment = ({ readonly assistantSegmentRef: Ref.Ref; }) => Ref.modify(assistantSegmentRef, (current) => { - if (!current.activeItemId) { + if (!current.active) { return [undefined, current] as const; } return [ - { - _tag: "AssistantItemCompleted", - itemId: current.activeItemId, - } satisfies AcpParsedSessionEvent, + completedSegmentEvent(current.active), { nextSegmentIndex: current.nextSegmentIndex, } satisfies AcpAssistantSegmentState, diff --git a/apps/web/src/session-logic.ts b/apps/web/src/session-logic.ts index 79ffa18ad7e..8f8074a3cc0 100644 --- a/apps/web/src/session-logic.ts +++ b/apps/web/src/session-logic.ts @@ -721,7 +721,9 @@ function toDerivedWorkLogEntry(activity: OrchestrationThreadActivity): DerivedWo turnId: activity.turnId, label: taskLabel || activity.summary, tone: - activity.kind === "task.progress" + // Reasoning rows (ACP agent_thought_chunk segments) share the thinking + // affordance with Codex task.progress reasoning updates. + activity.kind === "task.progress" || activity.kind === "reasoning" ? "thinking" : activity.tone === "approval" ? "info" diff --git a/packages/effect-acp/src/protocol.test.ts b/packages/effect-acp/src/protocol.test.ts index ece068dfc88..7f3cec01487 100644 --- a/packages/effect-acp/src/protocol.test.ts +++ b/packages/effect-acp/src/protocol.test.ts @@ -95,6 +95,14 @@ it.layer(NodeServices.layer)("effect-acp protocol", (it) => { sessionId: "session-1", }, }); + // JSON-RPC 2.0 notifications must omit `id` entirely; agents treat a + // present-but-empty id as a malformed request and drop the message. + const outboundText = + typeof outbound === "string" ? outbound : new TextDecoder().decode(outbound); + assert.equal( + outboundText, + '{"jsonrpc":"2.0","method":"session/cancel","params":{"sessionId":"session-1"}}\n', + ); yield* Queue.offer( input, @@ -211,7 +219,7 @@ it.layer(NodeServices.layer)("effect-acp protocol", (it) => { direction: "outgoing", stage: "raw", payload: - '{"jsonrpc":"2.0","method":"session/cancel","params":{"sessionId":"session-1"},"id":"","headers":[]}\n', + '{"jsonrpc":"2.0","method":"session/cancel","params":{"sessionId":"session-1"}}\n', }, ]); }), diff --git a/packages/effect-acp/src/protocol.ts b/packages/effect-acp/src/protocol.ts index 27c619296c0..9e6ebaa468d 100644 --- a/packages/effect-acp/src/protocol.ts +++ b/packages/effect-acp/src/protocol.ts @@ -120,8 +120,13 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi ? message.requestId : undefined; const requestId = encodedRequestId === "" ? undefined : encodedRequestId; + // JSON-RPC 2.0 notifications must omit the `id` member entirely. The RPC + // serializer would emit `"id": ""` for empty-id requests, which agents + // (e.g. the Cursor CLI) reject as a malformed request instead of handling + // the notification — so notifications are encoded by hand here. + const isNotification = message._tag === "Request" && message.id === ""; const encoded = yield* Effect.try({ - try: () => parser.encode(message), + try: () => (isNotification ? encodeJsonRpcNotificationWire(message) : parser.encode(message)), catch: (cause) => AcpError.AcpProtocolParseError.fromEncodingError(method, requestId, cause), }); @@ -558,6 +563,15 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi } satisfies AcpPatchedProtocol; }); +/** Encodes a JSON-RPC 2.0 notification (no `id`) with ndjson framing. */ +function encodeJsonRpcNotificationWire(message: RpcMessage.RequestEncoded): string { + return `${JSON.stringify({ + jsonrpc: "2.0", + method: message.tag, + ...(message.payload === undefined ? {} : { params: message.payload }), + })}\n`; +} + function isProtocolError( value: unknown, ): value is { code: number; message: string; data?: unknown } { From a9f731c0eb1cf3510344b78f34493bf7854f2762 Mon Sep 17 00:00:00 2001 From: taschaub Date: Fri, 3 Jul 2026 10:43:26 +0200 Subject: [PATCH 3/3] Address review: bail out of sendTurn after teardown, document reasoning deltas - sendTurn's drain race also settles when stopSessionInternal interrupts the notification fiber. Check ctx.stopped after the race and bail out before mutating turn state or emitting turn.completed, so a torn-down session can no longer produce events after session.exited. Covered by extending the stop-during-pending-approval test. - Document why reasoning_text content deltas are intentionally not appended to the assistant message by ingestion: the full thought text is accumulated in the runtime segment and persisted via item.completed, and segments close on prompt settlement so cancelled turns keep their reasoning. Co-authored-by: Cursor --- .../src/provider/Layers/CursorAdapter.test.ts | 22 ++++++++++++++++++- .../src/provider/Layers/CursorAdapter.ts | 14 ++++++++++++ .../src/provider/acp/AcpCoreRuntimeEvents.ts | 8 +++++++ 3 files changed, 43 insertions(+), 1 deletion(-) diff --git a/apps/server/src/provider/Layers/CursorAdapter.test.ts b/apps/server/src/provider/Layers/CursorAdapter.test.ts index 125f3ff41bd..f6510a53b76 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.test.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.test.ts @@ -1046,6 +1046,8 @@ cursorAdapterTestLayer("CursorAdapterLive", (it) => { const serverSettings = yield* ServerSettingsService; const threadId = ThreadId.make("cursor-stop-pending-approval"); const approvalRequested = yield* Deferred.make(); + const sessionExited = yield* Deferred.make(); + const observedEvents: Array = []; const wrapperPath = yield* Effect.promise(() => makeMockAgentWrapper({ T3_ACP_EMIT_TOOL_CALLS: "1" }), @@ -1053,7 +1055,14 @@ cursorAdapterTestLayer("CursorAdapterLive", (it) => { yield* serverSettings.updateSettings({ providers: { cursor: { binaryPath: wrapperPath } } }); yield* Stream.runForEach(adapter.streamEvents, (event) => { - if (String(event.threadId) !== String(threadId) || event.type !== "request.opened") { + if (String(event.threadId) !== String(threadId)) { + return Effect.void; + } + observedEvents.push(event); + if (event.type === "session.exited") { + return Deferred.succeed(sessionExited, undefined).pipe(Effect.ignore); + } + if (event.type !== "request.opened") { return Effect.void; } return Deferred.succeed(approvalRequested, undefined).pipe(Effect.ignore); @@ -1080,6 +1089,17 @@ cursorAdapterTestLayer("CursorAdapterLive", (it) => { yield* Fiber.await(sendTurnFiber); assert.equal(yield* adapter.hasSession(threadId), false); + + // Teardown interrupts the notification fiber, which settles sendTurn's + // drain race. sendTurn must then bail out: emitting turn.completed after + // session.exited would publish out-of-order events for a removed session. + // sendTurn already finished above, so after session.exited arrives we + // only need to let the stream consumer drain any remaining events. + yield* Deferred.await(sessionExited); + for (let i = 0; i < 10; i += 1) { + yield* Effect.yieldNow; + } + assert.isFalse(observedEvents.some((event) => event.type === "turn.completed")); }), ); diff --git a/apps/server/src/provider/Layers/CursorAdapter.ts b/apps/server/src/provider/Layers/CursorAdapter.ts index 8a65e67a0de..ca0e4b7f10f 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.ts @@ -1032,6 +1032,20 @@ export function makeCursorAdapter( ) : Effect.void; + // The notification fiber also terminates when stopSessionInternal + // interrupts it during teardown, which settles the race without the + // drain having happened. At that point the session is already gone + // (session.exited emitted, ctx removed from the map), so mutating + // turn state or emitting turn.completed here would produce + // out-of-order events against a removed session. Bail out instead. + if (ctx.stopped) { + return { + threadId: input.threadId, + turnId, + resumeCursor: ctx.session.resumeCursor, + }; + } + const turnRecord = ctx.turns.find((turn) => turn.id === turnId); if (turnRecord) { turnRecord.items.push({ prompt: promptParts, result }); diff --git a/apps/server/src/provider/acp/AcpCoreRuntimeEvents.ts b/apps/server/src/provider/acp/AcpCoreRuntimeEvents.ts index 14a90db85ab..20f0e68ff60 100644 --- a/apps/server/src/provider/acp/AcpCoreRuntimeEvents.ts +++ b/apps/server/src/provider/acp/AcpCoreRuntimeEvents.ts @@ -242,6 +242,14 @@ export function makeAcpContentDeltaEvent(input: { turnId: input.turnId, ...(input.itemId ? { itemId: RuntimeItemId.make(input.itemId) } : {}), payload: { + // reasoning_text deltas are intentionally not appended to the assistant + // message by ProviderRuntimeIngestion (same as the Codex/Claude/OpenCode + // adapters' reasoning deltas). The full thought text is accumulated in + // AcpSessionRuntime's active segment and delivered via item.completed + // (itemType "reasoning"), which ingestion persists as an expandable + // thinking row. Segments are closed on channel switches, tool calls, and + // prompt settlement (including cancellation), so accumulated reasoning + // survives an interrupted turn. streamKind: input.channel === "thought" ? "reasoning_text" : "assistant_text", delta: input.text, },