Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -478,6 +478,11 @@ function taskLinkageActivityFields(payload: Record<string, unknown>): Record<str
"outputFile",
"prompt",
"agentPath",
"cwd",
"worktreeBranch",
"isolation",
"linesAdded",
"linesRemoved",
"timelineBypass",
"typedUsage",
"status",
Expand Down
103 changes: 103 additions & 0 deletions apps/server/src/provider/Layers/ClaudeAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2324,6 +2324,109 @@ describe("ClaudeAdapterLive", () => {
);
});

it.effect("reports where each subagent works: session cwd, isolated worktree, or remote", () => {
const harness = makeHarness();
return Effect.gen(function* () {
const adapter = yield* ClaudeAdapter;
const taskEventsFiber = yield* adapter.streamEvents.pipe(
Stream.filter((event) => event.type === "task.started" || event.type === "task.updated"),
Stream.take(4),
Stream.runCollect,
Effect.forkChild,
);
const session = yield* adapter.startSession({
threadId: THREAD_ID,
provider: ProviderDriverKind.make("claudeAgent"),
runtimeMode: "full-access",
cwd: "/tmp/claude-workspace-project",
});
yield* adapter.sendTurn({ threadId: session.threadId, input: "delegate", attachments: [] });

const launch = (index: number, id: string, input: Record<string, unknown>) =>
harness.query.emit({
type: "stream_event",
session_id: CLAUDE_ORIGINAL_SESSION_ID,
uuid: `stream-${id}`,
parent_tool_use_id: null,
event: {
type: "content_block_start",
index,
content_block: { type: "tool_use", id, name: "Agent", input },
},
} as unknown as SDKMessage);
const started = (taskId: string, toolUseId: string) =>
harness.query.emit({
type: "system",
subtype: "task_started",
task_id: taskId,
description: taskId,
task_type: "local_agent",
tool_use_id: toolUseId,
uuid: `started-${taskId}`,
session_id: CLAUDE_ORIGINAL_SESSION_ID,
} as unknown as SDKMessage);

launch(0, "tool-shared", { description: "Shared", prompt: "look" });
started("agent-shared", "tool-shared");
launch(1, "tool-isolated", {
description: "Isolated",
prompt: "edit",
isolation: "worktree",
});
started("agent-isolated", "tool-isolated");
launch(2, "tool-remote", { description: "Remote", prompt: "e2e", isolation: "remote" });
started("agent-remote", "tool-remote");
harness.query.emit({
type: "user",
session_id: CLAUDE_ORIGINAL_SESSION_ID,
uuid: "user-isolated-result",
parent_tool_use_id: null,
message: {
role: "user",
content: [{ type: "tool_result", tool_use_id: "tool-isolated", content: "done" }],
},
tool_use_result: {
status: "completed",
prompt: "edit",
worktreePath: "/tmp/claude-workspace-project/.claude/worktrees/agent-a91f",
worktreeBranch: "worktree-agent-a91f",
toolStats: { linesAdded: 42, linesRemoved: 7 },
},
} as unknown as SDKMessage);

const [shared, isolated, remote, reported] = Array.from(yield* Fiber.join(taskEventsFiber));
assert.equal(shared?.type, "task.started");
if (shared?.type === "task.started") {
assert.equal(shared.payload.cwd, "/tmp/claude-workspace-project");
assert.notProperty(shared.payload, "isolation");
}
if (isolated?.type === "task.started") {
assert.equal(isolated.payload.isolation, "worktree");
assert.notProperty(isolated.payload, "cwd");
}
if (remote?.type === "task.started") {
assert.equal(remote.payload.isolation, "remote");
assert.notProperty(remote.payload, "cwd");
}
assert.equal(reported?.type, "task.updated");
if (reported?.type === "task.updated") {
assert.equal(reported.payload.taskId, "agent-isolated");
assert.notProperty(reported.payload, "status");
assert.equal(
reported.payload.cwd,
"/tmp/claude-workspace-project/.claude/worktrees/agent-a91f",
);
assert.equal(reported.payload.worktreeBranch, "worktree-agent-a91f");
assert.equal(reported.payload.isolation, "worktree");
assert.equal(reported.payload.linesAdded, 42);
assert.equal(reported.payload.linesRemoved, 7);
}
}).pipe(
Effect.provideService(Random.Random, makeDeterministicRandomService()),
Effect.provide(harness.layer),
);
});

