source.fact.ngo a coherence.ngo project

worker/index.js

raw ↗ · AGPL-3.0

// emergence-net Worker — the cloud-native cycle. Runs the exact same pipeline // cores as the local scripts (../pipeline/*), against Workers AI (binding, no // token) and R2 (binding) instead of REST + git storage. // // scheduled: every 2h UTC (cron trigger) // POST /run: manual cycle (Authorization: Bearer <RUN_SECRET>) // GET /latest.json : the live snapshot (same shape as the site's data.json) // GET /status.json : last cycle report // GET /partitions/<date>/<file> : raw partitions for the git sync // // R2 is the live primary; the git repo (emergence-data) remains the durable // archive, synced by scripts/sync-from-r2.sh. import { utcDate, prevDate, parseRegistryYaml, callCost, spendFromLedger, PROMPT_VERSION, nowIso } from "../pipeline/util.js"; import { SOURCES_YAML, TOPICS_YAML } from "./embeds.js"; import { DISTILL_PROMPT } from "../pipeline/prompts.js"; import { fetchSource, parseFeed, articleFromItem } from "../pipeline/feed.js"; import { distillUserMessage, topicListForPrompt, essenceFromResponse } from "../pipeline/distill-core.js"; import { runCluster } from "../pipeline/cluster-core.js"; import { buildTopicsTree, buildGis } from "../pipeline/aggregate-core.js"; const SOURCES = parseRegistryYaml(SOURCES_YAML); const TOPICS = parseRegistryYaml(TOPICS_YAML); const TOPIC_TEXT = topicListForPrompt(TOPICS); const CFG = { distillModel: "@cf/zai-org/glm-5.3-flash", budgetUsd: 1.0, maxAiCallsPerRun: 50, // slow-lane headroom; still bounded so a cycle always completes fastLaneAiCalls: 8, // per 5-minute fast lane: usually 0-2 new articles hotRecomputeAt: 10, // new essences since last full pass that trigger an immediate full cycle council: { lane: "verified", active: false }, }; /* ---------- R2 storage helpers ---------- */ const J = (o) => JSON.stringify(o); async function getJsonl(net, key) { const obj = await net.get(key); if (!obj) return []; const text = await obj.text(); return text.split("\n").filter(Boolean).map((l) => JSON.parse(l)); } async function putJsonl(net, key, lines) { await net.put(key, lines.length ? lines.map(J).join("\n") + "\n" : ""); } async function getJson(net, key) { const obj = await net.get(key); return obj ? obj.json() : null; } async function putJson(net, key, value) { await net.put(key, J(value)); } /* ---------- AI binding adapter (shapes vary by model family) ---------- */ let subrequests = 0; async function askAI(env, messages, { temperature = 0.2, maxTokens = 3000 } = {}) { if (subrequests >= CFG.maxAiCallsPerRun + 5) throw new Error("run subrequest cap reached"); subrequests++; const out = await env.AI.run(CFG.distillModel, { messages, temperature, max_tokens: maxTokens }); const text = out?.choices?.[0]?.message?.content ?? out?.response ?? out?.message?.content ?? String(out ?? ""); const usage = out?.usage || {}; const tokensIn = usage.prompt_tokens ?? Math.ceil(messages.reduce((n, m) => n + m.content.length, 0) / 4); const tokensOut = usage.completion_tokens ?? Math.ceil(text.length / 4); const cost = callCost(CFG.distillModel, tokensIn, tokensOut); return { text, tokensIn, tokensOut, cost }; } /* ---------- the cycle ---------- */ async function runCycle(env, mode = "full") { subrequests = 0; const started = Date.now(); const net = env.NET; const date = utcDate(); const status = { date, mode, started: nowIso(), ok: false, stages: {} }; try { /* ledger + budget */ const ledger = await getJsonl(net, "ledger.jsonl"); const spent = spendFromLedger(ledger, date); const budgetLeft = () => CFG.budgetUsd - spendFromLedger(ledger, date) > 0; const addSpend = (entry) => { ledger.push({ date, at: nowIso(), ...entry }); }; const persistLedger = () => putJsonl(net, "ledger.jsonl", ledger); /* scrape */ const articlesKey = `partitions/${date}/articles.jsonl`; const articles = await getJsonl(net, articlesKey); const seen = new Set(articles.map((a) => a.id)); const before = articles.length; const sourceStatus = (await getJson(net, `partitions/${date}/sources.json`)) || {}; const registry = Object.fromEntries(SOURCES.map((s) => [s.id, s])); const meta = (src) => ({ name: src.name, site: src.site || null, takes: src.takes || null, basis: src.basis || null, depth: src.depth }); const feedState = (await getJson(net, "feeds-state.json")) || {}; for (const src of Object.values(registry)) { if (src.via === "deferred") { sourceStatus[src.id] = { status: "deferred", http: null, fetched: 0, new: 0, ...meta(src) }; continue; } try { const r = await fetchSource(src, fetch, feedState[src.id]); if (r.validator) feedState[src.id] = r.validator; if (!r.notModified) { const items = parseFeed(r.xml); let fresh = 0; for (const it of items) { const article = articleFromItem(src, it, seen, date); if (!article) continue; articles.push(article); fresh++; } sourceStatus[src.id] = { status: "ok", http: r.http, fetched: items.length, new: fresh, ...meta(src) }; } else { const prev = sourceStatus[src.id] || { fetched: 0 }; sourceStatus[src.id] = { status: "ok", http: 304, fetched: prev.fetched, new: 0, ...meta(src) }; } } catch (e) { sourceStatus[src.id] = { status: "failed", http: null, error: String(e.message || e), fetched: 0, new: 0, ...meta(src) }; } } await putJson(net, "feeds-state.json", feedState); await putJsonl(net, articlesKey, articles); await putJson(net, `partitions/${date}/sources.json`, sourceStatus); status.stages.scrape = { new: articles.length - before, total: articles.length, sources: Object.values(sourceStatus).filter((s) => s.status === "ok").length }; /* distill */ const essenceKey = `partitions/${date}/essence.jsonl`; const essences = await getJsonl(net, essenceKey); const essenceBy = {}; for (const e of essences) { const prev = essenceBy[e.article.raw_id]; if (!prev || (prev.failure && !e.failure)) essenceBy[e.article.raw_id] = e; } const pending = articles.filter((a) => !essenceBy[a.id]); let distilled = 0, failures = 0; const distillCap = mode === "fast" ? CFG.fastLaneAiCalls : CFG.maxAiCallsPerRun; for (const a of pending) { if (distilled + failures >= distillCap) { status.stages.distillDeferred = pending.length - distilled - failures; break; } if (!budgetLeft()) break; try { const r = await askAI(env, [ { role: "system", content: DISTILL_PROMPT }, { role: "user", content: distillUserMessage(a, TOPIC_TEXT) }, ]); addSpend({ script: "distill", model: CFG.distillModel, tokens_in: r.tokensIn, tokens_out: r.tokensOut, cost_usd: r.cost }); const out = essenceFromResponse(r.text, a, { model: CFG.distillModel, promptVersion: PROMPT_VERSION, tokensIn: r.tokensIn, tokensOut: r.tokensOut, cost: r.cost, at: nowIso() }); essences.push(out.record); out.ok ? distilled++ : failures++; } catch (e) { status.stages.distillError = String(e.message || e); break; } } await putJsonl(net, essenceKey, essences); await persistLedger(); status.stages.distill = { distilled, failures }; /* fast lane: stop after distill; recompute fully only when enough is new */ if (mode === "fast") { await persistLedger(); status.stages.fast = { distilled, failures, deferred: status.stages.distillDeferred || 0 }; const lastFull = (await getJson(net, "status.json"))?.essences || 0; const nowOk = essences.filter((e) => !e.failure).length; status.hot = nowOk - lastFull >= CFG.hotRecomputeAt || distilled >= CFG.hotRecomputeAt; status.essences = nowOk; status.ok = true; status.finished = nowIso(); status.durationMs = Date.now() - started; await putJson(net, "status.json", status); if (status.hot) return runCycle(env, "full"); return status; } /* cluster (+ cross-day threads) */ const ok = essences.filter((e) => !e.failure); let events = null, clustering = null; if (ok.length) { const pd = prevDate(date); const prevData = await getJson(net, `partitions/${pd}/events.json`); const prevEvents = prevData?.events || null; let prevEssences = null; if (prevEvents) prevEssences = new Map((await getJsonl(net, `partitions/${pd}/essence.jsonl`)).filter((e) => !e.failure).map((e) => [e.id, e])); const result = await runCluster({ essences: ok, prevEvents, prevEssences, ask: (messages, opts) => askAI(env, messages, opts), onSpend: (s) => addSpend({ script: s.script, model: CFG.distillModel, tokens_in: s.tokensIn, tokens_out: s.tokensOut, cost_usd: s.cost }), budget: { exhausted: () => !budgetLeft() || subrequests >= CFG.maxAiCallsPerRun + 5 }, }); events = result.events; clustering = { ...result.clustering, previous_day: prevEvents ? pd : null }; await putJson(net, `partitions/${date}/events.json`, { date, generated_at: nowIso(), model: CFG.distillModel, prompt_version: PROMPT_VERSION, clustering, events }); await persistLedger(); status.stages.cluster = { events: events.length, ai: clustering.ai_calls, fallbacks: clustering.fallbacks, continuations: clustering.continuations }; } /* aggregate + gis + latest snapshot */ if (ok.length) { const { tree } = buildTopicsTree({ essences, events, topics: TOPICS }); const gis = buildGis({ essences, events }); const seenOk = new Set(), failedIds = new Set(); for (const e of essences) { if (e.failure) failedIds.add(e.article.raw_id); else if (!seenOk.has(e.article.raw_id)) seenOk.add(e.article.raw_id); } const dayTokensIn = ledger.filter((l) => l.date === date).reduce((n, l) => n + (l.tokens_in || 0), 0); const dayTokensOut = ledger.filter((l) => l.date === date).reduce((n, l) => n + (l.tokens_out || 0), 0); const daySpend = spendFromLedger(ledger, date); const day = { date, generated_at: nowIso(), counts: { articles: articles.length, essences: seenOk.size, failures: [...failedIds].filter((id) => !seenOk.has(id)).length, events: events ? events.length : null }, sources: sourceStatus, spend: { tokens_in: dayTokensIn, tokens_out: dayTokensOut, usd: daySpend }, }; await putJson(net, `partitions/${date}/topics.json`, { date, generated_at: nowIso(), tree }); await putJson(net, `partitions/${date}/day.json`, day); await putJson(net, `partitions/${date}/gis.json`, { date, generated_at: nowIso(), ...gis }); await putJson(net, "latest.json", { generated: nowIso(), date, day, topics: { tree }, gis: { date, ...gis }, events: events || [], essences: ok, council: CFG.council, license_note: "Summaries of newsroom teasers, attributed per record; not verified facts.", }); status.stages.aggregate = { topics: tree.length, countries: gis.countries.length, events: events ? events.length : null }; } status.ok = true; status.essences = essences.filter((e) => !e.failure).length; status.spend = Number(spendFromLedger(ledger, date).toFixed(4)); } catch (e) { status.error = String(e && e.message || e); } status.finished = nowIso(); status.durationMs = Date.now() - started; await putJson(net, "status.json", status); return status; } /* ---------- HTTP + scheduled ---------- */ const jsonResponse = (data, cache = "max-age=0") => new Response(JSON.stringify(data), { headers: { "content-type": "application/json", "cache-control": `public, ${cache}` } }); export default { async fetch(request, env, ctx) { const url = new URL(request.url); const key = decodeURIComponent(url.pathname.replace(/^\/+/, "")); if (request.method === "POST" && (key === "run" || key === "run/")) { const auth = request.headers.get("authorization") || ""; if (!env.RUN_SECRET || auth !== `Bearer ${env.RUN_SECRET}`) { return new Response("forbidden", { status: 403 }); } // Awaited: the cycle runs while the caller holds the connection open. // (waitUntil is capped ~30s after response — too short for a cycle; the // cron invocation needs no client and gets the full 15-minute window.) const status = await runCycle(env); return jsonResponse(status); } if (request.method === "GET") { if (key === "latest.json" || key === "" || key === "/") { const data = await getJson(env.NET, "latest.json"); return data ? jsonResponse(data, "max-age=60") : new Response("not ready yet", { status: 503, headers: { "cache-control": "max-age=0" } }); } if (key === "status.json") { const data = await getJson(env.NET, "status.json"); return jsonResponse(data || {}); } if (key.startsWith("partitions/")) { const obj = await env.NET.get(key); if (!obj) return new Response("not found", { status: 404 }); return new Response(obj.body, { headers: { "content-type": "application/json", "cache-control": "public, max-age=0" } }); } } return new Response("emergence-net: /latest.json /status.json POST /run", { status: 404 }); }, async scheduled(event, env) { // Two lanes: the 5-minute cron is the fast lane (scrape + distill, hot // recompute when enough is new); the 2-hour cron runs the full pipeline. await runCycle(env, event.cron === "*/5 * * * *" ? "fast" : "full"); }, };