Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
4.6 kB · 100 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101import fs from "node:fs/promises";import http from "node:http";import path from "node:path";import { afterEach, describe, expect, test } from "vitest";import { parseFeed, RssConnector } from "../src/connectors/rss.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import { temporaryProject, testStore } from "./helpers.js";
const stores: JazzThoughtStore[] = [];const roots: string[] = [];const servers: http.Server[] = [];
afterEach(async () => { await Promise.all(servers.splice(0).map((server) => new Promise<void>((resolve) => { server.closeAllConnections(); server.close(() => resolve()); }))); await Promise.all(stores.splice(0).map((store) => store.close())); await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true })));});
describe("RssConnector", () => { test("parses captured Atom identity, links, author, categories, and dates", async () => { const atom = await fs.readFile(path.join(process.cwd(), "fixtures", "rss", "feed.atom"), "utf8"); expect(parseFeed(atom, "https://example.test/feed.atom")).toEqual({ title: "thought stream Atom fixture", items: [{ identity: "tag:example.test,2026:local-evidence", title: "Local evidence fabric", canonicalUrl: "https://example.test/atom/local-evidence", author: "Co", updatedAt: "2026-07-14T12:00:00.000Z", categories: ["agents"], summary: "One captured Atom entry.", }], }); });
test("persists validators only after durable items and resumes with a conditional request", async () => { const body = await fs.readFile(path.join(process.cwd(), "fixtures", "rss", "feed.xml"), "utf8"); const seenHeaders: Array<{ etag: string | null; modified: string | null }> = []; let requestCount = 0; const server = http.createServer((request, response) => { requestCount += 1; seenHeaders.push({ etag: request.headers["if-none-match"] ?? null, modified: request.headers["if-modified-since"] ?? null, }); if (requestCount === 1) { response.writeHead(200, { "content-type": "application/rss+xml", etag: '"fixture-v1"', "last-modified": "Tue, 14 Jul 2026 11:00:00 GMT", }); response.end(body); } else if (requestCount === 2) { response.writeHead(304); response.end(); } else { response.writeHead(503, { "retry-after": "120" }); response.end(); } }); servers.push(server); await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve)); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing RSS fixture server address"); const project = await temporaryProject(); roots.push(project); const connector = new RssConnector({ id: "rss:fixture", feedUrl: `http://127.0.0.1:${address.port}/feed.xml` });
const firstStore = await testStore(project); const first = await connector.poll(firstStore); expect(first).toMatchObject({ status: "updated", inserted: 2, unchanged: 0 }); expect(first.events.map((event) => event.payload.identity)).toEqual(["fixture-guid-1", "http://127.0.0.1:PORT/cursor-discipline"].map((value) => value.replace("PORT", String(address.port)))); expect((await firstStore.getSourceCursor("cursor:rss:fixture"))?.cursor).toMatchObject({ etag: '"fixture-v1"', lastModified: "Tue, 14 Jul 2026 11:00:00 GMT", }); await firstStore.close();
const restartedStore = await testStore(project); stores.push(restartedStore); const second = await connector.poll(restartedStore); expect(second).toMatchObject({ status: "not-modified", inserted: 0, unchanged: 0 }); expect(seenHeaders).toEqual([ { etag: null, modified: null }, { etag: '"fixture-v1"', modified: "Tue, 14 Jul 2026 11:00:00 GMT" }, ]); expect(await restartedStore.listEvents({ types: ["stream.thought.source.rss.item"] })).toHaveLength(2);
await expect(connector.poll(restartedStore)).rejects.toThrow("HTTP 503; Retry-After 120"); const failedCursor = await restartedStore.getSourceCursor("cursor:rss:fixture"); expect(failedCursor?.cursor).toMatchObject({ etag: '"fixture-v1"' }); expect(failedCursor?.lastSuccessAt).toBe(second.cursor.lastSuccessAt); expect(failedCursor?.lastFailureAt).toBeDefined(); expect((await restartedStore.listEvents({ types: ["stream.thought.connector.failed"] })).at(-1)?.payload.error) .toBe("RSS poll failed with HTTP 503; Retry-After 120"); });});