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
8 changes: 8 additions & 0 deletions apps/server/scripts/record-claude-agent-sdk-replay-fixture.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@ import {
CLAUDE_BACKGROUND_SUBAGENT_LIFECYCLE_RESUME_PROMPT,
CLAUDE_BACKGROUND_SUBAGENT_LIFECYCLE_STOP_PROMPT,
} from "../src/orchestration-v2/testkit/fixtures/claude_background_subagent_lifecycle/input.ts";
import { CLAUDE_BACKGROUND_MONITOR_WAKE_PROMPT } from "../src/orchestration-v2/testkit/fixtures/claude_background_monitor_wake/input.ts";
import { CLAUDE_BACKGROUND_TASK_INTERRUPT_PROMPT } from "../src/orchestration-v2/testkit/fixtures/claude_background_task_interrupt/input.ts";
import { CLAUDE_BACKGROUND_WAKE_BEFORE_QUEUED_PROMPT_LAUNCH_PROMPT } from "../src/orchestration-v2/testkit/fixtures/claude_background_wake_before_queued_prompt/input.ts";
import {
Expand Down Expand Up @@ -192,6 +193,13 @@ const CLAUDE_RECORDINGS = {
enableTools: true,
backgroundWakeCounts: [1, 0],
},
claude_background_monitor_wake: {
prompts: [CLAUDE_BACKGROUND_MONITOR_WAKE_PROMPT],
defaultTranscriptFile: "fixtures/claude_background_monitor_wake/claude_transcript.ndjson",
queryMode: "streaming",
enableTools: true,
backgroundWakeCounts: [1],
},
claude_background_subagent_lifecycle: {
prompts: [
CLAUDE_BACKGROUND_SUBAGENT_LIFECYCLE_LAUNCH_PROMPT,
Expand Down
135 changes: 135 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3308,6 +3308,141 @@ describe("ClaudeAdapterV2 background wake turns", () => {
),
);

// Frame shapes follow the claude_background_monitor_wake recording: Claude
// runs a Monitor as a local_bash task, linked to its call by tool_use_id.
it.effect("keeps a running Claude monitor typed after many newer monitors end", () =>
Effect.scoped(
Effect.gen(function* () {
const harness = yield* makeWakeHarness;
const now = yield* DateTime.now;
let frameNumber = 0;
const nextUuid = () => `00000000-0000-4000-8000-${String(++frameNumber).padStart(12, "0")}`;
const monitorFrames = (index: number) => {
const taskId = `monitor-task-${index}`;
const toolUseId = `toolu_monitor_${index}`;
const description = `Monitor ${index}`;
return {
taskId,
start: [
claudeSdkFrame({
type: "assistant",
message: {
model: "claude-sonnet-4-6",
id: `msg_monitor_${index}`,
type: "message",
role: "assistant",
content: [
{
type: "tool_use",
id: toolUseId,
name: "Monitor",
input: { description, command: "sleep 8 && echo MONITOR_DONE" },
},
],
stop_reason: null,
stop_sequence: null,
usage: { input_tokens: 1, output_tokens: 1 },
},
parent_tool_use_id: null,
uuid: nextUuid(),
session_id: WAKE_NATIVE_SESSION,
}),
claudeSdkFrame({
type: "system",
subtype: "task_started",
task_id: taskId,
tool_use_id: toolUseId,
description,
is_backgrounded: true,
task_type: "local_bash",
uuid: nextUuid(),
session_id: WAKE_NATIVE_SESSION,
}),
claudeSdkFrame({
type: "user",
message: {
role: "user",
content: [
{
type: "tool_result",
tool_use_id: toolUseId,
content: `Monitor started (task ${taskId}).`,
},
],
},
parent_tool_use_id: null,
uuid: nextUuid(),
session_id: WAKE_NATIVE_SESSION,
tool_use_result: { taskId, persistent: false },
}),
],
end: claudeSdkFrame({
type: "system",
subtype: "task_notification",
task_id: taskId,
tool_use_id: toolUseId,
status: "completed",
output_file: `/tmp/claude-replay/tasks/${taskId}.output`,
summary: `Monitor "${description}" stream ended`,
uuid: nextUuid(),
session_id: WAKE_NATIVE_SESSION,
}),
};
};

yield* harness.runtime.startTurn(
makeClaudeTestTurnInput({
threadId: harness.threadId,
providerThread: harness.providerThread,
now,
attemptId: RunAttemptId.make("attempt-claude-many-monitors"),
text: "Watch the deploy, then re-arm short watches.",
attachments: [],
}),
);
const longRunning = monitorFrames(0);
for (const frame of longRunning.start) {
yield* Queue.offer(harness.sdkMessages, frame);
}
// More newer monitors start and end than any fixed id cap would hold.
for (let index = 1; index <= 65; index++) {
const monitor = monitorFrames(index);
for (const frame of monitor.start) {
yield* Queue.offer(harness.sdkMessages, frame);
}
yield* Queue.offer(harness.sdkMessages, monitor.end);
}
yield* Queue.offer(
harness.sdkMessages,
claudeSdkFrame({
type: "system",
subtype: "background_tasks_changed",
tasks: [
{
task_id: longRunning.taskId,
task_type: "local_bash",
description: "Monitor 0",
},
],
uuid: nextUuid(),
session_id: WAKE_NATIVE_SESSION,
}),
);
yield* Queue.offer(
harness.sdkMessages,
makeResultFrame({ uuid: nextUuid(), result: "Watching." }),
);
yield* awaitUntil(() => harness.terminalEvents().length === 1, "turn terminal");

const roster = providerThreadRosterEvents(harness.events).at(-1)?.providerThread
.pendingBackgroundTasks;
assert.deepEqual(roster, [
{ taskId: longRunning.taskId, kind: "monitor", description: "Monitor 0" },
]);
}).pipe(Effect.provide(Layer.merge(idAllocatorLayer, NodeServices.layer))),
),
);

it.effect("stops background work after the turn settled", () =>
Effect.scoped(
Effect.gen(function* () {
Expand Down
119 changes: 107 additions & 12 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1576,6 +1576,7 @@ const CLAUDE_KNOWN_TOOL_CLASSIFICATIONS: Record<
glob: { itemType: "dynamic_tool", requestKind: "file-read" },
grep: { itemType: "dynamic_tool", requestKind: "file-read" },
ls: { itemType: "dynamic_tool", requestKind: "file-read" },
monitor: { itemType: "dynamic_tool", requestKind: "command" },
multiedit: { itemType: "file_change", requestKind: "file-change" },
notebookedit: { itemType: "file_change", requestKind: "file-change" },
read: { itemType: "dynamic_tool", requestKind: "file-read" },
Expand Down Expand Up @@ -1700,14 +1701,18 @@ function isClaudeOpaqueBackgroundTaskType(taskType: string | null | undefined):
function claudePendingBackgroundTask(input: {
readonly taskId: string;
readonly taskType: string | null;
// Claude runs a Monitor as a local_bash task, so only the Monitor tool call
// that started it tells it apart from a background Bash command.
readonly startedByMonitor: boolean;
readonly description: string | undefined;
}): OrchestrationV2PendingBackgroundTask {
return {
taskId: input.taskId,
kind:
(input.taskType === null
? undefined
: CLAUDE_OPAQUE_BACKGROUND_TASK_KINDS.get(input.taskType)) ?? "background_task",
kind: input.startedByMonitor
? "monitor"
: ((input.taskType === null
? undefined
: CLAUDE_OPAQUE_BACKGROUND_TASK_KINDS.get(input.taskType)) ?? "background_task"),
...(input.description !== undefined && input.description.trim().length > 0
? { description: input.description }
: {}),
Expand Down Expand Up @@ -1742,6 +1747,7 @@ function claudePendingBackgroundTasksFromRoster(

function parseClaudeBackgroundTaskEntry(
entry: unknown,
monitorTasks: ReadonlyMap<string, unknown>,
): OrchestrationV2PendingBackgroundTask | null {
if (entry === null || typeof entry !== "object") {
return null;
Expand All @@ -1761,6 +1767,7 @@ function parseClaudeBackgroundTaskEntry(
return claudePendingBackgroundTask({
taskId,
taskType,
startedByMonitor: monitorTasks.has(taskId),
description: typeof description === "string" ? description : undefined,
});
}
Expand Down Expand Up @@ -2944,6 +2951,68 @@ export function makeClaudeAdapterV2(
const lastKnownOpaqueTasks = yield* Ref.make(
new Map<string, OrchestrationV2PendingBackgroundTask>(),
);
// Live Claude monitors. A Monitor call waits in `calls` until its
// task_started links it to a task, or its tool_result ends the call
// without one. A task stays in `tasks` until it leaves its native
// thread's roster or its task_notification arrives, so a monitor is
// never forgotten while it runs.
const claudeMonitors = yield* Ref.make<{
readonly calls: ReadonlySet<string>;
readonly tasks: ReadonlyMap<
string,
{ readonly toolUseId: string; readonly nativeThreadId: string }
>;
}>({ calls: new Set(), tasks: new Map() });
const trackClaudeMonitorCalls = (message: SDKMessage) => {
const started = claudeToolUseBlocksFromAssistantMessage(message).flatMap((toolUse) =>
toolUse.name === "Monitor" ? [toolUse.id] : [],
);
const returned = [
...claudeToolResultBlocksFromAssistantMessage(message),
...claudeToolResultBlocksFromUserMessage(message),
].map((toolResult) => toolResult.tool_use_id);
if (started.length === 0 && returned.length === 0) return Effect.void;
return Ref.update(claudeMonitors, (current) => {
// A replayed tool_use frame must not reopen a call whose task already started.
const linked = new Set([...current.tasks.values()].map((task) => task.toolUseId));
const opened = started.filter((id) => !linked.has(id) && !current.calls.has(id));
const closed = returned.filter((id) => current.calls.has(id));
if (opened.length === 0 && closed.length === 0) return current;
const calls = new Set([...current.calls, ...opened]);
for (const id of closed) calls.delete(id);
return { ...current, calls };
});
};
/** True when a Monitor call started this task; links the task to that call. */
const isClaudeMonitorTask = (input: {
readonly nativeThreadId: string;
readonly taskId: string;
readonly toolUseId: string | undefined;
}) =>
Ref.modify(claudeMonitors, (current) => {
if (current.tasks.has(input.taskId)) return [true, current] as const;
if (input.toolUseId === undefined || !current.calls.has(input.toolUseId)) {
return [false, current] as const;
}
const calls = new Set(current.calls);
calls.delete(input.toolUseId);
const tasks = new Map(current.tasks).set(input.taskId, {
toolUseId: input.toolUseId,
nativeThreadId: input.nativeThreadId,
});
return [true, { calls, tasks }] as const;
});
/** Drops monitor tasks that ended: gone from the thread's roster, or notified. */
const endClaudeMonitorTasks = (
ended: (taskId: string, task: { readonly nativeThreadId: string }) => boolean,
) =>
Ref.modify(claudeMonitors, (current) => {
const endedIds = [...current.tasks].filter(([taskId, task]) => ended(taskId, task));
if (endedIds.length === 0) return [current.tasks, current] as const;
const tasks = new Map(current.tasks);
for (const [taskId] of endedIds) tasks.delete(taskId);
return [tasks, { ...current, tasks }] as const;
});
const claudeTaskOutcome = (status: "completed" | "failed" | "stopped") =>
status === "completed" ? "completed" : status === "stopped" ? "cancelled" : "failed";
// Reads the roster, so call it before the notification clears the task from it.
Expand Down Expand Up @@ -3285,13 +3354,17 @@ export function makeClaudeAdapterV2(
});

const clearPendingBackgroundTasksForNativeThread = (nativeThreadId: string) =>
Ref.update(pendingBackgroundTasksByNativeThread, (current) => {
if (!current.has(nativeThreadId)) {
return current;
}
const updated = new Map(current);
updated.delete(nativeThreadId);
return updated;
Effect.gen(function* () {
yield* Ref.update(pendingBackgroundTasksByNativeThread, (current) => {
if (!current.has(nativeThreadId)) {
return current;
}
const updated = new Map(current);
updated.delete(nativeThreadId);
return updated;
});
// The thread's process died or its turn failed; those monitors never notify.
yield* endClaudeMonitorTasks((_taskId, task) => task.nativeThreadId === nativeThreadId);
});

// Drop idle wake traffic for a dead native process so it cannot pin
Expand Down Expand Up @@ -4969,8 +5042,21 @@ export function makeClaudeAdapterV2(
return false;
}
const nextTasks: OrchestrationV2PendingBackgroundTask[] = [];
const listedTaskIds = new Set(
roster.flatMap((entry) => {
const taskId =
entry !== null && typeof entry === "object"
? Reflect.get(entry, "task_id")
: undefined;
return typeof taskId === "string" ? [taskId] : [];
}),
);
const monitorTasks = yield* endClaudeMonitorTasks(
(taskId, task) =>
task.nativeThreadId === input.nativeThreadId && !listedTaskIds.has(taskId),
);
for (const entry of roster) {
const task = parseClaudeBackgroundTaskEntry(entry);
const task = parseClaudeBackgroundTaskEntry(entry, monitorTasks);
if (task !== null) {
nextTasks.push(task);
}
Expand All @@ -4992,12 +5078,18 @@ export function makeClaudeAdapterV2(
claudePendingBackgroundTask({
taskId: message.task_id,
taskType: claudeTaskTypeFromSdkMessage(message),
startedByMonitor: yield* isClaudeMonitorTask({
nativeThreadId: input.nativeThreadId,
taskId: message.task_id,
toolUseId: message.tool_use_id,
}),
description:
typeof message.description === "string" ? message.description : undefined,
}),
);
rosterChanged = true;
} else if (message.type === "system" && message.subtype === "task_notification") {
yield* endClaudeMonitorTasks((taskId) => taskId === message.task_id);
const removed = yield* clearPendingBackgroundTask(
input.nativeThreadId,
message.task_id,
Expand Down Expand Up @@ -5054,6 +5146,9 @@ export function makeClaudeAdapterV2(
}

const message = input.message;
// Before any routing: a Monitor started during an idle wake turn
// reports its task before the drain replays the tool call.
yield* trackClaudeMonitorCalls(message);
if (message.type === "rate_limit_event") {
const rateLimitInfo = message.rate_limit_info;
if (!rateLimitInfo) return;
Expand Down
Loading
Loading