it.effect("reads a subagent transcript from Claude history by its task id", () => {
const calls: Array<{ sessionId: string; agentId: string; dir: string | undefined }> = [];
const harness = makeHarness({
Expand Down
91 changes: 91 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,13 @@ interface ClaudeTaskAgentState {
* assistant snapshots (authoritative API model). */
model: string | undefined;
effort: string | undefined;
/** Agent tool isolation. A worktree agent's path and branch arrive only with
* the tool result, so its cwd stays unset until then. */
isolation?: "worktree" | "remote" | undefined;
cwd?: string | undefined;
worktreeBranch?: string | undefined;
linesAdded?: number | undefined;
linesRemoved?: number | undefined;
}

/**
Expand Down Expand Up @@ -1387,9 +1394,25 @@ function taskLinkageFor(
...(agent.toolUseId ? { toolUseId: agent.toolUseId } : {}),
...(agent.workflowName ? { workflowName: agent.workflowName } : {}),
...(agent.runHandles ? { runHandles: agent.runHandles } : {}),
...(agent.isolation ? { isolation: agent.isolation } : {}),
...(agent.cwd ? { cwd: agent.cwd } : {}),
...(agent.worktreeBranch ? { worktreeBranch: agent.worktreeBranch } : {}),
...(agent.linesAdded !== undefined ? { linesAdded: agent.linesAdded } : {}),
...(agent.linesRemoved !== undefined ? { linesRemoved: agent.linesRemoved } : {}),
};
}

function readAgentIsolation(
input: Record<string, unknown> | undefined,
): "worktree" | "remote" | undefined {
const isolation = input?.isolation;
return isolation === "worktree" || isolation === "remote" ? isolation : undefined;
}

function readLineCount(value: unknown): number | undefined {
return typeof value === "number" && Number.isInteger(value) && value >= 0 ? value : undefined;
}

const WORKFLOW_PHASE_CAP = 64;
const WORKFLOW_AGENT_CAP = 100;

Expand Down Expand Up @@ -3404,6 +3427,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
};
const existing = context.taskAgents.get(workflowTaskId);
context.taskAgents.set(workflowTaskId, {
...existing,
taskId: workflowTaskId,
toolUseId: existing?.toolUseId ?? tool.itemId,
description: existing?.description,
Expand All @@ -3419,6 +3443,51 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
}
}

// A synchronous Agent tool result reports the agent's isolated worktree
// and line counts. The task has usually settled by now, so re-announce
// its linkage in a status-free update.
const launchedAgent =
!toolResult.isError && toolUseResult
? Array.from(context.taskAgents.values()).find((agent) => agent.toolUseId === tool.itemId)
: undefined;
if (launchedAgent && toolUseResult) {
const worktreePath = trimmedString(toolUseResult.worktreePath);
const worktreeBranch = trimmedString(toolUseResult.worktreeBranch);
const toolStats =
typeof toolUseResult.toolStats === "object" && toolUseResult.toolStats !== null
? (toolUseResult.toolStats as Record<string, unknown>)
: undefined;
const linesAdded = readLineCount(toolStats?.linesAdded);
const linesRemoved = readLineCount(toolStats?.linesRemoved);
if (worktreePath || linesAdded !== undefined || linesRemoved !== undefined) {
if (worktreePath) {
launchedAgent.cwd = worktreePath;
launchedAgent.worktreeBranch = worktreeBranch;
}
launchedAgent.linesAdded = linesAdded ?? launchedAgent.linesAdded;
launchedAgent.linesRemoved = linesRemoved ?? launchedAgent.linesRemoved;
const workspaceStamp = yield* makeEventStamp();
yield* offerRuntimeEvent({
type: "task.updated",
eventId: workspaceStamp.eventId,
provider: PROVIDER,
createdAt: workspaceStamp.createdAt,
threadId: context.session.threadId,
...(context.turnState ? { turnId: asCanonicalTurnId(context.turnState.turnId) } : {}),
payload: {
taskId: RuntimeTaskId.make(launchedAgent.taskId),
...taskLinkageFor(context.taskAgents, launchedAgent.taskId),
},
providerRefs: nativeProviderRefs(context, { providerItemId: tool.itemId }),
raw: {
source: "claude.sdk.message",
method: "claude/user",
payload: message,
},
});
}
}

if (
!toolResult.isError &&
applyClaudeTaskToolResult(context.claudeTasks, tool, toolUseResult)
Expand Down Expand Up @@ -3847,9 +3916,28 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
(typeof rawLaunchEffort === "number" && Number.isFinite(rawLaunchEffort)
? String(rawLaunchEffort)
: context.currentEffort);
// Where the agent works: an isolated agent's worktree is only known
// from the Agent tool result; a nested agent inherits its owner's
// workspace; anything else runs in the session's cwd.
const isolation = readAgentIsolation(launchInput);
const owningAgent = owningAgentId
? Array.from(context.taskAgents.values()).find(
(agent) => agent.toolUseId === owningAgentId || agent.taskId === owningAgentId,
)
: undefined;
const workspace = isolation
? { isolation }
: owningAgentId
? {
isolation: owningAgent?.isolation,
cwd: owningAgent?.cwd,
worktreeBranch: owningAgent?.worktreeBranch,
}
: { cwd: trimmedString(context.session.cwd ?? undefined) };
// Remember the agent identity so every later task.* payload for this
// taskId is self-describing (identity must survive activity retention).
context.taskAgents.set(message.task_id, {
...workspace,
taskId: message.task_id,
toolUseId: message.tool_use_id,
description: message.description,
Expand Down Expand Up @@ -3883,6 +3971,9 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
...(effort ? { effort } : {}),
...(message.tool_use_id ? { toolUseId: message.tool_use_id } : {}),
...(message.workflow_name ? { workflowName: message.workflow_name } : {}),
...(workspace.isolation ? { isolation: workspace.isolation } : {}),
...(workspace.cwd ? { cwd: workspace.cwd } : {}),
...(workspace.worktreeBranch ? { worktreeBranch: workspace.worktreeBranch } : {}),
},
});
return;
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/provider/Layers/CodexAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1345,6 +1345,7 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => {
yield* runtime.emit(
childEvent("evt-child-start", "collabAgent/started", {
prompt: ` ${"x".repeat(5000)} `,
cwd: "/workspace/child-worktree",
}),
);
yield* runtime.emit(
Expand All @@ -1364,6 +1365,7 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => {
NodeAssert.ok(started?.type === "task.started");
NodeAssert.equal(started.payload.prompt?.length, 4000);
NodeAssert.ok(started.payload.prompt?.startsWith("xxx"));
NodeAssert.equal(started.payload.cwd, "/workspace/child-worktree");
NodeAssert.ok(errored?.type === "task.updated");
NodeAssert.deepStrictEqual(
{ status: errored.payload.status, error: errored.payload.error },
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/provider/Layers/CodexAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1241,6 +1241,7 @@ function mapCollabAgentEvent(
const model = typeof payload.model === "string" ? payload.model.trim() : "";
const effort = typeof payload.effort === "string" ? payload.effort.trim() : "";
const prompt = boundSubagentPrompt(payload.prompt, SUBAGENT_PROMPT_CHAR_LIMIT);
const cwd = typeof payload.cwd === "string" ? payload.cwd.trim() : "";
// Identity repeated on every status patch so rows are self-describing when
// the start row ages out of activity retention (review finding: a
// reconstructed agent had a UUID name and no role/path).
Expand All @@ -1250,6 +1251,7 @@ function mapCollabAgentEvent(
...(model ? { model } : {}),
...(effort ? { effort } : {}),
...(agentPath ? { agentPath } : {}),
...(cwd ? { cwd } : {}),
timelineBypass: true,
} as const;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -236,6 +236,51 @@ describe("CodexSessionRuntime collab integration", () => {
}).pipe(Effect.scoped, Effect.provide(NodeServices.layer)),
);

it.effect("carries the child thread's reported cwd on its agent identity", () =>
Effect.gen(function* () {
const spawned = capturedSpawnedThread();
const script = {
rootThreadId: ROOT,
notifications: [
{
...spawned,
params: { thread: { ...spawned.params.thread, cwd: "/workspace/child-worktree" } },
},
],
childResumeSnapshots: { [CHILD_A]: { model: "gpt-5.6-luna", reasoningEffort: "low" } },
};
// @effect-diagnostics-next-line preferSchemaOverJson:off
NodeFS.writeFileSync(scriptPath, JSON.stringify(script), "utf8");
yield* Effect.addFinalizer(() =>
Effect.sync(() => NodeFS.rmSync(scriptPath, { force: true })),
);

const runtime = yield* makeCodexSessionRuntime({
threadId: ThreadId.make("thread-collab-child-cwd"),
binaryPath: peerPath,
cwd: NodeOS.tmpdir(),
runtimeMode: "full-access",
environment: { ...process.env, T3_CODEX_COLLAB_SCRIPT: scriptPath },
});
const startedFiber = yield* runtime.events.pipe(
Stream.filter((event) => event.method === "collabAgent/started"),
Stream.take(1),
Stream.runCollect,
Effect.forkScoped,
);

yield* runtime.start();
yield* runtime.sendTurn({ input: "start one child" });
const [started] = Array.from(yield* Fiber.join(startedFiber));
assert.deepInclude(started?.payload, {
agentThreadId: CHILD_A,
cwd: "/workspace/child-worktree",
});

yield* runtime.close;
}).pipe(Effect.scoped, Effect.provide(NodeServices.layer)),
);

it.effect("keeps child settings and reroutes newer than the resume snapshot", () =>
Effect.gen(function* () {
const statusChanged = wireFixture.notifications.find(
Expand Down
5 changes: 5 additions & 0 deletions apps/server/src/provider/Layers/CodexSessionRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -970,6 +970,8 @@ interface CollabChildAgentState {
readonly agentPath: string | undefined;
readonly depth: number | undefined;
readonly parentThreadId: string | undefined;
/** Working directory the child thread reported in its thread/started. */
readonly cwd: string | undefined;
/**
* Parent canonical turn active when the child registered. Stamped on every
* synthetic collabAgent/* event so clients can batch a fleet by its spawn
Expand All @@ -995,6 +997,7 @@ function collabChildIdentity(
...(child.nickname ? { nickname: child.nickname } : {}),
...(child.role ? { role: child.role } : {}),
...(child.agentPath ? { agentPath: child.agentPath } : {}),
...(child.cwd ? { cwd: child.cwd } : {}),
...(metadata?.model ? { model: metadata.model } : {}),
...(metadata?.effort ? { effort: metadata.effort } : {}),
};
Expand Down Expand Up @@ -1665,6 +1668,7 @@ export const makeCodexSessionRuntime = (
depth: spawn.depth ?? existingChild?.depth,
parentThreadId:
spawn.parentThreadId ?? thread.parentThreadId ?? existingChild?.parentThreadId,
cwd: nonEmptyMetadataValue(thread.cwd) ?? existingChild?.cwd,
spawnTurnId,
};
yield* Ref.update(collabChildAgentsRef, (current) => {
Expand Down Expand Up @@ -1718,6 +1722,7 @@ export const makeCodexSessionRuntime = (
agentPath: existing?.agentPath ?? item.agentPath,
depth: existing?.depth,
parentThreadId: existing?.parentThreadId,
cwd: existing?.cwd,
spawnTurnId: existing ? existing.spawnTurnId : activitySpawnTurnId,
});
return next;
Expand Down
Loading