Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/bak/__tests__/archive.test.ts

Raw
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<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",
		);
	});
});