diff --git a/README.md b/README.md index ea368b1..155920e 100644 --- a/README.md +++ b/README.md @@ -123,16 +123,16 @@ Declarative, deterministic, transport-agnostic policy system operating on atprot Policies drive sync intervals, priority ordering, and `shouldReplicate` filtering in the replication manager. P2P policies are auto-generated from mutual offer records with `p2p:` prefixed IDs. -## Admin +## App -- **Dashboard**: Server-rendered HTML at `/xrpc/org.p2pds.admin.dashboard` (auto-refresh) +- **Dashboard**: Server-rendered HTML at `/` (auto-refresh) - **API**: Authenticated XRPC endpoints for overview, per-DID status, network status, policies, sync history - **DID management**: Add/remove DIDs at runtime via `addDid`/`removeDid` endpoints - **Rate limiting**: Per-IP and per-DID limits across HTTP, gossipsub, and libp2p ## Desktop App -Optional Tauri v2 wrapper at `apps/desktop/`. Spawns p2pds as a sidecar process and loads the admin dashboard in a webview. +Optional Tauri v2 wrapper at `apps/desktop/`. Spawns p2pds as a sidecar process and loads the dashboard in a webview. ``` cd apps/desktop diff --git a/apps/desktop/src/main.ts b/apps/desktop/src/main.ts index 35a9fea..d5fc4a0 100644 --- a/apps/desktop/src/main.ts +++ b/apps/desktop/src/main.ts @@ -1,7 +1,6 @@ import { Command, type Child } from "@tauri-apps/plugin-shell"; -const HEALTH_URL = "http://127.0.0.1:3000/xrpc/_health"; -const DASHBOARD_URL = "http://127.0.0.1:3000/xrpc/org.p2pds.admin.dashboard"; +const FALLBACK_PORT = 3000; const POLL_INTERVAL_MS = 500; const STARTUP_TIMEOUT_MS = 30_000; @@ -21,12 +20,34 @@ function showError(message: string): void { if (errorMsg) errorMsg.textContent = message; } -async function waitForServer(): Promise { +function healthUrl(port: number): string { + return `http://127.0.0.1:${port}/xrpc/_health`; +} + +function dashboardUrl(port: number): string { + return `http://127.0.0.1:${port}/`; +} + +/** + * Parse the P2PDS_READY line from sidecar stdout. + * Format: `P2PDS_READY {"port":12345,"url":"http://localhost:12345"}` + */ +function parseReadyLine(line: string): { port: number; url: string } | null { + const prefix = "P2PDS_READY "; + if (!line.startsWith(prefix)) return null; + try { + return JSON.parse(line.slice(prefix.length)); + } catch { + return null; + } +} + +async function waitForServer(port: number): Promise { const deadline = Date.now() + STARTUP_TIMEOUT_MS; while (Date.now() < deadline) { try { - const res = await fetch(HEALTH_URL); + const res = await fetch(healthUrl(port)); if (res.ok) { return; } @@ -42,11 +63,28 @@ async function waitForServer(): Promise { ); } -async function spawnSidecar(): Promise { +/** + * Spawn the sidecar and wait for the P2PDS_READY line to detect the port. + * Returns the detected port, or FALLBACK_PORT if detection fails. + */ +async function spawnSidecar(): Promise<{ child: Child; port: number }> { const command = Command.sidecar("binaries/p2pds"); + let detectedPort: number | null = null; + let resolvePort: ((port: number) => void) | null = null; + const portPromise = new Promise((resolve) => { + resolvePort = resolve; + }); + command.stdout.on("data", (line: string) => { console.log(`[p2pds] ${line}`); + if (detectedPort === null) { + const ready = parseReadyLine(line); + if (ready) { + detectedPort = ready.port; + resolvePort!(ready.port); + } + } }); command.stderr.on("data", (line: string) => { @@ -66,19 +104,32 @@ async function spawnSidecar(): Promise { }); const child = await command.spawn(); - return child; + + // Wait for port detection with timeout, fallback to default + const port = await Promise.race([ + portPromise, + new Promise((resolve) => + setTimeout(() => { + console.warn(`[p2pds] P2PDS_READY not detected, falling back to port ${FALLBACK_PORT}`); + resolve(FALLBACK_PORT); + }, STARTUP_TIMEOUT_MS) + ), + ]); + + return { child, port }; } async function start(): Promise { try { setStatus("Starting p2pds..."); - sidecarProcess = await spawnSidecar(); + const { child, port } = await spawnSidecar(); + sidecarProcess = child; setStatus("Waiting for server..."); - await waitForServer(); + await waitForServer(port); setStatus("Redirecting to dashboard..."); - window.location.href = DASHBOARD_URL; + window.location.href = dashboardUrl(port); } catch (err) { const message = err instanceof Error ? err.message : String(err); showError(message); diff --git a/memory/NEXT-STEPS.md b/memory/NEXT-STEPS.md new file mode 100644 index 0000000..28bd7c0 --- /dev/null +++ b/memory/NEXT-STEPS.md @@ -0,0 +1,5 @@ +# Next Steps + +## Reactive sync for watched accounts + +For tracked DIDs, respond to changes from whichever source fires first: firehose (centralized relay, has full blocks for direct apply) or gossipsub (peer notification, triggers peer-first fetch). Both paths exist independently with dedup. Next: ensure both are always active for all tracked DIDs, unify the signal, add per-DID metrics for which source won. diff --git a/scripts/demo-replication.ts b/scripts/demo-replication.ts index eaa1a79..73744f7 100644 --- a/scripts/demo-replication.ts +++ b/scripts/demo-replication.ts @@ -199,7 +199,7 @@ async function main() { } } - console.log(pc.bold(pc.green(`\n✓ Dashboard ready at: http://127.0.0.1:${PORT_B}/xrpc/org.p2pds.admin.dashboard`))); + console.log(pc.bold(pc.green(`\n✓ Dashboard ready at: http://127.0.0.1:${PORT_B}/`))); console.log(pc.dim(" Auth token: demo-token")); console.log(pc.dim(" Press Ctrl+C to stop\n")); diff --git a/scripts/real-bidir-test.ts b/scripts/real-bidir-test.ts new file mode 100644 index 0000000..04ad3a0 --- /dev/null +++ b/scripts/real-bidir-test.ts @@ -0,0 +1,485 @@ +/** + * Real OAuth Bidirectional Replication Smoke Test + * + * Starts two p2pds servers with IPFS networking enabled, guides user + * through OAuth in the browser, then automates: self-sync, peer dial, + * cross-sync via libp2p, verify cross-serving, cleanup. + * + * No polling or timeouts — uses property interception for OAuth identity + * events, promise capture for syncDid calls, and deterministic dial(). + * + * Usage: npx tsx scripts/real-bidir-test.ts alice.bsky.social bob.bsky.social [--clean] + * + * Sessions persist in stable data dirs — re-runs skip OAuth and self-sync. + * Use --clean to wipe data dirs and disconnect after the test. + */ + +import { mkdirSync, rmSync, existsSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { execSync } from "node:child_process"; +import { startServer, type ServerHandle } from "../src/start.js"; +import type { Config } from "../src/config.js"; +import type { ReplicationManager } from "../src/replication/replication-manager.js"; +import type { IpfsService } from "../src/ipfs.js"; + +// --------------------------------------------------------------------------- +// Shared state for signal-safe cleanup +// --------------------------------------------------------------------------- + +let serverA: ServerHandle | undefined; +let serverB: ServerHandle | undefined; +let tmpA: string | undefined; +let tmpB: string | undefined; +let cleanDirs = false; + +async function cleanup() { + log("Shutting down servers..."); + if (serverA) { await serverA.close().catch(() => {}); serverA = undefined; } + if (serverB) { await serverB.close().catch(() => {}); serverB = undefined; } + if (cleanDirs) { + if (tmpA) { log("Cleaning up data dirs..."); rmSync(tmpA, { recursive: true, force: true }); tmpA = undefined; } + if (tmpB) { rmSync(tmpB, { recursive: true, force: true }); tmpB = undefined; } + } + + // Brief settle for async gossipsub teardown + await new Promise((r) => setTimeout(r, 500)); + log("Done."); +} + +// Suppress gossipsub StreamStateError (async background streams during peer connect/disconnect) +process.on("uncaughtException", (err) => { + if (err?.constructor?.name === "StreamStateError") return; + console.error("Uncaught:", err); + cleanup().finally(() => process.exit(1)); +}); + +// Handle SIGINT/SIGTERM for clean shutdown when killed +for (const sig of ["SIGINT", "SIGTERM"] as const) { + process.on(sig, () => { + log(`Caught ${sig}, cleaning up...`); + cleanup().finally(() => process.exit(1)); + }); +} + +// --------------------------------------------------------------------------- +// Helpers +// --------------------------------------------------------------------------- + +function log(msg: string) { + console.log(`\x1b[36m[test]\x1b[0m ${msg}`); +} + +function fail(msg: string): never { + console.error(`\x1b[31m[FAIL]\x1b[0m ${msg}`); + cleanup().finally(() => process.exit(1)); + throw new Error("unreachable"); +} + +function ok(msg: string) { + console.log(`\x1b[32m[ OK ]\x1b[0m ${msg}`); +} + +/** + * Intercept config.DID setter so we get an event-driven promise + * that resolves when the OAuth callback establishes identity. + * No polling — the setter fires synchronously inside the callback. + */ +function withIdentityPromise(config: Config): Promise { + let _did = config.DID; + let _resolve: ((did: string) => void) | undefined; + + const promise = new Promise((resolve) => { + if (_did) { resolve(_did); return; } + _resolve = resolve; + }); + + Object.defineProperty(config, "DID", { + get: () => _did, + set: (v: string | undefined) => { + _did = v; + if (v && _resolve) { + _resolve(v); + _resolve = undefined; + } + }, + enumerable: true, + configurable: true, + }); + + return promise; +} + +/** + * Intercept syncDid on a ReplicationManager to capture per-DID promises. + * + * Returns an `awaitSync(did)` function that: + * - Returns the stored promise if syncDid was already called for that DID + * - Otherwise waits (event-driven, no polling) until syncDid is called + */ +function interceptSyncDid(rm: ReplicationManager) { + const captured = new Map>(); + const waiters = new Map) => void>(); + + const origSyncDid = rm.syncDid.bind(rm); + (rm as any).syncDid = (did: string) => { + const p = origSyncDid(did); + captured.set(did, p); + const waiter = waiters.get(did); + if (waiter) { + waiter(p); + waiters.delete(did); + } + return p; + }; + + return { + awaitSync(did: string): Promise { + if (captured.has(did)) return captured.get(did)!; + return new Promise((resolve, reject) => { + waiters.set(did, (p) => p.then(resolve, reject)); + }); + }, + }; +} + +function makeConfig(dataDir: string, port: number): Config { + return { + PDS_HOSTNAME: `localhost:${port}`, + AUTH_TOKEN: "smoke-test-token", + JWT_SECRET: "smoke-jwt-secret", + PASSWORD_HASH: "$2a$10$test", + DATA_DIR: dataDir, + PORT: port, + IPFS_ENABLED: true, + IPFS_NETWORKING: true, + REPLICATE_DIDS: [], + FIREHOSE_URL: "wss://bsky.network/xrpc/com.atproto.sync.subscribeRepos", + FIREHOSE_ENABLED: false, + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, + OAUTH_ENABLED: true, + }; +} + +async function fetchJson(url: string, opts?: RequestInit): Promise { + const res = await fetch(url, opts); + if (!res.ok) { + const text = await res.text(); + throw new Error(`HTTP ${res.status}: ${text}`); + } + return res.json() as Promise; +} + +// --------------------------------------------------------------------------- +// Main +// --------------------------------------------------------------------------- + +const RECORD_CEILING = 10; + +async function main() { + const args = process.argv.slice(2).filter((a) => !a.startsWith("--")); + const flags = new Set(process.argv.slice(2).filter((a) => a.startsWith("--"))); + cleanDirs = flags.has("--clean"); + + if (args.length < 2) { + console.log("Usage: npx tsx scripts/real-bidir-test.ts [--clean]"); + console.log(" --clean Remove data dirs after test (default: keep for fast re-runs)"); + console.log("Example: npx tsx scripts/real-bidir-test.ts alice.bsky.social bob.bsky.social"); + process.exit(1); + } + + const [handleA, handleB] = args; + const portA = 3100; + const portB = 3101; + + // 1. Stable data dirs (reused across runs to keep OAuth sessions) + tmpA = join(tmpdir(), "p2pds-smoke-a"); + tmpB = join(tmpdir(), "p2pds-smoke-b"); + mkdirSync(tmpA, { recursive: true }); + mkdirSync(tmpB, { recursive: true }); + log(`Data dirs: ${tmpA}, ${tmpB}`); + + try { + // 2. Start servers with identity interception + const configA = makeConfig(tmpA, portA); + const configB = makeConfig(tmpB, portB); + const identityA = withIdentityPromise(configA); + const identityB = withIdentityPromise(configB); + + log("Starting Node A on port 3100..."); + serverA = await startServer(configA); + log("Starting Node B on port 3101..."); + serverB = await startServer(configB); + ok("Both servers started"); + + const rmA = serverA.replicationManager; + const rmB = serverB.replicationManager; + const ipfsA = serverA.ipfsService; + const ipfsB = serverB.ipfsService; + if (!rmA || !rmB) fail("ReplicationManager not available on one or both servers"); + if (!ipfsA || !ipfsB) fail("IpfsService not available on one or both servers"); + + // 3. Stop periodic sync — we control sync explicitly in this test + // Set stopped=true to block the initial 5s delayed syncAll() too + (rmA as any).stopped = true; + (rmB as any).stopped = true; + if ((rmA as any).syncTimer) { clearInterval((rmA as any).syncTimer); (rmA as any).syncTimer = null; } + if ((rmB as any).syncTimer) { clearInterval((rmB as any).syncTimer); (rmB as any).syncTimer = null; } + + // Purge leftover cross-DIDs from previous runs + const selfDidA = configA.DID; + const selfDidB = configB.DID; + if (selfDidA) { + for (const did of rmA.getReplicateDids()) { + if (did !== selfDidA) { await rmA.removeDid(did, true); log(`Purged leftover ${did} from Node A`); } + } + } + if (selfDidB) { + for (const did of rmB.getReplicateDids()) { + if (did !== selfDidB) { await rmB.removeDid(did, true); log(`Purged leftover ${did} from Node B`); } + } + } + + // 4. Set up syncDid interception BEFORE OAuth — captures self-sync + cross-sync promises + const captureA = interceptSyncDid(rmA); + const captureB = interceptSyncDid(rmB); + + // 5. Skip blob sync — blocks-only is sufficient for smoke test + const noopBlobs = async () => ({ fetched: 0, skipped: 0, errors: 0, totalBytes: 0 }); + (rmA as any).syncBlobs = noopBlobs; + (rmB as any).syncBlobs = noopBlobs; + + // 5. Check if already authenticated (identity loaded from DB on startup) + const alreadyAuthedA = !!configA.DID; + const alreadyAuthedB = !!configB.DID; + + if (alreadyAuthedA && alreadyAuthedB) { + ok(`Reusing sessions: A=${configA.DID} (@${configA.HANDLE ?? "?"}), B=${configB.DID} (@${configB.HANDLE ?? "?"})`); + } else { + const urlA = `http://localhost:${portA}/oauth/login?handle=${encodeURIComponent(handleA!)}`; + const urlB = `http://localhost:${portB}/oauth/login?handle=${encodeURIComponent(handleB!)}`; + + console.log(); + console.log("\x1b[1m=== Opening OAuth login in browser ===\x1b[0m"); + console.log(); + if (!alreadyAuthedA) console.log(` Node A (${handleA}): ${urlA}`); + if (!alreadyAuthedB) console.log(` Node B (${handleB}): ${urlB}`); + console.log(); + + try { + if (!alreadyAuthedA) execSync(`open ${JSON.stringify(urlA)}`); + if (!alreadyAuthedB) execSync(`open ${JSON.stringify(urlB)}`); + log("Opened login URLs in your browser. Please authenticate."); + } catch { + log("Could not auto-open browser. Please open the URLs above manually."); + } + + log("Waiting for authentication (Ctrl+C to abort)..."); + } + + // 6. Await identity events (resolves immediately if already authed) + const [didA, didB] = await Promise.all([identityA, identityB]); + + ok(`Node A: ${didA} (@${configA.HANDLE ?? "unknown"})`); + ok(`Node B: ${didB} (@${configB.HANDLE ?? "unknown"})`); + + // 7. Self-sync: if fresh OAuth, wait for the auto-triggered sync. + // If session reused but data is missing (cleared), trigger sync explicitly. + const selfStateA = rmA.getSyncStorage().getState(didA); + const selfStateB = rmB.getSyncStorage().getState(didB); + const needSyncA = !selfStateA || selfStateA.status !== "synced"; + const needSyncB = !selfStateB || selfStateB.status !== "synced"; + + if (needSyncA || needSyncB) { + log("Self-sync needed..."); + const waits: Promise[] = []; + if (needSyncA && alreadyAuthedA) { + // Session reused but data cleared — trigger self-sync explicitly + rmA.addDid(didA).catch(() => {}); + } + if (needSyncB && alreadyAuthedB) { + rmB.addDid(didB).catch(() => {}); + } + if (needSyncA) waits.push(captureA.awaitSync(didA)); + if (needSyncB) waits.push(captureB.awaitSync(didB)); + await Promise.all(waits); + ok("Self-sync complete on both nodes"); + } else { + ok("Self-sync data already present — skipping"); + } + + // 8. Dial peers to establish libp2p connections (use local TCP, not relay) + const addrsA = ipfsA.getMultiaddrs(); + const addrsB = ipfsB.getMultiaddrs(); + const pickLocal = (addrs: typeof addrsA) => + addrs.find((a) => { + const s = a.toString(); + return s.includes("/ip4/127.0.0.1/tcp/") && !s.includes("/ws/") && !s.includes("/p2p-circuit/"); + }); + const localA = pickLocal(addrsA); + const localB = pickLocal(addrsB); + if (!localA || !localB) { + fail(`No local TCP addr (A: ${!!localA}, B: ${!!localB})`); + } + log(`Node A local: ${localA}`); + log(`Node B local: ${localB}`); + + await Promise.all([ + ipfsA.dial(localB), + ipfsB.dial(localA), + ]); + ok("Peer connections established via libp2p"); + + // 9. Add cross-DIDs — triggers syncDid (stopped=true only blocks syncAll, not syncDid) + log("Adding cross-DIDs and syncing via libp2p..."); + await rmA.addDid(didB); + ok(`Node A now tracking ${didB}`); + await rmB.addDid(didA); + ok(`Node B now tracking ${didA}`); + + // 10. Wait for cross-sync completion + await Promise.all([ + captureA.awaitSync(didB), + captureB.awaitSync(didA), + ]); + ok("Cross-sync complete"); + + // 11. Verify cross-sync used libp2p (not HTTP PDS fallback) + const historyA = rmA.getSyncStorage().getSyncHistory(didB, 1); + const historyB = rmB.getSyncStorage().getSyncHistory(didA, 1); + + if (historyA.length === 0) fail("No sync history on Node A for cross-sync"); + if (historyB.length === 0) fail("No sync history on Node B for cross-sync"); + + const sourceA = historyA[0]!.sourceType; + const sourceB = historyB[0]!.sourceType; + + if (sourceA !== "libp2p") { + fail(`Node A cross-sync used "${sourceA}" instead of "libp2p"`); + } + ok(`Node A cross-synced ${didB} via libp2p`); + + if (sourceB !== "libp2p") { + fail(`Node B cross-sync used "${sourceB}" instead of "libp2p"`); + } + ok(`Node B cross-synced ${didA} via libp2p`); + + // 12. Verify sync state + const syncStateA = rmA.getSyncStorage().getState(didB); + if (!syncStateA || syncStateA.status !== "synced") { + fail(`Node A sync status for ${didB}: ${syncStateA?.status ?? "missing"}`); + } + ok(`Node A synced ${didB}`); + + const syncStateB = rmB.getSyncStorage().getState(didA); + if (!syncStateB || syncStateB.status !== "synced") { + fail(`Node B sync status for ${didA}: ${syncStateB?.status ?? "missing"}`); + } + ok(`Node B synced ${didA}`); + + // 13. Verify cross-serving via HTTP reads (one-shot, not polling) + log("Verifying cross-serving..."); + + // Node A describes Node B's repo + const descA = await fetchJson<{ did: string; collections: string[] }>( + `http://localhost:${portA}/xrpc/com.atproto.repo.describeRepo?repo=${encodeURIComponent(didB)}`, + ); + if (descA.did !== didB) fail(`describeRepo DID mismatch: ${descA.did} !== ${didB}`); + ok(`Node A describes ${didB}: ${descA.collections.length} collections`); + + // Node B describes Node A's repo + const descB = await fetchJson<{ did: string; collections: string[] }>( + `http://localhost:${portB}/xrpc/com.atproto.repo.describeRepo?repo=${encodeURIComponent(didA)}`, + ); + if (descB.did !== didA) fail(`describeRepo DID mismatch: ${descB.did} !== ${didA}`); + ok(`Node B describes ${didA}: ${descB.collections.length} collections`); + + // Node A lists records from Node B (ceiling of 10 per collection, first 5 collections) + let totalRecsA = 0; + for (const coll of descA.collections.slice(0, 5)) { + const recs = await fetchJson<{ records: Array<{ uri: string }> }>( + `http://localhost:${portA}/xrpc/com.atproto.repo.listRecords?repo=${encodeURIComponent(didB)}&collection=${encodeURIComponent(coll)}&limit=${RECORD_CEILING}`, + ); + totalRecsA += recs.records.length; + if (recs.records.length > 0) { + log(` ${coll}: ${recs.records.length} records`); + } + } + if (totalRecsA === 0) fail("Node A served 0 records for Node B"); + ok(`Node A serves ${totalRecsA} records for ${didB}`); + + // Node B lists records from Node A (ceiling of 10 per collection, first 5 collections) + let totalRecsB = 0; + for (const coll of descB.collections.slice(0, 5)) { + const recs = await fetchJson<{ records: Array<{ uri: string }> }>( + `http://localhost:${portB}/xrpc/com.atproto.repo.listRecords?repo=${encodeURIComponent(didA)}&collection=${encodeURIComponent(coll)}&limit=${RECORD_CEILING}`, + ); + totalRecsB += recs.records.length; + if (recs.records.length > 0) { + log(` ${coll}: ${recs.records.length} records`); + } + } + if (totalRecsB === 0) fail("Node B served 0 records for Node A"); + ok(`Node B serves ${totalRecsB} records for ${didA}`); + + // Node A serves Node B's repo via getRepo + const getRepoA = await fetch( + `http://localhost:${portA}/xrpc/com.atproto.sync.getRepo?did=${encodeURIComponent(didB)}`, + ); + if (getRepoA.status !== 200) fail(`getRepo failed: ${getRepoA.status}`); + const carA = new Uint8Array(await getRepoA.arrayBuffer()); + ok(`Node A serves ${didB} repo CAR (${(carA.length / 1024).toFixed(1)} KB)`); + + // Node B serves Node A's repo via getRepo + const getRepoB = await fetch( + `http://localhost:${portB}/xrpc/com.atproto.sync.getRepo?did=${encodeURIComponent(didA)}`, + ); + if (getRepoB.status !== 200) fail(`getRepo failed: ${getRepoB.status}`); + const carB = new Uint8Array(await getRepoB.arrayBuffer()); + ok(`Node B serves ${didA} repo CAR (${(carB.length / 1024).toFixed(1)} KB)`); + + // 14. Summary + console.log(); + console.log("\x1b[1m=== Summary ===\x1b[0m"); + console.log(` Node A (@${configA.HANDLE ?? "?"}): verified ${totalRecsA} records served for @${configB.HANDLE ?? "?"}`); + console.log(` Node B (@${configB.HANDLE ?? "?"}): verified ${totalRecsB} records served for @${configA.HANDLE ?? "?"}`); + console.log(` Cross-sync transport: libp2p (both directions)`); + console.log(); + + // 15. Purge cross-sync data (keep self-sync for fast re-runs) + log("Purging cross-replicated data..."); + await rmA.removeDid(didB, true); + ok(`Node A purged ${didB}`); + await rmB.removeDid(didA, true); + ok(`Node B purged ${didA}`); + + // 16. Logout only if --clean (otherwise keep sessions for re-runs) + if (cleanDirs) { + log("Disconnecting identities..."); + await fetch(`http://localhost:${portA}/oauth/logout?disconnect=true`, { method: "POST" }); + ok("Node A disconnected"); + await fetch(`http://localhost:${portB}/oauth/logout?disconnect=true`, { method: "POST" }); + ok("Node B disconnected"); + } else { + log("Keeping sessions for re-runs (use --clean to disconnect)"); + } + + console.log(); + ok("All checks passed!"); + + } catch (err) { + const msg = err instanceof Error ? err.stack ?? err.message : String(err); + console.error(`\x1b[31m[FAIL]\x1b[0m ${msg}`); + } finally { + await cleanup(); + } +} + +main(); diff --git a/scripts/smoke-test.sh b/scripts/smoke-test.sh index ca1544a..15650cd 100755 --- a/scripts/smoke-test.sh +++ b/scripts/smoke-test.sh @@ -112,19 +112,19 @@ check "_health" \ "200" \ '"status":"ok"' -check "admin.dashboard" \ - "$BASE_URL/xrpc/org.p2pds.admin.dashboard" \ +check "app.dashboard" \ + "$BASE_URL/" \ "200" \ "P2PDS" -check "admin.getOverview" \ - "$BASE_URL/xrpc/org.p2pds.admin.getOverview" \ +check "app.getOverview" \ + "$BASE_URL/xrpc/org.p2pds.app.getOverview" \ "200" \ '"version"' \ "smoke-test-token" -check "admin.getNetworkStatus" \ - "$BASE_URL/xrpc/org.p2pds.admin.getNetworkStatus" \ +check "app.getNetworkStatus" \ + "$BASE_URL/xrpc/org.p2pds.app.getNetworkStatus" \ "200" \ "" \ "smoke-test-token" diff --git a/src/bidirectional-replication.test.ts b/src/bidirectional-replication.test.ts index 968dd52..7e25bc9 100644 --- a/src/bidirectional-replication.test.ts +++ b/src/bidirectional-replication.test.ts @@ -245,14 +245,14 @@ describe("Bidirectional replication E2E", () => { expect(handleB.replicationManager).toBeDefined(); // Node A adds Node B's DID, Node B adds Node A's DID - const addBToA = await fetch(`${handleA.url}/xrpc/org.p2pds.admin.addDid`, { + const addBToA = await fetch(`${handleA.url}/xrpc/org.p2pds.app.addDid`, { method: "POST", headers: { Authorization: `Bearer ${configA.AUTH_TOKEN}`, "Content-Type": "application/json" }, body: JSON.stringify({ did: DID_NODE_B }), }); expect(addBToA.status).toBe(200); - const addAToB = await fetch(`${handleB.url}/xrpc/org.p2pds.admin.addDid`, { + const addAToB = await fetch(`${handleB.url}/xrpc/org.p2pds.app.addDid`, { method: "POST", headers: { Authorization: `Bearer ${configB.AUTH_TOKEN}`, "Content-Type": "application/json" }, body: JSON.stringify({ did: DID_NODE_A }), @@ -274,7 +274,7 @@ describe("Bidirectional replication E2E", () => { // Verify Node A synced Bob's data const statusA = await fetch( - `${handleA.url}/xrpc/org.p2pds.admin.getDidStatus?did=${DID_NODE_B}`, + `${handleA.url}/xrpc/org.p2pds.app.getDidStatus?did=${DID_NODE_B}`, { headers: { Authorization: `Bearer ${configA.AUTH_TOKEN}` } }, ); const dsA = (await statusA.json()) as { did: string; blockCount: number; syncState: { status: string } }; @@ -284,7 +284,7 @@ describe("Bidirectional replication E2E", () => { // Verify Node B synced Alice's data const statusB = await fetch( - `${handleB.url}/xrpc/org.p2pds.admin.getDidStatus?did=${DID_NODE_A}`, + `${handleB.url}/xrpc/org.p2pds.app.getDidStatus?did=${DID_NODE_A}`, { headers: { Authorization: `Bearer ${configB.AUTH_TOKEN}` } }, ); const dsB = (await statusB.json()) as { did: string; blockCount: number; syncState: { status: string } }; diff --git a/src/config.ts b/src/config.ts index 9bc0148..7986287 100644 --- a/src/config.ts +++ b/src/config.ts @@ -86,7 +86,7 @@ function loadDotEnv(path: string): void { */ export function loadConfig(envPath?: string): Config { // Load .env file if it exists - const dotenvPath = envPath ?? resolve(process.cwd(), ".env"); + const dotenvPath = envPath ?? process.env.DOTENV_PATH ?? resolve(process.cwd(), ".env"); loadDotEnv(dotenvPath); // Validate required variables diff --git a/src/didless-startup.test.ts b/src/didless-startup.test.ts index b25f682..05cb87a 100644 --- a/src/didless-startup.test.ts +++ b/src/didless-startup.test.ts @@ -86,7 +86,7 @@ describe("DID-less server startup", () => { const config = didlessConfig(tmpDir); handle = await startServer(config); - const res = await fetch(`${handle.url}/xrpc/org.p2pds.admin.dashboard`); + const res = await fetch(`${handle.url}/`); expect(res.status).toBe(200); const html = await res.text(); expect(html).toContain("P2PDS"); @@ -126,7 +126,7 @@ describe("DID-less server startup", () => { handle = await startServer(config, { didResolver: mockResolver }); // Add DID via admin API - const addRes = await fetch(`${handle.url}/xrpc/org.p2pds.admin.addDid`, { + const addRes = await fetch(`${handle.url}/xrpc/org.p2pds.app.addDid`, { method: "POST", headers: { Authorization: `Bearer ${config.AUTH_TOKEN}`, @@ -147,7 +147,7 @@ describe("DID-less server startup", () => { await new Promise((r) => setTimeout(r, 2000)); // Verify replication state - const overviewRes = await fetch(`${handle.url}/xrpc/org.p2pds.admin.getOverview`, { + const overviewRes = await fetch(`${handle.url}/xrpc/org.p2pds.app.getOverview`, { headers: { Authorization: `Bearer ${config.AUTH_TOKEN}` }, }); expect(overviewRes.status).toBe(200); diff --git a/src/index.ts b/src/index.ts index ccddff8..e1f6b42 100644 --- a/src/index.ts +++ b/src/index.ts @@ -19,7 +19,7 @@ import type { ReplicatedRepoReader } from "./replication/replicated-repo-reader. import * as sync from "./xrpc/sync.js"; import * as repo from "./xrpc/repo.js"; import * as server from "./xrpc/server.js"; -import * as admin from "./xrpc/admin.js"; +import * as app_routes from "./xrpc/app.js"; import { respondToChallenge } from "./replication/challenge-response/challenge-responder.js"; import { serializeResponse } from "./replication/challenge-response/http-transport.js"; import type { StorageChallenge } from "./replication/challenge-response/types.js"; @@ -159,12 +159,12 @@ export function createApp( }), ); - // Admin endpoints - const adminRL = rateLimitMiddleware(rateLimiter, { - pool: "admin", + // App endpoints + const appRL = rateLimitMiddleware(rateLimiter, { + pool: "app", rule: { maxRequests: 300, windowMs: w }, }); - app.use("/xrpc/org.p2pds.admin.*", adminRL); + app.use("/xrpc/org.p2pds.app.*", appRL); } // Body size limits (always active, independent of rate limiting) @@ -258,49 +258,10 @@ export function createApp( } }); - // Homepage - app.get("/", (c) => { - const handleHtml = config.HANDLE - ? `` - : config.DID - ? `
${config.DID}
` - : ""; - const html = ` - - - - -P2PDS - - - -
P2PDS
-
a personal data server for the atmosphere
-${handleHtml} -
v${VERSION}
- -`; - return c.html(html); - }); + // Dashboard UI at root + app.get("/", (c) => + app_routes.getDashboard(c, networkService, replicationManager), + ); // ============================================ // Sync endpoints (federation) @@ -689,31 +650,28 @@ ${handleHtml} }); // ============================================ - // Admin monitoring + // App monitoring // ============================================ - app.get("/xrpc/org.p2pds.admin.getOverview", requireAuth, (c) => - admin.getOverview(c, configDid, networkService, replicationManager), - ); - app.get("/xrpc/org.p2pds.admin.getDidStatus", requireAuth, (c) => - admin.getDidStatus(c, replicationManager), + app.get("/xrpc/org.p2pds.app.getOverview", requireAuth, (c) => + app_routes.getOverview(c, configDid, networkService, replicationManager), ); - app.get("/xrpc/org.p2pds.admin.getNetworkStatus", requireAuth, (c) => - admin.getNetworkStatus(c, networkService), + app.get("/xrpc/org.p2pds.app.getDidStatus", requireAuth, (c) => + app_routes.getDidStatus(c, replicationManager), ); - app.get("/xrpc/org.p2pds.admin.getPolicies", requireAuth, (c) => - admin.getPolicies(c, replicationManager), + app.get("/xrpc/org.p2pds.app.getNetworkStatus", requireAuth, (c) => + app_routes.getNetworkStatus(c, networkService), ); - app.get("/xrpc/org.p2pds.admin.getSyncHistory", requireAuth, (c) => - admin.getSyncHistory(c, replicationManager), + app.get("/xrpc/org.p2pds.app.getPolicies", requireAuth, (c) => + app_routes.getPolicies(c, replicationManager), ); - app.post("/xrpc/org.p2pds.admin.addDid", requireAuth, (c) => - admin.addDid(c, configDid, replicationManager), + app.get("/xrpc/org.p2pds.app.getSyncHistory", requireAuth, (c) => + app_routes.getSyncHistory(c, replicationManager), ); - app.post("/xrpc/org.p2pds.admin.removeDid", requireAuth, (c) => - admin.removeDid(c, replicationManager), + app.post("/xrpc/org.p2pds.app.addDid", requireAuth, (c) => + app_routes.addDid(c, configDid, replicationManager), ); - app.get("/xrpc/org.p2pds.admin.dashboard", (c) => - admin.getDashboard(c, networkService, replicationManager), + app.post("/xrpc/org.p2pds.app.removeDid", requireAuth, (c) => + app_routes.removeDid(c, replicationManager), ); // ============================================ diff --git a/src/oauth/routes.ts b/src/oauth/routes.ts index d10a627..bb7af38 100644 --- a/src/oauth/routes.ts +++ b/src/oauth/routes.ts @@ -137,7 +137,7 @@ export function registerOAuthRoutes( }); } - return c.redirect("/xrpc/org.p2pds.admin.dashboard"); + return c.redirect("/"); } catch (err) { const message = err instanceof Error ? err.message : String(err); return c.html(errorPage("Authentication Failed", message), 500); @@ -241,7 +241,7 @@ export function registerOAuthRoutes( config.HANDLE = undefined; } - return c.redirect("/xrpc/org.p2pds.admin.dashboard"); + return c.redirect("/"); } catch (err) { const message = err instanceof Error ? err.message : String(err); return c.json({ error: "LogoutFailed", message }, 500); @@ -293,7 +293,7 @@ a:hover { background: #000; color: #fff; }
Connected
${escapeHtml(did)}
- Back to Dashboard + Back to Dashboard
`; @@ -325,7 +325,7 @@ a:hover { background: #000; color: #fff; }
${escapeHtml(title)}
${escapeHtml(message)}
- Back to Dashboard + Back to Dashboard
`; diff --git a/src/replication/replication-manager.ts b/src/replication/replication-manager.ts index 059ecec..5d72d49 100644 --- a/src/replication/replication-manager.ts +++ b/src/replication/replication-manager.ts @@ -1481,7 +1481,11 @@ export class ReplicationManager { clearInterval(this.notificationCleanupTimer); this.notificationCleanupTimer = null; } - this.networkService.unsubscribeIdentityTopics(this.getReplicateDids()); + try { + this.networkService.unsubscribeIdentityTopics(this.getReplicateDids()); + } catch { + // Gossipsub streams may already be closed during shutdown + } } /** diff --git a/src/server-startup.test.ts b/src/server-startup.test.ts index 15ca2d2..24cc3ad 100644 --- a/src/server-startup.test.ts +++ b/src/server-startup.test.ts @@ -87,7 +87,7 @@ describe("server startup integration", () => { const config = testConfig(tmpDir); handle = await startServer(config); - const res = await fetch(`${handle.url}/xrpc/org.p2pds.admin.dashboard`); + const res = await fetch(`${handle.url}/`); expect(res.status).toBe(200); const html = await res.text(); expect(html).toContain("P2PDS"); @@ -98,7 +98,7 @@ describe("server startup integration", () => { const config = testConfig(tmpDir); handle = await startServer(config); - const res = await fetch(`${handle.url}/xrpc/org.p2pds.admin.getOverview`, { + const res = await fetch(`${handle.url}/xrpc/org.p2pds.app.getOverview`, { headers: { Authorization: `Bearer ${config.AUTH_TOKEN}` }, }); expect(res.status).toBe(200); @@ -121,7 +121,7 @@ describe("server startup integration", () => { handle = await startServer(config, { didResolver: mockResolver }); // Add the DID via admin API - const addRes = await fetch(`${handle.url}/xrpc/org.p2pds.admin.addDid`, { + const addRes = await fetch(`${handle.url}/xrpc/org.p2pds.app.addDid`, { method: "POST", headers: { Authorization: `Bearer ${config.AUTH_TOKEN}`, @@ -142,7 +142,7 @@ describe("server startup integration", () => { await new Promise((r) => setTimeout(r, 2000)); // Check overview — the DID should appear in replication state - const overviewRes = await fetch(`${handle.url}/xrpc/org.p2pds.admin.getOverview`, { + const overviewRes = await fetch(`${handle.url}/xrpc/org.p2pds.app.getOverview`, { headers: { Authorization: `Bearer ${config.AUTH_TOKEN}` }, }); expect(overviewRes.status).toBe(200); diff --git a/src/sidecar-process.test.ts b/src/sidecar-process.test.ts new file mode 100644 index 0000000..f4b3593 --- /dev/null +++ b/src/sidecar-process.test.ts @@ -0,0 +1,166 @@ +/** + * Sidecar process integration test. + * + * Spawns `src/server.ts` as a child process (the same way Tauri does), + * validates the full sidecar contract: + * 1. Parse P2PDS_READY from stdout to get the port + * 2. GET /xrpc/_health returns 200 {"status":"ok"} + * 3. GET / returns 200 HTML + * 4. SIGTERM → process exits with code 0 + */ + +import { describe, it, expect, afterEach } from "vitest"; +import { spawn, type ChildProcess } from "node:child_process"; +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; + +const STARTUP_TIMEOUT_MS = 30_000; +const POLL_INTERVAL_MS = 200; + +interface ReadyInfo { + port: number; + url: string; +} + +function parseReadyLine(line: string): ReadyInfo | null { + const prefix = "P2PDS_READY "; + if (!line.startsWith(prefix)) return null; + try { + return JSON.parse(line.slice(prefix.length)) as ReadyInfo; + } catch { + return null; + } +} + +function waitForReady(proc: ChildProcess, timeoutMs: number): Promise { + return new Promise((resolve, reject) => { + const timer = setTimeout(() => { + reject(new Error(`P2PDS_READY not seen within ${timeoutMs}ms`)); + }, timeoutMs); + + let buffer = ""; + proc.stdout?.on("data", (chunk: Buffer) => { + buffer += chunk.toString(); + const lines = buffer.split("\n"); + buffer = lines.pop()!; // keep incomplete trailing line + for (const line of lines) { + const info = parseReadyLine(line); + if (info) { + clearTimeout(timer); + resolve(info); + return; + } + } + }); + + proc.on("error", (err) => { + clearTimeout(timer); + reject(err); + }); + + proc.on("exit", (code) => { + clearTimeout(timer); + reject(new Error(`Process exited with code ${code} before P2PDS_READY`)); + }); + }); +} + +async function pollHealth(url: string, timeoutMs: number): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + try { + const res = await fetch(`${url}/xrpc/_health`); + if (res.ok) return; + } catch { + // not ready yet + } + await new Promise((r) => setTimeout(r, POLL_INTERVAL_MS)); + } + throw new Error(`Health check did not pass within ${timeoutMs}ms`); +} + +describe("sidecar process", () => { + let proc: ChildProcess | undefined; + let tmpDir: string | undefined; + + afterEach(async () => { + if (proc && !proc.killed) { + proc.kill("SIGTERM"); + // Wait for exit (up to 5s) + await new Promise((resolve) => { + const timer = setTimeout(() => { + proc?.kill("SIGKILL"); + resolve(); + }, 5_000); + proc!.on("exit", () => { + clearTimeout(timer); + resolve(); + }); + }); + } + proc = undefined; + if (tmpDir) { + rmSync(tmpDir, { recursive: true, force: true }); + tmpDir = undefined; + } + }); + + it("full sidecar lifecycle: ready → health → dashboard → shutdown", async () => { + tmpDir = mkdtempSync(join(tmpdir(), "sidecar-test-")); + + proc = spawn("npx", ["tsx", "src/server.ts"], { + env: { + ...process.env, + PORT: "0", + DATA_DIR: tmpDir, + IPFS_ENABLED: "true", + IPFS_NETWORKING: "false", + FIREHOSE_ENABLED: "false", + RATE_LIMIT_ENABLED: "false", + OAUTH_ENABLED: "false", + PDS_HOSTNAME: "localhost", + AUTH_TOKEN: "test-token", + SIGNING_KEY: "0000000000000000000000000000000000000000000000000000000000000001", + SIGNING_KEY_PUBLIC: "zQ3shP2mWsZYWgvZM9GJ3EvMfRXQJwuTh6BdXLvJB9gFhT3Lr", + JWT_SECRET: "test-secret", + PASSWORD_HASH: "$2a$10$test", + DID: "did:plc:sidecartest", + }, + stdio: ["ignore", "pipe", "pipe"], + }); + + // Collect stderr for debugging + let stderr = ""; + proc.stderr?.on("data", (chunk: Buffer) => { + stderr += chunk.toString(); + }); + + // 1. Parse P2PDS_READY from stdout + const ready = await waitForReady(proc, STARTUP_TIMEOUT_MS); + expect(ready.port).toBeGreaterThan(0); + expect(ready.url).toMatch(/^http:\/\/localhost:\d+$/); + + // 2. Poll health endpoint + await pollHealth(ready.url, 5_000); + const healthRes = await fetch(`${ready.url}/xrpc/_health`); + expect(healthRes.status).toBe(200); + const healthBody = (await healthRes.json()) as { status: string }; + expect(healthBody.status).toBe("ok"); + + // 3. Dashboard returns HTML + const dashRes = await fetch(`${ready.url}/`); + expect(dashRes.status).toBe(200); + const html = await dashRes.text(); + expect(html).toContain("P2PDS"); + + // 4. Graceful shutdown via SIGTERM + const exitPromise = new Promise((resolve) => { + proc!.on("exit", (code) => resolve(code)); + }); + proc.kill("SIGTERM"); + const exitCode = await exitPromise; + expect(exitCode).toBe(0); + proc = undefined; // prevent double-kill in afterEach + }, STARTUP_TIMEOUT_MS + 10_000); +}); diff --git a/src/two-node-didless.test.ts b/src/two-node-didless.test.ts index d2fda41..9addca6 100644 --- a/src/two-node-didless.test.ts +++ b/src/two-node-didless.test.ts @@ -124,7 +124,7 @@ describe("Two-node DID-less replication", () => { dbB.close(); // Verify identity shows up in overview - const overviewA = await fetch(`${handleA.url}/xrpc/org.p2pds.admin.getOverview`, { + const overviewA = await fetch(`${handleA.url}/xrpc/org.p2pds.app.getOverview`, { headers: { Authorization: `Bearer ${configA.AUTH_TOKEN}` }, }); expect(overviewA.status).toBe(200); @@ -132,7 +132,7 @@ describe("Two-node DID-less replication", () => { expect(ovA.did).toBe("did:plc:node-a-identity"); // Node A replicates Alice's account - const addAlice = await fetch(`${handleA.url}/xrpc/org.p2pds.admin.addDid`, { + const addAlice = await fetch(`${handleA.url}/xrpc/org.p2pds.app.addDid`, { method: "POST", headers: { Authorization: `Bearer ${configA.AUTH_TOKEN}`, @@ -143,7 +143,7 @@ describe("Two-node DID-less replication", () => { expect(addAlice.status).toBe(200); // Node B replicates Bob's account - const addBob = await fetch(`${handleB.url}/xrpc/org.p2pds.admin.addDid`, { + const addBob = await fetch(`${handleB.url}/xrpc/org.p2pds.app.addDid`, { method: "POST", headers: { Authorization: `Bearer ${configB.AUTH_TOKEN}`, @@ -168,7 +168,7 @@ describe("Two-node DID-less replication", () => { // Verify Node A synced Alice's data const statusA = await fetch( - `${handleA.url}/xrpc/org.p2pds.admin.getDidStatus?did=${DID_ALICE}`, + `${handleA.url}/xrpc/org.p2pds.app.getDidStatus?did=${DID_ALICE}`, { headers: { Authorization: `Bearer ${configA.AUTH_TOKEN}` } }, ); expect(statusA.status).toBe(200); @@ -183,7 +183,7 @@ describe("Two-node DID-less replication", () => { // Verify Node B synced Bob's data const statusB = await fetch( - `${handleB.url}/xrpc/org.p2pds.admin.getDidStatus?did=${DID_BOB}`, + `${handleB.url}/xrpc/org.p2pds.app.getDidStatus?did=${DID_BOB}`, { headers: { Authorization: `Bearer ${configB.AUTH_TOKEN}` } }, ); expect(statusB.status).toBe(200); @@ -222,7 +222,7 @@ describe("Two-node DID-less replication", () => { db1.close(); // Add and sync - await fetch(`${handleA.url}/xrpc/org.p2pds.admin.addDid`, { + await fetch(`${handleA.url}/xrpc/org.p2pds.app.addDid`, { method: "POST", headers: { Authorization: `Bearer ${config1.AUTH_TOKEN}`, @@ -238,7 +238,7 @@ describe("Two-node DID-less replication", () => { // Verify sync worked const status1 = await fetch( - `${handleA.url}/xrpc/org.p2pds.admin.getDidStatus?did=${DID_ALICE}`, + `${handleA.url}/xrpc/org.p2pds.app.getDidStatus?did=${DID_ALICE}`, { headers: { Authorization: `Bearer ${config1.AUTH_TOKEN}` } }, ); const ds1 = (await status1.json()) as { blockCount: number }; @@ -258,7 +258,7 @@ describe("Two-node DID-less replication", () => { expect(config2.HANDLE).toBe("persistent.test"); // Replication state should persist — DID_ALICE is still tracked - const overview = await fetch(`${handleA.url}/xrpc/org.p2pds.admin.getOverview`, { + const overview = await fetch(`${handleA.url}/xrpc/org.p2pds.app.getOverview`, { headers: { Authorization: `Bearer ${config2.AUTH_TOKEN}` }, }); const ov = (await overview.json()) as { @@ -270,7 +270,7 @@ describe("Two-node DID-less replication", () => { // Blocks should still be in IPFS (persisted to disk) const status2 = await fetch( - `${handleA.url}/xrpc/org.p2pds.admin.getDidStatus?did=${DID_ALICE}`, + `${handleA.url}/xrpc/org.p2pds.app.getDidStatus?did=${DID_ALICE}`, { headers: { Authorization: `Bearer ${config2.AUTH_TOKEN}` } }, ); const ds2 = (await status2.json()) as { diff --git a/src/xrpc/admin-e2e.test.ts b/src/xrpc/app-e2e.test.ts similarity index 96% rename from src/xrpc/admin-e2e.test.ts rename to src/xrpc/app-e2e.test.ts index 98bb56f..49822d2 100644 --- a/src/xrpc/admin-e2e.test.ts +++ b/src/xrpc/app-e2e.test.ts @@ -207,7 +207,7 @@ describe("Admin E2E: two-node replication + dashboard", () => { } it("getOverview shows synced replication state with aggregate metrics", async () => { - const res = await fetchB("/xrpc/org.p2pds.admin.getOverview"); + const res = await fetchB("/xrpc/org.p2pds.app.getOverview"); expect(res.status).toBe(200); const json = (await res.json()) as Record; @@ -252,7 +252,7 @@ describe("Admin E2E: two-node replication + dashboard", () => { it("getDidStatus shows blocks, records, and sync history for replicated DID", async () => { const res = await fetchB( - `/xrpc/org.p2pds.admin.getDidStatus?did=${NODE_A_DID}`, + `/xrpc/org.p2pds.app.getDidStatus?did=${NODE_A_DID}`, ); expect(res.status).toBe(200); @@ -276,13 +276,13 @@ describe("Admin E2E: two-node replication + dashboard", () => { it("dashboard returns HTML with expected structure", async () => { // Dashboard doesn't require auth header const res = await fetch( - `http://127.0.0.1:${portB}/xrpc/org.p2pds.admin.dashboard`, + `http://127.0.0.1:${portB}/`, ); expect(res.status).toBe(200); expect(res.headers.get("content-type")).toContain("text/html"); const html = await res.text(); - expect(html).toContain("P2PDS Admin"); + expect(html).toContain("P2PDS"); expect(html).toContain('id="section-overview"'); expect(html).toContain('id="section-metrics"'); expect(html).toContain('id="section-replication"'); @@ -292,7 +292,7 @@ describe("Admin E2E: two-node replication + dashboard", () => { }); it("getSyncHistory returns sync events after replication", async () => { - const res = await fetchB("/xrpc/org.p2pds.admin.getSyncHistory"); + const res = await fetchB("/xrpc/org.p2pds.app.getSyncHistory"); expect(res.status).toBe(200); const json = (await res.json()) as { history: Array> }; @@ -308,7 +308,7 @@ describe("Admin E2E: two-node replication + dashboard", () => { }); it("getNetworkStatus returns valid response", async () => { - const res = await fetchB("/xrpc/org.p2pds.admin.getNetworkStatus"); + const res = await fetchB("/xrpc/org.p2pds.app.getNetworkStatus"); expect(res.status).toBe(200); const json = (await res.json()) as Record; diff --git a/src/xrpc/admin.test.ts b/src/xrpc/app.test.ts similarity index 91% rename from src/xrpc/admin.test.ts rename to src/xrpc/app.test.ts index 07b4e18..a4adcf2 100644 --- a/src/xrpc/admin.test.ts +++ b/src/xrpc/app.test.ts @@ -97,11 +97,11 @@ describe("Admin endpoints: auth required", () => { it("returns 401 for all admin endpoints without auth", async () => { const getEndpoints = [ - "/xrpc/org.p2pds.admin.getOverview", - "/xrpc/org.p2pds.admin.getDidStatus?did=did:plc:test", - "/xrpc/org.p2pds.admin.getNetworkStatus", - "/xrpc/org.p2pds.admin.getPolicies", - "/xrpc/org.p2pds.admin.getSyncHistory", + "/xrpc/org.p2pds.app.getOverview", + "/xrpc/org.p2pds.app.getDidStatus?did=did:plc:test", + "/xrpc/org.p2pds.app.getNetworkStatus", + "/xrpc/org.p2pds.app.getPolicies", + "/xrpc/org.p2pds.app.getSyncHistory", ]; for (const endpoint of getEndpoints) { @@ -110,8 +110,8 @@ describe("Admin endpoints: auth required", () => { } const postEndpoints = [ - "/xrpc/org.p2pds.admin.addDid", - "/xrpc/org.p2pds.admin.removeDid", + "/xrpc/org.p2pds.app.addDid", + "/xrpc/org.p2pds.app.removeDid", ]; for (const endpoint of postEndpoints) { @@ -146,7 +146,7 @@ describe("Admin: getOverview", () => { const firehose = new Firehose(repoManager); const app = createApp(config, firehose, undefined, undefined, undefined, undefined, undefined, repoManager); - const res = await authGet(app, "/xrpc/org.p2pds.admin.getOverview"); + const res = await authGet(app, "/xrpc/org.p2pds.app.getOverview"); expect(res.status).toBe(200); const json = await res.json() as Record; @@ -211,7 +211,7 @@ describe("Admin: getOverview", () => { repoManager, ); - const res = await authGet(app, "/xrpc/org.p2pds.admin.getOverview"); + const res = await authGet(app, "/xrpc/org.p2pds.app.getOverview"); expect(res.status).toBe(200); const json = await res.json() as Record; @@ -303,7 +303,7 @@ describe("Admin: getDidStatus", () => { }); it("returns 400 when did param is missing", async () => { - const res = await authGet(app, "/xrpc/org.p2pds.admin.getDidStatus"); + const res = await authGet(app, "/xrpc/org.p2pds.app.getDidStatus"); expect(res.status).toBe(400); const json = await res.json() as Record; expect(json.error).toBe("MissingParameter"); @@ -319,7 +319,7 @@ describe("Admin: getDidStatus", () => { const res = await authGet( app, - `/xrpc/org.p2pds.admin.getDidStatus?did=${trackedDid}`, + `/xrpc/org.p2pds.app.getDidStatus?did=${trackedDid}`, ); expect(res.status).toBe(200); @@ -341,7 +341,7 @@ describe("Admin: getDidStatus", () => { it("returns nulls for an untracked DID", async () => { const res = await authGet( app, - "/xrpc/org.p2pds.admin.getDidStatus?did=did:plc:unknown", + "/xrpc/org.p2pds.app.getDidStatus?did=did:plc:unknown", ); expect(res.status).toBe(200); @@ -378,7 +378,7 @@ describe("Admin: getNetworkStatus", () => { const firehose = new Firehose(repoManager); const app = createApp(config, firehose, undefined, undefined, undefined, undefined, undefined, repoManager); - const res = await authGet(app, "/xrpc/org.p2pds.admin.getNetworkStatus"); + const res = await authGet(app, "/xrpc/org.p2pds.app.getNetworkStatus"); expect(res.status).toBe(200); const json = await res.json() as Record; @@ -420,7 +420,7 @@ describe("Admin: getNetworkStatus", () => { repoManager, ); - const res = await authGet(app, "/xrpc/org.p2pds.admin.getNetworkStatus"); + const res = await authGet(app, "/xrpc/org.p2pds.app.getNetworkStatus"); expect(res.status).toBe(200); const json = await res.json() as Record; @@ -455,7 +455,7 @@ describe("Admin: getPolicies", () => { const firehose = new Firehose(repoManager); const app = createApp(config, firehose, undefined, undefined, undefined, undefined, undefined, repoManager); - const res = await authGet(app, "/xrpc/org.p2pds.admin.getPolicies"); + const res = await authGet(app, "/xrpc/org.p2pds.app.getPolicies"); expect(res.status).toBe(200); const json = await res.json() as Record; @@ -507,7 +507,7 @@ describe("Admin: getPolicies", () => { repoManager, ); - const res = await authGet(app, "/xrpc/org.p2pds.admin.getPolicies"); + const res = await authGet(app, "/xrpc/org.p2pds.app.getPolicies"); expect(res.status).toBe(200); const json = await res.json() as Record; @@ -577,7 +577,7 @@ describe("Admin: getPolicies", () => { repoManager, ); - const res = await authGet(app, "/xrpc/org.p2pds.admin.getPolicies"); + const res = await authGet(app, "/xrpc/org.p2pds.app.getPolicies"); expect(res.status).toBe(200); const json = await res.json() as Record; @@ -615,12 +615,12 @@ describe("Admin: getDashboard", () => { const firehose = new Firehose(repoManager); const app = createApp(config, firehose, undefined, undefined, undefined, undefined, undefined, repoManager); - const res = await noAuthGet(app, "/xrpc/org.p2pds.admin.dashboard"); + const res = await noAuthGet(app, "/"); expect(res.status).toBe(200); expect(res.headers.get("content-type")).toContain("text/html"); const html = await res.text(); - expect(html).toContain("P2PDS Admin"); + expect(html).toContain("P2PDS"); expect(html).toContain('id="section-overview"'); expect(html).toContain('id="section-metrics"'); expect(html).toContain('id="section-replication"'); @@ -705,26 +705,26 @@ describe("Admin: addDid / removeDid", () => { }); it("returns 400 when did is missing", async () => { - const res = await authPost(app, "/xrpc/org.p2pds.admin.addDid", {}); + const res = await authPost(app, "/xrpc/org.p2pds.app.addDid", {}); expect(res.status).toBe(400); const json = await res.json() as Record; expect(json.error).toBe("MissingParameter"); }); it("returns 400 for invalid DID format", async () => { - const res = await authPost(app, "/xrpc/org.p2pds.admin.addDid", { did: "not-a-did" }); + const res = await authPost(app, "/xrpc/org.p2pds.app.addDid", { did: "not-a-did" }); expect(res.status).toBe(400); const json = await res.json() as Record; expect(json.error).toBe("InvalidDid"); }); it("allows adding own DID for self-replication", async () => { - const res = await authPost(app, "/xrpc/org.p2pds.admin.addDid", { did: "did:plc:test123" }); + const res = await authPost(app, "/xrpc/org.p2pds.app.addDid", { did: "did:plc:test123" }); expect(res.status).toBe(200); }); it("reports already_tracked for config DID", async () => { - const res = await authPost(app, "/xrpc/org.p2pds.admin.addDid", { did: configDid }); + const res = await authPost(app, "/xrpc/org.p2pds.app.addDid", { did: configDid }); expect(res.status).toBe(200); const json = await res.json() as Record; expect(json.status).toBe("already_tracked"); @@ -733,7 +733,7 @@ describe("Admin: addDid / removeDid", () => { it("adds a new DID", async () => { const newDid = "did:plc:newdid123"; - const res = await authPost(app, "/xrpc/org.p2pds.admin.addDid", { did: newDid }); + const res = await authPost(app, "/xrpc/org.p2pds.app.addDid", { did: newDid }); expect(res.status).toBe(200); const json = await res.json() as Record; expect(json.status).toBe("added"); @@ -748,8 +748,8 @@ describe("Admin: addDid / removeDid", () => { it("idempotent: re-adding returns already_tracked", async () => { const newDid = "did:plc:idempotent1"; - await authPost(app, "/xrpc/org.p2pds.admin.addDid", { did: newDid }); - const res = await authPost(app, "/xrpc/org.p2pds.admin.addDid", { did: newDid }); + await authPost(app, "/xrpc/org.p2pds.app.addDid", { did: newDid }); + const res = await authPost(app, "/xrpc/org.p2pds.app.addDid", { did: newDid }); expect(res.status).toBe(200); const json = await res.json() as Record; expect(json.status).toBe("already_tracked"); @@ -757,7 +757,7 @@ describe("Admin: addDid / removeDid", () => { }); it("cannot remove a config DID", async () => { - const res = await authPost(app, "/xrpc/org.p2pds.admin.removeDid", { did: configDid }); + const res = await authPost(app, "/xrpc/org.p2pds.app.removeDid", { did: configDid }); expect(res.status).toBe(400); const json = await res.json() as Record; expect(json.error).toBe("CannotRemove"); @@ -765,9 +765,9 @@ describe("Admin: addDid / removeDid", () => { it("removes an admin-added DID", async () => { const newDid = "did:plc:removable1"; - await authPost(app, "/xrpc/org.p2pds.admin.addDid", { did: newDid }); + await authPost(app, "/xrpc/org.p2pds.app.addDid", { did: newDid }); - const res = await authPost(app, "/xrpc/org.p2pds.admin.removeDid", { did: newDid }); + const res = await authPost(app, "/xrpc/org.p2pds.app.removeDid", { did: newDid }); expect(res.status).toBe(200); const json = await res.json() as Record; expect(json.status).toBe("removed"); @@ -776,14 +776,14 @@ describe("Admin: addDid / removeDid", () => { it("removes with purgeData deletes all tracking data", async () => { const newDid = "did:plc:purgeable1"; - await authPost(app, "/xrpc/org.p2pds.admin.addDid", { did: newDid }); + await authPost(app, "/xrpc/org.p2pds.app.addDid", { did: newDid }); // Add some tracking data const syncStorage = replicationManager.getSyncStorage(); syncStorage.trackBlocks(newDid, ["bafyblock1"]); syncStorage.trackRecordPaths(newDid, ["app.bsky.feed.post/abc"]); - const res = await authPost(app, "/xrpc/org.p2pds.admin.removeDid", { did: newDid, purgeData: true }); + const res = await authPost(app, "/xrpc/org.p2pds.app.removeDid", { did: newDid, purgeData: true }); expect(res.status).toBe(200); const json = await res.json() as Record; expect(json.status).toBe("removed"); diff --git a/src/xrpc/admin.ts b/src/xrpc/app.ts similarity index 93% rename from src/xrpc/admin.ts rename to src/xrpc/app.ts index d72fd98..b7d2aae 100644 --- a/src/xrpc/admin.ts +++ b/src/xrpc/app.ts @@ -153,7 +153,7 @@ export function getDashboard( -P2PDS Admin +P2PDS