diff --git a/README.md b/README.md index 8937a2b..b49583c 100644 --- a/README.md +++ b/README.md @@ -119,6 +119,19 @@ Content-addressed retrieval is unforgeable: correct bytes for a CID = proof of s Challenge-response protocol: `StorageChallenge` → `StorageChallengeResponse` → `StorageChallengeResult`. Deterministic generation from epoch + DIDs + nonce. Transport-agnostic with libp2p primary and HTTP fallback. +### Lexicon index + +Every record path stored during sync includes a lexicon NSID (the collection portion of `collection/rkey`). P2PDS aggregates these into a local lexicon index — a catalog of every lexicon encountered across all replicated repos. + +- **Automatic population** — updated incrementally after every full sync and firehose event, rebuilt from scratch on startup +- **Public API** — three unauthenticated endpoints for querying: + - `GET /xrpc/org.p2pds.lexicon.search?q=app.bsky&limit=50` — prefix search + - `GET /xrpc/org.p2pds.lexicon.list?limit=100` — all NSIDs by record count + - `GET /xrpc/org.p2pds.lexicon.stats` — aggregate stats (unique NSIDs, total records) +- **UI** — searchable table in the web interface showing NSID, record count, repo count, first/last seen dates + +This is the foundation for distributed lexicon discovery — nodes can query each other's indexes to find who stores what. + ### Storage All persistent state in a single SQLite database (`pds.db`): @@ -126,8 +139,19 @@ All persistent state in a single SQLite database (`pds.db`): - **IPFS blocks/datastore** — SQLite-backed, no filesystem churn - **Replication state** — sync progress, peer info, block/blob tracking, firehose cursor - **Challenge history** — proof-of-storage results and peer reliability scores +- **Lexicon index** — aggregated NSID usage across all replicated repos +- **PLC mirror** — archived PLC operation logs for tracked DIDs - **Node identity** — DID + handle, established on first OAuth login +### PLC log archiving + +P2PDS mirrors PLC operation logs for all tracked `did:plc` DIDs. The PLC directory is the root of trust for DID resolution — if it goes down or loses data, DID documents become unresolvable. By archiving PLC logs locally, each node maintains an independent backup of the identity layer. + +- **Automatic** — logs are fetched on first sync and refreshed every 6 hours +- **Cross-node sharing** — public endpoint `GET /xrpc/org.p2pds.plc.getLog?did=...` lets nodes fetch PLC logs from each other, not just from the central PLC directory +- **Per-DID status** — the UI shows PLC archive status (op count, last fetch, tombstone state) for each tracked DID +- **Validation** — operation chains are validated on fetch + ### Policy engine Every replication relationship is backed by a policy object with lifecycle state, consent tracking, and merge rules. Policies are created automatically from offer negotiation or manually via config/UI. @@ -166,6 +190,8 @@ Your PDS holds rotation keys for your DID. If the PDS is gone, those keys are go **Handle survival:** If you used a custom domain handle (verified via DNS), it survives migration — the DNS record still points to your DID. If you used a `.bsky.social` handle, it's gone — that subdomain is controlled by Bluesky. +**Rotation key management in p2pds:** P2PDS includes a UI for adding rotation keys to your PLC document. The flow: request a PLC operation token (triggers an email from your PDS), enter the token, and submit your public key. This uses the standard `com.atproto.identity.signPlcOperation` API. You can also view your current rotation keys from the web interface. + **The honest reality for most Bluesky users today:** Bluesky has not yet shipped user-facing rotation key management. Most users don't independently hold a rotation key. If Bluesky's PDS infrastructure disappeared tomorrow, those users would have their data (thanks to p2pds) but could not recover their identity. This is an upstream gap in the atproto ecosystem, not something p2pds can solve — but it's important to understand. The data preservation is still valuable: your posts, social graph, and media survive, even if reattaching them to the same DID requires key management tooling that doesn't exist yet. ## Stack @@ -173,7 +199,7 @@ Your PDS holds rotation keys for your DID. If the PDS is gone, those keys are go - **Runtime**: Node.js, TypeScript (ES2022, strict) - **HTTP**: Hono - **Database**: better-sqlite3 -- **IPFS**: Helia with minimal libp2p (TCP + noise + yamux + autoNAT) +- **IPFS**: Helia with minimal libp2p (TCP + noise + yamux + autoNAT + kadDHT client mode) - **UI**: Lit web components, esbuild-bundled - **Identity**: AT Protocol DIDs via PLC directory - **Auth**: OAuth (primary) or legacy JWT (fallback) @@ -210,8 +236,9 @@ src/ ipfs.ts IpfsService (Helia wrapper, SQLite-backed) build-ui.ts esbuild bundler for Lit UI ui/ Lit web components (app shell, cards, state) - replication/ Sync, verification, challenges, offers + replication/ Sync, verification, challenges, offers, lexicon index policy/ Policy engine types, engine, presets + identity/ PLC mirror, rotation key management oauth/ OAuth client, routes, PdsClient xrpc/ XRPC endpoint handlers middleware/ Auth, rate limiting, body limits diff --git a/src/index.ts b/src/index.ts index fab4540..7c7688b 100644 --- a/src/index.ts +++ b/src/index.ts @@ -187,6 +187,15 @@ export function createApp( }), ); + // Lexicon index endpoints (unauthenticated, rate-limited) + const lexiconRL = rateLimitMiddleware(rateLimiter, { + pool: "lexicon", + rule: { maxRequests: 300, windowMs: w }, + }); + app.use("/xrpc/org.p2pds.lexicon.search", lexiconRL); + app.use("/xrpc/org.p2pds.lexicon.list", lexiconRL); + app.use("/xrpc/org.p2pds.lexicon.stats", lexiconRL); + // App endpoints const appRL = rateLimitMiddleware(rateLimiter, { pool: "app", @@ -703,6 +712,17 @@ export function createApp( app_routes.getPlcLogPublic(c, db), ); + // Lexicon index endpoints (unauthenticated) + app.get("/xrpc/org.p2pds.lexicon.search", (c) => + app_routes.searchLexicons(c, replicationManager), + ); + app.get("/xrpc/org.p2pds.lexicon.list", (c) => + app_routes.listLexicons(c, replicationManager), + ); + app.get("/xrpc/org.p2pds.lexicon.stats", (c) => + app_routes.getLexiconStats(c, replicationManager), + ); + app.post("/xrpc/org.p2pds.replication.revokeOffer", requireAuth, async (c) => { if (!replicationManager) { return c.json({ error: "ReplicationNotEnabled", message: "Replication is not enabled" }, 400); diff --git a/src/replication/replication-manager.ts b/src/replication/replication-manager.ts index 8a25b5d..f3e1186 100644 --- a/src/replication/replication-manager.ts +++ b/src/replication/replication-manager.ts @@ -59,6 +59,16 @@ export interface SyncProgressEvent { missingBlocks?: number; } +/** Extract unique NSIDs from record paths (collection/rkey format). */ +function extractNsids(paths: string[]): string[] { + const nsids = new Set(); + for (const p of paths) { + const slash = p.indexOf("/"); + if (slash > 0) nsids.add(p.slice(0, slash)); + } + return [...nsids]; +} + /** How old cached peer info can be before re-fetching (1 hour). */ const PEER_INFO_TTL_MS = 60 * 60 * 1000; @@ -1215,6 +1225,11 @@ export class ReplicationManager { ); this.syncStorage.clearRecordPaths(did); this.syncStorage.trackRecordPaths(did, recordPaths); + // Update lexicon index with NSIDs from these paths + const nsids = extractNsids(recordPaths); + if (nsids.length > 0) { + this.syncStorage.updateLexiconIndex(nsids); + } } catch { // Non-fatal: path extraction is best-effort } @@ -1513,6 +1528,13 @@ export class ReplicationManager { console.error("[replication] Initial PLC mirror fetch error:", err); }); + // Build lexicon index from existing record paths + try { + this.syncStorage.rebuildLexiconIndex(); + } catch (err) { + console.error("[replication] Lexicon index rebuild error:", err); + } + // Run verification once on startup, then on a timer this.runVerification().catch((err) => { console.error("Initial verification error:", err); @@ -1792,6 +1814,11 @@ export class ReplicationManager { .map((op) => op.path); if (createdPaths.length > 0) { this.syncStorage.trackRecordPaths(did, createdPaths); + // Update lexicon index with NSIDs from created paths + const nsids = extractNsids(createdPaths); + if (nsids.length > 0) { + this.syncStorage.updateLexiconIndex(nsids); + } } if (deletedPaths.length > 0) { this.syncStorage.removeRecordPaths(did, deletedPaths); @@ -2116,6 +2143,20 @@ export class ReplicationManager { return this.syncStorage.getAllStates(); } + /** + * Query the lexicon index with optional prefix filter. + */ + getLexiconIndex(prefix?: string, limit?: number) { + return this.syncStorage.getLexiconIndex(prefix, limit); + } + + /** + * Get aggregate lexicon stats. + */ + getLexiconStats() { + return this.syncStorage.getLexiconStats(); + } + /** * Get all offered DIDs (awaiting mutual consent). */ diff --git a/src/replication/sync-storage.ts b/src/replication/sync-storage.ts index e845223..f7a06f9 100644 --- a/src/replication/sync-storage.ts +++ b/src/replication/sync-storage.ts @@ -112,6 +112,17 @@ export class SyncStorage { ); `); + // Lexicon index table: aggregates NSID usage across all replicated repos. + this.db.exec(` + CREATE TABLE IF NOT EXISTS lexicon_index ( + nsid TEXT PRIMARY KEY, + first_seen_at TEXT NOT NULL, + last_seen_at TEXT NOT NULL, + record_count INTEGER NOT NULL DEFAULT 0, + repo_count INTEGER NOT NULL DEFAULT 0 + ); + `); + // Offered DIDs table: tracks DIDs we've offered to replicate // but don't yet have mutual consent for. this.db.exec(` @@ -1152,6 +1163,127 @@ export class SyncStorage { return row?.root_cid ?? null; } + // ============================================ + // Lexicon index + // ============================================ + + /** + * Upsert NSIDs into the lexicon index. + * Updates last_seen_at and recomputes record_count/repo_count from source data. + */ + updateLexiconIndex(nsids: string[]): void { + if (nsids.length === 0) return; + const now = new Date().toISOString(); + const upsert = this.db.prepare( + `INSERT INTO lexicon_index (nsid, first_seen_at, last_seen_at, record_count, repo_count) + VALUES (?, ?, ?, 0, 0) + ON CONFLICT(nsid) DO UPDATE SET last_seen_at = excluded.last_seen_at`, + ); + const countRecords = this.db.prepare( + `SELECT COUNT(*) as cnt FROM replication_record_paths + WHERE record_path LIKE ? || '/%'`, + ); + const countRepos = this.db.prepare( + `SELECT COUNT(DISTINCT did) as cnt FROM replication_record_paths + WHERE record_path LIKE ? || '/%'`, + ); + const updateCounts = this.db.prepare( + `UPDATE lexicon_index SET record_count = ?, repo_count = ? WHERE nsid = ?`, + ); + + const batch = this.db.transaction((items: string[]) => { + for (const nsid of items) { + upsert.run(nsid, now, now); + const rc = countRecords.get(nsid) as { cnt: number }; + const rp = countRepos.get(nsid) as { cnt: number }; + updateCounts.run(rc.cnt, rp.cnt, nsid); + } + }); + batch(nsids); + } + + /** + * Full rebuild of the lexicon index from replication_record_paths. + */ + rebuildLexiconIndex(): void { + const now = new Date().toISOString(); + this.db.exec("DELETE FROM lexicon_index"); + + const rows = this.db.prepare( + `SELECT + SUBSTR(record_path, 1, INSTR(record_path, '/') - 1) as nsid, + COUNT(*) as record_count, + COUNT(DISTINCT did) as repo_count, + MIN(?) as first_seen_at, + MAX(?) as last_seen_at + FROM replication_record_paths + WHERE INSTR(record_path, '/') > 0 + GROUP BY SUBSTR(record_path, 1, INSTR(record_path, '/') - 1)`, + ).all(now, now) as Array<{ + nsid: string; + record_count: number; + repo_count: number; + first_seen_at: string; + last_seen_at: string; + }>; + + if (rows.length === 0) return; + + const insert = this.db.prepare( + `INSERT INTO lexicon_index (nsid, first_seen_at, last_seen_at, record_count, repo_count) + VALUES (?, ?, ?, ?, ?)`, + ); + const batch = this.db.transaction(() => { + for (const r of rows) { + insert.run(r.nsid, r.first_seen_at, r.last_seen_at, r.record_count, r.repo_count); + } + }); + batch(); + } + + /** + * Query the lexicon index with optional NSID prefix filter. + */ + getLexiconIndex(prefix?: string, limit: number = 100): Array<{ + nsid: string; + firstSeenAt: string; + lastSeenAt: string; + recordCount: number; + repoCount: number; + }> { + const query = prefix + ? `SELECT * FROM lexicon_index WHERE nsid LIKE ? || '%' ORDER BY record_count DESC LIMIT ?` + : `SELECT * FROM lexicon_index ORDER BY record_count DESC LIMIT ?`; + const params = prefix ? [prefix, limit] : [limit]; + const rows = this.db.prepare(query).all(...params) as Array<{ + nsid: string; + first_seen_at: string; + last_seen_at: string; + record_count: number; + repo_count: number; + }>; + return rows.map((r) => ({ + nsid: r.nsid, + firstSeenAt: r.first_seen_at, + lastSeenAt: r.last_seen_at, + recordCount: r.record_count, + repoCount: r.repo_count, + })); + } + + /** + * Get aggregate lexicon stats. + */ + getLexiconStats(): { uniqueNsids: number; totalRecords: number } { + const row = this.db.prepare( + `SELECT COUNT(*) as unique_nsids, COALESCE(SUM(record_count), 0) as total_records FROM lexicon_index`, + ).get() as { unique_nsids: number; total_records: number }; + return { + uniqueNsids: row.unique_nsids, + totalRecords: row.total_records, + }; + } + /** * Delete all data from all replication tables in a single transaction. * Used during full disconnect to wipe the node clean. @@ -1169,6 +1301,7 @@ export class SyncStorage { this.db.prepare("DELETE FROM sync_history").run(); this.db.prepare("DELETE FROM firehose_cursor").run(); this.db.prepare("DELETE FROM plc_mirror").run(); + this.db.prepare("DELETE FROM lexicon_index").run(); }); purge(); } diff --git a/src/ui/app.ts b/src/ui/app.ts index a998057..cb2f2aa 100644 --- a/src/ui/app.ts +++ b/src/ui/app.ts @@ -21,3 +21,4 @@ import "./components/policies-card"; import "./components/verification-card"; import "./components/confirm-dialog"; import "./components/recovery-key-dialog"; +import "./components/lexicon-index"; diff --git a/src/ui/components/app-shell.ts b/src/ui/components/app-shell.ts index d499b94..1927c49 100644 --- a/src/ui/components/app-shell.ts +++ b/src/ui/components/app-shell.ts @@ -178,6 +178,8 @@ export class AppShell extends LitElement { + + `; } diff --git a/src/ui/components/lexicon-index.ts b/src/ui/components/lexicon-index.ts new file mode 100644 index 0000000..893c676 --- /dev/null +++ b/src/ui/components/lexicon-index.ts @@ -0,0 +1,121 @@ +import { LitElement, html } from "lit"; +import { customElement, state } from "lit/decorators.js"; + +interface LexiconEntry { + nsid: string; + firstSeenAt: string; + lastSeenAt: string; + recordCount: number; + repoCount: number; +} + +@customElement("p2p-lexicon-index") +export class LexiconIndex extends LitElement { + createRenderRoot() { return this; } + + @state() lexicons: LexiconEntry[] = []; + @state() query = ""; + @state() uniqueNsids = 0; + @state() loading = false; + + private debounceTimer: ReturnType | null = null; + + connectedCallback() { + super.connectedCallback(); + this.fetchData(); + } + + private async fetchData() { + this.loading = true; + try { + const endpoint = this.query + ? `/xrpc/org.p2pds.lexicon.search?q=${encodeURIComponent(this.query)}&limit=50` + : `/xrpc/org.p2pds.lexicon.list?limit=100`; + const res = await fetch(endpoint); + if (res.ok) { + const data = await res.json(); + this.lexicons = data.lexicons ?? []; + } + + const statsRes = await fetch("/xrpc/org.p2pds.lexicon.stats"); + if (statsRes.ok) { + const stats = await statsRes.json(); + this.uniqueNsids = stats.uniqueNsids ?? 0; + } + } catch { + // Fetch failed + } + this.loading = false; + } + + private handleInput(e: Event) { + const input = e.target as HTMLInputElement; + this.query = input.value; + if (this.debounceTimer) clearTimeout(this.debounceTimer); + this.debounceTimer = setTimeout(() => this.fetchData(), 300); + } + + private fmtDate(iso: string): string { + if (!iso) return "-"; + const d = new Date(iso); + return d.toLocaleDateString(undefined, { month: "short", day: "numeric" }); + } + + render() { + return html` +
+
+ + Lexicons + ${this.uniqueNsids} unique + +
+ +
+ ${this.loading + ? html`

Loading...

` + : this.lexicons.length === 0 + ? html`

No lexicons found${this.query ? ` matching "${this.query}"` : ""}.

` + : html` +
+ + + + + + + + + + + + ${this.lexicons.map((lex) => html` + + + + + + + + `)} + +
NSIDRecordsReposFirst SeenLast Seen
${lex.nsid}${lex.recordCount.toLocaleString()}${lex.repoCount}${this.fmtDate(lex.firstSeenAt)}${this.fmtDate(lex.lastSeenAt)}
+
+ `} +
+
+ `; + } +} + +declare global { + interface HTMLElementTagNameMap { + "p2p-lexicon-index": LexiconIndex; + } +} diff --git a/src/ui/components/recovery-key-dialog.ts b/src/ui/components/recovery-key-dialog.ts index 38ba35e..432f6af 100644 --- a/src/ui/components/recovery-key-dialog.ts +++ b/src/ui/components/recovery-key-dialog.ts @@ -1,7 +1,7 @@ import { LitElement, html, nothing } from "lit"; import { customElement, state } from "lit/decorators.js"; import { generateMnemonic, mnemonicToSeedSync } from "@scure/bip39"; -import { wordlist } from "@scure/bip39/wordlists/english"; +import { wordlist } from "@scure/bip39/wordlists/english.js"; import { HDKey } from "@scure/bip32"; import { secp256k1 } from "@noble/curves/secp256k1"; import { base58btc } from "multiformats/bases/base58"; diff --git a/src/ui/styles/theme.ts b/src/ui/styles/theme.ts index 6a4309b..e7085f8 100644 --- a/src/ui/styles/theme.ts +++ b/src/ui/styles/theme.ts @@ -222,8 +222,9 @@ button:disabled { opacity: 0.4; cursor: not-allowed; } .account-selected { display: flex; align-items: center; gap: 0.4rem; padding: 0.35rem 0.5rem; background: var(--selected-bg); border: 1px solid var(--selected-border); border-radius: 4px; - height: 100%; box-sizing: border-box; position: relative; + box-sizing: border-box; position: relative; } +.add-did-row .account-selected { height: 100%; } .account-selected-avatar { width: 22px; height: 22px; border-radius: 50%; flex-shrink: 0; } .account-selected-placeholder { width: 22px; height: 22px; border-radius: 50%; background: var(--acct-placeholder-bg); diff --git a/src/xrpc/app.ts b/src/xrpc/app.ts index 55d429a..0ef8dc3 100644 --- a/src/xrpc/app.ts +++ b/src/xrpc/app.ts @@ -1061,6 +1061,53 @@ export function getPlcLogPublic( }); } +// ============================================ +// Lexicon index +// ============================================ + +export function searchLexicons( + c: Context, + replicationManager: ReplicationManager | undefined, +): Response { + if (!replicationManager) { + return c.json({ error: "ReplicationNotEnabled", message: "Replication is not enabled" }, 400); + } + + const q = c.req.query("q") || ""; + const limitStr = c.req.query("limit"); + const limit = limitStr ? Math.min(Math.max(parseInt(limitStr, 10) || 50, 1), 200) : 50; + + const results = replicationManager.getLexiconIndex(q || undefined, limit); + return c.json({ lexicons: results, query: q }); +} + +export function listLexicons( + c: Context, + replicationManager: ReplicationManager | undefined, +): Response { + if (!replicationManager) { + return c.json({ error: "ReplicationNotEnabled", message: "Replication is not enabled" }, 400); + } + + const limitStr = c.req.query("limit"); + const limit = limitStr ? Math.min(Math.max(parseInt(limitStr, 10) || 100, 1), 500) : 100; + + const results = replicationManager.getLexiconIndex(undefined, limit); + return c.json({ lexicons: results }); +} + +export function getLexiconStats( + c: Context, + replicationManager: ReplicationManager | undefined, +): Response { + if (!replicationManager) { + return c.json({ error: "ReplicationNotEnabled", message: "Replication is not enabled" }, 400); + } + + const stats = replicationManager.getLexiconStats(); + return c.json(stats); +} + // ============================================ // PLC rotation key management // ============================================