repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/bak/__e2e__/r2.test.ts
Rawimport { 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(),
);
});
});