Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
4.8 kB · 152 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153import fs from "node:fs/promises";import { afterEach, describe, expect, test } from "vitest";import { buildSourceHealth } from "../src/projections/source-health.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import { temporaryProject, testStore } from "./helpers.js";
const stores: JazzThoughtStore[] = [];const roots: string[] = [];
afterEach(async () => { 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("source health projection", () => { test("distinguishes healthy, processing, and unknown durable source state", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store);
await store.upsertSourceCursor({ id: "cursor:rss:healthy", source: "rss:healthy", cursor: { etag: "v2" }, lastSuccessAt: "2026-07-14T00:00:03.000Z", lastFailureAt: "2026-07-14T00:00:01.000Z", updatedAt: "2026-07-14T00:00:03.000Z", }); await store.upsertSourceCursor({ id: "cursor:rss:unknown", source: "rss:unknown", cursor: {}, updatedAt: "2026-07-14T00:00:00.000Z", }); await appendConnectorEvent(store, "rss:healthy", "stream.thought.connector.poll.started", "healthy-poll"); await appendConnectorEvent(store, "rss:healthy", "stream.thought.connector.poll.completed", "healthy-poll"); await appendConnectorEvent(store, "rss:polling", "stream.thought.connector.poll.started", "open-poll");
expect(await buildSourceHealth(store)).toEqual([ expect.objectContaining({ source: "rss:healthy", status: "healthy", recordCount: 2, operationsStarted: 1, operationsCompleted: 1, inFlight: 0, cursor: { etag: "v2" }, }), expect.objectContaining({ source: "rss:polling", status: "processing", recordCount: 1, operationsStarted: 1, operationsCompleted: 0, inFlight: 1, cursor: {}, }), expect.objectContaining({ source: "rss:unknown", status: "unknown", recordCount: 0, operationsStarted: 0, inFlight: 0, }), ]); });
test("shows configured runtime and upstream state before the first activity event", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const now = "2026-08-10T23:12:00.000Z";
await store.registerSource({ id: "x:public-watch", kind: "x-webhook", enabled: true, updatedAt: now, config: { lane: "public-watchlist", control: { version: 1, source: "x:public-watch", kind: "x-webhook", enabled: true, lane: "public-watchlist", expectedSubscriptionCount: 12, eventTypes: ["post.create", "post.delete"], configuredAt: now, updatedAt: now, runtime: { state: "ready", readyAt: now, revision: "x-activity-v2-webhook-v1" }, upstream: { registered: true, valid: true, webhookId: "2086948647707852801", desiredSubscriptionCount: 12, liveSubscriptionCount: 12, subscriptionsConverged: true, checkedAt: now, }, }, }, });
expect(await buildSourceHealth(store, Date.parse(now))).toEqual([ expect.objectContaining({ source: "x:public-watch", kind: "x-webhook", configured: true, enabled: true, status: "unknown", hasActivityEvidence: false, recordCount: 0, control: { lane: "public-watchlist", runtime: expect.objectContaining({ state: "ready", current: true }), upstream: expect.objectContaining({ registered: true, valid: true, liveSubscriptionCount: 12, desiredSubscriptionCount: 12, subscriptionsConverged: true, current: true, }), }, }), ]); });});
async function appendConnectorEvent( store: JazzThoughtStore, source: string, type: "stream.thought.connector.poll.started" | "stream.thought.connector.poll.completed", correlationId: string,): Promise<void> { await store.appendEvent({ type, schemaVersion: 1, source, sourceKind: "rss", externalId: correlationId, idempotencyKey: `${correlationId}:${type}`, occurredAt: "2026-07-14T00:00:00.000Z", actor: source, correlationId, privacy: "public-source", payload: { status: type.endsWith("started") ? "started" : "completed" }, });}