From 6a6e1c413fe444fe253fa61bb88186c72bf358c0 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Thu, 6 Aug 2026 20:50:01 -0700 Subject: [PATCH] fix(server): settle stopped Claude subagents --- .../src/provider/Layers/ClaudeAdapter.test.ts | 19 ++++++++- .../src/provider/Layers/ClaudeAdapter.ts | 40 ++++++++++++++++--- 2 files changed, 53 insertions(+), 6 deletions(-) diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts index afa65ea39d6..d3d768b5384 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts @@ -1511,7 +1511,7 @@ describe("ClaudeAdapterLive", () => { ); }); - it.effect("interruptTurn stops every live task before interrupting the turn", () => { + it.effect("interruptTurn settles every acknowledged live task before interrupting", () => { const harness = makeHarness(); return Effect.gen(function* () { const adapter = yield* ClaudeAdapter; @@ -1568,11 +1568,28 @@ describe("ClaudeAdapterLive", () => { yield* Fiber.join(taskEventsFiber); + const stoppedTaskEventFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.type === "task.completed"), + Stream.take(1), + Stream.runCollect, + Effect.forkChild, + ); yield* adapter.interruptTurn(session.threadId); // Only the still-live task is stopped; interrupt always fires after. assert.deepEqual(harness.query.stopTaskCalls, ["task-live"]); assert.equal(harness.query.interruptCalls.length, 1); + + const stoppedTaskEvents = Array.from(yield* Fiber.join(stoppedTaskEventFiber)); + assert.equal(stoppedTaskEvents.length, 1); + const stoppedTaskEvent = stoppedTaskEvents[0]; + assert.equal(stoppedTaskEvent?.type, "task.completed"); + if (stoppedTaskEvent?.type === "task.completed") { + assert.equal(String(stoppedTaskEvent.payload.taskId), "task-live"); + assert.equal(stoppedTaskEvent.payload.status, "stopped"); + assert.equal(stoppedTaskEvent.payload.taskType, "local_agent"); + assert.equal(stoppedTaskEvent.payload.title, "Agent A"); + } }).pipe( Effect.provideService(Random.Random, makeDeterministicRandomService()), Effect.provide(harness.layer), diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.ts b/apps/server/src/provider/Layers/ClaudeAdapter.ts index f6f1c14420d..92445522cc4 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -65,6 +65,7 @@ import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as FileSystem from "effect/FileSystem"; import * as Fiber from "effect/Fiber"; +import * as Option from "effect/Option"; import * as Path from "effect/Path"; import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; @@ -4419,11 +4420,40 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( yield* Effect.forEach( liveIds, (taskId) => - Effect.tryPromise({ - // Invoke through the query object: SDK methods rely on `this`. - try: () => context.query.stopTask!(taskId), - catch: () => undefined, - }).pipe(Effect.timeoutOption("3 seconds"), Effect.ignore), + Effect.gen(function* () { + const stopAcknowledged = yield* Effect.tryPromise({ + // Invoke through the query object: SDK methods rely on `this`. + try: () => context.query.stopTask!(taskId), + catch: () => undefined, + }).pipe( + Effect.timeoutOption("3 seconds"), + Effect.orElseSucceed(() => Option.none()), + ); + if (Option.isNone(stopAcknowledged) || !context.liveTaskIds.delete(taskId)) { + return; + } + + // stopTask only acknowledges the control request. Its separate + // task_notification can lose the race with interrupt(), so make + // the acknowledged stop authoritative for the durable UI state. + const stamp = yield* makeEventStamp(); + yield* offerRuntimeEvent({ + type: "task.completed", + eventId: stamp.eventId, + provider: PROVIDER, + createdAt: stamp.createdAt, + threadId: context.session.threadId, + ...(context.turnState + ? { turnId: asCanonicalTurnId(context.turnState.turnId) } + : {}), + payload: { + taskId: RuntimeTaskId.make(taskId), + status: "stopped", + ...taskLinkageFor(context.taskAgents, taskId), + }, + providerRefs: nativeProviderRefs(context), + }); + }).pipe(Effect.ignore), { concurrency: 8, discard: true }, ).pipe(Effect.timeoutOption("10 seconds"), Effect.ignore); }