Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 10 additions & 5 deletions apps/server/src/orchestration-v2/ContextHandoffDelivery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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* <InjectError = never, PersistError = never, BudgetError = never>(input: {
readonly handoffs: ReadonlyArray<OrchestrationV2ContextHandoff>;
Expand All @@ -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(
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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.
Expand All @@ -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 }),
};
},
);
Expand Down
22 changes: 21 additions & 1 deletion apps/server/src/orchestration-v2/ProviderTurnStartService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -64,6 +65,11 @@ export class ProviderTurnStartError extends Schema.TaggedError<ProviderTurnStart

const isProviderTurnStartError = Schema.is(ProviderTurnStartError);

/** Claude refuses to replace a process running background work before it reads the prompt. */
const refusedBeforePrompt = (error: unknown): boolean =>
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
Expand Down Expand Up @@ -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(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -53,6 +56,7 @@ import {
type ProviderAdapterV2Event,
type ProviderAdapterV2HistoricalContext,
ProviderAdapterProtocolError,
ProviderAdapterTurnStartError,
type ProviderAdapterV2Shape,
type ProviderAdapterV2SessionRuntime,
} from "../ProviderAdapter.ts";
Expand Down Expand Up @@ -104,6 +108,8 @@ function makeTestAdapter(input: {
readonly capturedTurns: Ref.Ref<ReadonlyArray<CapturedTurn>>;
readonly injectedHistory?: Ref.Ref<ReadonlyArray<unknown>>;
readonly failStartOnce?: Ref.Ref<boolean>;
/** Starts left to refuse the way Claude does while background work runs. */
readonly refuseStarts?: Ref.Ref<number>;
readonly failInjectionOnce?: Ref.Ref<boolean>;
readonly nativeThreadGeneration?: Ref.Ref<number>;
readonly failResume?: boolean;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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<ReadonlyArray<CapturedTurn>>([]);
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) =>
Expand Down
Loading