This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

stigmergic / src / island-discover.ts
12 kB 332 lines
1/** 2 * island:discover 3 * 4 * Find connected components in the global network.cosmik.connection graph. 5 * Uses SQLite to cache network data so we don't rebuild each time. 6 * 7 * Usage: 8 * bun island:discover # find all connected components 9 * bun island:discover --refresh # force re-fetch all DIDs 10 * bun island:discover <island-id> # show full details for an island 11 * bun island:discover --min-size 5 # only show components with >= 5 nodes 12 * bun island:discover --json # machine-readable output 13 */ 14 15import { resolve, dirname } from 'node:path'; 16import { fileURLToPath } from 'node:url'; 17import { createHash } from 'node:crypto'; 18import { Database } from 'bun:sqlite'; 19import { loadDotEnv } from './cli-utils.js'; 20import { islandId, domainFromUrl } from './island-shared.js'; 21 22const __dirname = dirname(fileURLToPath(import.meta.url)); 23await loadDotEnv(resolve(__dirname, '..', '.env')); 24 25const RELAY = process.env.ATPROTO_RELAY ?? 'https://bsky.network'; 26const CARRY_DIR = process.env.CARRY_DIR ?? '/home/nandi/.local/share/carry-vault'; 27const DB_PATH = process.env.ISLAND_DETECT_DB ?? resolve(CARRY_DIR, '.island-detect.db'); 28const COLLECTION = 'network.cosmik.connection'; 29const CACHE_TTL_MS = 60 * 60 * 1000; // 1 hour 30 31interface Component { 32 nodes: Set<string>; 33 dids: Set<string>; 34 lexmin: string; 35 islandId: string; 36} 37 38// --- SQLite --- 39 40function openDb(): Database { 41 const db = new Database(DB_PATH, { create: true }); 42 db.exec('PRAGMA journal_mode = WAL'); 43 db.exec('PRAGMA synchronous = NORMAL'); 44 db.exec(` 45 CREATE TABLE IF NOT EXISTS edges ( 46 source TEXT NOT NULL, 47 target TEXT NOT NULL, 48 relation TEXT, 49 did TEXT NOT NULL, 50 fetched_at TEXT NOT NULL, 51 PRIMARY KEY (source, target, did) 52 ); 53 CREATE TABLE IF NOT EXISTS dids ( 54 did TEXT PRIMARY KEY, 55 handle TEXT, 56 pds TEXT, 57 edge_count INTEGER DEFAULT 0, 58 fetched_at TEXT 59 ); 60 CREATE INDEX IF NOT EXISTS idx_edges_source ON edges(source); 61 CREATE INDEX IF NOT EXISTS idx_edges_target ON edges(target); 62 `); 63 return db; 64} 65 66// --- Network discovery --- 67 68async function discoverDIDs(): Promise<string[]> { 69 const dids: string[] = []; 70 let cursor = ''; 71 while (true) { 72 const url = new URL(`${RELAY}/xrpc/com.atproto.sync.listReposByCollection`); 73 url.searchParams.set('collection', COLLECTION); 74 url.searchParams.set('limit', '1000'); 75 if (cursor) url.searchParams.set('cursor', cursor); 76 const res = await fetch(url.toString()); 77 if (!res.ok) break; 78 const data = await res.json() as { repos: Array<{ did: string }>; cursor?: string }; 79 for (const r of data.repos) dids.push(r.did); 80 cursor = data.cursor ?? ''; 81 if (!cursor) break; 82 await new Promise((r) => setTimeout(r, 200)); 83 } 84 return dids; 85} 86 87async function resolvePDS(did: string): Promise<string | null> { 88 try { 89 const res = await fetch(`https://plc.directory/${did}`); 90 if (!res.ok) return null; 91 const doc = await res.json() as { service?: Array<{ id: string; serviceEndpoint: string }> }; 92 return doc.service?.find((s) => s.id === '#atproto_pds')?.serviceEndpoint ?? null; 93 } catch { return null; } 94} 95 96async function resolveHandle(did: string): Promise<string | null> { 97 try { 98 const res = await fetch(`https://plc.directory/${did}`); 99 if (!res.ok) return null; 100 const doc = await res.json() as { alsoKnownAs?: string[] }; 101 return doc.alsoKnownAs?.find((h) => h.startsWith('at://'))?.replace('at://', '') ?? null; 102 } catch { return null; } 103} 104 105async function fetchConnectionRecords(pds: string, did: string): Promise<Array<{ source: string; target: string; relation?: string; atUri: string }>> { 106 const records: Array<{ source: string; target: string; relation?: string; atUri: string }> = []; 107 let cursor = ''; 108 while (true) { 109 const url = new URL(`${pds}/xrpc/com.atproto.repo.listRecords`); 110 url.searchParams.set('repo', did); 111 url.searchParams.set('collection', COLLECTION); 112 url.searchParams.set('limit', '100'); 113 if (cursor) url.searchParams.set('cursor', cursor); 114 const res = await fetch(url.toString()); 115 if (!res.ok) break; 116 const data = await res.json() as { 117 records: Array<{ uri: string; value: { source?: string; target?: string; connectionType?: string } }>; 118 cursor?: string; 119 }; 120 for (const r of data.records) { 121 if (r.value.source && r.value.target) { 122 records.push({ source: r.value.source, target: r.value.target, relation: r.value.connectionType, atUri: r.uri }); 123 } 124 } 125 cursor = data.cursor ?? ''; 126 if (!cursor) break; 127 await new Promise((r) => setTimeout(r, 100)); 128 } 129 return records; 130} 131 132// --- Cache operations --- 133 134function isCached(db: Database, did: string): boolean { 135 const row = db.query('SELECT fetched_at FROM dids WHERE did = ?').get(did) as { fetched_at: string } | null; 136 if (!row) return false; 137 return Date.now() - new Date(row.fetched_at).getTime() < CACHE_TTL_MS; 138} 139 140function upsertDid(db: Database, did: string, handle: string | null, pds: string | null, edgeCount: number): void { 141 const now = new Date().toISOString(); 142 db.query(` 143 INSERT INTO dids (did, handle, pds, edge_count, fetched_at) 144 VALUES (?, ?, ?, ?, ?) 145 ON CONFLICT(did) DO UPDATE SET handle=?, pds=?, edge_count=?, fetched_at=? 146 `).run(did, handle, pds, edgeCount, now, handle, pds, edgeCount, now); 147} 148 149function upsertEdges(db: Database, did: string, records: Array<{ source: string; target: string; relation?: string; atUri: string }>): void { 150 const now = new Date().toISOString(); 151 // Add at_uri column if missing 152 try { db.exec('ALTER TABLE edges ADD COLUMN at_uri TEXT'); } catch {} 153 const insert = db.query(` 154 INSERT INTO edges (source, target, relation, did, at_uri, fetched_at) 155 VALUES (?, ?, ?, ?, ?, ?) 156 ON CONFLICT(source, target, did) DO UPDATE SET relation=?, at_uri=?, fetched_at=? 157 `); 158 // Delete stale edges for this DID before inserting 159 db.query('DELETE FROM edges WHERE did = ?').run(did); 160 for (const r of records) { 161 insert.run(r.source, r.target, r.relation ?? null, did, r.atUri, now, r.relation ?? null, r.atUri, now); 162 } 163} 164 165// --- Connected components --- 166 167function findComponents(db: Database): Component[] { 168 // Build adjacency list from all edges 169 const adj = new Map<string, Set<string>>(); 170 const edgeDids = new Map<string, Set<string>>(); // node -> which DIDs have edges touching it 171 172 const rows = db.query('SELECT source, target, did FROM edges').all() as Array<{ source: string; target: string; did: string }>; 173 for (const row of rows) { 174 const s = row.source; 175 const t = row.target; 176 if (!adj.has(s)) adj.set(s, new Set()); 177 if (!adj.has(t)) adj.set(t, new Set()); 178 adj.get(s)!.add(t); 179 adj.get(t)!.add(s); 180 if (!edgeDids.has(s)) edgeDids.set(s, new Set()); 181 if (!edgeDids.has(t)) edgeDids.set(t, new Set()); 182 edgeDids.get(s)!.add(row.did); 183 edgeDids.get(t)!.add(row.did); 184 } 185 186 // BFS to find components 187 const visited = new Set<string>(); 188 const components: Component[] = []; 189 190 for (const node of adj.keys()) { 191 if (visited.has(node)) continue; 192 const compNodes = new Set<string>(); 193 const compDids = new Set<string>(); 194 const queue = [node]; 195 while (queue.length > 0) { 196 const current = queue.pop()!; 197 if (visited.has(current)) continue; 198 visited.add(current); 199 compNodes.add(current); 200 for (const did of edgeDids.get(current) ?? []) compDids.add(did); 201 for (const neighbor of adj.get(current) ?? []) { 202 if (!visited.has(neighbor)) queue.push(neighbor); 203 } 204 } 205 const lexmin = [...compNodes].sort()[0]; 206 components.push({ nodes: compNodes, dids: compDids, lexmin, islandId: islandId(lexmin) }); 207 } 208 209 return components.sort((a, b) => b.nodes.size - a.nodes.size); 210} 211 212// --- Main --- 213 214async function main(): Promise<void> { 215 const args = process.argv.slice(2); 216 const wantJson = args.includes('--json'); 217 const refresh = args.includes('--refresh'); 218 const targetId = args.find((a) => !a.startsWith('--')); 219 const minSizeArg = args.find((a) => a.startsWith('--min-size=')); 220 const minSize = minSizeArg ? Number(minSizeArg.split('=')[1]) : 3; 221 222 const db = openDb(); 223 224 console.error('Discovering network DIDs...'); 225 const networkDids = await discoverDIDs(); 226 console.error(`Found ${networkDids.length} DIDs`); 227 228 // Fetch edges for each DID (use cache unless --refresh) 229 let fetched = 0; 230 let cached = 0; 231 for (let i = 0; i < networkDids.length; i++) { 232 const did = networkDids[i]; 233 if (!refresh && isCached(db, did)) { 234 cached++; 235 continue; 236 } 237 console.error(`[${i + 1}/${networkDids.length}] Fetching ${did}...`); 238 const pds = await resolvePDS(did); 239 if (!pds) { console.error(' Could not resolve PDS'); continue; } 240 const handle = await resolveHandle(did); 241 const records = await fetchConnectionRecords(pds, did); 242 upsertDid(db, did, handle, pds, records.length); 243 upsertEdges(db, did, records); 244 fetched++; 245 } 246 console.error(`Fetched: ${fetched}, Cached: ${cached}`); 247 248 // Stats 249 const totalEdges = (db.query('SELECT COUNT(*) as c FROM edges').get() as { c: number }).c; 250 const totalDids = (db.query('SELECT COUNT(*) as c FROM dids').get() as { c: number }).c; 251 const uniqueUrls = new Set<string>(); 252 const urlRows = db.query('SELECT DISTINCT source FROM edges UNION SELECT DISTINCT target FROM edges').all() as Array<{ source: string }>; 253 for (const r of urlRows) uniqueUrls.add(r.source); 254 255 console.error(`Database: ${totalEdges} edges, ${totalDids} DIDs, ${uniqueUrls.size} unique URLs`); 256 257 // Find connected components 258 const components = findComponents(db); 259 const filtered = components.filter((c) => c.nodes.size >= minSize); 260 261 // Show a specific island by ID 262 if (targetId) { 263 const comp = components.find((c) => c.islandId === targetId); 264 if (!comp) { console.error(`Island ${targetId} not found`); process.exit(1); } 265 const didPlaceholders = [...comp.dids].map(() => '?').join(','); 266 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 }>; 267 if (wantJson) { 268 console.log(JSON.stringify({ 269 id: comp.islandId, 270 lexmin: comp.lexmin, 271 nodeCount: comp.nodes.size, 272 didCount: comp.dids.size, 273 nodes: [...comp.nodes].sort(), 274 dids: didRows, 275 }, null, 2)); 276 } else { 277 console.log(`Island ${comp.islandId}: ${comp.nodes.size} nodes, ${comp.dids.size} DIDs`); 278 console.log(`Lexmin: ${comp.lexmin}`); 279 console.log('\nNodes:'); 280 for (const url of [...comp.nodes].sort()) console.log(` ${url}`); 281 console.log('\nDIDs:'); 282 for (const d of didRows) console.log(` ${d.handle ?? d.did} (${d.edge_count} edges)`); 283 } 284 db.close(); 285 return; 286 } 287 288 if (wantJson) { 289 console.log(JSON.stringify({ 290 totalEdges, 291 totalDids, 292 uniqueUrls: uniqueUrls.size, 293 componentCount: filtered.length, 294 components: filtered.map((c) => ({ 295 id: c.islandId, 296 lexmin: c.lexmin, 297 nodeCount: c.nodes.size, 298 didCount: c.dids.size, 299 sampleNodes: [...c.nodes].sort().slice(0, 10), 300 dids: [...c.dids], 301 })), 302 }, null, 2)); 303 db.close(); 304 return; 305 } 306 307 console.log(`\nIsland discover: ${totalDids} DIDs, ${totalEdges} edges, ${uniqueUrls.size} unique URLs`); 308 console.log(`Connected components (>= ${minSize} nodes): ${filtered.length}\n`); 309 310 for (const comp of filtered.slice(0, 30)) { 311 const domain = domainFromUrl(comp.lexmin); 312 console.log(`${comp.islandId} ${comp.nodes.size} nodes, ${comp.dids.size} DIDs ${domain}`); 313 } 314 315 // Show which component we're in (if we can determine our DID) 316 const ourDid = process.env.BLUESKY_IDENTIFIER 317 ? (db.query('SELECT did FROM dids WHERE handle = ? OR did LIKE ?').get(process.env.BLUESKY_IDENTIFIER, `%${process.env.BLUESKY_IDENTIFIER}%`) as { did: string } | null) 318 : null; 319 if (ourDid) { 320 const ourComp = components.find((c) => c.dids.has(ourDid.did)); 321 if (ourComp) { 322 console.log(`\nOur island: ${ourComp.islandId} (${ourComp.nodes.size} nodes, shared with ${ourComp.dids.size - 1} other DIDs)`); 323 } 324 } 325 326 db.close(); 327} 328 329main().catch((err) => { 330 console.error('Fatal:', err instanceof Error ? err.message : String(err)); 331 process.exit(1); 332});