diff --git a/apps/server/src/orchestration-v2/ContextHandoffDelivery.ts b/apps/server/src/orchestration-v2/ContextHandoffDelivery.ts index 2ef3a5a5cd89..52a80448f8c3 100644 --- a/apps/server/src/orchestration-v2/ContextHandoffDelivery.ts +++ b/apps/server/src/orchestration-v2/ContextHandoffDelivery.ts @@ -7,7 +7,11 @@ import * as Effect from "effect/Effect"; import * as Schema from "effect/Schema"; import { historyCost, renderHistory, selectHistory } from "./ContextHandoffBudget.ts"; -/** Persist before/after injection: an ambiguous pending delivery requires a fresh native thread. */ +/** + * Persist before/after injection: an ambiguous pending delivery requires a fresh native thread. + * Run `unsent` when the provider refused the turn before reading the prompt, so that + * refusal does not leave an ambiguous marker behind. + */ export const deliverContextHandoffs = Effect.fn("orchestrationV2.deliverContextHandoffs")( function* (input: { readonly handoffs: ReadonlyArray; @@ -28,7 +32,7 @@ export const deliverContextHandoffs = Effect.fn("orchestrationV2.deliverContextH handoff.delivery.status === "pending", ); if (pending.length === 0 || (input.deferInline && input.inject === undefined)) - return { context: "", delivered: Effect.void }; + return { context: "", delivered: Effect.void, unsent: Effect.void }; const budget = typeof input.budget === "number" ? input.budget : yield* input.budget; let coverage = pending .map( @@ -71,7 +75,7 @@ export const deliverContextHandoffs = Effect.fn("orchestrationV2.deliverContextH budget, }); if (historyCost(selected.messages, selected.context) > budget) { - if (input.deferInline) return { context: "", delivered: Effect.void }; + if (input.deferInline) return { context: "", delivered: Effect.void, unsent: Effect.void }; return yield* new ContextHandoffBudgetError(); } const omittedItemIds = new Set(selected.omittedItemIds); @@ -122,7 +126,7 @@ export const deliverContextHandoffs = Effect.fn("orchestrationV2.deliverContextH }); if (injected) { yield* persist("injected"); - return { context: "", delivered: Effect.void }; + return { context: "", delivered: Effect.void, unsent: Effect.void }; } } else { // Text-only delivery can also be accepted before a connection drops. @@ -132,11 +136,12 @@ export const deliverContextHandoffs = Effect.fn("orchestrationV2.deliverContextH // Compaction APIs cannot accept an inline transcript. Keep it available // for the next ordinary turn when native injection is unsupported. yield* Effect.forEach(pending, input.persist, { discard: true }); - return { context: "", delivered: Effect.void }; + return { context: "", delivered: Effect.void, unsent: Effect.void }; } return { context: renderHistory(selected.messages, selected.context), delivered: persist("inline"), + unsent: Effect.forEach(pending, input.persist, { discard: true }), }; }, ); diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index fa95b5af8272..cf2f6c68417b 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -19,6 +19,7 @@ import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Predicate from "effect/Predicate"; import * as Schema from "effect/Schema"; import * as GitWorkflowService from "../git/GitWorkflowService.ts"; @@ -64,6 +65,11 @@ export class ProviderTurnStartError extends Schema.TaggedError + Predicate.isTagged(error, "ClaudeBackgroundWorkBlocksQueryReplacementError") || + (Predicate.hasProperty(error, "cause") && refusedBeforePrompt(error.cause)); + export interface ProviderTurnStartServiceV2Shape { /** * Starts the run's provider turn. When `willRetry` is true, a session open @@ -1176,7 +1182,21 @@ export const layer: Layer.Layer< ...turnInput.message, text: context === "" ? userText : `${context}\n\nUser message:\n${userText}`, }, - }); + }).pipe( + // A pending marker would make the next turn abandon this native + // session, though the refused prompt never reached it. + Effect.tapError((error) => + refusedBeforePrompt(error) + ? delivery.unsent.pipe( + Effect.catchCause(() => + Effect.logWarning("Failed to restore unsent context handoffs", { + runId: run.id, + }), + ), + ) + : Effect.void, + ), + ); // The provider already accepted the turn. A stale pending marker // can force a fresh thread later, but must not stop live ingestion. yield* delivery.delivered.pipe( diff --git a/apps/server/src/orchestration-v2/testkit/ProviderSwitch.integration.test.ts b/apps/server/src/orchestration-v2/testkit/ProviderSwitch.integration.test.ts index eef3a9420d16..2aadfa4b5cfa 100644 --- a/apps/server/src/orchestration-v2/testkit/ProviderSwitch.integration.test.ts +++ b/apps/server/src/orchestration-v2/testkit/ProviderSwitch.integration.test.ts @@ -34,7 +34,10 @@ import * as SqlClient from "effect/sql/SqlClient"; import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; import { CommandPolicyCapabilityUnsupportedError } from "../CommandPolicy.ts"; -import { ClaudeProviderCapabilitiesV2 } from "../Adapters/ClaudeAdapterV2.ts"; +import { + ClaudeBackgroundWorkBlocksQueryReplacementError, + ClaudeProviderCapabilitiesV2, +} from "../Adapters/ClaudeAdapterV2.ts"; import { CodexProviderCapabilitiesV2, canReuseCodexContextUsage, @@ -53,6 +56,7 @@ import { type ProviderAdapterV2Event, type ProviderAdapterV2HistoricalContext, ProviderAdapterProtocolError, + ProviderAdapterTurnStartError, type ProviderAdapterV2Shape, type ProviderAdapterV2SessionRuntime, } from "../ProviderAdapter.ts"; @@ -104,6 +108,8 @@ function makeTestAdapter(input: { readonly capturedTurns: Ref.Ref>; readonly injectedHistory?: Ref.Ref>; readonly failStartOnce?: Ref.Ref; + /** Starts left to refuse the way Claude does while background work runs. */ + readonly refuseStarts?: Ref.Ref; readonly failInjectionOnce?: Ref.Ref; readonly nativeThreadGeneration?: Ref.Ref; readonly failResume?: boolean; @@ -226,6 +232,17 @@ function makeTestAdapter(input: { (yield* Ref.getAndSet(input.failStartOnce, false)) ) return yield* unimplemented(input.driver, "turn start failed after injection"); + if ( + input.refuseStarts !== undefined && + (yield* Ref.getAndUpdate(input.refuseStarts, (left) => Math.max(0, left - 1))) > 0 + ) + return yield* new ProviderAdapterTurnStartError({ + driver: input.driver, + threadId: turnInput.threadId, + providerThreadId: turnInput.providerThread.id, + runId: turnInput.runId, + cause: new ClaudeBackgroundWorkBlocksQueryReplacementError(), + }); yield* Effect.yieldNow; yield* Ref.update(input.capturedTurns, (turns) => [ ...turns, @@ -1188,6 +1205,101 @@ describe("orchestration v2 provider switching", () => { ), ); + it.live("keeps the native session after turns refused before reaching the provider", () => + Effect.scoped( + Effect.gen(function* () { + const cwd = yield* checkpointWorkspace("refused-start-keeps-native"); + const capturedTurns = yield* Ref.make>([]); + const refuseStarts = yield* Ref.make(0); + const generation = yield* Ref.make(0); + const registry = ProviderAdapterRegistry.makeLayer([ + makeTestAdapter({ + instanceId: CLAUDE_MODEL_SELECTION.instanceId, + driver: CLAUDE_DRIVER, + capabilities: ClaudeProviderCapabilitiesV2, + modelSelection: CLAUDE_MODEL_SELECTION, + responseByRunOrdinal: {}, + capturedTurns, + refuseStarts, + nativeThreadGeneration: generation, + }), + ]); + yield* Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const worker = yield* EffectWorker.OrchestrationEffectWorkerV2; + const run = Effect.fn("run")(function* (ordinal: number) { + yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make(`refused-start:${ordinal}`), + threadId, + messageId: MessageId.make(`refused-start:${ordinal}`), + createdBy: "user", + creationSource: "web", + text: `Request ${ordinal}`, + attachments: [], + modelSelection: CLAUDE_MODEL_SELECTION, + dispatchMode: { type: "start_immediately" }, + }); + yield* orchestrator.streamStoredEvents.pipe( + Stream.filter( + ({ event }) => + event.type === "run.updated" && + event.payload.ordinal === ordinal && + (event.payload.status === "completed" || event.payload.status === "failed"), + ), + Stream.runHead, + ); + yield* worker.drain(); + return (yield* orchestrator.getThreadProjection(threadId)).runs.at(-1)?.status; + }); + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make("refused-start:create"), + threadId, + projectId, + createdBy: "user", + creationSource: "web", + title: "Refused start", + modelSelection: CLAUDE_MODEL_SELECTION, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + }); + assert.equal(yield* run(1), "completed"); + // Run 3 carries run 2's missed request as a handoff and is refused too. + yield* Ref.set(refuseStarts, 2); + assert.equal(yield* run(2), "failed"); + assert.equal(yield* run(3), "failed"); + assert.equal(yield* run(4), "completed"); + + const turns = yield* Ref.get(capturedTurns); + assert.equal(turns.length, 2); + // Nothing reached the provider, so the next turn continues the same + // native session instead of replacing it with a summary. + assert.equal(yield* Ref.get(generation), 1); + assert.equal(turns[1]?.nativeThreadId, turns[0]?.nativeThreadId); + assert.include(turns[1]?.text, "Request 2"); + assert.include(turns[1]?.text, "Request 3"); + }).pipe( + Effect.provide( + makeOrchestratorV2ReplayLayerWithRegistry( + { + name: "refused-start-keeps-native", + runtimePolicyOverride: { + cwd, + approvalPolicy: "never", + sandboxPolicy: { type: "readOnly" }, + }, + }, + registry, + ), + ), + ); + }), + ), + ); + it.live.each( (["failed", "interrupted"] as const).flatMap((status) => [false, true].flatMap((queued) =>