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
7 changes: 7 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,10 @@
## 3.5.38 - 2026-08-09 (Patch)

Release impact: Patch because this fixes provider turn attribution and queued-follow-up interruption without changing provider configuration.

- Claude no longer emits phantom turn completions for resume handshakes or late results, and ingestion rejects untargeted completions when no active turn can own them.
- Codex keeps the currently interruptible turn id when the app-server accepts a queued follow-up, then advances through normal lifecycle notifications.

## 3.5.37 - 2026-08-09 (Patch)

Release impact: Patch because this makes the existing Symphony review path fail closed for unsupported providers without changing review configuration.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1146,6 +1146,79 @@ describe("ProviderRuntimeIngestion", () => {
);
});

it("rejects an untargeted turn.completed when no turn is active", async () => {
const harness = await createHarness();
const seededAt = "2026-01-01T00:00:00.000Z";

await Effect.runPromise(
harness.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("cmd-session-seed-untargeted-completion"),
threadId: ThreadId.make("thread-1"),
session: {
threadId: ThreadId.make("thread-1"),
status: "starting",
providerName: "claudeAgent",
runtimeMode: "approval-required",
activeTurnId: null,
updatedAt: seededAt,
lastError: null,
},
createdAt: seededAt,
}),
);

harness.emit({
type: "turn.completed",
eventId: asEventId("evt-turn-completed-untargeted"),
provider: ProviderDriverKind.make("claudeAgent"),
createdAt: seededAt,
threadId: asThreadId("thread-1"),
status: "completed",
});

await harness.drain();
const readModel = await harness.readModel();
const thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1"));
expect(thread?.session?.status).toBe("starting");
expect(thread?.session?.activeTurnId).toBeNull();
});

it("accepts a targeted turn.completed when no turn is active", async () => {
const harness = await createHarness();
const seededAt = "2026-01-01T00:00:00.000Z";

await Effect.runPromise(
harness.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("cmd-session-seed-targeted-completion"),
threadId: ThreadId.make("thread-1"),
session: {
threadId: ThreadId.make("thread-1"),
status: "starting",
providerName: "claudeAgent",
runtimeMode: "approval-required",
activeTurnId: null,
updatedAt: seededAt,
lastError: null,
},
createdAt: seededAt,
}),
);

harness.emit({
type: "turn.completed",
eventId: asEventId("evt-turn-completed-targeted-late"),
provider: ProviderDriverKind.make("claudeAgent"),
createdAt: seededAt,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-late"),
status: "completed",
});

await waitForThread(harness.readModel, (thread) => thread.session?.status === "ready");
});

