Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/bak/__e2e__/r2.test.ts

Raw
import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises";
import { tmpdir } from "node:os";
import { dirname, join } from "node:path";
import { afterAll, beforeAll, describe, expect, it } from "vitest";
import { BakArchive, sessionInfoFromPath } from "../archive.js";
import { type S3Config, S3Store } from "../s3.js";
import type { StatePaths } from "../state.js";
import {
	appendMessage,
	removeHosts,
	sessionBody,
	storageConfig,
	testAlias,
	writeSession,
} from "./support.js";

const config = storageConfig;
const hostIds: string[] = [];
let workspace: string;

function statePathsIn(name: string): StatePaths {
	const root = join(workspace, name);
	return {
		identity: join(root, "identity.json"),
		cache: join(root, "catalog.json"),
		pendingDir: join(root, "pending"),
	};
}

function newArchive(paths: StatePaths): BakArchive {
	return new BakArchive(new S3Store(config as S3Config), paths);
}

async function initHost(name: string, alias: string) {
	const paths = statePathsIn(name);
	const archive = newArchive(paths);
	const identity = await archive.init(alias);
	hostIds.push(identity.hostId);
	return { archive, paths, identity };
}

async function sessionDir(name: string): Promise<string> {
	const directory = join(workspace, name);
	await mkdir(directory, { recursive: true });
	return directory;
}

describe.skipIf(!config)("bak archive against real S3 storage", () => {
	beforeAll(async () => {
		workspace = await mkdtemp(join(tmpdir(), "bak-e2e-"));
	});

	afterAll(async () => {
		if (config) await removeHosts(config, hostIds);
		await rm(workspace, { recursive: true, force: true });
	});

	it("registers a host and creates an empty index", async () => {
		const alias = testAlias();
		const { archive, identity } = await initHost("alpha", alias);

		expect(identity.alias).toBe(alias);
		const cache = await archive.refresh();
		const index = cache.indexes.find(
			(entry) => entry.hostId === identity.hostId,
		);
		expect(cache.catalog.hosts).toContainEqual(identity);
		expect(index?.sessions).toEqual([]);
	});

	it("rejects a duplicate alias and rolls back the local identity", async () => {
		const alias = testAlias();
		const first = await initHost("collide-a", alias);
		const paths = statePathsIn("collide-b");
		const archive = newArchive(paths);

		await expect(archive.init(alias)).rejects.toThrow(/Alias already exists/);
		await expect(readFile(paths.identity, "utf8")).rejects.toMatchObject({
			code: "ENOENT",
		});
		expect(first.identity.alias).toBe(alias);
	});

	it("uploads sessions, indexes them, and stays idempotent", async () => {
		const { archive, identity } = await initHost("backup", testAlias());
		const directory = await sessionDir("backup-sessions");
		const first = await writeSession(directory, { name: "first session" });
		const second = await writeSession(directory, { messages: ["a", "b"] });

		expect(await archive.backup([first, second])).toEqual({
			uploaded: 2,
			indexed: 2,
		});
		expect(await archive.backup([first, second])).toEqual({
			uploaded: 2,
			indexed: 2,
		});

		const cache = await archive.refresh();
		const index = cache.indexes.find(
			(entry) => entry.hostId === identity.hostId,
		);
		expect(index?.sessions).toHaveLength(2);
		expect(index?.sessions.find((record) => record.id === first.id)?.name).toBe(
			"first session",
		);
		expect(index?.sessions.map((record) => record.id).sort()).toEqual(
			[first.id, second.id].sort(),
		);
	});

	it("extends an appended session and rejects divergent rewrites", async () => {
		const { archive } = await initHost("append", testAlias());
		const directory = await sessionDir("append-sessions");
		const info = await writeSession(directory, { messages: ["one"] });
		await archive.backup([info]);

		const grown = await appendMessage(info.path, "two");
		expect(await archive.backup([grown])).toEqual({ uploaded: 1, indexed: 1 });

		await writeFile(
			info.path,
			sessionBody(info.id, info.created.toISOString(), info.cwd, ["other"]),
		);
		const diverged = await sessionInfoFromPath(info.path);
		await expect(archive.backup([diverged])).rejects.toThrow(
			/Archived session diverged/,
		);
	});

	it("restores, skips identical files, and reports conflicts", async () => {
		const { archive, identity } = await initHost("restore", testAlias());
		const directory = await sessionDir("restore-source");
		const info = await writeSession(directory, { messages: ["restore me"] });
		await archive.backup([info]);
		const cache = await archive.refresh();
		const records =
			cache.indexes.find((entry) => entry.hostId === identity.hostId)
				?.sessions ?? [];
		expect(records).toHaveLength(1);

		const target = join(workspace, "restore-target");
		expect(await archive.restore(records, target)).toEqual({
			restored: 1,
			skipped: 0,
			conflicts: 0,
		});
		expect(await archive.restore(records, target)).toEqual({
			restored: 0,
			skipped: 1,
			conflicts: 0,
		});

		const conflictDir = await sessionDir("restore-conflict");
		await writeFile(
			join(conflictDir, `conflict_${info.id}.jsonl`),
			sessionBody(info.id, info.created.toISOString(), info.cwd, ["local"]),
		);
		expect(await archive.restore(records, conflictDir)).toEqual({
			restored: 0,
			skipped: 0,
			conflicts: 1,
		});
	});

	it("survives concurrent index updates from two archives", async () => {
		const { archive, identity, paths } = await initHost("cas-a", testAlias());
		const secondPaths = statePathsIn("cas-b");
		await mkdir(dirname(secondPaths.identity), { recursive: true });
		await writeFile(secondPaths.identity, `${JSON.stringify(identity)}\n`);
		const other = newArchive(secondPaths);
		const directory = await sessionDir("cas-sessions");
		const left = await writeSession(directory, { messages: ["left"] });
		const right = await writeSession(directory, { messages: ["right"] });

		const [a, b] = await Promise.all([
			archive.backup([left]),
			other.backup([right]),
		]);

		expect(a.uploaded).toBe(1);
		expect(b.uploaded).toBe(1);
		const cache = await newArchive(paths).refresh();
		const index = cache.indexes.find(
			(entry) => entry.hostId === identity.hostId,
		);
		expect(index?.sessions.map((record) => record.id).sort()).toEqual(
			[left.id, right.id].sort(),
		);
	});
});