import { 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 { 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; 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, 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 => 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 | undefined { return value && typeof value === "object" && !Array.isArray(value) ? value as Record : 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): JsonObject { return Object.fromEntries(Object.entries(value).filter(([, item]) => item !== undefined)) as JsonObject; }