/** * Nanopub Appview — indexes org.latha.nanopub records from the * nanopub server network, serves them at nanopub.latha.org. * * Stack: Cloudflare Workers + D1 (Contrail), TypeScript, Hono */ import { Hono } from "hono"; import { cors } from "hono/cors"; import { createWorker } from "@atmo-dev/contrail/worker"; import { Contrail } from "@atmo-dev/contrail"; import { lexicons } from "../lexicons/generated/index.js"; import { config } from "./contrail.config.js"; import { LANDING_PAGE } from "./landing-page.js"; import { NANOPUB_PAGE } from "./nanopub-page.js"; const REGISTRY = "https://registry.knowledgepixels.com"; const PDS_ORIGIN = "https://pds.latha.org"; const RESEARCHER_DID = "did:plc:3kkhul7jznlb6ba7rprzawnj"; const STIGMERGIC = "https://stigmergic.latha.org"; const RESEARCHER_HANDLE = "researcher.pds.latha.org"; interface Env { DB: D1Database; RESEARCHER_KEY_PEM: string; } const contrailWorker = createWorker(config, { lexicons }); const contrail = new Contrail(config); let contrailReady = false; async function ensureContrailReady(db: D1Database): Promise { if (contrailReady) return; await contrail.init(db); contrailReady = true; } // --- Nanopub conversion --- interface NanopubGraph { "@id": string; "@graph": Array<{ "@id": string; [predicate: string]: Array<{ "@id"?: string; "@value"?: string; "@type"?: string }> | undefined; }>; } function extractValue(arr: Array<{ "@value"?: string }> | undefined): string | undefined { if (!arr || !arr.length) return undefined; return arr[0]["@value"]; } function convertNanopub(graphs: NanopubGraph[], nanopubId: string) { const nanopubUri = `https://w3id.org/np/${nanopubId}`; const assertionGraph = graphs.find(g => g["@id"]?.endsWith("/assertion")); const provenanceGraph = graphs.find(g => g["@id"]?.endsWith("/provenance")); const pubInfoGraph = graphs.find(g => g["@id"]?.endsWith("/pubinfo")); const headGraph = graphs.find(g => g["@id"]?.endsWith("/Head")); const triples: Array<{ subject: string; predicate: string; object: string; objectType?: string }> = []; let subjectIri: string | undefined; let predicateIri: string | undefined; let objectIri: string | undefined; if (assertionGraph) { for (const node of assertionGraph["@graph"]) { const nodeId = node["@id"]; for (const [pred, vals] of Object.entries(node)) { if (pred === "@id") continue; if (!Array.isArray(vals)) continue; for (const v of vals) { const obj = v["@id"] || v["@value"] || ""; if (!obj) continue; triples.push({ subject: nodeId || "", predicate: pred, object: obj, objectType: v["@id"] ? "uri" : "literal", }); } } } if (triples.length > 0) { subjectIri = triples[0].subject; predicateIri = triples[0].predicate; objectIri = triples[0].object; } } const provenanceItems: Array<{ predicate: string; object: string }> = []; let generatedBy: string | undefined; let derivedFrom: string | undefined; let authoredOn: string | undefined; if (provenanceGraph) { for (const node of provenanceGraph["@graph"]) { for (const [pred, vals] of Object.entries(node)) { if (pred === "@id") continue; if (!Array.isArray(vals)) continue; for (const v of vals) { const obj = v["@id"] || v["@value"] || ""; if (!obj) continue; provenanceItems.push({ predicate: pred, object: obj }); if ((pred.includes("wasGeneratedBy") || pred.includes("providedBy")) && v["@id"]) generatedBy = generatedBy || v["@id"]; if (pred.includes("wasDerivedFrom") || pred.includes("derivedFrom")) derivedFrom = derivedFrom || (v["@id"] || v["@value"]); if (pred.includes("authoredOn") || pred.includes("createdOn")) authoredOn = authoredOn || v["@value"]; } } } } const pubInfoItems: Array<{ predicate: string; object: string }> = []; const creators: string[] = []; let license: string | undefined; let createdOn: string | undefined; let nanopubType: string | undefined; let label: string | undefined; if (pubInfoGraph) { for (const node of pubInfoGraph["@graph"]) { for (const [pred, vals] of Object.entries(node)) { if (pred === "@id") continue; if (!Array.isArray(vals)) continue; for (const v of vals) { const obj = v["@id"] || v["@value"] || ""; if (!obj) continue; pubInfoItems.push({ predicate: pred, object: obj }); if (pred.includes("terms/creator") && v["@id"]) creators.push(v["@id"]); if (pred.includes("terms/license") && v["@id"]) license = v["@id"]; if (pred.includes("terms/created") && v["@value"]) createdOn = v["@value"]; if (pred.includes("hasNanopubType") && v["@id"]) nanopubType = v["@id"]; if ((pred.includes("rdf-schema#label") || pred.includes("rdfs/label")) && v["@value"]) { if (!label || v["@value"].length > label.length) label = v["@value"]; } if (pred.includes("hasLabelFromApi") && v["@value"] && !label) label = v["@value"]; } } } } let typeEnum: string | undefined; if (nanopubType) { if (nanopubType.includes("Retraction")) typeEnum = "org.latha.nanopub.defs#retraction"; else if (nanopubType.includes("NewVersion")) typeEnum = "org.latha.nanopub.defs#newVersion"; else typeEnum = "org.latha.nanopub.defs#assertion"; } let signerPubkey: string | undefined; if (headGraph) { for (const node of headGraph["@graph"]) { const pkVals = node["http://purl.org/nanopub/x/hasPublicKey"]; if (pkVals && pkVals.length) { signerPubkey = extractValue(pkVals as any); break; } } } const record: Record = { $type: "org.latha.nanopub", assertion: { triples: triples.slice(0, 100) }, pubInfo: { creators: creators.slice(0, 20) }, sourceNanopubUri: nanopubUri, trustUri: nanopubUri, createdAt: new Date().toISOString(), }; if (label) record.assertion.label = label; if (subjectIri) record.assertion.subjectIri = subjectIri; if (predicateIri) record.assertion.predicateIri = predicateIri; if (objectIri) record.assertion.objectIri = objectIri; if (typeEnum) record.nanopubType = typeEnum; if (signerPubkey) record.signerPubkey = signerPubkey; if (license) record.pubInfo.license = license; if (createdOn) record.pubInfo.createdOn = createdOn; if (provenanceItems.length) record.provenance = { items: provenanceItems.slice(0, 50) }; if (generatedBy) { record.provenance = record.provenance || {}; record.provenance.generatedBy = generatedBy; } if (derivedFrom) { record.provenance = record.provenance || {}; record.provenance.derivedFrom = derivedFrom; } if (authoredOn) { record.provenance = record.provenance || {}; record.provenance.authoredOn = authoredOn; } return record; } function generateRkey(nanopubId: string): string { // Use a simple hash — crypto.subtle is async, so we use a deterministic approach let hash = 0; for (let i = 0; i < nanopubId.length; i++) { const c = nanopubId.charCodeAt(i); hash = ((hash << 5) - hash + c) | 0; } // Mix in more bits for better distribution const buf = new ArrayBuffer(8); const view = new DataView(buf); view.setFloat64(0, Math.abs(hash) * 2654435761); const bytes = new Uint8Array(buf); const chars = "0123456789abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ_-"; let rkey = ""; for (let i = 0; i < 13; i++) { rkey += chars[bytes[i] % chars.length]; } return rkey; } async function islandHash(centroid: string): Promise { const data = new TextEncoder().encode(centroid); const hashBuffer = await crypto.subtle.digest("SHA-256", data); return Array.from(new Uint8Array(hashBuffer)).map(b => b.toString(16).padStart(2, "0")).join("").slice(0, 12); } // --- PDS write via DPoP (Web Crypto) --- function pemToArrayBuffer(pem: string): ArrayBuffer { const b64 = pem.replace(/-----BEGIN.*?-----/g, "").replace(/-----END.*?-----/g, "").replace(/\s/g, ""); const bin = atob(b64); const buf = new Uint8Array(bin.length); for (let i = 0; i < bin.length; i++) buf[i] = bin.charCodeAt(i); return buf.buffer; } function arrayBufferToBase64url(buf: ArrayBuffer): string { const bytes = new Uint8Array(buf); let binary = ""; for (let i = 0; i < bytes.byteLength; i++) binary += String.fromCharCode(bytes[i]); return btoa(binary).replace(/\+/g, "-").replace(/\//g, "_").replace(/=+$/, ""); } function stringToBase64url(str: string): string { return btoa(str).replace(/\+/g, "-").replace(/\//g, "_").replace(/=+$/, ""); } async function importPrivateKey(pem: string): Promise { return crypto.subtle.importKey( "pkcs8", pemToArrayBuffer(pem), { name: "RSASSA-PKCS1-v1_5", hash: "SHA-256" }, true, // exportable so we can get the JWK ["sign"] ); } async function computeThumbprint(jwk: { kty: string; n: string; e: string }): Promise { const canonical = JSON.stringify({ e: jwk.e, kty: "RSA", n: jwk.n }); const hash = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(canonical)); return arrayBufferToBase64url(hash); } async function createJwt(header: object, payload: object, privateKey: CryptoKey): Promise { const encHeader = stringToBase64url(JSON.stringify(header)); const encPayload = stringToBase64url(JSON.stringify(payload)); const signingInput = `${encHeader}.${encPayload}`; const sig = await crypto.subtle.sign("RSASSA-PKCS1-v1_5", privateKey, new TextEncoder().encode(signingInput)); return `${signingInput}.${arrayBufferToBase64url(sig)}`; } async function createDpopProof(method: string, url: string, accessToken: string, publicJwk: JsonWebKey, privateKey: CryptoKey): Promise { const athBuf = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(accessToken)); const ath = arrayBufferToBase64url(athBuf); return createJwt( { typ: "dpop+jwt", alg: "RS256", jwk: publicJwk }, { jti: crypto.randomUUID(), htm: method, htu: url, iat: Math.floor(Date.now() / 1000), ath }, privateKey ); } interface DpopSession { did: string; handle: string; wmJwt: string; privateKey: CryptoKey; publicJwk: JsonWebKey; } async function getPublicJwkFromPrivate(pem: string): Promise { const privateKey = await crypto.subtle.importKey( "pkcs8", pemToArrayBuffer(pem), { name: "RSASSA-PKCS1-v1_5", hash: "SHA-256" }, true, ["sign"] ); const jwk = await crypto.subtle.exportKey("jwk", privateKey); return { kty: jwk.kty!, n: jwk.n!, e: jwk.e! }; } async function createDpopSession(pem: string): Promise { const privateKey = await importPrivateKey(pem); const publicJwk = await getPublicJwkFromPrivate(pem); const thumbprint = await computeThumbprint(publicJwk as any); const tosRes = await fetch(`${PDS_ORIGIN}/tos`); if (!tosRes.ok) throw new Error(`Failed to fetch ToS: ${tosRes.status}`); const tosText = await tosRes.text(); const tosHashBuf = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(tosText)); const tosHash = arrayBufferToBase64url(tosHashBuf); const wmJwt = await createJwt( { typ: "wm+jwt", alg: "RS256" }, { tos_hash: tosHash, aud: PDS_ORIGIN, cnf: { jkt: thumbprint }, iat: Math.floor(Date.now() / 1000) }, privateKey ); const sessionUrl = `${PDS_ORIGIN}/xrpc/com.atproto.server.createSession`; const dpop = await createDpopProof("POST", sessionUrl, wmJwt, publicJwk, privateKey); const res = await fetch(sessionUrl, { method: "POST", headers: { "Content-Type": "application/json", Authorization: `DPoP ${wmJwt}`, DPoP: dpop }, body: JSON.stringify({}), }); if (!res.ok) throw new Error(`createSession failed (${res.status}): ${await res.text()}`); const data = await res.json() as { did: string; handle: string }; return { did: data.did, handle: data.handle, wmJwt, privateKey, publicJwk }; } async function dpopPutRecord(session: DpopSession, collection: string, rkey: string, record: any): Promise<{ ok: boolean; uri?: string }> { const url = `${PDS_ORIGIN}/xrpc/com.atproto.repo.putRecord`; const dpop = await createDpopProof("POST", url, session.wmJwt, session.publicJwk, session.privateKey); const res = await fetch(url, { method: "POST", headers: { "Content-Type": "application/json", Authorization: `DPoP ${session.wmJwt}`, DPoP: dpop }, body: JSON.stringify({ repo: session.did, collection, rkey, record }), }); if (!res.ok) { console.error(`putRecord failed: ${res.status} ${await res.text()}`); return { ok: false }; } const result = await res.json() as { uri: string }; return { ok: true, uri: result.uri }; } async function dpopCreateRecord(session: DpopSession, collection: string, record: any): Promise<{ ok: boolean; uri?: string }> { const url = `${PDS_ORIGIN}/xrpc/com.atproto.repo.createRecord`; const dpop = await createDpopProof("POST", url, session.wmJwt, session.publicJwk, session.privateKey); const res = await fetch(url, { method: "POST", headers: { "Content-Type": "application/json", Authorization: `DPoP ${session.wmJwt}`, DPoP: dpop }, body: JSON.stringify({ repo: session.did, collection, record }), }); if (!res.ok) { console.error(`createRecord failed: ${res.status}`); return { ok: false }; } const result = await res.json() as { uri: string }; return { ok: true, uri: result.uri }; } // --- Cron: poll registry, ingest new nanopubs, create islands --- async function runCron(db: D1Database, keyPem: string): Promise { await ensureContrailReady(db); // Ensure state tables await db.prepare("CREATE TABLE IF NOT EXISTS cron_state (key TEXT PRIMARY KEY, value TEXT NOT NULL)").run(); 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(); // Create DPoP session const session = await createDpopSession(keyPem); // Get last known count and cursor const lastCount = (await db.prepare("SELECT value FROM cron_state WHERE key = 'last_nanopub_count'").first<{ value: string }>())?.value || "0"; const lastFirstId = (await db.prepare("SELECT value FROM cron_state WHERE key = 'last_first_id'").first<{ value: string }>())?.value || ""; // Check current count via HEAD const headRes = await fetch(`${REGISTRY}/nanopubs.json`, { method: "HEAD" }); const currentCount = headRes.headers.get("nanopub-registry-nanopub-count") || ""; if (currentCount === lastCount) { console.log(`No new nanopubs (count unchanged: ${currentCount})`); return; } console.log(`Nanopub count changed: ${lastCount} → ${currentCount}`); // Fetch the list const listRes = await fetch(`${REGISTRY}/nanopubs.json`); if (!listRes.ok) { console.error(`Failed to fetch nanopub list: ${listRes.status}`); return; } const allIds: string[] = await listRes.json() as string[]; // Find new IDs by comparing with cursor let newIds: string[] = []; if (lastFirstId) { const lastIdx = allIds.indexOf(lastFirstId); if (lastIdx > 0) { newIds = allIds.slice(0, lastIdx); } else { newIds = allIds.slice(0, 20); } } else { newIds = allIds.slice(0, 20); } // Cap at 20 per cron tick newIds = newIds.slice(0, 20); if (newIds.length === 0) { console.log("No new IDs to process"); await db.prepare("INSERT OR REPLACE INTO cron_state (key, value) VALUES (?, ?)").bind("last_nanopub_count", currentCount).run(); return; } console.log(`Processing ${newIds.length} new nanopubs`); // Get existing rkeys to skip duplicates const existingRkeys = new Set(); for (const id of newIds) { const rkey = generateRkey(id); const row = await db.prepare("SELECT rkey FROM records_nanopub WHERE rkey = ?").bind(rkey).first(); if (row) existingRkeys.add(rkey); } let ingested = 0; let islandsCreated = 0; for (const nanopubId of newIds) { const rkey = generateRkey(nanopubId); if (existingRkeys.has(rkey)) continue; try { const npRes = await fetch(`${REGISTRY}/np/${nanopubId}`, { headers: { Accept: "application/ld+json" }, signal: AbortSignal.timeout(10000), }); if (!npRes.ok) { console.error(`SKIP ${nanopubId.slice(0, 10)}: HTTP ${npRes.status}`); continue; } const graphs: NanopubGraph[] = await npRes.json() as NanopubGraph[]; const record = convertNanopub(graphs, nanopubId); // Write to PDS via DPoP const putResult = await dpopPutRecord(session, "org.latha.nanopub", rkey, record); if (!putResult.ok) { console.error(`PDS write failed for ${nanopubId.slice(0, 10)}`); continue; } // Insert into D1 const now = Date.now() * 1000; await db.prepare( "INSERT OR REPLACE INTO records_nanopub (uri, did, rkey, cid, record, time_us, indexed_at) VALUES (?, ?, ?, ?, ?, ?, ?)" ).bind(putResult.uri!, RESEARCHER_DID, rkey, "", JSON.stringify(record), now, now).run(); ingested++; // Create connections + island from URI triples const uriTriples = (record.assertion?.triples || []).filter( (t: any) => t.objectType === "uri" && t.object?.startsWith("https://") ); if (uriTriples.length >= 2) { try { const connectionUris: string[] = []; for (const t of uriTriples) { const connRecord = { $type: "network.cosmik.connection", source: t.subject, target: t.object, connectionType: t.predicate.split("/").pop() || t.predicate, note: `${record.assertion?.label || "Nanopub assertion"}`, createdAt: new Date().toISOString(), }; const connResult = await dpopCreateRecord(session, "network.cosmik.connection", connRecord); if (connResult.ok && connResult.uri) connectionUris.push(connResult.uri); } // Compute centroid const degreeCounts = new Map(); for (const t of uriTriples) { degreeCounts.set(t.subject, (degreeCounts.get(t.subject) || 0) + 1); degreeCounts.set(t.object, (degreeCounts.get(t.object) || 0) + 1); } const centroid = [...degreeCounts.entries()].sort((a, b) => b[1] - a[1])[0]?.[0] || uriTriples[0].subject; const iHash = await islandHash(centroid); // Build island record const label = record.assertion?.label || record.assertion?.subjectIri || "Untitled"; const creators = (record.pubInfo?.creators || []).map((c: string) => { const m = c.match(/orcid\.org\/([\d-]+)/); return m ? m[1] : c.split("/").pop(); }).join(", "); const islandRecord = { $type: "org.latha.island", source: { uri: centroid, collection: "network.cosmik.connection" }, connections: connectionUris.map(uri => ({ uri })), analysis: { title: label, themes: [label.split(" ").slice(0, 3).join(" "), "nanopublication", "semantic graph"], highlights: [], tensions: [], openQuestions: [], synthesis: `Nanopublication asserting ${uriTriples.length} semantic relationships. ${label}. Created by ${creators || "unknown"}. Source: ${record.sourceNanopubUri}`, }, createdAt: new Date().toISOString(), }; // Publish island const islandResult = await dpopPutRecord(session, "org.latha.island", iHash, islandRecord); if (islandResult.ok) { await db.prepare( "INSERT OR REPLACE INTO island_map (nanopub_rkey, island_rkey, created_at) VALUES (?, ?, ?)" ).bind(rkey, iHash, Date.now()).run(); islandsCreated++; } // Notify stigmergic try { await fetch(`${STIGMERGIC}/xrpc/org.latha.island.notifyOfUpdate`, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ did: RESEARCHER_DID }), }); } catch {} } catch (err: any) { console.error(`Island creation failed for ${rkey}: ${err?.message}`); } } await new Promise(r => setTimeout(r, 100)); } catch (err: any) { console.error(`Error processing ${nanopubId.slice(0, 10)}: ${err?.message}`); } } // Update cursor await db.prepare("INSERT OR REPLACE INTO cron_state (key, value) VALUES (?, ?)").bind("last_nanopub_count", currentCount).run(); await db.prepare("INSERT OR REPLACE INTO cron_state (key, value) VALUES (?, ?)").bind("last_first_id", allIds[0] || "").run(); console.log(`Cron done: ${ingested} ingested, ${islandsCreated} islands created`); } // --- HTTP app --- let app: Hono<{ Bindings: Env }> | null = null; function buildApp(env: Env): Hono<{ Bindings: Env }> { const app = new Hono<{ Bindings: Env }>(); const db = env.DB; app.use("*", cors()); // Landing page app.get("/", async (c) => { return c.html(LANDING_PAGE); }); // Nanopub detail page app.get("/nanopub/:rkey", async (c) => { const rkey = c.req.param("rkey"); await ensureContrailReady(db); let nanopub: any = null; let islandRkey: string | null = null; try { const row = await db .prepare("SELECT uri, record, did, rkey FROM records_nanopub WHERE rkey = ? LIMIT 1") .bind(rkey) .first<{ uri: string; record: string; did: string; rkey: string }>(); if (row) { nanopub = { uri: row.uri, rkey: row.rkey, ...JSON.parse(row.record) }; } // Check island mapping const islandRow = await db.prepare("SELECT island_rkey FROM island_map WHERE nanopub_rkey = ?").bind(rkey).first<{ island_rkey: string }>(); if (islandRow) islandRkey = islandRow.island_rkey; } catch {} if (!nanopub) { return c.html("

