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
191 changes: 177 additions & 14 deletions apps/server/src/provider/Layers/ClaudeAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -54,12 +54,11 @@ class FakeClaudeQuery implements AsyncIterable<SDKMessage> {
private done = false;
private failure: unknown | undefined;

public readonly interruptCalls: Array<void> = [];
public readonly stopTaskCalls: Array<string> = [];
public readonly setModelCalls: Array<string | undefined> = [];
public readonly setPermissionModeCalls: Array<string> = [];
public readonly setMaxThinkingTokensCalls: Array<number | null> = [];
public closeCalls = 0;
public closeError: unknown | undefined;

emit(message: SDKMessage): void {
if (this.done) {
Expand Down Expand Up @@ -95,14 +94,6 @@ class FakeClaudeQuery implements AsyncIterable<SDKMessage> {
}
}

readonly interrupt = async (): Promise<void> => {
this.interruptCalls.push(undefined);
};

readonly stopTask = async (taskId: string): Promise<void> => {
this.stopTaskCalls.push(taskId);
};

readonly setModel = async (model?: string): Promise<void> => {
this.setModelCalls.push(model);
};
Expand All @@ -117,6 +108,9 @@ class FakeClaudeQuery implements AsyncIterable<SDKMessage> {

readonly close = (): void => {
this.closeCalls += 1;
if (this.closeError !== undefined) {
throw this.closeError;
}
this.finish();
};

Expand Down Expand Up @@ -1580,7 +1574,7 @@ describe("ClaudeAdapterLive", () => {
);
});

it.effect("interruptTurn settles every acknowledged live task before interrupting", () => {
it.effect("interruptTurn settles live tasks and closes the provider session", () => {
const harness = makeHarness();
return Effect.gen(function* () {
const adapter = yield* ClaudeAdapter;
Expand Down Expand Up @@ -1645,9 +1639,12 @@ describe("ClaudeAdapterLive", () => {
);
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);
// Closing the session is the hard stop because SDK interrupt can leave
// resumed background work alive.
assert.equal(harness.query.closeCalls, 1);

const sessions = yield* adapter.listSessions();
assert.equal(sessions.length, 0);

const stoppedTaskEvents = Array.from(yield* Fiber.join(stoppedTaskEventFiber));
assert.equal(stoppedTaskEvents.length, 1);
Expand All @@ -1665,6 +1662,172 @@ describe("ClaudeAdapterLive", () => {
);
});

it.effect("keeps the session available when process close fails", () => {
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",
});
harness.query.closeError = new Error("close failed");

const result = yield* adapter.interruptTurn(session.threadId).pipe(Effect.result);

assert.equal(result._tag, "Failure");
if (result._tag === "Failure") {
assert.equal(result.failure._tag, "ProviderAdapterProcessError");
}
assert.equal(harness.query.closeCalls, 1);
assert.equal(yield* adapter.hasSession(session.threadId), true);
assert.equal((yield* adapter.listSessions())[0]?.status, "ready");
}).pipe(
Effect.provideService(Random.Random, makeDeterministicRandomService()),
Effect.provide(harness.layer),
);
});

it.effect("stopAll attempts every session when one process close fails", () => {
const queries: FakeClaudeQuery[] = [];
const layer = Layer.effect(
ClaudeAdapter,
Effect.gen(function* () {
const claudeConfig = decodeClaudeSettings({});
return yield* makeClaudeAdapter(claudeConfig, {
createQuery: () => {
const query = new FakeClaudeQuery();
queries.push(query);
return query;
},
});
}),
).pipe(
Layer.provideMerge(ServerConfig.layerTest("/tmp/claude-adapter-test", "/tmp")),
Layer.provideMerge(ServerSettingsService.layerTest()),
Layer.provideMerge(NodeServices.layer),
);

return Effect.gen(function* () {
const adapter = yield* ClaudeAdapter;
yield* adapter.startSession({
threadId: THREAD_ID,
provider: ProviderDriverKind.make("claudeAgent"),
runtimeMode: "full-access",
});
yield* adapter.startSession({
threadId: RESUME_THREAD_ID,
provider: ProviderDriverKind.make("claudeAgent"),
runtimeMode: "full-access",
});
const firstQuery = queries[0];
if (!firstQuery) {
return;
}
firstQuery.closeError = new Error("close failed");

const result = yield* adapter.stopAll().pipe(Effect.result);

assert.equal(result._tag, "Failure");
assert.equal(queries[0]?.closeCalls, 1);
assert.equal(queries[1]?.closeCalls, 1);
assert.equal(yield* adapter.hasSession(THREAD_ID), true);
assert.equal(yield* adapter.hasSession(RESUME_THREAD_ID), false);
}).pipe(
Effect.provideService(Random.Random, makeDeterministicRandomService()),
Effect.provide(layer),
);
});

it.effect("keeps a resumed replacement session during slow stop cleanup", () => {
const queries: FakeClaudeQuery[] = [];
let signalUsageStarted: () => void = () => undefined;
const usageStarted = new Promise<void>((resolve) => {
signalUsageStarted = resolve;
});
const layer = Layer.effect(
ClaudeAdapter,
Effect.gen(function* () {
const claudeConfig = decodeClaudeSettings({});
return yield* makeClaudeAdapter(claudeConfig, {
createQuery: () => {
const query = new FakeClaudeQuery();
if (queries.length === 0) {
Object.assign(query, {
getContextUsage: async () => {
signalUsageStarted();
return await new Promise<never>(() => undefined);
},
});
}
queries.push(query);
return query;
},
});
}),
).pipe(
Layer.provideMerge(ServerConfig.layerTest("/tmp/claude-adapter-test", "/tmp")),
Layer.provideMerge(ServerSettingsService.layerTest()),
Layer.provideMerge(NodeServices.layer),
);

return Effect.gen(function* () {
const adapter = yield* ClaudeAdapter;
const runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 8).pipe(
Stream.runCollect,
Effect.forkChild,
);
const firstSession = yield* adapter.startSession({
threadId: THREAD_ID,
provider: ProviderDriverKind.make("claudeAgent"),
runtimeMode: "full-access",
});
yield* adapter.sendTurn({
threadId: firstSession.threadId,
input: "hello",
attachments: [],
});

const interruptFiber = yield* adapter
.interruptTurn(firstSession.threadId)
.pipe(Effect.forkChild);
yield* Effect.promise(() => usageStarted);
assert.equal(queries[0]?.closeCalls, 1);

const replacement = yield* adapter.startSession({
threadId: THREAD_ID,
provider: ProviderDriverKind.make("claudeAgent"),
runtimeMode: "full-access",
resumeCursor: firstSession.resumeCursor,
});
yield* TestClock.adjust("1 second");
yield* Fiber.join(interruptFiber);

const activeSessions = yield* adapter.listSessions();
const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber));
assert.equal(queries.length, 2);
assert.equal(queries[1]?.closeCalls, 0);
assert.equal(activeSessions.length, 1);
assert.deepEqual(activeSessions[0]?.resumeCursor, replacement.resumeCursor);
assert.deepEqual(
runtimeEvents
.filter((event) => event.type.startsWith("session."))
.map((event) => event.type),
[
"session.started",
"session.configured",
"session.state.changed",
"session.started",
"session.configured",
"session.state.changed",
],
);
}).pipe(
Effect.provideService(Random.Random, makeDeterministicRandomService()),
Effect.provide(layer),
);
});

it.effect("workflow member coalescing: identical snapshots suppress, changes emit", () => {
const harness = makeHarness();
return Effect.gen(function* () {
Expand Down
Loading
Loading