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
109 changes: 109 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2106,6 +2106,7 @@ describe("ClaudeAdapterV2 background wake turns", () => {
const makeWakeHarnessWithOptions = (options?: {
readonly close?: (sdkMessages: Queue.Queue<SDKMessage>) => Effect.Effect<void>;
readonly interrupt?: Effect.Effect<void>;
readonly stopTask?: (taskId: string) => Effect.Effect<void>;
readonly environment?: NodeJS.ProcessEnv;
// A CLI process opened after the first streams from its own queue, so the
// first one can exit (Queue.shutdown) and a later turn can start another.
Expand All @@ -2130,6 +2131,10 @@ describe("ClaudeAdapterV2 background wake turns", () => {
const permissionModeChanges: Array<string> = [];
const continuationRequests: Array<ProviderContinuationRequests.ProviderContinuationRequest> =
[];
const subagentReceipts =
yield* Queue.unbounded<
Extract<ProviderAdapter.ProviderAdapterV2Event, { type: "subagent.updated" }>
>();
const terminalReceipts =
yield* Queue.unbounded<
Extract<ProviderAdapter.ProviderAdapterV2Event, { type: "turn.terminal" }>
Expand Down Expand Up @@ -2202,6 +2207,7 @@ describe("ClaudeAdapterV2 background wake turns", () => {
Effect.sync(() => {
permissionModeChanges.push(mode);
}),
...(options?.stopTask === undefined ? {} : { stopTask: options.stopTask }),
interrupt: options?.interrupt ?? Effect.void,
close: options?.close?.(sdkMessages) ?? Effect.void,
};
Expand All @@ -2228,6 +2234,7 @@ describe("ClaudeAdapterV2 background wake turns", () => {
Stream.runForEach((event) =>
Effect.gen(function* () {
events.push(event);
if (event.type === "subagent.updated") yield* Queue.offer(subagentReceipts, event);
if (event.type === "turn.terminal") {
yield* Queue.offer(terminalReceipts, event);
}
Expand Down Expand Up @@ -2267,6 +2274,7 @@ describe("ClaudeAdapterV2 background wake turns", () => {
continuationRequests,
events,
terminalReceipts,
subagentReceipts,
systemNoticeReceipts,
getOpenedOptions: () => openedOptions,
terminalEvents,
Expand Down Expand Up @@ -8051,6 +8059,107 @@ describe("ClaudeAdapterV2 background wake turns", () => {
),
);

it.effect.each([false, true])(
"stops one native subagent while its sibling keeps running, owner settled %s",
(ownerSettled) =>
Effect.scoped(
Effect.gen(function* () {
const stopped: string[] = [];
let closes = 0;
const harness = yield* makeWakeHarnessWithOptions({
stopTask: (taskId) =>
Effect.sync(() => {
stopped.push(taskId);
}),
interrupt: Effect.die("A child stop must not interrupt its owner"),
close: () =>
Effect.sync(() => {
closes++;
}),
});
yield* harness.runtime.startTurn(
makeClaudeTestTurnInput({
threadId: harness.threadId,
providerThread: harness.providerThread,
now: yield* DateTime.now,
attemptId: RunAttemptId.make("attempt-native-child-stop"),
text: "Start two agents.",
attachments: [],
}),
);
for (const [taskId, toolUseId, uuid] of [
["agent-stop", "toolu-stop", "00000000-0000-4000-8000-000000000801"],
["agent-keep", "toolu-keep", "00000000-0000-4000-8000-000000000802"],
]) {
yield* harness.offerAndWait(
makeSubagentTaskStartedFrame({ taskId: taskId!, toolUseId: toolUseId!, uuid: uuid! }),
);
yield* Queue.take(harness.subagentReceipts);
}
if (ownerSettled) {
yield* harness.offerAndWait(
makeResultFrame({
uuid: "00000000-0000-4000-8000-000000000803",
result: "Agents are working.",
}),
);
yield* Queue.take(harness.terminalReceipts);
}
const stop = harness.runtime.stopSubagent;
assert.isDefined(stop);
yield* stop!({ providerThread: harness.providerThread, nativeTaskId: "agent-stop" });
assert.deepEqual(stopped, ["agent-stop"]);
yield* harness.offerAndWait(
claudeSdkFrame({
...makeSubagentNotificationFrame({
taskId: "agent-stop",
toolUseId: "toolu-stop",
summary: "Stopped by user",
uuid: "00000000-0000-4000-8000-000000000804",
}),
status: "stopped",
}),
);
if (ownerSettled) {
yield* harness.runtime.startTurn(
makeClaudeTestTurnInput({
threadId: harness.threadId,
providerThread: harness.providerThread,
now: yield* DateTime.now,
attemptId: RunAttemptId.make("attempt-native-child-stop-wake"),
text: "Background task stopped.",
attachments: [],
providerTurnOrdinal: 2,
messageCreatedBy: "agent",
messageCreationSource: "provider",
}),
);
}
const stoppedEvent = yield* Queue.take(harness.subagentReceipts);
assert.equal(stoppedEvent.subagent.status, "cancelled");
const tasks = harness.events.filter((event) => event.type === "subagent.updated");
assert.equal(
tasks.findLast((event) => event.subagent.nativeTaskRef?.nativeId === "agent-stop")
?.subagent.status,
"cancelled",
);
assert.equal(
tasks.findLast((event) => event.subagent.nativeTaskRef?.nativeId === "agent-keep")
?.subagent.status,
"running",
);
assert.lengthOf(harness.terminalEvents(), ownerSettled ? 1 : 0);
yield* stop!({ providerThread: harness.providerThread, nativeTaskId: "agent-stop" });
assert.deepEqual(stopped, ["agent-stop"]);
assert.equal(closes, 0);
}).pipe(
Effect.provide(
Layer.mergeAll(IdAllocator.layer, McpProviderSessions.layer, NodeServices.layer),
),
),
),
);

it.effect.each(["requested", "observed-before", "observed-after", "inherit", "unknown"] as const)(
"records the subagent model from %s without inheriting the parent override",
(source) =>
Expand Down
34 changes: 34 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -338,6 +338,7 @@ export interface ClaudeAgentSdkQuerySession {
readonly setPermissionMode: (
mode: PermissionMode,
) => Effect.Effect<void, ClaudeAgentSdkQueryRunnerError>;
readonly stopTask?: (taskId: string) => Effect.Effect<void, ClaudeAgentSdkQueryRunnerError>;
readonly interrupt: Effect.Effect<void, ClaudeAgentSdkQueryRunnerError>;
readonly close: Effect.Effect<void, ClaudeAgentSdkQueryRunnerError>;
}
Expand Down Expand Up @@ -730,6 +731,11 @@ export const layerQueryRunner: Layer.Layer<
}),
),
),
stopTask: (taskId) =>
Effect.tryPromise({
try: () => queryRuntime.stopTask(taskId),
catch: (cause) => queryRunnerError(cause, "stopTask"),
}),
interrupt: Effect.tryPromise({
try: () => queryRuntime.interrupt(),
catch: (cause) => queryRunnerError(cause, "interrupt"),
Expand Down Expand Up @@ -8013,6 +8019,34 @@ export const makeClaudeAdapterV2 = Effect.fn("makeClaudeAdapterV2")(function* (
}),
steerTurn,
interruptTurn,
stopSubagent: Effect.fn("ClaudeAdapterV2.stopSubagent")(function* (input) {
const existing = yield* Ref.get(queryContext);
const nativeThreadId = input.providerThread.nativeThreadRef?.nativeId;
const subagent = (yield* Ref.get(sessionSubagentsByTaskId)).get(input.nativeTaskId);
if (subagent === undefined || subagent.task.status !== "running") return;
if (
existing === null ||
existing.nativeThreadId !== nativeThreadId ||
subagent.task.threadId !== input.providerThread.appThreadId ||
existing.subagentsFromEarlierProcesses.has(subagent) ||
existing.query.stopTask === undefined
) {
return yield* new ProviderAdapter.ProviderAdapterProtocolError({
driver: CLAUDE_PROVIDER,
detail: "Claude subagent has no live owning query with task stopping support.",
});
}
yield* existing.query.stopTask(input.nativeTaskId).pipe(
Effect.mapError(
(cause) =>
new ProviderAdapter.ProviderAdapterProtocolError({
driver: CLAUDE_PROVIDER,
detail: "Claude native subagent stop failed.",
cause,
}),
),
);
}),
respondToRuntimeRequest: Effect.fn("ClaudeAdapterV2.respondToRuntimeRequest")(
function* (requestInput) {
const pending = (yield* Ref.get(pendingRuntimeRequests)).get(
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/orchestration-v2/EffectOutbox.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import {
CheckpointScopeId,
CommandId,
MessageId,
NodeId,
ProviderSessionId,
RunAttemptId,
ProviderApprovalDecision,
Expand Down Expand Up @@ -48,6 +49,7 @@ export const OrchestrationEffectRequestV2 = Schema.Union([
providerSessionId: ProviderSessionId,
providerThreadId: ProviderThreadId,
providerTurnId: ProviderTurnId,
subagent: Schema.optional(Schema.Struct({ id: NodeId, nativeTaskId: Schema.String })),
}),
Schema.Struct({
type: Schema.Literal("provider-turn.steer"),
Expand Down
30 changes: 30 additions & 0 deletions apps/server/src/orchestration-v2/EffectWorker.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { assert, it } from "@effect/vitest";
import {
CommandId,
NodeId,
ProviderSessionId,
ProviderThreadId,
ProviderTurnId,
Expand Down Expand Up @@ -837,3 +838,32 @@ it.effect("settles a delegated child once its restart continuation fails for goo
}).pipe(Effect.provide(layer));
}),
);

it.effect("a native subagent stop does not settle its owner's background work", () =>
Effect.gen(function* () {
const now = yield* DateTime.now;
const events = yield* Ref.make<ReadonlyArray<string>>([]);
const layer = layerExecutorFor({
events,
interrupt: (input) => {
assert.equal(input.subagent?.nativeTaskId, "native-task");
return Ref.update(events, (current) => [...current, "stop-child"]);
},
threads: { dispatch: () => Effect.die("The owner's background work must stay running") },
});
yield* Effect.gen(function* () {
const executor = yield* EffectWorker.OrchestrationEffectExecutorV2;
yield* executor.execute({
...restartEffect(now, { type: "detach" }),
request: {
type: "provider-turn.interrupt",
providerSessionId: oldSessionId,
providerThreadId,
providerTurnId,
subagent: { id: NodeId.make("subagent"), nativeTaskId: "native-task" },
},
});
}).pipe(Effect.provide(layer));
assert.deepEqual(yield* Ref.get(events), ["stop-child"]);
}),
);
19 changes: 12 additions & 7 deletions apps/server/src/orchestration-v2/EffectWorker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -169,6 +169,9 @@ export const layerExecutor: Layer.Layer<
providerSessionId: effect.request.providerSessionId,
providerThreadId: effect.request.providerThreadId,
providerTurnId: effect.request.providerTurnId,
...(effect.request.subagent === undefined
? {}
: { subagent: effect.request.subagent }),
})
.pipe(
Effect.catch((cause) =>
Expand All @@ -185,13 +188,15 @@ export const layerExecutor: Layer.Layer<
// One Stop can interrupt several provider threads, so the
// settle is keyed by effect, not by the Stop command.
Effect.andThen(
threads.dispatch({
type: "thread.background-work.settle",
commandId: CommandId.make(`${effect.id}:background-work-settled`),
threadId: effect.threadId,
providerThreadId: effect.request.providerThreadId,
providerTurnId: effect.request.providerTurnId,
}),
effect.request.subagent !== undefined
? Effect.void
: threads.dispatch({
type: "thread.background-work.settle",
commandId: CommandId.make(`${effect.id}:background-work-settled`),
threadId: effect.threadId,
providerThreadId: effect.request.providerThreadId,
providerTurnId: effect.request.providerTurnId,
}),
),
Effect.mapError(
(cause) =>
Expand Down
Loading
Loading