Skip to content

Commit 67d63c0

Browse files
committed
feat(server): read Codex and Claude agent history
1 parent 859304b commit 67d63c0

26 files changed

Lines changed: 1337 additions & 6 deletions

‎apps/server/integration/orphanedProviderSessionStartup.integration.test.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -123,6 +123,7 @@ const startupDependencies = Layer.mergeAll(
123123
assertConversationRollbackSupported: () => Effect.die("unused"),
124124
getInstanceInfo: () => Effect.die("unused"),
125125
rollbackConversation: () => Effect.die("unused"),
126+
getAgentHistory: () => Effect.die("unused"),
126127
uploadFeedback: () => Effect.die("unused"),
127128
streamEvents: Stream.empty,
128129
}),

‎apps/server/src/auth/RpcAuthorization.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ type WsRpcMethod = RpcGroup.Rpcs<typeof WsRpcGroup>["_tag"];
2323
export const RPC_REQUIRED_SCOPES = {
2424
[ORCHESTRATION_WS_METHODS.dispatchCommand]: AuthOrchestrationOperateScope,
2525
[ORCHESTRATION_WS_METHODS.getWorkflowScript]: AuthOrchestrationReadScope,
26+
[ORCHESTRATION_WS_METHODS.getAgentHistory]: AuthOrchestrationReadScope,
2627
[ORCHESTRATION_WS_METHODS.getTurnDiff]: AuthOrchestrationReadScope,
2728
[ORCHESTRATION_WS_METHODS.getFullThreadDiff]: AuthOrchestrationReadScope,
2829
[ORCHESTRATION_WS_METHODS.searchThreads]: AuthOrchestrationReadScope,

‎apps/server/src/orchestration/Layers/CheckpointReactor.test.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -137,6 +137,7 @@ function createProviderServiceHarness(
137137
},
138138
}),
139139
rollbackConversation,
140+
getAgentHistory: () => unsupported(),
140141
uploadFeedback: () => unsupported(),
141142
get streamEvents() {
142143
return Stream.fromPubSub(runtimeEventPubSub);

‎apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -390,6 +390,7 @@ describe("ProviderCommandReactor", () => {
390390
});
391391
},
392392
rollbackConversation: () => unsupported(),
393+
getAgentHistory: () => unsupported(),
393394
uploadFeedback: () => unsupported(),
394395
get streamEvents() {
395396
return Stream.fromPubSub(runtimeEventPubSub);

‎apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -141,6 +141,7 @@ function createProviderServiceHarness() {
141141
});
142142
},
143143
rollbackConversation: () => unsupported(),
144+
getAgentHistory: () => unsupported(),
144145
uploadFeedback: () => unsupported(),
145146
get streamEvents() {
146147
return Stream.fromPubSub(runtimeEventPubSub).pipe(

‎apps/server/src/provider/Layers/ClaudeAdapter.test.ts‎

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7279,3 +7279,48 @@ describe("ClaudeAdapterLive", () => {
72797279
);
72807280
});
72817281
});
7282+
7283+
for (const source of ["configured", "inherited"] as const) {
7284+
it.effect(`reads child history from the ${source} Claude home without starting a session`, () => {
7285+
const home = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "claude-history-adapter-"));
7286+
const sessionId = "0945fcd9-8a8d-450a-b8cf-4ebd0ef17468";
7287+
const directory = NodePath.join(home, "projects", "-workspace", sessionId, "subagents");
7288+
NodeFS.mkdirSync(directory, { recursive: true });
7289+
NodeFS.writeFileSync(
7290+
NodePath.join(directory, "agent-child.jsonl"),
7291+
JSON.stringify({
7292+
type: "user",
7293+
uuid: "9eb84969-819b-4d5b-854f-327474810b34",
7294+
parentUuid: null,
7295+
isSidechain: true,
7296+
sessionId,
7297+
agentId: "child",
7298+
message: { role: "user", content: "Saved task" },
7299+
}) + "\n",
7300+
);
7301+
const harness = makeHarness({
7302+
claudeConfig: { homePath: source === "configured" ? home : "" },
7303+
environment: {
7304+
...process.env,
7305+
CLAUDE_CONFIG_DIR: source === "configured" ? "/unused-claude-home" : home,
7306+
},
7307+
});
7308+
return Effect.gen(function* () {
7309+
const adapter = yield* ClaudeAdapter;
7310+
assert.isDefined(adapter.getAgentHistory);
7311+
const result = yield* adapter.getAgentHistory!({
7312+
threadId: ThreadId.make("stopped-history-thread"),
7313+
agentId: "child",
7314+
offset: 0,
7315+
resumeCursor: { resume: sessionId },
7316+
});
7317+
assert.equal(result.status, "ready");
7318+
assert.equal(result.entries[0]?.detail, "Saved task");
7319+
assert.isUndefined(harness.getLastCreateQueryInput());
7320+
assert.deepStrictEqual(yield* adapter.listSessions(), []);
7321+
}).pipe(
7322+
Effect.provide(harness.layer),
7323+
Effect.ensuring(Effect.sync(() => NodeFS.rmSync(home, { recursive: true, force: true }))),
7324+
);
7325+
});
7326+
}

