Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
4.3 kB · 90 lines
TypeScript
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091import fs from "node:fs/promises";import path from "node:path";import { afterEach, describe, expect, test } from "vitest";import { FastmailConnector, parseFastmailCapture } from "../src/connectors/fastmail.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("FastmailConnector", () => { test("parses captured JMAP changes while excluding full bodies", async () => { const capture = await fixture(); const parsed = parseFastmailCapture(capture); expect(parsed.changes).toMatchObject({ accountId: "account-fixture", oldState: "email-state-1", newState: "email-state-2", created: ["email-created"], updated: ["email-updated"], destroyed: ["email-destroyed"], }); expect(parsed.emails.get("email-created")).not.toHaveProperty("bodyValues"); });
test("persists source events before cursor and absorbs exact capture replay", async () => { const project = await temporaryProject(); roots.push(project); const connector = new FastmailConnector({ id: "fastmail:fixture", accountId: "account-fixture" }); const firstStore = await testStore(project); const first = await connector.ingestCapture(firstStore, await fixture()); expect(first).toMatchObject({ inserted: 3, unchanged: 0, created: 1, updated: 1, destroyed: 1, hasMoreChanges: false }); expect(first.events.map((event) => event.payload.operation)).toEqual(["created", "updated", "destroyed"]); expect(first.events[0]).toMatchObject({ privacy: "sensitive", sourceKind: "fastmail", payload: { accountId: "account-fixture", emailState: "email-state-2", subject: "Synthetic fixture message", preview: "A synthetic preview with no private correspondence.", attachments: [{ blobId: "blob-fixture-1", name: "fixture.pdf", size: 2048 }], }, }); expect(first.events[0]?.payload).not.toHaveProperty("bodyValues"); await firstStore.close();
const restartedStore = await testStore(project); stores.push(restartedStore); const replay = await connector.ingestCapture(restartedStore, await fixture()); expect(replay).toMatchObject({ inserted: 0, unchanged: 3 }); expect(await restartedStore.listEvents({ types: ["stream.thought.source.email.observed"] })).toHaveLength(3); expect((await restartedStore.getSourceCursor("cursor:fastmail:fixture"))?.cursor).toEqual({ accountIdHash: expect.any(String), emailState: "email-state-2", queryState: "query-state-2", }); });
test("fails closed on state gaps and preserves the durable cursor", async () => { const project = await temporaryProject(); roots.push(project); const connector = new FastmailConnector({ id: "fastmail:fixture", accountId: "account-fixture" }); const store = await testStore(project); stores.push(store); await connector.ingestCapture(store, await fixture()); const before = await store.getSourceCursor("cursor:fastmail:fixture"); const gap = await fixture() as { methodResponses: Array<[string, Record<string, unknown>, string]> }; const changes = gap.methodResponses.find(([name]) => name === "Email/changes")?.[1]; if (!changes) throw new Error("Missing fixture changes response"); changes.oldState = "wrong-old-state"; changes.newState = "email-state-3"; const get = gap.methodResponses.find(([name]) => name === "Email/get")?.[1]; if (!get) throw new Error("Missing fixture get response"); get.state = "email-state-3"; await expect(connector.ingestCapture(store, gap)).rejects.toThrow("oldState does not match the durable cursor"); const after = await store.getSourceCursor("cursor:fastmail:fixture"); expect(after?.cursor).toEqual(before?.cursor); expect(after?.lastFailureAt).toBeDefined(); });});
async function fixture(): Promise<unknown> { return JSON.parse(await fs.readFile(path.join(process.cwd(), "fixtures", "fastmail", "changes.json"), "utf8")) as unknown;}