Skip to content

Commit 91d9cc0

Browse files
committed
perf(subagent): overlap session store initialization
1 parent b3a95ad commit 91d9cc0

2 files changed

Lines changed: 219 additions & 70 deletions

File tree

Lines changed: 214 additions & 67 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
import { expect, test } from "bun:test";
2-
import { mkdtemp } from "node:fs/promises";
2+
import { mkdtemp, rm, stat } from "node:fs/promises";
33
import { tmpdir } from "node:os";
4-
import { join } from "node:path";
4+
import { basename, join } from "node:path";
55
import type { AuditStore, ContextStore } from "@intx/types/runtime";
66

77
import { withMockedModuleDuring } from "../../tests/helpers/mock-module.js";
@@ -15,81 +15,228 @@ const permissionGate = createPermissionGate({
1515
reactorGated: false,
1616
});
1717

18-
test("runSubAgent threads the isogit audit store and session id into createAgent", async () => {
19-
const cwd = await mkdtemp(join(tmpdir(), "corbits-run-audit-"));
20-
const fakeStore = {
18+
function fakeStore(): ContextStore & AuditStore {
19+
return {
2120
readBlob: async () => new Uint8Array(),
2221
} as unknown as ContextStore & AuditStore;
22+
}
23+
24+
function stubAgent() {
25+
return {
26+
send: async () => ({
27+
type: "reply" as const,
28+
reply: "ok",
29+
turn: { role: "assistant" as const, content: [] },
30+
}),
31+
stream: () =>
32+
(async function* () {
33+
yield* [];
34+
})(),
35+
deliver: () => undefined,
36+
close: async () => undefined,
37+
setSource: () => undefined,
38+
setSources: () => undefined,
39+
history: async () => [],
40+
checkpoints: async () => [],
41+
readAt: async () => [],
42+
blobReader: {},
43+
};
44+
}
45+
46+
function runParams(cwd: string, id: string) {
47+
return {
48+
cwd,
49+
workdirBase: join(cwd, ".ctx"),
50+
permissionGate,
51+
provider: {
52+
providerName: "test",
53+
baseURL: "http://localhost",
54+
model: "test-model",
55+
},
56+
description: "audit wiring",
57+
prompt: "noop",
58+
id,
59+
};
60+
}
61+
62+
async function eventuallyExists(path: string): Promise<boolean> {
63+
for (let attempt = 0; attempt < 50; attempt++) {
64+
if (
65+
await stat(path)
66+
.then(() => true)
67+
.catch(() => false)
68+
)
69+
return true;
70+
await Bun.sleep(1);
71+
}
72+
return false;
73+
}
74+
75+
test("runSubAgent threads the isogit audit store and session id into createAgent", async () => {
76+
const cwd = await mkdtemp(join(tmpdir(), "corbits-run-audit-"));
77+
const store = fakeStore();
2378
let seen:
2479
| { audit: AuditStore; sessionId?: string; storage: ContextStore }
2580
| undefined;
2681

27-
await withMockedModuleDuring(
28-
import.meta.resolve("../session/optimized-context-store.js"),
29-
(real: typeof import("../session/optimized-context-store.js")) => ({
30-
...real,
31-
createSessionStores: async () => ({
32-
storage: fakeStore,
33-
audit: fakeStore,
82+
try {
83+
await withMockedModuleDuring(
84+
import.meta.resolve("../session/optimized-context-store.js"),
85+
(real: typeof import("../session/optimized-context-store.js")) => ({
86+
...real,
87+
createSessionStores: async () => ({ storage: store, audit: store }),
3488
}),
35-
}),
36-
async () => {
37-
await withMockedModuleDuring(
38-
import.meta.resolve("../agent/live-tool-dispatch.js"),
39-
(real: typeof import("../agent/live-tool-dispatch.js")) => ({
40-
...real,
41-
createAgentWithLiveToolDispatch: async (
42-
_def: unknown,
43-
env: {
44-
storage: ContextStore;
45-
audit: AuditStore;
46-
sessionId?: string;
89+
async () => {
90+
await withMockedModuleDuring(
91+
import.meta.resolve("../agent/live-tool-dispatch.js"),
92+
(real: typeof import("../agent/live-tool-dispatch.js")) => ({
93+
...real,
94+
createAgentWithLiveToolDispatch: async (
95+
_def: unknown,
96+
env: {
97+
storage: ContextStore;
98+
audit: AuditStore;
99+
sessionId?: string;
100+
},
101+
) => {
102+
seen = env;
103+
return stubAgent() as unknown as Awaited<
104+
ReturnType<typeof real.createAgentWithLiveToolDispatch>
105+
>;
47106
},
48-
) => {
49-
seen = env;
50-
return {
51-
send: async () => ({
52-
type: "reply" as const,
53-
reply: "ok",
54-
turn: { role: "assistant" as const, content: [] },
55-
}),
56-
stream: () =>
57-
(async function* () {
58-
yield* [];
59-
})(),
60-
deliver: () => undefined,
61-
close: async () => undefined,
62-
setSource: () => undefined,
63-
setSources: () => undefined,
64-
history: async () => [],
65-
checkpoints: async () => [],
66-
readAt: async () => [],
67-
blobReader: {},
68-
};
107+
}),
108+
async () => {
109+
const { runSubAgent } = await import("./run.js");
110+
await runSubAgent(runParams(cwd, "child-session-1"));
69111
},
70-
}),
71-
async () => {
72-
const { runSubAgent } = await import("./run.js");
73-
await runSubAgent({
74-
cwd,
75-
workdirBase: join(cwd, ".ctx"),
76-
permissionGate,
77-
provider: {
78-
providerName: "test",
79-
baseURL: "http://localhost",
80-
model: "test-model",
112+
);
113+
},
114+
);
115+
116+
const seenStores = defined(seen);
117+
expect(seenStores.storage).toBe(store);
118+
expect(seenStores.audit).toBe(store);
119+
expect(seenStores.sessionId).toBe("child-session-1");
120+
} finally {
121+
await rm(cwd, { recursive: true, force: true });
122+
}
123+
});
124+
125+
test("runSubAgent overlaps store creation with workdir setup", async () => {
126+
const cwd = await mkdtemp(join(tmpdir(), "corbits-run-store-overlap-"));
127+
const store = fakeStore();
128+
let resolveStores:
129+
| ((stores: { storage: ContextStore; audit: AuditStore }) => void)
130+
| undefined;
131+
let signalStoreStarted: (() => void) | undefined;
132+
const storeStarted = new Promise<void>((resolve) => {
133+
signalStoreStarted = resolve;
134+
});
135+
const pendingStores = new Promise<{
136+
storage: ContextStore;
137+
audit: AuditStore;
138+
}>((resolve) => {
139+
resolveStores = resolve;
140+
});
141+
let agentConstructed = false;
142+
143+
try {
144+
await withMockedModuleDuring(
145+
import.meta.resolve("../session/optimized-context-store.js"),
146+
(real: typeof import("../session/optimized-context-store.js")) => ({
147+
...real,
148+
createSessionStores: () => {
149+
defined(signalStoreStarted, "store start signal")();
150+
return pendingStores;
151+
},
152+
}),
153+
async () => {
154+
await withMockedModuleDuring(
155+
import.meta.resolve("../agent/live-tool-dispatch.js"),
156+
(real: typeof import("../agent/live-tool-dispatch.js")) => ({
157+
...real,
158+
createAgentWithLiveToolDispatch: async () => {
159+
agentConstructed = true;
160+
return stubAgent() as unknown as Awaited<
161+
ReturnType<typeof real.createAgentWithLiveToolDispatch>
162+
>;
81163
},
82-
description: "audit wiring",
83-
prompt: "noop",
84-
id: "child-session-1",
85-
});
164+
}),
165+
async () => {
166+
const { runSubAgent } = await import("./run.js");
167+
const run = runSubAgent(runParams(cwd, "overlap-child"));
168+
await storeStarted;
169+
170+
const workdir = join(cwd, ".ctx", "subagents", "overlap-child");
171+
expect(await eventuallyExists(workdir)).toBe(true);
172+
expect(agentConstructed).toBe(false);
173+
174+
defined(
175+
resolveStores,
176+
"store resolver",
177+
)({
178+
storage: store,
179+
audit: store,
180+
});
181+
await run;
182+
expect(agentConstructed).toBe(true);
183+
},
184+
);
185+
},
186+
);
187+
} finally {
188+
await rm(cwd, { recursive: true, force: true });
189+
}
190+
});
191+
192+
test("runSubAgent keeps session stores isolated between workers", async () => {
193+
const cwd = await mkdtemp(join(tmpdir(), "corbits-run-store-isolation-"));
194+
const storesBySession = new Map<string, ContextStore & AuditStore>();
195+
const seenBySession = new Map<string, ContextStore>();
196+
197+
try {
198+
await withMockedModuleDuring(
199+
import.meta.resolve("../session/optimized-context-store.js"),
200+
(real: typeof import("../session/optimized-context-store.js")) => ({
201+
...real,
202+
createSessionStores: async (dir: string) => {
203+
const store = fakeStore();
204+
storesBySession.set(basename(dir), store);
205+
return { storage: store, audit: store };
86206
},
87-
);
88-
},
89-
);
207+
}),
208+
async () => {
209+
await withMockedModuleDuring(
210+
import.meta.resolve("../agent/live-tool-dispatch.js"),
211+
(real: typeof import("../agent/live-tool-dispatch.js")) => ({
212+
...real,
213+
createAgentWithLiveToolDispatch: async (
214+
_def: unknown,
215+
env: { storage: ContextStore; workdir: string },
216+
) => {
217+
seenBySession.set(basename(env.workdir), env.storage);
218+
return stubAgent() as unknown as Awaited<
219+
ReturnType<typeof real.createAgentWithLiveToolDispatch>
220+
>;
221+
},
222+
}),
223+
async () => {
224+
const { runSubAgent } = await import("./run.js");
225+
await Promise.all([
226+
runSubAgent(runParams(cwd, "isolated-a")),
227+
runSubAgent(runParams(cwd, "isolated-b")),
228+
]);
229+
},
230+
);
231+
},
232+
);
90233

91-
const seenStores = defined(seen);
92-
expect(seenStores.storage).toBe(fakeStore);
93-
expect(seenStores.audit).toBe(fakeStore);
94-
expect(seenStores.sessionId).toBe("child-session-1");
234+
const storeA = defined(storesBySession.get("isolated-a"), "worker A store");
235+
const storeB = defined(storesBySession.get("isolated-b"), "worker B store");
236+
expect(storeA).not.toBe(storeB);
237+
expect(seenBySession.get("isolated-a")).toBe(storeA);
238+
expect(seenBySession.get("isolated-b")).toBe(storeB);
239+
} finally {
240+
await rm(cwd, { recursive: true, force: true });
241+
}
95242
});

‎src/subagent/run.ts‎

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1145,6 +1145,8 @@ async function runSubAgentInner(
11451145
: undefined;
11461146
const sessionId = safeRequestedId ?? generateSessionId();
11471147
const workdir = join(params.workdirBase, "subagents", sessionId);
1148+
const sessionStoresPromise = createSessionStores(workdir);
1149+
void sessionStoresPromise.catch(() => undefined);
11481150
await mkdir(workdir, { recursive: true });
11491151
childContextDir = workdir;
11501152
// One record per stop/nudge, with its measured value beside its
@@ -1173,9 +1175,6 @@ async function runSubAgentInner(
11731175
},
11741176
});
11751177

1176-
const { storage, audit } = await createSessionStores(workdir);
1177-
childBlobWriter = (key, bytes, contentType) =>
1178-
storage.writeBlob(key, bytes, contentType);
11791178
const authorize = createWorkerAuthorize(params.permissionGate);
11801179

11811180
const head = {
@@ -1202,6 +1201,9 @@ async function runSubAgentInner(
12021201
bundle.sources[0];
12031202
if (workerSource === undefined)
12041203
throw new Error("sub-agent source bundle is empty");
1204+
const { storage, audit } = await sessionStoresPromise;
1205+
childBlobWriter = (key, bytes, contentType) =>
1206+
storage.writeBlob(key, bytes, contentType);
12051207
agent = await createAgentWithLiveToolDispatch(def, {
12061208
sources: bundle.sources,
12071209
defaultSource: bundle.defaultSource,

0 commit comments

Comments
 (0)