diff --git a/scripts/measure-subagent-cold-start.ts b/scripts/measure-subagent-cold-start.ts new file mode 100644 index 000000000..4a13e9a8c --- /dev/null +++ b/scripts/measure-subagent-cold-start.ts @@ -0,0 +1,216 @@ +#!/usr/bin/env bun + +import { mkdir, mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join, resolve } from "node:path"; +import { performance } from "node:perf_hooks"; + +import { + createSkillSearchTool, + workerSkillSearchDefinition, +} from "../src/agent/skill-search.js"; +import { + createUseSkillTool, + workerUseSkillDefinition, +} from "../src/agent/use-skill.js"; +import { discoverSkills } from "../src/extensions/skills.js"; +import { assembleInferenceBase } from "../src/session/assemble-runtime.js"; +import { createSessionStores } from "../src/session/optimized-context-store.js"; + +type Mode = "sequential" | "overlap"; +type Measurement = { + a: number; + b: number; + c: number; + readiness: number; +}; +type Sample = { + sample: number; + sequential: Measurement; + overlap: Measurement; + savings: number; +}; + +const DEFAULT_SAMPLES = 7; +const repositoryRoot = resolve(import.meta.dir, ".."); +const skillDirs = [join(repositoryRoot, "plugins/corbits-skills")]; +const allowedSkillNames = [ + "style", + "philosophy", + "native-runtime", + "idiot-proof", + "ponytail", +]; +const attachedSkills = ["style", "philosophy"]; + +function parseSampleCount(args: readonly string[]): number { + const raw = args + .find((arg) => arg.startsWith("--samples=")) + ?.slice("--samples=".length); + if (raw === undefined) return DEFAULT_SAMPLES; + + const samples = Number(raw); + if (!Number.isSafeInteger(samples) || samples < 1) { + throw new Error("--samples must be a positive integer"); + } + return samples; +} + +function elapsedSince(startedAt: number): number { + return performance.now() - startedAt; +} + +async function measureLeaf(mode: Mode, workdir: string): Promise { + const readinessStartedAt = performance.now(); + + let phaseStartedAt = performance.now(); + await assembleInferenceBase(); + const a = elapsedSince(phaseStartedAt); + + phaseStartedAt = performance.now(); + const skills = await discoverSkills(repositoryRoot, skillDirs); + createSkillSearchTool({ + skills, + allowedNames: allowedSkillNames, + definition: workerSkillSearchDefinition, + }); + createUseSkillTool( + repositoryRoot, + skillDirs, + undefined, + allowedSkillNames, + workerUseSkillDefinition, + attachedSkills, + ); + const b = elapsedSince(phaseStartedAt); + + phaseStartedAt = performance.now(); + if (mode === "sequential") { + await mkdir(workdir, { recursive: true }); + await createSessionStores(workdir); + } else { + const storesPromise = createSessionStores(workdir); + void storesPromise.catch(() => undefined); + await mkdir(workdir, { recursive: true }); + await storesPromise; + } + const c = elapsedSince(phaseStartedAt); + + return { a, b, c, readiness: elapsedSince(readinessStartedAt) }; +} + +function summary(values: readonly number[]) { + const sorted = [...values].sort((left, right) => left - right); + const middle = Math.floor(sorted.length / 2); + const median = + sorted.length % 2 === 0 + ? ((sorted[middle - 1] ?? 0) + (sorted[middle] ?? 0)) / 2 + : (sorted[middle] ?? 0); + return { + median, + mean: values.reduce((sum, value) => sum + value, 0) / values.length, + }; +} + +function round(value: number): number { + return Number(value.toFixed(3)); +} + +function roundedSummary(values: readonly number[]) { + const result = summary(values); + return { median: round(result.median), mean: round(result.mean) }; +} + +async function main(): Promise { + const sampleCount = parseSampleCount(process.argv.slice(2)); + const root = await mkdtemp(join(tmpdir(), "corbits-subagent-cold-start-")); + const samples: Sample[] = []; + + try { + await measureLeaf("sequential", join(root, "warmup-sequential")); + await measureLeaf("overlap", join(root, "warmup-overlap")); + + for (let sample = 1; sample <= sampleCount; sample++) { + const modes: readonly Mode[] = + sample % 2 === 1 + ? ["sequential", "overlap"] + : ["overlap", "sequential"]; + const measurements = new Map(); + + for (const mode of modes) { + measurements.set( + mode, + await measureLeaf(mode, join(root, `${sample}-${mode}`)), + ); + } + + const sequential = measurements.get("sequential"); + const overlap = measurements.get("overlap"); + if (sequential === undefined || overlap === undefined) { + throw new Error("both benchmark modes must complete"); + } + samples.push({ + sample, + sequential, + overlap, + savings: sequential.readiness - overlap.readiness, + }); + } + + const roundedSamples = samples.map((sample) => ({ + sample: sample.sample, + sequential: Object.fromEntries( + Object.entries(sample.sequential).map(([key, value]) => [ + key, + round(value), + ]), + ), + overlap: Object.fromEntries( + Object.entries(sample.overlap).map(([key, value]) => [ + key, + round(value), + ]), + ), + savings: round(sample.savings), + })); + const sequentialReadiness = samples.map( + (sample) => sample.sequential.readiness, + ); + const overlapReadiness = samples.map((sample) => sample.overlap.readiness); + const savings = samples.map((sample) => sample.savings); + const overlapResidual = samples.map((sample) => sample.overlap.c); + + console.log( + JSON.stringify( + { + samplesPerMode: sampleCount, + warmupRunsPerMode: 1, + units: "milliseconds", + boundaries: { + a: "assemble inference dependencies", + b: "discover skills and construct skill tools", + c: "make the workdir and initialize session stores", + readiness: "a + b + c, immediately before agent construction", + sequential: "await mkdir, then initialize stores", + overlap: "start stores, await mkdir, then await stores", + residual: + "overlap c latency still visible on the leaf-readiness critical path", + }, + samples: roundedSamples, + summary: { + beforeSequentialReadiness: roundedSummary(sequentialReadiness), + afterOverlapReadiness: roundedSummary(overlapReadiness), + pairedSavings: roundedSummary(savings), + overlapResidual: roundedSummary(overlapResidual), + }, + }, + null, + 2, + ), + ); + } finally { + await rm(root, { recursive: true, force: true }); + } +} + +await main(); diff --git a/src/subagent/run-audit-store.test.ts b/src/subagent/run-audit-store.test.ts index 6215e4046..6b1348863 100644 --- a/src/subagent/run-audit-store.test.ts +++ b/src/subagent/run-audit-store.test.ts @@ -1,7 +1,7 @@ import { expect, test } from "bun:test"; -import { mkdtemp } from "node:fs/promises"; +import { mkdtemp, rm } from "node:fs/promises"; import { tmpdir } from "node:os"; -import { join } from "node:path"; +import { basename, join } from "node:path"; import type { AuditStore, ContextStore } from "@intx/types/runtime"; import { withMockedModuleDuring } from "../../tests/helpers/mock-module.js"; @@ -15,81 +15,327 @@ const permissionGate = createPermissionGate({ reactorGated: false, }); -test("runSubAgent threads the isogit audit store and session id into createAgent", async () => { - const cwd = await mkdtemp(join(tmpdir(), "corbits-run-audit-")); - const fakeStore = { +function fakeStore(): ContextStore & AuditStore { + return { readBlob: async () => new Uint8Array(), } as unknown as ContextStore & AuditStore; +} + +function stubAgent() { + return { + send: async () => ({ + type: "reply" as const, + reply: "ok", + turn: { role: "assistant" as const, content: [] }, + }), + stream: () => + (async function* () { + yield* []; + })(), + deliver: () => undefined, + close: async () => undefined, + setSource: () => undefined, + setSources: () => undefined, + history: async () => [], + checkpoints: async () => [], + readAt: async () => [], + blobReader: {}, + }; +} + +function runParams(cwd: string, id: string) { + return { + cwd, + workdirBase: join(cwd, ".ctx"), + permissionGate, + provider: { + providerName: "test", + baseURL: "http://localhost", + model: "test-model", + }, + description: "audit wiring", + prompt: "noop", + id, + }; +} + +test("runSubAgent threads the isogit audit store and session id into createAgent", async () => { + const cwd = await mkdtemp(join(tmpdir(), "corbits-run-audit-")); + const store = fakeStore(); let seen: | { audit: AuditStore; sessionId?: string; storage: ContextStore } | undefined; - await withMockedModuleDuring( - import.meta.resolve("../session/optimized-context-store.js"), - (real: typeof import("../session/optimized-context-store.js")) => ({ - ...real, - createSessionStores: async () => ({ - storage: fakeStore, - audit: fakeStore, + try { + await withMockedModuleDuring( + import.meta.resolve("../session/optimized-context-store.js"), + (real: typeof import("../session/optimized-context-store.js")) => ({ + ...real, + createSessionStores: async () => ({ storage: store, audit: store }), }), - }), - async () => { - await withMockedModuleDuring( - import.meta.resolve("../agent/live-tool-dispatch.js"), - (real: typeof import("../agent/live-tool-dispatch.js")) => ({ - ...real, - createAgentWithLiveToolDispatch: async ( - _def: unknown, - env: { - storage: ContextStore; - audit: AuditStore; - sessionId?: string; + async () => { + await withMockedModuleDuring( + import.meta.resolve("../agent/live-tool-dispatch.js"), + (real: typeof import("../agent/live-tool-dispatch.js")) => ({ + ...real, + createAgentWithLiveToolDispatch: async ( + _def: unknown, + env: { + storage: ContextStore; + audit: AuditStore; + sessionId?: string; + }, + ) => { + seen = env; + return stubAgent() as unknown as Awaited< + ReturnType + >; }, - ) => { - seen = env; - return { - send: async () => ({ - type: "reply" as const, - reply: "ok", - turn: { role: "assistant" as const, content: [] }, - }), - stream: () => - (async function* () { - yield* []; - })(), - deliver: () => undefined, - close: async () => undefined, - setSource: () => undefined, - setSources: () => undefined, - history: async () => [], - checkpoints: async () => [], - readAt: async () => [], - blobReader: {}, - }; + }), + async () => { + const { runSubAgent } = await import("./run.js"); + await runSubAgent(runParams(cwd, "child-session-1")); }, - }), - async () => { - const { runSubAgent } = await import("./run.js"); - await runSubAgent({ - cwd, - workdirBase: join(cwd, ".ctx"), - permissionGate, - provider: { - providerName: "test", - baseURL: "http://localhost", - model: "test-model", - }, - description: "audit wiring", - prompt: "noop", - id: "child-session-1", + ); + }, + ); + + const seenStores = defined(seen); + expect(seenStores.storage).toBe(store); + expect(seenStores.audit).toBe(store); + expect(seenStores.sessionId).toBe("child-session-1"); + } finally { + await rm(cwd, { recursive: true, force: true }); + } +}); + +test("runSubAgent overlaps store creation with workdir setup", async () => { + const cwd = await mkdtemp(join(tmpdir(), "corbits-run-store-overlap-")); + const store = fakeStore(); + let releaseMkdir: (() => void) | undefined; + let signalMkdirStarted: (() => void) | undefined; + let signalMkdirFinished: (() => void) | undefined; + const pendingMkdir = new Promise((resolve) => { + releaseMkdir = resolve; + }); + const mkdirStarted = new Promise((resolve) => { + signalMkdirStarted = resolve; + }); + const mkdirFinished = new Promise((resolve) => { + signalMkdirFinished = resolve; + }); + let resolveStores: + | ((stores: { storage: ContextStore; audit: AuditStore }) => void) + | undefined; + const pendingStores = new Promise<{ + storage: ContextStore; + audit: AuditStore; + }>((resolve) => { + resolveStores = resolve; + }); + let storeStarted = false; + let agentConstructed = false; + + try { + await withMockedModuleDuring( + import.meta.resolve("node:fs/promises"), + (real: typeof import("node:fs/promises")) => ({ + ...real, + mkdir: () => { + defined(signalMkdirStarted, "mkdir start signal")(); + return pendingMkdir.then(() => { + defined(signalMkdirFinished, "mkdir finish signal")(); + return undefined; }); }, - ); - }, - ); + }), + async () => { + await withMockedModuleDuring( + import.meta.resolve("../session/optimized-context-store.js"), + (real: typeof import("../session/optimized-context-store.js")) => ({ + ...real, + createSessionStores: () => { + storeStarted = true; + return pendingStores; + }, + }), + async () => { + await withMockedModuleDuring( + import.meta.resolve("../agent/live-tool-dispatch.js"), + (real: typeof import("../agent/live-tool-dispatch.js")) => ({ + ...real, + createAgentWithLiveToolDispatch: async () => { + agentConstructed = true; + return stubAgent() as unknown as Awaited< + ReturnType + >; + }, + }), + async () => { + const { runSubAgent } = await import("./run.js"); + const run = runSubAgent(runParams(cwd, "overlap-child")); + await mkdirStarted; + + expect(storeStarted).toBe(true); + expect(agentConstructed).toBe(false); + + defined(releaseMkdir, "mkdir resolver")(); + await mkdirFinished; + expect(agentConstructed).toBe(false); + + defined( + resolveStores, + "store resolver", + )({ + storage: store, + audit: store, + }); + await run; + expect(agentConstructed).toBe(true); + }, + ); + }, + ); + }, + ); + } finally { + await rm(cwd, { recursive: true, force: true }); + } +}); + +test("runSubAgent waits for workdir setup after stores resolve", async () => { + const cwd = await mkdtemp(join(tmpdir(), "corbits-run-mkdir-barrier-")); + const store = fakeStore(); + let releaseMkdir: (() => void) | undefined; + let signalMkdirStarted: (() => void) | undefined; + const pendingMkdir = new Promise((resolve) => { + releaseMkdir = resolve; + }); + const mkdirStarted = new Promise((resolve) => { + signalMkdirStarted = resolve; + }); + let resolveStores: + | ((stores: { storage: ContextStore; audit: AuditStore }) => void) + | undefined; + const pendingStores = new Promise<{ + storage: ContextStore; + audit: AuditStore; + }>((resolve) => { + resolveStores = resolve; + }); + let agentStores: { storage: ContextStore; audit: AuditStore } | undefined; + + try { + await withMockedModuleDuring( + import.meta.resolve("node:fs/promises"), + (real: typeof import("node:fs/promises")) => ({ + ...real, + mkdir: () => { + defined(signalMkdirStarted, "mkdir start signal")(); + return pendingMkdir; + }, + }), + async () => { + await withMockedModuleDuring( + import.meta.resolve("../session/optimized-context-store.js"), + (real: typeof import("../session/optimized-context-store.js")) => ({ + ...real, + createSessionStores: () => pendingStores, + }), + async () => { + await withMockedModuleDuring( + import.meta.resolve("../agent/live-tool-dispatch.js"), + (real: typeof import("../agent/live-tool-dispatch.js")) => ({ + ...real, + createAgentWithLiveToolDispatch: async ( + _def: unknown, + env: { storage: ContextStore; audit: AuditStore }, + ) => { + agentStores = env; + return stubAgent() as unknown as Awaited< + ReturnType + >; + }, + }), + async () => { + const { runSubAgent } = await import("./run.js"); + const run = runSubAgent(runParams(cwd, "mkdir-barrier-child")); + await mkdirStarted; + + defined( + resolveStores, + "store resolver", + )({ + storage: store, + audit: store, + }); + await pendingStores; + expect(agentStores).toBeUndefined(); + + defined(releaseMkdir, "mkdir resolver")(); + await run; + expect(defined(agentStores).storage).toBe(store); + expect(defined(agentStores).audit).toBe(store); + }, + ); + }, + ); + }, + ); + } finally { + releaseMkdir?.(); + await rm(cwd, { recursive: true, force: true }); + } +}); + +test("runSubAgent keeps session stores isolated between workers", async () => { + const cwd = await mkdtemp(join(tmpdir(), "corbits-run-store-isolation-")); + const storesBySession = new Map(); + const seenBySession = new Map(); + + try { + await withMockedModuleDuring( + import.meta.resolve("../session/optimized-context-store.js"), + (real: typeof import("../session/optimized-context-store.js")) => ({ + ...real, + createSessionStores: async (dir: string) => { + const store = fakeStore(); + storesBySession.set(basename(dir), store); + return { storage: store, audit: store }; + }, + }), + async () => { + await withMockedModuleDuring( + import.meta.resolve("../agent/live-tool-dispatch.js"), + (real: typeof import("../agent/live-tool-dispatch.js")) => ({ + ...real, + createAgentWithLiveToolDispatch: async ( + _def: unknown, + env: { storage: ContextStore; workdir: string }, + ) => { + seenBySession.set(basename(env.workdir), env.storage); + return stubAgent() as unknown as Awaited< + ReturnType + >; + }, + }), + async () => { + const { runSubAgent } = await import("./run.js"); + await Promise.all([ + runSubAgent(runParams(cwd, "isolated-a")), + runSubAgent(runParams(cwd, "isolated-b")), + ]); + }, + ); + }, + ); - const seenStores = defined(seen); - expect(seenStores.storage).toBe(fakeStore); - expect(seenStores.audit).toBe(fakeStore); - expect(seenStores.sessionId).toBe("child-session-1"); + const storeA = defined(storesBySession.get("isolated-a"), "worker A store"); + const storeB = defined(storesBySession.get("isolated-b"), "worker B store"); + expect(storeA).not.toBe(storeB); + expect(seenBySession.get("isolated-a")).toBe(storeA); + expect(seenBySession.get("isolated-b")).toBe(storeB); + } finally { + await rm(cwd, { recursive: true, force: true }); + } }); diff --git a/src/subagent/run.ts b/src/subagent/run.ts index e35b395b5..e8094e3f0 100644 --- a/src/subagent/run.ts +++ b/src/subagent/run.ts @@ -1145,6 +1145,8 @@ async function runSubAgentInner( : undefined; const sessionId = safeRequestedId ?? generateSessionId(); const workdir = join(params.workdirBase, "subagents", sessionId); + const sessionStoresPromise = createSessionStores(workdir); + void sessionStoresPromise.catch(() => undefined); await mkdir(workdir, { recursive: true }); childContextDir = workdir; // One record per stop/nudge, with its measured value beside its @@ -1173,9 +1175,6 @@ async function runSubAgentInner( }, }); - const { storage, audit } = await createSessionStores(workdir); - childBlobWriter = (key, bytes, contentType) => - storage.writeBlob(key, bytes, contentType); const authorize = createWorkerAuthorize(params.permissionGate); const head = { @@ -1202,6 +1201,9 @@ async function runSubAgentInner( bundle.sources[0]; if (workerSource === undefined) throw new Error("sub-agent source bundle is empty"); + const { storage, audit } = await sessionStoresPromise; + childBlobWriter = (key, bytes, contentType) => + storage.writeBlob(key, bytes, contentType); agent = await createAgentWithLiveToolDispatch(def, { sources: bundle.sources, defaultSource: bundle.defaultSource,