diff --git a/package.json b/package.json index 8b8099c4..6967a252 100644 --- a/package.json +++ b/package.json @@ -46,6 +46,7 @@ "test:e2e:cli:p2p": "npm run test:e2e:p2p --workspace self-hosted-livesync-cli", "test:e2e:cli:adaptive-s3": "npm run test:e2e:adaptive-s3 --workspace self-hosted-livesync-cli", "test:e2e:cli:adaptive-webdav": "npm run test:e2e:adaptive-webdav --workspace self-hosted-livesync-cli", + "test:e2e:cli:adaptive-postgrest": "npm run test:e2e:adaptive-postgrest --workspace self-hosted-livesync-cli", "test:e2e:cli:all": "npm run test:e2e:all --workspace self-hosted-livesync-cli", "test:integration": "npx dotenv-cli -e .env -e .test.env -- vitest run --config vitest.config.integration.ts", "test:unit:coverage": "vitest run --config vitest.config.unit.ts --coverage", @@ -96,6 +97,8 @@ "test:docker-s3:start": "npm run test:docker-s3:up && sleep 3 && npm run test:docker-s3:init", "test:docker-s3:down": "npx dotenv-cli -e .env -e .test.env -- ./test/shell/minio-stop.sh", "test:docker-s3:stop": "npm run test:docker-s3:down", + "test:docker-postgrest:start": "npx dotenv-cli -e .env -e .test.env -- deno run --env-file=src/apps/cli/testdeno/.test.env -A --no-check src/apps/cli/testdeno/manage-postgrest-fixture.ts start", + "test:docker-postgrest:stop": "npx dotenv-cli -e .env -e .test.env -- deno run --env-file=src/apps/cli/testdeno/.test.env -A --no-check src/apps/cli/testdeno/manage-postgrest-fixture.ts stop", "test:docker-all:up": "npm run test:docker-couchdb:up ; npm run test:docker-s3:up", "test:docker-all:init": "npm run test:docker-couchdb:init ; npm run test:docker-s3:init", "test:docker-all:down": "npm run test:docker-couchdb:down ; npm run test:docker-s3:down", diff --git a/src/apps/cli/package.json b/src/apps/cli/package.json index ac6255a9..cb959812 100644 --- a/src/apps/cli/package.json +++ b/src/apps/cli/package.json @@ -26,6 +26,8 @@ "test:e2e:adaptive-s3": "deno task --cwd testdeno test:adaptive-journal-s3", "pretest:e2e:adaptive-webdav": "npm run build", "test:e2e:adaptive-webdav": "deno task --cwd testdeno test:adaptive-journal-webdav", + "pretest:e2e:adaptive-postgrest": "npm run build", + "test:e2e:adaptive-postgrest": "deno task --cwd testdeno test:adaptive-journal-postgrest", "test:e2e:mirror": "bash test/test-mirror-linux.sh", "test:e2e:remote-commands": "bash test/test-remote-commands-linux.sh", "pretest:e2e:all": "npm run build", diff --git a/src/apps/cli/testdeno/deno.json b/src/apps/cli/testdeno/deno.json index def59c9e..8a8e55f8 100644 --- a/src/apps/cli/testdeno/deno.json +++ b/src/apps/cli/testdeno/deno.json @@ -36,7 +36,8 @@ "test:e2e-matrix:minio-enc0": "deno test --env-file=.test.env -A --no-check --filter='e2e matrix: MINIO-enc0' test-e2e-two-vaults-matrix.ts", "test:e2e-matrix:minio-enc1": "deno test --env-file=.test.env -A --no-check --filter='e2e matrix: MINIO-enc1' test-e2e-two-vaults-matrix.ts", "test:adaptive-journal-s3": "deno test --env-file=.test.env -A --no-check test-adaptive-journal-s3.ts", - "test:adaptive-journal-webdav": "deno test --env-file=.test.env -A --no-check test-adaptive-journal-webdav.ts" + "test:adaptive-journal-webdav": "deno test --env-file=.test.env -A --no-check test-adaptive-journal-webdav.ts", + "test:adaptive-journal-postgrest": "deno test --env-file=.test.env -A --no-check test-adaptive-journal-postgrest.ts" }, "imports": { "@std/assert": "jsr:@std/assert@^1.0.13", diff --git a/src/apps/cli/testdeno/fixtures/postgrest/00_roles.sql b/src/apps/cli/testdeno/fixtures/postgrest/00_roles.sql new file mode 100644 index 00000000..968cab2f --- /dev/null +++ b/src/apps/cli/testdeno/fixtures/postgrest/00_roles.sql @@ -0,0 +1,12 @@ +do $$ +begin + if not exists (select 1 from pg_roles where rolname = 'livesync_postgrest_anon') then + create role livesync_postgrest_anon nologin; + end if; + if not exists (select 1 from pg_roles where rolname = 'livesync_postgrest_authenticator') then + create role livesync_postgrest_authenticator noinherit login password 'integration-password'; + end if; +end +$$; + +grant livesync_postgrest_anon to livesync_postgrest_authenticator; diff --git a/src/apps/cli/testdeno/fixtures/postgrest/02_vault.sql b/src/apps/cli/testdeno/fixtures/postgrest/02_vault.sql new file mode 100644 index 00000000..83dc46be --- /dev/null +++ b/src/apps/cli/testdeno/fixtures/postgrest/02_vault.sql @@ -0,0 +1,6 @@ +insert into livesync_private.adaptive_vaults (vault_id, credential_digest) +values ( + 'adaptive-cli-vault-01', + livesync_private.sha256(convert_to('adaptive-cli-vault-credential-0000000000001', 'UTF8')) +) +on conflict (vault_id) do update set credential_digest = excluded.credential_digest; diff --git a/src/apps/cli/testdeno/helpers/cli.ts b/src/apps/cli/testdeno/helpers/cli.ts index 85aca27f..7bf963f6 100644 --- a/src/apps/cli/testdeno/helpers/cli.ts +++ b/src/apps/cli/testdeno/helpers/cli.ts @@ -40,13 +40,15 @@ function concatChunks(chunks: Uint8Array[]): Uint8Array { return out; } +export function redactCliSensitiveText(value: string): string { + return value + .replace(/(obsidian:\/\/setuplivesync\?settings=)[^\s"']+/gu, "$1") + .replace(/(sls\+[^:\s]+:\/\/)[^/?#@\s]*@/gu, "$1@"); +} + export function formatTeeCommand(args: string[]): string { - const redactArgument = (argument: string): string => { - if (argument.startsWith(SETUP_URI_PREFIX)) { - return `${SETUP_URI_PREFIX}`; - } - return argument.replace(/^(sls\+[^:]+:\/\/)[^/?#@]*@/u, "$1@"); - }; + const redactArgument = (argument: string): string => + argument.startsWith(SETUP_URI_PREFIX) ? `${SETUP_URI_PREFIX}` : redactCliSensitiveText(argument); return ["node", CLI_DIST, ...args.map(redactArgument)].map((part) => JSON.stringify(part)).join(" "); } @@ -64,7 +66,7 @@ export function createLineTeeWriter( writer(enc.encode(`[CLI tee pid=${pid}:${streamName}]\n`)); headerWritten = true; } - writer(enc.encode(`[CLI tee pid=${pid}:${streamName}] ${line}\n`)); + writer(enc.encode(`[CLI tee pid=${pid}:${streamName}] ${redactCliSensitiveText(line)}\n`)); }; const flush = (final = false) => { diff --git a/src/apps/cli/testdeno/helpers/docker.ts b/src/apps/cli/testdeno/helpers/docker.ts index 3b50ff0f..23c7de07 100644 --- a/src/apps/cli/testdeno/helpers/docker.ts +++ b/src/apps/cli/testdeno/helpers/docker.ts @@ -6,6 +6,8 @@ * available — including Windows — without needing bash. */ +import { join } from "@std/path"; + type DockerInvoker = { bin: string; prefix: string[]; @@ -15,6 +17,7 @@ type DockerInvoker = { let dockerInvokerPromise: Promise | null = null; const DOCKER_TEE = Deno.env.get("LIVESYNC_DOCKER_TEE") === "1" || Deno.env.get("LIVESYNC_TEST_TEE") === "1"; const trackedContainers = new Set(); +const trackedNetworks = new Set(); const CLEANUP_SIGNALS: Deno.Signal[] = ["SIGINT", "SIGTERM"]; let signalCleanupHandlersInstalled = false; let signalCleanupInProgress = false; @@ -62,17 +65,28 @@ async function collectStream( return out; } -async function runCommand(bin: string, args: string[]): Promise<{ code: number; stdout: string; stderr: string }> { +async function runCommand( + bin: string, + args: string[], + stdinData?: Uint8Array +): Promise<{ code: number; stdout: string; stderr: string }> { try { const child = new Deno.Command(bin, { args, - stdin: "null", + stdin: stdinData ? "piped" : "null", stdout: "piped", stderr: "piped", }).spawn(); const stdoutPromise = collectStream(child.stdout, DOCKER_TEE ? (chunk) => Deno.stdout.writeSync(chunk) : null); const stderrPromise = collectStream(child.stderr, DOCKER_TEE ? (chunk) => Deno.stderr.writeSync(chunk) : null); - const [status, stdout, stderr] = await Promise.all([child.status, stdoutPromise, stderrPromise]); + const stdinPromise = stdinData + ? (async () => { + const writer = child.stdin.getWriter(); + await writer.write(stdinData); + await writer.close(); + })() + : Promise.resolve(); + const [status, stdout, stderr] = await Promise.all([child.status, stdoutPromise, stderrPromise, stdinPromise]); const dec = new TextDecoder(); const result = { code: status.code, @@ -163,6 +177,20 @@ async function getDockerInvoker(): Promise { } async function docker(...args: string[]): Promise<{ code: number; stdout: string; stderr: string }> { + return await dockerCommand(undefined, ...args); +} + +async function dockerWithInput( + input: Uint8Array, + ...args: string[] +): Promise<{ code: number; stdout: string; stderr: string }> { + return await dockerCommand(input, ...args); +} + +async function dockerCommand( + input: Uint8Array | undefined, + ...args: string[] +): Promise<{ code: number; stdout: string; stderr: string }> { const invoker = await getDockerInvoker(); // Either: @@ -176,7 +204,7 @@ async function docker(...args: string[]): Promise<{ code: number; stdout: string : args : [...invoker.prefix, ...args]; - const r = await runCommand(invoker.bin, finalArgs); + const r = await runCommand(invoker.bin, finalArgs, input); return { code: r.code, stdout: r.stdout, stderr: r.stderr }; } @@ -195,12 +223,21 @@ async function stopAndRemoveContainer(container: string): Promise { async function cleanupTrackedContainers(reason: string): Promise { const names = [...trackedContainers]; - if (names.length === 0) return; + if (names.length > 0) { + console.warn(`[WARN] cleaning up tracked containers on ${reason}: ${names.join(", ")}`); + for (const container of names.reverse()) { + await stopAndRemoveContainer(container); + trackedContainers.delete(container); + } + } - console.warn(`[WARN] cleaning up tracked containers on ${reason}: ${names.join(", ")}`); - for (const container of names.reverse()) { - await stopAndRemoveContainer(container); - trackedContainers.delete(container); + const networks = [...trackedNetworks]; + if (networks.length > 0) { + console.warn(`[WARN] cleaning up tracked networks on ${reason}: ${networks.join(", ")}`); + for (const network of networks.reverse()) { + await docker("network", "rm", network).catch(() => {}); + trackedNetworks.delete(network); + } } } @@ -250,7 +287,19 @@ function trackContainer(container: string): void { function untrackContainer(container: string): void { trackedContainers.delete(container); - if (trackedContainers.size === 0) { + if (trackedContainers.size === 0 && trackedNetworks.size === 0) { + removeSignalCleanupHandlers(); + } +} + +function trackNetwork(network: string): void { + ensureSignalCleanupHandlers(); + trackedNetworks.add(network); +} + +function untrackNetwork(network: string): void { + trackedNetworks.delete(network); + if (trackedContainers.size === 0 && trackedNetworks.size === 0) { removeSignalCleanupHandlers(); } } @@ -332,6 +381,14 @@ const MINIO_MC_IMAGE = "minio/mc:RELEASE.2025-04-16T18-13-26Z"; const WEBDAV_CONTAINER = "webdav-test"; const WEBDAV_IMAGE = "httpd:2.4.68"; +const POSTGRES_CONTAINER = "postgrest-postgres-test"; +const POSTGRES_IMAGE = "postgres:15-alpine"; +const POSTGREST_CONTAINER = "postgrest-test"; +const POSTGREST_IMAGE = "postgrest/postgrest:v14.16"; +const POSTGREST_NETWORK = "postgrest-test"; +const POSTGREST_DATABASE_PASSWORD = "integration-password"; +const POSTGREST_VAULT_ID = "adaptive-cli-vault-01"; +const POSTGREST_VAULT_CREDENTIAL = "adaptive-cli-vault-credential-0000000000001"; const WEBDAV_HTTPD_CONFIG = `ServerRoot "/usr/local/apache2" Listen 80 @@ -765,6 +822,220 @@ export async function readWebDAVObjectText(collectionEndpoint: string, key: stri return await response.text(); } +// --------------------------------------------------------------------------- +// PostgREST +// --------------------------------------------------------------------------- + +export interface PostgRESTFixture { + endpoint: string; + vaultCredential: string; + vaultId: string; +} + +export interface PostgRESTAdaptiveRowCounts { + chunks: number; + commits: number; + manifests: number; + writers: number; +} + +async function waitForPostgres(): Promise { + let consecutiveReadyChecks = 0; + for (let attempt = 0; attempt < 60; attempt += 1) { + const result = await docker("exec", POSTGRES_CONTAINER, "pg_isready", "-U", "postgres"); + if (result.code === 0) { + consecutiveReadyChecks += 1; + if (consecutiveReadyChecks >= 3) return; + } else { + consecutiveReadyChecks = 0; + } + await sleep(500); + } + throw new Error("PostgreSQL did not become ready in time"); +} + +async function applyPostgresSql(sql: string): Promise { + const result = await dockerWithInput( + new TextEncoder().encode(sql), + "exec", + "-i", + POSTGRES_CONTAINER, + "psql", + "-v", + "ON_ERROR_STOP=1", + "-U", + "postgres", + "-d", + "postgres" + ); + if (result.code !== 0) { + throw new Error(`Could not initialise PostgreSQL for PostgREST: ${result.stderr.trim()}`); + } +} + +async function waitForPostgREST(fixture: PostgRESTFixture): Promise { + const capabilityUrl = `${fixture.endpoint.replace(/\/+$/u, "")}/rpc/livesync_adaptive_capabilities`; + for (let attempt = 0; attempt < 30; attempt += 1) { + try { + const response = await fetch(capabilityUrl, { + method: "POST", + headers: { + Accept: "application/json", + "Accept-Profile": "livesync_api", + "Content-Profile": "livesync_api", + "Content-Type": "application/json", + "X-LiveSync-Vault-Credential": fixture.vaultCredential, + "X-LiveSync-Vault-ID": fixture.vaultId, + }, + body: "{}", + signal: AbortSignal.timeout(3000), + }); + await response.body?.cancel().catch(() => {}); + if (response.ok) return; + } catch { + // PostgREST or its schema cache is still starting. + } + await sleep(500); + } + throw new Error(`PostgREST did not become ready: ${fixture.endpoint}`); +} + +export async function stopPostgREST(): Promise { + await stopAndRemoveContainer(POSTGREST_CONTAINER); + untrackContainer(POSTGREST_CONTAINER); + await stopAndRemoveContainer(POSTGRES_CONTAINER); + untrackContainer(POSTGRES_CONTAINER); + await docker("network", "rm", POSTGREST_NETWORK).catch(() => {}); + untrackNetwork(POSTGREST_NETWORK); +} + +export async function startPostgREST(endpoint: string): Promise { + const endpointUrl = new URL(endpoint); + if ( + endpointUrl.protocol !== "http:" || + (endpointUrl.hostname !== "127.0.0.1" && endpointUrl.hostname !== "localhost") + ) { + throw new Error(`Managed PostgREST requires a local HTTP endpoint, received: ${endpoint}`); + } + if (endpointUrl.pathname !== "/" || endpointUrl.search || endpointUrl.hash) { + throw new Error(`Managed PostgREST requires a root endpoint without query parameters, received: ${endpoint}`); + } + + const fixture: PostgRESTFixture = { + endpoint: endpoint.replace(/\/+$/u, ""), + vaultCredential: POSTGREST_VAULT_CREDENTIAL, + vaultId: POSTGREST_VAULT_ID, + }; + const fixtureDirectory = join(import.meta.dirname!, "..", "fixtures", "postgrest"); + const commonlibSqlPath = join( + import.meta.dirname!, + "..", + "..", + "..", + "..", + "..", + "node_modules", + "@vrtmrz", + "livesync-commonlib", + "sql", + "postgrest", + "adaptive_journal_v1.sql" + ); + + console.log("[INFO] stopping leftover PostgREST test services if present"); + await stopPostgREST().catch(() => {}); + + try { + await dockerOrFail("network", "create", POSTGREST_NETWORK); + trackNetwork(POSTGREST_NETWORK); + + console.log("[INFO] starting PostgreSQL test container"); + await dockerOrFail( + "run", + "-d", + "--name", + POSTGRES_CONTAINER, + "--network", + POSTGREST_NETWORK, + "--network-alias", + "postgres", + "-e", + `POSTGRES_PASSWORD=${POSTGREST_DATABASE_PASSWORD}`, + POSTGRES_IMAGE + ); + trackContainer(POSTGRES_CONTAINER); + await waitForPostgres(); + + for (const sqlPath of [ + join(fixtureDirectory, "00_roles.sql"), + commonlibSqlPath, + join(fixtureDirectory, "02_vault.sql"), + ]) { + await applyPostgresSql(await Deno.readTextFile(sqlPath)); + } + + console.log("[INFO] starting PostgREST test container"); + await dockerOrFail( + "run", + "-d", + "--name", + POSTGREST_CONTAINER, + "--network", + POSTGREST_NETWORK, + "-p", + `${endpointUrl.port || "80"}:3000`, + "-e", + "PGRST_DB_ANON_ROLE=livesync_postgrest_anon", + "-e", + "PGRST_DB_SCHEMAS=livesync_api", + "-e", + `PGRST_DB_URI=postgres://livesync_postgrest_authenticator:${POSTGREST_DATABASE_PASSWORD}@postgres:5432/postgres`, + "-e", + "PGRST_SERVER_PORT=3000", + POSTGREST_IMAGE + ); + trackContainer(POSTGREST_CONTAINER); + await waitForPostgREST(fixture); + return fixture; + } catch (error) { + await stopPostgREST().catch(() => {}); + throw error; + } +} + +export async function readPostgRESTAdaptiveRowCounts(vaultId: string): Promise { + if (!/^[A-Za-z0-9_-]{16,128}$/u.test(vaultId)) { + throw new Error(`Invalid PostgREST test Vault ID: ${vaultId}`); + } + const query = `select + (select count(*) from livesync_private.adaptive_v1_manifests where vault_id = '${vaultId}'), + (select count(*) from livesync_private.adaptive_v1_chunks where vault_id = '${vaultId}'), + (select count(*) from livesync_private.adaptive_v1_writers where vault_id = '${vaultId}'), + (select count(*) from livesync_private.adaptive_v1_commits where vault_id = '${vaultId}')`; + const output = await dockerOrFail( + "exec", + POSTGRES_CONTAINER, + "psql", + "-At", + "-F", + ",", + "-U", + "postgres", + "-d", + "postgres", + "-c", + query + ); + const [manifests, chunks, writers, commits] = output + .trim() + .split(",") + .map((value) => Number(value)); + if (![manifests, chunks, writers, commits].every(Number.isSafeInteger)) { + throw new Error(`Could not parse PostgREST row counts: ${output}`); + } + return { chunks, commits, manifests, writers }; +} + // --------------------------------------------------------------------------- // P2P relay (strfry) // --------------------------------------------------------------------------- diff --git a/src/apps/cli/testdeno/helpers/settings.ts b/src/apps/cli/testdeno/helpers/settings.ts index 76481aae..9d833f17 100644 --- a/src/apps/cli/testdeno/helpers/settings.ts +++ b/src/apps/cli/testdeno/helpers/settings.ts @@ -127,7 +127,7 @@ export async function applyCouchdbSettings( export async function applyRemoteSyncSettings( settingsFile: string, options: { - remoteType: "COUCHDB" | "MINIO" | "WEBDAV"; + remoteType: "COUCHDB" | "MINIO" | "POSTGREST" | "WEBDAV"; couchdbUri?: string; couchdbUser?: string; couchdbPassword?: string; @@ -136,6 +136,7 @@ export async function applyRemoteSyncSettings( minioEndpoint?: string; minioAccessKey?: string; minioSecretKey?: string; + postgRESTConnectionURI?: string; webDAVConnectionURI?: string; encrypt?: boolean; passphrase?: string; @@ -162,6 +163,9 @@ export async function applyRemoteSyncSettings( data.secretKey = options.minioSecretKey; data.region = "auto"; data.forcePathStyle = true; + } else if (options.remoteType === "POSTGREST") { + data.remoteType = "POSTGREST"; + data.postgrestActiveConnectionURI = options.postgRESTConnectionURI; } else { data.remoteType = "WEBDAV"; data.webDAVactiveConnectionURI = options.webDAVConnectionURI; diff --git a/src/apps/cli/testdeno/manage-postgrest-fixture.ts b/src/apps/cli/testdeno/manage-postgrest-fixture.ts new file mode 100644 index 00000000..5cb7fb27 --- /dev/null +++ b/src/apps/cli/testdeno/manage-postgrest-fixture.ts @@ -0,0 +1,26 @@ +import { startPostgREST, stopPostgREST } from "./helpers/docker.ts"; + +const action = Deno.args[0]; +const endpoint = ( + Deno.env.get("POSTGREST_ENDPOINT") ?? + Deno.env.get("postgrestEndpoint") ?? + "http://127.0.0.1:3001" +).replace(/\/+$/u, ""); + +try { + if (action === "start") { + await startPostgREST(endpoint); + } else if (action === "stop") { + await stopPostgREST(); + } else { + throw new Error("Usage: manage-postgrest-fixture.ts "); + } +} catch (error) { + if (action === "start") await stopPostgREST().catch(() => undefined); + console.error(error instanceof Error ? error.stack : error); + Deno.exit(1); +} + +// The one-shot command intentionally leaves successful services running after +// 'start'; the matching 'stop' command owns their removal. +Deno.exit(0); diff --git a/src/apps/cli/testdeno/run-ci-suite.ts b/src/apps/cli/testdeno/run-ci-suite.ts index ba3a0562..95f312a5 100644 --- a/src/apps/cli/testdeno/run-ci-suite.ts +++ b/src/apps/cli/testdeno/run-ci-suite.ts @@ -13,6 +13,7 @@ const TASKS = [ "test:e2e-matrix:minio-enc1", "test:adaptive-journal-s3", "test:adaptive-journal-webdav", + "test:adaptive-journal-postgrest", ] as const; for (const [index, task] of TASKS.entries()) { diff --git a/src/apps/cli/testdeno/test-adaptive-journal-postgrest.ts b/src/apps/cli/testdeno/test-adaptive-journal-postgrest.ts new file mode 100644 index 00000000..7a618536 --- /dev/null +++ b/src/apps/cli/testdeno/test-adaptive-journal-postgrest.ts @@ -0,0 +1,203 @@ +import { assert, assertEquals } from "@std/assert"; +import { TempDir } from "./helpers/temp.ts"; +import { + assertFilesEqual, + redactCliSensitiveText, + runCli, + runCliOrFail, + runCliWithInputOrFail, + sanitiseCatStdout, +} from "./helpers/cli.ts"; +import { applyRemoteSyncSettings, generateSetupUriFromSettings, initSettingsFile } from "./helpers/settings.ts"; +import { + type PostgRESTFixture, + readPostgRESTAdaptiveRowCounts, + startPostgREST, + stopPostgREST, +} from "./helpers/docker.ts"; + +const BINARY_TEST_BYTES = 2 * 1024 * 1024; + +function deterministicBytes(length: number, seed: number): Uint8Array { + const bytes = new Uint8Array(length); + let state = seed; + for (let index = 0; index < bytes.byteLength; index += 1) { + state ^= state << 13; + state ^= state >>> 17; + state ^= state << 5; + bytes[index] = state & 0xff; + } + return bytes; +} + +function postgRESTConnectionURI(fixture: PostgRESTFixture): string { + const endpointUrl = new URL(fixture.endpoint); + const proxyUrl = new URL(`https://${endpointUrl.host}${endpointUrl.pathname}`); + proxyUrl.username = fixture.vaultId; + proxyUrl.password = fixture.vaultCredential; + if (endpointUrl.protocol === "http:") proxyUrl.searchParams.set("insecure", "true"); + return `sls+postgrest:${proxyUrl.toString().slice("https:".length)}`; +} + +function remoteIdFromListing(listing: string): string { + const line = listing + .split(/\r?\n/u) + .find((candidate) => candidate.includes("\tPostgREST E2E\t") || candidate.includes("\tPostgREST ")); + const id = line?.split("\t", 1)[0]; + if (!id) throw new Error(`PostgREST remote profile was not listed:\n${listing}`); + return id; +} + +function externalFixture(endpoint: string): PostgRESTFixture { + return { + endpoint, + vaultCredential: + Deno.env.get("POSTGREST_VAULT_CREDENTIAL") ?? + Deno.env.get("postgrestVaultCredential") ?? + "adaptive-cli-vault-credential-0000000000001", + vaultId: Deno.env.get("POSTGREST_VAULT_ID") ?? Deno.env.get("postgrestVaultId") ?? "adaptive-cli-vault-01", + }; +} + +Deno.test("CLI tee output redacts a PostgREST Vault credential", () => { + assertEquals( + redactCliSensitiveText( + "exported sls+postgrest://adaptive-cli-vault-01:adaptive-cli-vault-credential@example.test/rest/v1" + ), + "exported sls+postgrest://@example.test/rest/v1" + ); +}); + +Deno.test("e2e: two CLI vaults synchronise through Adaptive Journal PostgREST", async () => { + const suffix = `${Date.now()}-${Math.floor(Math.random() * 100000)}`; + const endpoint = ( + Deno.env.get("POSTGREST_ENDPOINT") ?? + Deno.env.get("postgrestEndpoint") ?? + "http://127.0.0.1:3001" + ).replace(/\/+$/u, ""); + const vaultPassphrase = "adaptive-journal-postgrest-cli-e2ee"; + const setupPassphrase = "adaptive-journal-postgrest-cli-setup"; + const shouldStartDocker = Deno.env.get("LIVESYNC_START_DOCKER") !== "0"; + const keepDocker = Deno.env.get("LIVESYNC_DEBUG_KEEP_DOCKER") === "1"; + const fixture = shouldStartDocker ? await startPostgREST(endpoint) : externalFixture(endpoint); + const connectionURI = postgRESTConnectionURI(fixture); + + await using workDir = await TempDir.create("livesync-cli-adaptive-journal-postgrest"); + const vaultA = workDir.join("vault-a"); + const vaultB = workDir.join("vault-b"); + const settingsA = workDir.join("settings-a.json"); + const settingsB = workDir.join("settings-b.json"); + const binarySourceA = workDir.join("source-a.bin"); + const binarySourceB = workDir.join("source-b.bin"); + const binaryDestinationA = workDir.join("destination-a.bin"); + const binaryDestinationB = workDir.join("destination-b.bin"); + await Deno.mkdir(vaultA, { recursive: true }); + await Deno.mkdir(vaultB, { recursive: true }); + + try { + await initSettingsFile(settingsA); + await applyRemoteSyncSettings(settingsA, { + remoteType: "POSTGREST", + postgRESTConnectionURI: connectionURI, + encrypt: true, + passphrase: vaultPassphrase, + enableCompression: false, + journalFormat: "adaptive-v1", + packReadPolicy: "whole-pack", + }); + const addedRemote = await runCliOrFail( + vaultA, + "--settings", + settingsA, + "remote-add", + "PostgREST E2E", + connectionURI + ); + const remoteId = addedRemote.trim().split("\t", 1)[0]; + assert(remoteId, `remote-add did not return a profile ID:\n${addedRemote}`); + await runCliOrFail(vaultA, "--settings", settingsA, "remote-activate", remoteId); + + const textPath = "adaptive/text.md"; + const binaryPath = "adaptive/data.bin"; + await runCliWithInputOrFail(`created-by-a-${suffix}\n`, vaultA, "--settings", settingsA, "put", textPath); + await Deno.writeFile(binarySourceA, deterministicBytes(BINARY_TEST_BYTES, 0x1a2b3c4d)); + await runCliOrFail(vaultA, "--settings", settingsA, "push", binarySourceA, binaryPath); + await runCliOrFail(vaultA, "--settings", settingsA, "sync"); + + const remoteListing = await runCliOrFail(vaultA, "--settings", settingsA, "remote-ls"); + assertEquals(remoteIdFromListing(remoteListing), remoteId); + assert( + remoteListing + .split(/\r?\n/u) + .some((line) => line.startsWith(`${remoteId}\t`) && line.includes("\tactive\t")), + `Activated PostgREST profile was not listed as active:\n${remoteListing}` + ); + const exportedConnection = ( + await runCliOrFail(vaultA, "--settings", settingsA, "remote-export", remoteId) + ).trim(); + assert(exportedConnection.startsWith("sls+postgrest://")); + assert(exportedConnection.includes("journalFormat=adaptive-v1")); + assert(!exportedConnection.includes("packReadPolicy=")); + + const setupURI = await generateSetupUriFromSettings(settingsA, setupPassphrase, true, vaultPassphrase); + await initSettingsFile(settingsB); + await runCliWithInputOrFail(`${setupPassphrase}\n`, vaultB, "--settings", settingsB, "setup", setupURI); + const settingsAfterSetup = JSON.parse(await Deno.readTextFile(settingsB)) as { + encryptedPassphrase?: string; + postgrestActiveConnectionURI?: string; + }; + assert( + typeof settingsAfterSetup.encryptedPassphrase === "string" && + settingsAfterSetup.encryptedPassphrase.length > 0, + "setup did not persist the encrypted Vault passphrase" + ); + assertEquals(settingsAfterSetup.postgrestActiveConnectionURI, connectionURI); + + await runCliOrFail(vaultB, "--settings", settingsB, "sync"); + assertEquals( + sanitiseCatStdout(await runCliOrFail(vaultB, "--settings", settingsB, "cat", textPath)).trimEnd(), + `created-by-a-${suffix}` + ); + await runCliOrFail(vaultB, "--settings", settingsB, "pull", binaryPath, binaryDestinationB); + await assertFilesEqual(binarySourceA, binaryDestinationB, "Adaptive Journal PostgREST transfer differs"); + + await runCliWithInputOrFail(`updated-by-b-${suffix}\n`, vaultB, "--settings", settingsB, "put", textPath); + await Deno.writeFile(binarySourceB, deterministicBytes(BINARY_TEST_BYTES, 0x5e6f7788)); + await runCliOrFail(vaultB, "--settings", settingsB, "push", binarySourceB, binaryPath); + await runCliOrFail(vaultB, "--settings", settingsB, "sync"); + await runCliOrFail(vaultA, "--settings", settingsA, "sync"); + assertEquals( + sanitiseCatStdout(await runCliOrFail(vaultA, "--settings", settingsA, "cat", textPath)).trimEnd(), + `updated-by-b-${suffix}` + ); + await runCliOrFail(vaultA, "--settings", settingsA, "pull", binaryPath, binaryDestinationA); + await assertFilesEqual(binarySourceB, binaryDestinationA, "Adaptive Journal PostgREST return transfer differs"); + + await runCliOrFail(vaultA, "--settings", settingsA, "rm", binaryPath); + await runCliOrFail(vaultA, "--settings", settingsA, "sync"); + await runCliOrFail(vaultB, "--settings", settingsB, "sync"); + const deleted = await runCli(vaultB, "--settings", settingsB, "cat", binaryPath); + assert(deleted.code !== 0, `Deleted binary remained readable:\n${deleted.combined}`); + + const statusOutput = await runCliOrFail(vaultA, "--settings", settingsA, "remote-status", remoteId); + const statusJsonStart = statusOutput.indexOf("{"); + assert(statusJsonStart >= 0, `PostgREST remote status did not contain JSON:\n${statusOutput}`); + const status = JSON.parse(statusOutput.slice(statusJsonStart)) as { estimatedSize?: unknown }; + assert( + typeof status.estimatedSize === "number" && status.estimatedSize > 0, + `PostgREST remote status did not report a positive size: ${JSON.stringify(status)}` + ); + + if (shouldStartDocker) { + const counts = await readPostgRESTAdaptiveRowCounts(fixture.vaultId); + assertEquals(counts.manifests, 1); + assert(counts.chunks > 0, `PostgREST did not persist Chunk rows: ${JSON.stringify(counts)}`); + assert(counts.writers >= 2, `PostgREST did not preserve both Writer streams: ${JSON.stringify(counts)}`); + assert(counts.commits >= 2, `PostgREST did not persist Commit Bundles: ${JSON.stringify(counts)}`); + } + } finally { + if (shouldStartDocker && !keepDocker) { + await stopPostgREST().catch(() => {}); + } + } +});