From 3a56dcd9e598ca861dbfa14d31703f0669ece357 Mon Sep 17 00:00:00 2001 From: Codinget Date: Wed, 22 Jul 2026 23:36:21 +0000 Subject: [PATCH] feat: add export/import to the postgres backend; import auto-detects backend MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The export/import plan predated the postgres backend. Bring it up to parity: - PostgresInterface implements Exportable + Importable, mirroring the sqlite backend. Export is a signal-checked async generator (one query per table); import drives a transaction via `WrappedPool.multi`, which rolls back on any error/abort and commits only after the whole stream is consumed cleanly. - The `note` kind is skipped on postgres (its note CRUD is still unimplemented, so a pg database holds none) — a full dump from sqlite imports its user/abode/resident/apikey records and drops notes. - Unlike sqlite (PRAGMA foreign_keys=off), postgres keeps FK enforcement; the FK-safe insertion order keeps a full dump valid, and truly-dangling partial dumps will (correctly) fail. - abode-import now resolves the backend via getDbInterface instead of constructing SqliteInterface directly; isImportable keeps it from ever running over the remote (api) interface, which has no import. - Add isImportable(); make postgres selectClientApikeys' where optional. - Tests: exercise the real PostgresInterface export/import paths against an in-memory fake WrappedPgClient (no pg service in CI) — NDJSON shape, meta.source, filtering, note-skip, insert dispatch, and error/abort rollback. Co-Authored-By: Claude Opus 4.8 --- src/bin/abode-import.ts | 25 +- src/db/postgres/PostgresInterface.ts | 208 ++++++++++++- src/db/postgres/query.ts | 4 +- src/db/types/ExportImport.ts | 4 + test/backends/postgres/export-import.test.ts | 296 +++++++++++++++++++ 5 files changed, 524 insertions(+), 13 deletions(-) create mode 100644 test/backends/postgres/export-import.test.ts diff --git a/src/bin/abode-import.ts b/src/bin/abode-import.ts index d8e0c2e..501892b 100644 --- a/src/bin/abode-import.ts +++ b/src/bin/abode-import.ts @@ -1,9 +1,7 @@ import { createReadStream } from "node:fs"; -import { getWrappedDb } from "../db/sqlite/impl/index.js"; -import { SqliteInterface } from "../db/sqlite/SqliteInterface.js"; -import { parseSqliteUrl } from "../db/sqlite/url.js"; +import { getDbInterface } from "../db/index.js"; import { isExportKind } from "../db/export/filter.js"; -import type { ExportFilter } from "../db/types/ExportImport.js"; +import { isImportable, type ExportFilter } from "../db/types/ExportImport.js"; const args = process.argv.slice(2); @@ -15,14 +13,15 @@ function printUsage(err: boolean | string = false): never { } log("Usage:"); log("\tabode-import --help"); - log("\tabode-import [--kinds=...] \\"); + log("\tabode-import [--kinds=...] \\"); log( "\t [--exclude-kinds=...] [--abodes=aid,...] [--users=uid,...]", ); log(""); log( - "The target database must already be migrated (run abode-migrate first).", + "The target must be a local backend (sqlite or postgres), already migrated", ); + log("(run abode-migrate first). Remote (api) targets are not importable."); process.exit(err ? 1 : 0); } @@ -42,7 +41,7 @@ for (const arg of args) { const url = positional[0]; const input = positional[1]; -if (!url) printUsage("missing "); +if (!url) printUsage("missing "); if (!input) printUsage("missing "); if (positional.length > 2) printUsage("too many arguments"); @@ -73,9 +72,15 @@ if (abodes) filter.abodes = abodes; const users = parseList(flags.get("users")); if (users) filter.users = users; -// Import is sqlite-only: construct the backend directly rather than resolving -// generically, so it can never be pointed at a remote (api) target. -const db = new SqliteInterface(getWrappedDb(...parseSqliteUrl(url))); +// Resolve the backend generically. Import lives on the local backends (sqlite, +// postgres); the remote (api) interface has no `import`, so `isImportable` +// keeps it from ever running over HTTP. +const db = await getDbInterface(url); +if (!isImportable(db)) { + console.error(`Error: backend '${db.name}' does not support import`); + await db.close().catch(() => {}); + process.exit(1); +} const source = input === "-" ? process.stdin : createReadStream(input); diff --git a/src/db/postgres/PostgresInterface.ts b/src/db/postgres/PostgresInterface.ts index a92aa21..3be560c 100644 --- a/src/db/postgres/PostgresInterface.ts +++ b/src/db/postgres/PostgresInterface.ts @@ -43,8 +43,21 @@ import type { UpdateNote, } from "../types/Note.js"; import type { WrappedPgClient } from "./pool.js"; +import { Readable } from "node:stream"; +import readline from "node:readline"; +import type { + Exportable, + ExportKind, + ExportOptions, + Importable, + ImportOptions, + ImportResult, +} from "../types/ExportImport.js"; +import { isExportKind, kindAllowed, recordAllowed } from "../export/filter.js"; -export class PostgresInterface implements BackendDbInterface { +export class PostgresInterface + implements BackendDbInterface, Exportable, Importable +{ #db: WrappedPgClient; constructor(db: WrappedPgClient) { @@ -508,4 +521,197 @@ export class PostgresInterface implements BackendDbInterface { async listNotesByUserId(_uid: string): Promise { throw new Error("Unimplemented"); } + + export(options: ExportOptions = {}): NodeJS.ReadableStream { + const { filter, signal } = options; + const db = this.#db; + const source = this.name; + + // `note` is omitted: the postgres backend has no note CRUD yet, so a pg + // database can hold none. Each `load()` is a single query, run lazily and + // skipped once the destination aborts. + const tables: [ + ExportKind, + () => Promise<{ uid?: string; aid?: string }[]>, + ][] = [ + ["user", () => selectClientUsers(db)], + ["abode", () => selectAbodes(db)], + ["resident", () => selectResidents(db)], + ["apikey", () => selectClientApikeys(db)], + ]; + + async function* generate(): AsyncGenerator { + if (signal?.aborted) return; + yield JSON.stringify({ + kind: "meta", + data: { + v: 1, + exportedAt: new Date().toISOString(), + source, + filter: filter ?? {}, + }, + }) + "\n"; + try { + for (const [kind, load] of tables) { + if (signal?.aborted) return; + if (!kindAllowed(filter, kind)) continue; + for (const row of await load()) { + if (signal?.aborted) return; + if (recordAllowed(filter, kind, row)) { + yield JSON.stringify({ kind, data: row }) + "\n"; + } + } + } + } catch (e) { + if (signal?.aborted) return; + yield JSON.stringify({ + kind: "error", + data: { + message: e instanceof Error ? e.message : String(e), + code: e instanceof Error ? e.name : undefined, + }, + }) + "\n"; + } + } + + return Readable.from(generate()); + } + + async import( + source: NodeJS.ReadableStream, + options: ImportOptions = {}, + ): Promise { + this.#checkReadonly(); + const { filter, signal } = options; + const counts: Partial> = {}; + + const rl = readline.createInterface({ + input: source, + crlfDelay: Infinity, + signal, + }); + + // Postgres holds FK enforcement (unlike sqlite, which we toggle off): the + // FK-safe insertion order — user, abode, resident, apikey — keeps a full + // dump valid. `notes` are skipped entirely (unsupported on this backend). + // The transaction commits only if the whole stream is consumed cleanly; an + // error/abort rolls it back via `multi`. + try { + await this.#db.rethrow(() => + this.#db.multi(async (tx) => { + for await (const raw of rl) { + signal?.throwIfAborted(); + const line = raw.trim(); + if (!line) continue; + let parsed: { kind?: unknown; data?: unknown }; + try { + parsed = JSON.parse(line); + } catch { + throw new InvalidAbodeError(); + } + if (parsed.kind === "meta") continue; + if (parsed.kind === "error") { + throw new Error( + `export stream reported an error: ${ + (parsed.data as { message?: string })?.message ?? "unknown" + }`, + ); + } + if (!isExportKind(parsed.kind)) continue; + if (parsed.kind === "note") continue; // unsupported on postgres + if (!kindAllowed(filter, parsed.kind)) continue; + const data = parsed.data as { uid?: string; aid?: string }; + if (!recordAllowed(filter, parsed.kind, data)) continue; + await this.#importRecord(tx, parsed.kind, parsed.data); + counts[parsed.kind] = (counts[parsed.kind] ?? 0) + 1; + } + // An abort while blocked on the source closes readline without + // throwing, so re-check before the transaction commits. + signal?.throwIfAborted(); + }), + ); + } finally { + rl.close(); + } + + return { counts }; + } + + async #importRecord( + tx: WrappedPgClient, + kind: ExportKind, + data: unknown, + ): Promise { + switch (kind) { + case "user": { + const u = data as ClientUser; + // `password` is never exported; imported users land on the schema + // default ('#unset') and must reset before they can log in. + await tx.run(sql` + INSERT INTO "users"("uid", "email", "name", "flags", "created_at", "updated_at") + VALUES( + ${{ uuid: u.uid }}, + ${{ text: u.email }}, + ${{ text: u.name }}, + ${{ jsonb: u.flags }}, + ${{ date: u.created_at }}, + ${{ date: u.updated_at }} + ) + `); + break; + } + case "abode": { + const a = data as Abode; + await tx.run(sql` + INSERT INTO "abodes"("aid", "name", "created_at", "created_by", "updated_at", "updated_by") + VALUES( + ${{ uuid: a.aid }}, + ${{ text: a.name }}, + ${{ date: a.created_at }}, + ${a.created_by ? { uuid: a.created_by } : { null: true }}, + ${{ date: a.updated_at }}, + ${a.updated_by ? { uuid: a.updated_by } : { null: true }} + ) + `); + break; + } + case "resident": { + const r = data as Resident; + await tx.run(sql` + INSERT INTO "residents"("uid", "aid", "flags", "created_at", "created_by", "updated_at", "updated_by") + VALUES( + ${{ uuid: r.uid }}, + ${{ uuid: r.aid }}, + ${{ jsonb: r.flags }}, + ${{ date: r.created_at }}, + ${r.created_by ? { uuid: r.created_by } : { null: true }}, + ${{ date: r.updated_at }}, + ${r.updated_by ? { uuid: r.updated_by } : { null: true }} + ) + `); + break; + } + case "apikey": { + const k = data as ClientApikey; + // `token` is never exported; mint a fresh unique one so the record's + // metadata (kid/permissions/expiry) survives even though the original + // secret cannot. + await tx.run(sql` + INSERT INTO "apikeys"("uid", "kid", "token", "name", "permissions", "created_at", "expires_at") + VALUES( + ${{ uuid: k.uid }}, + ${{ uuid: k.kid }}, + ${{ text: createApikeyToken() }}, + ${{ text: k.name }}, + ${{ jsonb: k.permissions }}, + ${{ date: k.created_at }}, + ${k.expires_at ? { date: k.expires_at } : { null: true }} + ) + `); + break; + } + case "note": + break; // unsupported on postgres; skipped before reaching here + } + } } diff --git a/src/db/postgres/query.ts b/src/db/postgres/query.ts index 766bb04..d74f4d8 100644 --- a/src/db/postgres/query.ts +++ b/src/db/postgres/query.ts @@ -130,10 +130,10 @@ export async function selectClientApikey( } export async function selectClientApikeys( db: WrappedPgClient, - where: SqlCode, + where?: SqlCode, ): Promise { const rows = await db.all( - sql`${sqlClientApikey} WHERE ${where}`, + where ? sql`${sqlClientApikey} WHERE ${where}` : sqlClientApikey, ); return rows.map(pgToClientApikey); } diff --git a/src/db/types/ExportImport.ts b/src/db/types/ExportImport.ts index e1afa39..812defe 100644 --- a/src/db/types/ExportImport.ts +++ b/src/db/types/ExportImport.ts @@ -74,3 +74,7 @@ export type ExportMeta = { export function isExportable(db: DbInterface): db is DbInterface & Exportable { return typeof (db as Partial).export === "function"; } + +export function isImportable(db: DbInterface): db is DbInterface & Importable { + return typeof (db as Partial).import === "function"; +} diff --git a/test/backends/postgres/export-import.test.ts b/test/backends/postgres/export-import.test.ts new file mode 100644 index 0000000..3d1b789 --- /dev/null +++ b/test/backends/postgres/export-import.test.ts @@ -0,0 +1,296 @@ +import { describe, it } from "node:test"; +import assert from "node:assert/strict"; +import { Readable } from "node:stream"; +import { PostgresInterface } from "../../../src/db/postgres/PostgresInterface.js"; +import type { WrappedPgClient } from "../../../src/db/postgres/pool.js"; +import type { SqlCode } from "../../../src/db/postgres/sql.js"; + +// The postgres backend has no CI database, so these tests drive the real +// PostgresInterface.export/import code paths against an in-memory fake client. +// They verify control flow (NDJSON shape, meta.source, filtering, note +// skipping, counts, insert dispatch, abort -> rollback); the SQL-arg forms are +// the same {uuid}/{text}/{jsonb}/{date} patterns the pg backend's own CRUD +// already exercises against real Postgres. + +interface Rows { + users?: Record[]; + abodes?: Record[]; + residents?: Record[]; + apikeys?: Record[]; +} + +class FakePg implements WrappedPgClient { + readonly = false; + inserts: { table: string; vars: unknown[] }[] = []; + committed = false; + rolledBack = false; + #rows: Rows; + + constructor(rows: Rows = {}) { + this.#rows = rows; + } + + async destroy(): Promise {} + + async all(stmt: SqlCode): Promise { + const s = stmt._sql; + if (s.includes('FROM "users"')) return (this.#rows.users ?? []) as R[]; + if (s.includes('FROM "abodes"')) return (this.#rows.abodes ?? []) as R[]; + if (s.includes('FROM "residents"')) + return (this.#rows.residents ?? []) as R[]; + if (s.includes('FROM "apikeys"')) return (this.#rows.apikeys ?? []) as R[]; + return [] as R[]; + } + + async get(stmt: SqlCode): Promise { + const rows = await this.all(stmt); + return rows[0] ?? null; + } + + async run(stmt: SqlCode): Promise<{ changes: number }> { + const table = stmt._sql.match(/INSERT INTO "(\w+)"/)?.[1] ?? "?"; + this.inserts.push({ table, vars: stmt._vars }); + return { changes: 1 }; + } + + async multi(fn: (tx: WrappedPgClient) => Promise): Promise { + try { + const r = await fn(this); + this.committed = true; + return r; + } catch (e) { + this.rolledBack = true; + throw e; + } + } + + async rethrow(fn: () => Promise): Promise { + return fn(); + } +} + +const iso = "2026-01-02T03:04:05.000Z"; + +function seededRows(): Rows { + return { + users: [ + { + uid: "11111111-1111-1111-1111-111111111111", + email: "u1@test.example", + name: "User One", + flags: {}, + created_at: new Date(iso), + updated_at: new Date(iso), + }, + ], + abodes: [ + { + aid: "aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaa1", + name: "Abode One", + created_at: new Date(iso), + created_by: "11111111-1111-1111-1111-111111111111", + updated_at: new Date(iso), + updated_by: null, + }, + { + aid: "aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaa2", + name: "Abode Two", + created_at: new Date(iso), + created_by: null, + updated_at: new Date(iso), + updated_by: null, + }, + ], + residents: [ + { + uid: "11111111-1111-1111-1111-111111111111", + aid: "aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaa1", + flags: {}, + created_at: new Date(iso), + created_by: null, + updated_at: new Date(iso), + updated_by: null, + }, + ], + apikeys: [ + { + uid: "11111111-1111-1111-1111-111111111111", + kid: "kkkkkkkk-kkkk-kkkk-kkkk-kkkkkkkkkkk1", + name: "key one", + permissions: {}, + created_at: new Date(iso), + expires_at: null, + }, + ], + }; +} + +async function streamToString(s: NodeJS.ReadableStream): Promise { + let out = ""; + for await (const chunk of s) out += chunk; + return out; +} +function parseLines(ndjson: string): { kind: string; data: any }[] { + return ndjson + .split("\n") + .filter(Boolean) + .map((l) => JSON.parse(l)); +} +function line(kind: string, data: unknown): string { + return JSON.stringify({ kind, data }) + "\n"; +} + +describe("postgres export", () => { + it("streams meta + the four supported kinds, never note", async () => { + const db = new PostgresInterface(new FakePg(seededRows())); + const lines = parseLines(await streamToString(db.export())); + + const meta = lines.find((l) => l.kind === "meta"); + assert.ok(meta); + assert.equal(meta!.data.source, "postgres"); + assert.equal(meta!.data.v, 1); + + const kinds = new Set(lines.map((l) => l.kind)); + assert.ok(kinds.has("user")); + assert.ok(kinds.has("abode")); + assert.ok(kinds.has("resident")); + assert.ok(kinds.has("apikey")); + assert.ok(!kinds.has("note"), "note is unsupported on postgres"); + }); + + it("applies the kinds filter", async () => { + const db = new PostgresInterface(new FakePg(seededRows())); + const lines = parseLines( + await streamToString(db.export({ filter: { kinds: ["abode"] } })), + ); + const kinds = new Set( + lines.filter((l) => l.kind !== "meta").map((l) => l.kind), + ); + assert.deepEqual(kinds, new Set(["abode"])); + }); + + it("applies the abodes allowlist to abode/resident records", async () => { + const db = new PostgresInterface(new FakePg(seededRows())); + const lines = parseLines( + await streamToString( + db.export({ + filter: { abodes: ["aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaa1"] }, + }), + ), + ); + const abodeAids = lines + .filter((l) => l.kind === "abode") + .map((l) => l.data.aid); + assert.deepEqual(abodeAids, ["aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaa1"]); + }); +}); + +describe("postgres import", () => { + it("dispatches inserts per kind, skips note, and commits", async () => { + const fake = new FakePg(); + const db = new PostgresInterface(fake); + const source = Readable.from([ + line("meta", { v: 1 }), + line("user", { + uid: "11111111-1111-1111-1111-111111111111", + email: "u1@test.example", + name: "User One", + flags: {}, + created_at: iso, + updated_at: iso, + }), + line("abode", { + aid: "aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaa1", + name: "Abode One", + created_at: iso, + created_by: "11111111-1111-1111-1111-111111111111", + updated_at: iso, + updated_by: null, + }), + line("apikey", { + uid: "11111111-1111-1111-1111-111111111111", + kid: "kkkkkkkk-kkkk-kkkk-kkkk-kkkkkkkkkkk1", + name: "key one", + permissions: {}, + created_at: iso, + expires_at: null, + }), + // note lines are silently skipped on postgres + line("note", { + nid: "nnnnnnnn-nnnn-nnnn-nnnn-nnnnnnnnnnn1", + aid: "aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaa1", + name: "Note", + content: "x", + properties: {}, + created_at: iso, + created_by: null, + updated_at: iso, + updated_by: null, + }), + ]); + + const result = await db.import(source); + assert.deepEqual(result.counts, { user: 1, abode: 1, apikey: 1 }); + assert.ok(fake.committed); + assert.deepEqual(fake.inserts.map((i) => i.table).sort(), [ + "abodes", + "apikeys", + "users", + ]); + assert.ok(!fake.inserts.some((i) => i.table === "notes"), "no note insert"); + // apikey gets a freshly-minted token (never exported) + const apikeyInsert = fake.inserts.find((i) => i.table === "apikeys")!; + assert.ok( + apikeyInsert.vars.some( + (v) => typeof v === "string" && v.startsWith("at_"), + ), + "apikey insert carries a fresh token", + ); + }); + + it("rejects and rolls back when the stream carries an error sentinel", async () => { + const fake = new FakePg(); + const db = new PostgresInterface(fake); + const source = Readable.from([ + line("meta", { v: 1 }), + line("user", { + uid: "11111111-1111-1111-1111-111111111111", + email: "u1@test.example", + name: "User One", + flags: {}, + created_at: iso, + updated_at: iso, + }), + line("error", { message: "boom" }), + ]); + + await assert.rejects(() => db.import(source), /boom/); + assert.ok(fake.rolledBack); + assert.ok(!fake.committed); + }); + + it("rejects and rolls back on signal abort", async () => { + const fake = new FakePg(); + const db = new PostgresInterface(fake); + const ac = new AbortController(); + const source = Readable.from( + (async function* () { + yield line("meta", { v: 1 }); + yield line("user", { + uid: "11111111-1111-1111-1111-111111111111", + email: "u1@test.example", + name: "User One", + flags: {}, + created_at: iso, + updated_at: iso, + }); + ac.abort(); + await new Promise((r) => setTimeout(r, 1000)); + })(), + ); + + await assert.rejects(() => db.import(source, { signal: ac.signal })); + assert.ok(fake.rolledBack); + assert.ok(!fake.committed); + }); +});