Skip to content
Closed
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
136 changes: 76 additions & 60 deletions apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ import {
import * as DateTime from "effect/DateTime";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
Expand Down Expand Up @@ -208,6 +209,8 @@ function makeDeterministicAdapter(input: {
});
const runOrdinals = new Map<ProviderTurnId, number>();
const turnInputs = new Map<ProviderTurnId, ProviderAdapterV2TurnInput>();
const turnFibers = new Map<ProviderTurnId, Fiber.Fiber<void>>();
const sessionScope = yield* Effect.scope;

return {
instanceId: input.instanceId,
Expand Down Expand Up @@ -279,76 +282,89 @@ function makeDeterministicAdapter(input: {
},
},
]);
const terminalGate = input.terminalGate?.(turnInput);
if (terminalGate !== undefined) {
yield* Deferred.await(terminalGate);
} else if (!input.shouldComplete(turnInput)) {
return;
}
const response = input.response(turnInput);
yield* publish([
{
type: "provider_turn.updated",
driver: input.driver,
providerTurn: {
id: providerTurnId,
providerThreadId: turnInput.providerThread.id,
nodeId: turnInput.rootNodeId,
runAttemptId: turnInput.attemptId,
nativeTurnRef: {
driver: input.driver,
nativeId: `native-turn:${turnInput.threadId}:${turnInput.runOrdinal}`,
strength: "strong",
// Acquire-then-return like the production adapters: the call
// returns after dispatch and the gated completion runs in a
// session-scoped fiber. Holding the call open until the gate
// released would leave an admitted adapter op parked until scope
// close, where the release drain would wait on it forever.
const fiber = yield* Effect.gen(function* () {
const terminalGate = input.terminalGate?.(turnInput);
if (terminalGate !== undefined) {
yield* Deferred.await(terminalGate);
} else if (!input.shouldComplete(turnInput)) {
return;
}
const response = input.response(turnInput);
yield* publish([
{
type: "provider_turn.updated",
driver: input.driver,
providerTurn: {
id: providerTurnId,
providerThreadId: turnInput.providerThread.id,
nodeId: turnInput.rootNodeId,
runAttemptId: turnInput.attemptId,
nativeTurnRef: {
driver: input.driver,
nativeId: `native-turn:${turnInput.threadId}:${turnInput.runOrdinal}`,
strength: "strong",
},
ordinal: turnInput.providerTurnOrdinal,
status: "completed",
startedAt: eventTime,
completedAt: eventTime,
},
ordinal: turnInput.providerTurnOrdinal,
status: "completed",
startedAt: eventTime,
completedAt: eventTime,
},
},
{
type: "turn_item.updated",
driver: input.driver,
turnItem: {
id: TurnItemId.make(
`turn-item:${input.instanceId}:${turnInput.threadId}:${turnInput.runOrdinal}:assistant`,
),
threadId: turnInput.threadId,
runId: turnInput.runId,
nodeId: turnInput.rootNodeId,
{
type: "turn_item.updated",
driver: input.driver,
turnItem: {
id: TurnItemId.make(
`turn-item:${input.instanceId}:${turnInput.threadId}:${turnInput.runOrdinal}:assistant`,
),
threadId: turnInput.threadId,
runId: turnInput.runId,
nodeId: turnInput.rootNodeId,
providerThreadId: turnInput.providerThread.id,
providerTurnId,
nativeItemRef: null,
parentItemId: null,
ordinal: turnInput.runOrdinal * 100 + 1,
status: "completed",
title: null,
startedAt: eventTime,
completedAt: eventTime,
updatedAt: eventTime,
type: "assistant_message",
messageId: MessageId.make(
`message:${input.instanceId}:${turnInput.threadId}:${turnInput.runOrdinal}:assistant`,
),
text: response,
streaming: false,
},
},
{
type: "turn.terminal",
driver: input.driver,
providerThreadId: turnInput.providerThread.id,
providerTurnId,
nativeItemRef: null,
parentItemId: null,
ordinal: turnInput.runOrdinal * 100 + 1,
runOrdinal: turnInput.runOrdinal,
status: "completed",
title: null,
startedAt: eventTime,
completedAt: eventTime,
updatedAt: eventTime,
type: "assistant_message",
messageId: MessageId.make(
`message:${input.instanceId}:${turnInput.threadId}:${turnInput.runOrdinal}:assistant`,
),
text: response,
streaming: false,
failure: null,
threadDisposition: "reusable",
},
},
{
type: "turn.terminal",
driver: input.driver,
providerThreadId: turnInput.providerThread.id,
providerTurnId,
runOrdinal: turnInput.runOrdinal,
status: "completed",
failure: null,
threadDisposition: "reusable",
},
]);
]);
}).pipe(Effect.forkIn(sessionScope));
turnFibers.set(providerTurnId, fiber);
}),
steerTurn: () => Effect.void,
interruptTurn: ({ providerThread, providerTurnId }) =>
Effect.gen(function* () {
const fiber = turnFibers.get(providerTurnId);
if (fiber !== undefined) {
yield* Fiber.interrupt(fiber);
turnFibers.delete(providerTurnId);
}
const turnInput = turnInputs.get(providerTurnId);
const completedAt = yield* DateTime.now;
if (turnInput !== undefined) {
Expand Down
9 changes: 7 additions & 2 deletions apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7014,7 +7014,12 @@ export function makeAcpAdapterV2(options: AcpAdapterV2Options): ProviderAdapterV
!restartRequired &&
!transportHadOutstandingResponses
) {
yield* runtime.closeSession().pipe(Effect.ignore);
// A failed native close means the ACP subprocess may still
// own the thread — the failure must reach Scope.close so
// the manager keeps cleanup ownership instead of releasing
// a live process. Finalizers cannot carry typed errors,
// so the failure surfaces as a defect.
yield* runtime.closeSession().pipe(Effect.orDie);
}
}),
);
Expand All @@ -7023,7 +7028,7 @@ export function makeAcpAdapterV2(options: AcpAdapterV2Options): ProviderAdapterV
yield* flavor.assertComplete.pipe(Effect.orDie);
}
if (runtimeScope !== undefined) {
yield* Scope.close(runtimeScope, Exit.void).pipe(Effect.ignore);
yield* Scope.close(runtimeScope, Exit.void);
}
}),
);
Expand Down
183 changes: 183 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,9 +40,11 @@ import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Path from "effect/Path";
import * as Queue from "effect/Queue";
import * as Ref from "effect/Ref";
import * as Schema from "effect/Schema";
import * as Scope from "effect/Scope";
import * as Stream from "effect/Stream";
import { TestClock } from "effect/testing";
import { Tool } from "effect/unstable/ai";
import { formatClaudeResumeCompactionQuestion } from "@t3tools/shared/claudeCompaction";

