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
79 changes: 79 additions & 0 deletions apps/server/src/provider/Layers/ClaudeAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,8 @@ class FakeClaudeQuery implements AsyncIterable<SDKMessage> {
public readonly setMaxThinkingTokensCalls: Array<number | null> = [];
public closeCalls = 0;
public closeError: unknown | undefined;
/** Set by tests that exercise Claude's graceful interrupt. */
public interrupt?: () => Promise<unknown>;

emit(message: SDKMessage): void {
if (this.done) {
Expand Down Expand Up @@ -3255,6 +3257,83 @@ describe("ClaudeAdapterLive", () => {
);
});

it.effect("interruptTurn lets Claude abort the turn before closing the session", () => {
const harness = makeHarness();
return Effect.gen(function* () {
const adapter = yield* ClaudeAdapter;
const session = yield* adapter.startSession({
threadId: THREAD_ID,
provider: ProviderDriverKind.make("claudeAgent"),
runtimeMode: "full-access",
});
yield* adapter.sendTurn({
threadId: session.threadId,
input: "hello",
attachments: [],
});

const turnCompletedFiber = yield* adapter.streamEvents.pipe(
Stream.filter((event) => event.type === "turn.completed"),
Stream.take(1),
Stream.runCollect,
Effect.forkChild,
);
let closeCallsAtInterrupt: number | undefined;
harness.query.interrupt = async () => {
closeCallsAtInterrupt = harness.query.closeCalls;
harness.query.emit({
type: "result",
subtype: "error_during_execution",
is_error: false,
errors: ["Error: Request was aborted."],
session_id: "sdk-session",
uuid: "result-interrupted",
} as unknown as SDKMessage);
};

yield* adapter.interruptTurn(session.threadId);

assert.equal(closeCallsAtInterrupt, 0);
assert.equal(harness.query.closeCalls, 1);
const [turnCompleted] = Array.from(yield* Fiber.join(turnCompletedFiber));
assert.equal(turnCompleted?.type, "turn.completed");
if (turnCompleted?.type === "turn.completed") {
assert.equal(turnCompleted.payload.state, "interrupted");
}
}).pipe(
Effect.provideService(Random.Random, makeDeterministicRandomService()),
Effect.provide(harness.layer),
);
});

it.effect("interruptTurn closes the session when Claude never aborts the turn", () => {
const harness = makeHarness();
return Effect.gen(function* () {
const adapter = yield* ClaudeAdapter;
const session = yield* adapter.startSession({
threadId: THREAD_ID,
provider: ProviderDriverKind.make("claudeAgent"),
runtimeMode: "full-access",
});
yield* adapter.sendTurn({
threadId: session.threadId,
input: "hello",
attachments: [],
});
harness.query.interrupt = () => new Promise(() => {});

const interruptFiber = yield* adapter.interruptTurn(session.threadId).pipe(Effect.forkChild);
yield* TestClock.adjust("3 seconds");
yield* Fiber.join(interruptFiber);

assert.equal(harness.query.closeCalls, 1);
assert.equal(yield* adapter.hasSession(session.threadId), false);
}).pipe(
Effect.provideService(Random.Random, makeDeterministicRandomService()),
Effect.provide(harness.layer),
);
});

it.effect("keeps the session available when process close fails", () => {
const harness = makeHarness();
return Effect.gen(function* () {
Expand Down
30 changes: 30 additions & 0 deletions apps/server/src/provider/Layers/ClaudeAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -385,6 +385,8 @@ interface ClaudeTaskAgentState {
* lifetime; oldest entries evict first.
*/
const PENDING_TASK_MODEL_CAP = 64;
/** How long Stop waits for Claude to abort a turn before killing the process. */
const CLAUDE_INTERRUPT_GRACE = "3 seconds";

/**
* Buffers a subagent snapshot's authoritative model under its
Expand Down Expand Up @@ -453,10 +455,14 @@ interface ClaudeSessionContext {
lastThreadStartedId: string | undefined;
/** Limits already announced for the running turn, keyed `window:resetsAt`. */
announcedUsageLimits: { turnId: string; keys: Set<string> } | undefined;
/** Resolved by completeTurn while Stop waits for Claude to abort the turn. */
interruptedTurnSettled: Deferred.Deferred<void> | undefined;
stopped: boolean;
}

interface ClaudeQueryRuntime extends AsyncIterable<SDKMessage> {
/** SDK Query.interrupt — present on real queries; optional for test doubles. */
readonly interrupt?: () => Promise<unknown>;
readonly setModel: (model?: string) => Promise<void>;
readonly setPermissionMode: (mode: PermissionMode) => Promise<void>;
readonly setMaxThinkingTokens: (maxThinkingTokens: number | null) => Promise<void>;
Expand Down Expand Up @@ -2836,6 +2842,9 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (

const updatedAt = yield* nowIso;
context.turnState = undefined;
if (context.interruptedTurnSettled) {
yield* Deferred.succeed(context.interruptedTurnSettled, undefined);
}
context.session = {
...context.session,
status: "ready",
Expand Down Expand Up @@ -5045,6 +5054,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
lastAssistantUuid: resumeState?.resumeSessionAt,
lastThreadStartedId: undefined,
announcedUsageLimits: undefined,
interruptedTurnSettled: undefined,
stopped: false,
};
yield* Ref.set(contextRef, context);
Expand Down Expand Up @@ -5276,13 +5286,33 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
const interruptTurn: ClaudeAdapterShape["interruptTurn"] = Effect.fn("interruptTurn")(
function* (threadId, _turnId) {
const context = yield* requireSession(threadId);
yield* settleInterruptedTurn(context);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
// interrupt() can acknowledge while resumed background tasks keep the
// CLI alive. Stop is a hard session boundary for Claude, so close the
// query and let the SDK escalate to SIGKILL when graceful exit fails.
yield* stopSessionInternal(context);
},
);

// Lets Claude abort the running turn through its own path before Stop kills
// the process, so the prompt reaches the transcript. Killing a first turn
// before Claude writes it leaves a resume cursor for a session Claude never
// saved, and every later message fails with "No conversation found".
const settleInterruptedTurn = Effect.fn("settleInterruptedTurn")(function* (
context: ClaudeSessionContext,
) {
const interrupt = context.query.interrupt?.bind(context.query);
if (context.stopped || !context.turnState || !interrupt) return;
const settled = yield* Deferred.make<void>();
context.interruptedTurnSettled = settled;
yield* Effect.tryPromise(interrupt).pipe(
Effect.ignore,
Effect.andThen(Deferred.await(settled)),
Effect.timeoutOption(CLAUDE_INTERRUPT_GRACE),
);
context.interruptedTurnSettled = undefined;
});

const readThread: ClaudeAdapterShape["readThread"] = Effect.fn("readThread")(
function* (threadId) {
const context = yield* requireSession(threadId);
Expand Down
Loading