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
2 changes: 1 addition & 1 deletion apps/server/src/claudeModelOptions.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ export function compileClaudeModelSelection(
const resolvedEffort = resolveClaudeCatalogEffort(catalog, selection.model, rawEffort);
const effort = normalizeClaudeCatalogEffort(catalog, resolvedEffort, selection.model);
const fastMode = supportsBoolean("fastMode")
? getModelSelectionBooleanOptionValue(selection, "fastMode")
? (getModelSelectionBooleanOptionValue(selection, "fastMode") ?? false)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
: undefined;
const thinking = supportsBoolean("thinking")
? getModelSelectionBooleanOptionValue(selection, "thinking")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,7 @@ function makeClaudeTestTurnInput(input: {
readonly text: string;
readonly attachments: ProviderAdapterV2TurnInput["message"]["attachments"];
readonly providerTurnOrdinal?: number;
readonly nativeThreadHasTurns?: boolean;
readonly messageCreatedBy?: ProviderAdapterV2TurnInput["message"]["createdBy"];
readonly messageCreationSource?: ProviderAdapterV2TurnInput["message"]["creationSource"];
readonly modelSelection?: ModelSelection;
Expand All @@ -137,6 +138,9 @@ function makeClaudeTestTurnInput(input: {
runId: RunId.make(`run-${input.attemptId}`),
runOrdinal: 1,
providerTurnOrdinal: input.providerTurnOrdinal ?? 1,
...(input.nativeThreadHasTurns === undefined
? {}
: { nativeThreadHasTurns: input.nativeThreadHasTurns }),
attemptId: input.attemptId,
rootNodeId: NodeId.make(`node-${input.attemptId}`),
providerThread: input.providerThread,
Expand Down Expand Up @@ -1783,7 +1787,7 @@ describe("ClaudeAdapterV2 native fork", () => {
});

describe("ClaudeAdapterV2 native session identity", () => {
const openTurnWithOrdinal = (providerTurnOrdinal: number) =>
const openTurnWithOrdinal = (providerTurnOrdinal: number, nativeThreadHasTurns?: boolean) =>
Effect.scoped(
Effect.gen(function* () {
const fileSystem = yield* FileSystem.FileSystem;
Expand Down Expand Up @@ -1842,6 +1846,7 @@ describe("ClaudeAdapterV2 native session identity", () => {
text: "Respond with identity ok",
attachments: [],
providerTurnOrdinal,
...(nativeThreadHasTurns === undefined ? {} : { nativeThreadHasTurns }),
}),
);
return openedQueries;
Expand All @@ -1867,6 +1872,14 @@ describe("ClaudeAdapterV2 native session identity", () => {
assert.equal(openedQueries[0]?.options.sessionId, undefined);
}),
);

it.effect("creates a fresh native session despite earlier provider-thread turns", () =>
Effect.gen(function* () {
const openedQueries = yield* openTurnWithOrdinal(4, false);
assert.equal(openedQueries[0]?.options.sessionId, "native-session-identity");
assert.equal(openedQueries[0]?.options.resume, undefined);
}),
);
});

const encodeJsonString = Schema.encodeSync(Schema.fromJsonString(Schema.Unknown));
Expand Down Expand Up @@ -2182,6 +2195,68 @@ describe("ClaudeAdapterV2 background wake turns", () => {
});
const makeWakeHarness = makeWakeHarnessWithOptions();

it.effect(
"reuses a background shell's query for omitted and explicit Normal, but blocks Fast",
() =>
Effect.scoped(
Effect.gen(function* () {
const harness = yield* makeWakeHarness;
const now = yield* DateTime.now;
const normal = {
instanceId: ClaudeAdapterV2.CLAUDE_DEFAULT_INSTANCE_ID,
model: "claude-opus-5-5",
} satisfies ModelSelection;
const turn = (ordinal: number, modelSelection: ModelSelection) =>
makeClaudeTestTurnInput({
threadId: harness.threadId,
providerThread: harness.providerThread,
now,
attemptId: RunAttemptId.make(`attempt-normal-background:${ordinal}`),
text: `Request ${ordinal}`,
attachments: [],
providerTurnOrdinal: ordinal,
modelSelection,
});
yield* harness.runtime.startTurn(turn(1, normal));
const originalOptions = harness.getOpenedOptions();
yield* harness.offerAndWait(wakeTaskStarted);
yield* harness.offerAndWait(turnOneResult);
yield* Queue.take(harness.terminalReceipts);
assert.isTrue(yield* harness.hasPendingBackgroundWork);

yield* harness.runtime.startTurn(
turn(2, {
...normal,
options: [{ id: "fastMode", value: false }],
}),
);
assert.strictEqual(harness.getOpenedOptions(), originalOptions);
assert.lengthOf(harness.offeredMessages, 2);
yield* harness.offerAndWait(turnOneResult);
yield* Queue.take(harness.terminalReceipts);

const refused = yield* harness.runtime
.startTurn(
turn(3, {
...normal,
options: [{ id: "fastMode", value: true }],
}),
)
.pipe(Effect.result);
assert.equal(refused._tag, "Failure");
if (refused._tag === "Failure") {
assert.instanceOf(
refused.failure.cause,
ClaudeAdapterV2.ClaudeBackgroundWorkBlocksQueryReplacementError,
);
}
assert.strictEqual(harness.getOpenedOptions(), originalOptions);
assert.lengthOf(harness.offeredMessages, 2);
assert.isTrue(yield* harness.hasPendingBackgroundWork);
}).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))),
),
);

it.effect.each(["completed", "interrupted"] as const)(
"projects Claude thinking blocks when %s",
(status) =>
Expand Down
11 changes: 6 additions & 5 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6969,11 +6969,12 @@ export function makeClaudeAdapterV2(

const openedWithResume = (yield* Ref.get(openedNativeThreads)).has(nativeThreadId);
// openedNativeThreads is per session instance and is lost when the
// provider session is idle-released. A prior persisted provider turn
// proves the native session already exists, so the query must resume
// it; reopening with a fixed session id makes the CLI fail fast with
// "Session ID ... is already in use".
const hasPersistedProviderTurn = turnInput.providerTurnOrdinal > 1;
// provider session is idle-released. A prior turn on this native id
// requires resume; sessionId would fail with "already in use".
// A fresh-session fallback keeps provider-thread history but binds
// a new native id, which must be created before it can be resumed.
const hasPersistedProviderTurn =
turnInput.nativeThreadHasTurns ?? turnInput.providerTurnOrdinal > 1;
const shouldResume =
resumeSessionAt !== undefined || openedWithResume || hasPersistedProviderTurn;
const queryOptions = makeClaudeQueryOptions({
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ const modelSelection = { instanceId, model: "test-model" };
// Codex turns leave commands running, then the thread moves to another
// provider thread (a provider switch). Stop on the newer, settled run must
// reach both provider threads and end all of the Codex work.
it.effect("Stop reaches background work an earlier provider thread still runs", () =>
const stopEarlierBackgroundWork = (failedStart: boolean) =>
Effect.scoped(
Effect.gen(function* () {
const cwd = yield* checkpointWorkspace("background-work-stop");
Expand Down Expand Up @@ -233,7 +233,10 @@ it.effect("Stop reaches background work an earlier provider thread still runs",
const settledRun = (input: {
readonly ordinal: number;
readonly providerThreadId: ProviderThreadId;
readonly runningItem?: { readonly id: TurnItemId; readonly kind: "command" | "subagent" };
readonly runningItem?: {
readonly id: TurnItemId;
readonly kind: "command" | "subagent";
};
}) => {
const runId = RunId.make(`run:${input.ordinal}`);
const attemptId = RunAttemptId.make(`attempt:${input.ordinal}`);
Expand Down Expand Up @@ -407,11 +410,41 @@ it.effect("Stop reaches background work an earlier provider thread still runs",
],
});

const failedRun = settledRun({ ordinal: 5, providerThreadId: otherProviderThreadId });
if (failedStart) {
yield* sink.write({
events: failedRun.events.flatMap((event): Array<OrchestrationV2DomainEvent> => {
switch (event.type) {
case "provider-turn.updated":
return [];
case "run.created":
return [{ ...event, payload: { ...event.payload, status: "failed" } }];
case "node.updated":
return [
{
...event,
payload: { ...event.payload, status: "failed", providerTurnId: null },
},
];
case "run-attempt.created":
return [
{
...event,
payload: { ...event.payload, status: "failed", providerTurnId: null },
},
];
default:
return [event];
}
}),
});
}

yield* orchestrator.dispatch({
type: "run.interrupt",
commandId: CommandId.make("stop-background-work"),
threadId,
runId: latestRun.runId,
runId: failedStart ? failedRun.runId : latestRun.runId,
});
yield* worker.drain();

Expand Down Expand Up @@ -442,5 +475,13 @@ it.effect("Stop reaches background work an earlier provider thread still runs",
),
);
}),
),
);

it.effect("Stop reaches background work an earlier provider thread still runs", () =>
stopEarlierBackgroundWork(false),
);

it.effect(
"Stop reaches earlier background work after the newest run fails before provider start",
() => stopEarlierBackgroundWork(true),
);
18 changes: 13 additions & 5 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8040,11 +8040,19 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
activeProviderThreadId: projection.thread.activeProviderThreadId,
runs: projection.runs,
}).length > 0;
const providerTurn = projection.providerTurns.findLast(
(candidate) =>
candidate.runAttemptId === run?.activeAttemptId &&
(candidate.status === "running" || hasBackgroundWork),
);
// A failed start has no provider turn. Background work still belongs
// to the provider thread, so Stop reaches its latest accepted turn.
const providerTurn =
projection.providerTurns.findLast(
(candidate) =>
candidate.runAttemptId === run?.activeAttemptId &&
(candidate.status === "running" || hasBackgroundWork),
) ??
(hasBackgroundWork
? projection.providerTurns.findLast(
(candidate) => candidate.providerThreadId === run?.providerThreadId,
)
: undefined);
if (run === undefined || rootNode === undefined || providerThread === undefined) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/orchestration-v2/ProviderAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -398,6 +398,8 @@ export interface ProviderAdapterV2TurnInput {
readonly threadId: ThreadId;
readonly runId: RunId;
readonly runOrdinal: number;
/** Whether the current native session has an accepted turn; omitted when unknown. */
readonly nativeThreadHasTurns?: boolean;
readonly providerTurnOrdinal: number;
readonly restartContinuationOfRunId?: RunId;
readonly attemptId: RunAttemptId;
Expand Down
7 changes: 7 additions & 0 deletions apps/server/src/orchestration-v2/ProviderTurnStartService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1231,6 +1231,13 @@ export const layer: Layer.Layer<
.filter((turn) => turn.providerThreadId === providerThread.id)
.map((turn) => turn.ordinal),
) + 1,
// Legacy accepted attempts have no native id. They count only before
// a replacement, while no accepted attempt records a native identity.
nativeThreadHasTurns:
nativeInputRunIds.size > 0 ||
(legacyInputRunIds.size > 0 &&
sameNativeThread &&
!acceptedAttempts.some((source) => source.nativeThreadId !== undefined)),
shouldStartProviderTurn: runControls.shouldStartProviderTurn,
shouldFinalizeRun: runControls.shouldFinalizeRun,
hasUnpairedRunInterruptRequest: runControls.hasUnpairedRunInterruptRequest,
Expand Down
4 changes: 4 additions & 0 deletions apps/server/src/orchestration-v2/RunExecutionService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -507,6 +507,7 @@ export interface RunExecutionServiceV2StartRootRunInput {
readonly attempt: OrchestrationV2RunAttempt;
readonly attemptId: RunAttemptId;
readonly providerTurnOrdinal: number;
readonly nativeThreadHasTurns?: boolean;
readonly loadInheritedBackgroundTurnItems?: () => Effect.Effect<
ReadonlyArray<InheritedBackgroundTurnItemRoute>,
unknown
Expand Down Expand Up @@ -1348,6 +1349,9 @@ export const layer: Layer.Layer<
runId: input.run.id,
runOrdinal: input.run.ordinal,
providerTurnOrdinal: input.providerTurnOrdinal,
...(input.nativeThreadHasTurns === undefined
? {}
: { nativeThreadHasTurns: input.nativeThreadHasTurns }),
...(input.run.restartContinuationOfRunId === undefined
? {}
: {
Expand Down
Loading
Loading