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
33 changes: 22 additions & 11 deletions apps/server/src/claudeHistoryWorker.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,8 @@
import { forkSession, getSessionMessages } from "@anthropic-ai/claude-agent-sdk";
import {
forkSession,
getSessionMessages,
getSubagentMessages,
} from "@anthropic-ai/claude-agent-sdk";
import * as Schema from "effect/Schema";

// A separate process gives SDK history helpers the provider's environment without
Expand All @@ -11,8 +15,10 @@ import * as Schema from "effect/Schema";
const decodeHistoryOptions = Schema.decodeSync(
Schema.fromJsonString(
Schema.Struct({
agentId: Schema.optionalKey(Schema.String),
dir: Schema.optionalKey(Schema.String),
includeSystemMessages: Schema.optionalKey(Schema.Boolean),
limit: Schema.optionalKey(Schema.Number),
upToMessageId: Schema.optionalKey(Schema.String),
}),
),
Expand All @@ -23,15 +29,20 @@ export async function runClaudeHistoryWorker(
sessionId: string | undefined,
rawOptions: string | undefined,
): Promise<void> {
const options = decodeHistoryOptions(rawOptions ?? "{}");
const { agentId, ...options } = decodeHistoryOptions(rawOptions ?? "{}");
if (!sessionId) throw new Error("Claude history session id is required.");
const result =
method === "getSessionMessages"
? await getSessionMessages(sessionId, options)
: method === "forkSession"
? await forkSession(sessionId, options)
: (() => {
throw new Error("Unknown Claude history operation.");
})();
process.stdout.write(JSON.stringify(result));
const run = () => {
switch (method) {
case "getSessionMessages":
return getSessionMessages(sessionId, options);
case "forkSession":
return forkSession(sessionId, options);
case "getSubagentMessages":
if (!agentId) throw new Error("Claude subagent id is required.");
return getSubagentMessages(sessionId, agentId, options);
default:
throw new Error("Unknown Claude history operation.");
}
};
process.stdout.write(JSON.stringify(await run()));
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,125 @@
import * as NodeServices from "@effect/platform-node/NodeServices";
import { ProviderSessionId, ThreadId } from "@t3tools/contracts";
import { assert, describe, it } from "@effect/vitest";
import * as Effect from "effect/Effect";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as Path from "effect/Path";
import * as Schema from "effect/Schema";

import * as ProviderEventLoggers from "../../provider/ProviderEventLoggers.ts";
import { ClaudeAgentSdkQueryRunner, layerQueryRunner } from "./ClaudeAdapterV2.ts";

const encodeJson = Schema.encodeSync(Schema.fromJsonString(Schema.Unknown));

const sessionId = "6f0c6a52-1d5e-4f39-9a51-3c1f6f0b2a11";
const userMessageId = "0b8c2a47-5f0e-4a8f-9a55-0c6f7f6d1e01";
const assistantMessageId = "4a1d9e3c-7b2f-4c6a-8e10-2d5b9f3c7a02";
const subagentId = "a1b2c3d4e5f6";
const subagentToolUseId = "toolu_01SubagentLaunch";

// The SDK's session helpers find transcripts under CLAUDE_CONFIG_DIR, read from
// the process environment and cached on first use. A provider instance with its
// own Claude home keeps its transcripts there, not wherever the server's own
// environment points, so the runner must reach them through that instance's
// environment. The server's home here is an empty directory.
const writeHomes = Effect.gen(function* () {
const fs = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const serverConfigDir = yield* fs.makeTempDirectory({ prefix: "t3-claude-server-" });
process.env.CLAUDE_CONFIG_DIR = serverConfigDir;
const configDir = yield* fs.makeTempDirectory({ prefix: "t3-claude-provider-" });
const projectDir = path.join(configDir, "projects", "-tmp-t3-claude-project");
const subagentsDir = path.join(projectDir, sessionId, "subagents");
const base = { sessionId, cwd: "/tmp/t3-claude-project", timestamp: "2026-10-08T00:00:00.000Z" };
yield* fs.makeDirectory(subagentsDir, { recursive: true });
yield* fs.writeFileString(
path.join(projectDir, `${sessionId}.jsonl`),
[
{
...base,
type: "user",
uuid: userMessageId,
parentUuid: null,
message: { role: "user", content: "hello" },
},
{
...base,
type: "assistant",
uuid: assistantMessageId,
parentUuid: userMessageId,
message: { role: "assistant", content: [{ type: "text", text: "hi" }] },
},
]
.map((line) => encodeJson(line))
.join("\n") + "\n",
);
yield* fs.writeFileString(
path.join(subagentsDir, `agent-${subagentId}.jsonl`),
encodeJson({
...base,
type: "user",
uuid: "9c3e1f2a-6d4b-4e8a-b7c1-5f2e8d9a0b03",
parentUuid: null,
isSidechain: true,
agentId: subagentId,
message: { role: "user", content: "subagent task" },
}) + "\n",
);
yield* fs.writeFileString(
path.join(subagentsDir, `agent-${subagentId}.meta.json`),
encodeJson({ agentType: "general-purpose", toolUseId: subagentToolUseId }),
);
return { serverConfigDir, configDir, projectDir };
});

const runnerLayer = layerQueryRunner.pipe(
Layer.provide(
Layer.succeed(
ProviderEventLoggers.ProviderEventLoggers,
ProviderEventLoggers.NoOpProviderEventLoggers,
),
),
Layer.provideMerge(NodeServices.layer),
);

describe("ClaudeAgentSdkQueryRunner provider Claude home", () => {
it.effect("forks a session stored in the provider's Claude home", () =>
Effect.gen(function* () {
const fs = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const { serverConfigDir, configDir, projectDir } = yield* writeHomes;
const runner = yield* ClaudeAgentSdkQueryRunner;

const forked = yield* runner.forkSession({
sessionId,
options: { upToMessageId: assistantMessageId },
environment: { ...process.env, CLAUDE_CONFIG_DIR: configDir },
threadId: ThreadId.make("thread-claude-fork-provider-home"),
providerSessionId: ProviderSessionId.make("provider-session-claude-fork"),
});

assert.notStrictEqual(forked.sessionId, sessionId);
assert.isTrue(yield* fs.exists(path.join(projectDir, `${forked.sessionId}.jsonl`)));
assert.deepStrictEqual(yield* fs.readDirectory(serverConfigDir), []);
}).pipe(Effect.provide(runnerLayer)),
);

it.effect("finds a subagent's launch in the provider's Claude home", () =>
Effect.gen(function* () {
const { configDir } = yield* writeHomes;
const runner = yield* ClaudeAgentSdkQueryRunner;

const toolUseId = yield* runner.subagentLaunchToolUseId({
sessionId,
agentId: subagentId,
dir: null,
environment: { ...process.env, CLAUDE_CONFIG_DIR: configDir },
threadId: ThreadId.make("thread-claude-subagent-provider-home"),
providerSessionId: ProviderSessionId.make("provider-session-claude-subagent"),
});

assert.strictEqual(toolUseId, subagentToolUseId);
}).pipe(Effect.provide(runnerLayer)),
);
});
Original file line number Diff line number Diff line change
Expand Up @@ -1641,12 +1641,7 @@ describe("ClaudeAdapterV2 native fork", () => {
prefix: "t3-claude-v2-fork-attachments-",
});
const openedQueries: Array<ClaudeAdapterV2.ClaudeAgentSdkQueryOpenInput> = [];
const forkCalls: Array<{
readonly sessionId: string;
readonly options: unknown;
readonly threadId: ThreadId;
readonly providerSessionId: ProviderSessionId;
}> = [];
const forkCalls: Array<ClaudeAdapterV2.ClaudeAgentSdkSessionForkInput> = [];
const adapter = ClaudeAdapterV2.makeClaudeAdapterV2({
instanceId: ClaudeAdapterV2.CLAUDE_DEFAULT_INSTANCE_ID,
settings: DEFAULT_CLAUDE_SETTINGS,
Expand Down Expand Up @@ -1739,6 +1734,7 @@ describe("ClaudeAdapterV2 native fork", () => {
dir: "/workspace",
upToMessageId: "assistant-message-cursor",
},
environment: {},
threadId: targetThreadId,
providerSessionId,
},
Expand Down
102 changes: 90 additions & 12 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,10 +8,8 @@ import { isWorkspaceImagePreviewPath } from "@t3tools/shared/filePreview";
import { normalizeClaudeTurnTokenUsage } from "../../provider/ClaudeTurnTokenUsage.ts";
import {
type CanUseTool,
forkSession as forkClaudeSession,
type ForkSessionOptions,
type ForkSessionResult,
getSubagentMessages,
query,
type Options as ClaudeQueryOptions,
type PermissionMode,
Expand All @@ -32,7 +30,7 @@ import type {
WebSearchOutput,
} from "@anthropic-ai/claude-agent-sdk/sdk-tools";
import { parseCliArgs } from "@t3tools/shared/cliArgs";
import { HostProcessEnvironment } from "@t3tools/shared/hostProcess";
import { HostProcessEnvironment, HostProcessIsExecutable } from "@t3tools/shared/hostProcess";
import { applyClaudePromptEffortPrefix } from "@t3tools/shared/model";
import {
CLAUDE_RESUME_COMPACTION_NEVER_ANSWER,
Expand Down Expand Up @@ -85,8 +83,10 @@ import * as Queue from "effect/Queue";
import * as Ref from "effect/Ref";
import * as Schema from "effect/Schema";
import * as Stream from "effect/Stream";
import { ChildProcess, ChildProcessSpawner } from "effect/process";

import { resolveAttachmentPath } from "../../attachmentStore.ts";
import { spawnAndCollect } from "../../provider/providerSnapshot.ts";
import { resolveClaudeSdkExecutablePath } from "../../provider/Drivers/ClaudeExecutable.ts";
import { planClaudeSkillDispatch } from "../../provider/Drivers/ClaudeSkillDispatch.ts";
import { discoverClaudeSkills } from "../../provider/Drivers/ClaudeSkills.ts";
Expand Down Expand Up @@ -388,6 +388,8 @@ export class ClaudeAgentSdkQueryRunner extends Context.Service<
export interface ClaudeAgentSdkSessionForkInput {
readonly sessionId: string;
readonly options: ForkSessionOptions;
/** The provider instance's environment; its CLAUDE_CONFIG_DIR locates the transcript. */
readonly environment: NodeJS.ProcessEnv;
readonly threadId: ThreadId;
readonly providerSessionId: OrchestrationV2ProviderSession["id"];
}
Expand All @@ -396,6 +398,8 @@ export interface ClaudeAgentSdkSubagentLookupInput {
readonly sessionId: string;
readonly agentId: string;
readonly dir: string | null;
/** The provider instance's environment; its CLAUDE_CONFIG_DIR locates the transcript. */
readonly environment: NodeJS.ProcessEnv;
readonly threadId: ThreadId;
readonly providerSessionId: OrchestrationV2ProviderSession["id"];
}
Expand Down Expand Up @@ -613,15 +617,75 @@ export function makeClaudeAgentSdkProtocolLogger(input: {
};
}

const encodeHistoryOptions = Schema.encodeSync(Schema.fromJsonString(Schema.Unknown));
const decodeHistoryFork = Schema.decodeSync(
Schema.fromJsonString(Schema.Struct({ sessionId: Schema.String })),
);
const decodeHistorySubagentMessages = Schema.decodeSync(
Schema.fromJsonString(
Schema.Array(Schema.Struct({ parent_tool_use_id: Schema.NullOr(Schema.String) })),
),
);

export const layerQueryRunner: Layer.Layer<
ClaudeAgentSdkQueryRunner,
never,
Crypto.Crypto | ProviderEventLoggers.ProviderEventLoggers
| Crypto.Crypto
| ProviderEventLoggers.ProviderEventLoggers
| ChildProcessSpawner.ChildProcessSpawner
| Path.Path
> = Layer.effect(
ClaudeAgentSdkQueryRunner,
Effect.gen(function* () {
const crypto = yield* Crypto.Crypto;
const { native: nativeEventLogger } = yield* ProviderEventLoggers.ProviderEventLoggers;
const spawner = yield* ChildProcessSpawner.ChildProcessSpawner;
const path = yield* Path.Path;
const hostIsExecutable = yield* HostProcessIsExecutable;

// The SDK's session helpers resolve the Claude home from process.env and
// cache it, so in the server they would read the server's home, not the
// provider instance's. They run in a child given the instance's
// environment instead. The single-executable has no sibling script and no
// Node to run one with, so it hosts the worker as a hidden subcommand.
const runHistoryWorker = (
method: "forkSession" | "getSubagentMessages",
sessionId: string,
options: object,
environment: NodeJS.ProcessEnv,
) =>
Effect.gen(function* () {
const workerArguments = hostIsExecutable
? ["__claude-history"]
: [
yield* path.fromFileUrl(
new URL(
import.meta.url.endsWith(".ts")
? "../../claude-history-worker.ts"
: "./claude-history-worker.mjs",
import.meta.url,
),
),
];
const result = yield* spawnAndCollect(
process.execPath,
ChildProcess.make(
process.execPath,
[...workerArguments, method, sessionId, encodeHistoryOptions(options)],
{ env: { ...environment, ELECTRON_RUN_AS_NODE: "1" } },
),
).pipe(
Effect.timeout("30 seconds"),
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
);
if (result.code !== 0) {
return yield* queryRunnerError(
result.stderr.trim() || `Claude history worker exited with ${result.code}.`,
method,
);
}
return result.stdout;
}).pipe(Effect.mapError((cause) => queryRunnerError(cause, method)));

return ClaudeAgentSdkQueryRunner.of({
allocateSessionId: crypto.randomUUIDv4.pipe(
Expand Down Expand Up @@ -800,8 +864,14 @@ export const layerQueryRunner: Layer.Layer<
options: input.options,
},
});
const result = yield* Effect.tryPromise({
try: () => forkClaudeSession(input.sessionId, input.options),
const stdout = yield* runHistoryWorker(
"forkSession",
input.sessionId,
input.options,
input.environment,
);
const result = yield* Effect.try({
try: () => decodeHistoryFork(stdout),
catch: (cause) => queryRunnerError(cause, "forkSession"),
});
yield* logProtocolEvent({
Expand Down Expand Up @@ -834,12 +904,18 @@ export const layerQueryRunner: Layer.Layer<
});
// The CLI stamps every message of a subagent's transcript with the
// tool call that launched it; one message is enough.
const messages = yield* Effect.tryPromise({
try: () =>
getSubagentMessages(input.sessionId, input.agentId, {
...(input.dir === null ? {} : { dir: input.dir }),
limit: 1,
}),
const stdout = yield* runHistoryWorker(
"getSubagentMessages",
input.sessionId,
{
agentId: input.agentId,
...(input.dir === null ? {} : { dir: input.dir }),
limit: 1,
},
input.environment,
);
const messages = yield* Effect.try({
try: () => decodeHistorySubagentMessages(stdout),
catch: (cause) => queryRunnerError(cause, "getSubagentMessages"),
});
const toolUseId = messages[0]?.parent_tool_use_id ?? null;
Expand Down Expand Up @@ -4106,6 +4182,7 @@ export function makeClaudeAdapterV2(
sessionId: resume.nativeThreadId,
agentId: resume.taskId,
dir: resume.context.input.runtimePolicy.cwd,
environment: adapterOptions.environment,
threadId: resume.context.input.threadId,
providerSessionId: input.providerSessionId,
})
Expand Down Expand Up @@ -8004,6 +8081,7 @@ export function makeClaudeAdapterV2(
const forked = yield* queryRunner.forkSession({
sessionId: sourceNativeThreadId,
options: forkOptions,
environment: adapterOptions.environment,
threadId: forkInput.targetThreadId,
providerSessionId: input.providerSessionId,
});
Expand Down
Loading