Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
12 kB · 288 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289import { XMLParser } from "fast-xml-parser";import { sha256, type JsonObject } from "../core/json.js";import { newId } from "../core/ids.js";import type { EventCandidate, ThoughtEvent } from "../events/types.js";import type { JazzThoughtStore } from "../jazz/store.js";import type { SourceCursor } from "../store/types.js";
export interface RssConnectorOptions { id: string; feedUrl: string; privacy?: "public-source" | "private" | "sensitive"; timeoutMs?: number; maxItems?: number;}
export interface RssPollResult { status: "updated" | "not-modified"; inserted: number; unchanged: number; events: ThoughtEvent[]; cursor: SourceCursor;}
interface ParsedFeedItem { identity: string; title?: string; canonicalUrl?: string; author?: string; publishedAt?: string; updatedAt?: string; categories: string[]; summary?: string;}
export class RssConnector { readonly kind = "rss" as const; readonly id: string; private readonly feedUrl: string; private readonly privacy: "public-source" | "private" | "sensitive"; private readonly timeoutMs: number; private readonly maxItems: number;
constructor(options: RssConnectorOptions) { const url = new URL(options.feedUrl); if (url.protocol !== "http:" && url.protocol !== "https:") throw new Error("RSS feed URL must use http or https"); this.id = options.id; this.feedUrl = url.href; this.privacy = options.privacy ?? "public-source"; this.timeoutMs = options.timeoutMs ?? 15_000; this.maxItems = options.maxItems ?? 200; }
describe(): JsonObject { return { id: this.id, kind: this.kind, feedUrl: this.feedUrl, privacy: this.privacy, maxItems: this.maxItems }; }
async poll(store: JazzThoughtStore, fetchImpl: typeof fetch = fetch): Promise<RssPollResult> { const cursorId = `cursor:${this.id}`; const prior = await store.getSourceCursor(cursorId); const startedAt = new Date().toISOString(); const correlationId = newId("rss_poll"); await store.appendEvent(this.connectorEvent("started", correlationId, startedAt, { status: "started", feedUrl: this.feedUrl, })); const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), this.timeoutMs); try { const headers = new Headers({ accept: "application/atom+xml, application/rss+xml, application/xml, text/xml;q=0.9" }); if (typeof prior?.cursor.etag === "string") headers.set("if-none-match", prior.cursor.etag); if (typeof prior?.cursor.lastModified === "string") headers.set("if-modified-since", prior.cursor.lastModified); const response = await fetchImpl(this.feedUrl, { method: "GET", headers, signal: controller.signal, redirect: "follow" }); if (response.status === 304) { const completedAt = new Date().toISOString(); const cursor = this.successCursor(cursorId, prior?.cursor ?? {}, response, completedAt); await store.appendProducerBatch([ this.cursorEvent(correlationId, completedAt, cursor), this.connectorEvent("completed", correlationId, completedAt, { status: "not-modified", inserted: 0, unchanged: 0, httpStatus: 304, }), ], cursor); return { status: "not-modified", inserted: 0, unchanged: 0, events: [], cursor }; } if (!response.ok) { const retryAfter = response.headers.get("retry-after"); throw new Error(`RSS poll failed with HTTP ${response.status}${retryAfter ? `; Retry-After ${retryAfter}` : ""}`); } const body = await response.text(); const { title: feedTitle, items } = parseFeed(body, this.feedUrl); const candidates: EventCandidate[] = []; for (const item of items.slice(0, this.maxItems)) { const occurredAt = item.publishedAt ?? item.updatedAt ?? startedAt; candidates.push({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: this.id, sourceKind: "rss", externalId: item.identity, idempotencyKey: sha256(`${this.feedUrl}\n${item.identity}`), occurredAt, actor: this.id, correlationId, privacy: this.privacy, payload: compact({ feedUrl: this.feedUrl, feedTitle, identity: item.identity, title: item.title, canonicalUrl: item.canonicalUrl, author: item.author, publishedAt: item.publishedAt, updatedAt: item.updatedAt, categories: item.categories, summary: item.summary, }), }); } const completedAt = new Date().toISOString(); const cursor = this.successCursor(cursorId, prior?.cursor ?? {}, response, completedAt); const sourceCount = candidates.length; const batch = await store.appendProducerBatch([ ...candidates, this.cursorEvent(correlationId, completedAt, cursor), this.connectorEvent("completed", correlationId, completedAt, { status: "updated", offered: sourceCount, httpStatus: response.status, }), ], cursor); const sourceEvents = batch.events.slice(0, sourceCount); const insertedIds = new Set(batch.inserted.map((event) => event.id)); const events = sourceEvents.filter((event) => insertedIds.has(event.id)); const inserted = events.length; const unchanged = sourceCount - inserted; return { status: "updated", inserted, unchanged, events, cursor }; } catch (error) { const failedAt = new Date().toISOString(); const message = error instanceof Error ? error.message : String(error); const failedCursor: SourceCursor = { id: cursorId, source: this.id, cursor: prior?.cursor ?? {}, ...(prior?.lastSuccessAt ? { lastSuccessAt: prior.lastSuccessAt } : {}), lastFailureAt: failedAt, lastError: message, updatedAt: failedAt, }; await store.appendProducerBatch([this.connectorEvent("failed", correlationId, failedAt, { status: "failed", error: message, feedUrl: this.feedUrl, })], failedCursor); throw error; } finally { clearTimeout(timer); } }
private successCursor(id: string, previous: JsonObject, response: Response, at: string): SourceCursor { const notModified = response.status === 304; return { id, source: this.id, cursor: compact({ ...previous, etag: response.headers.get("etag") ?? (notModified ? previous.etag : undefined), lastModified: response.headers.get("last-modified") ?? (notModified ? previous.lastModified : undefined), lastPollAt: at, }), lastSuccessAt: at, updatedAt: at, }; }
private cursorEvent(correlationId: string, at: string, cursor: SourceCursor): EventCandidate { return { type: "stream.thought.connector.cursor.advanced", schemaVersion: 1, source: this.id, sourceKind: "rss", externalId: cursor.id, idempotencyKey: `${correlationId}:cursor`, occurredAt: at, actor: this.id, correlationId, privacy: this.privacy, payload: { source: this.id, cursorId: cursor.id, cursor: cursor.cursor }, }; }
private connectorEvent( phase: "started" | "completed" | "failed", correlationId: string, at: string, payload: JsonObject, ): EventCandidate { return { type: phase === "failed" ? "stream.thought.connector.failed" : `stream.thought.connector.poll.${phase}`, schemaVersion: 1, source: this.id, sourceKind: "rss", externalId: correlationId, idempotencyKey: `${correlationId}:${phase}`, occurredAt: at, actor: this.id, correlationId, privacy: this.privacy, payload, }; }}
export function parseFeed(xml: string, feedUrl: string): { title?: string; items: ParsedFeedItem[] } { const parser = new XMLParser({ ignoreAttributes: false, trimValues: true, parseTagValue: false }); const document = parser.parse(xml) as Record<string, unknown>; const rssChannel = object(object(document.rss)?.channel); const atomFeed = object(document.feed); const rdfFeed = object(document["rdf:RDF"]); const container = rssChannel ?? atomFeed ?? rdfFeed; if (!container) throw new Error("RSS/Atom response has no recognized feed root"); const rawItems = rssChannel ? array(rssChannel.item) : atomFeed ? array(atomFeed.entry) : array(rdfFeed?.item); const title = text(container.title); return { ...(title ? { title } : {}), items: rawItems.map((raw) => parseItem(object(raw) ?? {}, feedUrl, Boolean(atomFeed))), };}
function parseItem(item: Record<string, unknown>, feedUrl: string, atom: boolean): ParsedFeedItem { const canonicalUrl = atom ? atomLink(item.link, feedUrl) : resolveUrl(text(item.link), feedUrl); const rawIdentity = text(item.guid) ?? text(item.id) ?? canonicalUrl; const summary = text(item.summary) ?? text(item.description) ?? text(item.content) ?? text(item["content:encoded"]); const identity = rawIdentity ?? `content:${sha256(JSON.stringify({ title: text(item.title), summary }))}`; const authorValue = object(item.author); const author = text(authorValue?.name) ?? text(item.author) ?? text(item["dc:creator"]); const categories = array(item.category).map((category) => text(object(category)?.["@_term"]) ?? text(category)).filter(isString); return compact({ identity, title: text(item.title), canonicalUrl, author, publishedAt: normalizeDate(text(item.published) ?? text(item.pubDate) ?? text(item["dc:date"])), updatedAt: normalizeDate(text(item.updated)), categories: [...new Set(categories)], summary, }) as unknown as ParsedFeedItem;}
function atomLink(value: unknown, base: string): string | undefined { const links = array(value).map(object).filter((link): link is Record<string, unknown> => Boolean(link)); const selected = links.find((link) => !link["@_rel"] || link["@_rel"] === "alternate") ?? links[0]; return resolveUrl(text(selected?.["@_href"]) ?? text(value), base);}
function normalizeDate(value?: string): string | undefined { if (!value) return undefined; const timestamp = Date.parse(value); return Number.isFinite(timestamp) ? new Date(timestamp).toISOString() : undefined;}
function resolveUrl(value: string | undefined, base: string): string | undefined { if (!value) return undefined; try { const url = new URL(value, base); return url.protocol === "http:" || url.protocol === "https:" ? url.href : undefined; } catch { return undefined; }}
function text(value: unknown): string | undefined { if (typeof value === "string") return value.trim() || undefined; if (typeof value === "number") return String(value); const record = object(value); return record ? text(record["#text"]) : undefined;}
function object(value: unknown): Record<string, unknown> | undefined { return value && typeof value === "object" && !Array.isArray(value) ? value as Record<string, unknown> : undefined;}
function array(value: unknown): unknown[] { return value === undefined ? [] : Array.isArray(value) ? value : [value];}
function isString(value: string | undefined): value is string { return Boolean(value);}
function compact(value: Record<string, unknown>): JsonObject { return Object.fromEntries(Object.entries(value).filter(([, item]) => item !== undefined)) as JsonObject;}