From 8dd73ca37c13a0f0716924b8b255d346902adc32 Mon Sep 17 00:00:00 2001 From: Yash Singh Date: Sun, 4 Oct 2026 23:18:10 -0500 Subject: [PATCH] fix(server): queue background notifications during active tools - Preserve Claude's native tool cancellation metadata and avoid interrupting active tools for automatic deliveries - Queue scheduled prompts for their bound thread --- .../Adapters/ClaudeAdapterV2.test.ts | 85 ++++ .../Adapters/ClaudeAdapterV2.ts | 35 +- ...laudeAutomaticDelivery.integration.test.ts | 420 ++++++++++++++++++ .../src/orchestration-v2/Orchestrator.ts | 3 +- .../SteeringCompletion.integration.test.ts | 7 +- .../scheduledTasks/ScheduledTaskService.ts | 3 +- .../contracts/src/orchestrationV2.test.ts | 29 ++ packages/contracts/src/orchestrationV2.ts | 3 + 8 files changed, 581 insertions(+), 4 deletions(-) create mode 100644 apps/server/src/orchestration-v2/ClaudeAutomaticDelivery.integration.test.ts diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts index 2e237bf3bb4c..5a732985603d 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts @@ -867,6 +867,20 @@ describe("ClaudeAdapterV2 context usage", () => { }); describe("ClaudeAdapterV2 session permissions", () => { + it("keeps explicit user refusals classified as user_reject", () => { + const result = ClaudeAdapterV2.permissionResultFromDecision({ + toolName: "Bash", + decision: "decline", + toolInput: { command: "make" }, + toolUseID: "denied-build", + }); + assert.equal(result.behavior, "deny"); + if (result.behavior !== "deny") return; + assert.equal(result.decisionClassification, "user_reject"); + assert.equal(result.message, "User declined tool execution."); + assert.equal(result.interrupt, undefined); + }); + it("forces suggested permission updates to session scope", () => { const result = ClaudeAdapterV2.permissionResultFromDecision({ toolName: "Bash", @@ -2300,6 +2314,77 @@ describe("ClaudeAdapterV2 background wake turns", () => { }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), ); + it.effect.each(["cancelled", "denied", "permission_denied", undefined])( + "preserves native tool non-execution metadata %s without inferring a denial from text", + (kind) => + Effect.scoped( + Effect.gen(function* () { + const harness = yield* makeWakeHarness; + yield* harness.runtime.startTurn( + makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("attempt-tool-non-execution"), + text: "Run the tool.", + attachments: [], + }), + ); + // Same error text can describe a cancellation or a real refusal. + // Each result must use its own metadata, even in a multi-result frame. + yield* harness.offerAndWait( + claudeSdkFrame({ + type: "user", + uuid: "tool-non-execution", + session_id: WAKE_NATIVE_SESSION, + parent_tool_use_id: null, + message: { + role: "user", + content: [ + { + type: "tool_result", + tool_use_id: "tool-error", + is_error: true, + content: "STOP and wait for the user.", + }, + { type: "tool_result", tool_use_id: "tool-ok", is_error: false, content: "OK" }, + ], + }, + ...(kind === undefined + ? {} + : { + tool_result_meta: [ + { id: "tool-error", non_execution_kind: kind }, + { id: "tool-ok", non_execution_kind: null }, + ], + }), + }), + ); + yield* Queue.offer( + harness.sdkMessages, + makeResultFrame({ uuid: "result-non-execution", result: "Done" }), + ); + yield* Queue.take(harness.terminalReceipts); + const items = harness.events.flatMap((event) => + event.type === "turn_item.updated" && event.turnItem.type === "dynamic_tool" + ? [event.turnItem] + : [], + ); + const failed = items.findLast((item) => item.nativeItemRef?.nativeId === "tool-error")!; + assert.equal(failed.status, kind === "cancelled" ? "cancelled" : "failed"); + assert.equal(failed.toolNonExecutionKind, kind); + const ok = items.findLast((item) => item.nativeItemRef?.nativeId === "tool-ok")!; + assert.equal(ok.status, "completed"); + assert.equal(ok.toolNonExecutionKind, undefined); + const node = harness.events.findLast( + (event) => + event.type === "node.updated" && event.node.nativeItemRef?.nativeId === "tool-error", + ); + assert.equal(node?.type === "node.updated" ? node.node.status : undefined, failed.status); + }).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ), + ); + it.effect.each( (["aborted_tools", "aborted_streaming"] as const).flatMap((terminalReason) => [true, false].map((steered) => ({ terminalReason, steered })), diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts index ba0cfbd8fbd5..5e8f0e81bb0e 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts @@ -198,6 +198,7 @@ export const ClaudeProviderCapabilitiesV2 = { emitsTurnCompleted: true, supportsInterrupt: true, supportsActiveSteering: true, + activeSteeringInterruptsTools: true, supportsSteeringByInterruptRestart: false, supportsQueuedMessages: true, terminalStatusQuality: "strong", @@ -2192,6 +2193,24 @@ function claudeToolResultEntriesFromMessage(message: SDKMessage): ReadonlyArray< ]; } +// These wire fields are not yet declared by the SDK's SDKUserMessage type. +const ClaudeToolResultMetadata = Schema.Struct({ + tool_result_meta: Schema.Array( + Schema.Struct({ + id: Schema.String, + non_execution_kind: Schema.optional(Schema.NullOr(Schema.String)), + }), + ), +}); +const isClaudeToolResultMetadata = Schema.is(ClaudeToolResultMetadata); + +function claudeToolNonExecutionKind(message: SDKMessage, toolUseId: string) { + return isClaudeToolResultMetadata(message) + ? (message.tool_result_meta.find((meta) => meta.id === toolUseId)?.non_execution_kind ?? + undefined) + : undefined; +} + function parentToolUseIdFromSdkMessage(message: SDKMessage): string | null { return message.type === "assistant" || message.type === "user" ? message.parent_tool_use_id @@ -3734,6 +3753,7 @@ export function makeClaudeAdapterV2( readonly startedAt: DateTime.Utc; readonly updatedAt: DateTime.Utc; readonly presentation: ClaudeToolPresentation | undefined; + readonly toolNonExecutionKind?: string; }) => { const completedAt = input.status === "running" ? null : input.updatedAt; const nodeId = idAllocator.derive.nodeFromProviderItem({ @@ -3781,6 +3801,9 @@ export function makeClaudeAdapterV2( }) : undefined; const itemBase = { + ...(input.toolNonExecutionKind === undefined + ? {} + : { toolNonExecutionKind: input.toolNonExecutionKind }), id: turnItemId, threadId: input.threadId, runId: input.runId, @@ -6102,6 +6125,10 @@ export function makeClaudeAdapterV2( parentToolUseId, })); const completedAt = yield* DateTime.now; + const toolNonExecutionKind = claudeToolNonExecutionKind( + message, + toolResult.tool_use_id, + ); const artifacts = buildToolCallArtifacts({ context, nativeItemId: toolCall.nativeItemId, @@ -6114,7 +6141,13 @@ export function makeClaudeAdapterV2( parentNodeId: toolCall.parentNodeId, ordinal: toolCall.ordinal, output, - status: isClaudeToolResultError(toolResult) ? "failed" : "completed", + status: + toolNonExecutionKind === "cancelled" + ? "cancelled" + : isClaudeToolResultError(toolResult) + ? "failed" + : "completed", + ...(toolNonExecutionKind === undefined ? {} : { toolNonExecutionKind }), startedAt: toolCall.startedAt, updatedAt: completedAt, presentation: toolCall.presentation, diff --git a/apps/server/src/orchestration-v2/ClaudeAutomaticDelivery.integration.test.ts b/apps/server/src/orchestration-v2/ClaudeAutomaticDelivery.integration.test.ts new file mode 100644 index 000000000000..b13a8b5d90c1 --- /dev/null +++ b/apps/server/src/orchestration-v2/ClaudeAutomaticDelivery.integration.test.ts @@ -0,0 +1,420 @@ +import type { SDKMessage, SDKUserMessage } from "@anthropic-ai/claude-agent-sdk"; +import * as NodeCrypto from "@effect/platform-node/NodeCrypto"; +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { assert, it } from "@effect/vitest"; +import { + ClaudeSettings, + CommandId, + EventId, + MessageId, + NodeId, + ProjectId, + ProviderInstanceId, + ScheduledTaskUpsertInput, + ThreadId, + type OrchestrationV2DomainEvent, +} from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Fiber from "effect/Fiber"; +import * as Layer from "effect/Layer"; +import * as Path from "effect/Path"; +import * as Queue from "effect/Queue"; +import * as Schema from "effect/Schema"; +import * as Stream from "effect/Stream"; +import * as ScheduledTaskService from "../scheduledTasks/ScheduledTaskService.ts"; +import * as Scheduler from "../scheduling/Scheduler.ts"; +import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; +import * as ClaudeAdapterV2 from "./Adapters/ClaudeAdapterV2.ts"; +import * as EffectWorker from "./EffectWorker.ts"; +import * as EventSink from "./EventSink.ts"; +import * as IdAllocator from "./IdAllocator.ts"; +import * as LegacyV1ThreadImporter from "./legacy/LegacyV1ThreadImporter.ts"; +import * as Orchestrator from "./Orchestrator.ts"; +import * as ProviderAdapterRegistry from "./ProviderAdapterRegistry.ts"; +import * as ThreadLaunchService from "./ThreadLaunchService.ts"; +import * as ThreadManagementService from "./ThreadManagementService.ts"; +import { makeOrchestratorV2ReplayLayerWithRegistry } from "./testkit/ProviderReplayHarness.ts"; +import { checkpointWorkspace } from "./testkit/ReplayFixtureWorkspace.ts"; + +const sessionId = "automatic-delivery-session"; +const settings = Schema.decodeSync(ClaudeSettings)({}); +const decodeScheduledTask = Schema.decodeEffect(ScheduledTaskUpsertInput); +const refusal = + "The user doesn't want to take this action right now. STOP what you are doing and wait for the user to tell you how to proceed."; +const statusIds = ["status-1", "status-2", "status-3", "status-4"]; +const modelSelection = { + instanceId: ProviderInstanceId.make("claudeAgent"), + model: "claude-sonnet-4-6", +}; + +// Native frames match the captured Bash/status batch and tool_result_meta shape. +// The external SDK boundary is replayed; the adapter and durable dispatch are real. +function frame(value: unknown): SDKMessage { + return value as SDKMessage; +} + +const batch = frame({ + type: "assistant", + uuid: "batch", + session_id: sessionId, + parent_tool_use_id: null, + message: { + id: "tool-batch", + role: "assistant", + type: "message", + model: modelSelection.model, + content: [ + { type: "tool_use", id: "build", name: "Bash", input: { command: "make" } }, + ...statusIds.map((id) => ({ + type: "tool_use", + id, + name: "mcp__t3-code__task_status", + input: { taskId: id }, + })), + ], + stop_reason: "tool_use", + stop_sequence: null, + usage: { input_tokens: 1, output_tokens: 1 }, + }, +}); + +const toolResult = (id: string, cancelled: boolean) => + frame({ + type: "user", + uuid: `result-${id}`, + session_id: sessionId, + parent_tool_use_id: null, + message: { + role: "user", + content: [ + { + type: "tool_result", + tool_use_id: id, + is_error: cancelled, + content: cancelled ? refusal : id === "build" ? "build OK" : "Task completed normally", + }, + ], + }, + ...(cancelled ? { tool_result_meta: [{ id, non_execution_kind: "cancelled" }] } : {}), + }); + +const result = (uuid: string, aborted: boolean) => + frame({ + type: "result", + subtype: "success", + uuid, + session_id: sessionId, + is_error: false, + num_turns: 1, + result: "", + stop_reason: "end_turn", + terminal_reason: aborted ? "aborted_tools" : "completed", + permission_denials: [], + duration_ms: 1, + duration_api_ms: 1, + total_cost_usd: 0, + usage: { input_tokens: 1, output_tokens: 1 }, + modelUsage: {}, + }); + +it.effect.each(["child completion", "scheduled message", "user steering"] as const)( + "delivers %s with Claude's native pending-tool cancellation behavior", + (delivery) => + Effect.scoped( + Effect.gen(function* () { + const cwd = yield* checkpointWorkspace("claude-automatic-delivery"); + const sdkMessages = yield* Queue.unbounded(); + const offers: SDKUserMessage[] = []; + const nativeQueue: SDKUserMessage[] = []; + const batchAbort = new AbortController(); + const adapter = ClaudeAdapterV2.makeClaudeAdapterV2({ + instanceId: modelSelection.instanceId, + settings, + environment: {}, + attachmentsDir: cwd, + fileSystem: yield* FileSystem.FileSystem, + path: yield* Path.Path, + idAllocator: yield* IdAllocator.IdAllocatorV2, + queryRunner: { + allocateSessionId: Effect.succeed(sessionId), + open: () => + Effect.succeed({ + messages: Stream.fromQueue(sdkMessages), + offer: (message) => + Effect.sync(() => { + offers.push(message); + nativeQueue.push(message); + // Claude Code 2.1.289's queue watcher aborts when a now-priority + // command arrives, including while Bash blocks pending reads. + if (nativeQueue.some((command) => command.priority === "now")) { + batchAbort.abort({ kind: "interrupt" }); + } + }), + setModel: () => Effect.void, + setPermissionMode: () => Effect.void, + interrupt: Effect.die("automatic delivery must never interrupt"), + close: Effect.void, + }), + forkSession: () => Effect.die("unused"), + subagentLaunchToolUseId: () => Effect.succeed(null), + assertComplete: Effect.void, + }, + }); + yield* Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const worker = yield* EffectWorker.OrchestrationEffectWorkerV2; + const threadId = ThreadId.make("thread:automatic-delivery"); + const projectId = ProjectId.make("project:automatic-delivery"); + const watch = (predicate: (event: OrchestrationV2DomainEvent) => boolean) => + orchestrator.streamDomainEvents.pipe( + Stream.filter(predicate), + Stream.take(1), + Stream.runDrain, + Effect.forkScoped, + ); + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make("create"), + threadId, + projectId, + title: "Automatic delivery", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: cwd, + createdBy: "user", + creationSource: "web", + }); + const running = yield* watch( + (event) => event.type === "provider-turn.updated" && event.payload.status === "running", + ); + yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make("first"), + threadId, + messageId: MessageId.make("first"), + text: "Build and check all four tasks.", + attachments: [], + dispatchMode: { type: "start_immediately" }, + createdBy: "user", + creationSource: "web", + }); + yield* worker.drain(); + yield* Fiber.join(running); + const pending = yield* watch( + (event) => + event.type === "turn-item.updated" && + event.payload.nativeItemRef?.nativeId === "status-4", + ); + yield* Queue.offer(sdkMessages, batch); + yield* Fiber.join(pending); + const before = yield* orchestrator.getThreadProjection(threadId); + assert.equal(before.turnItems.filter((item) => item.status === "running").length, 5); + const parent = before.runs[0]!; + if (parent.rootNodeId === null) return yield* Effect.die("parent has no root node"); + const messageId = MessageId.make("automatic-notice"); + if (delivery === "child completion") { + const taskId = NodeId.make("completed-child"); + const now = yield* DateTime.now; + const sink = yield* EventSink.EventSinkV2; + yield* sink.write({ + events: [ + { + id: EventId.make("cohort"), + type: "run.updated", + threadId, + runId: parent.id, + occurredAt: now, + payload: { + ...parent, + delegatedCompletion: { + disposition: "open", + nextGeneration: 2, + delivery: { generation: 1, messageId, taskIds: [taskId] }, + }, + }, + }, + { + id: EventId.make("child"), + type: "subagent.updated", + threadId, + runId: parent.id, + nodeId: taskId, + occurredAt: now, + payload: { + id: taskId, + threadId, + runId: parent.id, + parentNodeId: parent.rootNodeId, + origin: "app_owned", + createdBy: "agent", + driver: ClaudeAdapterV2.CLAUDE_PROVIDER, + providerInstanceId: modelSelection.instanceId, + providerThreadId: null, + childThreadId: null, + nativeTaskRef: null, + prompt: "Background work", + title: "Child", + model: null, + completionWake: "always", + completionDelivery: { state: "claimed", observedByRunId: null }, + status: "completed", + result: "done", + startedAt: now, + completedAt: now, + updatedAt: now, + }, + }, + ], + }); + const command = { + type: "message.dispatch" as const, + commandId: CommandId.make("completion"), + threadId, + messageId, + text: "Child completed", + attachments: [], + dispatchMode: { type: "queue_after_active" as const }, + createdBy: "agent" as const, + creationSource: "server" as const, + delegatedCompletion: { parentRunId: parent.id, generation: 1, taskIds: [taskId] }, + }; + yield* orchestrator.dispatch(command); + yield* orchestrator.dispatch(command); + yield* orchestrator.recoverDelegatedTasks; + } else if (delivery === "scheduled message") { + yield* Effect.gen(function* () { + const service = yield* ScheduledTaskService.ScheduledTaskService; + const { task } = yield* service.upsert( + yield* decodeScheduledTask({ + title: "Status check", + prompt: "Check the four tasks.", + enabled: true, + schedule: { type: "interval", everyMs: 60_000 }, + projectId, + threadId, + workspaceStrategy: { type: "root" }, + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + }), + ); + const ran = yield* service.runNow({ id: task.id }); + assert.equal(ran.task.lastRunStatus, "succeeded"); + }).pipe( + Effect.provide( + ScheduledTaskService.layer.pipe( + Layer.provide(ThreadManagementService.layer), + Layer.provide( + Layer.mock(LegacyV1ThreadImporter.LegacyV1ThreadImporter)({ + ensureTranscript: () => + Effect.succeed({ importedThreadCount: 0, importedMessageCount: 0 }), + }), + ), + Layer.provide(Layer.mock(ThreadLaunchService.ThreadLaunchService)({})), + Layer.provide( + Layer.mergeAll(NodeCrypto.layer, Scheduler.layer, SqlitePersistenceMemory), + ), + ), + ), + ); + } else { + yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make("steer"), + threadId, + messageId, + text: "Change direction now.", + attachments: [], + dispatchMode: { type: "steer_active", targetRunId: parent.id }, + createdBy: "user", + creationSource: "web", + }); + } + yield* worker.drain(); + const explicitSteer = delivery === "user steering"; + assert.equal(batchAbort.signal.aborted, explicitSteer); + assert.equal(offers.length, explicitSteer ? 2 : 1); + const finished = yield* watch( + (event) => + event.type === "run.updated" && + event.payload.id === parent.id && + event.payload.status === "waiting", + ); + // Bash finishes before the pending status reads, as in the capture. + yield* Queue.offer(sdkMessages, toolResult("build", false)); + for (const id of statusIds) + yield* Queue.offer(sdkMessages, toolResult(id, batchAbort.signal.aborted)); + if (explicitSteer) yield* Queue.offer(sdkMessages, result("aborted", true)); + yield* Queue.offer(sdkMessages, result("completed", false)); + yield* Fiber.join(finished); + yield* worker.drain(); + yield* orchestrator.resumeQueuedRuns; + yield* worker.drain(); + if (!explicitSteer) { + const queued = (yield* orchestrator.getThreadProjection(threadId)).runs[1]!; + const noticeFinished = yield* watch( + (event) => + event.type === "run.updated" && + event.payload.id === queued.id && + event.payload.status === "waiting", + ); + const delivered = + delivery === "child completion" + ? yield* watch( + (event) => + event.type === "subagent.updated" && + event.payload.completionDelivery?.state === "delivered", + ) + : null; + yield* Queue.offer(sdkMessages, result("notice-completed", false)); + yield* Fiber.join(noticeFinished); + yield* worker.drain(); + if (delivered !== null) yield* Fiber.join(delivered); + } + const after = yield* orchestrator.getThreadProjection(threadId); + const reads = after.turnItems.filter((item) => + statusIds.includes(item.nativeItemRef?.nativeId ?? ""), + ); + assert.equal(reads.length, 4); + for (const read of reads) { + assert.equal(read.status, explicitSteer ? "cancelled" : "completed"); + assert.equal(read.toolNonExecutionKind, explicitSteer ? "cancelled" : undefined); + assert.equal(read.type === "dynamic_tool" && read.output === refusal, explicitSteer); + } + assert.equal(offers.length, 2); + assert.equal( + offers.filter((offer) => offer.priority === "now").length, + explicitSteer ? 1 : 0, + ); + if (!explicitSteer) { + assert.equal(after.runs.length, 2); + if (delivery === "child completion") { + assert.equal(after.messages.filter((message) => message.id === messageId).length, 1); + assert.equal(after.subagents[0]?.completionDelivery?.state, "delivered"); + } else { + assert.equal( + after.messages.filter((message) => message.scheduledTaskId !== undefined).length, + 1, + ); + } + } + yield* orchestrator.recoverDelegatedTasks; + yield* orchestrator.resumeQueuedRuns; + yield* worker.drain(); + assert.equal(offers.length, 2); + }).pipe( + Effect.provide( + makeOrchestratorV2ReplayLayerWithRegistry( + { name: "claude-automatic-delivery" }, + ProviderAdapterRegistry.makeSingleLayer(adapter), + { runEffectWorker: false }, + ), + ), + ); + }).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ), +); diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index abdfc29579cd..8172ec5bcd35 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -4596,7 +4596,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio .pipe(Effect.orElseSucceed(() => Option.none())); if ( Option.isSome(session) && - session.value.providerSession.capabilities.turns.supportsActiveSteering + session.value.providerSession.capabilities.turns.supportsActiveSteering && + session.value.providerSession.capabilities.turns.activeSteeringInterruptsTools !== true ) { dispatchMode = { type: "steer_active", targetRunId: active.id }; } diff --git a/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts b/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts index 6a8dded63caa..d72b036316c9 100644 --- a/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts +++ b/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts @@ -47,6 +47,7 @@ it.effect.each( "before dispatch", "after delivery", "without native steering", + "with interrupting native steering", "settled only", ] as const ).map((timing) => ({ @@ -57,7 +58,10 @@ it.effect.each( ) .filter( ({ mailbox, timing }) => - mailbox || (timing !== "without native steering" && timing !== "settled only"), + mailbox || + (timing !== "without native steering" && + timing !== "with interrupting native steering" && + timing !== "settled only"), ), )("delivers $label when completion wins $timing", ({ mailbox, timing }) => Effect.scoped( @@ -73,6 +77,7 @@ it.effect.each( turns: { ...CodexProviderCapabilitiesV2.turns, supportsActiveSteering: timing !== "without native steering", + activeSteeringInterruptsTools: timing === "with interrupting native steering", }, }; const adapter: ProviderAdapterV2Shape = { diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.ts index 400ee5d0b111..d2a6dedc46f5 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.ts @@ -545,7 +545,8 @@ export const layer = Layer.effect( text: prompt, attachments: [], modelSelection: active.modelSelection, - mode: "auto", + // Scheduled prompts must not interrupt tools in the bound thread. + mode: "queue", createdBy: active.createdBy, creationSource: active.creationSource, }), diff --git a/packages/contracts/src/orchestrationV2.test.ts b/packages/contracts/src/orchestrationV2.test.ts index 601597fb568a..aae9cc40bd8c 100644 --- a/packages/contracts/src/orchestrationV2.test.ts +++ b/packages/contracts/src/orchestrationV2.test.ts @@ -1038,6 +1038,35 @@ describe("orchestration V2 contracts", () => { }); }); +it("preserves tool cancellation and denial metadata through persisted and wire schemas", () => { + for (const kind of ["cancelled", "denied", undefined]) { + const item = decodeOrchestrationV2TurnItem({ + id: "tool-result", + threadId: "thread", + runId: null, + nodeId: null, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + type: "dynamic_tool", + toolName: "task_status", + input: { taskId: "child" }, + status: kind === "cancelled" ? "cancelled" : "failed", + ...(kind === undefined ? {} : { toolNonExecutionKind: kind }), + title: null, + startedAt: now, + completedAt: now, + updatedAt: now, + output: "Tool did not execute", + }); + const wire = encodeOrchestrationV2TurnItemJson(item); + expect(wire.toolNonExecutionKind).toBe(kind); + expect(decodeOrchestrationV2TurnItemJson(wire)).toEqual(item); + } +}); + it("round-trips typed notifications and keeps work outcome separate from item status", () => { const now = DateTime.makeUnsafe("2026-09-09T00:00:00Z"); const base = { diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 220fec2b15fb..675faa4b7f4f 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -211,6 +211,8 @@ export const OrchestrationV2TurnCapabilities = Schema.Struct({ emitsTurnCompleted: Schema.Boolean, supportsInterrupt: Schema.Boolean, supportsActiveSteering: Schema.Boolean, + // Some native steering mechanisms cancel pending tools before consuming the message. + activeSteeringInterruptsTools: Schema.optional(Schema.Boolean), supportsSteeringByInterruptRestart: Schema.Boolean, supportsQueuedMessages: Schema.Boolean, terminalStatusQuality: Schema.Literals(["strong", "weak", "none"]), @@ -1253,6 +1255,7 @@ export type OrchestrationV2UserMessageInputIntent = typeof OrchestrationV2UserMessageInputIntent.Type; const OrchestrationV2TurnItemBaseFields = { + toolNonExecutionKind: Schema.optional(Schema.String), toolSurface: Schema.optional(ToolActivitySurface), toolIcon: Schema.optional(ToolActivityIcon), toolSource: Schema.optional(ToolActivitySource),