diff --git a/packages/agent/src/harness/session/fork.ts b/packages/agent/src/harness/session/fork.ts index 347b1b570..76d0bb897 100644 --- a/packages/agent/src/harness/session/fork.ts +++ b/packages/agent/src/harness/session/fork.ts @@ -40,7 +40,7 @@ export function createForkSnapshot(source: ForkSourceSnapshot, options: ForkOpti const sourceTips = storedValuesInNamespace(source.scalarValues, branchTip("")); validateForkSourceSnapshot(source, sourceEntries, sourceTips, options); - const { entryIds, sourceToDestinationTip } = selectForkContents(sourceEntries, sourceTips, options); + const { entryIds, destinationTips } = selectForkContents(sourceEntries, sourceTips, options); const entries = new Map(); for (const id of entryIds) entries.set(id, sourceEntries.get(id)!); @@ -53,12 +53,12 @@ export function createForkSnapshot(source: ForkSourceSnapshot, options: ForkOpti seq: nextSeq++, }); }; - for (const [sourceName, destination] of sourceToDestinationTip) { - const configuration = findStoredValue(source.scalarValues, laneConfig(sourceName)); - store(branchTip(destination.name), destination.tipId); + for (const [branch, tipId] of destinationTips) { + const configuration = findStoredValue(source.scalarValues, laneConfig(branch)); + store(branchTip(branch), tipId); if (configuration !== undefined) { - store(laneConfig(destination.name), configuration.value); - store(laneState(destination.name), { currentOperationId: null, lastOperationId: null, inbox: [] }); + store(laneConfig(branch), configuration.value); + store(laneState(branch), { currentOperationId: null, lastOperationId: null, inbox: [] }); } } const name = findStoredValue(source.scalarValues, sessionName); @@ -93,22 +93,17 @@ function selectForkContents( options: ForkOptions, ): { entryIds: Set; - sourceToDestinationTip: Map; + destinationTips: Map; } { const entryIds = new Set(); - const sourceToDestinationTip = new Map(); + const destinationTips = new Map(); if (options.scope === "tree") { for (const id of sourceEntries.keys()) entryIds.add(id); - for (const stored of sourceTips) { - sourceToDestinationTip.set(stored.address.key, { - name: stored.address.key, - tipId: stored.value, - }); - } + for (const stored of sourceTips) destinationTips.set(stored.address.key, stored.value); } else { - const mainTip = sourceTips.find((stored) => stored.address.key === "main"); - if (mainTip === undefined) throw new Error("Source session is missing main branch"); - const requested = options.entryId ?? mainTip.value; + const sourceTip = sourceTips.find((stored) => stored.address.key === options.branch); + if (sourceTip === undefined) throw new Error(`Unknown source branch: ${options.branch}`); + const requested = options.entryId ?? sourceTip.value; let tipId = requested; if (requested !== null) { const target = sourceEntries.get(requested); @@ -123,9 +118,9 @@ function selectForkContents( entryIds.add(entryId); entryId = entry.parentId; } - sourceToDestinationTip.set("main", { name: "main", tipId }); + destinationTips.set(options.branch, tipId); } - return { entryIds, sourceToDestinationTip }; + return { entryIds, destinationTips }; } function validateForkSourceSnapshot( @@ -136,9 +131,6 @@ function validateForkSourceSnapshot( ): void { const sourceTipKeys = new Set(sourceTips.map((stored) => stored.address.key)); - if (options.scope !== "tree" && !sourceTipKeys.has("main")) { - throw new Error("Source session is missing main branch"); - } for (const stored of source.scalarValues) { if ( (stored.address.namespace === laneConfig("").namespace || diff --git a/packages/agent/src/harness/session/testing/benchmark/session-repo.ts b/packages/agent/src/harness/session/testing/benchmark/session-repo.ts index 562082033..7a848f23e 100644 --- a/packages/agent/src/harness/session/testing/benchmark/session-repo.ts +++ b/packages/agent/src/harness/session/testing/benchmark/session-repo.ts @@ -156,7 +156,11 @@ export const SESSION_REPO_FORK_WRITE_BENCHMARK_SCENARIOS: readonly SessionRepoFo return dataset.entryCount; }, async run(repo, source, dataset) { - const fork = await repo.fork(source, { id: FORK_DESTINATION_SESSION_ID }, BACKGROUND_CONTEXT); + const fork = await repo.fork( + source, + { id: FORK_DESTINATION_SESSION_ID, scope: "branch", branch: "main" }, + BACKGROUND_CONTEXT, + ); return fork.metadata.id === FORK_DESTINATION_SESSION_ID && fork.metadata.parentSessionId === source.id ? dataset.entryCount : 0; diff --git a/packages/agent/src/harness/session/testing/conformance/session-repo.ts b/packages/agent/src/harness/session/testing/conformance/session-repo.ts index 4e037f70a..8e396963f 100644 --- a/packages/agent/src/harness/session/testing/conformance/session-repo.ts +++ b/packages/agent/src/harness/session/testing/conformance/session-repo.ts @@ -336,7 +336,7 @@ export function createSessionRepoForkBehaviorConformance { const source = await repo.create({ id: "source" }, BACKGROUND_CONTEXT); await source.mutate( @@ -361,9 +361,10 @@ export function createSessionRepoForkBehaviorConformance id), [ROOT_ID, CHILD_ID], ); - strictEqual(await getBranchTip(fork), CHILD_ID); - deepStrictEqual((await fork.getValue(laneConfig("main"), BACKGROUND_CONTEXT))?.value, configuration); - deepStrictEqual((await fork.getValue(laneState("main"), BACKGROUND_CONTEXT))?.value, idleLaneState); + strictEqual(await fork.branch("main", BACKGROUND_CONTEXT), undefined); + strictEqual(await getBranchTip(fork, "review"), CHILD_ID); + deepStrictEqual((await fork.getValue(laneConfig("review"), BACKGROUND_CONTEXT))?.value, configuration); + deepStrictEqual((await fork.getValue(laneState("review"), BACKGROUND_CONTEXT))?.value, idleLaneState); + strictEqual(await fork.getValue(laneConfig("main"), BACKGROUND_CONTEXT), undefined); + strictEqual(await fork.getValue(laneState("main"), BACKGROUND_CONTEXT), undefined); strictEqual(await fork.getName(BACKGROUND_CONTEXT), "source name"); strictEqual(await fork.getValue(applicationValue, BACKGROUND_CONTEXT), undefined); deepStrictEqual(await fork.readList(applicationList, undefined, BACKGROUND_CONTEXT), []); @@ -503,7 +507,7 @@ export function createSessionRepoForkBehaviorConformance id), [ROOT_ID], ); - await rejects(repo.fork(source.metadata, { id: "failed", entryId: SIBLING_ID }, BACKGROUND_CONTEXT)); + await rejects( + repo.fork( + source.metadata, + { id: "failed", scope: "branch", branch: "main", entryId: SIBLING_ID }, + BACKGROUND_CONTEXT, + ), + ); deepStrictEqual((await repo.list(undefined, BACKGROUND_CONTEXT)).map(({ id }) => id).sort(), [ "before", "source", @@ -541,7 +551,11 @@ export function createSessionRepoForkBehaviorConformance mutator.commit( diff --git a/packages/agent/src/harness/session/types.ts b/packages/agent/src/harness/session/types.ts index 4c9cefeaa..2c8fee6ec 100644 --- a/packages/agent/src/harness/session/types.ts +++ b/packages/agent/src/harness/session/types.ts @@ -560,13 +560,10 @@ export interface SessionCreateOptions { export type ForkOptions = | { - /** - * Copy one path into destination Branch main. A configured source AgentLane - * copies its configuration plus fresh idle state; a data-only Branch remains - * data-only. Operation/pending/result/usage state is excluded. - */ - scope?: "branch"; - /** Entry to fork from. Defaults to the source main Branch's current tip. */ + scope: "branch"; + /** Source Branch to copy under the same name in the destination. */ + branch: string; + /** Entry to fork from. Defaults to the source Branch's current tip. */ entryId?: string; /** * Whether the fork includes the selected entry or stops at its parent. diff --git a/packages/agent/test/harness/jsonl-v3-migration.test.ts b/packages/agent/test/harness/jsonl-v3-migration.test.ts index 28a573139..8028c2720 100644 --- a/packages/agent/test/harness/jsonl-v3-migration.test.ts +++ b/packages/agent/test/harness/jsonl-v3-migration.test.ts @@ -276,7 +276,11 @@ describe("JSONL v3 migration", () => { it("forks a closed source into a complete v4 destination without rewriting it", async () => { const { path, content, metadata } = await writeForkFixture(); - const fork = await repo.fork(metadata, { id: "closed-fork" }, BACKGROUND_CONTEXT); + const fork = await repo.fork( + metadata, + { id: "closed-fork", scope: "branch", branch: "main" }, + BACKGROUND_CONTEXT, + ); expect(getOrThrow(await fileSystem.readTextFile(path, BACKGROUND_CONTEXT))).toBe(content); expect(fork.metadata.id).toBe("closed-fork"); @@ -309,7 +313,11 @@ describe("JSONL v3 migration", () => { const sourceEntries = await source.findEntries({ order: "asc" }, BACKGROUND_CONTEXT); const sourceStats = await source.getStats(BACKGROUND_CONTEXT); - const fork = await repo.fork(metadata, { id: "open-fork" }, BACKGROUND_CONTEXT); + const fork = await repo.fork( + metadata, + { id: "open-fork", scope: "branch", branch: "main" }, + BACKGROUND_CONTEXT, + ); const forkEntries = await expectForkedState(fork); expect(fork.metadata.id).toBe("open-fork"); diff --git a/packages/agent/test/harness/memory-session-repo.test.ts b/packages/agent/test/harness/memory-session-repo.test.ts index 9ab925b45..fcf8fca77 100644 --- a/packages/agent/test/harness/memory-session-repo.test.ts +++ b/packages/agent/test/harness/memory-session-repo.test.ts @@ -103,6 +103,8 @@ describe("MemorySessionRepo metadata", () => { ); const options = { id: "fork", + scope: "branch" as const, + branch: "main", entryId: childId, position: "before" as "before" | "at", }; diff --git a/packages/session-backends/sqlite-node/src/sqlite/repo.ts b/packages/session-backends/sqlite-node/src/sqlite/repo.ts index b78ec09e3..5168e9d48 100644 --- a/packages/session-backends/sqlite-node/src/sqlite/repo.ts +++ b/packages/session-backends/sqlite-node/src/sqlite/repo.ts @@ -101,6 +101,7 @@ function buildForkSnapshot(source: SqliteStorageSnapshot, options: ForkOptions): }; } +// TODO(WP08): Remove this snapshot path when SQLite forks use streaming staging. function readForkSourceEntries( db: SqliteDatabase, sessionId: string, @@ -108,12 +109,12 @@ function readForkSourceEntries( options: ForkOptions, ): Entry[] { if (options.scope === "tree") return readSourceEntries(db, sessionId); - const mainAddress = branchTip("main"); - const mainTip = scalarValues.find( - (stored) => stored.address.namespace === mainAddress.namespace && stored.address.key === mainAddress.key, + const sourceAddress = branchTip(options.branch); + const sourceTip = scalarValues.find( + (stored) => stored.address.namespace === sourceAddress.namespace && stored.address.key === sourceAddress.key, ) as StoredValue | undefined; - if (mainTip === undefined) throw new Error("Source session is missing main branch"); - const requested = options.entryId ?? mainTip.value; + if (sourceTip === undefined) throw new Error(`Unknown source branch: ${options.branch}`); + const requested = options.entryId ?? sourceTip.value; return requested === null ? [] : scanBranchEntries(db, sessionId, { start: requested, order: "oldestFirst" }); } diff --git a/packages/session-backends/sqlite-node/src/sqlite/storage.ts b/packages/session-backends/sqlite-node/src/sqlite/storage.ts index 3e9586c3e..a9c24e0d7 100644 --- a/packages/session-backends/sqlite-node/src/sqlite/storage.ts +++ b/packages/session-backends/sqlite-node/src/sqlite/storage.ts @@ -129,8 +129,7 @@ export class SqliteStorage implements Storage { return Promise.resolve(readSessionStats(this.db, this.sessionId)); } - snapshot(options: ForkOptions | undefined, _context: Context): Promise { - options ??= {}; + snapshot(options: ForkOptions, _context: Context): Promise { if (this.state !== "open") return Promise.reject(new Error("SqliteStorage is closed")); const result = this.commitQueue.then(() => this.readSnapshot(options)); this.commitQueue = result.then( @@ -151,12 +150,12 @@ export class SqliteStorage implements Storage { private readSnapshotEntries(options: ForkOptions, scalarValues: readonly StoredValue[]): Entry[] { if (options.scope === "tree") return readAllEntryRows(this.db, this.sessionId).map(decodeEntryRow); - const mainAddress = branchTip("main"); - const mainTip = scalarValues.find( - (stored) => stored.address.namespace === mainAddress.namespace && stored.address.key === mainAddress.key, + const sourceAddress = branchTip(options.branch); + const sourceTip = scalarValues.find( + (stored) => stored.address.namespace === sourceAddress.namespace && stored.address.key === sourceAddress.key, ) as StoredValue | undefined; - if (mainTip === undefined) throw new Error("Source session is missing main branch"); - const requested = options.entryId ?? mainTip.value; + if (sourceTip === undefined) throw new Error(`Unknown source branch: ${options.branch}`); + const requested = options.entryId ?? sourceTip.value; return requested === null ? [] : scanBranchEntries(this.db, this.sessionId, { start: requested, order: "oldestFirst" }); diff --git a/packages/session-backends/sqlite-node/test/repo.test.ts b/packages/session-backends/sqlite-node/test/repo.test.ts index ec574ee65..b92d21ca0 100644 --- a/packages/session-backends/sqlite-node/test/repo.test.ts +++ b/packages/session-backends/sqlite-node/test/repo.test.ts @@ -577,7 +577,11 @@ describe("SqliteSessionRepo", () => { ); await source.close(BACKGROUND_CONTEXT); - const fork = await repo.fork(source.metadata, { id: "fork", entryId: "child" }, BACKGROUND_CONTEXT); + const fork = await repo.fork( + source.metadata, + { id: "fork", scope: "branch", branch: "main", entryId: "child" }, + BACKGROUND_CONTEXT, + ); await withDb(fork.metadata.path, (db) => { expect( @@ -776,7 +780,11 @@ describe("SqliteSessionRepo", () => { BACKGROUND_CONTEXT, ); - const fork = await rightRepo.fork(left.metadata, { id: "fork-left" }, BACKGROUND_CONTEXT); + const fork = await rightRepo.fork( + left.metadata, + { id: "fork-left", scope: "branch", branch: "main" }, + BACKGROUND_CONTEXT, + ); expect((await fork.findEntries({ order: "asc" }, BACKGROUND_CONTEXT)).map(({ id }) => id)).toEqual([ "left-root", ]); @@ -827,7 +835,11 @@ describe("SqliteSessionRepo", () => { now: () => 2, }); - const first = await serverRepo.fork(source.metadata, { id: "first-fork" }, BACKGROUND_CONTEXT); + const first = await serverRepo.fork( + source.metadata, + { id: "first-fork", scope: "branch", branch: "main" }, + BACKGROUND_CONTEXT, + ); expect(laterCommitCompleted).toBe(true); expect(databaseFactory.readOnlyOpenCount).toBe(1); @@ -841,7 +853,11 @@ describe("SqliteSessionRepo", () => { ); expect((await first.getStats(BACKGROUND_CONTEXT)).messageCount).toBe(0); - const second = await serverRepo.fork(source.metadata, { id: "second-fork" }, BACKGROUND_CONTEXT); + const second = await serverRepo.fork( + source.metadata, + { id: "second-fork", scope: "branch", branch: "main" }, + BACKGROUND_CONTEXT, + ); expect(databaseFactory.readOnlyOpenCount).toBe(2); expect((await second.findEntries({ order: "asc" }, BACKGROUND_CONTEXT)).map(({ id }) => id)).toEqual([ "root",