This repository has no description
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};