From 2319e0686bce8a5f0a6821e4664446505f8d930f Mon Sep 17 00:00:00 2001 From: Wout Stiens <71498452+StiensWout@users.noreply.github.com> Date: Sun, 4 Oct 2026 15:02:57 +0200 Subject: [PATCH 1/2] fix(server): delegated-task notices no longer cancel queued Claude tool calls When an async delegated task finished, T3 steered the completion notice into the running Claude turn with SDK priority "now". Claude ends the turn to deliver a "now" message, so tool calls it had issued but not started came back as "The user doesn't want to take this action right now. STOP". The agent read that as an instruction from the user and stopped. The Claude adapter now sends a server-created delegated completion with priority "next". Claude reads it at the turn's next tool boundary and every issued call still runs. A notice that lands during the final reply gets a native turn of its own afterwards, which reaches the thread through the existing continuation run. User steers, runtime-question answers and agent-to-agent steers keep "now". The steer input now carries the persisted message's delegatedCompletion so the adapter can tell a notice from other steers. Only a "now" steer marks the turn as expecting an abort result, so an abort after a notice still ends the turn. Fixes #15351 Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Adapters/ClaudeAdapterV2.test.ts | 293 +++++++++++++++--- .../Adapters/ClaudeAdapterV2.ts | 32 +- .../src/orchestration-v2/ProviderAdapter.ts | 1 + .../ProviderTurnControlService.ts | 3 + .../SteeringCompletion.integration.test.ts | 8 + 5 files changed, 286 insertions(+), 51 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts index 2e237bf3bb4c..0f8752cde989 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts @@ -2300,43 +2300,125 @@ describe("ClaudeAdapterV2 background wake turns", () => { }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), ); + // What a steered message looks like to the adapter, by who sent it. + const userSteer = { createdBy: "user", creationSource: "web" } as const; + const completionNotice = { + createdBy: "agent", + creationSource: "server", + delegatedCompletion: { + parentRunId: RunId.make("run-completion-parent"), + generation: 1, + taskIds: [NodeId.make("task-completion-child")], + }, + } as const; + const steerActiveTurn = Effect.fnUntraced(function* (input: { + readonly harness: Effect.Success; + readonly turn: ProviderAdapterV2TurnInput; + readonly message: Pick< + ProviderAdapterV2TurnInput["message"], + "createdBy" | "creationSource" | "delegatedCompletion" + >; + }) { + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const steerOrdinal = input.harness.offeredMessages.length; + yield* input.harness.runtime.steerTurn({ + threadId: input.harness.threadId, + runId: input.turn.runId, + providerThread: input.harness.providerThread, + providerTurnId: idAllocator.derive.providerTurn({ + driver: ClaudeAdapterV2.CLAUDE_PROVIDER, + nativeTurnId: `turn:${input.turn.attemptId}`, + }), + message: { + ...input.message, + messageId: MessageId.make(`message-steer-${steerOrdinal}`), + text: "Check the delegated work.", + attachments: [], + }, + }); + return input.harness.offeredMessages[steerOrdinal]?.priority; + }); + + it.effect("steers only a server-created delegated completion with next priority", () => + Effect.scoped( + Effect.gen(function* () { + const harness = yield* makeWakeHarness; + const turn = makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("attempt-steer-priority"), + text: "Audit the settings pages.", + attachments: [], + }); + yield* harness.runtime.startTurn(turn); + const { delegatedCompletion } = completionNotice; + const steers = [ + { message: completionNotice, priority: "next" }, + { message: userSteer, priority: "now" }, + // A runtime-question answer. + { message: { createdBy: "user", creationSource: "server" }, priority: "now" }, + // One agent steering another through MCP. + { message: { createdBy: "agent", creationSource: "mcp" }, priority: "now" }, + { message: { createdBy: "agent", creationSource: "server" }, priority: "now" }, + { message: { ...userSteer, delegatedCompletion }, priority: "now" }, + { + message: { createdBy: "agent", creationSource: "mcp", delegatedCompletion }, + priority: "now", + }, + ] as const; + for (const steer of steers) { + assert.equal( + yield* steerActiveTurn({ harness, turn, message: steer.message }), + steer.priority, + `${steer.message.createdBy}/${steer.message.creationSource}`, + ); + } + }).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ), + ); + + // Claude aborts a turn to deliver a `now` steer, so that abort is expected + // and the turn carries on. A `next` notice never aborts: an abort after one + // is real, and it leaves an abort a user steer caused expected. it.effect.each( (["aborted_tools", "aborted_streaming"] as const).flatMap((terminalReason) => - [true, false].map((steered) => ({ terminalReason, steered })), + ( + [ + { steers: [], continues: false }, + { steers: ["user"], continues: true }, + { steers: ["notice"], continues: false }, + { steers: ["user", "notice"], continues: true }, + { steers: ["notice", "user"], continues: true }, + ] as const + ).map((scenario) => ({ + ...scenario, + terminalReason, + steered: scenario.steers.join(" then ") || "nothing", + })), ), - )("handles $terminalReason with active steering=$steered", ({ terminalReason, steered }) => + )("handles $terminalReason after steering $steered", ({ terminalReason, steers, continues }) => Effect.scoped( Effect.gen(function* () { const harness = yield* makeWakeHarness; - const idAllocator = yield* IdAllocator.IdAllocatorV2; - const attemptId = RunAttemptId.make("attempt-steering-abort"); - const input = makeClaudeTestTurnInput({ + const turn = makeClaudeTestTurnInput({ threadId: harness.threadId, providerThread: harness.providerThread, now: yield* DateTime.now, - attemptId, + attemptId: RunAttemptId.make("attempt-steering-abort"), text: "Audit the settings pages.", attachments: [], }); - yield* harness.runtime.startTurn(input); - if (steered) { - yield* harness.runtime.steerTurn({ - threadId: harness.threadId, - runId: input.runId, - providerThread: harness.providerThread, - providerTurnId: idAllocator.derive.providerTurn({ - driver: ClaudeAdapterV2.CLAUDE_PROVIDER, - nativeTurnId: `turn:${attemptId}`, + yield* harness.runtime.startTurn(turn); + for (const steer of steers) { + assert.equal( + yield* steerActiveTurn({ + harness, + turn, + message: steer === "user" ? userSteer : completionNotice, }), - message: { - createdBy: "user", - creationSource: "web", - messageId: MessageId.make("message-steering-abort"), - text: "Include the hierarchy mock.", - attachments: [], - }, - }); - assert.equal(harness.offeredMessages[1]?.priority, "now"); + steer === "user" ? "now" : "next", + ); } yield* Queue.offer( harness.sdkMessages, @@ -2346,28 +2428,82 @@ describe("ClaudeAdapterV2 background wake turns", () => { terminalReason, }), ); - if (steered) { - yield* Queue.offer(harness.sdkMessages, wakeAssistant); - yield* Queue.offer( - harness.sdkMessages, - makeResultFrame({ - uuid: "00000000-0000-4000-8000-000000000902", - result: "Audit finished after the steer.", - }), - ); - } + // What Claude sends next belongs to this turn only if the abort did not end it. + yield* Queue.offer(harness.sdkMessages, wakeAssistant); + yield* Queue.offer( + harness.sdkMessages, + makeResultFrame({ + uuid: "00000000-0000-4000-8000-000000000902", + result: "Audit finished after the steer.", + }), + ); const terminal = yield* Queue.take(harness.terminalReceipts); - assert.equal(terminal.status, steered ? "completed" : "interrupted"); - if (steered) { - assert.isTrue( - harness.events.some( - (event) => - event.type === "turn_item.updated" && - event.turnItem.type === "assistant_message" && - event.turnItem.text === WAKE_ASSISTANT_TEXT, - ), - ); - } + assert.equal(terminal.status, continues ? "completed" : "interrupted"); + assert.equal( + harness.events.some( + (event) => + event.type === "turn_item.updated" && + event.turnItem.type === "assistant_message" && + event.turnItem.text === WAKE_ASSISTANT_TEXT, + ), + continues, + ); + assert.lengthOf(harness.terminalEvents(), 1); + }).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ), + ); + + it.effect("ends a steered turn on its abort result once the user stops it", () => + Effect.scoped( + Effect.gen(function* () { + const closeGate = yield* Deferred.make(); + const testScope = yield* Scope.Scope; + yield* Scope.addFinalizer(testScope, Deferred.succeed(closeGate, undefined)); + const interruptStarted = yield* Deferred.make(); + const harness = yield* makeWakeHarnessWithOptions({ + close: (sdkMessages) => + Deferred.await(closeGate).pipe(Effect.andThen(Queue.shutdown(sdkMessages))), + interrupt: Deferred.succeed(interruptStarted, undefined), + }); + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const turn = makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("attempt-steering-stop"), + text: "Audit the settings pages.", + attachments: [], + }); + yield* harness.runtime.startTurn(turn); + yield* steerActiveTurn({ harness, turn, message: userSteer }); + yield* steerActiveTurn({ harness, turn, message: completionNotice }); + yield* harness.runtime + .interruptTurn({ + providerThread: harness.providerThread, + providerTurnId: idAllocator.derive.providerTurn({ + driver: ClaudeAdapterV2.CLAUDE_PROVIDER, + nativeTurnId: `turn:${turn.attemptId}`, + }), + }) + .pipe(Effect.forkScoped); + yield* Deferred.await(interruptStarted); + yield* Queue.offer( + harness.sdkMessages, + makeResultFrame({ + uuid: "00000000-0000-4000-8000-000000000903", + result: "", + terminalReason: "aborted_tools", + }), + ); + const terminalized = Exit.isSuccess( + yield* awaitUntil( + () => harness.terminalEvents().length === 1, + "interrupted terminal", + ).pipe(Effect.exit), + ); + yield* Deferred.succeed(closeGate, undefined); + assert.isTrue(terminalized); + assert.equal(harness.terminalEvents()[0]?.status, "interrupted"); assert.lengthOf(harness.terminalEvents(), 1); }).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), ), @@ -4818,6 +4954,73 @@ describe("ClaudeAdapterV2 background wake turns", () => { ), ); + // A `next` notice that arrives during the final reply misses the turn's last + // tool boundary. Claude answers it in a turn of its own right after the + // result: `init`, the reply, then a result with no task-notification origin. + it.effect("shows a completion notice Claude answers after the turn in one continuation", () => + Effect.scoped( + Effect.gen(function* () { + const harness = yield* makeWakeHarness; + const now = yield* DateTime.now; + const noticeReply = "Read the delegated result."; + const turn = makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now, + attemptId: RunAttemptId.make("attempt-claude-notice-tail-1"), + text: "Summarize the audit.", + attachments: [], + }); + const continuation = makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now, + attemptId: RunAttemptId.make("attempt-claude-notice-tail-2"), + text: "Background task completed.", + attachments: [], + providerTurnOrdinal: 2, + messageCreatedBy: "agent", + messageCreationSource: "provider", + }); + + yield* harness.runtime.startTurn(turn); + assert.equal(yield* steerActiveTurn({ harness, turn, message: completionNotice }), "next"); + yield* Queue.offer(harness.sdkMessages, turnOneResult); + assert.equal((yield* Queue.take(harness.terminalReceipts)).status, "completed"); + + yield* harness.offerAndWait(wakeTurnInit); + assert.lengthOf(harness.continuationRequests, 1); + yield* harness.offerAndWait( + makeAssistantTextFrame({ + uuid: "00000000-0000-4000-8000-000000000120", + text: noticeReply, + }), + ); + yield* harness.offerAndWait( + makeResultFrame({ uuid: "00000000-0000-4000-8000-000000000121", result: noticeReply }), + ); + assert.lengthOf(harness.continuationRequests, 1); + assert.isTrue(yield* harness.hasPendingBackgroundWork); + assert.lengthOf(harness.terminalEvents(), 1); + + yield* harness.runtime.startTurn(continuation); + assert.equal((yield* Queue.take(harness.terminalReceipts)).status, "completed"); + const replies = harness.events.flatMap((event) => + event.type === "message.updated" && event.message.text === noticeReply + ? [event.message] + : [], + ); + assert.equal(new Set(replies.map((reply) => reply.id)).size, 1); + assert.deepEqual([...new Set(replies.map((reply) => reply.runId))], [continuation.runId]); + // The prompt and the steered notice; the continuation prompts nothing. + assert.lengthOf(harness.offeredMessages, 2); + assert.lengthOf(harness.continuationRequests, 1); + assert.lengthOf(harness.terminalEvents(), 2); + assert.isFalse(yield* harness.hasPendingBackgroundWork); + }).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ), + ); + it.effect("does not offer a continuation for notification-only opaque work", () => Effect.scoped( Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts index ba0cfbd8fbd5..354b25ffbd93 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts @@ -2373,6 +2373,21 @@ function isClaudeActiveSteeringAbortResult(message: SDKResultMessage): boolean { ); } +// The SDK priority for a message steered into a running turn. Claude ends the +// turn to deliver a `now` message, and the model sees tool calls it had issued +// but not started as refused by the user. A delegated-task completion the +// server steers in interrupts nothing, so it goes in as `next`. Claude reads it +// at the turn's next tool boundary and every issued call still runs. A notice +// that lands during the final reply gets a native turn of its own afterwards, +// which reaches the thread as a continuation run. +function claudeSteerPriority(message: ProviderAdapter.ProviderAdapterV2TurnMessage) { + return message.createdBy === "agent" && + message.creationSource === "server" && + message.delegatedCompletion !== undefined + ? "next" + : "now"; +} + function isClaudeProviderContinuationTurn( input: ProviderAdapter.ProviderAdapterV2TurnInput, ): boolean { @@ -7380,22 +7395,27 @@ export function makeClaudeAdapterV2( detail: `Claude provider turn ${turnInput.providerTurnId} is not the active turn.`, }); } + const priority = claudeSteerPriority(turnInput.message); const userMessage = yield* makeClaudeUserMessageWithAttachments({ text: applyClaudePromptEffortPrefix( turnInput.message.text, compileClaudeModelSelection(currentTurn.input.modelSelection).promptEffort, ), attachments: turnInput.message.attachments, - priority: "now", + priority, attachmentsDir, fileSystem, skillNames: yield* userInvocableSkillNames(currentTurn.input.runtimePolicy.cwd), }); - yield* Ref.update(steeredTurns, (current) => { - const next = new Set(current); - next.add(turnInput.providerTurnId); - return next; - }); + // Only a `now` message aborts the turn, so only it expects the + // abort result that follows. + if (priority === "now") { + yield* Ref.update(steeredTurns, (current) => { + const next = new Set(current); + next.add(turnInput.providerTurnId); + return next; + }); + } yield* existing.query.offer(userMessage); }, (effect, turnInput) => diff --git a/apps/server/src/orchestration-v2/ProviderAdapter.ts b/apps/server/src/orchestration-v2/ProviderAdapter.ts index ac300004322e..07747a57c7d3 100644 --- a/apps/server/src/orchestration-v2/ProviderAdapter.ts +++ b/apps/server/src/orchestration-v2/ProviderAdapter.ts @@ -62,6 +62,7 @@ export const ProviderAdapterV2TurnMessage = Schema.Struct({ creationSource: OrchestrationV2ConversationMessage.fields.creationSource, scheduledTaskId: OrchestrationV2ConversationMessage.fields.scheduledTaskId, senderThreadId: OrchestrationV2ConversationMessage.fields.senderThreadId, + delegatedCompletion: OrchestrationV2ConversationMessage.fields.delegatedCompletion, }); export type ProviderAdapterV2TurnMessage = typeof ProviderAdapterV2TurnMessage.Type; diff --git a/apps/server/src/orchestration-v2/ProviderTurnControlService.ts b/apps/server/src/orchestration-v2/ProviderTurnControlService.ts index 47f7aa80bdc4..1ecc065df4d8 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnControlService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnControlService.ts @@ -314,6 +314,9 @@ export const layer: Layer.Layer< ...(message.senderThreadId === undefined ? {} : { senderThreadId: message.senderThreadId }), + ...(message.delegatedCompletion === undefined + ? {} + : { delegatedCompletion: message.delegatedCompletion }), }, }) .pipe( diff --git a/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts b/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts index 6a8dded63caa..0b07b2b7f3b6 100644 --- a/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts +++ b/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts @@ -27,6 +27,7 @@ import { ProviderAdapterSteerRunError, type ProviderAdapterV2Event, type ProviderAdapterV2Shape, + type ProviderAdapterV2SteerInput, type ProviderAdapterV2TurnInput, } from "./ProviderAdapter.ts"; import * as ProviderAdapterRegistry from "./ProviderAdapterRegistry.ts"; @@ -68,6 +69,7 @@ it.effect.each( const steerEntered = yield* Deferred.make(); const rejectSteer = yield* Deferred.make(); let steerCalls = 0; + const steered: ProviderAdapterV2SteerInput["message"][] = []; const capabilities = { ...CodexProviderCapabilitiesV2, turns: { @@ -145,6 +147,7 @@ it.effect.each( steerTurn: (turn) => Effect.gen(function* () { steerCalls += 1; + steered.push(turn.message); if (timing === "after delivery") return; yield* Deferred.succeed(steerEntered, undefined); yield* Deferred.await(rejectSteer); @@ -333,6 +336,11 @@ it.effect.each( if (timing === "after delivery") { assert.equal(steerCalls, 1); assert.equal(started.length, 1); + // The adapter is told whose completion it steers in, and nothing for a user steer. + assert.deepEqual( + steered[0]?.delegatedCompletion, + mailbox ? { parentRunId: first.runId, generation: 1, taskIds: [taskId] } : undefined, + ); if (mailbox) { const delivered = yield* orchestrator.getThreadProjection(threadId); assert.equal(delivered.subagents[0]?.completionDelivery?.state, "delivered"); From cf9bdfa05194b8d34e102fed6e813ff121163459 Mon Sep 17 00:00:00 2001 From: Wout Stiens <71498452+StiensWout@users.noreply.github.com> Date: Sun, 4 Oct 2026 15:19:37 +0200 Subject: [PATCH 2/2] fix(server): a late completion notice no longer takes over the next Claude turn A "next" completion notice that lands during the final reply is answered by Claude in a turn of its own after the result. That turn echoed no prompt uuid, so when T3 had already started a queued run, the adapter took the notice's result for that run's own. The run ended on the notice's reply and its real reply showed up on a continuation run. The Claude adapter now stamps a "next" steer with a uuid derived from the steered message. Claude echoes it on the notice's turn, so the existing prompt-echo gate sends that output to a continuation run and leaves the queued run to end on its own result. "now" steers stay unstamped. Refs #15351 Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Adapters/ClaudeAdapterV2.test.ts | 186 +++++++++++++++--- .../Adapters/ClaudeAdapterV2.ts | 15 +- 2 files changed, 168 insertions(+), 33 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts index 0f8752cde989..f4bcac7489cb 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts @@ -2336,8 +2336,11 @@ describe("ClaudeAdapterV2 background wake turns", () => { attachments: [], }, }); - return input.harness.offeredMessages[steerOrdinal]?.priority; + return input.harness.offeredMessages[steerOrdinal]; }); + // A frame of the turn that answers the offered message with this uuid. + const echoing = (frame: SDKMessage, uuid: SDKUserMessage["uuid"]) => + claudeSdkFrame({ ...frame, user_message_uuid: uuid, user_message_uuids: [uuid] }); it.effect("steers only a server-created delegated completion with next priority", () => Effect.scoped( @@ -2368,11 +2371,13 @@ describe("ClaudeAdapterV2 background wake turns", () => { }, ] as const; for (const steer of steers) { - assert.equal( - yield* steerActiveTurn({ harness, turn, message: steer.message }), - steer.priority, - `${steer.message.createdBy}/${steer.message.creationSource}`, - ); + const offered = yield* steerActiveTurn({ harness, turn, message: steer.message }); + const sender = `${steer.message.createdBy}/${steer.message.creationSource}`; + assert.equal(offered?.priority, steer.priority, sender); + // Only a notice can be answered after the turn, so only it carries + // a uuid, and never the prompt's. + assert.equal(offered?.uuid !== undefined, steer.priority === "next", sender); + assert.notEqual(offered?.uuid, harness.offeredMessages[0]?.uuid, sender); } }).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), ), @@ -2411,14 +2416,12 @@ describe("ClaudeAdapterV2 background wake turns", () => { }); yield* harness.runtime.startTurn(turn); for (const steer of steers) { - assert.equal( - yield* steerActiveTurn({ - harness, - turn, - message: steer === "user" ? userSteer : completionNotice, - }), - steer === "user" ? "now" : "next", - ); + const offered = yield* steerActiveTurn({ + harness, + turn, + message: steer === "user" ? userSteer : completionNotice, + }); + assert.equal(offered?.priority, steer === "user" ? "now" : "next"); } yield* Queue.offer( harness.sdkMessages, @@ -2495,15 +2498,9 @@ describe("ClaudeAdapterV2 background wake turns", () => { terminalReason: "aborted_tools", }), ); - const terminalized = Exit.isSuccess( - yield* awaitUntil( - () => harness.terminalEvents().length === 1, - "interrupted terminal", - ).pipe(Effect.exit), - ); + const terminal = yield* Queue.take(harness.terminalReceipts); yield* Deferred.succeed(closeGate, undefined); - assert.isTrue(terminalized); - assert.equal(harness.terminalEvents()[0]?.status, "interrupted"); + assert.equal(terminal.status, "interrupted"); assert.lengthOf(harness.terminalEvents(), 1); }).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), ), @@ -4956,7 +4953,8 @@ describe("ClaudeAdapterV2 background wake turns", () => { // A `next` notice that arrives during the final reply misses the turn's last // tool boundary. Claude answers it in a turn of its own right after the - // result: `init`, the reply, then a result with no task-notification origin. + // result: `init`, then a reply and a result that echo the notice's uuid and + // carry no task-notification origin. it.effect("shows a completion notice Claude answers after the turn in one continuation", () => Effect.scoped( Effect.gen(function* () { @@ -4984,20 +4982,27 @@ describe("ClaudeAdapterV2 background wake turns", () => { }); yield* harness.runtime.startTurn(turn); - assert.equal(yield* steerActiveTurn({ harness, turn, message: completionNotice }), "next"); + const notice = yield* steerActiveTurn({ harness, turn, message: completionNotice }); + assert.equal(notice?.priority, "next"); yield* Queue.offer(harness.sdkMessages, turnOneResult); assert.equal((yield* Queue.take(harness.terminalReceipts)).status, "completed"); yield* harness.offerAndWait(wakeTurnInit); assert.lengthOf(harness.continuationRequests, 1); yield* harness.offerAndWait( - makeAssistantTextFrame({ - uuid: "00000000-0000-4000-8000-000000000120", - text: noticeReply, - }), + echoing( + makeAssistantTextFrame({ + uuid: "00000000-0000-4000-8000-000000000120", + text: noticeReply, + }), + notice?.uuid, + ), ); yield* harness.offerAndWait( - makeResultFrame({ uuid: "00000000-0000-4000-8000-000000000121", result: noticeReply }), + echoing( + makeResultFrame({ uuid: "00000000-0000-4000-8000-000000000121", result: noticeReply }), + notice?.uuid, + ), ); assert.lengthOf(harness.continuationRequests, 1); assert.isTrue(yield* harness.hasPendingBackgroundWork); @@ -5021,6 +5026,129 @@ describe("ClaudeAdapterV2 background wake turns", () => { ), ); + // The same tail turn, with a user message queued behind the turn the notice + // was steered into. T3 starts that run while Claude still answers the notice, + // so the notice's output arrives with the user's prompt pending. + it.effect("keeps a completion notice's reply off the user turn queued behind it", () => + Effect.scoped( + Effect.gen(function* () { + const harness = yield* makeWakeHarness; + const now = yield* DateTime.now; + const noticeReply = "Read the delegated result."; + const userReply = "The build passed."; + const shown = (text: string) => + harness.events.flatMap((event) => + event.type === "message.updated" && event.message.text === text ? [event.message] : [], + ); + const runsShowing = (text: string) => [...new Set(shown(text).map((reply) => reply.runId))]; + const turn = makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now, + attemptId: RunAttemptId.make("attempt-claude-notice-race-1"), + text: "Summarize the audit.", + attachments: [], + }); + const queued = makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now, + attemptId: RunAttemptId.make("attempt-claude-notice-race-2"), + text: "How is the build going?", + attachments: [], + providerTurnOrdinal: 2, + }); + const continuation = makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now, + attemptId: RunAttemptId.make("attempt-claude-notice-race-3"), + text: "Background task completed.", + attachments: [], + providerTurnOrdinal: 3, + messageCreatedBy: "agent", + messageCreationSource: "provider", + }); + + yield* harness.runtime.startTurn(turn); + const prompt = harness.offeredMessages[0]; + // This CLI echoes a prompt on the first frame of the turn answering it. + yield* harness.offerAndWait( + echoing( + makeAssistantTextFrame({ + uuid: "00000000-0000-4000-8000-000000000130", + text: "Audit summary.", + }), + prompt?.uuid, + ), + ); + const notice = yield* steerActiveTurn({ harness, turn, message: completionNotice }); + yield* Queue.offer( + harness.sdkMessages, + echoing( + makeResultFrame({ + uuid: "00000000-0000-4000-8000-000000000131", + result: "Audit summary.", + }), + prompt?.uuid, + ), + ); + assert.equal((yield* Queue.take(harness.terminalReceipts)).status, "completed"); + + yield* harness.offerAndWait(wakeTurnInit); + yield* harness.runtime.startTurn(queued); + const queuedPrompt = harness.offeredMessages[2]; + yield* harness.offerAndWait( + echoing( + makeAssistantTextFrame({ + uuid: "00000000-0000-4000-8000-000000000132", + text: noticeReply, + }), + notice?.uuid, + ), + ); + yield* harness.offerAndWait( + echoing( + makeResultFrame({ uuid: "00000000-0000-4000-8000-000000000133", result: noticeReply }), + notice?.uuid, + ), + ); + yield* harness.offerAndWait( + echoing( + makeAssistantTextFrame({ + uuid: "00000000-0000-4000-8000-000000000134", + text: userReply, + }), + queuedPrompt?.uuid, + ), + ); + yield* harness.offerAndWait( + echoing( + makeResultFrame({ uuid: "00000000-0000-4000-8000-000000000135", result: userReply }), + queuedPrompt?.uuid, + ), + ); + + // The user run ends on its own result, showing its own reply alone. + assert.equal((yield* Queue.take(harness.terminalReceipts)).status, "completed"); + assert.deepEqual(runsShowing(userReply), [queued.runId]); + assert.deepEqual(runsShowing(noticeReply), []); + assert.isTrue(yield* harness.hasPendingBackgroundWork); + + yield* harness.runtime.startTurn(continuation); + assert.equal((yield* Queue.take(harness.terminalReceipts)).status, "completed"); + assert.deepEqual(runsShowing(noticeReply), [continuation.runId]); + assert.equal(new Set(shown(noticeReply).map((reply) => reply.id)).size, 1); + assert.deepEqual(runsShowing(userReply), [queued.runId]); + // Two prompts and the steered notice; the continuation prompts nothing. + assert.lengthOf(harness.offeredMessages, 3); + assert.lengthOf(harness.continuationRequests, 1); + assert.lengthOf(harness.terminalEvents(), 3); + assert.isFalse(yield* harness.hasPendingBackgroundWork); + }).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ), + ); + it.effect("does not offer a continuation for notification-only opaque work", () => Effect.scoped( Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts index 354b25ffbd93..ac655993846c 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts @@ -1356,10 +1356,11 @@ const makeClaudeUserMessageWithAttachments = Effect.fnUntraced(function* (input: } satisfies SDKUserMessage; }); -// Stable per run attempt, so a replayed prompt offer matches its recording. -// Claude echoes it back as user_message_uuid on the turn that answers it. -export function claudePromptUuid(attemptId: string): NonNullable { - const hex = NodeCrypto.createHash("sha256").update(`t3-claude-prompt:${attemptId}`).digest("hex"); +// Stable per seed (a prompt's run attempt, a steered notice's message), so a +// replayed offer matches its recording. Claude echoes it back as +// user_message_uuid on the turn that answers it. +export function claudePromptUuid(seed: string): NonNullable { + const hex = NodeCrypto.createHash("sha256").update(`t3-claude-prompt:${seed}`).digest("hex"); const variant = ((Number.parseInt(hex[16]!, 16) & 0x3) | 0x8).toString(16); return `${hex.slice(0, 8)}-${hex.slice(8, 12)}-4${hex.slice(13, 16)}-${variant}${hex.slice(17, 20)}-${hex.slice(20, 32)}`; } @@ -7403,6 +7404,12 @@ export function makeClaudeAdapterV2( ), attachments: turnInput.message.attachments, priority, + // Claude may answer a `next` notice in a turn of its own after + // this one. That turn echoes this uuid, so it never passes for + // the turn answering whichever prompt is offered next. + ...(priority === "next" + ? { uuid: claudePromptUuid(`steer:${turnInput.message.messageId}`) } + : {}), attachmentsDir, fileSystem, skillNames: yield* userInvocableSkillNames(currentTurn.input.runtimePolicy.cwd),