mirror of
https://github.com/earendil-works/pi.git
synced 2026-09-29 17:19:06 +08:00
feat(agent): require explicit named branches for session forks
Require fork scope and source branch selection, preserve the branch name in the destination, and cover named-branch behavior across backends.
This commit is contained in:
@@ -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<string, Entry>();
|
||||
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<string>;
|
||||
sourceToDestinationTip: Map<string, { name: string; tipId: string | null }>;
|
||||
destinationTips: Map<string, string | null>;
|
||||
} {
|
||||
const entryIds = new Set<string>();
|
||||
const sourceToDestinationTip = new Map<string, { name: string; tipId: string | null }>();
|
||||
const destinationTips = new Map<string, string | null>();
|
||||
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 ||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -336,7 +336,7 @@ export function createSessionRepoForkBehaviorConformance<TMetadata extends Sessi
|
||||
createCase(
|
||||
factory,
|
||||
"forks",
|
||||
"forks one configured branch with scoped values and a zero ledger",
|
||||
"forks one named configured branch with scoped values and a zero ledger",
|
||||
async ({ repo }) => {
|
||||
const source = await repo.create({ id: "source" }, BACKGROUND_CONTEXT);
|
||||
await source.mutate(
|
||||
@@ -361,9 +361,10 @@ export function createSessionRepoForkBehaviorConformance<TMetadata extends Sessi
|
||||
type: "custom",
|
||||
customType: "sibling",
|
||||
}),
|
||||
setValue(branchTip("main"), CHILD_ID),
|
||||
setValue(laneConfig("main"), configuration),
|
||||
setValue(laneState("main"), {
|
||||
setValue(branchTip("main"), SIBLING_ID),
|
||||
setValue(branchTip("review"), CHILD_ID),
|
||||
setValue(laneConfig("review"), configuration),
|
||||
setValue(laneState("review"), {
|
||||
currentOperationId: OPERATION_ID,
|
||||
lastOperationId: "previous",
|
||||
inbox: [{ entryId: PENDING_ID, kind: "write" }],
|
||||
@@ -388,7 +389,7 @@ export function createSessionRepoForkBehaviorConformance<TMetadata extends Sessi
|
||||
}),
|
||||
setValue(operationMeta(OPERATION_ID), {
|
||||
operationId: OPERATION_ID,
|
||||
lane: "main",
|
||||
lane: "review",
|
||||
sourceTipId: CHILD_ID,
|
||||
startedAt: 1,
|
||||
intent: { kind: "compaction" },
|
||||
@@ -431,7 +432,7 @@ export function createSessionRepoForkBehaviorConformance<TMetadata extends Sessi
|
||||
|
||||
const fork = await repo.fork(
|
||||
source.metadata,
|
||||
{ id: "fork", entryId: CHILD_ID, position: "at" },
|
||||
{ id: "fork", scope: "branch", branch: "review", entryId: CHILD_ID, position: "at" },
|
||||
BACKGROUND_CONTEXT,
|
||||
);
|
||||
|
||||
@@ -439,9 +440,12 @@ export function createSessionRepoForkBehaviorConformance<TMetadata extends Sessi
|
||||
(await fork.findEntries({ order: "asc" }, BACKGROUND_CONTEXT)).map(({ id }) => 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<TMetadata extends Sessi
|
||||
|
||||
const before = await repo.fork(
|
||||
source.metadata,
|
||||
{ id: "before", entryId: CHILD_ID, position: "before" },
|
||||
{ id: "before", scope: "branch", branch: "main", entryId: CHILD_ID, position: "before" },
|
||||
BACKGROUND_CONTEXT,
|
||||
);
|
||||
strictEqual(await getBranchTip(before), ROOT_ID);
|
||||
@@ -511,7 +515,13 @@ export function createSessionRepoForkBehaviorConformance<TMetadata extends Sessi
|
||||
(await before.findEntries({ order: "asc" }, BACKGROUND_CONTEXT)).map(({ id }) => 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<TMetadata extends Sessi
|
||||
);
|
||||
await source.close(BACKGROUND_CONTEXT);
|
||||
|
||||
const fork = await repo.fork(source.metadata, { id: "fork" }, BACKGROUND_CONTEXT);
|
||||
const fork = await repo.fork(
|
||||
source.metadata,
|
||||
{ id: "fork", scope: "branch", branch: "main" },
|
||||
BACKGROUND_CONTEXT,
|
||||
);
|
||||
strictEqual(await getBranchTip(fork), ROOT_ID);
|
||||
strictEqual(await fork.getValue(laneConfig("main"), BACKGROUND_CONTEXT), undefined);
|
||||
strictEqual(await fork.getValue(laneState("main"), BACKGROUND_CONTEXT), undefined);
|
||||
@@ -675,7 +689,11 @@ export function createSessionRepoForkSourceSnapshotConformance<TMetadata extends
|
||||
],
|
||||
BACKGROUND_CONTEXT,
|
||||
);
|
||||
const fork = repo.fork(source.metadata, { id: "fork" }, BACKGROUND_CONTEXT);
|
||||
const fork = repo.fork(
|
||||
source.metadata,
|
||||
{ id: "fork", scope: "branch", branch: "main" },
|
||||
BACKGROUND_CONTEXT,
|
||||
);
|
||||
const secondCommit = source.mutate(
|
||||
(mutator) =>
|
||||
mutator.commit(
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -103,6 +103,8 @@ describe("MemorySessionRepo metadata", () => {
|
||||
);
|
||||
const options = {
|
||||
id: "fork",
|
||||
scope: "branch" as const,
|
||||
branch: "main",
|
||||
entryId: childId,
|
||||
position: "before" as "before" | "at",
|
||||
};
|
||||
|
||||
@@ -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<string | null> | 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" });
|
||||
}
|
||||
|
||||
|
||||
@@ -129,8 +129,7 @@ export class SqliteStorage implements Storage {
|
||||
return Promise.resolve(readSessionStats(this.db, this.sessionId));
|
||||
}
|
||||
|
||||
snapshot(options: ForkOptions | undefined, _context: Context): Promise<SqliteStorageSnapshot> {
|
||||
options ??= {};
|
||||
snapshot(options: ForkOptions, _context: Context): Promise<SqliteStorageSnapshot> {
|
||||
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<unknown>[]): 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<string | null> | 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" });
|
||||
|
||||
@@ -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",
|
||||
|
||||
Reference in New Issue
Block a user