Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
15 kB · 384 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385import { canonicalJson, 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";import { z } from "zod";
const nullableOptionalName = z.string().max(1_024).nullable().optional().transform((value) => value ?? undefined);const nullableOptionalMediaType = z.string().max(255).nullable().optional().transform((value) => value ?? undefined);const nullableOptionalDisposition = z.string().max(100).nullable().optional().transform((value) => value ?? undefined);const nullableOptionalCid = z.string().max(1_024).nullable().optional().transform((value) => value ?? undefined);
const addressSchema = z.object({ name: nullableOptionalName, email: z.string().min(1).max(500),});
const attachmentSchema = z.object({ blobId: z.string().min(1).max(500), name: nullableOptionalName, type: nullableOptionalMediaType, size: z.number().int().nonnegative(), disposition: nullableOptionalDisposition, cid: nullableOptionalCid,});
const addressListSchema = z.array(addressSchema).max(100).nullable().optional().transform((value) => value ?? undefined);const attachmentListSchema = z.array(attachmentSchema).max(100).nullable().optional().transform((value) => value ?? undefined);const booleanRecordSchema = z.record(z.string().min(1).max(500), z.boolean()).superRefine((value, context) => { if (Object.keys(value).length > 500) context.addIssue({ code: "custom", message: "Fastmail boolean record exceeds 500 entries" });});
export const fastmailEmailSchema = z.object({ id: z.string().min(1).max(500), threadId: z.string().min(1).max(500), mailboxIds: booleanRecordSchema, keywords: booleanRecordSchema, from: addressListSchema, to: addressListSchema, cc: addressListSchema, bcc: addressListSchema, replyTo: addressListSchema, subject: z.string().max(8_192), receivedAt: z.iso.datetime({ offset: true }), preview: z.string().max(8_192), hasAttachment: z.boolean().optional(), attachments: attachmentListSchema,});
const changesSchema = z.object({ accountId: z.string().min(1), oldState: z.string(), newState: z.string().min(1), hasMoreChanges: z.boolean(), created: z.array(z.string().min(1)), updated: z.array(z.string().min(1)), destroyed: z.array(z.string().min(1)),}).passthrough();
const queryChangesSchema = z.object({ accountId: z.string().min(1), oldQueryState: z.string(), newQueryState: z.string().min(1), removed: z.array(z.string().min(1)), added: z.array(z.object({ id: z.string().min(1), index: z.number().int().nonnegative() }).passthrough()), total: z.number().int().nonnegative().optional(),}).passthrough();
const getSchema = z.object({ accountId: z.string().min(1), state: z.string().min(1), list: z.array(fastmailEmailSchema), notFound: z.array(z.string()).optional(),}).passthrough();
const methodResponseSchema = z.tuple([ z.string().min(1), z.record(z.string(), z.unknown()), z.string().min(1),]);
const captureSchema = z.object({ methodResponses: z.array(methodResponseSchema).min(2),}).passthrough();
export type FastmailEmail = z.infer<typeof fastmailEmailSchema>;type Changes = z.infer<typeof changesSchema>;type QueryChanges = z.infer<typeof queryChangesSchema>;
export interface FastmailConnectorOptions { id: string; accountId: string;}
export interface FastmailCaptureResult { inserted: number; unchanged: number; created: number; updated: number; destroyed: number; hasMoreChanges: boolean; events: ThoughtEvent[]; cursor: SourceCursor;}
interface ParsedCapture { changes: Changes; queryChanges?: QueryChanges; emails: Map<string, FastmailEmail>;}
export class FastmailConnector { readonly kind = "fastmail" as const; readonly id: string; readonly accountId: string;
constructor(options: FastmailConnectorOptions) { this.id = required(options.id, "Fastmail connector id"); this.accountId = required(options.accountId, "Fastmail account id"); }
describe(): JsonObject { return { id: this.id, kind: this.kind, accountIdHash: sha256(this.accountId), transport: "captured-jmap-response", authority: "read-only", fullBodies: false, }; }
async ingestCapture(store: JazzThoughtStore, raw: unknown): Promise<FastmailCaptureResult> { const cursorId = `cursor:${this.id}`; const prior = await store.getSourceCursor(cursorId); const correlationId = newId("fastmail_capture"); const startedAt = new Date().toISOString(); await this.appendConnectorEvent(store, "started", correlationId, startedAt, { status: "started", transport: "captured-jmap-response", accountIdHash: sha256(this.accountId), });
try { const capture = parseFastmailCapture(raw); this.assertCaptureCompatible(capture, prior); const operations = [ ...capture.changes.created.map((id) => ({ id, operation: "created" as const })), ...capture.changes.updated.map((id) => ({ id, operation: "updated" as const })), ...capture.changes.destroyed.map((id) => ({ id, operation: "destroyed" as const })), ]; const candidates: EventCandidate[] = []; for (const item of operations) { const email = item.operation === "destroyed" ? undefined : capture.emails.get(item.id); if (item.operation !== "destroyed" && !email) { throw new Error(`JMAP Email/get response is missing changed email ${item.id}`); } candidates.push(this.eventCandidate( item.id, item.operation, capture.changes.newState, email, correlationId, startedAt, )); }
const completedAt = new Date().toISOString(); const cursor: SourceCursor = { id: cursorId, source: this.id, cursor: compact({ accountIdHash: sha256(this.accountId), emailState: capture.changes.newState, queryState: capture.queryChanges?.newQueryState ?? prior?.cursor.queryState, }), lastSuccessAt: completedAt, ...(prior?.lastFailureAt ? { lastFailureAt: prior.lastFailureAt } : {}), updatedAt: completedAt, }; const sourceCount = candidates.length; candidates.push(this.cursorEvent(correlationId, completedAt, cursor)); if (prior?.lastFailureAt && (!prior.lastSuccessAt || prior.lastFailureAt > prior.lastSuccessAt)) { candidates.push(this.connectorEvent("recovered", correlationId, completedAt, { status: "recovered", previousFailureAt: prior.lastFailureAt, })); } candidates.push(this.connectorEvent("completed", correlationId, completedAt, { status: "updated", created: capture.changes.created.length, updated: capture.changes.updated.length, destroyed: capture.changes.destroyed.length, hasMoreChanges: capture.changes.hasMoreChanges, })); const batch = await store.appendProducerBatch(candidates, cursor); const offeredEvents = batch.events.slice(0, sourceCount); const insertedIds = new Set(batch.inserted.map((event) => event.id)); const events = offeredEvents.filter((event) => insertedIds.has(event.id)); const inserted = events.length; const unchanged = sourceCount - inserted; return { inserted, unchanged, created: capture.changes.created.length, updated: capture.changes.updated.length, destroyed: capture.changes.destroyed.length, hasMoreChanges: capture.changes.hasMoreChanges, 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 ?? { accountIdHash: sha256(this.accountId) }, ...(prior?.lastSuccessAt ? { lastSuccessAt: prior.lastSuccessAt } : {}), lastFailureAt: failedAt, lastError: message, updatedAt: failedAt, }; await store.appendProducerBatch([this.connectorEvent("failed", correlationId, failedAt, { status: "failed", error: message, phase: "captured-jmap-ingest", accountIdHash: sha256(this.accountId), })], failedCursor); throw error; } }
private assertCaptureCompatible(capture: ParsedCapture, prior: SourceCursor | undefined): void { if (capture.changes.accountId !== this.accountId) throw new Error("JMAP Email/changes accountId does not match connector configuration"); if (capture.queryChanges && capture.queryChanges.accountId !== this.accountId) throw new Error("JMAP Email/queryChanges accountId does not match connector configuration"); const accountIdHash = prior?.cursor.accountIdHash; if (accountIdHash !== undefined && accountIdHash !== sha256(this.accountId)) throw new Error("Fastmail cursor belongs to a different account"); const emailState = stringCursor(prior, "emailState"); const replay = emailState === capture.changes.newState; if (emailState !== undefined && !replay && capture.changes.oldState !== emailState) { throw new Error("JMAP Email/changes oldState does not match the durable cursor"); } const queryState = stringCursor(prior, "queryState"); if (capture.queryChanges && queryState !== undefined && queryState !== capture.queryChanges.newQueryState && capture.queryChanges.oldQueryState !== queryState) { throw new Error("JMAP Email/queryChanges oldQueryState does not match the durable cursor"); } }
private eventCandidate( id: string, operation: "created" | "updated" | "destroyed", emailState: string, email: FastmailEmail | undefined, correlationId: string, observedAt: string, ): EventCandidate { return { type: "stream.thought.source.email.observed", schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: id, idempotencyKey: sha256(canonicalJson({ accountId: this.accountId, emailId: id, emailState, operation })), occurredAt: email?.receivedAt ?? observedAt, actor: email?.from?.[0]?.email ?? this.id, correlationId, privacy: "sensitive", payload: compact({ accountId: this.accountId, emailId: id, emailState, operation, threadId: email?.threadId, mailboxIds: email?.mailboxIds, keywords: email?.keywords, from: email?.from, to: email?.to, cc: email?.cc, bcc: email?.bcc, replyTo: email?.replyTo, subject: email?.subject, receivedAt: email?.receivedAt, preview: email?.preview, hasAttachment: email?.hasAttachment, attachments: email?.attachments?.map((attachment) => compact({ blobId: attachment.blobId, name: attachment.name, type: attachment.type, size: attachment.size, disposition: attachment.disposition, cid: attachment.cid, })), }), }; }
private cursorEvent(correlationId: string, at: string, cursor: SourceCursor): EventCandidate { return { type: "stream.thought.connector.cursor.advanced", schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: cursor.id, idempotencyKey: `${correlationId}:cursor`, occurredAt: at, actor: this.id, correlationId, privacy: "sensitive", payload: { source: this.id, cursorId: cursor.id, cursor: cursor.cursor }, }; }
private connectorEvent( phase: "started" | "completed" | "failed" | "recovered", correlationId: string, at: string, payload: JsonObject, ): EventCandidate { const type = phase === "failed" ? "stream.thought.connector.failed" : phase === "recovered" ? "stream.thought.connector.recovered" : `stream.thought.connector.poll.${phase}`; return { type, schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: correlationId, idempotencyKey: `${correlationId}:${phase}`, occurredAt: at, actor: this.id, correlationId, privacy: "sensitive", payload, }; }
private async appendConnectorEvent( store: JazzThoughtStore, phase: "started" | "completed" | "failed" | "recovered", correlationId: string, at: string, payload: JsonObject, ): Promise<void> { await store.appendEvent(this.connectorEvent(phase, correlationId, at, payload)); }}
export function parseFastmailCapture(raw: unknown): ParsedCapture { const capture = captureSchema.parse(raw); const byName = new Map(capture.methodResponses.map((response) => [response[0], response[1]])); const changesRaw = byName.get("Email/changes"); const getRaw = byName.get("Email/get"); if (!changesRaw || !getRaw) throw new Error("Captured JMAP response requires Email/changes and Email/get method responses"); const changes = changesSchema.parse(changesRaw); const get = getSchema.parse(getRaw); if (get.accountId !== changes.accountId) throw new Error("JMAP Email/get and Email/changes accountId values differ"); if (get.state !== changes.newState) throw new Error("JMAP Email/get state does not match Email/changes newState"); const queryRaw = byName.get("Email/queryChanges"); const queryChanges = queryRaw ? queryChangesSchema.parse(queryRaw) : undefined; const emails = new Map(get.list.map((email) => [email.id, email])); if (emails.size !== get.list.length) throw new Error("JMAP Email/get response contains duplicate email ids"); return { changes, ...(queryChanges ? { queryChanges } : {}), emails };}
function stringCursor(cursor: SourceCursor | undefined, key: string): string | undefined { const value = cursor?.cursor[key]; if (value === undefined) return undefined; if (typeof value !== "string" || !value) throw new Error(`Fastmail cursor ${key} must be a nonempty string`); return value;}
function required(value: string, label: string): string { const normalized = value.trim(); if (!normalized) throw new Error(`${label} is required`); return normalized;}
function compact(value: Record<string, unknown>): JsonObject { return Object.fromEntries(Object.entries(value).filter(([, child]) => child !== undefined)) as JsonObject;}