Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
34 kB · 859 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860import { canonicalJson, sha256, type JsonObject, type JsonValue } 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 { fastmailEmailSchema, type FastmailEmail } from "./fastmail.js";import { z } from "zod";
const CORE_CAPABILITY = "urn:ietf:params:jmap:core";const MAIL_CAPABILITY = "urn:ietf:params:jmap:mail";const FASTMAIL_SESSION_URL = "https://api.fastmail.com/jmap/session";const FASTMAIL_SESSION_ORIGIN = "https://api.fastmail.com";const CURSOR_REVISION = "fastmail-jmap-v1";const ALLOWED_METHODS = new Set(["Email/get", "Email/query", "Email/changes", "Email/queryChanges"]);const EMAIL_PROPERTIES = [ "id", "threadId", "mailboxIds", "keywords", "from", "to", "cc", "bcc", "replyTo", "subject", "receivedAt", "preview", "hasAttachment", "attachments",] as const;
const sessionSchema = z.object({ apiUrl: z.url(), capabilities: z.record(z.string(), z.unknown()), accounts: z.record(z.string(), z.object({ accountCapabilities: z.record(z.string(), z.unknown()), }).passthrough()), primaryAccounts: z.record(z.string(), z.string()),}).passthrough();
const methodResponseSchema = z.tuple([ z.string().min(1), z.record(z.string(), z.unknown()), z.string().min(1),]);
const jmapResponseSchema = z.object({ methodResponses: z.array(methodResponseSchema), sessionState: z.string().optional(),}).passthrough();
const emailChangesSchema = z.object({ accountId: z.string().min(1), oldState: z.string().min(1), 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 emailQueryChangesSchema = z.object({ accountId: z.string().min(1), oldQueryState: z.string().min(1), 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()),}).passthrough();
const emailQuerySchema = z.object({ accountId: z.string().min(1), queryState: z.string().min(1), canCalculateChanges: z.boolean(), position: z.number().int().nonnegative(), ids: z.array(z.string().min(1)),}).passthrough();
const emailGetSchema = z.object({ accountId: z.string().min(1), state: z.string().min(1), list: z.array(fastmailEmailSchema), notFound: z.array(z.string()).optional(),}).passthrough();
type EmailChanges = z.infer<typeof emailChangesSchema>;type EmailQueryChanges = z.infer<typeof emailQueryChangesSchema>;type JmapResponse = z.infer<typeof jmapResponseSchema>;type FetchLike = typeof globalThis.fetch;
interface MethodCall { name: string; arguments: Record<string, JsonValue>; id: string;}
interface FastmailSession { accountId: string; apiUrl: string; maxObjectsInGet: number;}
export interface FastmailJmapClientOptions { token: string; timeoutMs?: number; maxResponseBytes?: number; sessionUrl?: string; fetchImpl?: FetchLike;}
export interface FastmailJmapConnectorOptions { id: string; client: FastmailJmapClient; maxChanges: number; maxPages: number; resnapshotLimit: number; credentialCustody: "dedicated-mail-ingress" | "shared-operator-accepted";}
export interface FastmailJmapPollResult { source: string; initialized: boolean; resnapshot: boolean; inserted: number; unchanged: number; created: number; updated: number; destroyed: number; events: ThoughtEvent[]; cursor: SourceCursor;}
interface FastmailOperation { id: string; operation: "created" | "updated" | "destroyed"; email?: FastmailEmail;}
interface CurrentSnapshot { accountId: string; emailState: string; queryState: string; emails: FastmailEmail[];}
interface ChangeSet { accountId: string; emailState: string; queryState: string; operations: FastmailOperation[];}
export class FastmailJmapError extends Error { readonly code: string;
constructor(code: string, message: string) { super(message); this.name = "FastmailJmapError"; this.code = code; }}
export class FastmailJmapCannotCalculateChanges extends FastmailJmapError { constructor() { super("cannot-calculate-changes", "Fastmail can no longer calculate changes from the durable cursor"); this.name = "FastmailJmapCannotCalculateChanges"; }}
export class FastmailJmapClient { readonly #token: string; private readonly timeoutMs: number; private readonly maxResponseBytes: number; private readonly sessionUrl: string; private readonly fetchImpl: FetchLike; private session: FastmailSession | undefined;
constructor(options: FastmailJmapClientOptions) { this.#token = required(options.token, "Fastmail JMAP token"); this.timeoutMs = boundedPositive(options.timeoutMs ?? 15_000, 120_000, "Fastmail timeout"); this.maxResponseBytes = boundedPositive(options.maxResponseBytes ?? 4 * 1024 * 1024, 64 * 1024 * 1024, "Fastmail response bound"); this.sessionUrl = validateCredentialEndpoint(options.sessionUrl ?? FASTMAIL_SESSION_URL, "Fastmail session URL"); if (new URL(this.sessionUrl).origin !== FASTMAIL_SESSION_ORIGIN) { throw new FastmailJmapError("invalid-session-origin", "Fastmail session URL must use the pinned Fastmail API origin"); } this.fetchImpl = options.fetchImpl ?? globalThis.fetch; }
async discover(signal?: AbortSignal): Promise<FastmailSession> { if (this.session) return this.session; const raw = await this.fetchJson(this.sessionUrl, { method: "GET" }, signal); const parsed = parseWithCode(sessionSchema, raw, "invalid-session"); if (!(CORE_CAPABILITY in parsed.capabilities) || !(MAIL_CAPABILITY in parsed.capabilities)) { throw new FastmailJmapError("missing-mail-capability", "Fastmail session does not advertise required JMAP capabilities"); } const accountId = parsed.primaryAccounts[MAIL_CAPABILITY]; if (!accountId || !(MAIL_CAPABILITY in (parsed.accounts[accountId]?.accountCapabilities ?? {}))) { throw new FastmailJmapError("missing-primary-account", "Fastmail session does not expose one primary mail account"); } const apiUrl = validateCredentialEndpoint(parsed.apiUrl, "Fastmail API URL"); if (!isAllowedFastmailApiUrl(apiUrl)) { throw new FastmailJmapError("cross-origin-api", "Fastmail API URL is outside the pinned Fastmail API origins"); } const maxObjectsInGet = capabilityInteger(parsed.capabilities[CORE_CAPABILITY], "maxObjectsInGet", 1, 10_000); this.session = { accountId, apiUrl, maxObjectsInGet }; return this.session; }
async currentSnapshot(limit: number, signal?: AbortSignal): Promise<CurrentSnapshot> { const session = await this.discover(signal); const boundedLimit = boundedNonnegative(limit, 1_000, "Fastmail snapshot limit"); const effectiveLimit = Math.min(boundedLimit, session.maxObjectsInGet); const response = await this.request(session, [ { name: "Email/query", id: "snapshot-query", arguments: { accountId: session.accountId, position: 0, limit: effectiveLimit, calculateTotal: false, sort: [{ property: "receivedAt", isAscending: false }], }, }, { name: "Email/get", id: "snapshot-get", arguments: { accountId: session.accountId, ...(effectiveLimit === 0 ? { ids: [] } : { "#ids": { resultOf: "snapshot-query", name: "Email/query", path: "/ids" } }), properties: [...EMAIL_PROPERTIES], }, }, ], signal); const query = methodResult(response, "snapshot-query", "Email/query", emailQuerySchema); const get = methodResult(response, "snapshot-get", "Email/get", emailGetSchema); assertAccount(session.accountId, query.accountId); assertAccount(session.accountId, get.accountId); if (!query.canCalculateChanges) { throw new FastmailJmapError("query-changes-unavailable", "Fastmail cannot calculate query changes for the configured email query"); } const requestedIds = new Set(query.ids.slice(0, effectiveLimit)); const accountedIds = new Set([...get.list.map((email) => email.id), ...(get.notFound ?? [])]); for (const email of get.list) { if (!requestedIds.has(email.id)) throw new FastmailJmapError("snapshot-id-mismatch", "Fastmail snapshot returned an email outside the requested query result"); } for (const id of requestedIds) { if (!accountedIds.has(id)) throw new FastmailJmapError("snapshot-incomplete", "Fastmail snapshot omitted a requested email without notFound evidence"); } const admittedIds = new Set(query.ids.slice(0, effectiveLimit)); return { accountId: session.accountId, emailState: get.state, queryState: query.queryState, emails: effectiveLimit === 0 ? [] : get.list.filter((email) => admittedIds.has(email.id)).slice(0, effectiveLimit), }; }
async changes( emailState: string, queryState: string, maxChanges: number, maxPages: number, signal?: AbortSignal, ): Promise<ChangeSet> { const session = await this.discover(signal); const pageSize = boundedPositive(maxChanges, 1_000, "Fastmail max changes"); const pageLimit = boundedPositive(maxPages, 100, "Fastmail max pages"); const email = await this.collectEmailChanges(session, emailState, pageSize, pageLimit, signal); const query = await this.collectQueryChanges(session, queryState, pageSize, signal); const operations = collapseOperations(email.pages); const changedIds = [...operations.entries()].filter(([, operation]) => operation !== "destroyed").map(([id]) => id).sort(); const details = await this.getEmails(session, changedIds, pageSize, signal); for (const id of details.notFound) operations.set(id, "destroyed"); for (const id of changedIds) { if (!details.emails.has(id) && !details.notFound.has(id)) { throw new FastmailJmapError("email-get-incomplete", "Fastmail Email/get omitted a changed id without notFound evidence"); } } return { accountId: session.accountId, emailState: email.state, queryState: query.state, operations: [...operations.entries()].sort(([left], [right]) => left.localeCompare(right)).map(([id, operation]) => ({ id, operation, ...(operation !== "destroyed" && details.emails.has(id) ? { email: details.emails.get(id)! } : {}), })), }; }
private async collectEmailChanges( session: FastmailSession, initialState: string, maxChanges: number, maxPages: number, signal?: AbortSignal, ): Promise<{ state: string; pages: EmailChanges[] }> { let state = required(initialState, "Fastmail email state"); const pages: EmailChanges[] = []; for (let page = 0; page < maxPages; page += 1) { const id = `email-changes-${page}`; const response = await this.request(session, [{ name: "Email/changes", id, arguments: { accountId: session.accountId, sinceState: state, maxChanges }, }], signal); const changes = methodResult(response, id, "Email/changes", emailChangesSchema); assertAccount(session.accountId, changes.accountId); if (changes.oldState !== state) throw new FastmailJmapError("email-state-gap", "Fastmail Email/changes returned a different old state"); if (changes.created.length + changes.updated.length + changes.destroyed.length > maxChanges) { throw new FastmailJmapError("email-change-bound", "Fastmail Email/changes exceeded the requested change bound"); } pages.push(changes); state = changes.newState; if (!changes.hasMoreChanges) return { state, pages }; } throw new FastmailJmapError("email-page-bound", "Fastmail Email/changes exceeded the configured page bound"); }
private async collectQueryChanges( session: FastmailSession, initialState: string, maxChanges: number, signal?: AbortSignal, ): Promise<{ state: string; response: EmailQueryChanges }> { const state = required(initialState, "Fastmail query state"); const id = "query-changes"; const response = await this.request(session, [{ name: "Email/queryChanges", id, arguments: { accountId: session.accountId, sinceQueryState: state, maxChanges, calculateTotal: false, sort: [{ property: "receivedAt", isAscending: false }], }, }], signal); const changes = methodResult(response, id, "Email/queryChanges", emailQueryChangesSchema); assertAccount(session.accountId, changes.accountId); if (changes.oldQueryState !== state) throw new FastmailJmapError("query-state-gap", "Fastmail Email/queryChanges returned a different old state"); if (changes.removed.length + changes.added.length > maxChanges) { throw new FastmailJmapError("query-change-bound", "Fastmail Email/queryChanges exceeded the requested change bound"); } return { state: changes.newQueryState, response: changes }; }
private async getEmails( session: FastmailSession, ids: string[], chunkSize: number, signal?: AbortSignal, ): Promise<{ emails: Map<string, FastmailEmail>; notFound: Set<string> }> { const emails = new Map<string, FastmailEmail>(); const notFound = new Set<string>(); const effectiveChunkSize = Math.min(chunkSize, session.maxObjectsInGet); for (let index = 0; index < ids.length; index += effectiveChunkSize) { const chunk = ids.slice(index, index + effectiveChunkSize); const id = `email-get-${index / effectiveChunkSize}`; const response = await this.request(session, [{ name: "Email/get", id, arguments: { accountId: session.accountId, ids: chunk, properties: [...EMAIL_PROPERTIES] }, }], signal); const get = methodResult(response, id, "Email/get", emailGetSchema); assertAccount(session.accountId, get.accountId); for (const email of get.list) emails.set(email.id, email); for (const missing of get.notFound ?? []) notFound.add(missing); } return { emails, notFound }; }
private async request(session: FastmailSession, calls: MethodCall[], signal?: AbortSignal): Promise<JmapResponse> { for (const call of calls) { if (!ALLOWED_METHODS.has(call.name)) { throw new FastmailJmapError("method-not-allowed", `Fastmail read client does not allow ${call.name}`); } } const raw = await this.fetchJson(session.apiUrl, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ using: [CORE_CAPABILITY, MAIL_CAPABILITY], methodCalls: calls.map((call) => [call.name, call.arguments, call.id]), }), }, signal); const response = parseWithCode(jmapResponseSchema, raw, "invalid-jmap-response"); const expectedIds = new Set(calls.map((call) => call.id)); if (response.methodResponses.length !== calls.length || response.methodResponses.some((entry) => !expectedIds.has(entry[2]))) { throw new FastmailJmapError("unexpected-method-response", "Fastmail returned an unexpected JMAP method response set"); } return response; }
private async fetchJson(url: string, init: RequestInit, externalSignal?: AbortSignal): Promise<unknown> { const timeout = AbortSignal.timeout(this.timeoutMs); const signal = externalSignal ? AbortSignal.any([externalSignal, timeout]) : timeout; let response: Response; try { response = await this.fetchImpl(url, { ...init, redirect: "error", signal, headers: { ...Object.fromEntries(new Headers(init.headers).entries()), authorization: `Bearer ${this.#token}`, accept: "application/json", }, }); } catch (error) { if (signal.aborted) throw new FastmailJmapError("request-timeout", "Fastmail request timed out or was aborted"); throw new FastmailJmapError("transport-failed", "Fastmail request failed before a response was received"); } if (!response.ok) throw new FastmailJmapError(`http-${response.status}`, `Fastmail returned HTTP ${response.status}`); const contentType = response.headers.get("content-type")?.split(";", 1)[0]?.trim().toLowerCase(); if (contentType !== "application/json") throw new FastmailJmapError("invalid-content-type", "Fastmail response was not JSON"); const declared = Number(response.headers.get("content-length")); if (Number.isFinite(declared) && declared > this.maxResponseBytes) { throw new FastmailJmapError("response-too-large", "Fastmail response exceeded the configured byte bound"); } const bytes = await readBoundedBody(response, this.maxResponseBytes); try { return JSON.parse(new TextDecoder().decode(bytes)) as unknown; } catch { throw new FastmailJmapError("invalid-json", "Fastmail response contained invalid JSON"); } }}
export class FastmailJmapConnector { readonly id: string; readonly kind = "fastmail" as const; private readonly client: FastmailJmapClient; private readonly maxChanges: number; private readonly maxPages: number; private readonly resnapshotLimit: number; private readonly credentialCustody: "dedicated-mail-ingress" | "shared-operator-accepted";
constructor(options: FastmailJmapConnectorOptions) { this.id = required(options.id, "Fastmail source id"); this.client = options.client; this.maxChanges = boundedPositive(options.maxChanges, 1_000, "Fastmail max changes"); this.maxPages = boundedPositive(options.maxPages, 100, "Fastmail max pages"); this.resnapshotLimit = boundedPositive(options.resnapshotLimit, 1_000, "Fastmail resnapshot limit"); this.credentialCustody = options.credentialCustody; }
describe(): JsonObject { return { id: this.id, kind: this.kind, transport: "authenticated-jmap-poll", authority: "read-only", fullBodies: false, initialReplay: "now", maxChanges: this.maxChanges, maxPages: this.maxPages, resnapshotLimit: this.resnapshotLimit, cursorRevision: CURSOR_REVISION, credentialCustody: this.credentialCustody, }; }
async poll(store: JazzThoughtStore, signal?: AbortSignal): Promise<FastmailJmapPollResult> { const cursorId = `cursor:${this.id}`; const prior = await store.getSourceCursor(cursorId); const correlationId = newId("fastmail_jmap_poll"); const startedAt = new Date().toISOString(); await store.registerSource({ id: this.id, kind: this.kind, enabled: true, config: { ...this.describe(), control: { version: 1, kind: "fastmail-jmap", enabled: true, lane: "mail-ingress", }, }, updatedAt: startedAt, }); await store.appendEvent(connectorEvent(this.id, "started", correlationId, startedAt, { status: "started", transport: "authenticated-jmap-poll", cursorRevision: CURSOR_REVISION, }));
try { const session = await this.client.discover(signal); const accountIdHash = sha256(session.accountId); assertCursorCompatible(prior, accountIdHash); const readyAt = new Date().toISOString(); await store.registerSource({ id: this.id, kind: this.kind, enabled: true, config: { ...this.describe(), control: { version: 1, kind: "fastmail-jmap", enabled: true, lane: "mail-ingress", runtime: { state: "ready", readyAt, revision: CURSOR_REVISION }, }, }, updatedAt: readyAt, }); const priorEmailState = stringCursor(prior, "emailState"); const priorQueryState = stringCursor(prior, "queryState"); let initialized = false; let resnapshot = false; let emailState: string; let queryState: string; let operations: FastmailOperation[] = [];
if (!priorEmailState && !priorQueryState) { const current = await this.client.currentSnapshot(0, signal); assertAccount(session.accountId, current.accountId); emailState = current.emailState; queryState = current.queryState; initialized = true; } else if (!priorEmailState || !priorQueryState) { throw new FastmailJmapError("incomplete-cursor", "Fastmail durable cursor is missing one required state"); } else { try { const changes = await this.client.changes(priorEmailState, priorQueryState, this.maxChanges, this.maxPages, signal); assertAccount(session.accountId, changes.accountId); emailState = changes.emailState; queryState = changes.queryState; operations = changes.operations; } catch (error) { if (!(error instanceof FastmailJmapCannotCalculateChanges)) throw error; const current = await this.client.currentSnapshot(this.resnapshotLimit, signal); assertAccount(session.accountId, current.accountId); emailState = current.emailState; queryState = current.queryState; operations = current.emails.map((email) => ({ id: email.id, operation: "updated" as const, email })); resnapshot = true; } }
const candidates = operations.map((operation) => emailEventCandidate({ sourceId: this.id, accountIdHash, operation, correlationId, observedAt: startedAt, resnapshot, })); const completedAt = new Date().toISOString(); const cursor: SourceCursor = { id: cursorId, source: this.id, cursor: { revision: CURSOR_REVISION, accountIdHash, emailState, queryState }, lastSuccessAt: completedAt, ...(prior?.lastFailureAt ? { lastFailureAt: prior.lastFailureAt } : {}), updatedAt: completedAt, }; const sourceCount = candidates.length; candidates.push(cursorEvent(this.id, correlationId, completedAt, cursor)); if (prior?.lastFailureAt && (!prior.lastSuccessAt || prior.lastFailureAt > prior.lastSuccessAt)) { candidates.push(connectorEvent(this.id, "recovered", correlationId, completedAt, { status: "recovered", previousFailureAt: prior.lastFailureAt, })); } candidates.push(connectorEvent(this.id, "completed", correlationId, completedAt, { status: initialized ? "initialized" : resnapshot ? "resnapshot" : "updated", initialized, resnapshot, created: operations.filter((operation) => operation.operation === "created").length, updated: operations.filter((operation) => operation.operation === "updated").length, destroyed: operations.filter((operation) => operation.operation === "destroyed").length, })); 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)); return { source: this.id, initialized, resnapshot, inserted: events.length, unchanged: sourceCount - events.length, created: operations.filter((operation) => operation.operation === "created").length, updated: operations.filter((operation) => operation.operation === "updated").length, destroyed: operations.filter((operation) => operation.operation === "destroyed").length, events, cursor, }; } catch (error) { const failedAt = new Date().toISOString(); const code = fastmailJmapErrorCode(error); const failedCursor: SourceCursor = { id: cursorId, source: this.id, cursor: prior?.cursor ?? { revision: CURSOR_REVISION }, ...(prior?.lastSuccessAt ? { lastSuccessAt: prior.lastSuccessAt } : {}), lastFailureAt: failedAt, lastError: code, updatedAt: failedAt, }; await store.appendProducerBatch([connectorEvent(this.id, "failed", correlationId, failedAt, { status: "failed", phase: "fastmail-jmap-poll", code, })], failedCursor); throw error; } }}
export function fastmailJmapErrorCode(error: unknown): string { return error instanceof FastmailJmapError ? error.code : "fastmail-jmap-failed";}
function methodResult<T>(response: JmapResponse, id: string, name: string, schema: z.ZodType<T>): T { const matches = response.methodResponses.filter((entry) => entry[2] === id); if (matches.length !== 1) throw new FastmailJmapError("method-response-count", `Fastmail returned ${matches.length} responses for ${id}`); const [actualName, payload] = matches[0]!; if (actualName === "error") { const type = typeof payload.type === "string" ? payload.type : "unknown"; if (type === "cannotCalculateChanges" || type === "tooManyChanges") throw new FastmailJmapCannotCalculateChanges(); throw new FastmailJmapError(`method-${safeCode(type)}`, `Fastmail ${name} returned ${safeCode(type)}`); } if (actualName !== name) throw new FastmailJmapError("method-name-mismatch", `Fastmail returned ${actualName} for ${id}`); return parseWithCode(schema, payload, `invalid-${name.toLowerCase().replaceAll("/", "-")}`);}
function collapseOperations(pages: EmailChanges[]): Map<string, FastmailOperation["operation"]> { const operations = new Map<string, FastmailOperation["operation"]>(); for (const page of pages) { for (const id of page.created) operations.set(id, "created"); for (const id of page.updated) { if (!operations.has(id)) operations.set(id, "updated"); } for (const id of page.destroyed) operations.set(id, "destroyed"); } return operations;}
function emailEventCandidate(options: { sourceId: string; accountIdHash: string; operation: FastmailOperation; correlationId: string; observedAt: string; resnapshot: boolean;}): EventCandidate { const email = options.operation.email; const observedMetadata = compact({ 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, })), }); return { type: "stream.thought.source.email.observed", schemaVersion: 1, source: options.sourceId, sourceKind: "fastmail", externalId: options.operation.id, idempotencyKey: sha256(canonicalJson({ accountIdHash: options.accountIdHash, emailId: options.operation.id, operation: options.operation.operation, observedMetadata, })), occurredAt: options.operation.operation === "created" && email?.receivedAt ? email.receivedAt : options.observedAt, actor: email?.from?.[0]?.email ?? options.sourceId, correlationId: options.correlationId, privacy: "sensitive", payload: compact({ accountIdHash: options.accountIdHash, emailId: options.operation.id, operation: options.operation.operation, resnapshot: options.resnapshot || undefined, ...observedMetadata, }), };}
function cursorEvent(sourceId: string, correlationId: string, at: string, cursor: SourceCursor): EventCandidate { return { type: "stream.thought.connector.cursor.advanced", schemaVersion: 1, source: sourceId, sourceKind: "fastmail", externalId: cursor.id, idempotencyKey: `${correlationId}:cursor`, occurredAt: at, actor: sourceId, correlationId, privacy: "sensitive", payload: { source: sourceId, cursorId: cursor.id, cursor: cursor.cursor }, };}
function connectorEvent( sourceId: string, 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: sourceId, sourceKind: "fastmail", externalId: correlationId, idempotencyKey: `${correlationId}:${phase}`, occurredAt: at, actor: sourceId, correlationId, privacy: "sensitive", payload, };}
function assertCursorCompatible(cursor: SourceCursor | undefined, accountIdHash: string): void { if (!cursor) return; const revision = cursor.cursor.revision; if (revision !== CURSOR_REVISION) throw new FastmailJmapError("cursor-revision", "Fastmail cursor revision is incompatible with live JMAP polling"); if (cursor.cursor.accountIdHash !== accountIdHash) throw new FastmailJmapError("cursor-account", "Fastmail cursor belongs to a different primary account");}
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 FastmailJmapError("invalid-cursor", `Fastmail cursor ${key} must be a nonempty string`); return value;}
function assertAccount(expected: string, actual: string): void { if (actual !== expected) throw new FastmailJmapError("account-mismatch", "Fastmail method response belongs to a different account");}
function capabilityInteger(capability: unknown, key: string, minimum: number, maximum: number): number { if (!capability || typeof capability !== "object" || Array.isArray(capability)) { throw new FastmailJmapError("invalid-core-capability", "Fastmail session has an invalid JMAP Core capability"); } const value = (capability as Record<string, unknown>)[key]; if (!Number.isSafeInteger(value) || (value as number) < minimum || (value as number) > maximum) { throw new FastmailJmapError("invalid-core-capability", `Fastmail JMAP Core ${key} is outside the supported bound`); } return value as number;}
function validateCredentialEndpoint(raw: string, label: string): string { let url: URL; try { url = new URL(raw); } catch { throw new FastmailJmapError("invalid-endpoint", `${label} is invalid`); } if (url.protocol !== "https:" || url.username || url.password || url.search || url.hash || url.port) { throw new FastmailJmapError("invalid-endpoint", `${label} must be credential-free HTTPS on the default port`); } return url.toString();}
function isAllowedFastmailApiUrl(raw: string): boolean { const url = new URL(raw); return url.hostname === "api.fastmail.com" || url.hostname.endsWith(".api.fastmail.com") || url.hostname === "jmap.fastmail.com";}
async function readBoundedBody(response: Response, maximum: number): Promise<Uint8Array> { if (!response.body) return new Uint8Array(); const reader = response.body.getReader(); const chunks: Uint8Array[] = []; let total = 0; while (true) { const item = await reader.read(); if (item.done) break; total += item.value.byteLength; if (total > maximum) { await reader.cancel(); throw new FastmailJmapError("response-too-large", "Fastmail response exceeded the configured byte bound"); } chunks.push(item.value); } const output = new Uint8Array(total); let offset = 0; for (const chunk of chunks) { output.set(chunk, offset); offset += chunk.byteLength; } return output;}
function parseWithCode<T>(schema: z.ZodType<T>, value: unknown, code: string): T { const result = schema.safeParse(value); if (!result.success) throw new FastmailJmapError(code, `Fastmail response failed ${code}`); return result.data;}
function required(value: string, label: string): string { const normalized = value.trim(); if (!normalized) throw new FastmailJmapError("required", `${label} is required`); return normalized;}
function boundedPositive(value: number, maximum: number, label: string): number { if (!Number.isSafeInteger(value) || value <= 0 || value > maximum) { throw new FastmailJmapError("invalid-bound", `${label} must be a positive integer no greater than ${maximum}`); } return value;}
function boundedNonnegative(value: number, maximum: number, label: string): number { if (!Number.isSafeInteger(value) || value < 0 || value > maximum) { throw new FastmailJmapError("invalid-bound", `${label} must be a nonnegative integer no greater than ${maximum}`); } return value;}
function safeCode(value: string): string { const normalized = value.toLowerCase().replace(/[^a-z0-9-]+/g, "-").replace(/^-+|-+$/g, ""); return normalized.slice(0, 80) || "unknown";}
function compact(value: Record<string, unknown>): JsonObject { return Object.fromEntries(Object.entries(value).filter(([, child]) => child !== undefined)) as JsonObject;}