Skip to content
Closed
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
Original file line number Diff line number Diff line change
Expand Up @@ -3542,12 +3542,15 @@ describe("ClaudeAdapterV2 background wake turns", () => {

// The Waiting strip's Stop reaches the adapter as an interrupt of the
// settled turn with requestRuntimeRestart.
yield* harness.runtime.interruptTurn({
const outcome = yield* harness.runtime.interruptTurn({
providerThread: settledThread ?? harness.providerThread,
providerTurnId: harness.terminalEvents()[0]!.providerTurnId,
requestRuntimeRestart: true,
});

// The turn's terminal went out when it settled; Stop emits no other.
assert.equal(outcome, "turn_not_active");
assert.lengthOf(harness.terminalEvents(), 1);
assert.equal(closes, 1, "Stop must close the CLI process that owns the task");
yield* awaitUntil(
() =>
Expand Down Expand Up @@ -3589,11 +3592,12 @@ describe("ClaudeAdapterV2 background wake turns", () => {
let quietYields = 0;
yield* awaitUntil(() => quietYields++ >= 50, "query exit");

yield* harness.runtime.interruptTurn({
const outcome = yield* harness.runtime.interruptTurn({
providerThread: settledThread ?? harness.providerThread,
providerTurnId: harness.terminalEvents()[0]!.providerTurnId,
requestRuntimeRestart: true,
});
assert.equal(outcome, "turn_not_active");
yield* awaitUntil(
() =>
(providerThreadRosterEvents(harness.events).at(-1)?.providerThread
Expand Down
9 changes: 5 additions & 4 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7243,14 +7243,15 @@ export function makeClaudeAdapterV2(
// 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;
// settles the items the thread still shows. No terminal follows:
// the turn's was emitted when it settled.
if (nativeThreadId === null) return "turn_not_active" as const;
if (existing === null || existing.nativeThreadId !== nativeThreadId) {
yield* clearWakeStateForNativeThread(nativeThreadId);
yield* resetBackgroundTaskStateForNativeThreadProcess(nativeThreadId, {
status: "idle",
});
return;
return "turn_not_active" as const;
}
// The background shells belong to the CLI process, so closing
// its query is what stops them.
Expand All @@ -7264,7 +7265,7 @@ export function makeClaudeAdapterV2(
status: "idle",
});
}
return;
return "turn_not_active" as const;
}
if (existing === null) {
return yield* new ProviderAdapter.ProviderAdapterProtocolError({
Expand Down
27 changes: 24 additions & 3 deletions apps/server/src/orchestration-v2/EventSink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -82,14 +82,21 @@ export interface EventSinkV2Shape {
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
readonly effects: ReadonlyArray<EffectOutbox.PendingOrchestrationEffectV2>;
}) => Effect.Effect<ReadonlyArray<OrchestrationV2StoredEvent>, EventSinkV2Error>;
/**
* Atomically commit only while the run's active attempt and status are the
* expected ones; effects are enqueued in the same transaction.
*/
readonly writeIfRunCurrent: (input: {
readonly guardPendingUserInputCancellations?: boolean;
readonly commandId?: CommandId;
readonly threadId: ThreadId;
readonly runId: RunId;
readonly activeAttemptId: RunAttemptId;
readonly expectedStatus: OrchestrationV2Run["status"];
readonly expectedStatus:
| OrchestrationV2Run["status"]
| ReadonlyArray<OrchestrationV2Run["status"]>;
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
readonly effects?: ReadonlyArray<EffectOutbox.PendingOrchestrationEffectV2>;
}) => Effect.Effect<
{
readonly committed: boolean;
Expand Down Expand Up @@ -416,9 +423,13 @@ const baseLayer: Layer.Layer<
LIMIT 1
`;
const current = rows[0];
const expectedStatuses: ReadonlyArray<string> =
typeof input.expectedStatus === "string"
? [input.expectedStatus]
: input.expectedStatus;
if (
current === undefined ||
current.status !== input.expectedStatus ||
!expectedStatuses.includes(current.status) ||
current.active_attempt_id !== input.activeAttemptId
) {
return {
Expand All @@ -437,9 +448,19 @@ const baseLayer: Layer.Layer<
events: normalized,
});
yield* applyStoredEvents(storedEvents);
yield* effectOutbox.enqueue(input.effects ?? []);
return { committed: true as const, storedEvents };
}),
(result) => (result.committed ? publishStoredEvents(result.storedEvents) : Effect.void),
(result) =>
result.committed
? Effect.gen(function* () {
const effectCount = input.effects?.length ?? 0;
if (effectCount > 0) {
yield* effectOutbox.notifyAvailable(effectCount);
}
yield* publishStoredEvents(result.storedEvents);
})
: Effect.void,
);
},
);
Expand Down
13 changes: 12 additions & 1 deletion apps/server/src/orchestration-v2/ProjectionControlReads.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,10 @@ import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as SqlClient from "effect/unstable/sql/SqlClient";
import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts";
import * as EventSink from "./EventSink.ts";
import * as IdAllocator from "./IdAllocator.ts";
import * as ProjectionStore from "./ProjectionStore.ts";
import * as ProviderEventIngestor from "./ProviderEventIngestor.ts";
import * as ProviderSessionManager from "./ProviderSessionManager.ts";
import * as ProviderTurnControlService from "./ProviderTurnControlService.ts";
import * as RuntimeRequestService from "./RuntimeRequestService.ts";
Expand Down Expand Up @@ -346,7 +349,15 @@ it.effect.each(storageCases)(
}).pipe(
Effect.provide(
Layer.merge(ProviderTurnControlService.layer, RuntimeRequestService.layer).pipe(
Layer.provide(sessions),
Layer.provide(
Layer.mergeAll(
sessions,
IdAllocator.layer,
// No Stop here ends an orphaned run, which is what writes.
Layer.mock(EventSink.EventSinkV2)({}),
Layer.mock(ProviderEventIngestor.ProviderEventIngestorV2)({}),
),
),
),
),
);
Expand Down
10 changes: 9 additions & 1 deletion apps/server/src/orchestration-v2/ProviderAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -423,6 +423,14 @@ export interface ProviderAdapterV2InterruptInput {
readonly requestRuntimeRestart?: boolean;
}

/**
* What an interrupt found, for adapters that can tell. `turn_not_active`: the
* adapter held no live turn with this id and emitted nothing for it, so any
* terminal it will ever emit for that turn was emitted before the interrupt.
* Adapters that cannot tell return nothing.
*/
export type ProviderAdapterV2InterruptOutcome = "turn_not_active";

export interface ProviderAdapterV2RuntimeRequestResponseInput {
readonly requestId: RuntimeRequestId;
readonly decision?: ProviderApprovalDecision;
Expand Down Expand Up @@ -547,7 +555,7 @@ export interface ProviderAdapterV2SessionRuntime {
) => Effect.Effect<void, ProviderAdapterV2Error>;
readonly interruptTurn: (
input: ProviderAdapterV2InterruptInput,
) => Effect.Effect<void, ProviderAdapterV2Error>;
) => Effect.Effect<ProviderAdapterV2InterruptOutcome | void, ProviderAdapterV2Error>;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟠 High orchestration-v2/ProviderAdapter.ts:558

Codex interruptTurn returns undefined when no active or settled context is retained, so ProviderTurnControlService.interrupt leaves turnOrphaned false and thread.background-work.settle cannot terminalize the still-running run. Return "turn_not_active" from the Codex Stop path as Claude does, or otherwise propagate an equivalent orphan signal.

Also found in 1 other location(s)

apps/server/src/orchestration-v2/ProviderTurnControlService.ts:229

The new orphan detection only waits/marks the turn when interruptTurn returns &#34;turn_not_active&#34;. Codex's interruptTurn returns undefined when it has no retained active turn for a Stop (CodexAdapterV2.ts lines 5664–5668, including its settled/released-turn path), so this branch returns turnOrphaned: false immediately. Consequently thread.background-work.settle is dispatched without providerTurnOrphaned and cannot terminalize the still-running Codex run—the same permanently “Thinking” state this change is intended to recover.

🤖 Copy this AI Prompt to have your agent fix this:
In file @apps/server/src/orchestration-v2/ProviderAdapter.ts around line 558:

Codex `interruptTurn` returns `undefined` when no active or settled context is retained, so `ProviderTurnControlService.interrupt` leaves `turnOrphaned` false and `thread.background-work.settle` cannot terminalize the still-running run. Return `"turn_not_active"` from the Codex Stop path as Claude does, or otherwise propagate an equivalent orphan signal.

Also found in 1 other location(s):
- apps/server/src/orchestration-v2/ProviderTurnControlService.ts:229 -- The new orphan detection only waits/marks the turn when `interruptTurn` returns `"turn_not_active"`. Codex's `interruptTurn` returns `undefined` when it has no retained active turn for a Stop (`CodexAdapterV2.ts` lines 5664–5668, including its settled/released-turn path), so this branch returns `turnOrphaned: false` immediately. Consequently `thread.background-work.settle` is dispatched without `providerTurnOrphaned` and cannot terminalize the still-running Codex run—the same permanently “Thinking” state this change is intended to recover.

/**
* Lets a runtime shared by several app threads unload one provider thread's
* native state (and its MCP servers) when that app thread detaches, while
Expand Down
Loading
Loading