it("ignores non-active turn completion when runtime omits thread id", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
Expand Down
10 changes: 8 additions & 2 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1584,8 +1584,14 @@ const make = Effect.gen(function* () {
if (activeTurnId !== null && eventTurnId !== undefined) {
return sameId(activeTurnId, eventTurnId);
}
// If no active turn is tracked, accept completion scoped to this thread.
return true;
// No active turn tracked: accept only completions that name their
// turn (covers a real completion whose turn.started was lost). An
// untargeted completion cannot prove it belongs to any turn this
// thread ran — the known emitter was the Claude resume handshake
// (system/init + result(num_turns: 0)), which is not a turn at
// all — and applying it here stomps the "starting" lifecycle
// state while a turn start is pending.
return eventTurnId !== undefined;
default:
return true;
}
Expand Down
56 changes: 56 additions & 0 deletions apps/server/src/provider/Layers/ClaudeAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1045,6 +1045,62 @@ describe("ClaudeAdapterLive", () => {
);
});

it.effect("does not emit turn.completed for a result with no active turn", () => {
const harness = makeHarness();
return Effect.gen(function* () {
const adapter = yield* ClaudeAdapter;
const runtimeEventsFiber = yield* adapter.streamEvents.pipe(
Stream.takeUntil((event) => event.type === "session.exited"),
Stream.runCollect,
Effect.forkChild,
);

const session = yield* adapter.startSession({
threadId: THREAD_ID,
provider: ProviderDriverKind.make("claudeAgent"),
runtimeMode: "full-access",
});
const turn = yield* adapter.sendTurn({
threadId: session.threadId,
input: "hello",
attachments: [],
});

harness.query.emit({
type: "result",
subtype: "success",
is_error: false,
errors: [],
num_turns: 1,
session_id: "sdk-session-1",
uuid: "result-real",
} as unknown as SDKMessage);
harness.query.emit({
type: "result",
subtype: "success",
is_error: false,
errors: [],
num_turns: 0,
usage: { input_tokens: 0, output_tokens: 0 },
session_id: "sdk-session-1",
uuid: "result-handshake",
} as unknown as SDKMessage);
harness.query.finish();

const completions = Array.from(yield* Fiber.join(runtimeEventsFiber)).filter(
(event) => event.type === "turn.completed",
);
assert.equal(completions.length, 1);
const completed = completions[0];
if (completed?.type === "turn.completed") {
assert.equal(String(completed.turnId), String(turn.turnId));
}
}).pipe(
Effect.provideService(Random.Random, makeDeterministicRandomService()),
Effect.provide(harness.layer),
);
});

it.effect("steers a running turn instead of opening a new one on mid-turn sendTurn", () => {
const harness = makeHarness();
return Effect.gen(function* () {
Expand Down
34 changes: 17 additions & 17 deletions apps/server/src/provider/Layers/ClaudeAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2075,24 +2075,24 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
rawPayload: result ?? { status },
});

const stamp = yield* makeEventStamp();
yield* offerRuntimeEvent({
type: "turn.completed",
eventId: stamp.eventId,
provider: PROVIDER,
createdAt: stamp.createdAt,
// A result with no local turn is never a turn this adapter started:
// real turns get turnState in sendTurn, and assistant messages that
// arrive outside a turn auto-start a synthetic one. What lands here is
// the resume handshake (system/init + result(num_turns: 0)), a late
// result for a turn already completed locally (steer auto-close,
// stream teardown), or a stream failure with no turn in flight. The
// untargeted turn.completed this branch used to emit carried no turnId,
// so ingestion could not attribute it — and whenever the projection had
// no active turn (a pending turn start included) it flipped the session
// lifecycle for a turn that never existed. Keep the usage emission
// above, drop the lifecycle event, and leave a tripwire so the trigger
// stays measurable in the field.
yield* Effect.logInfo("claude.turn.result-without-active-turn", {
threadId: context.session.threadId,
payload: {
state: status,
...(result?.stop_reason !== undefined ? { stopReason: result.stop_reason } : {}),
...(result?.usage ? { usage: result.usage } : {}),
...(result?.modelUsage ? { modelUsage: result.modelUsage } : {}),
...(typeof result?.total_cost_usd === "number"
? { totalCostUsd: result.total_cost_usd }
: {}),
...(errorMessage ? { errorMessage } : {}),
},
providerRefs: {},
status,
numTurns: result?.num_turns,
hasUsage: result?.usage !== undefined,
...(errorMessage ? { errorMessage } : {}),
});
return;
}
Expand Down
13 changes: 12 additions & 1 deletion apps/server/src/provider/Layers/CodexSessionRuntime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ import { it } from "@effect/vitest";
import * as Effect from "effect/Effect";
import * as Schema from "effect/Schema";
import { describe } from "vite-plus/test";
import { ThreadId } from "@neokod/contracts";
import { ThreadId, TurnId } from "@neokod/contracts";
import * as CodexErrors from "effect-codex-app-server/errors";
import * as CodexRpc from "effect-codex-app-server/rpc";

Expand All @@ -13,6 +13,7 @@ import {
CODEX_PLAN_MODE_DEVELOPER_INSTRUCTIONS,
} from "../CodexDeveloperInstructions.ts";
import {
activeTurnIdAfterStart,
buildTurnStartParams,
hasConfiguredMcpServer,
isRecoverableThreadResumeError,
Expand All @@ -37,6 +38,16 @@ describe("CodexSessionRuntimeIdentifierGenerationError", () => {
});
});

describe("activeTurnIdAfterStart", () => {
it("keeps the active id when Codex accepts a queued follow-up", () => {
const active = TurnId.make("turn-active");
const queued = TurnId.make("turn-queued");

NodeAssert.equal(activeTurnIdAfterStart(active, queued), active);
NodeAssert.equal(activeTurnIdAfterStart(undefined, queued), queued);
});
});

function makeThreadOpenResponse(
threadId: string,
): CodexRpc.ClientRequestResponsesByMethod["thread/start"] {
Expand Down
18 changes: 13 additions & 5 deletions apps/server/src/provider/Layers/CodexSessionRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -680,18 +680,21 @@ function currentProviderThreadId(session: ProviderSession): string | undefined {

function updateSession(
sessionRef: Ref.Ref<ProviderSession>,
updates: Partial<ProviderSession>,
updates: Partial<ProviderSession> | ((session: ProviderSession) => Partial<ProviderSession>),
): Effect.Effect<void> {
return Effect.gen(function* () {
const updatedAt = DateTime.formatIso(yield* DateTime.now);
yield* Ref.update(sessionRef, (session) => ({
...session,
...updates,
...(typeof updates === "function" ? updates(session) : updates),
updatedAt,
}));
});
}

export const activeTurnIdAfterStart = (current: TurnId | undefined, started: TurnId): TurnId =>
current ?? started;

function parseThreadSnapshot(
response: EffectCodexSchema.V2ThreadReadResponse | EffectCodexSchema.V2ThreadRollbackResponse,
): CodexThreadSnapshot {
Expand Down Expand Up @@ -1313,11 +1316,16 @@ export const makeCodexSessionRuntime = (
),
);
const turnId = TurnId.make(response.turn.id);
yield* updateSession(sessionRef, {
yield* updateSession(sessionRef, (session) => ({
status: "running",
activeTurnId: turnId,
// Codex accepts follow-ups while the current turn is still
// running. The response carries the queued turn id, but
// turn/interrupt only accepts the id that is active now, so keep
// the existing active id until turn-lifecycle notifications
// advance it.
activeTurnId: activeTurnIdAfterStart(session.activeTurnId, turnId),
...(normalizedModel ? { model: normalizedModel } : {}),
});
}));
const resumedProviderThreadId = currentProviderThreadId(yield* Ref.get(sessionRef));
return {
threadId: options.threadId,
Expand Down
Loading