import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; import type { SessionInfo } from "@earendil-works/pi-coding-agent"; import { afterEach, describe, expect, it } from "vitest"; import { BakArchive, PartialBackupError } from "../archive.js"; import { S3Store } from "../s3.js"; class MemoryS3 { readonly objects = new Map(); failCatalogPuts = 0; readonly failSessionIds = new Set(); private revision = 0; fetch = async (url: string, init?: RequestInit): Promise => { const key = decodeURIComponent( new URL(url).pathname.slice("/pi-bak/".length), ); const current = this.objects.get(key); if (init?.method === "PUT") { if ([...this.failSessionIds].some((id) => key.endsWith(`/${id}.jsonl`))) return new Response("rejected", { status: 400 }); if (key === "catalog/hosts.json" && this.failCatalogPuts-- > 0) throw new Error("network down"); const headers = new Headers(init.headers); if (headers.get("if-none-match") === "*" && current) return new Response("precondition", { status: 412 }); if (headers.has("if-match") && headers.get("if-match") !== current?.etag) return new Response("precondition", { status: 412 }); const body = new Uint8Array( await new Response(init.body as BodyInit | null).arrayBuffer(), ); const etag = `"${++this.revision}"`; this.objects.set(key, { body, etag }); return new Response(null, { status: 200, headers: { etag } }); } if (!current) return new Response("missing", { status: 404 }); return new Response(current.body, { status: 200, headers: { etag: current.etag }, }); }; } function info(path: string, id: string): SessionInfo { return { path, id, cwd: "/work/project", name: "Archive test", created: new Date("2026-01-02T03:04:05.000Z"), modified: new Date("2026-01-02T03:05:05.000Z"), messageCount: 2, firstMessage: "hello", allMessagesText: "hello world", }; } describe("bak archive", () => { let directory: string | undefined; afterEach(async () => { if (directory) await rm(directory, { recursive: true, force: true }); directory = undefined; }); it("initializes, backs up, refreshes, and restores without listing", async () => { directory = await mkdtemp(join(tmpdir(), "pi-bak-")); const memory = new MemoryS3(); const store = new S3Store( { endpoint: "https://example.invalid", bucket: "pi-bak", region: "auto", accessKeyId: "id", secretAccessKey: "secret", }, memory as never, ); const archive = new BakArchive(store, { identity: join(directory, "state", "identity.json"), cache: join(directory, "cache", "catalog.json"), pendingDir: join(directory, "state", "pending"), }); const identity = await archive.init("work-laptop"); const sessionPath = join(directory, "session.jsonl"); const sessionId = "12345678-1234-1234-1234-123456789abc"; const body = `${JSON.stringify({ type: "session", version: 3, id: sessionId, timestamp: "2026-01-02T03:04:05.000Z", cwd: "/work/project", })}\n${JSON.stringify({ type: "message", id: "a", parentId: null, timestamp: "2026-01-02T03:04:06.000Z", message: { role: "user", content: "hello" } })}\n`; await writeFile(sessionPath, body); await expect( archive.backup([info(sessionPath, sessionId)]), ).resolves.toEqual({ uploaded: 1, indexed: 1, }); const catalog = await archive.refresh(); expect(catalog.indexes[0]).toMatchObject({ hostId: identity.hostId, alias: "work-laptop", sessions: [{ id: sessionId, name: "Archive test" }], }); const destination = join(directory, "restore"); await expect( archive.restore(catalog.indexes[0]?.sessions ?? [], destination), ).resolves.toEqual({ restored: 1, skipped: 0, conflicts: 0 }); await expect( archive.restore(catalog.indexes[0]?.sessions ?? [], destination), ).resolves.toEqual({ restored: 0, skipped: 1, conflicts: 0 }); const restored = join( destination, `2026-01-02T03-04-05-000Z_${sessionId}.jsonl`, ); expect(await readFile(restored, "utf8")).toBe(body); const divergent = `${body}${JSON.stringify({ type: "message", id: "b", parentId: "a", timestamp: "2026-01-02T03:04:07.000Z", message: { role: "assistant", content: [] } })}\n`; await writeFile(restored, divergent); await expect( archive.restore(catalog.indexes[0]?.sessions ?? [], destination), ).resolves.toEqual({ restored: 0, skipped: 0, conflicts: 1 }); expect(await readFile(restored, "utf8")).toBe(divergent); expect([...memory.objects.keys()].sort()).toEqual([ "catalog/hosts.json", `hosts/${identity.hostId}.json`, `sessions/${identity.hostId}/${sessionId}.jsonl`, ]); }); it("rejects unsafe remote session records", async () => { directory = await mkdtemp(join(tmpdir(), "pi-bak-")); const memory = new MemoryS3(); const archive = new BakArchive( new S3Store( { endpoint: "https://example.invalid", bucket: "pi-bak", region: "auto", accessKeyId: "id", secretAccessKey: "secret", }, memory as never, ), { identity: join(directory, "state", "identity.json"), cache: join(directory, "cache", "catalog.json"), pendingDir: join(directory, "state", "pending"), }, ); await expect( archive.restore( [ { id: "x/../../escaped", key: "sessions/host/x/../../escaped.jsonl", cwd: "/work", created: "2026-01-02T03:04:05.000Z", }, ], join(directory, "restore"), ), ).rejects.toThrow("Invalid archived session record"); }); it("rejects a session whose first line is blank", async () => { directory = await mkdtemp(join(tmpdir(), "pi-bak-")); const memory = new MemoryS3(); const store = new S3Store( { endpoint: "https://example.invalid", bucket: "pi-bak", region: "auto", accessKeyId: "id", secretAccessKey: "secret", }, memory as never, ); const archive = new BakArchive(store, { identity: join(directory, "state", "identity.json"), cache: join(directory, "cache", "catalog.json"), pendingDir: join(directory, "state", "pending"), }); await archive.init("work-laptop"); const sessionId = "12345678-1234-1234-1234-123456789abc"; const sessionPath = join(directory, "blank-header.jsonl"); await writeFile( sessionPath, `\n${JSON.stringify({ type: "session", version: 3, id: sessionId, timestamp: "2026-01-02T03:04:05.000Z", cwd: "/work/project", })}\n`, ); await expect( archive.backup([info(sessionPath, sessionId)]), ).rejects.toThrow("Invalid Pi session header"); }); it("merges concurrent host-index updates", async () => { directory = await mkdtemp(join(tmpdir(), "pi-bak-")); const memory = new MemoryS3(); const archive = new BakArchive( new S3Store( { endpoint: "https://example.invalid", bucket: "pi-bak", region: "auto", accessKeyId: "id", secretAccessKey: "secret", }, memory as never, ), { identity: join(directory, "state", "identity.json"), cache: join(directory, "cache", "catalog.json"), pendingDir: join(directory, "state", "pending"), }, ); await archive.init("work-laptop"); const ids = [ "12345678-1234-1234-1234-123456789abc", "abcdefab-1234-1234-1234-123456789abc", ]; const paths = await Promise.all( ids.map(async (id, index) => { const path = join(directory, `${index}.jsonl`); await writeFile( path, `${JSON.stringify({ type: "session", version: 3, id, timestamp: `2026-01-02T03:04:0${index}.000Z`, cwd: "/work/project" })}\n`, ); return path; }), ); await Promise.all([ archive.backup([info(paths[0] ?? "", ids[0] ?? "")]), archive.backup([info(paths[1] ?? "", ids[1] ?? "")]), ]); const catalog = await archive.refresh(); expect( catalog.indexes[0]?.sessions.map((session) => session.id).sort(), ).toEqual([...ids].sort()); }); it("drains failed batches and indexes successful transfers", async () => { directory = await mkdtemp(join(tmpdir(), "pi-bak-")); const memory = new MemoryS3(); const archive = new BakArchive( new S3Store( { endpoint: "https://example.invalid", bucket: "pi-bak", region: "auto", accessKeyId: "id", secretAccessKey: "secret", }, memory as never, ), { identity: join(directory, "state", "identity.json"), cache: join(directory, "cache", "catalog.json"), pendingDir: join(directory, "state", "pending"), }, ); await archive.init("work-laptop"); const ids = Array.from( { length: 12 }, (_, index) => `00000000-0000-0000-0000-${String(index).padStart(12, "0")}`, ); for (const id of ids.slice(0, 6)) memory.failSessionIds.add(id); const infos = await Promise.all( ids.map(async (id, index) => { const path = join(directory, `batch-${index}.jsonl`); await writeFile( path, `${JSON.stringify({ type: "session", version: 3, id, timestamp: "2026-01-02T03:04:05.000Z", cwd: "/work/project" })}\n`, ); return info(path, id); }), ); let failure: unknown; try { await archive.backup(infos); } catch (error) { failure = error; } expect(failure).toBeInstanceOf(PartialBackupError); expect((failure as PartialBackupError).succeededPaths).toHaveLength(6); const indexed = (await archive.refresh()).indexes[0]?.sessions ?? []; expect(indexed).toHaveLength(6); expect( indexed.every((record) => !memory.failSessionIds.has(record.id)), ).toBe(true); }); it("reuses durable identity after interrupted initialization", async () => { directory = await mkdtemp(join(tmpdir(), "pi-bak-")); const memory = new MemoryS3(); memory.failCatalogPuts = 3; const archive = new BakArchive( new S3Store( { endpoint: "https://example.invalid", bucket: "pi-bak", region: "auto", accessKeyId: "id", secretAccessKey: "secret", }, memory as never, ), { identity: join(directory, "state", "identity.json"), cache: join(directory, "cache", "catalog.json"), pendingDir: join(directory, "state", "pending"), }, ); await expect(archive.init("work-laptop")).rejects.toThrow("network down"); const pendingIdentity = await archive.identity(); expect(pendingIdentity).toBeDefined(); const sessionId = "12345678-1234-1234-1234-123456789abc"; const sessionPath = join(directory, "recovered.jsonl"); await writeFile( sessionPath, `${JSON.stringify({ type: "session", version: 3, id: sessionId, timestamp: "2026-01-02T03:04:05.000Z", cwd: "/work/project" })}\n`, ); await expect( archive.backup([info(sessionPath, sessionId)]), ).resolves.toMatchObject({ uploaded: 1 }); expect((await archive.refresh()).catalog.hosts).toEqual([pendingIdentity]); }); it("rejects duplicate aliases through the shared catalog", async () => { directory = await mkdtemp(join(tmpdir(), "pi-bak-")); const memory = new MemoryS3(); const config = { endpoint: "https://example.invalid", bucket: "pi-bak", region: "auto", accessKeyId: "id", secretAccessKey: "secret", }; const first = new BakArchive(new S3Store(config, memory as never), { identity: join(directory, "one", "identity.json"), cache: join(directory, "one", "catalog.json"), pendingDir: join(directory, "one", "pending"), }); const second = new BakArchive(new S3Store(config, memory as never), { identity: join(directory, "two", "identity.json"), cache: join(directory, "two", "catalog.json"), pendingDir: join(directory, "two", "pending"), }); await first.init("work-laptop"); await expect(second.init("work-laptop")).rejects.toThrow( "Alias already exists: work-laptop", ); }); });