‎apps/server/src/provider/Layers/ClaudeAdapter.ts‎

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,9 @@
66
*
77
* @module ClaudeAdapterLive
88
*/
9+
import * as NodeOS from "node:os";
10+
import { readClaudeAgentHistory } from "./claudeAgentHistory.ts";
11+
import { expandHomePath } from "../../pathExpansion.ts";
912
import {
1013
type CanUseTool,
1114
query,
@@ -5028,6 +5031,39 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
50285031
},
50295032
);
50305033

5034+
/** Read the configured instance’s saved child transcript without creating a live SDK query. */
5035+
const getAgentHistory: ClaudeAdapterShape["getAgentHistory"] = Effect.fn("getAgentHistory")(
5036+
function* (input) {
5037+
const sessionId = readClaudeResumeState(input.resumeCursor)?.resume;
5038+
if (!sessionId)
5039+
return {
5040+
status: "unavailable",
5041+
entries: [],
5042+
nextOffset: null,
5043+
message: "No saved Claude session is available for this thread.",
5044+
};
5045+
return yield* Effect.tryPromise({
5046+
try: () =>
5047+
readClaudeAgentHistory({
5048+
sessionId,
5049+
agentId: input.agentId,
5050+
offset: input.offset,
5051+
view: input.view,
5052+
configDir: claudeEnvironment.CLAUDE_CONFIG_DIR
5053+
? expandHomePath(claudeEnvironment.CLAUDE_CONFIG_DIR)
5054+
: path.join(NodeOS.homedir(), ".claude"),
5055+
}),
5056+
catch: (cause) =>
5057+
new ProviderAdapterRequestError({
5058+
provider: PROVIDER,
5059+
method: "getSubagentMessages",
5060+
detail: "Could not read saved Claude agent history.",
5061+
cause,
5062+
}),
5063+
});
5064+
},
5065+
);
5066+
50315067
const readThread: ClaudeAdapterShape["readThread"] = Effect.fn("readThread")(
50325068
function* (threadId) {
50335069
const context = yield* requireSession(threadId);
@@ -5135,6 +5171,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
51355171
sendTurn,
51365172
interruptTurn,
51375173
readThread,
5174+
getAgentHistory,
51385175
rollbackThread,
51395176
respondToRequest,
51405177
respondToUserInput,
Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,74 @@
1+
// @effect-diagnostics nodeBuiltinImport:off
2+
import * as NodeFS from "node:fs";
3+
import * as NodeOS from "node:os";
4+
import * as NodePath from "node:path";
5+
import * as NodeServices from "@effect/platform-node/NodeServices";
6+
import { it, expect } from "@effect/vitest";
7+
import { CodexSettings, ThreadId } from "@t3tools/contracts";
8+
import * as Effect from "effect/Effect";
9+
import * as Layer from "effect/Layer";
10+
import * as Schema from "effect/Schema";
11+
import { ServerConfig } from "../../config.ts";
12+
import { writeFakeCli } from "../../testUtils/fakeCli.ts";
13+
import { makeCodexAdapter } from "./CodexAdapter.ts";
14+
15+
const decodeCodexSettings = Schema.decodeEffect(CodexSettings);
16+
17+
it.effect(
18+
"reuses one initialize-only Codex transport across concurrent and repeated history reads",
19+
() =>
20+
Effect.gen(function* () {
21+
const dir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "codex-history-client-"));
22+
yield* Effect.addFinalizer(() =>
23+
Effect.sync(() => NodeFS.rmSync(dir, { recursive: true, force: true })),
24+
);
25+
const log = NodePath.join(dir, "methods.txt");
26+
const binaryPath = writeFakeCli({
27+
directory: dir,
28+
name: "codex-history",
29+
env: { HISTORY_TEST_LOG: log },
30+
source: `
31+
import { appendFileSync } from "node:fs";
32+
import { createInterface } from "node:readline";
33+
createInterface({ input: process.stdin }).on("line", (line) => {
34+
const request = JSON.parse(line);
35+
appendFileSync(process.env.HISTORY_TEST_LOG, request.method + "\\n");
36+
if (request.id === undefined) return;
37+
const result = request.method === "initialize" ? { userAgent: "codex/1.0", codexHome: process.cwd(), platformFamily: "unix", platformOs: "linux" }
38+
: { thread: { id: request.params.threadId, cliVersion: "test", createdAt: 0, updatedAt: 0,
39+
cwd: process.cwd(), ephemeral: false, modelProvider: "openai", preview: "", sessionId: "session",
40+
status: { type: "idle" }, source: { subAgent: { thread_spawn: { parent_thread_id: "parent", depth: 1 } } },
41+
turns: [{ id: "turn", status: "completed", error: null, items: [{ id: "message", type: "agentMessage", text: "Saved answer" }] }] } };
42+
process.stdout.write(JSON.stringify({ id: request.id, result }) + "\\n");
43+
});
44+
`,
45+
});
46+
const adapter = yield* makeCodexAdapter(yield* decodeCodexSettings({ binaryPath }));
47+
const input = {
48+
threadId: ThreadId.make("stopped"),
49+
agentId: "child",
50+
cwd: dir,
51+
offset: 0,
52+
resumeCursor: { threadId: "parent" },
53+
};
54+
const results = yield* Effect.all(
55+
[adapter.getAgentHistory!(input), adapter.getAgentHistory!(input)],
56+
{ concurrency: "unbounded" },
57+
);
58+
expect(results.every((result) => result.entries[0]?.detail === "Saved answer")).toBe(true);
59+
yield* adapter.getAgentHistory!(input);
60+
const methods = NodeFS.readFileSync(log, "utf8").trim().split("\n");
61+
expect(methods.filter((method) => method === "initialize")).toHaveLength(1);
62+
expect(methods.filter((method) => method === "thread/read")).toHaveLength(6);
63+
expect(
64+
methods.every((method) => ["initialize", "initialized", "thread/read"].includes(method)),
65+
).toBe(true);
66+
expect(yield* adapter.listSessions()).toEqual([]);
67+
}).pipe(
68+
Effect.provide(
69+
ServerConfig.layerTest(process.cwd(), process.cwd()).pipe(
70+
Layer.provideMerge(NodeServices.layer),
71+
),
72+
),
73+
),
74+
);

