/** * island:discover * * Find connected components in the global network.cosmik.connection graph. * Uses SQLite to cache network data so we don't rebuild each time. * * Usage: * bun island:discover # find all connected components * bun island:discover --refresh # force re-fetch all DIDs * bun island:discover # show full details for an island * bun island:discover --min-size 5 # only show components with >= 5 nodes * bun island:discover --json # machine-readable output */ import { resolve, dirname } from 'node:path'; import { fileURLToPath } from 'node:url'; import { createHash } from 'node:crypto'; import { Database } from 'bun:sqlite'; import { loadDotEnv } from './cli-utils.js'; import { islandId, domainFromUrl } from './island-shared.js'; const __dirname = dirname(fileURLToPath(import.meta.url)); await loadDotEnv(resolve(__dirname, '..', '.env')); const RELAY = process.env.ATPROTO_RELAY ?? 'https://bsky.network'; const CARRY_DIR = process.env.CARRY_DIR ?? '/home/nandi/.local/share/carry-vault'; const DB_PATH = process.env.ISLAND_DETECT_DB ?? resolve(CARRY_DIR, '.island-detect.db'); const COLLECTION = 'network.cosmik.connection'; const CACHE_TTL_MS = 60 * 60 * 1000; // 1 hour interface Component { nodes: Set; dids: Set; lexmin: string; islandId: string; } // --- SQLite --- function openDb(): Database { const db = new Database(DB_PATH, { create: true }); db.exec('PRAGMA journal_mode = WAL'); db.exec('PRAGMA synchronous = NORMAL'); db.exec(` CREATE TABLE IF NOT EXISTS edges ( source TEXT NOT NULL, target TEXT NOT NULL, relation TEXT, did TEXT NOT NULL, fetched_at TEXT NOT NULL, PRIMARY KEY (source, target, did) ); CREATE TABLE IF NOT EXISTS dids ( did TEXT PRIMARY KEY, handle TEXT, pds TEXT, edge_count INTEGER DEFAULT 0, fetched_at TEXT ); CREATE INDEX IF NOT EXISTS idx_edges_source ON edges(source); CREATE INDEX IF NOT EXISTS idx_edges_target ON edges(target); `); return db; } // --- Network discovery --- async function discoverDIDs(): Promise { const dids: string[] = []; let cursor = ''; while (true) { const url = new URL(`${RELAY}/xrpc/com.atproto.sync.listReposByCollection`); url.searchParams.set('collection', COLLECTION); url.searchParams.set('limit', '1000'); if (cursor) url.searchParams.set('cursor', cursor); const res = await fetch(url.toString()); if (!res.ok) break; const data = await res.json() as { repos: Array<{ did: string }>; cursor?: string }; for (const r of data.repos) dids.push(r.did); cursor = data.cursor ?? ''; if (!cursor) break; await new Promise((r) => setTimeout(r, 200)); } return dids; } async function resolvePDS(did: string): Promise { try { const res = await fetch(`https://plc.directory/${did}`); if (!res.ok) return null; const doc = await res.json() as { service?: Array<{ id: string; serviceEndpoint: string }> }; return doc.service?.find((s) => s.id === '#atproto_pds')?.serviceEndpoint ?? null; } catch { return null; } } async function resolveHandle(did: string): Promise { try { const res = await fetch(`https://plc.directory/${did}`); if (!res.ok) return null; const doc = await res.json() as { alsoKnownAs?: string[] }; return doc.alsoKnownAs?.find((h) => h.startsWith('at://'))?.replace('at://', '') ?? null; } catch { return null; } } async function fetchConnectionRecords(pds: string, did: string): Promise> { const records: Array<{ source: string; target: string; relation?: string; atUri: string }> = []; let cursor = ''; while (true) { const url = new URL(`${pds}/xrpc/com.atproto.repo.listRecords`); url.searchParams.set('repo', did); url.searchParams.set('collection', COLLECTION); url.searchParams.set('limit', '100'); if (cursor) url.searchParams.set('cursor', cursor); const res = await fetch(url.toString()); if (!res.ok) break; const data = await res.json() as { records: Array<{ uri: string; value: { source?: string; target?: string; connectionType?: string } }>; cursor?: string; }; for (const r of data.records) { if (r.value.source && r.value.target) { records.push({ source: r.value.source, target: r.value.target, relation: r.value.connectionType, atUri: r.uri }); } } cursor = data.cursor ?? ''; if (!cursor) break; await new Promise((r) => setTimeout(r, 100)); } return records; } // --- Cache operations --- function isCached(db: Database, did: string): boolean { const row = db.query('SELECT fetched_at FROM dids WHERE did = ?').get(did) as { fetched_at: string } | null; if (!row) return false; return Date.now() - new Date(row.fetched_at).getTime() < CACHE_TTL_MS; } function upsertDid(db: Database, did: string, handle: string | null, pds: string | null, edgeCount: number): void { const now = new Date().toISOString(); db.query(` INSERT INTO dids (did, handle, pds, edge_count, fetched_at) VALUES (?, ?, ?, ?, ?) ON CONFLICT(did) DO UPDATE SET handle=?, pds=?, edge_count=?, fetched_at=? `).run(did, handle, pds, edgeCount, now, handle, pds, edgeCount, now); } function upsertEdges(db: Database, did: string, records: Array<{ source: string; target: string; relation?: string; atUri: string }>): void { const now = new Date().toISOString(); // Add at_uri column if missing try { db.exec('ALTER TABLE edges ADD COLUMN at_uri TEXT'); } catch {} const insert = db.query(` INSERT INTO edges (source, target, relation, did, at_uri, fetched_at) VALUES (?, ?, ?, ?, ?, ?) ON CONFLICT(source, target, did) DO UPDATE SET relation=?, at_uri=?, fetched_at=? `); // Delete stale edges for this DID before inserting db.query('DELETE FROM edges WHERE did = ?').run(did); for (const r of records) { insert.run(r.source, r.target, r.relation ?? null, did, r.atUri, now, r.relation ?? null, r.atUri, now); } } // --- Connected components --- function findComponents(db: Database): Component[] { // Build adjacency list from all edges const adj = new Map>(); const edgeDids = new Map>(); // node -> which DIDs have edges touching it const rows = db.query('SELECT source, target, did FROM edges').all() as Array<{ source: string; target: string; did: string }>; for (const row of rows) { const s = row.source; const t = row.target; if (!adj.has(s)) adj.set(s, new Set()); if (!adj.has(t)) adj.set(t, new Set()); adj.get(s)!.add(t); adj.get(t)!.add(s); if (!edgeDids.has(s)) edgeDids.set(s, new Set()); if (!edgeDids.has(t)) edgeDids.set(t, new Set()); edgeDids.get(s)!.add(row.did); edgeDids.get(t)!.add(row.did); } // BFS to find components const visited = new Set(); const components: Component[] = []; for (const node of adj.keys()) { if (visited.has(node)) continue; const compNodes = new Set(); const compDids = new Set(); const queue = [node]; while (queue.length > 0) { const current = queue.pop()!; if (visited.has(current)) continue; visited.add(current); compNodes.add(current); for (const did of edgeDids.get(current) ?? []) compDids.add(did); for (const neighbor of adj.get(current) ?? []) { if (!visited.has(neighbor)) queue.push(neighbor); } } const lexmin = [...compNodes].sort()[0]; components.push({ nodes: compNodes, dids: compDids, lexmin, islandId: islandId(lexmin) }); } return components.sort((a, b) => b.nodes.size - a.nodes.size); } // --- Main --- async function main(): Promise { const args = process.argv.slice(2); const wantJson = args.includes('--json'); const refresh = args.includes('--refresh'); const targetId = args.find((a) => !a.startsWith('--')); const minSizeArg = args.find((a) => a.startsWith('--min-size=')); const minSize = minSizeArg ? Number(minSizeArg.split('=')[1]) : 3; const db = openDb(); console.error('Discovering network DIDs...'); const networkDids = await discoverDIDs(); console.error(`Found ${networkDids.length} DIDs`); // Fetch edges for each DID (use cache unless --refresh) let fetched = 0; let cached = 0; for (let i = 0; i < networkDids.length; i++) { const did = networkDids[i]; if (!refresh && isCached(db, did)) { cached++; continue; } console.error(`[${i + 1}/${networkDids.length}] Fetching ${did}...`); const pds = await resolvePDS(did); if (!pds) { console.error(' Could not resolve PDS'); continue; } const handle = await resolveHandle(did); const records = await fetchConnectionRecords(pds, did); upsertDid(db, did, handle, pds, records.length); upsertEdges(db, did, records); fetched++; } console.error(`Fetched: ${fetched}, Cached: ${cached}`); // Stats const totalEdges = (db.query('SELECT COUNT(*) as c FROM edges').get() as { c: number }).c; const totalDids = (db.query('SELECT COUNT(*) as c FROM dids').get() as { c: number }).c; const uniqueUrls = new Set(); const urlRows = db.query('SELECT DISTINCT source FROM edges UNION SELECT DISTINCT target FROM edges').all() as Array<{ source: string }>; for (const r of urlRows) uniqueUrls.add(r.source); console.error(`Database: ${totalEdges} edges, ${totalDids} DIDs, ${uniqueUrls.size} unique URLs`); // Find connected components const components = findComponents(db); const filtered = components.filter((c) => c.nodes.size >= minSize); // Show a specific island by ID if (targetId) { const comp = components.find((c) => c.islandId === targetId); if (!comp) { console.error(`Island ${targetId} not found`); process.exit(1); } const didPlaceholders = [...comp.dids].map(() => '?').join(','); const didRows = db.query(`SELECT did, handle, edge_count FROM dids WHERE did IN (${didPlaceholders})`).all(...comp.dids) as Array<{ did: string; handle: string | null; edge_count: number }>; if (wantJson) { console.log(JSON.stringify({ id: comp.islandId, lexmin: comp.lexmin, nodeCount: comp.nodes.size, didCount: comp.dids.size, nodes: [...comp.nodes].sort(), dids: didRows, }, null, 2)); } else { console.log(`Island ${comp.islandId}: ${comp.nodes.size} nodes, ${comp.dids.size} DIDs`); console.log(`Lexmin: ${comp.lexmin}`); console.log('\nNodes:'); for (const url of [...comp.nodes].sort()) console.log(` ${url}`); console.log('\nDIDs:'); for (const d of didRows) console.log(` ${d.handle ?? d.did} (${d.edge_count} edges)`); } db.close(); return; } if (wantJson) { console.log(JSON.stringify({ totalEdges, totalDids, uniqueUrls: uniqueUrls.size, componentCount: filtered.length, components: filtered.map((c) => ({ id: c.islandId, lexmin: c.lexmin, nodeCount: c.nodes.size, didCount: c.dids.size, sampleNodes: [...c.nodes].sort().slice(0, 10), dids: [...c.dids], })), }, null, 2)); db.close(); return; } console.log(`\nIsland discover: ${totalDids} DIDs, ${totalEdges} edges, ${uniqueUrls.size} unique URLs`); console.log(`Connected components (>= ${minSize} nodes): ${filtered.length}\n`); for (const comp of filtered.slice(0, 30)) { const domain = domainFromUrl(comp.lexmin); console.log(`${comp.islandId} ${comp.nodes.size} nodes, ${comp.dids.size} DIDs ${domain}`); } // Show which component we're in (if we can determine our DID) const ourDid = process.env.BLUESKY_IDENTIFIER ? (db.query('SELECT did FROM dids WHERE handle = ? OR did LIKE ?').get(process.env.BLUESKY_IDENTIFIER, `%${process.env.BLUESKY_IDENTIFIER}%`) as { did: string } | null) : null; if (ourDid) { const ourComp = components.find((c) => c.dids.has(ourDid.did)); if (ourComp) { console.log(`\nOur island: ${ourComp.islandId} (${ourComp.nodes.size} nodes, shared with ${ourComp.dids.size - 1} other DIDs)`); } } db.close(); } main().catch((err) => { console.error('Fatal:', err instanceof Error ? err.message : String(err)); process.exit(1); });