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
45 changes: 45 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3492,6 +3492,51 @@ describe("ClaudeAdapterV2 background wake turns", () => {
),
);

// The CLI process exits on its own after the turn settled (idle, crash),
// leaving its background task on the roster. Stop must succeed so the
// orchestrator goes on to settle what the thread still shows.
it.effect("a settled Stop with no CLI process left succeeds and clears the roster", () =>
Effect.scoped(
Effect.gen(function* () {
const harness = yield* makeWakeHarnessWithOptions();
const now = yield* DateTime.now;
yield* harness.runtime.startTurn(
makeClaudeTestTurnInput({
threadId: harness.threadId,
providerThread: harness.providerThread,
now,
attemptId: RunAttemptId.make("attempt-claude-settled-stop-no-process"),
text: "Run the build in the background.",
attachments: [],
}),
);
yield* Queue.offer(harness.sdkMessages, wakeTaskStarted);
yield* Queue.offer(harness.sdkMessages, turnOneResult);
yield* awaitUntil(() => harness.terminalEvents().length === 1, "first turn terminal");
const settledThread = providerThreadRosterEvents(harness.events).at(-1)?.providerThread;
assert.equal(settledThread?.pendingBackgroundTasks?.[0]?.taskId, WAKE_TASK_ID);

yield* Queue.shutdown(harness.sdkMessages);
let quietYields = 0;
yield* awaitUntil(() => quietYields++ >= 50, "query exit");

yield* harness.runtime.interruptTurn({
providerThread: settledThread ?? harness.providerThread,
providerTurnId: harness.terminalEvents()[0]!.providerTurnId,
requestRuntimeRestart: true,
});
yield* awaitUntil(
() =>
(providerThreadRosterEvents(harness.events).at(-1)?.providerThread
.pendingBackgroundTasks?.length ?? 0) === 0,
"roster clear after Stop",
);
assert.isFalse(yield* harness.hasPendingBackgroundWork);
assert.lengthOf(harness.continuationRequests, 0);
}).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))),
),
);