‎apps/server/src/provider/Layers/CodexAdapter.ts‎

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,6 @@
1+
import { makeAgentHistoryClient } from "./agentHistoryClient.ts";
2+
import { withCodexAppServerClient } from "./CodexProvider.ts";
3+
import { readCodexAgentHistory } from "./codexAgentHistory.ts";
14
/**
25
* CodexAdapterLive - Scoped live implementation for the Codex provider adapter.
36
*
@@ -2234,6 +2237,16 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* (
22342237
const runtimeEventQueue = yield* Queue.unbounded<ProviderRuntimeEvent>();
22352238
const sessions = new Map<ThreadId, CodexAdapterSessionContext>();
22362239

2240+
const withHistoryClient = yield* makeAgentHistoryClient((cwd) =>
2241+
withCodexAppServerClient({
2242+
binaryPath: codexConfig.binaryPath,
2243+
homePath: codexConfig.homePath,
2244+
launchArgs: resolveCodexLaunchArgs(codexConfig.launchArgs, options?.environment),
2245+
environment: options?.environment,
2246+
cwd,
2247+
}).pipe(Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, childProcessSpawner)),
2248+
);
2249+
22372250
const startSession: CodexAdapterShape["startSession"] = (input) =>
22382251
Effect.scoped(
22392252
Effect.gen(function* () {
@@ -2573,6 +2586,41 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* (
25732586
);
25742587
});
25752588

2589+
/** Read saved descendants through the shared read-only transport without resuming a coding session. */
2590+
const getAgentHistory: CodexAdapterShape["getAgentHistory"] = Effect.fn("getAgentHistory")(
2591+
function* (input) {
2592+
if (!isCodexResumeCursorSchema(input.resumeCursor)) {
2593+
return {
2594+
status: "unavailable",
2595+
entries: [],
2596+
nextOffset: null,
2597+
message: "No saved Codex session is available for this thread.",
2598+
};
2599+
}
2600+
const parentThreadId = input.resumeCursor.threadId;
2601+
return yield* withHistoryClient(input.cwd ?? process.cwd(), ({ client }) =>
2602+
readCodexAgentHistory({
2603+
parentThreadId,
2604+
agentId: input.agentId,
2605+
offset: input.offset,
2606+
view: input.view,
2607+
readThread: (threadId, includeTurns) =>
2608+
client.request("thread/read", { threadId, includeTurns }),
2609+
}),
2610+
).pipe(
2611+
Effect.mapError(
2612+
(cause) =>
2613+
new ProviderAdapterRequestError({
2614+
provider: PROVIDER,
2615+
method: "thread/read",
2616+
detail: "Could not read saved agent history.",
2617+
cause,
2618+
}),
2619+
),
2620+
);
2621+
},
2622+
);
2623+
25762624
const readThread: CodexAdapterShape["readThread"] = (threadId) =>
25772625
requireSession(threadId).pipe(
25782626
Effect.flatMap((session) => session.runtime.readThread),
@@ -2721,6 +2769,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* (
27212769
compaction: { type: "native", start: compactThread },
27222770
interruptTurn,
27232771
readThread,
2772+
getAgentHistory,
27242773
rollbackThread,
27252774
uploadFeedback,
27262775
respondToRequest,

‎apps/server/src/provider/Layers/ProviderService.test.ts‎

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -257,6 +257,14 @@ function makeFakeCodexAdapter(
257257
Effect.succeed({ threadId, turns: [] }),
258258
);
259259

260+
const getAgentHistory = vi.fn(
261+
(
262+
_input: Parameters<
263+
NonNullable<ProviderAdapterShape<ProviderAdapterError>["getAgentHistory"]>
264+
>[0],
265+
) => Effect.succeed({ status: "ready" as const, entries: [], nextOffset: null, message: null }),
266+
);
267+
260268
const uploadFeedback = vi.fn(
261269
(
262270
input: ProviderUploadFeedbackInput,
@@ -295,6 +303,7 @@ function makeFakeCodexAdapter(
295303
readThread,
296304
rollbackThread,
297305
...(provider === CODEX_DRIVER ? { uploadFeedback } : {}),
306+
...(provider === CODEX_DRIVER || provider === CLAUDE_AGENT_DRIVER ? { getAgentHistory } : {}),
298307
stopAll,
299308
get streamEvents() {
300309
return Stream.fromPubSub(runtimeEventPubSub);
@@ -332,6 +341,7 @@ function makeFakeCodexAdapter(
332341
readThread,
333342
rollbackThread,
334343
uploadFeedback,
344+
getAgentHistory,
335345
stopAll,
336346
};
337347
}
@@ -1974,6 +1984,63 @@ routing.layer("ProviderServiceLive routing", (it) => {
19741984
}),
19751985
);
19761986

1987+
it.effect("reads saved agent history without recovering a stopped session", () =>
1988+
Effect.gen(function* () {
1989+
const provider = yield* ProviderService.ProviderService;
1990+
const threadId = asThreadId("thread-agent-history-stopped");
1991+
const cwd = fixtureCwd("agent-history");
1992+
yield* provider.startSession(threadId, {
1993+
provider: CODEX_DRIVER,
1994+
providerInstanceId: codexInstanceId,
1995+
threadId,
1996+
cwd,
1997+
resumeCursor: { threadId: "native-parent" },
1998+
runtimeMode: "full-access",
1999+
});
2000+
yield* routing.codex.stopSession(threadId);
2001+
routing.codex.startSession.mockClear();
2002+
routing.codex.sendTurn.mockClear();
2003+
routing.codex.getAgentHistory.mockClear();
2004+
const result = yield* provider.getAgentHistory({
2005+
threadId,
2006+
agentId: "native-child",
2007+
offset: 50,
2008+
});
2009+
assert.equal(result.status, "ready");
2010+
assert.equal(routing.codex.startSession.mock.calls.length, 0);
2011+
assert.equal(routing.codex.sendTurn.mock.calls.length, 0);
2012+
assert.deepStrictEqual(routing.codex.getAgentHistory.mock.calls, [
2013+
[
2014+
{
2015+
threadId,
2016+
agentId: "native-child",
2017+
offset: 50,
2018+
cwd,
2019+
resumeCursor: { threadId: "native-parent" },
2020+
},
2021+
],
2022+
]);
2023+
}),
2024+
);
2025+
2026+
it.effect("reports unsupported history without restarting the provider", () =>
2027+
Effect.gen(function* () {
2028+
const provider = yield* ProviderService.ProviderService;
2029+
const threadId = asThreadId("thread-agent-history-unsupported");
2030+
yield* provider.startSession(threadId, {
2031+
provider: CURSOR_DRIVER,
2032+
providerInstanceId: ProviderInstanceId.make("cursor"),
2033+
threadId,
2034+
runtimeMode: "full-access",
2035+
});
2036+
yield* routing.cursor.stopSession(threadId);
2037+
routing.cursor.startSession.mockClear();
2038+
const result = yield* provider.getAgentHistory({ threadId, agentId: "child", offset: 0 });
2039+
assert.equal(result.status, "unsupported");
2040+
assert.equal(routing.cursor.startSession.mock.calls.length, 0);
2041+
}),
2042+
);
2043+
19772044
it.effect("routes feedback to the Codex adapter and returns its feedback ID", () =>
19782045
Effect.gen(function* () {
19792046
const provider = yield* ProviderService.ProviderService;

0 commit comments

Comments
 (0)