Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
16 kB · 395 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396import { stableKey } from "../core/ids.js";import { promises as fs } from "node:fs";import os from "node:os";import path from "node:path";import { canonicalJson, sha256, type JsonObject } from "../core/json.js";import type { EventCandidate, PrivacyClass, ThoughtEvent } from "../events/types.js";import type { JazzThoughtStore } from "../jazz/store.js";import type { BatchDeclarationManifest } from "../runtime/manifest.js";import type { ConsumerProgress } from "../store/types.js";
export interface BatchCycleResult { declarationId: string; examined: number; emitted: number; memberCount: number; waiting: boolean; flushReason?: "quiet-window" | "max-age" | "max-items"; batchEventId?: string;}
export class DeterministicBatcher { constructor( private readonly store: JazzThoughtStore, private readonly now: () => Date = () => new Date(), ) {}
async cycle(declaration: BatchDeclarationManifest): Promise<BatchCycleResult> { if (!declaration.enabled) return emptyResult(declaration.id); const sources = [...new Set(declaration.input.sourceIds)].sort(); const progressBySource = new Map<string, ConsumerProgress | undefined>(); const pendingBySource = new Map<string, ThoughtEvent[]>(); await Promise.all(sources.map(async (source) => { const progress = await this.progressForStart(declaration, source); progressBySource.set(source, progress); pendingBySource.set(source, await this.store.queryConsumerEvents({ consumerId: declaration.id, consumerVersion: declaration.version, source, eventTypes: declaration.input.eventTypes, acceptedPrivacy: ["public-source", "private", "sensitive"], afterSequence: progress?.lastSequence ?? 0, limit: declaration.maxItems, })); })); const examined = [...pendingBySource.values()].reduce((total, events) => total + events.length, 0); const pending = mergeSourcePrefixes(pendingBySource, declaration.maxItems); if (pending.length === 0) return emptyResult(declaration.id);
const nowMs = this.now().getTime(); const firstMs = eventTime(pending[0]!, nowMs); const lastMs = eventTime(pending[pending.length - 1]!, nowMs); const maxItems = pending.length >= declaration.maxItems; const maxAge = nowMs - firstMs >= declaration.maxAgeMs; const quiet = nowMs - lastMs >= declaration.quietWindowMs; const flushReason = maxItems ? "max-items" : maxAge ? "max-age" : quiet ? "quiet-window" : undefined; if (!flushReason) { return { declarationId: declaration.id, examined, emitted: 0, memberCount: 0, waiting: true }; }
if (!maxItems) { const selectedLastBySource = new Map<string, ThoughtEvent>(); for (const event of pending) selectedLastBySource.set(event.source, event); const advanced = await Promise.all([...selectedLastBySource].map(([source, last]) => this.store.queryConsumerEvents({ consumerId: declaration.id, consumerVersion: declaration.version, source, eventTypes: declaration.input.eventTypes, acceptedPrivacy: ["public-source", "private", "sensitive"], afterSequence: last.sourceSequence, limit: 1, }))); if (advanced.some((events) => events.length > 0)) { return { declarationId: declaration.id, examined, emitted: 0, memberCount: 0, waiting: true }; } }
const fingerprint = batchDeclarationFingerprint(declaration); const memberIds = pending.map((event) => event.id); const batchId = stableKey("derived-event-batch", fingerprint, ...memberIds); const privacy = batchPrivacy(declaration, pending); const candidate: EventCandidate = { type: declaration.output.eventType, schemaVersion: 1, source: declaration.output.sourceId, sourceKind: "system", externalId: batchId, idempotencyKey: batchId, occurredAt: pending[pending.length - 1]!.occurredAt, actor: `batch:${declaration.id}`, ...(sources.length === 1 ? { rootEventId: pending[0]!.rootEventId, parentEventId: pending[pending.length - 1]!.id, } : {}), correlationId: batchId, privacy, payload: { declaration: { id: declaration.id, version: declaration.version, fingerprint }, flushReason, firstOccurredAt: pending[0]!.occurredAt, lastOccurredAt: pending[pending.length - 1]!.occurredAt, members: pending.map(memberReference), }, }; const selectedSources = [...new Set(pending.map((event) => event.source))]; const sourceSettlements = selectedSources.map((source) => { const sourceMembers = pending.filter((event) => event.source === source) .sort((left, right) => left.sourceSequence - right.sourceSequence); return { progress: progressFor(declaration, sourceMembers[sourceMembers.length - 1]!, this.now().toISOString()), priorFilteredSequence: progressBySource.get(source)?.lastSequence ?? 0, declaration: { inputEventTypes: declaration.input.eventTypes, acceptedPrivacy: ["public-source", "private", "sensitive"] as PrivacyClass[], source, fingerprint, }, }; }); const settled = sourceSettlements.length === 1 ? await this.store.settleDerivedBatch({ candidate, members: pending, ...sourceSettlements[0]!, }) : await this.store.settleDerivedBatch({ candidate, members: pending, sourceSettlements }); return { declarationId: declaration.id, examined, emitted: settled.inserted ? 1 : 0, memberCount: pending.length, waiting: false, flushReason, batchEventId: settled.event.id, }; }
private async progressForStart( declaration: BatchDeclarationManifest, source: string, ): Promise<ConsumerProgress | undefined> { const id = batchProgressId(declaration, source); const existing = await this.store.getConsumerProgress(id); if (existing || declaration.replay !== "now") return existing; const sourceState = (await this.store.listSources()).find((candidate) => candidate.id === source); const lastSequence = sourceState?.lastSequence ?? 0; const lastEvent = lastSequence > 0 ? await this.store.latestSourceEvent(source) : undefined; if (lastSequence > 0 && (!lastEvent || lastEvent.sourceSequence !== lastSequence)) { throw new Error(`Source head evidence mismatch for ${source}`); } return this.store.initializeConsumerProgress({ id, consumerId: declaration.id, consumerVersion: declaration.version, source, lastSequence, lastEventId: lastEvent?.id ?? "", updatedAt: this.now().toISOString(), }); }}
export interface BatchServiceHandle { stop(): Promise<void>;}
export async function startBatchService( store: JazzThoughtStore, declarations: BatchDeclarationManifest[], options: { now?: () => Date; onCycle?: (result: BatchCycleResult) => void; projectRoot?: string } = {},): Promise<BatchServiceHandle> { const enabled = declarations.filter((declaration) => declaration.enabled); const owner = await acquireBatchOwnerLocks(options.projectRoot ?? process.cwd(), enabled); const batcher = new DeterministicBatcher(store, options.now); let stopped = false; let active = Promise.resolve(); const run = () => { active = active.then(async () => { for (const declaration of enabled) options.onCycle?.(await batcher.cycle(declaration)); }); }; run(); const interval = Math.min(...enabled.map((declaration) => declaration.pollIntervalMs), 1_000); const timer = setInterval(run, Math.max(100, interval)); timer.unref(); return { stop: async () => { if (stopped) return; stopped = true; clearInterval(timer); try { await active; } finally { await owner.release(); } }, };}
export function batchDeclarationFingerprint(declaration: BatchDeclarationManifest): string { return sha256(canonicalJson(JSON.parse(JSON.stringify(declaration)) as JsonObject));}
export function batchProgressId(declaration: BatchDeclarationManifest, source: string): string { return stableKey("batch-progress", declaration.id, String(declaration.version), source);}
function progressFor(declaration: BatchDeclarationManifest, event: ThoughtEvent, at: string): ConsumerProgress { return { id: batchProgressId(declaration, event.source), consumerId: declaration.id, consumerVersion: declaration.version, source: event.source, lastSequence: event.sourceSequence, lastEventId: event.id, updatedAt: at, };}
function memberReference(event: ThoughtEvent): JsonObject { return { eventId: event.id, source: event.source, sourceSequence: event.sourceSequence, type: event.type, schemaVersion: event.schemaVersion, privacy: event.privacy, occurredAt: event.occurredAt, observedAt: event.observedAt, payloadHash: event.payloadHash, };}
function batchPrivacy(declaration: BatchDeclarationManifest, members: ThoughtEvent[]): PrivacyClass { const values = new Set(members.map((member) => member.privacy)); if (declaration.privacy === "preserve") { if (values.size !== 1) throw new Error("Preserved batch privacy cannot combine different privacy classes"); return members[0]!.privacy; } const rank: Record<PrivacyClass, number> = { "public-source": 0, private: 1, sensitive: 2 }; return members.reduce<PrivacyClass>( (current, member) => rank[member.privacy] > rank[current] ? member.privacy : current, "public-source", );}
function mergeSourcePrefixes(pendingBySource: Map<string, ThoughtEvent[]>, limit: number): ThoughtEvent[] { const offsets = new Map([...pendingBySource.keys()].map((source) => [source, 0])); const selected: ThoughtEvent[] = []; while (selected.length < limit) { let next: ThoughtEvent | undefined; for (const [source, events] of pendingBySource) { const candidate = events[offsets.get(source) ?? 0]; if (!candidate) continue; if (!next || compareBatchMembers(candidate, next) < 0) next = candidate; } if (!next) break; selected.push(next); offsets.set(next.source, (offsets.get(next.source) ?? 0) + 1); } return selected;}
function compareBatchMembers(left: ThoughtEvent, right: ThoughtEvent): number { return left.observedAt.localeCompare(right.observedAt) || left.occurredAt.localeCompare(right.occurredAt) || left.source.localeCompare(right.source) || left.sourceSequence - right.sourceSequence || left.id.localeCompare(right.id);}
function eventTime(event: ThoughtEvent, nowMs: number): number { const occurred = Date.parse(event.occurredAt); const observed = Date.parse(event.observedAt); if (!Number.isFinite(occurred) || !Number.isFinite(observed)) throw new Error("Batch member timestamps must be valid ISO timestamps"); // Future source timestamps cannot defer batching indefinitely. Clock rollback can // still delay a batch by at most the configured quiet/max-age windows. return Math.min(nowMs, Math.max(occurred, observed));}
function emptyResult(id: string): BatchCycleResult { return { declarationId: id, examined: 0, emitted: 0, memberCount: 0, waiting: false };}
interface BatchOwnerRecord { bootId: string; pid: number; processStart: string; declarationKey: string;}
export interface BatchOwnerLockOptions { bootId?: () => Promise<string>; processStart?: (pid: number) => Promise<string | undefined>; pid?: number;}
export async function acquireBatchOwnerLocks( projectRoot: string, declarations: BatchDeclarationManifest[], options: BatchOwnerLockOptions = {},): Promise<{ release(): Promise<void> }> { const lockRoot = path.join(projectRoot, ".thoughtstream", "runtime", "batch-owners"); await fs.mkdir(lockRoot, { recursive: true, mode: 0o700 }); await fs.chmod(lockRoot, 0o700); const pid = options.pid ?? process.pid; const bootId = await (options.bootId ?? readBootId)(); const processStart = await (options.processStart ?? readProcessStart)(pid); if (!processStart) throw new Error("Cannot establish batch owner process identity"); const held: Array<{ path: string; record: BatchOwnerRecord }> = []; try { for (const declaration of declarations) { const declarationKey = `${declaration.id}@${declaration.version}`; const lockPath = path.join(lockRoot, `${sha256(declarationKey)}.json`); const record = { bootId, pid, processStart, declarationKey }; await claimOwnerLock(lockPath, record, options.processStart ?? readProcessStart, options.bootId ?? readBootId); held.push({ path: lockPath, record }); } } catch (error) { await Promise.all(held.map((lock) => releaseOwnerLock(lock.path, lock.record))); throw error; } return { release: async () => { await Promise.all(held.map((lock) => releaseOwnerLock(lock.path, lock.record))); }, };}
async function claimOwnerLock( lockPath: string, record: BatchOwnerRecord, processStart: (pid: number) => Promise<string | undefined>, bootId: () => Promise<string>,): Promise<void> { for (let attempt = 0; attempt < 3; attempt += 1) { try { const handle = await fs.open(lockPath, "wx", 0o600); try { await handle.writeFile(canonicalJson(record as unknown as JsonObject)); await handle.sync(); } finally { await handle.close(); } return; } catch (error) { if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; const existing = await readOwnerRecord(lockPath); const currentBoot = await bootId(); const currentStart = existing && existing.bootId === currentBoot ? await processStart(existing.pid) : undefined; if (existing && existing.bootId === currentBoot && currentStart === existing.processStart) { throw new Error(`A live batch owner already holds ${record.declarationKey}`); } await fs.rm(lockPath, { force: true }); } } throw new Error(`Could not acquire batch owner lock for ${record.declarationKey}`);}
async function releaseOwnerLock(lockPath: string, record: BatchOwnerRecord): Promise<void> { const existing = await readOwnerRecord(lockPath); if (existing && canonicalJson(existing as unknown as JsonObject) === canonicalJson(record as unknown as JsonObject)) await fs.rm(lockPath, { force: true });}
async function readOwnerRecord(lockPath: string): Promise<BatchOwnerRecord | undefined> { try { const value = JSON.parse(await fs.readFile(lockPath, "utf8")) as Partial<BatchOwnerRecord>; if (typeof value.bootId !== "string" || !Number.isInteger(value.pid) || Number(value.pid) <= 0 || typeof value.processStart !== "string" || typeof value.declarationKey !== "string") return undefined; return value as BatchOwnerRecord; } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return undefined; return undefined; }}
async function readBootId(): Promise<string> { try { return (await fs.readFile("/proc/sys/kernel/random/boot_id", "utf8")).trim(); } catch { return `${os.hostname()}:${Math.floor(Date.now() - os.uptime() * 1_000)}`; }}
async function readProcessStart(pid: number): Promise<string | undefined> { try { const stat = await fs.readFile(`/proc/${pid}/stat`, "utf8"); const close = stat.lastIndexOf(")"); const fields = stat.slice(close + 2).split(" "); return fields[19]; } catch { return undefined; }}