it.effect("a settled Stop leaves a turn that replaced the closing process alone", () =>
Effect.scoped(
Effect.gen(function* () {
Expand Down
35 changes: 21 additions & 14 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6940,22 +6940,23 @@ export function makeClaudeAdapterV2(
const interruptTurn = Effect.fn("ClaudeAdapterV2.interruptTurn")(
function* (turnInput: ProviderAdapter.ProviderAdapterV2InterruptInput) {
const existing = yield* Ref.get(queryContext);
if (existing === null) {
return yield* new ProviderAdapter.ProviderAdapterProtocolError({
driver: CLAUDE_PROVIDER,
detail: `Claude provider thread ${turnInput.providerThread.id} has no live query.`,
});
}
const currentTurn = yield* Ref.get(activeTurn);
const nativeThreadId = turnInput.providerThread.nativeThreadRef?.nativeId ?? null;
if (
currentTurn === null &&
turnInput.requestRuntimeRestart === true &&
nativeThreadId !== null &&
existing.nativeThreadId === nativeThreadId
) {
// Stop after the turn settled: the background shells belong to
// the CLI process, so closing its query is what stops them.
if (currentTurn === null && turnInput.requestRuntimeRestart === true) {
// Stop after the turn settled. With no CLI process of this
// native thread left, nothing it started is still running: its
// roster is not authoritative any more, and the orchestrator
// settles the items the thread still shows.
if (nativeThreadId === null) return;
if (existing === null || existing.nativeThreadId !== nativeThreadId) {
yield* clearWakeStateForNativeThread(nativeThreadId);
yield* resetBackgroundTaskStateForNativeThreadProcess(nativeThreadId, {
status: "idle",
});
return;
}
// The background shells belong to the CLI process, so closing
// its query is what stops them.
yield* closeLiveQueryForNativeThread(nativeThreadId);
// A turn started while the close was pending may have opened a
// replacement process. Its Waiting and wake state are its own.
Expand All @@ -6968,6 +6969,12 @@ export function makeClaudeAdapterV2(
}
return;
}
if (existing === null) {
return yield* new ProviderAdapter.ProviderAdapterProtocolError({
driver: CLAUDE_PROVIDER,
detail: `Claude provider thread ${turnInput.providerThread.id} has no live query.`,
});
}
if (currentTurn?.providerTurnId !== turnInput.providerTurnId) {
return yield* new ProviderAdapter.ProviderAdapterProtocolError({
driver: CLAUDE_PROVIDER,
Expand Down
120 changes: 120 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3987,6 +3987,126 @@ describe("CodexAdapterV2 post-settle continuation", () => {
}
}

// The app-server exits after the root turn, before the command's own
// item/completed (Codex always sends one, so only a lost notification or a
// gone process leaves it running). Nothing tracks the command any more, yet
// the thread still shows it, and Stop is the only way to clear it.
it.effect("Stop ends a background command no Codex process tracks any more", () =>
Effect.scoped(
Effect.gen(function* () {
const fs = yield* FileSystem.FileSystem;
const cwd = yield* fs.makeTempDirectoryScoped({ prefix: "t3-bg-stale-workspace-" });
const staleTranscript = makeCodexReplayTranscript({
scenario: "codex-bg-stop-untracked",
entries: [
...backgroundExecTranscript.entries.slice(0, -1),
{ type: "runtime_exit", status: "success" },
],
});
const localTranscript = yield* decodeReplayTranscriptJson(
(yield* encodeReplayTranscriptJson(staleTranscript)).replaceAll(
yield* encodeStringJson("/workspace"),
yield* encodeStringJson(cwd),
),
);
const spawner = yield* ChildProcessSpawner.ChildProcessSpawner;
for (const args of [
["init", "--quiet"],
[
"-c",
"user.name=Test",
"-c",
"user.email=test@example.com",
"commit",
"--allow-empty",
"--quiet",
"-m",
"Initial commit",
],
]) {
assert.equal(Number(yield* spawner.exitCode(ChildProcess.make("git", args, { cwd }))), 0);
}
yield* Effect.gen(function* () {
const orchestrator = yield* Orchestrator.OrchestratorV2;
const worker = yield* EffectWorker.OrchestrationEffectWorkerV2;
const threadId = ThreadId.make("thread:background-stop-untracked");
yield* orchestrator.dispatch({
type: "thread.create",
commandId: CommandId.make("create-background-stop-untracked"),
threadId,
projectId: ProjectId.make("project:background-stop-untracked"),
title: "Background stop untracked",
modelSelection: CODEX_TEST_MODEL_SELECTION,
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: cwd,
createdBy: "user",
creationSource: "web",
});
const settled = yield* orchestrator.streamDomainEvents.pipe(
Stream.filter(
(event) =>
event.type === "run.updated" &&
(event.payload.status === "waiting" || event.payload.status === "completed"),
),
Stream.runHead,
Effect.forkChild({ startImmediately: true }),
);
yield* orchestrator.dispatch({
type: "message.dispatch",
commandId: CommandId.make("start-background-stop-untracked"),
threadId,
messageId: MessageId.make("message:background-stop-untracked"),
text: BG_PROMPT,
attachments: [],
createdBy: "user",
creationSource: "web",
dispatchMode: { type: "start_immediately" },
});
yield* worker.drain();
yield* Fiber.join(settled);
yield* worker.drain();
const before = yield* orchestrator.getThreadShell(threadId);
assert.deepEqual(
before?.pendingBackgroundTasks?.map((task) => task.kind),
["command"],
"the thread still shows the command the gone process never finished",
);
const run = (yield* orchestrator.getThreadProjection(threadId)).runs.at(-1)!;
yield* orchestrator.dispatch({
type: "run.interrupt",
commandId: CommandId.make("stop-background-untracked"),
threadId,
runId: run.id,
holdQueue: true,
});
yield* worker.drain();
const projection = yield* orchestrator.getThreadProjection(threadId);
assert.deepEqual(
projection.turnItems.flatMap((item) =>
item.type === "command_execution" ? [item.status] : [],
),
["interrupted"],
);
assert.equal(projection.runs.at(-1)?.status, "completed");
assert.deepEqual(
(yield* orchestrator.getThreadShell(threadId))?.pendingBackgroundTasks,
[],
);
}).pipe(
Effect.provide(
makeOrchestratorV2ReplayLayerWithRegistry(
{ name: "codex-background-stop-untracked", runtimePolicyOverride: { cwd } },
makeCodexProviderAdapterRegistryReplayLayer({ transcript: localTranscript }),
{ runEffectWorker: false },
),
),
);
}).pipe(Effect.provide(NodeServices.layer)),
),
);

const PRE_SETTLE_SCENARIO = "codex-bg-exec-pre-settle";
const PRE_SETTLE_NATIVE_THREAD = "native-codex-pre-settle-thread";
const PRE_SETTLE_NATIVE_TURN = "native-codex-pre-settle-turn";
Expand Down
20 changes: 16 additions & 4 deletions apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2042,10 +2042,19 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi

const terminateBackgroundTerminal = Effect.fn("CodexAdapterV2.terminateBackgroundTerminal")(
function* (nativeThreadId: string, processId: string) {
const response = yield* client.raw.request("thread/backgroundTerminals/terminate", {
threadId: nativeThreadId,
processId,
});
const response = yield* client.raw
.request("thread/backgroundTerminals/terminate", {
threadId: nativeThreadId,
processId,
})
.pipe(
// The app-server that ran the terminal is gone, and with it
// the only handle to the terminal: nothing is left to stop.
Effect.catchTags({
CodexAppServerProcessExitedError: () => Effect.succeed({ terminated: true }),
CodexAppServerInputStreamEndedError: () => Effect.succeed({ terminated: true }),
}),
);
const result = yield* decodeCodexBackgroundTerminalTerminateResponse(response);
if (result.terminated) return;
let cursor: string | null = null;
Expand Down Expand Up @@ -5653,6 +5662,9 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi
)
: undefined);
if (activeTurn === undefined) {
// Stop on a settled turn this process retains nothing for
// (released, restarted, or every command already reported).
if (turnInput.requestRuntimeRestart === true) return;
return yield* toProtocolError(
`Provider turn ${turnInput.providerTurnId} is not active and cannot be interrupted.`,
);
Expand Down
3 changes: 3 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/CursorAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2429,6 +2429,9 @@ export function makeCursorAdapterV2(
interruptTurn: Effect.fn("CursorAdapterV2.interruptTurn")(
function* (turnInput: ProviderAdapter.ProviderAdapterV2InterruptInput) {
const context = yield* Ref.get(activeTurn);
// Stop on a settled turn: finalization already ended its tools
// and subagents, so nothing of it is left running to stop.
if (context === null && turnInput.requestRuntimeRestart === true) return;
if (context?.providerTurnId !== turnInput.providerTurnId) {
return yield* new ProviderAdapter.ProviderAdapterProtocolError({
driver: CursorAgentSdk.CURSOR_PROVIDER,
Expand Down
11 changes: 11 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3385,6 +3385,17 @@ export function makeOpenCodeAdapterV2(
const sessionId = nativeThreadId(interruptInput.providerThread);
const state = threads.get(sessionId);
const turn = state?.activeTurn;
// Stop on a settled turn: a background task child can outlive
// it (busySessionIds), and aborting the descendants ends it.
if (
(turn === undefined || turn === null) &&
interruptInput.requestRuntimeRestart === true
) {
for (const controller of commandControllers.get(sessionId) ?? [])
controller.abort();
yield* abortDescendants(sessionId);
return;
}
if (
turn === undefined ||
turn === null ||
Expand Down
3 changes: 3 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/PiAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2466,6 +2466,9 @@ export function makePiAdapterV2(
interruptTurn: (interruptInput) =>
Effect.gen(function* () {
const turn = threadState?.activeTurn ?? null;
// Stop on a settled turn: Pi runs nothing between prompts, so
// nothing of that turn is left to stop.
if (turn === null && interruptInput.requestRuntimeRestart === true) return;
if (turn === null || turn.providerTurn.id !== interruptInput.providerTurnId) {
return yield* protocolError(`Pi turn ${interruptInput.providerTurnId} is not active`);
}
Expand Down
12 changes: 12 additions & 0 deletions apps/server/src/orchestration-v2/EffectWorker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -166,6 +166,18 @@ export const executorLayer: Layer.Layer<
providerTurnId: effect.request.providerTurnId,
})
.pipe(
// The provider has stopped what it still ran and reported it.
// Whatever the thread still shows on that provider thread is
// work no process will report on, so the Stop ends it too.
Effect.andThen(
threads.dispatch({
type: "thread.background-work.settle",
commandId: CommandId.make(`${effect.commandId}:background-work-settled`),
threadId: effect.threadId,
providerThreadId: effect.request.providerThreadId,
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
providerTurnId: effect.request.providerTurnId,
}),
),
Effect.mapError(
(cause) =>
new OrchestrationEffectExecutionError({
Expand Down
Loading
Loading