From 3d536a43612111d4079f4310ad3c722a5ac91ca3 Mon Sep 17 00:00:00 2001 From: Robert Nisipeanu Date: Sat, 22 Aug 2026 19:07:43 +0300 Subject: [PATCH 1/2] fix(claude): report busy heartbeats only while a turn owns them The status and api_retry system messages map to busy session states that only a turn completion can clear in the projector. When one lands after its turn already completed -- compaction routinely outlives the turn that requested it -- the session sticks at running with no active turn and the thread shows Working forever. The adapter reports these heartbeats only while a turn is active; session_state_changed stays unconditional because it is authoritative and reports idle as well as busy. Co-Authored-By: Claude Fable 5 --- .../src/provider/Layers/ClaudeAdapter.test.ts | 91 +++++++++++++++++++ .../src/provider/Layers/ClaudeAdapter.ts | 14 ++- 2 files changed, 104 insertions(+), 1 deletion(-) diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts index 7917b8360946..c610b6d97557 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts @@ -4407,6 +4407,13 @@ describe("ClaudeAdapterLive", () => { provider: ProviderDriverKind.make("claudeAgent"), runtimeMode: "full-access", }); + // status/api_retry heartbeats only report while a turn owns them, so + // the whole batch runs inside one. + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "hello", + attachments: [], + }); // Undeclared wire-only roster snapshot + every typed UX-internal // subtype and top-level type consumed silently: none may surface as @@ -5129,6 +5136,90 @@ describe("ClaudeAdapterLive", () => { ); }); + it.effect("drops status and api_retry heartbeats that arrive with no turn to own them", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEvents: Array = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => runtimeEvents.push(event)), + ).pipe(Effect.forkChild); + + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ + threadId: session.threadId, + input: "/compact", + attachments: [], + }); + + // The sequence a /compact produces when compaction outlives its turn: + // an in-turn compacting status, the turn's result, then the + // post-compaction status clear and a transport retry landing on a + // thread whose turn already completed. Both heartbeats map to busy + // states that only a turn can clear, so reporting either would leave + // the session at running with no active turn forever. + harness.query.emit({ + type: "system", + subtype: "status", + status: "compacting", + session_id: "sdk-session-1", + uuid: "status-compacting", + } as unknown as SDKMessage); + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + errors: [], + num_turns: 0, + session_id: "sdk-session-1", + uuid: "compact-result", + } as unknown as SDKMessage); + harness.query.emit({ + type: "system", + subtype: "status", + status: null, + compact_result: "success", + session_id: "sdk-session-1", + uuid: "status-clear", + } as unknown as SDKMessage); + harness.query.emit({ + type: "system", + subtype: "api_retry", + attempt: 3, + max_retries: 10, + retry_delay_ms: 1000, + error_status: 502, + error: { type: "api_error" }, + session_id: "sdk-session-1", + uuid: "post-turn-retry", + } as unknown as SDKMessage); + yield* Effect.yieldNow; + yield* Effect.yieldNow; + + const heartbeats = runtimeEvents + .filter((event) => event.type === "session.state.changed") + .map((event) => + event.type === "session.state.changed" + ? `${event.payload.state}:${event.payload.reason ?? ""}` + : "", + ) + .filter((entry) => entry.includes(":status:") || entry.includes(":api_retry:")); + // Only the in-turn compacting heartbeat reports; the post-turn pair is + // dropped, leaving the turn completion's ready state in charge. + assert.deepEqual(heartbeats, ["waiting:status:compacting"]); + const turnCompleted = runtimeEvents.find((event) => event.type === "turn.completed"); + assert.equal(turnCompleted?.type, "turn.completed"); + runtimeEventsFiber.interruptUnsafe(); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + it.effect("consumes Claude command lifecycle notifications silently", () => { const harness = makeHarness(); return Effect.gen(function* () { diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.ts b/apps/server/src/provider/Layers/ClaudeAdapter.ts index d50b3cbb9248..93108f2d95ca 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -3624,6 +3624,14 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( }); return; case "status": + // Busy-only heartbeat: every status maps to a busy state, and only + // the owning turn moves the projected session out of it again. A + // status with no turn to own it (compaction routinely outlives the + // turn that requested it) would strand the thread on Working, so it + // is only reported while a turn is active. + if (context.turnState === undefined) { + return; + } yield* offerRuntimeEvent({ ...base, type: "session.state.changed", @@ -3896,7 +3904,11 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( // Transport-level retry heartbeat. Surfacing each attempt as a // warning row spammed the work log (10 rows during a 502 storm); // the terminal result/error path reports the actual failure. Keep - // the session visibly alive instead. + // the session visibly alive instead. Busy-only like `status`, so it + // is only reported while a turn owns it. + if (context.turnState === undefined) { + return; + } yield* offerRuntimeEvent({ ...base, type: "session.state.changed", From 861a3a6ce4a0ec09a9adc9600cf346cb742dc11a Mon Sep 17 00:00:00 2001 From: Robert Nisipeanu Date: Fri, 11 Sep 2026 09:27:18 +0300 Subject: [PATCH 2/2] test(claude): await receipts when draining post-turn messages --- .../src/provider/Layers/ClaudeAdapter.test.ts | 36 ++++++++----------- 1 file changed, 14 insertions(+), 22 deletions(-) diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts index c610b6d97557..6851b24eea5e 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts @@ -3600,7 +3600,7 @@ describe("ClaudeAdapterLive", () => { const firstTurnSettled = yield* Deferred.make(); const completedFiber = yield* adapter.streamEvents.pipe( Stream.tap((event) => - event.type === "session.state.changed" && event.payload.reason === "api_retry:1/2" + event.type === "session.state.changed" && event.payload.reason === "session_state:idle" ? Deferred.succeed(firstTurnSettled, undefined).pipe(Effect.asVoid) : Effect.void, ), @@ -3651,12 +3651,8 @@ describe("ClaudeAdapterLive", () => { } as unknown as SDKMessage); harness.query.emit({ type: "system", - subtype: "api_retry", - attempt: 1, - max_retries: 2, - retry_delay_ms: 1, - error_status: 502, - error: { type: "api_error" }, + subtype: "session_state_changed", + state: "idle", session_id: "sdk-session-consecutive-usage", uuid: "consecutive-usage-barrier", } as unknown as SDKMessage); @@ -4623,7 +4619,7 @@ describe("ClaudeAdapterLive", () => { if ( receipt && event.type === "session.state.changed" && - event.payload.reason === "api_retry:1/1" + event.payload.reason === "session_state:idle" ) { yield* Deferred.succeed(receipt, undefined); } @@ -4631,15 +4627,12 @@ describe("ClaudeAdapterLive", () => { ).pipe(Effect.forkChild); const drainSdkMessages = Effect.gen(function* () { receipt = yield* Deferred.make(); - // The heartbeat follows queued SDK messages without adding a warning. + // Session-state notifications still report between turns, so this + // receipt drains queued SDK messages even after a turn's result. query.emit({ type: "system", - subtype: "api_retry", - attempt: 1, - max_retries: 1, - retry_delay_ms: 0, - error_status: 429, - error: { type: "rate_limit_error" }, + subtype: "session_state_changed", + state: "idle", session_id: "sdk-session-limit", uuid: "usage-limit-drain", } as unknown as SDKMessage); @@ -5140,10 +5133,10 @@ describe("ClaudeAdapterLive", () => { const harness = makeHarness(); return Effect.gen(function* () { const adapter = yield* ClaudeAdapter; - const runtimeEvents: Array = []; - const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => - Effect.sync(() => runtimeEvents.push(event)), - ).pipe(Effect.forkChild); + const runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "session.exited", + ).pipe(Stream.runCollect, Effect.forkChild); const session = yield* adapter.startSession({ threadId: THREAD_ID, @@ -5197,8 +5190,8 @@ describe("ClaudeAdapterLive", () => { session_id: "sdk-session-1", uuid: "post-turn-retry", } as unknown as SDKMessage); - yield* Effect.yieldNow; - yield* Effect.yieldNow; + harness.query.finish(); + const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); const heartbeats = runtimeEvents .filter((event) => event.type === "session.state.changed") @@ -5213,7 +5206,6 @@ describe("ClaudeAdapterLive", () => { assert.deepEqual(heartbeats, ["waiting:status:compacting"]); const turnCompleted = runtimeEvents.find((event) => event.type === "turn.completed"); assert.equal(turnCompleted?.type, "turn.completed"); - runtimeEventsFiber.interruptUnsafe(); }).pipe( Effect.provideService(Random.Random, makeDeterministicRandomService()), Effect.provide(harness.layer),