repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/bak/__tests__/archive.test.ts
Rawimport { 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<string, { body: Uint8Array; etag: string }>();
failCatalogPuts = 0;
readonly failSessionIds = new Set<string>();
private revision = 0;
fetch = async (url: string, init?: RequestInit): Promise<Response> => {
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",
);
});
});