worker/index.js
// 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");
},
};