Not Found

Back

", 404); } // Inject island link into nanopub data const dataWithIsland = { ...nanopub, _islandRkey: islandRkey, _islandUrl: islandRkey ? `${STIGMERGIC}/island/${RESEARCHER_HANDLE}/${islandRkey}` : null }; const page = NANOPUB_PAGE.replace( "", `` ); return c.html(page); }); // API: list nanopubs sorted by publish time app.get("/api/nanopubs", async (c) => { await ensureContrailReady(db); const limit = Math.min(Number(c.req.query("limit") || 50), 200); const cursor = c.req.query("cursor"); // Sort by nanopub publish time (pubInfo.createdOn) extracted from record let rows: D1Result; if (cursor) { rows = await db .prepare( "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 ?" ) .bind(cursor, limit) .all(); } else { rows = await db .prepare( "SELECT uri, record, did, rkey, time_us FROM records_nanopub ORDER BY json_extract(record, '$.pubInfo.createdOn') DESC LIMIT ?" ) .bind(limit) .all(); } const nanopubs = (rows.results || []).map((r: any) => { try { const value = JSON.parse(r.record); const { assertion, pubInfo, nanopubType, sourceNanopubUri, signerPubkey, createdAt, trustUri, retracts, supersedes } = value; // Check island mapping let islandUrl: string | null = null; // We'll add island info separately after the loop return { uri: r.uri, rkey: r.rkey, assertion: assertion ? { label: assertion.label, subjectIri: assertion.subjectIri, predicateIri: assertion.predicateIri, objectIri: assertion.objectIri } : undefined, pubInfo: pubInfo ? { creators: pubInfo.creators, createdOn: pubInfo.createdOn } : undefined, nanopubType, sourceNanopubUri, signerPubkey, createdAt, trustUri, retracts, supersedes, _rkeyForIsland: r.rkey, }; } catch { return null; } }).filter(Boolean); // Batch-fetch island mappings const rkeys = nanopubs.map((n: any) => n._rkeyForIsland).filter(Boolean); if (rkeys.length > 0) { const placeholders = rkeys.map(() => "?").join(","); const islandRows = await db.prepare( `SELECT nanopub_rkey, island_rkey FROM island_map WHERE nanopub_rkey IN (${placeholders})` ).bind(...rkeys).all(); const islandMap = new Map(); for (const row of (islandRows.results || [])) { islandMap.set((row as any).nanopub_rkey, (row as any).island_rkey); } for (const n of nanopubs) { const rkey = (n as any)._rkeyForIsland; const islandRkey = islandMap.get(rkey); if (islandRkey) { (n as any).islandUrl = `${STIGMERGIC}/island/${RESEARCHER_HANDLE}/${islandRkey}`; } delete (n as any)._rkeyForIsland; } } const lastRow = rows.results?.[rows.results.length - 1]; const nextCursor = lastRow ? JSON.parse((lastRow as any).record)?.pubInfo?.createdOn : undefined; return c.json({ nanopubs, cursor: nextCursor }); }); // API: health check app.get("/api/health", async (c) => { await ensureContrailReady(db); let count = 0; try { const row = await db.prepare("SELECT COUNT(*) as cnt FROM records_nanopub").first<{ cnt: number }>(); count = row?.cnt || 0; } catch {} return c.json({ ok: true, nanopubs: count }); }); // Manual cron trigger app.post("/api/cron", async (c) => { const keyPem = c.env.RESEARCHER_KEY_PEM; if (!keyPem) return c.json({ error: "RESEARCHER_KEY_PEM not set" }, 500); try { await runCron(c.env.DB, keyPem); return c.json({ ok: true }); } catch (err: any) { return c.json({ error: err?.message }, 500); } }); // Crawl endpoint app.post("/xrpc/com.atproto.sync.requestCrawl", async (c) => { let body: { did?: string }; try { body = await c.req.json(); } catch { body = {}; } const did = body.did; if (!did) return c.json({ error: "BadRequest", message: "Provide did" }, 400); await ensureContrailReady(db); try { await db.prepare("INSERT OR IGNORE INTO identities (did, handle, pds, resolved_at) VALUES (?, ?, ?, ?)") .bind(did, "", "pds.latha.org", Date.now() * 1000).run(); const pdsUrl = `https://pds.latha.org/xrpc/com.atproto.repo.listRecords?repo=${did}&collection=org.latha.nanopub&limit=100`; const res = await fetch(pdsUrl); if (res.ok) { const data = await res.json() as { records: Array<{ uri: string; cid: string; value: any }> }; for (const rec of data.records) { const rkey = rec.uri.split("/").pop() || ""; const now = Date.now() * 1000; await db.prepare( "INSERT OR REPLACE INTO records_nanopub (uri, did, rkey, cid, record, time_us, indexed_at) VALUES (?, ?, ?, ?, ?, ?, ?)" ).bind(rec.uri, did, rkey, rec.cid, JSON.stringify(rec.value), now, now).run(); } } } catch (err: any) { console.error(`Crawl error: ${err?.message ?? err}`); } return c.json({ ok: true, did }); }); // Contrail passthrough app.all("*", async (c) => { const response = await contrailWorker.fetch(c.req.raw, c.env as unknown as Record); return response; }); return app; } export default { fetch(request: Request, env: Env): Response | Promise { app ??= buildApp(env); return app.fetch(request, env); }, async scheduled(event: ScheduledEvent, env: Env, ctx: ExecutionContext) { const keyPem = env.RESEARCHER_KEY_PEM; if (!keyPem) { console.error("RESEARCHER_KEY_PEM not set — skipping cron"); return; } ctx.waitUntil(runCron(env.DB, keyPem)); }, };