Expand Down Expand Up @@ -7109,3 +7111,184 @@ describe("ClaudeAdapterV2 query message stream", () => {
}),
);
});

describe("ClaudeAdapterV2 session cleanup", () => {
it.effect("propagates a native query close failure through the session scope", () =>
Effect.gen(function* () {
const fileSystem = yield* FileSystem.FileSystem;
const attachmentsDir = yield* fileSystem.makeTempDirectoryScoped({
prefix: "t3-claude-close-",
});
const scope = yield* Scope.make();
const adapter = makeClaudeAdapterV2({
instanceId: CLAUDE_DEFAULT_INSTANCE_ID,
settings: DEFAULT_CLAUDE_SETTINGS,
environment: {},
attachmentsDir,
fileSystem,
path: yield* Path.Path,
idAllocator: yield* IdAllocatorV2,
queryRunner: {
allocateSessionId: Effect.succeed("native-thread-claude-close"),
open: () =>
Effect.succeed({
messages: Stream.never,
offer: () => Effect.void,
setModel: () => Effect.void,
interrupt: Effect.void,
// A typed runner failure is what `Effect.ignore` used to
// swallow — a defect would propagate through it either way, so
// only the error channel distinguishes the fixed behavior.
close: Effect.fail(
new ClaudeAgentSdkQueryRunnerError({
method: "query.close",
cause: new Error("claude query close failed"),
}),
),
}),
forkSession: () => Effect.die("unused"),
subagentLaunchToolUseId: () => Effect.succeed(null),
assertComplete: Effect.void,
},
});
const threadId = ThreadId.make("thread-claude-close");
const runtime = yield* adapter
.openSession({
threadId,
providerSessionId: ProviderSessionId.make("provider-session-claude-close"),
modelSelection: CLAUDE_TEST_MODEL_SELECTION,
runtimePolicy: CLAUDE_TEST_RUNTIME_POLICY,
})
.pipe(Effect.provideService(Scope.Scope, scope));
const providerThread = yield* runtime.ensureThread({
threadId,
modelSelection: CLAUDE_TEST_MODEL_SELECTION,
runtimePolicy: CLAUDE_TEST_RUNTIME_POLICY,
});
yield* runtime.startTurn(
makeClaudeTestTurnInput({
threadId,
providerThread,
now: yield* DateTime.now,
attemptId: RunAttemptId.make("attempt-claude-close"),
text: "hello",
attachments: [],
}),
);

// The native close failure must reach the scope close so the session
// manager records a failed cleanup instead of a clean release.
const closeExit = yield* Scope.close(scope, Exit.void).pipe(Effect.exit);
assert.isTrue(Exit.isFailure(closeExit));
}).pipe(Effect.provide(Layer.mergeAll(NodeServices.layer, idAllocatorLayer))),
);

it.effect("keeps a failed interrupt close tracked so the session scope retries it", () =>
Effect.gen(function* () {
const fileSystem = yield* FileSystem.FileSystem;
const attachmentsDir = yield* fileSystem.makeTempDirectoryScoped({
prefix: "t3-claude-interrupt-close-",
});
const scope = yield* Scope.make();
const closeCalls = yield* Ref.make(0);
const interrupted = yield* Deferred.make<void>();
const adapter = makeClaudeAdapterV2({
instanceId: CLAUDE_DEFAULT_INSTANCE_ID,
settings: DEFAULT_CLAUDE_SETTINGS,
environment: {},
attachmentsDir,
fileSystem,
path: yield* Path.Path,
idAllocator: yield* IdAllocatorV2,
queryRunner: {
allocateSessionId: Effect.succeed("native-thread-claude-interrupt-close"),
open: () =>
Effect.succeed({
// Never ends, so interruptTurn's closed wait can only resolve
// through its own timeout path.
messages: Stream.never,
offer: () => Effect.void,
setModel: () => Effect.void,
interrupt: Deferred.succeed(interrupted, undefined).pipe(Effect.asVoid),
close: Ref.update(closeCalls, (count) => count + 1).pipe(
Effect.andThen(
Effect.fail(
new ClaudeAgentSdkQueryRunnerError({
method: "query.close",
cause: new Error("claude query close failed"),
}),
),
),
),
}),
forkSession: () => Effect.die("unused"),
subagentLaunchToolUseId: () => Effect.succeed(null),
assertComplete: Effect.void,
},
});
const threadId = ThreadId.make("thread-claude-interrupt-close");
const runtime = yield* adapter
.openSession({
threadId,
providerSessionId: ProviderSessionId.make("provider-session-claude-interrupt-close"),
modelSelection: CLAUDE_TEST_MODEL_SELECTION,
runtimePolicy: CLAUDE_TEST_RUNTIME_POLICY,
})
.pipe(Effect.provideService(Scope.Scope, scope));
const providerThread = yield* runtime.ensureThread({
threadId,
modelSelection: CLAUDE_TEST_MODEL_SELECTION,
runtimePolicy: CLAUDE_TEST_RUNTIME_POLICY,
});
const events: Array<ProviderAdapterV2Event> = [];
yield* runtime.events.pipe(
Stream.runForEach((event) =>
Effect.sync(() => {
events.push(event);
}),
),
Effect.forkIn(scope),
);
yield* runtime.startTurn(
makeClaudeTestTurnInput({
threadId,
providerThread,
now: yield* DateTime.now,
attemptId: RunAttemptId.make("attempt-claude-interrupt-close"),
text: "hello",
attachments: [],
}),
);
for (let attempt = 0; attempt < 5000; attempt++) {
if (events.some((event) => event.type === "provider_turn.updated")) {
break;
}
yield* Effect.yieldNow;
}
const providerTurnId = events.find(
(event): event is Extract<ProviderAdapterV2Event, { type: "provider_turn.updated" }> =>
event.type === "provider_turn.updated",
)?.providerTurn.id;
assert.isDefined(providerTurnId);

const interrupting = yield* runtime
.interruptTurn({ providerThread, providerTurnId: providerTurnId! })
.pipe(Effect.exit, Effect.forkChild);
yield* Deferred.await(interrupted);
yield* Effect.yieldNow;
yield* Effect.yieldNow;
yield* TestClock.adjust("10 seconds");
const interruptExit = yield* Fiber.join(interrupting);
// The failed close propagates through interruptTurn...
assert.equal(interruptExit._tag, "Failure");
assert.equal(yield* Ref.get(closeCalls), 1);

// ...and the query stays tracked, so the scope close retries the
// native close and surfaces the failure instead of reporting a clean
// release over a CLI that may still be alive.
const closeExit = yield* Scope.close(scope, Exit.void).pipe(Effect.exit);
assert.isTrue(Exit.isFailure(closeExit));
assert.equal(yield* Ref.get(closeCalls), 2);
}).pipe(Effect.provide(Layer.mergeAll(NodeServices.layer, idAllocatorLayer))),
);
});
Loading
Loading