/** * Papers Appview — watches the Paper Skygest feed, detects linked papers, * generates summaries, and writes org.latha.papers.summary records to PDS. * * Hosted at papers.latha.org */ 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 { pollNewPapers } from "./feed-watcher.js"; import { buildSummary, enrichWithLetta } from "./paper-summarizer.js"; import { writeSummaryToPds } from "./pds-writer.js"; import { LANDING_PAGE } from "./landing-page.js"; import { SUMMARY_PAGE } from "./summary-page.js"; interface Env { DB: D1Database; FEED_CURSOR?: string; RESEARCHER_KEY_PEM?: string; LETTA_API_KEY?: string; LETTA_AGENT_ID?: 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; } 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 — browse all indexed summaries app.get("/", async (c) => { await ensureContrailReady(db); let summaries: any[] = []; try { const rows = await db .prepare( "SELECT uri, record, did, rkey, time_us FROM records_summary ORDER BY time_us DESC LIMIT 100" ) .all(); summaries = (rows.results || []).map((r: any) => { try { const value = JSON.parse(r.record); return { uri: r.uri, rkey: r.rkey, title: value?.title || "Untitled", paperUrl: value?.paperUrl || "", venue: value?.venue || "", year: value?.year || null, authors: value?.authors || [], summary: value?.summary || "", domains: value?.domains || [], posterHandle: value?.posterHandle || "", indexedAt: value?.indexedAt || "", }; } catch { return null; } }).filter(Boolean); } catch { // empty list is fine } const page = LANDING_PAGE.replace( "", `` ); return c.html(page); }); // Summary detail page app.get("/paper/:rkey", async (c) => { const rkey = c.req.param("rkey"); await ensureContrailReady(db); let summary: any = null; try { const row = await db .prepare("SELECT uri, record, did, rkey FROM records_summary WHERE rkey = ? LIMIT 1") .bind(rkey) .first<{ uri: string; record: string; did: string; rkey: string }>(); if (row) { const value = JSON.parse(row.record); summary = { uri: row.uri, rkey: row.rkey, ...value, }; } } catch { // not found } if (!summary) { return c.html("

Not Found

", 404); } const page = SUMMARY_PAGE.replace( "", `` ); return c.html(page); }); // API: list summaries as JSON app.get("/api/summaries", async (c) => { await ensureContrailReady(db); const limit = Math.min(Number(c.req.query("limit") || 50), 200); const cursor = c.req.query("cursor"); let rows: D1Result; if (cursor) { rows = await db .prepare("SELECT uri, record, did, rkey, time_us FROM records_summary WHERE time_us < ? ORDER BY time_us DESC LIMIT ?") .bind(cursor, limit) .all(); } else { rows = await db .prepare("SELECT uri, record, did, rkey, time_us FROM records_summary ORDER BY time_us DESC LIMIT ?") .bind(limit) .all(); } const summaries = (rows.results || []).map((r: any) => { try { const value = JSON.parse(r.record); return { uri: r.uri, rkey: r.rkey, ...value }; } catch { return null; } }).filter(Boolean); const lastRow = rows.results?.[rows.results.length - 1]; const nextCursor = lastRow?.time_us || undefined; return c.json({ summaries, cursor: nextCursor }); }); // API: trigger feed poll + summary generation app.post("/api/poll", async (c) => { await ensureContrailReady(db); // Get last cursor from D1 let lastCursor: string | undefined; try { const row = await db .prepare("SELECT value FROM instance_settings WHERE key = 'feed_cursor' LIMIT 1") .first<{ value: string }>(); lastCursor = row?.value || undefined; } catch { // first run } const { posts, newCursor } = await pollNewPapers(lastCursor); // Deduplicate against existing records const existingUrls = new Set(); if (posts.length > 0) { const placeholders = posts.map(() => "?").join(","); const urls = posts.map(p => p.paperUrl); const rows = await db .prepare( `SELECT DISTINCT json_extract(record, '$.paperUrl') as paperUrl FROM records_summary WHERE paperUrl IN (${placeholders})` ) .bind(...urls) .all<{ paperUrl: string }>(); for (const r of rows.results || []) { if (r.paperUrl) existingUrls.add(r.paperUrl); } } const newPosts = posts.filter(p => !existingUrls.has(p.paperUrl)); let written = 0; let errors = 0; for (const post of newPosts) { try { const summary = await buildSummary(post); const result = await writeSummaryToPds(post, summary); // Notify contrail to index the new record await contrail.notify(result.uri); written++; } catch (err: any) { console.error(`Failed to write summary for ${post.paperUrl}: ${err?.message ?? err}`); errors++; } } // Save cursor if (newCursor) { await db .prepare("INSERT OR REPLACE INTO instance_settings (key, value) VALUES ('feed_cursor', ?)") .bind(newCursor) .run(); } return c.json({ polled: posts.length, new: newPosts.length, written, errors, }); }); // 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_summary") .first<{ cnt: number }>(); count = row?.cnt || 0; } catch { // table might not exist yet } return c.json({ ok: true, summaries: count }); }); // Enrich existing summaries using Letta API app.post("/api/enrich", async (c) => { const apiKey = process.env.LETTA_API_KEY; const agentId = process.env.LETTA_AGENT_ID; if (!apiKey || !agentId) { return c.json({ error: "Missing LETTA_API_KEY or LETTA_AGENT_ID" }, 500); } await ensureContrailReady(db); // Get records that need enrichment (empty/placeholder summaries, no takeaway, or Untitled) const rows = await db .prepare( `SELECT rkey, record FROM records_summary WHERE json_extract(record, '$.summary') = '' OR json_extract(record, '$.summary') LIKE 'Published in%' OR json_extract(record, '$.title') = 'Untitled' OR json_extract(record, '$.takeaway') IS NULL OR json_extract(record, '$.takeaway') = '' ORDER BY time_us DESC LIMIT 20` ) .all<{ rkey: string; record: string }>(); if (!rows.results?.length) { return c.json({ enriched: 0, message: "No records need enrichment" }); } let enriched = 0; let errors = 0; for (const row of rows.results) { try { const record = JSON.parse(row.record) as { title: string; paperUrl: string; abstract: string; domains: string[]; venue: string; summary: string; takeaway?: string; }; const result = await enrichWithLetta( { title: record.title, paperUrl: record.paperUrl, abstract: record.abstract, domains: record.domains || [], venue: record.venue || "", }, apiKey, agentId, ); // Update the record in D1 record.summary = result.summary; record.domains = result.domains; if (result.takeaway) record.takeaway = result.takeaway; await db .prepare("UPDATE records_summary SET record = ? WHERE rkey = ?") .bind(JSON.stringify(record), row.rkey) .run(); enriched++; console.log(`Enriched: ${record.title?.substring(0, 50)}...`); } catch (err: any) { console.error(`Enrichment failed for ${row.rkey}: ${err?.message ?? err}`); errors++; } } return c.json({ enriched, errors, total: rows.results.length }); }); // Crawl endpoint: register a DID and trigger immediate ingest app.post("/xrpc/com.atproto.sync.requestCrawl", async (c) => { let body: { did?: string; hostname?: 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 { // Ensure identities table has this DID await db.exec( "CREATE TABLE IF NOT EXISTS identities (did TEXT PRIMARY KEY, handle TEXT, time_us INTEGER)" ); await db .prepare("INSERT OR IGNORE INTO identities (did, handle, time_us) VALUES (?, ?, ?)") .bind(did, "researcher.pds.latha.org", Date.now() * 1000) .run(); // Ensure records_summary table exists await db.exec( "CREATE TABLE IF NOT EXISTS records_summary (uri TEXT PRIMARY KEY, cid TEXT, did TEXT, rkey TEXT, record TEXT, time_us INTEGER)" ); // Fetch records from PDS and insert directly const pdsUrl = `https://pds.latha.org/xrpc/com.atproto.repo.listRecords?repo=${did}&collection=org.latha.papers.summary&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() || ""; await db .prepare( "INSERT OR REPLACE INTO records_summary (uri, cid, did, rkey, record, time_us) VALUES (?, ?, ?, ?, ?, ?)" ) .bind(rec.uri, rec.cid, did, rkey, JSON.stringify(rec.value), Date.now() * 1000) .run(); } } } catch (err: any) { console.error(`Crawl error: ${err?.message ?? err}`); } return c.json({ ok: true, did }); }); // All other routes pass through to contrail 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) { // Cron: poll feed and generate summaries const db = env.DB; await ensureContrailReady(db); let lastCursor: string | undefined; try { const row = await db .prepare("SELECT value FROM instance_settings WHERE key = 'feed_cursor' LIMIT 1") .first<{ value: string }>(); lastCursor = row?.value || undefined; } catch { // first run } const { posts, newCursor } = await pollNewPapers(lastCursor); // Deduplicate const existingUrls = new Set(); if (posts.length > 0) { const placeholders = posts.map(() => "?").join(","); const urls = posts.map(p => p.paperUrl); const rows = await db .prepare( `SELECT DISTINCT json_extract(record, '$.paperUrl') as paperUrl FROM records_summary WHERE paperUrl IN (${placeholders})` ) .bind(...urls) .all<{ paperUrl: string }>(); for (const r of rows.results || []) { if (r.paperUrl) existingUrls.add(r.paperUrl); } } const newPosts = posts.filter(p => !existingUrls.has(p.paperUrl)); for (const post of newPosts) { try { const summary = await buildSummary(post); const result = await writeSummaryToPds(post, summary); await contrail.notify(result.uri); } catch (err: any) { console.error(`Cron: failed to write summary for ${post.paperUrl}: ${err?.message ?? err}`); } } if (newCursor) { await db .prepare("INSERT OR REPLACE INTO instance_settings (key, value) VALUES ('feed_cursor', ?)") .bind(newCursor) .run(); } }, };