This repository has no description
0

Configure Feed

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

nanopub / src / worker.ts
28 kB 727 lines
1/** 2 * Nanopub Appview — indexes org.latha.nanopub records from the 3 * nanopub server network, serves them at nanopub.latha.org. 4 * 5 * Stack: Cloudflare Workers + D1 (Contrail), TypeScript, Hono 6 */ 7 8import { Hono } from "hono"; 9import { cors } from "hono/cors"; 10import { createWorker } from "@atmo-dev/contrail/worker"; 11import { Contrail } from "@atmo-dev/contrail"; 12import { lexicons } from "../lexicons/generated/index.js"; 13import { config } from "./contrail.config.js"; 14import { LANDING_PAGE } from "./landing-page.js"; 15import { NANOPUB_PAGE } from "./nanopub-page.js"; 16 17const REGISTRY = "https://registry.knowledgepixels.com"; 18const PDS_ORIGIN = "https://pds.latha.org"; 19const RESEARCHER_DID = "did:plc:3kkhul7jznlb6ba7rprzawnj"; 20const STIGMERGIC = "https://stigmergic.latha.org"; 21const RESEARCHER_HANDLE = "researcher.pds.latha.org"; 22 23interface Env { 24 DB: D1Database; 25 RESEARCHER_KEY_PEM: string; 26} 27 28const contrailWorker = createWorker(config, { lexicons }); 29const contrail = new Contrail(config); 30let contrailReady = false; 31 32async function ensureContrailReady(db: D1Database): Promise<void> { 33 if (contrailReady) return; 34 await contrail.init(db); 35 contrailReady = true; 36} 37 38// --- Nanopub conversion --- 39 40interface NanopubGraph { 41 "@id": string; 42 "@graph": Array<{ 43 "@id": string; 44 [predicate: string]: Array<{ "@id"?: string; "@value"?: string; "@type"?: string }> | undefined; 45 }>; 46} 47 48function extractValue(arr: Array<{ "@value"?: string }> | undefined): string | undefined { 49 if (!arr || !arr.length) return undefined; 50 return arr[0]["@value"]; 51} 52 53function convertNanopub(graphs: NanopubGraph[], nanopubId: string) { 54 const nanopubUri = `https://w3id.org/np/${nanopubId}`; 55 const assertionGraph = graphs.find(g => g["@id"]?.endsWith("/assertion")); 56 const provenanceGraph = graphs.find(g => g["@id"]?.endsWith("/provenance")); 57 const pubInfoGraph = graphs.find(g => g["@id"]?.endsWith("/pubinfo")); 58 const headGraph = graphs.find(g => g["@id"]?.endsWith("/Head")); 59 60 const triples: Array<{ subject: string; predicate: string; object: string; objectType?: string }> = []; 61 let subjectIri: string | undefined; 62 let predicateIri: string | undefined; 63 let objectIri: string | undefined; 64 65 if (assertionGraph) { 66 for (const node of assertionGraph["@graph"]) { 67 const nodeId = node["@id"]; 68 for (const [pred, vals] of Object.entries(node)) { 69 if (pred === "@id") continue; 70 if (!Array.isArray(vals)) continue; 71 for (const v of vals) { 72 const obj = v["@id"] || v["@value"] || ""; 73 if (!obj) continue; 74 triples.push({ 75 subject: nodeId || "", 76 predicate: pred, 77 object: obj, 78 objectType: v["@id"] ? "uri" : "literal", 79 }); 80 } 81 } 82 } 83 if (triples.length > 0) { 84 subjectIri = triples[0].subject; 85 predicateIri = triples[0].predicate; 86 objectIri = triples[0].object; 87 } 88 } 89 90 const provenanceItems: Array<{ predicate: string; object: string }> = []; 91 let generatedBy: string | undefined; 92 let derivedFrom: string | undefined; 93 let authoredOn: string | undefined; 94 95 if (provenanceGraph) { 96 for (const node of provenanceGraph["@graph"]) { 97 for (const [pred, vals] of Object.entries(node)) { 98 if (pred === "@id") continue; 99 if (!Array.isArray(vals)) continue; 100 for (const v of vals) { 101 const obj = v["@id"] || v["@value"] || ""; 102 if (!obj) continue; 103 provenanceItems.push({ predicate: pred, object: obj }); 104 if ((pred.includes("wasGeneratedBy") || pred.includes("providedBy")) && v["@id"]) generatedBy = generatedBy || v["@id"]; 105 if (pred.includes("wasDerivedFrom") || pred.includes("derivedFrom")) derivedFrom = derivedFrom || (v["@id"] || v["@value"]); 106 if (pred.includes("authoredOn") || pred.includes("createdOn")) authoredOn = authoredOn || v["@value"]; 107 } 108 } 109 } 110 } 111 112 const pubInfoItems: Array<{ predicate: string; object: string }> = []; 113 const creators: string[] = []; 114 let license: string | undefined; 115 let createdOn: string | undefined; 116 let nanopubType: string | undefined; 117 let label: string | undefined; 118 119 if (pubInfoGraph) { 120 for (const node of pubInfoGraph["@graph"]) { 121 for (const [pred, vals] of Object.entries(node)) { 122 if (pred === "@id") continue; 123 if (!Array.isArray(vals)) continue; 124 for (const v of vals) { 125 const obj = v["@id"] || v["@value"] || ""; 126 if (!obj) continue; 127 pubInfoItems.push({ predicate: pred, object: obj }); 128 if (pred.includes("terms/creator") && v["@id"]) creators.push(v["@id"]); 129 if (pred.includes("terms/license") && v["@id"]) license = v["@id"]; 130 if (pred.includes("terms/created") && v["@value"]) createdOn = v["@value"]; 131 if (pred.includes("hasNanopubType") && v["@id"]) nanopubType = v["@id"]; 132 if ((pred.includes("rdf-schema#label") || pred.includes("rdfs/label")) && v["@value"]) { 133 if (!label || v["@value"].length > label.length) label = v["@value"]; 134 } 135 if (pred.includes("hasLabelFromApi") && v["@value"] && !label) label = v["@value"]; 136 } 137 } 138 } 139 } 140 141 let typeEnum: string | undefined; 142 if (nanopubType) { 143 if (nanopubType.includes("Retraction")) typeEnum = "org.latha.nanopub.defs#retraction"; 144 else if (nanopubType.includes("NewVersion")) typeEnum = "org.latha.nanopub.defs#newVersion"; 145 else typeEnum = "org.latha.nanopub.defs#assertion"; 146 } 147 148 let signerPubkey: string | undefined; 149 if (headGraph) { 150 for (const node of headGraph["@graph"]) { 151 const pkVals = node["http://purl.org/nanopub/x/hasPublicKey"]; 152 if (pkVals && pkVals.length) { signerPubkey = extractValue(pkVals as any); break; } 153 } 154 } 155 156 const record: Record<string, any> = { 157 $type: "org.latha.nanopub", 158 assertion: { triples: triples.slice(0, 100) }, 159 pubInfo: { creators: creators.slice(0, 20) }, 160 sourceNanopubUri: nanopubUri, 161 trustUri: nanopubUri, 162 createdAt: new Date().toISOString(), 163 }; 164 if (label) record.assertion.label = label; 165 if (subjectIri) record.assertion.subjectIri = subjectIri; 166 if (predicateIri) record.assertion.predicateIri = predicateIri; 167 if (objectIri) record.assertion.objectIri = objectIri; 168 if (typeEnum) record.nanopubType = typeEnum; 169 if (signerPubkey) record.signerPubkey = signerPubkey; 170 if (license) record.pubInfo.license = license; 171 if (createdOn) record.pubInfo.createdOn = createdOn; 172 if (provenanceItems.length) record.provenance = { items: provenanceItems.slice(0, 50) }; 173 if (generatedBy) { record.provenance = record.provenance || {}; record.provenance.generatedBy = generatedBy; } 174 if (derivedFrom) { record.provenance = record.provenance || {}; record.provenance.derivedFrom = derivedFrom; } 175 if (authoredOn) { record.provenance = record.provenance || {}; record.provenance.authoredOn = authoredOn; } 176 177 return record; 178} 179 180function generateRkey(nanopubId: string): string { 181 // Use a simple hash — crypto.subtle is async, so we use a deterministic approach 182 let hash = 0; 183 for (let i = 0; i < nanopubId.length; i++) { 184 const c = nanopubId.charCodeAt(i); 185 hash = ((hash << 5) - hash + c) | 0; 186 } 187 // Mix in more bits for better distribution 188 const buf = new ArrayBuffer(8); 189 const view = new DataView(buf); 190 view.setFloat64(0, Math.abs(hash) * 2654435761); 191 const bytes = new Uint8Array(buf); 192 const chars = "0123456789abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ_-"; 193 let rkey = ""; 194 for (let i = 0; i < 13; i++) { 195 rkey += chars[bytes[i] % chars.length]; 196 } 197 return rkey; 198} 199 200async function islandHash(centroid: string): Promise<string> { 201 const data = new TextEncoder().encode(centroid); 202 const hashBuffer = await crypto.subtle.digest("SHA-256", data); 203 return Array.from(new Uint8Array(hashBuffer)).map(b => b.toString(16).padStart(2, "0")).join("").slice(0, 12); 204} 205 206// --- PDS write via DPoP (Web Crypto) --- 207 208function pemToArrayBuffer(pem: string): ArrayBuffer { 209 const b64 = pem.replace(/-----BEGIN.*?-----/g, "").replace(/-----END.*?-----/g, "").replace(/\s/g, ""); 210 const bin = atob(b64); 211 const buf = new Uint8Array(bin.length); 212 for (let i = 0; i < bin.length; i++) buf[i] = bin.charCodeAt(i); 213 return buf.buffer; 214} 215 216function arrayBufferToBase64url(buf: ArrayBuffer): string { 217 const bytes = new Uint8Array(buf); 218 let binary = ""; 219 for (let i = 0; i < bytes.byteLength; i++) binary += String.fromCharCode(bytes[i]); 220 return btoa(binary).replace(/\+/g, "-").replace(/\//g, "_").replace(/=+$/, ""); 221} 222 223function stringToBase64url(str: string): string { 224 return btoa(str).replace(/\+/g, "-").replace(/\//g, "_").replace(/=+$/, ""); 225} 226 227async function importPrivateKey(pem: string): Promise<CryptoKey> { 228 return crypto.subtle.importKey( 229 "pkcs8", 230 pemToArrayBuffer(pem), 231 { name: "RSASSA-PKCS1-v1_5", hash: "SHA-256" }, 232 true, // exportable so we can get the JWK 233 ["sign"] 234 ); 235} 236 237async function computeThumbprint(jwk: { kty: string; n: string; e: string }): Promise<string> { 238 const canonical = JSON.stringify({ e: jwk.e, kty: "RSA", n: jwk.n }); 239 const hash = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(canonical)); 240 return arrayBufferToBase64url(hash); 241} 242 243async function createJwt(header: object, payload: object, privateKey: CryptoKey): Promise<string> { 244 const encHeader = stringToBase64url(JSON.stringify(header)); 245 const encPayload = stringToBase64url(JSON.stringify(payload)); 246 const signingInput = `${encHeader}.${encPayload}`; 247 const sig = await crypto.subtle.sign("RSASSA-PKCS1-v1_5", privateKey, new TextEncoder().encode(signingInput)); 248 return `${signingInput}.${arrayBufferToBase64url(sig)}`; 249} 250 251async function createDpopProof(method: string, url: string, accessToken: string, publicJwk: JsonWebKey, privateKey: CryptoKey): Promise<string> { 252 const athBuf = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(accessToken)); 253 const ath = arrayBufferToBase64url(athBuf); 254 return createJwt( 255 { typ: "dpop+jwt", alg: "RS256", jwk: publicJwk }, 256 { jti: crypto.randomUUID(), htm: method, htu: url, iat: Math.floor(Date.now() / 1000), ath }, 257 privateKey 258 ); 259} 260 261interface DpopSession { 262 did: string; 263 handle: string; 264 wmJwt: string; 265 privateKey: CryptoKey; 266 publicJwk: JsonWebKey; 267} 268 269async function getPublicJwkFromPrivate(pem: string): Promise<JsonWebKey> { 270 const privateKey = await crypto.subtle.importKey( 271 "pkcs8", 272 pemToArrayBuffer(pem), 273 { name: "RSASSA-PKCS1-v1_5", hash: "SHA-256" }, 274 true, 275 ["sign"] 276 ); 277 const jwk = await crypto.subtle.exportKey("jwk", privateKey); 278 return { kty: jwk.kty!, n: jwk.n!, e: jwk.e! }; 279} 280 281async function createDpopSession(pem: string): Promise<DpopSession> { 282 const privateKey = await importPrivateKey(pem); 283 const publicJwk = await getPublicJwkFromPrivate(pem); 284 const thumbprint = await computeThumbprint(publicJwk as any); 285 286 const tosRes = await fetch(`${PDS_ORIGIN}/tos`); 287 if (!tosRes.ok) throw new Error(`Failed to fetch ToS: ${tosRes.status}`); 288 const tosText = await tosRes.text(); 289 const tosHashBuf = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(tosText)); 290 const tosHash = arrayBufferToBase64url(tosHashBuf); 291 292 const wmJwt = await createJwt( 293 { typ: "wm+jwt", alg: "RS256" }, 294 { tos_hash: tosHash, aud: PDS_ORIGIN, cnf: { jkt: thumbprint }, iat: Math.floor(Date.now() / 1000) }, 295 privateKey 296 ); 297 298 const sessionUrl = `${PDS_ORIGIN}/xrpc/com.atproto.server.createSession`; 299 const dpop = await createDpopProof("POST", sessionUrl, wmJwt, publicJwk, privateKey); 300 const res = await fetch(sessionUrl, { 301 method: "POST", 302 headers: { "Content-Type": "application/json", Authorization: `DPoP ${wmJwt}`, DPoP: dpop }, 303 body: JSON.stringify({}), 304 }); 305 if (!res.ok) throw new Error(`createSession failed (${res.status}): ${await res.text()}`); 306 const data = await res.json() as { did: string; handle: string }; 307 return { did: data.did, handle: data.handle, wmJwt, privateKey, publicJwk }; 308} 309 310async function dpopPutRecord(session: DpopSession, collection: string, rkey: string, record: any): Promise<{ ok: boolean; uri?: string }> { 311 const url = `${PDS_ORIGIN}/xrpc/com.atproto.repo.putRecord`; 312 const dpop = await createDpopProof("POST", url, session.wmJwt, session.publicJwk, session.privateKey); 313 const res = await fetch(url, { 314 method: "POST", 315 headers: { "Content-Type": "application/json", Authorization: `DPoP ${session.wmJwt}`, DPoP: dpop }, 316 body: JSON.stringify({ repo: session.did, collection, rkey, record }), 317 }); 318 if (!res.ok) { console.error(`putRecord failed: ${res.status} ${await res.text()}`); return { ok: false }; } 319 const result = await res.json() as { uri: string }; 320 return { ok: true, uri: result.uri }; 321} 322 323async function dpopCreateRecord(session: DpopSession, collection: string, record: any): Promise<{ ok: boolean; uri?: string }> { 324 const url = `${PDS_ORIGIN}/xrpc/com.atproto.repo.createRecord`; 325 const dpop = await createDpopProof("POST", url, session.wmJwt, session.publicJwk, session.privateKey); 326 const res = await fetch(url, { 327 method: "POST", 328 headers: { "Content-Type": "application/json", Authorization: `DPoP ${session.wmJwt}`, DPoP: dpop }, 329 body: JSON.stringify({ repo: session.did, collection, record }), 330 }); 331 if (!res.ok) { console.error(`createRecord failed: ${res.status}`); return { ok: false }; } 332 const result = await res.json() as { uri: string }; 333 return { ok: true, uri: result.uri }; 334} 335 336// --- Cron: poll registry, ingest new nanopubs, create islands --- 337 338async function runCron(db: D1Database, keyPem: string): Promise<void> { 339 await ensureContrailReady(db); 340 341 // Ensure state tables 342 await db.prepare("CREATE TABLE IF NOT EXISTS cron_state (key TEXT PRIMARY KEY, value TEXT NOT NULL)").run(); 343 await db.prepare("CREATE TABLE IF NOT EXISTS island_map (nanopub_rkey TEXT PRIMARY KEY, island_rkey TEXT NOT NULL, created_at INTEGER NOT NULL)").run(); 344 345 // Create DPoP session 346 const session = await createDpopSession(keyPem); 347 348 // Get last known count and cursor 349 const lastCount = (await db.prepare("SELECT value FROM cron_state WHERE key = 'last_nanopub_count'").first<{ value: string }>())?.value || "0"; 350 const lastFirstId = (await db.prepare("SELECT value FROM cron_state WHERE key = 'last_first_id'").first<{ value: string }>())?.value || ""; 351 352 // Check current count via HEAD 353 const headRes = await fetch(`${REGISTRY}/nanopubs.json`, { method: "HEAD" }); 354 const currentCount = headRes.headers.get("nanopub-registry-nanopub-count") || ""; 355 356 if (currentCount === lastCount) { 357 console.log(`No new nanopubs (count unchanged: ${currentCount})`); 358 return; 359 } 360 361 console.log(`Nanopub count changed: ${lastCount}${currentCount}`); 362 363 // Fetch the list 364 const listRes = await fetch(`${REGISTRY}/nanopubs.json`); 365 if (!listRes.ok) { 366 console.error(`Failed to fetch nanopub list: ${listRes.status}`); 367 return; 368 } 369 const allIds: string[] = await listRes.json() as string[]; 370 371 // Find new IDs by comparing with cursor 372 let newIds: string[] = []; 373 if (lastFirstId) { 374 const lastIdx = allIds.indexOf(lastFirstId); 375 if (lastIdx > 0) { 376 newIds = allIds.slice(0, lastIdx); 377 } else { 378 newIds = allIds.slice(0, 20); 379 } 380 } else { 381 newIds = allIds.slice(0, 20); 382 } 383 384 // Cap at 20 per cron tick 385 newIds = newIds.slice(0, 20); 386 387 if (newIds.length === 0) { 388 console.log("No new IDs to process"); 389 await db.prepare("INSERT OR REPLACE INTO cron_state (key, value) VALUES (?, ?)").bind("last_nanopub_count", currentCount).run(); 390 return; 391 } 392 393 console.log(`Processing ${newIds.length} new nanopubs`); 394 395 // Get existing rkeys to skip duplicates 396 const existingRkeys = new Set<string>(); 397 for (const id of newIds) { 398 const rkey = generateRkey(id); 399 const row = await db.prepare("SELECT rkey FROM records_nanopub WHERE rkey = ?").bind(rkey).first(); 400 if (row) existingRkeys.add(rkey); 401 } 402 403 let ingested = 0; 404 let islandsCreated = 0; 405 406 for (const nanopubId of newIds) { 407 const rkey = generateRkey(nanopubId); 408 if (existingRkeys.has(rkey)) continue; 409 410 try { 411 const npRes = await fetch(`${REGISTRY}/np/${nanopubId}`, { 412 headers: { Accept: "application/ld+json" }, 413 signal: AbortSignal.timeout(10000), 414 }); 415 if (!npRes.ok) { 416 console.error(`SKIP ${nanopubId.slice(0, 10)}: HTTP ${npRes.status}`); 417 continue; 418 } 419 420 const graphs: NanopubGraph[] = await npRes.json() as NanopubGraph[]; 421 const record = convertNanopub(graphs, nanopubId); 422 423 // Write to PDS via DPoP 424 const putResult = await dpopPutRecord(session, "org.latha.nanopub", rkey, record); 425 if (!putResult.ok) { 426 console.error(`PDS write failed for ${nanopubId.slice(0, 10)}`); 427 continue; 428 } 429 430 // Insert into D1 431 const now = Date.now() * 1000; 432 await db.prepare( 433 "INSERT OR REPLACE INTO records_nanopub (uri, did, rkey, cid, record, time_us, indexed_at) VALUES (?, ?, ?, ?, ?, ?, ?)" 434 ).bind(putResult.uri!, RESEARCHER_DID, rkey, "", JSON.stringify(record), now, now).run(); 435 436 ingested++; 437 438 // Create connections + island from URI triples 439 const uriTriples = (record.assertion?.triples || []).filter( 440 (t: any) => t.objectType === "uri" && t.object?.startsWith("https://") 441 ); 442 443 if (uriTriples.length >= 2) { 444 try { 445 const connectionUris: string[] = []; 446 for (const t of uriTriples) { 447 const connRecord = { 448 $type: "network.cosmik.connection", 449 source: t.subject, 450 target: t.object, 451 connectionType: t.predicate.split("/").pop() || t.predicate, 452 note: `${record.assertion?.label || "Nanopub assertion"}`, 453 createdAt: new Date().toISOString(), 454 }; 455 const connResult = await dpopCreateRecord(session, "network.cosmik.connection", connRecord); 456 if (connResult.ok && connResult.uri) connectionUris.push(connResult.uri); 457 } 458 459 // Compute centroid 460 const degreeCounts = new Map<string, number>(); 461 for (const t of uriTriples) { 462 degreeCounts.set(t.subject, (degreeCounts.get(t.subject) || 0) + 1); 463 degreeCounts.set(t.object, (degreeCounts.get(t.object) || 0) + 1); 464 } 465 const centroid = [...degreeCounts.entries()].sort((a, b) => b[1] - a[1])[0]?.[0] || uriTriples[0].subject; 466 const iHash = await islandHash(centroid); 467 468 // Build island record 469 const label = record.assertion?.label || record.assertion?.subjectIri || "Untitled"; 470 const creators = (record.pubInfo?.creators || []).map((c: string) => { 471 const m = c.match(/orcid\.org\/([\d-]+)/); 472 return m ? m[1] : c.split("/").pop(); 473 }).join(", "); 474 475 const islandRecord = { 476 $type: "org.latha.island", 477 source: { uri: centroid, collection: "network.cosmik.connection" }, 478 connections: connectionUris.map(uri => ({ uri })), 479 analysis: { 480 title: label, 481 themes: [label.split(" ").slice(0, 3).join(" "), "nanopublication", "semantic graph"], 482 highlights: [], 483 tensions: [], 484 openQuestions: [], 485 synthesis: `Nanopublication asserting ${uriTriples.length} semantic relationships. ${label}. Created by ${creators || "unknown"}. Source: ${record.sourceNanopubUri}`, 486 }, 487 createdAt: new Date().toISOString(), 488 }; 489 490 // Publish island 491 const islandResult = await dpopPutRecord(session, "org.latha.island", iHash, islandRecord); 492 if (islandResult.ok) { 493 await db.prepare( 494 "INSERT OR REPLACE INTO island_map (nanopub_rkey, island_rkey, created_at) VALUES (?, ?, ?)" 495 ).bind(rkey, iHash, Date.now()).run(); 496 islandsCreated++; 497 } 498 499 // Notify stigmergic 500 try { 501 await fetch(`${STIGMERGIC}/xrpc/org.latha.island.notifyOfUpdate`, { 502 method: "POST", 503 headers: { "Content-Type": "application/json" }, 504 body: JSON.stringify({ did: RESEARCHER_DID }), 505 }); 506 } catch {} 507 } catch (err: any) { 508 console.error(`Island creation failed for ${rkey}: ${err?.message}`); 509 } 510 } 511 512 await new Promise(r => setTimeout(r, 100)); 513 514 } catch (err: any) { 515 console.error(`Error processing ${nanopubId.slice(0, 10)}: ${err?.message}`); 516 } 517 } 518 519 // Update cursor 520 await db.prepare("INSERT OR REPLACE INTO cron_state (key, value) VALUES (?, ?)").bind("last_nanopub_count", currentCount).run(); 521 await db.prepare("INSERT OR REPLACE INTO cron_state (key, value) VALUES (?, ?)").bind("last_first_id", allIds[0] || "").run(); 522 523 console.log(`Cron done: ${ingested} ingested, ${islandsCreated} islands created`); 524} 525 526// --- HTTP app --- 527 528let app: Hono<{ Bindings: Env }> | null = null; 529 530function buildApp(env: Env): Hono<{ Bindings: Env }> { 531 const app = new Hono<{ Bindings: Env }>(); 532 const db = env.DB; 533 534 app.use("*", cors()); 535 536 // Landing page 537 app.get("/", async (c) => { 538 return c.html(LANDING_PAGE); 539 }); 540 541 // Nanopub detail page 542 app.get("/nanopub/:rkey", async (c) => { 543 const rkey = c.req.param("rkey"); 544 await ensureContrailReady(db); 545 546 let nanopub: any = null; 547 let islandRkey: string | null = null; 548 try { 549 const row = await db 550 .prepare("SELECT uri, record, did, rkey FROM records_nanopub WHERE rkey = ? LIMIT 1") 551 .bind(rkey) 552 .first<{ uri: string; record: string; did: string; rkey: string }>(); 553 if (row) { 554 nanopub = { uri: row.uri, rkey: row.rkey, ...JSON.parse(row.record) }; 555 } 556 // Check island mapping 557 const islandRow = await db.prepare("SELECT island_rkey FROM island_map WHERE nanopub_rkey = ?").bind(rkey).first<{ island_rkey: string }>(); 558 if (islandRow) islandRkey = islandRow.island_rkey; 559 } catch {} 560 561 if (!nanopub) { 562 return c.html("<h1>Not Found</h1><p><a href='/'>Back</a></p>", 404); 563 } 564 565 // Inject island link into nanopub data 566 const dataWithIsland = { ...nanopub, _islandRkey: islandRkey, _islandUrl: islandRkey ? `${STIGMERGIC}/island/${RESEARCHER_HANDLE}/${islandRkey}` : null }; 567 568 const page = NANOPUB_PAGE.replace( 569 "</head>", 570 `<script>window.__NANOPUB__=${JSON.stringify(dataWithIsland)};</script></head>` 571 ); 572 return c.html(page); 573 }); 574 575 // API: list nanopubs sorted by publish time 576 app.get("/api/nanopubs", async (c) => { 577 await ensureContrailReady(db); 578 const limit = Math.min(Number(c.req.query("limit") || 50), 200); 579 const cursor = c.req.query("cursor"); 580 581 // Sort by nanopub publish time (pubInfo.createdOn) extracted from record 582 let rows: D1Result<any>; 583 if (cursor) { 584 rows = await db 585 .prepare( 586 "SELECT uri, record, did, rkey, time_us FROM records_nanopub WHERE json_extract(record, '$.pubInfo.createdOn') < ? ORDER BY json_extract(record, '$.pubInfo.createdOn') DESC LIMIT ?" 587 ) 588 .bind(cursor, limit) 589 .all(); 590 } else { 591 rows = await db 592 .prepare( 593 "SELECT uri, record, did, rkey, time_us FROM records_nanopub ORDER BY json_extract(record, '$.pubInfo.createdOn') DESC LIMIT ?" 594 ) 595 .bind(limit) 596 .all(); 597 } 598 599 const nanopubs = (rows.results || []).map((r: any) => { 600 try { 601 const value = JSON.parse(r.record); 602 const { assertion, pubInfo, nanopubType, sourceNanopubUri, signerPubkey, createdAt, trustUri, retracts, supersedes } = value; 603 // Check island mapping 604 let islandUrl: string | null = null; 605 // We'll add island info separately after the loop 606 return { 607 uri: r.uri, 608 rkey: r.rkey, 609 assertion: assertion ? { label: assertion.label, subjectIri: assertion.subjectIri, predicateIri: assertion.predicateIri, objectIri: assertion.objectIri } : undefined, 610 pubInfo: pubInfo ? { creators: pubInfo.creators, createdOn: pubInfo.createdOn } : undefined, 611 nanopubType, 612 sourceNanopubUri, 613 signerPubkey, 614 createdAt, 615 trustUri, 616 retracts, 617 supersedes, 618 _rkeyForIsland: r.rkey, 619 }; 620 } catch { 621 return null; 622 } 623 }).filter(Boolean); 624 625 // Batch-fetch island mappings 626 const rkeys = nanopubs.map((n: any) => n._rkeyForIsland).filter(Boolean); 627 if (rkeys.length > 0) { 628 const placeholders = rkeys.map(() => "?").join(","); 629 const islandRows = await db.prepare( 630 `SELECT nanopub_rkey, island_rkey FROM island_map WHERE nanopub_rkey IN (${placeholders})` 631 ).bind(...rkeys).all(); 632 const islandMap = new Map<string, string>(); 633 for (const row of (islandRows.results || [])) { 634 islandMap.set((row as any).nanopub_rkey, (row as any).island_rkey); 635 } 636 for (const n of nanopubs) { 637 const rkey = (n as any)._rkeyForIsland; 638 const islandRkey = islandMap.get(rkey); 639 if (islandRkey) { 640 (n as any).islandUrl = `${STIGMERGIC}/island/${RESEARCHER_HANDLE}/${islandRkey}`; 641 } 642 delete (n as any)._rkeyForIsland; 643 } 644 } 645 646 const lastRow = rows.results?.[rows.results.length - 1]; 647 const nextCursor = lastRow ? JSON.parse((lastRow as any).record)?.pubInfo?.createdOn : undefined; 648 649 return c.json({ nanopubs, cursor: nextCursor }); 650 }); 651 652 // API: health check 653 app.get("/api/health", async (c) => { 654 await ensureContrailReady(db); 655 let count = 0; 656 try { 657 const row = await db.prepare("SELECT COUNT(*) as cnt FROM records_nanopub").first<{ cnt: number }>(); 658 count = row?.cnt || 0; 659 } catch {} 660 return c.json({ ok: true, nanopubs: count }); 661 }); 662 663 // Manual cron trigger 664 app.post("/api/cron", async (c) => { 665 const keyPem = c.env.RESEARCHER_KEY_PEM; 666 if (!keyPem) return c.json({ error: "RESEARCHER_KEY_PEM not set" }, 500); 667 try { 668 await runCron(c.env.DB, keyPem); 669 return c.json({ ok: true }); 670 } catch (err: any) { 671 return c.json({ error: err?.message }, 500); 672 } 673 }); 674 675 // Crawl endpoint 676 app.post("/xrpc/com.atproto.sync.requestCrawl", async (c) => { 677 let body: { did?: string }; 678 try { body = await c.req.json(); } catch { body = {}; } 679 const did = body.did; 680 if (!did) return c.json({ error: "BadRequest", message: "Provide did" }, 400); 681 682 await ensureContrailReady(db); 683 try { 684 await db.prepare("INSERT OR IGNORE INTO identities (did, handle, pds, resolved_at) VALUES (?, ?, ?, ?)") 685 .bind(did, "", "pds.latha.org", Date.now() * 1000).run(); 686 687 const pdsUrl = `https://pds.latha.org/xrpc/com.atproto.repo.listRecords?repo=${did}&collection=org.latha.nanopub&limit=100`; 688 const res = await fetch(pdsUrl); 689 if (res.ok) { 690 const data = await res.json() as { records: Array<{ uri: string; cid: string; value: any }> }; 691 for (const rec of data.records) { 692 const rkey = rec.uri.split("/").pop() || ""; 693 const now = Date.now() * 1000; 694 await db.prepare( 695 "INSERT OR REPLACE INTO records_nanopub (uri, did, rkey, cid, record, time_us, indexed_at) VALUES (?, ?, ?, ?, ?, ?, ?)" 696 ).bind(rec.uri, did, rkey, rec.cid, JSON.stringify(rec.value), now, now).run(); 697 } 698 } 699 } catch (err: any) { 700 console.error(`Crawl error: ${err?.message ?? err}`); 701 } 702 return c.json({ ok: true, did }); 703 }); 704 705 // Contrail passthrough 706 app.all("*", async (c) => { 707 const response = await contrailWorker.fetch(c.req.raw, c.env as unknown as Record<string, unknown>); 708 return response; 709 }); 710 711 return app; 712} 713 714export default { 715 fetch(request: Request, env: Env): Response | Promise<Response> { 716 app ??= buildApp(env); 717 return app.fetch(request, env); 718 }, 719 async scheduled(event: ScheduledEvent, env: Env, ctx: ExecutionContext) { 720 const keyPem = env.RESEARCHER_KEY_PEM; 721 if (!keyPem) { 722 console.error("RESEARCHER_KEY_PEM not set — skipping cron"); 723 return; 724 } 725 ctx.waitUntil(runCron(env.DB, keyPem)); 726 }, 727};