Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
12 kB · 321 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322import { z } from "zod";import { 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";
export interface JetstreamConnectorOptions { id: string; collections: string[]; dids?: string[];}
export interface JetstreamBatchResult { inserted: number; unchanged: number; ignored: number; stale: number; events: ThoughtEvent[]; offeredEvents: ThoughtEvent[]; cursor?: SourceCursor;}
export interface JetstreamBatchOptions { replayFromTimeUs?: number;}
interface JetstreamMessage { did: string; timeUs: number; kind: string; commit?: { rev: string; operation: "create" | "update" | "delete"; collection: string; rkey: string; cid?: string; record?: JsonObject; };}
const jsonValueSchema: z.ZodType<JsonValue> = z.lazy(() => z.union([ z.string(), z.number(), z.boolean(), z.null(), z.array(jsonValueSchema), z.record(z.string(), jsonValueSchema),]));const recordSchema = z.record(z.string(), jsonValueSchema) as z.ZodType<JsonObject>;const CREATE_ONLY_COLLECTIONS = new Set(["network.cosmik.collectionLink"]);
const rawMessageSchema = z.object({ did: z.string().startsWith("did:"), time_us: z.number().int().nonnegative().safe(), kind: z.string().min(1), commit: z.object({ rev: z.string().min(1), operation: z.enum(["create", "update", "delete"]), collection: z.string().min(1), rkey: z.string().min(1), cid: z.string().min(1).optional(), record: recordSchema.optional(), }).optional(),}).passthrough();
export class JetstreamConnector { readonly kind = "jetstream" as const; readonly id: string; readonly collections: string[]; readonly dids: string[]; readonly createOnlyCollections: string[]; readonly filterRevision: string;
constructor(options: JetstreamConnectorOptions) { this.id = required(options.id, "Jetstream connector id"); this.collections = normalized(options.collections); if (this.collections.length === 0) throw new Error("Jetstream requires at least one collection filter"); if (this.collections.length > 100) throw new Error("Jetstream supports at most 100 collection filters"); if (this.collections.some((collection) => !validCollectionFilter(collection))) { throw new Error("Jetstream collection filters must be exact NSIDs or prefixes ending in .*"); } this.dids = normalized(options.dids ?? []); if (this.dids.length > 10_000) throw new Error("Jetstream supports at most 10000 DID filters"); if (this.dids.some((did) => !did.startsWith("did:"))) throw new Error("Jetstream DID filters must start with did:"); this.createOnlyCollections = this.collections.filter((collection) => CREATE_ONLY_COLLECTIONS.has(collection)); this.filterRevision = sha256(canonicalJson({ collections: this.collections, dids: this.dids, ...(this.createOnlyCollections.length > 0 ? { createOnlyCollections: this.createOnlyCollections } : {}), })); }
describe(): JsonObject { return { id: this.id, kind: this.kind, collections: this.collections, dids: this.dids, ...(this.createOnlyCollections.length > 0 ? { createOnlyCollections: this.createOnlyCollections } : {}), filterRevision: this.filterRevision, transport: "captured-batch-and-live-websocket", }; }
normalize(raw: unknown, correlationId: string): EventCandidate | undefined { const message = parseJetstreamMessage(raw); if (message.kind !== "commit") return undefined; if (!message.commit) throw new Error("Jetstream commit message is missing commit data"); if (!this.collections.some((filter) => collectionMatches(message.commit!.collection, filter))) return undefined; if (this.dids.length > 0 && !this.dids.includes(message.did)) return undefined; const { commit } = message; if (CREATE_ONLY_COLLECTIONS.has(commit.collection) && commit.operation !== "create") return undefined; const atUri = `at://${message.did}/${commit.collection}/${commit.rkey}`; return { type: "stream.thought.source.atproto.commit", schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: atUri, idempotencyKey: `${message.did}:${commit.collection}:${commit.rkey}:${commit.rev}:${commit.operation}`, occurredAt: new Date(Math.floor(message.timeUs / 1_000)).toISOString(), actor: message.did, correlationId, privacy: "public-source", payload: compact({ did: message.did, operation: commit.operation, collection: commit.collection, rkey: commit.rkey, rev: commit.rev, atUri, cid: commit.cid, record: commit.record, }), }; }
async ingestBatch( store: JazzThoughtStore, rawMessages: Iterable<unknown>, options: JetstreamBatchOptions = {}, ): Promise<JetstreamBatchResult> { const cursorId = `cursor:${this.id}`; const prior = await store.getSourceCursor(cursorId); const correlationId = newId("jetstream_batch"); const startedAt = new Date().toISOString(); try { this.assertCursorCompatible(prior); const messages = [...rawMessages].map(parseJetstreamMessage) .sort((left, right) => left.timeUs - right.timeUs); const priorTimeUs = cursorTimeUs(prior); const replayFromTimeUs = options.replayFromTimeUs ?? priorTimeUs; if (!Number.isSafeInteger(replayFromTimeUs) || replayFromTimeUs < 0 || replayFromTimeUs > priorTimeUs) { throw new Error("Jetstream replayFromTimeUs must be a nonnegative safe integer at or before the durable cursor"); } let maxSeenTimeUs = priorTimeUs; let ignored = 0; let stale = 0; const candidates: EventCandidate[] = [];
for (const message of messages) { if (message.timeUs <= replayFromTimeUs) { stale += 1; continue; } maxSeenTimeUs = Math.max(maxSeenTimeUs, message.timeUs); const candidate = this.normalize(messageToRaw(message), correlationId); if (!candidate) { ignored += 1; continue; } candidates.push(candidate); }
const completedAt = new Date().toISOString(); const cursor = this.successCursor(cursorId, prior, maxSeenTimeUs, completedAt); const sourceCount = candidates.length; if (cursor && maxSeenTimeUs > priorTimeUs) candidates.push(this.cursorEvent(correlationId, completedAt, cursor)); 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, ignored, stale, events, offeredEvents, ...(cursor ? { cursor } : {}) }; } catch (error) { const failedAt = new Date().toISOString(); const message = error instanceof Error ? error.message : String(error); const failedEvent: EventCandidate = { type: "stream.thought.connector.failed", schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: correlationId, idempotencyKey: `${correlationId}:failed`, occurredAt: failedAt, actor: this.id, correlationId, privacy: "public-source", payload: { status: "failed", error: message, phase: "captured-batch-ingest" }, }; const failedCursor: SourceCursor = { id: cursorId, source: this.id, cursor: prior?.cursor ?? { filterRevision: this.filterRevision }, ...(prior?.lastSuccessAt ? { lastSuccessAt: prior.lastSuccessAt } : {}), lastFailureAt: failedAt, lastError: message, updatedAt: failedAt, }; await store.appendProducerBatch([failedEvent], failedCursor); throw error; } }
async readCursor(store: JazzThoughtStore): Promise<SourceCursor | undefined> { const cursor = await store.getSourceCursor(`cursor:${this.id}`); this.assertCursorCompatible(cursor); return cursor; }
private assertCursorCompatible(cursor: SourceCursor | undefined): void { const revision = cursor?.cursor.filterRevision; if (revision !== undefined && revision !== this.filterRevision) { throw new Error("Jetstream cursor filter revision does not match connector configuration"); } cursorTimeUs(cursor); }
private successCursor(id: string, prior: SourceCursor | undefined, timeUs: number, at: string): SourceCursor | undefined { if (!prior && timeUs === 0) return undefined; return { id, source: this.id, cursor: { filterRevision: this.filterRevision, timeUs }, lastSuccessAt: at, updatedAt: at, }; }
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: "public-source", payload: { source: this.id, cursorId: cursor.id, cursor: cursor.cursor }, }; }}
export function parseJetstreamMessage(raw: unknown): JetstreamMessage { const value = rawMessageSchema.parse(raw); return { did: value.did, timeUs: value.time_us, kind: value.kind, ...(value.commit ? { commit: { rev: value.commit.rev, operation: value.commit.operation, collection: value.commit.collection, rkey: value.commit.rkey, ...(value.commit.cid ? { cid: value.commit.cid } : {}), ...(value.commit.record ? { record: value.commit.record } : {}), } } : {}), };}
export function parseJetstreamNdjson(value: string): unknown[] { return value.split(/\r?\n/).map((line) => line.trim()).filter(Boolean).map((line, index) => { try { return JSON.parse(line) as unknown; } catch (error) { throw new Error(`Invalid Jetstream NDJSON line ${index + 1}: ${error instanceof Error ? error.message : String(error)}`); } });}
export function cursorTimeUs(cursor: SourceCursor | undefined): number { const value = cursor?.cursor.timeUs; if (value === undefined) return 0; if (typeof value !== "number" || !Number.isSafeInteger(value) || value < 0) { throw new Error("Jetstream cursor timeUs must be a nonnegative safe integer"); } return value;}
function messageToRaw(message: JetstreamMessage): JsonObject { return compact({ did: message.did, time_us: message.timeUs, kind: message.kind, commit: message.commit, });}
function normalized(values: string[]): string[] { return [...new Set(values.map((value) => value.trim()).filter(Boolean))].sort();}
function validCollectionFilter(value: string): boolean { if (!value.includes("*")) return value.includes("."); return value.endsWith(".*") && !value.slice(0, -2).includes("*") && value.slice(0, -2).includes(".");}
function collectionMatches(collection: string, filter: string): boolean { if (!filter.endsWith(".*")) return collection === filter; return collection.startsWith(`${filter.slice(0, -2)}.`);}
function required(value: string, label: string): string { const normalizedValue = value.trim(); if (!normalizedValue) throw new Error(`${label} is required`); return normalizedValue;}
function compact<T extends Record<string, unknown>>(value: T): JsonObject { return Object.fromEntries(Object.entries(value).filter(([, child]) => child !== undefined)) as JsonObject;}