This repository has no description
0

Configure Feed

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

papers / src / worker.ts
20 kB 603 lines
1/** 2 * Papers Appview — watches the Paper Skygest feed, detects linked papers, 3 * generates summaries, and writes org.latha.papers.summary records to PDS. 4 * 5 * Hosted at papers.latha.org 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 { pollNewPapers, pollFollowedAccounts, type PaperPost } from "./feed-watcher.js"; 15import { buildSummary, enrichWithLetta, type PaperSummary } from "./paper-summarizer.js"; 16import { writeSummaryToPds } from "./pds-writer.js"; 17import { LANDING_PAGE } from "./landing-page.js"; 18import { SUMMARY_PAGE } from "./summary-page.js"; 19 20interface Env { 21 DB: D1Database; 22 FEED_CURSOR?: string; 23 RESEARCHER_KEY_PEM?: string; 24 LETTA_API_KEY?: string; 25 LETTA_AGENT_ID?: 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 38let app: Hono<{ Bindings: Env }> | null = null; 39 40function buildApp(env: Env): Hono<{ Bindings: Env }> { 41 const app = new Hono<{ Bindings: Env }>(); 42 const db = env.DB; 43 44 app.use("*", cors()); 45 46 // Landing page — browse all indexed summaries 47 app.get("/", async (c) => { 48 await ensureContrailReady(db); 49 50 let summaries: any[] = []; 51 try { 52 const rows = await db 53 .prepare( 54 "SELECT uri, record, did, rkey, time_us FROM records_summary ORDER BY time_us DESC LIMIT 100" 55 ) 56 .all(); 57 summaries = (rows.results || []).map((r: any) => { 58 try { 59 const value = JSON.parse(r.record); 60 return { 61 uri: r.uri, 62 rkey: r.rkey, 63 title: value?.title || "Untitled", 64 paperUrl: value?.paperUrl || "", 65 venue: value?.venue || "", 66 year: value?.year || null, 67 authors: value?.authors || [], 68 summary: value?.summary || "", 69 domains: value?.domains || [], 70 posterHandle: value?.posterHandle || "", 71 indexedAt: value?.indexedAt || "", 72 }; 73 } catch { 74 return null; 75 } 76 }).filter(Boolean); 77 } catch { 78 // empty list is fine 79 } 80 81 const page = LANDING_PAGE.replace( 82 "</head>", 83 `<script>window.__SUMMARIES__=${JSON.stringify(summaries).replace(/<\//g, "<\\/")};</script></head>` 84 ); 85 return c.html(page); 86 }); 87 88 // Summary detail page 89 app.get("/paper/:rkey", async (c) => { 90 const rkey = c.req.param("rkey"); 91 await ensureContrailReady(db); 92 93 let summary: any = null; 94 try { 95 const row = await db 96 .prepare("SELECT uri, record, did, rkey FROM records_summary WHERE rkey = ? LIMIT 1") 97 .bind(rkey) 98 .first<{ uri: string; record: string; did: string; rkey: string }>(); 99 if (row) { 100 const value = JSON.parse(row.record); 101 summary = { 102 uri: row.uri, 103 rkey: row.rkey, 104 ...value, 105 }; 106 } 107 } catch { 108 // not found 109 } 110 111 if (!summary) { 112 return c.html("<h1>Not Found</h1>", 404); 113 } 114 115 const page = SUMMARY_PAGE.replace( 116 "</head>", 117 `<script>window.__SUMMARY__=${JSON.stringify(summary)};</script></head>` 118 ); 119 return c.html(page); 120 }); 121 122 // API: list summaries as JSON 123 app.get("/api/summaries", async (c) => { 124 await ensureContrailReady(db); 125 const limit = Math.min(Number(c.req.query("limit") || 50), 200); 126 const cursor = c.req.query("cursor"); 127 128 let rows: D1Result<any>; 129 if (cursor) { 130 rows = await db 131 .prepare("SELECT uri, record, did, rkey, time_us FROM records_summary WHERE time_us < ? ORDER BY time_us DESC LIMIT ?") 132 .bind(cursor, limit) 133 .all(); 134 } else { 135 rows = await db 136 .prepare("SELECT uri, record, did, rkey, time_us FROM records_summary ORDER BY time_us DESC LIMIT ?") 137 .bind(limit) 138 .all(); 139 } 140 141 const summaries = (rows.results || []).map((r: any) => { 142 try { 143 const value = JSON.parse(r.record); 144 return { uri: r.uri, rkey: r.rkey, ...value }; 145 } catch { 146 return null; 147 } 148 }).filter(Boolean); 149 150 const lastRow = rows.results?.[rows.results.length - 1]; 151 const nextCursor = lastRow?.time_us || undefined; 152 153 return c.json({ summaries, cursor: nextCursor }); 154 }); 155 156 // API: trigger feed poll + summary generation 157 app.post("/api/poll", async (c) => { 158 await ensureContrailReady(db); 159 160 // Get last cursor from D1 161 let lastCursor: string | undefined; 162 try { 163 const row = await db 164 .prepare("SELECT value FROM instance_settings WHERE key = 'feed_cursor' LIMIT 1") 165 .first<{ value: string }>(); 166 lastCursor = row?.value || undefined; 167 } catch { 168 // first run 169 } 170 171 const { posts: feedPosts, newCursor } = await pollNewPapers(lastCursor); 172 173 // Also poll followed accounts for paper links 174 let followedPosts: PaperPost[] = []; 175 try { 176 followedPosts = await pollFollowedAccounts(5); 177 } catch (err: any) { 178 console.error(`Followed accounts poll failed: ${err?.message ?? err}`); 179 } 180 181 // Merge and deduplicate 182 const seenPaperUrls = new Set<string>(); 183 const posts: PaperPost[] = []; 184 for (const p of [...feedPosts, ...followedPosts]) { 185 if (!seenPaperUrls.has(p.paperUrl)) { 186 seenPaperUrls.add(p.paperUrl); 187 posts.push(p); 188 } 189 } 190 191 // Deduplicate against existing records 192 const existingUrls = new Set<string>(); 193 if (posts.length > 0) { 194 const placeholders = posts.map(() => "?").join(","); 195 const urls = posts.map(p => p.paperUrl); 196 const rows = await db 197 .prepare( 198 `SELECT DISTINCT json_extract(record, '$.paperUrl') as paperUrl FROM records_summary WHERE paperUrl IN (${placeholders})` 199 ) 200 .bind(...urls) 201 .all<{ paperUrl: string }>(); 202 for (const r of rows.results || []) { 203 if (r.paperUrl) existingUrls.add(r.paperUrl); 204 } 205 } 206 207 const newPosts = posts.filter(p => !existingUrls.has(p.paperUrl)); 208 let written = 0; 209 let errors = 0; 210 const errorDetails: string[] = []; 211 212 const apiKey = process.env.LETTA_API_KEY; 213 const agentId = process.env.LETTA_AGENT_ID; 214 215 // Build heuristic summaries in parallel 216 const heuristicResults = await Promise.all(newPosts.map(async (post) => { 217 try { 218 return { post, summary: await buildSummary(post) }; 219 } catch (err: any) { 220 return { post, error: err?.message ?? err }; 221 } 222 })); 223 224 // Enrich in parallel (batch of 5 to avoid rate limits) 225 const enrichedResults: Array<{ post: PaperPost; summary: PaperSummary; error?: string }> = []; 226 for (let i = 0; i < heuristicResults.length; i += 5) { 227 const batch = heuristicResults.slice(i, i + 5); 228 const batchResults = await Promise.all(batch.map(async ({ post, summary, error }) => { 229 if (error || !summary) return { post, summary: summary!, error }; 230 if (!apiKey || !agentId) return { post, summary }; 231 try { 232 const enrichment = await enrichWithLetta( 233 { title: summary.title, paperUrl: summary.paperUrl, abstract: summary.abstract, domains: summary.domains, venue: summary.venue }, 234 apiKey, 235 agentId, 236 ); 237 if (enrichment.title) summary.title = enrichment.title; 238 summary.summary = enrichment.summary; 239 summary.domains = enrichment.domains; 240 if (enrichment.takeaway) (summary as any).takeaway = enrichment.takeaway; 241 if (enrichment.year) summary.year = enrichment.year; 242 if (enrichment.publishedDate) (summary as any).publishedDate = enrichment.publishedDate; 243 return { post, summary }; 244 } catch (err: any) { 245 console.error(`Enrichment failed for ${post.paperUrl}: ${err?.message ?? err}, using heuristic`); 246 return { post, summary }; 247 } 248 })); 249 enrichedResults.push(...batchResults); 250 } 251 252 // Write all to PDS and D1 253 for (const { post, summary, error } of enrichedResults) { 254 if (error || !summary) { 255 const msg = `${post.paperUrl}: ${error}`; 256 console.error(`Failed to build summary for ${msg}`); 257 if (errors < 3) errorDetails.push(msg.substring(0, 200)); 258 errors++; 259 continue; 260 } 261 try { 262 const result = await writeSummaryToPds(post, summary); 263 264 const rkey = result.uri.split("/").pop() || ""; 265 await db 266 .prepare("INSERT OR REPLACE INTO records_summary (uri, cid, did, rkey, record, time_us, indexed_at) VALUES (?, ?, ?, ?, ?, ?, ?)") 267 .bind(result.uri, result.cid, "did:plc:3kkhul7jznlb6ba7rprzawnj", rkey, JSON.stringify({ 268 $type: "org.latha.papers.summary", 269 ...summary, 270 sourceUri: post.postUri, 271 posterDid: post.posterDid, 272 posterHandle: post.posterHandle, 273 postText: post.postText, 274 indexedAt: new Date().toISOString(), 275 }), Date.now() * 1000, Date.now() * 1000) 276 .run(); 277 written++; 278 } catch (err: any) { 279 const msg = `${post.paperUrl}: ${err?.message ?? err}`; 280 console.error(`Failed to write summary for ${msg}`); 281 if (errors < 3) errorDetails.push(msg.substring(0, 200)); 282 errors++; 283 } 284 } 285 286 // Save cursor 287 if (newCursor) { 288 await db 289 .prepare("INSERT OR REPLACE INTO instance_settings (key, value) VALUES ('feed_cursor', ?)") 290 .bind(newCursor) 291 .run(); 292 } 293 294 return c.json({ 295 polled: posts.length, 296 fromFeed: feedPosts.length, 297 fromFollows: followedPosts.length, 298 new: newPosts.length, 299 written, 300 errors, 301 lastError: errors > 0 ? errorDetails[0] : undefined, 302 }); 303 }); 304 305 // Health check 306 app.get("/api/health", async (c) => { 307 await ensureContrailReady(db); 308 let count = 0; 309 try { 310 const row = await db 311 .prepare("SELECT COUNT(*) as cnt FROM records_summary") 312 .first<{ cnt: number }>(); 313 count = row?.cnt || 0; 314 } catch { 315 // table might not exist yet 316 } 317 const hasKey = !!process.env.RESEARCHER_KEY_PEM; 318 return c.json({ ok: true, summaries: count, hasPdsKey: hasKey }); 319 }); 320 321 // Enrich existing summaries using Letta API 322 app.post("/api/enrich", async (c) => { 323 const apiKey = process.env.LETTA_API_KEY; 324 const agentId = process.env.LETTA_AGENT_ID; 325 if (!apiKey || !agentId) { 326 return c.json({ error: "Missing LETTA_API_KEY or LETTA_AGENT_ID" }, 500); 327 } 328 329 await ensureContrailReady(db); 330 331 // Get records that need enrichment 332 const rows = await db 333 .prepare( 334 `SELECT rkey, record FROM records_summary 335 WHERE json_extract(record, '$.summary') = '' 336 OR json_extract(record, '$.summary') LIKE 'Published in%' 337 OR json_extract(record, '$.title') = 'Untitled' 338 OR json_extract(record, '$.takeaway') IS NULL 339 OR json_extract(record, '$.takeaway') = '' 340 OR json_extract(record, '$.year') IS NULL 341 OR json_extract(record, '$.year') = 0 342 OR json_extract(record, '$.publishedDate') IS NULL 343 ORDER BY time_us DESC LIMIT 20` 344 ) 345 .all<{ rkey: string; record: string }>(); 346 347 if (!rows.results?.length) { 348 return c.json({ enriched: 0, message: "No records need enrichment" }); 349 } 350 351 let enriched = 0; 352 let errors = 0; 353 354 // Enrich in parallel batches of 5 355 for (let i = 0; i < rows.results.length; i += 5) { 356 const batch = rows.results.slice(i, i + 5); 357 const results = await Promise.allSettled(batch.map(async (row) => { 358 const record = JSON.parse(row.record) as { 359 title: string; 360 paperUrl: string; 361 abstract: string; 362 domains: string[]; 363 venue: string; 364 summary: string; 365 takeaway?: string; 366 year?: number; 367 }; 368 369 const result = await enrichWithLetta( 370 { 371 title: record.title, 372 paperUrl: record.paperUrl, 373 abstract: record.abstract, 374 domains: record.domains || [], 375 venue: record.venue || "", 376 }, 377 apiKey, 378 agentId, 379 ); 380 381 record.summary = result.summary; 382 record.domains = result.domains; 383 if (result.takeaway) record.takeaway = result.takeaway; 384 if (result.title) record.title = result.title; 385 if (result.year) record.year = result.year; 386 if (result.publishedDate) record.publishedDate = result.publishedDate; 387 388 await db 389 .prepare("UPDATE records_summary SET record = ? WHERE rkey = ?") 390 .bind(JSON.stringify(record), row.rkey) 391 .run(); 392 393 return record.title; 394 })); 395 396 for (const r of results) { 397 if (r.status === "fulfilled") { 398 enriched++; 399 console.log(`Enriched: ${r.value?.substring(0, 50)}...`); 400 } else { 401 console.error(`Enrichment failed: ${r.reason}`); 402 errors++; 403 } 404 } 405 } 406 407 return c.json({ enriched, errors, total: rows.results.length }); 408 }); 409 410 // Crawl endpoint: register a DID and trigger immediate ingest 411 app.post("/xrpc/com.atproto.sync.requestCrawl", async (c) => { 412 let body: { did?: string; hostname?: string }; 413 try { 414 body = await c.req.json(); 415 } catch { 416 body = {}; 417 } 418 const did = body.did; 419 if (!did) { 420 return c.json({ error: "BadRequest", message: "Provide did" }, 400); 421 } 422 423 await ensureContrailReady(db); 424 try { 425 // Ensure identities table has this DID 426 await db.exec( 427 "CREATE TABLE IF NOT EXISTS identities (did TEXT PRIMARY KEY, handle TEXT, time_us INTEGER)" 428 ); 429 await db 430 .prepare("INSERT OR IGNORE INTO identities (did, handle, time_us) VALUES (?, ?, ?)") 431 .bind(did, "researcher.pds.latha.org", Date.now() * 1000) 432 .run(); 433 434 // Ensure records_summary table exists 435 await db.exec( 436 "CREATE TABLE IF NOT EXISTS records_summary (uri TEXT PRIMARY KEY, cid TEXT, did TEXT, rkey TEXT, record TEXT, time_us INTEGER)" 437 ); 438 439 // Fetch records from PDS and insert directly 440 const pdsUrl = `https://pds.latha.org/xrpc/com.atproto.repo.listRecords?repo=${did}&collection=org.latha.papers.summary&limit=100`; 441 const res = await fetch(pdsUrl); 442 if (res.ok) { 443 const data = await res.json() as { 444 records: Array<{ uri: string; cid: string; value: any }>; 445 }; 446 for (const rec of data.records) { 447 const rkey = rec.uri.split("/").pop() || ""; 448 await db 449 .prepare( 450 "INSERT OR REPLACE INTO records_summary (uri, cid, did, rkey, record, time_us) VALUES (?, ?, ?, ?, ?, ?)" 451 ) 452 .bind(rec.uri, rec.cid, did, rkey, JSON.stringify(rec.value), Date.now() * 1000) 453 .run(); 454 } 455 } 456 } catch (err: any) { 457 console.error(`Crawl error: ${err?.message ?? err}`); 458 } 459 460 return c.json({ ok: true, did }); 461 }); 462 463 // All other routes pass through to contrail 464 app.all("*", async (c) => { 465 const response = await contrailWorker.fetch( 466 c.req.raw, 467 c.env as unknown as Record<string, unknown> 468 ); 469 return response; 470 }); 471 472 return app; 473} 474 475export default { 476 fetch(request: Request, env: Env): Response | Promise<Response> { 477 app ??= buildApp(env); 478 return app.fetch(request, env); 479 }, 480 async scheduled(event: ScheduledEvent, env: Env, ctx: ExecutionContext) { 481 // Cron: poll feed and generate summaries 482 const db = env.DB; 483 await ensureContrailReady(db); 484 485 let lastCursor: string | undefined; 486 try { 487 const row = await db 488 .prepare("SELECT value FROM instance_settings WHERE key = 'feed_cursor' LIMIT 1") 489 .first<{ value: string }>(); 490 lastCursor = row?.value || undefined; 491 } catch { 492 // first run 493 } 494 495 const { posts: feedPosts, newCursor } = await pollNewPapers(lastCursor); 496 497 // Also poll followed accounts 498 let followedPosts: PaperPost[] = []; 499 try { 500 followedPosts = await pollFollowedAccounts(5); 501 } catch {} 502 503 // Merge and deduplicate 504 const seenPaperUrls = new Set<string>(); 505 const posts: PaperPost[] = []; 506 for (const p of [...feedPosts, ...followedPosts]) { 507 if (!seenPaperUrls.has(p.paperUrl)) { 508 seenPaperUrls.add(p.paperUrl); 509 posts.push(p); 510 } 511 } 512 513 // Deduplicate against existing records 514 const existingUrls = new Set<string>(); 515 if (posts.length > 0) { 516 const placeholders = posts.map(() => "?").join(","); 517 const urls = posts.map(p => p.paperUrl); 518 const rows = await db 519 .prepare( 520 `SELECT DISTINCT json_extract(record, '$.paperUrl') as paperUrl FROM records_summary WHERE paperUrl IN (${placeholders})` 521 ) 522 .bind(...urls) 523 .all<{ paperUrl: string }>(); 524 for (const r of rows.results || []) { 525 if (r.paperUrl) existingUrls.add(r.paperUrl); 526 } 527 } 528 529 const newPosts = posts.filter(p => !existingUrls.has(p.paperUrl)); 530 const apiKey = process.env.LETTA_API_KEY; 531 const agentId = process.env.LETTA_AGENT_ID; 532 533 // Build heuristic summaries in parallel 534 const heuristicResults = await Promise.all(newPosts.map(async (post) => { 535 try { 536 return { post, summary: await buildSummary(post) }; 537 } catch (err: any) { 538 return { post, error: err?.message ?? err }; 539 } 540 })); 541 542 // Enrich in parallel (batch of 5) 543 const enrichedResults: Array<{ post: PaperPost; summary: PaperSummary; error?: string }> = []; 544 for (let i = 0; i < heuristicResults.length; i += 5) { 545 const batch = heuristicResults.slice(i, i + 5); 546 const batchResults = await Promise.all(batch.map(async ({ post, summary, error }) => { 547 if (error || !summary) return { post, summary: summary!, error }; 548 if (!apiKey || !agentId) return { post, summary }; 549 try { 550 const enrichment = await enrichWithLetta( 551 { title: summary.title, paperUrl: summary.paperUrl, abstract: summary.abstract, domains: summary.domains, venue: summary.venue }, 552 apiKey, 553 agentId, 554 ); 555 if (enrichment.title) summary.title = enrichment.title; 556 summary.summary = enrichment.summary; 557 summary.domains = enrichment.domains; 558 if (enrichment.takeaway) (summary as any).takeaway = enrichment.takeaway; 559 if (enrichment.year) summary.year = enrichment.year; 560 if (enrichment.publishedDate) (summary as any).publishedDate = enrichment.publishedDate; 561 return { post, summary }; 562 } catch { 563 return { post, summary }; 564 } 565 })); 566 enrichedResults.push(...batchResults); 567 } 568 569 // Write all to PDS and D1 570 for (const { post, summary, error } of enrichedResults) { 571 if (error || !summary) { 572 console.error(`Cron: failed to build summary for ${post.paperUrl}: ${error}`); 573 continue; 574 } 575 try { 576 const result = await writeSummaryToPds(post, summary); 577 578 const rkey = result.uri.split("/").pop() || ""; 579 await db 580 .prepare("INSERT OR REPLACE INTO records_summary (uri, cid, did, rkey, record, time_us, indexed_at) VALUES (?, ?, ?, ?, ?, ?, ?)") 581 .bind(result.uri, result.cid, "did:plc:3kkhul7jznlb6ba7rprzawnj", rkey, JSON.stringify({ 582 $type: "org.latha.papers.summary", 583 ...summary, 584 sourceUri: post.postUri, 585 posterDid: post.posterDid, 586 posterHandle: post.posterHandle, 587 postText: post.postText, 588 indexedAt: new Date().toISOString(), 589 }), Date.now() * 1000, Date.now() * 1000) 590 .run(); 591 } catch (err: any) { 592 console.error(`Cron: failed to write summary for ${post.paperUrl}: ${err?.message ?? err}`); 593 } 594 } 595 596 if (newCursor) { 597 await db 598 .prepare("INSERT OR REPLACE INTO instance_settings (key, value) VALUES ('feed_cursor', ?)") 599 .bind(newCursor) 600 .run(); 601 } 602 }, 603};