Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
102 kB · 2308 lines
TypeScript
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234223522362237223822392240224122422243224422452246224722482249225022512252225322542255225622572258225922602261226222632264226522662267226822692270227122722273227422752276227722782279228022812282228322842285228622872288228922902291229222932294229522962297229822992300230123022303230423052306230723082309import { mkdirSync } from "node:fs";import { promises as fs } from "node:fs";import { createHash, randomBytes } from "node:crypto";import path from "node:path";import { type Db, type TransactionScope, type JazzSession } from "jazz-tools";import { createJazzSession, type JazzClient } from "jazz-tools/backend";import { startLocalJazzServer, type LocalJazzServerHandle } from "jazz-tools/dev";import { canonicalJson, hashJson, parseJsonObject, sha256, type JsonObject, type JsonValue } from "../core/json.js";import { stableKey } from "../core/ids.js";import { modelAdapterIdentitySchema, type ModelAdapterIdentity } from "../adapters/model-adapters.js";import { createDefaultRegistry, type EventRegistry } from "../events/registry.js";import type { AppendEventResult, CurrentDocument, DocumentVersion, EventCandidate, ThoughtEvent,} from "../events/types.js";import type { AgentRun, AgentRunStatus, ConsumerEventQuery, ConsumerFailureSettlement, ConsumerProgress, ConsumerSuccessSettlement, InferenceAccountingRecord, InferenceBudgetLimit, InferenceBudgetPolicy, InferenceCharge, InferenceReservationDecision, InferenceReservationEstimate, InferenceReservationRequest, InferenceUsage, LettaConversationBinding, ProducerBatchResult, Projection, SourceState, SourceCursor, StoredAgent, TraceChunk,} from "../store/types.js";import thoughtstreamPermissions from "./permissions.js";import { inferenceAccountingEnabled, thoughtstreamApp } from "./schema.js";
export interface JazzThoughtStoreOptions { projectRoot: string; dataPath?: string; appId?: string; serverUrl?: string; backendSecret?: string; adminSecret?: string; runtimeRevision?: string; registry?: EventRegistry; /** * Storage-owner topology for this process. `"auto"` (the default) is * TEST-ONLY: it lets the first process to open a root become the storage * owner (in-process server + discovery-file attach), which is fine for * isolated test roots but is a production anti-pattern — the owner is an * arbitrary short-lived process whose exit tears the server down under its * peers. `"client"` is the production mode: it requires an explicit * `serverUrl` (and, outside tests, a backend secret) and never starts or * stops a server. `"owner"` is used by the dedicated storage-server process * (`scripts/serve-jazz-storage.ts`): it starts the server, publishes the * discovery file, and is the only process allowed to own storage. */ storageMode?: "auto" | "client" | "owner";}
/** * Whether `storageMode: "auto"` is currently permitted. Auto mode exists so * test roots (fresh temp directories, single test process or vitest workers) * keep working without a server supervisor. It is never enabled by a * deployment: nothing in the shipped configuration sets * `THOUGHTSTREAM_JAZZ_AUTO=1`, and production processes must run as * `client` against the dedicated storage server. */function autoStorageModeAllowed(): boolean { return process.env.THOUGHTSTREAM_JAZZ_AUTO === "1" || process.env.VITEST === "true";}
function resolveStorageMode(options: JazzThoughtStoreOptions): "auto" | "client" | "owner" { const mode = options.storageMode ?? process.env.THOUGHTSTREAM_JAZZ_STORAGE_MODE; if (mode === "client" || mode === "owner") return mode; if (mode === "auto" || mode === undefined) { if (!autoStorageModeAllowed()) { throw new Error( "Jazz storage mode is not configured. Production processes must set THOUGHTSTREAM_JAZZ_STORAGE_MODE=client with JAZZ_SERVER_URL and JAZZ_BACKEND_SECRET (run scripts/serve-jazz-storage.ts as the dedicated storage owner). Test roots may set THOUGHTSTREAM_JAZZ_AUTO=1 for auto ownership.", ); } return "auto"; } throw new Error(`Unknown Jazz storage mode: ${mode}`);}
/** * Cryptographically strong secrets for locally-owned servers. The upstream * `startLocalJazzServer` dev helper defaults its backend secret to * `jazz-test-backend-<8 hex chars>` (32 bits of a random UUID) — a test * convenience, never acceptable for a listening server. We always pass our * own 256-bit secrets so the upstream default is unreachable. */function strongSecret(prefix: string): string { return `${prefix}-${randomBytes(32).toString("hex")}`;}
type JazzTable = { where(input: Record<string, unknown>): { limit(count: number): unknown } };
// Process-wide rather than store-instance-wide: a process may construct more than one// JazzThoughtStore, but all of them must share the same per-account reservation queue.// Independent processes still enforce independently and may briefly overshoot an aggregate cap.const inferenceAccountQueues = new Map<string, Promise<void>>();const producerSourceQueues = new Map<string, Promise<void>>();
export class JazzThoughtStore { private readonly session: JazzSession<JazzClient>; private readonly client: JazzClient; private readonly db: Db; private readonly durabilityTier: "local" | "edge" | "global"; private readonly registry: EventRegistry; private readonly runtimeRevision: string; private readonly localServer: LocalJazzServerHandle | undefined; private producerStorageRevision = 3; private consumerStorageRevision = 5; private inferenceStorageRevision = 2;
private constructor( private readonly options: JazzThoughtStoreOptions, session: JazzSession<JazzClient>, localServer?: LocalJazzServerHandle, ) { this.registry = options.registry ?? createDefaultRegistry(); this.runtimeRevision = options.runtimeRevision ?? "thoughtstream-dev"; this.session = session; this.client = this.requireReadyClient(session); this.db = this.client.db; // Durability mapping (alpha55): the client never owns persistent storage // (memory driver), so a "local"-tier wait would only settle its in-memory // cache and a write could be lost on close. The Jazz server — in-process // or remote — is the storage owner, so every required write and read waits // at "edge": acknowledged as persisted by the server. This is the alpha55 // equivalent of the old local-file durability guarantee, and it is also // the old server-configured behavior; the two cases now share one tier. this.durabilityTier = "edge"; this.localServer = localServer; this.db.onMutationError(() => console.error("thought stream Jazz mutation rejected")); }
/** * Open a store against the configured Jazz server. * * Storage-owner topology (jazz-tools 2.0.0-alpha.55): exactly one process * owns the on-disk database — the Jazz server. The production topology is * `"client"`: a dedicated long-lived storage-server process * (`scripts/serve-jazz-storage.ts`, mode `"owner"`) owns the database and * publishes a stable loopback endpoint plus credentials in a 0600 discovery * file; every other process attaches as a backend client and never starts * or stops a server. `"auto"` (test-only) lets the first opener become the * owner so isolated test roots need no supervisor; it is gated on * `THOUGHTSTREAM_JAZZ_AUTO=1` / vitest and fails closed otherwise. * * alpha55 persistent storage takes an exclusive RocksDB lock, so a second * owner cannot open the same `dataPath`. In auto mode the owner records its * server URL and backend secret in a discovery file next to the database; * a concurrent `open` probes it and attaches as a client instead of * starting a competing server. In client mode a stale or missing server is * a hard failure — clients never silently promote themselves to owner. * * `createJazzSession` is asynchronous, so construction is an async factory. * A failed open releases the session and any locally-owned server before * propagating the error. */ static async open(options: JazzThoughtStoreOptions): Promise<JazzThoughtStore> { const storageMode = resolveStorageMode(options); const serverUrl = options.serverUrl ?? process.env.JAZZ_SERVER_URL; const backendSecret = options.backendSecret ?? process.env.JAZZ_BACKEND_SECRET; const dataPath = options.dataPath ?? path.join(options.projectRoot, ".thoughtstream", "state", "jazz.sqlite"); let localServer: LocalJazzServerHandle | undefined; let resolvedServerUrl = serverUrl; let resolvedBackendSecret = backendSecret; if (storageMode === "client") { // Fail closed: a client without a server is a configuration error, not // an invitation to become the owner. if (!resolvedServerUrl) { throw new Error("Jazz storage mode is client but JAZZ_SERVER_URL is not set; refusing to start a server implicitly"); } if (!resolvedBackendSecret) { throw new Error("Jazz storage mode is client but JAZZ_BACKEND_SECRET is not set; refusing to connect without backend admission credentials"); } } else if (!resolvedServerUrl) { mkdirSync(path.dirname(dataPath), { recursive: true }); const discovery = await readServerDiscoveryFile(dataPath); if (discovery && await isJazzServerReachable(discovery.url)) { if (storageMode === "owner") { // An owner never attaches to a foreign server: either a previous // owner is still running (this one must not start) or the discovery // file is stale-but-reachable (a server we do not control). throw new Error(`Refusing to start a storage-owning server: another server is already reachable at ${discovery.url} for ${dataPath}`); } // Another process owns storage for this dataPath; attach as a client. resolvedServerUrl = discovery.url; resolvedBackendSecret = backendSecret ?? discovery.backendSecret; } else { if (discovery) await clearServerDiscoveryFile(dataPath); // Never rely on the upstream dev-helper secret defaults (32 bits); // generate 256-bit secrets here so the defaults are unreachable. const generatedBackendSecret = backendSecret ?? strongSecret("thoughtstream-backend"); const generatedAdminSecret = options.adminSecret ?? process.env.JAZZ_ADMIN_SECRET ?? strongSecret("thoughtstream-admin"); localServer = await startLocalJazzServer({ appId: options.appId ?? process.env.JAZZ_APP_ID ?? "thoughtstream-local", dataDir: dataPath, schema: thoughtstreamApp, backendSecret: generatedBackendSecret, adminSecret: generatedAdminSecret, }); resolvedServerUrl = localServer.url; resolvedBackendSecret = localServer.backendSecret; await writeServerDiscoveryFile(dataPath, localServer.url, localServer.backendSecret); } } // The client never opens the server's storage. Only the server process // touches the on-disk database (alpha55 persistent driver is a RocksDB // directory at `dataPath`); this client keeps its query cache in memory // and reaches durability through the sync transport. This is what makes // the single-storage-owner guarantee structural rather than conventional. let session: JazzSession<JazzClient>; try { session = await createJazzSession({ appId: options.appId ?? process.env.JAZZ_APP_ID ?? "thoughtstream-local", app: thoughtstreamApp, permissions: thoughtstreamPermissions, driver: { type: "memory" }, env: process.env.JAZZ_ENV ?? "dev", serverUrl: resolvedServerUrl, ...(resolvedBackendSecret ? { initial: { backendSecret: resolvedBackendSecret } } : {}), }); } catch (error) { await localServer?.stop().catch(() => undefined); throw error; } try { return new JazzThoughtStore(options, session, localServer); } catch (error) { await session.close().catch(() => undefined); await localServer?.stop().catch(() => undefined); throw error; } }
private requireReadyClient(session: JazzSession<JazzClient>): JazzClient { const snapshot = session.getSnapshot(); if (snapshot.status !== "ready" || !snapshot.client) { void session.close().catch(() => undefined); throw new Error(`Jazz session did not open ready (status: ${snapshot.status}${snapshot.error ? `: ${snapshot.error.message}` : ""})`); } return snapshot.client; }
getArtifactRoot(): string { return path.join(this.options.projectRoot, ".thoughtstream", "artifacts"); }
async appendEvent(candidate: EventCandidate): Promise<AppendEventResult> { const result = await this.appendProducerBatch([candidate]); const event = result.events[0]; if (!event) throw new Error("Atomic event append returned no event"); return { event, inserted: result.inserted.length === 1 }; }
async appendProducerBatch(candidates: EventCandidate[], cursor?: SourceCursor): Promise<ProducerBatchResult> { if (candidates.length === 0 && !cursor) return { events: [], inserted: [], unchanged: [] }; const source = candidates[0]?.source ?? cursor!.source; return this.withProducerSourceLock(source, () => this.appendProducerBatchUnlocked(source, candidates, cursor)); }
private async appendProducerBatchUnlocked( source: string, candidates: EventCandidate[], cursor?: SourceCursor, ): Promise<ProducerBatchResult> { if (candidates.some((candidate) => candidate.source !== source)) { throw new Error("A producer transaction may settle events for only one source"); } if (candidates.some((candidate) => candidate.sourceKind !== candidates[0]?.sourceKind)) { throw new Error("A producer transaction may not mix source kinds"); } if (cursor && cursor.source !== source) throw new Error("Producer cursor source does not match event source"); const sourceKind = candidates[0]?.sourceKind ?? "system"; for (let attempt = 0; attempt < 8; attempt += 1) { const storageRevision = this.producerStorageRevision; const observedAt = new Date().toISOString(); await this.ensureTransactionReady(source); const [sourceSnapshot] = await this.db.all(thoughtstreamApp.sources.where({ key: source }).limit(1), { tier: this.durabilityTier, }); const eventSnapshots = new Map<string, Record<string, unknown>>(); for (const candidate of candidates) { const id = eventId(candidate); if (eventSnapshots.has(id)) continue; const [eventSnapshot] = await this.db.all(thoughtstreamApp.events.where({ key: id }).limit(1), { tier: this.durabilityTier, }); if (eventSnapshot) eventSnapshots.set(id, eventSnapshot); } try { const result = await this.db.transaction(async (tx) => { const sourceRow = sourceSnapshot; if (sourceRow && candidates.length > 0 && String(sourceRow.kind) !== sourceKind) { throw new Error(`Source kind changed for ${source}: ${String(sourceRow.kind)} -> ${sourceKind}`); } let sequence = Number(sourceRow?.lastSequence ?? 0); const events: ThoughtEvent[] = []; const inserted: ThoughtEvent[] = []; const unchanged: ThoughtEvent[] = []; const seen = new Map<string, ThoughtEvent>(); for (const candidate of candidates) { const id = eventId(candidate); const repeated = seen.get(id); if (repeated) { assertSameEventIdentity(repeated, candidate); events.push(repeated); unchanged.push(repeated); continue; } const existing = eventSnapshots.get(id); if (existing) { const event = eventFromJazz(existing); assertSameEventIdentity(event, candidate); events.push(event); unchanged.push(event); seen.set(id, event); continue; } sequence += 1; const event = buildEvent(candidate, id, sequence, observedAt, this.runtimeRevision, this.registry); tx.insert(thoughtstreamApp.events, eventToJazz(event), { id: jazzRowId(`producer-event-v${storageRevision}`, id), }); events.push(event); inserted.push(event); seen.set(id, event); } const sourceData = { key: source, kind: sourceRow ? String(sourceRow.kind) : sourceKind, enabled: true, configJson: sourceRow ? String(sourceRow.configJson) : "{}", lastSequence: sequence, createdAt: sourceRow ? String(sourceRow.createdAt) : observedAt, updatedAt: observedAt, }; if (sourceRow) tx.update(thoughtstreamApp.sources, String(sourceRow.id), sourceData); else { tx.insert(thoughtstreamApp.sources, sourceData, { id: jazzRowId(`producer-source-v${storageRevision}`, source), }); } if (cursor) await upsertCursorInTransaction(tx, cursor, storageRevision); return { events, inserted, unchanged }; }); await result.wait({ tier: this.durabilityTier }); return result.value; } catch (error) { if (!isRecoverableStorageError(error) || attempt === 7) throw error; this.producerStorageRevision = storageRevision + 1; await new Promise((resolve) => setTimeout(resolve, Math.min(500, 20 * (2 ** attempt)))); } } throw new Error("Producer transaction exhausted storage recovery attempts"); }
private async withProducerSourceLock<T>(source: string, operation: () => Promise<T>): Promise<T> { const previous = producerSourceQueues.get(source) ?? Promise.resolve(); let release!: () => void; const lock = new Promise<void>((resolve) => { release = resolve; }); const queued = previous.then(() => lock); producerSourceQueues.set(source, queued); await previous; try { return await operation(); } finally { release(); if (producerSourceQueues.get(source) === queued) producerSourceQueues.delete(source); } }
async settleDerivedBatch(settlement: { candidate: EventCandidate; progress: ConsumerProgress; members: ThoughtEvent[]; priorFilteredSequence: number; declaration: { inputEventTypes: string[]; acceptedPrivacy: Array<ThoughtEvent["privacy"]>; source: string; fingerprint: string; }; } | { candidate: EventCandidate; members: ThoughtEvent[]; sourceSettlements: Array<{ progress: ConsumerProgress; priorFilteredSequence: number; declaration: { inputEventTypes: string[]; acceptedPrivacy: Array<ThoughtEvent["privacy"]>; source: string; fingerprint: string; }; }>; }): Promise<AppendEventResult> { const { candidate, members } = settlement; if (members.length === 0) throw new Error("Derived batch requires at least one member"); const sourceSettlements = "sourceSettlements" in settlement ? settlement.sourceSettlements : [{ progress: settlement.progress, priorFilteredSequence: settlement.priorFilteredSequence, declaration: settlement.declaration, }]; if (sourceSettlements.length === 0) throw new Error("Derived batch requires at least one source settlement"); const fingerprints = new Set(sourceSettlements.map((group) => group.declaration.fingerprint)); if (fingerprints.size !== 1) throw new Error("Derived batch source settlements must share one declaration fingerprint"); const payloadDeclaration = candidate.payload.declaration; if (!payloadDeclaration || typeof payloadDeclaration !== "object" || Array.isArray(payloadDeclaration) || !fingerprints.has(String(payloadDeclaration.fingerprint))) { throw new Error("Derived batch declaration fingerprint does not match settlement authority"); } const sourceNames = sourceSettlements.map((group) => group.declaration.source); if (new Set(sourceNames).size !== sourceNames.length) throw new Error("Derived batch source settlements must be unique"); if (members.some((member) => !sourceNames.includes(member.source))) { throw new Error("Derived batch member has no declared source settlement"); } const payloadMembers = candidate.payload.members; if (!Array.isArray(payloadMembers) || payloadMembers.length !== members.length || payloadMembers.some((reference, index) => ( !reference || typeof reference !== "object" || Array.isArray(reference) || reference.eventId !== members[index]!.id || reference.source !== members[index]!.source || reference.sourceSequence !== members[index]!.sourceSequence || reference.type !== members[index]!.type || reference.schemaVersion !== members[index]!.schemaVersion || reference.privacy !== members[index]!.privacy || reference.occurredAt !== members[index]!.occurredAt || reference.observedAt !== members[index]!.observedAt || reference.payloadHash !== members[index]!.payloadHash ))) { throw new Error("Derived batch payload must name every filtered progress-covered member in order"); } for (const group of sourceSettlements) { const { progress, priorFilteredSequence, declaration } = group; const sourceMembers = members.filter((member) => member.source === declaration.source) .sort((left, right) => left.sourceSequence - right.sourceSequence); if (sourceMembers.length === 0) throw new Error("Derived batch source settlement has no members"); const last = sourceMembers[sourceMembers.length - 1]!; if (progress.source !== declaration.source || progress.lastSequence !== last.sourceSequence || progress.lastEventId !== last.id) { throw new Error("Derived batch progress must settle at the final named member for its source"); } const authoritativeRows = await this.db.all(thoughtstreamApp.events.where({ source: declaration.source, sourceSequence: { gt: priorFilteredSequence, lte: last.sourceSequence }, type: { in: declaration.inputEventTypes }, privacy: { in: declaration.acceptedPrivacy }, }).orderBy("sourceSequence", "asc"), { tier: this.durabilityTier }); const authoritative = authoritativeRows.map(eventFromJazz); if (authoritative.length !== sourceMembers.length || authoritative.some((event, index) => event.id !== sourceMembers[index]!.id)) { throw new Error("Derived batch settlement does not cover the complete matching filtered interval"); } for (const [index, event] of authoritative.entries()) { const member = sourceMembers[index]!; if (event.sourceSequence !== member.sourceSequence || event.type !== member.type || event.schemaVersion !== member.schemaVersion || event.privacy !== member.privacy || event.payloadHash !== member.payloadHash || event.occurredAt !== member.occurredAt || event.observedAt !== member.observedAt) { throw new Error("Derived batch member differs from its authoritative Jazz row"); } } }
const outputSource = candidate.source; const id = eventId(candidate); await this.ensureTransactionReady(outputSource); const [sourceSnapshot] = await this.db.all(thoughtstreamApp.sources.where({ key: outputSource }).limit(1), { tier: this.durabilityTier }); const [eventSnapshot] = await this.db.all(thoughtstreamApp.events.where({ key: id }).limit(1), { tier: this.durabilityTier }); const progressSnapshots = new Map<string, Record<string, unknown> | undefined>(); for (const group of sourceSettlements) { const [snapshot] = await this.db.all(thoughtstreamApp.consumerProgress.where({ key: group.progress.id }).limit(1), { tier: this.durabilityTier }); progressSnapshots.set(group.progress.id, snapshot as Record<string, unknown> | undefined); } const observedAt = new Date().toISOString(); const result = await this.db.transaction(async (tx) => { let event: ThoughtEvent; let inserted = false; if (eventSnapshot) { event = eventFromJazz(eventSnapshot); assertSameEventIdentity(event, candidate); } else { const sequence = Number(sourceSnapshot?.lastSequence ?? 0) + 1; event = buildEvent(candidate, id, sequence, observedAt, this.runtimeRevision, this.registry); tx.insert(thoughtstreamApp.events, eventToJazz(event), { id: jazzRowId("batch-event-v2", id) }); const sourceData = { key: outputSource, kind: candidate.sourceKind, enabled: true, configJson: sourceSnapshot ? String(sourceSnapshot.configJson) : "{}", lastSequence: sequence, createdAt: sourceSnapshot ? String(sourceSnapshot.createdAt) : observedAt, updatedAt: observedAt, }; if (sourceSnapshot) tx.update(thoughtstreamApp.sources, String(sourceSnapshot.id), sourceData); else tx.insert(thoughtstreamApp.sources, sourceData, { id: jazzRowId("batch-source-v2", outputSource) }); inserted = true; } for (const group of sourceSettlements) { await upsertProgressInTransaction(tx, group.progress, progressSnapshots.get(group.progress.id)); } return { event, inserted }; }); await result.wait({ tier: this.durabilityTier }); return result.value; }
async getEvent(id: string): Promise<ThoughtEvent | undefined> { const row = await this.db.one(thoughtstreamApp.events.where({ key: id }).limit(1)); return row ? eventFromJazz(row) : undefined; }
async getEvents(ids: string[]): Promise<ThoughtEvent[]> { const keys = [...new Set(ids)].filter(Boolean); if (keys.length === 0) return []; if (keys.length > 10_000) throw new Error("Bulk event lookup exceeds 10000 ids"); const rows = await this.db.all(thoughtstreamApp.events.where({ key: { in: keys } })); return rows.map(eventFromJazz); }
async listEvents(options: { types?: string[]; source?: string; limit?: number } = {}): Promise<ThoughtEvent[]> { const where = { ...(options.source ? { source: options.source } : {}), ...(options.types && options.types.length > 0 ? { type: { in: [...new Set(options.types)] } } : {}), }; const query = thoughtstreamApp.events.where(where); const rows = options.limit ? await this.db.all(query .orderBy("observedAt", "desc") .orderBy("key", "desc") .limit(options.limit)) : await this.db.all(query.orderBy("observedAt", "asc")); return rows .map(eventFromJazz) .sort((left, right) => left.observedAt.localeCompare(right.observedAt) || left.id.localeCompare(right.id)); }
async listRecentRootEvents(options: { types: string[]; limit: number; maxScanned?: number }): Promise<{ events: ThoughtEvent[]; complete: boolean; scannedEvents: number; }> { if (options.limit <= 0) return { events: [], complete: true, scannedEvents: 0 }; const types = [...new Set(options.types)].filter(Boolean); if (types.length === 0) return { events: [], complete: true, scannedEvents: 0 }; const maxScanned = Math.max(options.limit, options.maxScanned ?? 2_000); const rows = await this.db.all(thoughtstreamApp.events.where({ type: { in: types } }) .orderBy("observedAt", "desc") .orderBy("key", "desc") .limit(maxScanned)); const roots = new Map<string, ThoughtEvent>(); for (const row of rows) { if (String(row.key) !== String(row.rootEventKey)) continue; const event = eventFromJazz(row); roots.set(event.id, event); if (roots.size === options.limit) break; } const events = [...roots.values()].sort((left, right) => ( left.observedAt.localeCompare(right.observedAt) || left.id.localeCompare(right.id) )); return { events, complete: events.length === options.limit || rows.length < maxScanned, scannedEvents: rows.length, }; }
async listEventsForRoots(rootEventIds: string[], options: { limit?: number } = {}): Promise<ThoughtEvent[]> { const keys = [...new Set(rootEventIds)].filter(Boolean); if (keys.length === 0) return []; if (keys.length > 10_000) throw new Error("Root event lookup exceeds 10000 ids"); const query = thoughtstreamApp.events.where({ rootEventKey: { in: keys } }); const rows = options.limit ? await this.db.all(query .orderBy("observedAt", "desc") .orderBy("key", "desc") .limit(options.limit)) : await this.db.all(query.orderBy("observedAt", "asc")); return rows .map(eventFromJazz) .sort((left, right) => left.observedAt.localeCompare(right.observedAt) || left.id.localeCompare(right.id)); }
async latestSourceEvent(source: string): Promise<ThoughtEvent | undefined> { return (await this.listEvents({ source })).reduce<ThoughtEvent | undefined>((latest, event) => ( !latest || event.sourceSequence > latest.sourceSequence ? event : latest ), undefined); }
async listChildEvents(parentEventId: string): Promise<ThoughtEvent[]> { return (await this.db.all(thoughtstreamApp.events.where({ parentEventKey: parentEventId }))) .map(eventFromJazz) .sort((left, right) => left.observedAt.localeCompare(right.observedAt) || left.id.localeCompare(right.id)); }
async listSources(): Promise<SourceState[]> { return (await this.db.all(thoughtstreamApp.sources.where({}))) .map(sourceStateFromJazz) .sort((left, right) => left.id.localeCompare(right.id)); }
async registerSource(input: { id: string; kind: string; enabled: boolean; config: JsonObject; updatedAt?: string; }): Promise<SourceState> { const existingRow = await this.db.one(thoughtstreamApp.sources.where({ key: input.id }).limit(1)); const existing = existingRow ? sourceStateFromJazz(existingRow) : undefined; if (existing && existing.kind !== input.kind) { throw new Error(`Source kind changed for ${input.id}: ${existing.kind} -> ${input.kind}`); } const updatedAt = input.updatedAt ?? new Date().toISOString(); const source: SourceState = { id: input.id, kind: input.kind, enabled: input.enabled, config: input.config, lastSequence: existing?.lastSequence ?? 0, createdAt: existing?.createdAt ?? updatedAt, updatedAt, }; await this.upsert(thoughtstreamApp.sources, { key: input.id }, { key: source.id, kind: source.kind, enabled: source.enabled, configJson: canonicalJson(source.config), lastSequence: source.lastSequence, createdAt: source.createdAt, updatedAt: source.updatedAt, }); return source; }
async appendDocumentVersion(version: DocumentVersion): Promise<boolean> { const existing = await this.db.one(thoughtstreamApp.documentVersions.where({ key: version.id }).limit(1)); if (existing) return false; await this.wait(this.db.insert(thoughtstreamApp.documentVersions, { key: version.id, source: version.source, documentKey: version.documentId, path: version.path, contentType: version.contentType, sha256: version.sha256, content: version.content, sizeBytes: version.sizeBytes, mtimeMs: version.mtimeMs, createdAt: version.createdAt, })); return true; }
async getDocumentVersion(id: string): Promise<DocumentVersion | undefined> { const row = await this.db.one(thoughtstreamApp.documentVersions.where({ key: id }).limit(1)); return row ? documentVersionFromJazz(row) : undefined; }
async listCurrentDocuments(source: string): Promise<CurrentDocument[]> { const rows = await this.db.all(thoughtstreamApp.documents.where({ source })); return rows.map(currentDocumentFromJazz); }
async upsertCurrentDocument(document: CurrentDocument): Promise<void> { await this.upsert(thoughtstreamApp.documents, { key: document.id }, { key: document.id, source: document.source, documentKey: document.documentId, path: document.path, versionKey: document.versionId, sha256: document.sha256, contentType: document.contentType, sizeBytes: document.sizeBytes, mtimeMs: document.mtimeMs, deleted: document.deleted, updatedAt: document.updatedAt, }); }
async upsertAgent(agent: StoredAgent): Promise<void> { await this.upsert(thoughtstreamApp.agents, { key: agent.id }, { key: agent.id, version: agent.version, enabled: agent.enabled, specJson: canonicalJson(agent.spec), specHash: agent.specHash, updatedAt: agent.updatedAt, }); }
async listAgents(): Promise<StoredAgent[]> { return (await this.db.all(thoughtstreamApp.agents.where({}))).map(agentFromJazz); }
async upsertRun(run: AgentRun): Promise<void> { await this.upsert(thoughtstreamApp.runs, { key: run.id }, runToJazz(run)); }
async getRun(id: string): Promise<AgentRun | undefined> { const row = await this.db.one(thoughtstreamApp.runs.where({ key: id }).limit(1)); return row ? runFromJazz(row) : undefined; }
async getRuns(ids: string[]): Promise<AgentRun[]> { const keys = [...new Set(ids)].filter(Boolean); if (keys.length === 0) return []; if (keys.length > 10_000) throw new Error("Bulk run lookup exceeds 10000 ids"); const rows = await this.db.all(thoughtstreamApp.runs.where({ key: { in: keys } })); return rows.map(runFromJazz); }
async getRunsForTriggerEvents(eventIds: string[], options: { limit?: number } = {}): Promise<AgentRun[]> { const keys = [...new Set(eventIds)].filter(Boolean); if (keys.length === 0) return []; if (keys.length > 10_000) throw new Error("Trigger run lookup exceeds 10000 event ids"); const query = thoughtstreamApp.runs.where({ triggerEventKey: { in: keys } }); const rows = options.limit ? await this.db.all(query .orderBy("createdAt", "desc") .orderBy("key", "desc") .limit(options.limit)) : await this.db.all(query); return rows .map(runFromJazz) .sort((left, right) => left.createdAt.localeCompare(right.createdAt) || left.id.localeCompare(right.id)); }
async listRecentRuns(limit = 1_000): Promise<AgentRun[]> { const rows = await this.db.all(thoughtstreamApp.runs.where({}) .orderBy("createdAt", "desc") .orderBy("key", "desc") .limit(limit)); return rows .map(runFromJazz) .sort((left, right) => left.createdAt.localeCompare(right.createdAt) || left.id.localeCompare(right.id)); }
async latestRunForExecution(executionKey: string): Promise<AgentRun | undefined> { const row = await this.db.one(thoughtstreamApp.runs.where({ executionKey }).orderBy("attempt", "desc").limit(1)); return row ? runFromJazz(row) : undefined; }
async listRuns(query: { status?: AgentRunStatus; completedAfter?: string } = {}): Promise<AgentRun[]> { const rows = query.status ? await this.db.all(thoughtstreamApp.runs.where({ status: query.status })) : await this.db.all(thoughtstreamApp.runs.where({})); return rows .map(runFromJazz) .filter((run) => !query.completedAfter || (run.completedAt ?? run.updatedAt) > query.completedAfter) .sort((left, right) => left.createdAt.localeCompare(right.createdAt) || left.id.localeCompare(right.id)); }
async reserveInference(request: InferenceReservationRequest): Promise<InferenceReservationDecision> { const accountId = inferenceAccountId(request.scopeType, request.scopeKey); return this.withInferenceAccountLock(accountId, () => this.reserveInferenceUnlocked(request)); }
private async reserveInferenceUnlocked(request: InferenceReservationRequest): Promise<InferenceReservationDecision> { assertReservationRequest(request); const accountId = inferenceAccountId(request.scopeType, request.scopeKey); const recordRowId = jazzRowId("inference-accounting", request.reservationId); const accountSnapshot = await this.ensureInferenceBudgetAccount( request.scopeType, request.scopeKey, accountId, request.reservedAt, ); await this.db.all(thoughtstreamApp.inferenceAccounting.where({ id: recordRowId }).limit(1), { tier: this.durabilityTier }); const result = await this.db.transaction(async (tx) => { // TransactionScope reads do not share Db's readiness barrier. Use the // durable pre-transaction snapshot instead of asking a fresh transaction // to rediscover an account created immediately before reservation. const accountRow = accountSnapshot; const account = budgetAccountFromJazz(accountRow); await expireStaleReservations(tx, account, request.reservedAt);
const existingRow = await tx.one(thoughtstreamApp.inferenceAccounting.where({ id: recordRowId }).limit(1)); if (existingRow) { const existing = inferenceAccountingFromJazz(existingRow); assertSameReservation(existing, request); upsertBudgetAccountInTransaction(tx, accountRow, accountId, account, request.reservedAt); return { approved: existing.status !== "denied", acquired: false, record: existing, } satisfies InferenceReservationDecision; }
account.windows = reconcileBudgetWindows(account.windows, request.policy, request.reservedAt); const limitingWindowKeys = account.windows .filter((window) => exceedsBudgetLimit(window, request.estimate)) .map((window) => window.key); const leaseExpiresAt = new Date(Date.parse(request.reservedAt) + request.policy.leaseMs).toISOString(); const denied = limitingWindowKeys.length > 0; const record: InferenceAccountingRecord = { id: request.reservationId, scopeType: request.scopeType, scopeKey: request.scopeKey, runId: request.runId, agentId: request.agentId, agentVersion: request.agentVersion, provider: request.provider, model: request.model, status: denied ? "denied" : "reserved", usageStatus: denied ? "unavailable" : "pending", estimate: request.estimate, charged: denied ? zeroInferenceCharge(request.estimate.costMicrousd !== undefined) : request.estimate, windowKeys: denied ? [] : account.windows.map((window) => window.key), reservedAt: request.reservedAt, leaseExpiresAt, ...(denied ? { settledAt: request.reservedAt, denialReason: "inference-budget-exhausted", limitingWindowKeys, } : {}), updatedAt: request.reservedAt, }; if (!denied) { account.windows = account.windows.map((window) => chargeBudgetWindow( window, request.estimate, request.reservationId, request.reservedAt, )); account.activeLeases.push({ reservationId: request.reservationId, leaseExpiresAt }); } tx.insert(thoughtstreamApp.inferenceAccounting, inferenceAccountingToJazz(record), { id: recordRowId }); upsertBudgetAccountInTransaction(tx, accountRow, accountId, account, request.reservedAt); return { approved: !denied, acquired: !denied, record } satisfies InferenceReservationDecision; }); await result.wait({ tier: this.durabilityTier }); return result.value; }
async settleInferenceReservation( reservationId: string, usage: InferenceUsage | undefined, settledAt = new Date().toISOString(), ): Promise<InferenceAccountingRecord> { const snapshot = await this.getInferenceAccountingRecord(reservationId); if (!snapshot) throw new Error(`Cannot settle missing inference reservation ${reservationId}`); const accountId = inferenceAccountId(snapshot.scopeType, snapshot.scopeKey); return this.withInferenceAccountLock(accountId, () => this.settleInferenceReservationUnlocked(reservationId, usage, settledAt)); }
private async settleInferenceReservationUnlocked( reservationId: string, usage: InferenceUsage | undefined, settledAt: string, ): Promise<InferenceAccountingRecord> { const snapshot = await this.getInferenceAccountingRecord(reservationId); if (!snapshot) throw new Error(`Cannot settle missing inference reservation ${reservationId}`); const accountId = inferenceAccountId(snapshot.scopeType, snapshot.scopeKey); const [recordSnapshot] = await this.db.all(thoughtstreamApp.inferenceAccounting.where({ key: reservationId, }).limit(1), { tier: this.durabilityTier }); if (!recordSnapshot) throw new Error(`Cannot settle missing inference reservation ${reservationId}`); const [accountSnapshot] = await this.db.all(thoughtstreamApp.inferenceBudgetAccounts.where({ key: accountId, }).limit(1), { tier: this.durabilityTier }); if (!accountSnapshot) throw new Error(`Cannot settle inference reservation without account ${accountId}`); const result = await this.db.transaction(async (tx) => { // Use durable snapshots rather than visibility-sensitive TransactionScope // reads for rows created by the immediately preceding reservation. const recordRow = recordSnapshot; const record = inferenceAccountingFromJazz(recordRow); if (record.status !== "reserved") return record; const accountRow = accountSnapshot; const account = budgetAccountFromJazz(accountRow); const normalizedUsage = normalizeInferenceUsage(usage); const charged = chargedInferenceUsage(record.estimate, normalizedUsage); const settled: InferenceAccountingRecord = { ...record, status: "settled", usageStatus: inferenceUsageStatus(normalizedUsage, record.estimate), charged, ...(normalizedUsage ? { actualUsage: normalizedUsage } : {}), settledAt, updatedAt: settledAt, }; account.windows = account.windows.map((window) => record.windowKeys.includes(window.key) ? adjustBudgetWindowCharge(window, record.estimate, charged, reservationId) : window); account.activeLeases = account.activeLeases.filter((lease) => lease.reservationId !== reservationId); tx.update(thoughtstreamApp.inferenceAccounting, String(recordRow.id), inferenceAccountingToJazz(settled)); upsertBudgetAccountInTransaction(tx, accountRow, accountId, account, settledAt); return settled; }); await result.wait({ tier: this.durabilityTier }); return result.value; }
async getInferenceAccountingRecord(id: string): Promise<InferenceAccountingRecord | undefined> { const row = await this.db.one(thoughtstreamApp.inferenceAccounting.where({ key: id }).limit(1)); return row ? inferenceAccountingFromJazz(row) : undefined; }
async listInferenceAccounting(options: { agentId?: string; status?: InferenceAccountingRecord["status"] } = {}): Promise<InferenceAccountingRecord[]> { const rows = options.agentId ? await this.db.all(thoughtstreamApp.inferenceAccounting.where({ agentKey: options.agentId })) : await this.db.all(thoughtstreamApp.inferenceAccounting.where({})); return rows .map(inferenceAccountingFromJazz) .filter((record) => !options.status || record.status === options.status) .sort((left, right) => left.reservedAt.localeCompare(right.reservedAt) || left.id.localeCompare(right.id)); }
async appendTrace(chunk: TraceChunk): Promise<boolean> { const existing = await this.db.one(thoughtstreamApp.traceChunks.where({ key: chunk.id }).limit(1)); if (existing) return false; await this.wait(this.db.insert(thoughtstreamApp.traceChunks, { key: chunk.id, runKey: chunk.runId, sequence: chunk.sequence, type: chunk.type, payloadJson: canonicalJson(chunk.payload), createdAt: chunk.createdAt, })); return true; }
async listTrace(runId: string): Promise<TraceChunk[]> { return (await this.db.all(thoughtstreamApp.traceChunks.where({ runKey: runId }))) .map(traceFromJazz) .sort((left, right) => left.sequence - right.sequence); }
async upsertProjection(projection: Projection): Promise<void> { await this.upsert(thoughtstreamApp.projections, { key: projection.id }, { key: projection.id, payloadJson: canonicalJson(projection.payload), lastEventKey: projection.lastEventId, projectionVersion: projection.projectionVersion, updatedAt: projection.updatedAt, }); }
async getProjection(id: string): Promise<Projection | undefined> { const row = await this.db.one(thoughtstreamApp.projections.where({ key: id }).limit(1)); return row ? projectionFromJazz(row) : undefined; }
async getProjections(ids: string[]): Promise<Projection[]> { const keys = [...new Set(ids)].filter(Boolean); if (keys.length === 0) return []; if (keys.length > 10_000) throw new Error("Bulk projection lookup exceeds 10000 ids"); const rows = await this.db.all(thoughtstreamApp.projections.where({ key: { in: keys } })); return rows.map(projectionFromJazz); }
async upsertSourceCursor(cursor: SourceCursor): Promise<void> { await this.ensureTransactionReady(cursor.source); const result = await this.db.transaction(async (tx) => { await upsertCursorInTransaction(tx, cursor); }); await result.wait({ tier: this.durabilityTier }); }
async getSourceCursor(id: string): Promise<SourceCursor | undefined> { const row = await this.db.one(thoughtstreamApp.sourceCursors.where({ key: id }).limit(1)); return row ? sourceCursorFromJazz(row) : undefined; }
async listSourceCursors(): Promise<SourceCursor[]> { return (await this.db.all(thoughtstreamApp.sourceCursors.where({}))) .map(sourceCursorFromJazz) .sort((left, right) => left.source.localeCompare(right.source) || left.id.localeCompare(right.id)); }
async getConsumerProgress(id: string): Promise<ConsumerProgress | undefined> { const row = await this.db.one(thoughtstreamApp.consumerProgress.where({ key: id }).limit(1)); return row ? consumerProgressFromJazz(row) : undefined; }
async initializeConsumerProgress(progress: ConsumerProgress): Promise<ConsumerProgress> { const result = await this.db.transaction(async (tx) => { const row = await tx.one(thoughtstreamApp.consumerProgress.where({ id: jazzRowId("progress", progress.id), }).limit(1)); if (row) return consumerProgressFromJazz(row); tx.insert(thoughtstreamApp.consumerProgress, { key: progress.id, consumerKey: progress.consumerId, consumerVersion: progress.consumerVersion, source: progress.source, lastSequence: progress.lastSequence, lastEventKey: progress.lastEventId, updatedAt: progress.updatedAt, }, { id: jazzRowId("progress", progress.id) }); return progress; }); await result.wait({ tier: this.durabilityTier }); return result.value; }
async listConsumerProgress(): Promise<ConsumerProgress[]> { return (await this.db.all(thoughtstreamApp.consumerProgress.where({}))) .map(consumerProgressFromJazz) .sort((left, right) => left.consumerId.localeCompare(right.consumerId) || left.consumerVersion - right.consumerVersion || left.source.localeCompare(right.source)); }
async getLettaConversationBinding(id: string): Promise<LettaConversationBinding | undefined> { const events = await this.listLettaConversationBindingEvents(); const event = [...events].reverse().find((candidate) => candidate.payload.id === id); if (event) { const binding = lettaConversationBindingFromPayload(event.payload); const projectionId = lettaConversationProjectionId(id); const projection = await this.getProjection(projectionId); if (!projection || projection.lastEventId !== event.id) { await this.upsertProjection(lettaConversationProjection(binding, event.id)); } return binding; } const projection = await this.getProjection(lettaConversationProjectionId(id)); return projection ? lettaConversationBindingFromPayload(projection.payload) : undefined; }
async listLettaConversationBindings(): Promise<LettaConversationBinding[]> { const latest = new Map<string, { binding: LettaConversationBinding; eventId: string }>(); for (const event of await this.listLettaConversationBindingEvents()) { const binding = lettaConversationBindingFromPayload(event.payload); latest.set(binding.id, { binding, eventId: event.id }); } return [...latest.values()].map((value) => value.binding).sort((left, right) => left.id.localeCompare(right.id)); }
async upsertLettaConversationBinding(binding: LettaConversationBinding): Promise<LettaConversationBinding> { const bindings = await this.listLettaConversationBindings(); const existing = bindings.find((candidate) => candidate.id === binding.id); if (bindings.some((candidate) => candidate.remoteMarker === binding.remoteMarker && candidate.id !== binding.id)) { throw new Error(`Letta conversation marker is already bound to another scope: ${binding.remoteMarker}`); } const next = existing ? { ...existing, declarationVersion: binding.declarationVersion, currentPath: binding.currentPath, updatedAt: binding.updatedAt, } : binding; if (existing) { assertSameLettaConversationBinding(existing, binding); if (binding.declarationVersion < existing.declarationVersion) { throw new Error(`Letta conversation declaration version regression for ${binding.id}`); } if (binding.declarationVersion === existing.declarationVersion && binding.currentPath === existing.currentPath) { return existing; } } const appended = await this.appendEvent({ type: "stream.thought.runtime.letta-conversation.binding", schemaVersion: 1, source: "runtime:letta-conversations", sourceKind: "system", externalId: next.id, idempotencyKey: [ next.id, String(next.declarationVersion), sha256(next.currentPath), sha256(next.updatedAt), next.remoteMarker, ].join(":"), occurredAt: next.updatedAt, actor: "thoughtstream", correlationId: next.id, privacy: "sensitive", payload: lettaConversationBindingPayload(next), }); await this.upsertProjection(lettaConversationProjection(next, appended.event.id)); return next; }
private async listLettaConversationBindingEvents(): Promise<ThoughtEvent[]> { return this.listEvents({ source: "runtime:letta-conversations", types: ["stream.thought.runtime.letta-conversation.binding"], }); }
subscribeSources(onSources: (sources: SourceState[]) => void): () => void { // alpha55 `Db.subscribe` delivers the complete current result on every // change (the `subscribeAll` delta API is gone), so the full snapshot is // passed through directly. const unsubscribe = this.db.subscribe( thoughtstreamApp.sources.where({ enabled: true }), (rows) => onSources(rows.map(sourceStateFromJazz)), ); this.subscriptions.push(unsubscribe); return () => { const index = this.subscriptions.indexOf(unsubscribe); if (index >= 0) this.subscriptions.splice(index, 1); unsubscribe(); }; }
subscribeInspectorActivity(onChange: () => void): () => void { const unsubscribers = [ this.db.subscribe( thoughtstreamApp.events.where({}).orderBy("observedAt", "desc").limit(1), () => onChange(), ), this.db.subscribe( thoughtstreamApp.sources.where({}), () => onChange(), ), ]; this.subscriptions.push(...unsubscribers); return () => { for (const unsubscribe of unsubscribers) { const index = this.subscriptions.indexOf(unsubscribe); if (index >= 0) this.subscriptions.splice(index, 1); unsubscribe(); } }; }
async queryConsumerEvents(query: ConsumerEventQuery): Promise<ThoughtEvent[]> { if (query.eventTypes.length === 0 || query.acceptedPrivacy.length === 0) return []; const rows = await this.db.all(thoughtstreamApp.events.where({ source: query.source, sourceSequence: { gt: query.afterSequence }, type: { in: query.eventTypes }, privacy: { in: query.acceptedPrivacy }, }).orderBy("sourceSequence", "asc").limit(query.limit ?? 1_000)); return rows.map(eventFromJazz); }
subscribeConsumerEvents(query: ConsumerEventQuery, onEvents: (events: ThoughtEvent[]) => void): () => void { if (query.eventTypes.length === 0 || query.acceptedPrivacy.length === 0) return () => undefined; // alpha55 `Db.subscribe` delivers the complete current result on every // change, so the added-row delta is reconstructed by diffing consecutive // full snapshots by row key. Insertions and updates to existing ids are // both surfaced: a newly inserted id is delivered as an added row, and an // update to an already-delivered id is re-delivered so durable consumer // progress, not an accumulated in-memory snapshot, decides what remains // after restart. Rows present in the initial snapshot are delivered as the // first delta, matching the old `subscribeAll` initial-delivery behavior. let seen = new Set<string>(); const unsubscribe = this.db.subscribe(thoughtstreamApp.events.where({ source: query.source, sourceSequence: { gt: query.afterSequence }, type: { in: query.eventTypes }, privacy: { in: query.acceptedPrivacy }, }).orderBy("sourceSequence", "asc"), (rows) => { const added = rows.filter((row) => !seen.has(String(row.key))); for (const row of rows) seen.add(String(row.key)); if (added.length > 0) onEvents(added.map(eventFromJazz)); }); this.subscriptions.push(unsubscribe); return () => { const index = this.subscriptions.indexOf(unsubscribe); if (index >= 0) this.subscriptions.splice(index, 1); unsubscribe(); }; }
async settleConsumerSuccess( settlement: ConsumerSuccessSettlement, ): Promise<{ outputEvent: ThoughtEvent; sideEffectEvents: ThoughtEvent[]; completedEvent: ThoughtEvent }> { const events = await this.settleConsumerTerminal( settlement.run, settlement.inputEvent, [settlement.output, ...(settlement.sideEffects ?? []), settlement.completed], settlement.progress, ); const outputEvent = events[0]; const completedEvent = events.at(-1); if (!outputEvent || !completedEvent) throw new Error("Consumer success settlement did not produce both receipts"); return { outputEvent, sideEffectEvents: events.slice(1, -1), completedEvent }; }
async settleConsumerFailure(settlement: ConsumerFailureSettlement): Promise<ThoughtEvent> { const [terminal] = await this.settleConsumerTerminal( settlement.run, settlement.inputEvent, [settlement.terminal], settlement.progress, ); if (!terminal) throw new Error("Consumer failure settlement did not produce a terminal receipt"); return terminal; }
async settleConsumerAbandonment(run: AgentRun, terminal: EventCandidate): Promise<ThoughtEvent> { const [event] = await this.settleRunAndEvents(run, [terminal]); if (!event) throw new Error("Consumer abandonment did not produce a terminal receipt"); return event; }
async flush(): Promise<void> { // alpha55: the backend JazzClient exposes flush() (local settlement // forwarding), replacing the old JazzContext.flush(). It is not a // durability barrier; writes that must be durable before continuing use // `wait({ tier })`, which this store applies on every required write. this.client.flush(); }
async close(): Promise<void> { // Release order matters: closing the session shuts down its client (and // its subscriptions), then a locally-owned server is stopped so a failed // or completed open never leaks the storage-owner process. The discovery // file is removed first so no new process attaches to a server that is // about to stop; a concurrent attacher that already read it fails its // reachability check or its session open and falls back to ownership. if (this.localServer) { const dataPath = this.options.dataPath ?? path.join(this.options.projectRoot, ".thoughtstream", "state", "jazz.sqlite"); await clearServerDiscoveryFile(dataPath); } // alpha55 native shutdown quiesces the foreground transaction clock via a // synchronous native `block_on` (jazz-napi `foreground_tx_time_high_water`). // That native future can depend on core work the runtime only advances from // JS-scheduled ticks (`setTickScheduler` → microtasks/timers), so closing // with unsettled foreground transactions can leave the main thread spinning // in `block_on` with no way to deliver the tick — a self-deadlock under // load. Retire store-owned subscriptions and drain pending writes through // the JS-side poll path (`shutdown({ waitForSync: true })` → // `waitForPendingWrites`, which yields via `sleep(0)`) BEFORE the session // begins its teardown, so the later native quiesce finds an idle core. for (const unsubscribe of this.subscriptions.splice(0)) unsubscribe(); await new Promise<void>((resolve) => setImmediate(resolve)); try { await this.client.shutdown({ waitForSync: true }); } catch (error) { if (!isGracefulShutdownSyncError(error)) throw error; // A cancelled graceful drain leaves the client usable; the session // close below still performs the full teardown. } await this.session.close(); await this.localServer?.stop(); }
/** Unsubscribe functions created by this store's `subscribe*` methods. */ private readonly subscriptions: Array<() => void> = [];
private async settleConsumerTerminal( run: AgentRun, inputEvent: ThoughtEvent, candidates: EventCandidate[], progress: ConsumerProgress, ): Promise<ThoughtEvent[]> { if (progress.lastSequence !== inputEvent.sourceSequence || progress.lastEventId !== inputEvent.id) { throw new Error("Consumer progress must settle exactly at the terminal input event"); } if (progress.source !== inputEvent.source) throw new Error("Consumer progress source does not match input event source"); return this.settleRunAndEvents(run, candidates, progress); }
private async settleRunAndEvents( run: AgentRun, candidates: EventCandidate[], progress?: ConsumerProgress, ): Promise<ThoughtEvent[]> { const outputSource = candidates[0]?.source; if (!outputSource || candidates.some((candidate) => candidate.source !== outputSource)) { throw new Error("Consumer settlement receipts must share one output source"); } for (let attempt = 0; attempt < 8; attempt += 1) { const storageRevision = this.consumerStorageRevision; const settledAt = new Date().toISOString(); await this.ensureTransactionReady(outputSource); // Persistent/edge Jazz transactions do not inherit the readiness barrier of // ordinary Db reads. Prime every pre-existing row the transaction may update // or deduplicate, not only the output source. Otherwise a restarted consumer // can miss its durable run/progress rows and leave a partially allocated // lifecycle object when the transaction fails. const [sourceSnapshot] = await this.db.all(thoughtstreamApp.sources.where({ key: outputSource, }).limit(1), { tier: this.durabilityTier }); const [runSnapshot] = await this.db.all(thoughtstreamApp.runs.where({ key: run.id }).limit(1), { tier: this.durabilityTier, }); if (!runSnapshot) throw new Error(`Cannot settle missing run ${run.id}`); const eventSnapshots = new Map<string, Record<string, unknown>>(); for (const candidate of candidates) { const id = eventId(candidate); const [eventSnapshot] = await this.db.all(thoughtstreamApp.events.where({ key: id, }).limit(1), { tier: this.durabilityTier }); if (eventSnapshot) eventSnapshots.set(id, eventSnapshot); } let progressSnapshot: Record<string, unknown> | undefined; if (progress) { [progressSnapshot] = await this.db.all(thoughtstreamApp.consumerProgress.where({ key: progress.id, }).limit(1), { tier: this.durabilityTier }); } try { const result = await this.db.transaction(async (tx) => { const sourceRow = sourceSnapshot; if (sourceRow && String(sourceRow.kind) !== candidates[0]!.sourceKind) { throw new Error(`Source kind changed for ${outputSource}`); } let sequence = Number(sourceRow?.lastSequence ?? 0); const events: ThoughtEvent[] = []; for (const candidate of candidates) { const id = eventId(candidate); const existing = eventSnapshots.get(id); if (existing) { const event = eventFromJazz(existing); assertSameEventIdentity(event, candidate); events.push(event); continue; } sequence += 1; const event = buildEvent(candidate, id, sequence, settledAt, this.runtimeRevision, this.registry); tx.insert(thoughtstreamApp.events, eventToJazz(event), { id: jazzRowId(`consumer-event-v${storageRevision}`, id), }); events.push(event); } const sourceData = { key: outputSource, kind: candidates[0]!.sourceKind, enabled: true, configJson: sourceRow ? String(sourceRow.configJson) : "{}", lastSequence: sequence, createdAt: sourceRow ? String(sourceRow.createdAt) : settledAt, updatedAt: settledAt, }; if (sourceRow) tx.update(thoughtstreamApp.sources, String(sourceRow.id), sourceData); else { tx.insert(thoughtstreamApp.sources, sourceData, { id: jazzRowId(`consumer-source-v${storageRevision}`, outputSource), }); }
tx.update(thoughtstreamApp.runs, String(runSnapshot.id), runToJazz(run)); if (progress) await upsertProgressInTransaction(tx, progress, progressSnapshot, storageRevision); return events; }); await result.wait({ tier: this.durabilityTier }); return result.value; } catch (error) { if (!isRecoverableStorageError(error) || attempt === 7) throw error; this.consumerStorageRevision = storageRevision + 1; await new Promise((resolve) => setTimeout(resolve, Math.min(500, 20 * (2 ** attempt)))); } } throw new Error("Consumer transaction exhausted storage recovery attempts"); }
// Jazz transactions commit the account and reservation together but do not serialize // concurrent read-modify-write callbacks. The process-wide queue provides a strict // per-process cap; independent processes provide only best-effort aggregate enforcement. private async withInferenceAccountLock<T>(accountId: string, operation: () => Promise<T>): Promise<T> { const previous = inferenceAccountQueues.get(accountId) ?? Promise.resolve(); let release!: () => void; const lock = new Promise<void>((resolve) => { release = resolve; }); const queued = previous.then(() => lock); inferenceAccountQueues.set(accountId, queued); await previous; try { return await operation(); } finally { release(); if (inferenceAccountQueues.get(accountId) === queued) inferenceAccountQueues.delete(accountId); } }
private async ensureInferenceBudgetAccount( scopeType: InferenceReservationRequest["scopeType"], scopeKey: string, accountId: string, at: string, ): Promise<Record<string, unknown>> { const existing = await this.db.all(thoughtstreamApp.inferenceBudgetAccounts.where({ key: accountId, }).limit(1), { tier: this.durabilityTier }); if (existing[0]) return existing[0]; const account = emptyBudgetAccount(scopeType, scopeKey); const data = { key: accountId, scopeType, scopeKey, windowsJson: canonicalJson(asJsonValue(account.windows)), activeLeasesJson: canonicalJson(asJsonValue(account.activeLeases)), updatedAt: at, }; for (let attempt = 0; attempt < 8; attempt += 1) { const storageRevision = this.inferenceStorageRevision; try { await this.wait(this.db.insert(thoughtstreamApp.inferenceBudgetAccounts, data, { id: jazzRowId(`inference-budget-account-v${storageRevision}`, accountId), })); } catch (error) { if (!isRecoverableStorageError(error) || attempt === 7) throw error; this.inferenceStorageRevision = storageRevision + 1; } const rows = await this.db.all(thoughtstreamApp.inferenceBudgetAccounts.where({ key: accountId, }).limit(1), { tier: this.durabilityTier }); if (rows[0]) return rows[0]; this.inferenceStorageRevision = Math.max(this.inferenceStorageRevision, storageRevision + 1); await new Promise((resolve) => setTimeout(resolve, Math.min(500, 20 * (2 ** attempt)))); } throw new Error(`Inference budget account exhausted storage recovery attempts: ${accountId}`); }
private async upsert(table: JazzTable, query: Record<string, unknown>, data: Record<string, unknown>): Promise<void> { const existing = await this.db.one(table.where(query).limit(1) as never); if (existing && typeof existing === "object" && "id" in existing) { await this.wait(this.db.update(table as never, String(existing.id), data)); } else { await this.wait(this.db.insert(table as never, data)); } }
private async wait(write: { wait(options: { tier: "local" | "edge" | "global" }): Promise<unknown> }): Promise<void> { await write.wait({ tier: this.durabilityTier }); }
private async ensureTransactionReady(source: string): Promise<void> { // TransactionScope reads do not pass through Db's readiness barrier. Prime the // persistent runtime before a transaction queries deterministic rows, or a // fresh process can miss an existing source and collide while reinserting it. await this.db.all(thoughtstreamApp.sources.where({ key: source }).limit(1), { tier: this.durabilityTier }); }}
interface BudgetWindowState { key: string; window: InferenceBudgetLimit["window"]; startedAt: string; endsAt: string; limit: InferenceBudgetLimit; calls: number; inputTokens: number; outputTokens: number; costMicrousd: number; charges?: RollingInferenceCharge[] | undefined;}
interface RollingInferenceCharge { reservationId: string; chargedAt: string; charge: InferenceCharge;}
interface ActiveInferenceLease { reservationId: string; leaseExpiresAt: string;}
interface BudgetAccountState { scopeType: InferenceReservationRequest["scopeType"]; scopeKey: string; windows: BudgetWindowState[]; activeLeases: ActiveInferenceLease[];}
function emptyBudgetAccount( scopeType: InferenceReservationRequest["scopeType"], scopeKey: string,): BudgetAccountState { return { scopeType, scopeKey, windows: [], activeLeases: [] };}
function inferenceAccountId(scopeType: InferenceReservationRequest["scopeType"], scopeKey: string): string { return `${scopeType}:${scopeKey}`;}
function budgetWindowKey(limit: InferenceBudgetLimit): string { return limit.window === "rolling" ? `rolling:${limit.durationMs}` : limit.window;}
function reconcileBudgetWindows( current: BudgetWindowState[], policy: InferenceBudgetPolicy, at: string,): BudgetWindowState[] { const atMs = Date.parse(at); if (!Number.isFinite(atMs)) throw new Error("Inference reservation time must be an ISO timestamp"); const prior = new Map(current.map((window) => [window.key, window])); return policy.limits.map((limit) => { const key = budgetWindowKey(limit); const existing = prior.get(key); if (limit.window === "rolling") { const durationMs = limit.durationMs!; const cutoffMs = atMs - durationMs; const charges = (existing?.charges ?? []).filter((entry) => Date.parse(entry.chargedAt) > cutoffMs); const totals = totalRollingCharges(charges); return { key, window: limit.window, startedAt: new Date(cutoffMs).toISOString(), endsAt: at, limit: { ...limit }, ...totals, charges, }; } const bounds = budgetWindowBounds(limit, atMs); if (existing && existing.startedAt === bounds.startedAt && existing.endsAt === bounds.endsAt) { return { ...existing, limit: { ...limit } }; } return { key, window: limit.window, startedAt: bounds.startedAt, endsAt: bounds.endsAt, limit: { ...limit }, calls: 0, inputTokens: 0, outputTokens: 0, costMicrousd: 0, }; });}
function budgetWindowBounds( limit: InferenceBudgetLimit, atMs: number,): { startedAt: string; endsAt: string } { const durationMs = limit.window === "hour" ? 3_600_000 : 86_400_000; const startedMs = Math.floor(atMs / durationMs) * durationMs; return { startedAt: new Date(startedMs).toISOString(), endsAt: new Date(startedMs + durationMs).toISOString() };}
function exceedsBudgetLimit(window: BudgetWindowState, estimate: InferenceReservationEstimate): boolean { return window.calls + estimate.calls > window.limit.maxCalls || window.inputTokens + estimate.inputTokens > window.limit.maxInputTokens || window.outputTokens + estimate.outputTokens > window.limit.maxOutputTokens || (window.limit.maxCostMicrousd !== undefined && window.costMicrousd + (estimate.costMicrousd ?? 0) > window.limit.maxCostMicrousd);}
function chargeBudgetWindow( window: BudgetWindowState, charge: InferenceCharge, reservationId: string, chargedAt: string,): BudgetWindowState { if (window.window === "rolling") { const charges = [...(window.charges ?? []), { reservationId, chargedAt, charge }]; return { ...window, ...totalRollingCharges(charges), charges }; } return { ...window, calls: window.calls + charge.calls, inputTokens: window.inputTokens + charge.inputTokens, outputTokens: window.outputTokens + charge.outputTokens, costMicrousd: window.costMicrousd + (charge.costMicrousd ?? 0), };}
function adjustBudgetWindowCharge( window: BudgetWindowState, estimate: InferenceReservationEstimate, charged: InferenceCharge, reservationId: string,): BudgetWindowState { if (window.window === "rolling") { const charges = (window.charges ?? []).map((entry) => entry.reservationId === reservationId ? { ...entry, charge: charged } : entry); return { ...window, ...totalRollingCharges(charges), charges }; } return { ...window, calls: Math.max(0, window.calls - estimate.calls + charged.calls), inputTokens: Math.max(0, window.inputTokens - estimate.inputTokens + charged.inputTokens), outputTokens: Math.max(0, window.outputTokens - estimate.outputTokens + charged.outputTokens), costMicrousd: Math.max( 0, window.costMicrousd - (estimate.costMicrousd ?? 0) + (charged.costMicrousd ?? 0), ), };}
function totalRollingCharges(charges: RollingInferenceCharge[]): Pick< BudgetWindowState, "calls" | "inputTokens" | "outputTokens" | "costMicrousd"> { return charges.reduce((total, entry) => ({ calls: total.calls + entry.charge.calls, inputTokens: total.inputTokens + entry.charge.inputTokens, outputTokens: total.outputTokens + entry.charge.outputTokens, costMicrousd: total.costMicrousd + (entry.charge.costMicrousd ?? 0), }), { calls: 0, inputTokens: 0, outputTokens: 0, costMicrousd: 0 });}
function zeroInferenceCharge(trackCost: boolean): InferenceCharge { return { calls: 0, inputTokens: 0, outputTokens: 0, ...(trackCost ? { costMicrousd: 0 } : {}), };}
function normalizeInferenceUsage(usage: InferenceUsage | undefined): InferenceUsage | undefined { if (!usage) return undefined; const normalized: InferenceUsage = {}; for (const key of ["inputTokens", "outputTokens", "costMicrousd"] as const) { const value = usage[key]; if (value === undefined) continue; if (!Number.isSafeInteger(value) || value < 0) throw new Error(`Inference usage ${key} must be a nonnegative integer`); normalized[key] = value; } return Object.keys(normalized).length > 0 ? normalized : undefined;}
function chargedInferenceUsage( estimate: InferenceReservationEstimate, usage: InferenceUsage | undefined,): InferenceCharge { return { calls: 1, inputTokens: usage?.inputTokens ?? estimate.inputTokens, outputTokens: usage?.outputTokens ?? estimate.outputTokens, ...(estimate.costMicrousd !== undefined ? { costMicrousd: usage?.costMicrousd ?? estimate.costMicrousd } : {}), };}
function inferenceUsageStatus( usage: InferenceUsage | undefined, estimate: InferenceReservationEstimate,): InferenceAccountingRecord["usageStatus"] { if (!usage) return "unavailable"; const tracked = [usage.inputTokens, usage.outputTokens]; if (estimate.costMicrousd !== undefined) tracked.push(usage.costMicrousd); if (tracked.every((value) => value !== undefined)) return "reported"; return tracked.some((value) => value !== undefined) ? "partial" : "unavailable";}
function assertReservationRequest(request: InferenceReservationRequest): void { if (!request.reservationId || !request.scopeKey || !request.runId || !request.agentId || !request.provider || !request.model) { throw new Error("Inference reservation identity is incomplete"); } if (request.estimate.calls !== 1) throw new Error("Inference reservations must reserve exactly one provider call"); for (const [key, value] of Object.entries(request.estimate)) { if (!Number.isSafeInteger(value) || value <= 0) throw new Error(`Inference reservation ${key} must be a positive integer`); } if (request.policy.reservation.inputTokens !== request.estimate.inputTokens || request.policy.reservation.outputTokens !== request.estimate.outputTokens || request.policy.reservation.costMicrousd !== request.estimate.costMicrousd) { throw new Error("Inference estimate must match the policy reservation"); } if (request.policy.limits.length === 0) throw new Error("Inference policy must declare at least one budget limit"); for (const limit of request.policy.limits) { for (const key of ["maxCalls", "maxInputTokens", "maxOutputTokens"] as const) { if (!Number.isSafeInteger(limit[key]) || limit[key] <= 0) { throw new Error(`Inference budget ${key} must be a positive integer`); } } if (limit.maxCostMicrousd !== undefined && (!Number.isSafeInteger(limit.maxCostMicrousd) || limit.maxCostMicrousd <= 0)) { throw new Error("Inference budget maxCostMicrousd must be a positive integer"); } if (limit.maxCostMicrousd !== undefined && (request.policy.reservation.costMicrousd === undefined || request.estimate.costMicrousd === undefined)) { throw new Error("Inference cost limits require a cost reservation"); } } if (!Number.isFinite(Date.parse(request.reservedAt))) throw new Error("Inference reservation time must be an ISO timestamp");}
function assertSameReservation(existing: InferenceAccountingRecord, request: InferenceReservationRequest): void { if (existing.scopeType !== request.scopeType || existing.scopeKey !== request.scopeKey || existing.runId !== request.runId || existing.agentId !== request.agentId || existing.agentVersion !== request.agentVersion || existing.provider !== request.provider || existing.model !== request.model || canonicalJson(asJsonValue(existing.estimate)) !== canonicalJson(asJsonValue(request.estimate))) { throw new Error(`Inference reservation identity conflict for ${request.reservationId}`); }}
async function expireStaleReservations( tx: TransactionScope, account: BudgetAccountState, at: string,): Promise<void> { const atMs = Date.parse(at); const retained: ActiveInferenceLease[] = []; for (const lease of account.activeLeases) { if (Date.parse(lease.leaseExpiresAt) > atMs) { retained.push(lease); continue; } const recordRow = await tx.one(thoughtstreamApp.inferenceAccounting.where({ id: jazzRowId("inference-accounting", lease.reservationId), }).limit(1)); if (!recordRow) throw new Error(`Inference lease has no accounting record: ${lease.reservationId}`); const record = inferenceAccountingFromJazz(recordRow); if (record.status === "reserved") { const expired: InferenceAccountingRecord = { ...record, status: "expired", usageStatus: "unavailable", settledAt: at, updatedAt: at, }; tx.update(thoughtstreamApp.inferenceAccounting, String(recordRow.id), inferenceAccountingToJazz(expired)); } } account.activeLeases = retained;}
function upsertBudgetAccountInTransaction( tx: TransactionScope, row: Record<string, unknown> | null | undefined, accountId: string, account: BudgetAccountState, at: string,): void { const data = { key: accountId, scopeType: account.scopeType, scopeKey: account.scopeKey, windowsJson: canonicalJson(asJsonValue(account.windows)), activeLeasesJson: canonicalJson(asJsonValue(account.activeLeases)), updatedAt: at, }; if (row) tx.update(thoughtstreamApp.inferenceBudgetAccounts, String(row.id), data); else tx.insert(thoughtstreamApp.inferenceBudgetAccounts, data, { id: jazzRowId("inference-budget-account", accountId) });}
function budgetAccountFromJazz(row: Record<string, unknown>): BudgetAccountState { return { scopeType: String(row.scopeType) as BudgetAccountState["scopeType"], scopeKey: String(row.scopeKey), windows: JSON.parse(String(row.windowsJson)) as BudgetWindowState[], activeLeases: JSON.parse(String(row.activeLeasesJson)) as ActiveInferenceLease[], };}
function inferenceAccountingToJazz(record: InferenceAccountingRecord): Record<string, unknown> { return { key: record.id, scopeType: record.scopeType, scopeKey: record.scopeKey, runKey: record.runId, agentKey: record.agentId, agentVersion: record.agentVersion, provider: record.provider, model: record.model, status: record.status, usageStatus: record.usageStatus, estimateJson: canonicalJson(asJsonValue(record.estimate)), chargedJson: canonicalJson(asJsonValue(record.charged)), actualUsageJson: record.actualUsage ? canonicalJson(asJsonValue(record.actualUsage)) : "", windowKeysJson: canonicalJson(record.windowKeys), reservedAt: record.reservedAt, leaseExpiresAt: record.leaseExpiresAt, settledAt: record.settledAt ?? "", denialReason: record.denialReason ?? "", limitingWindowKeysJson: canonicalJson(record.limitingWindowKeys ?? []), updatedAt: record.updatedAt, };}
function inferenceAccountingFromJazz(row: Record<string, unknown>): InferenceAccountingRecord { const actualUsage = String(row.actualUsageJson ?? ""); const settledAt = String(row.settledAt ?? ""); const denialReason = String(row.denialReason ?? ""); const limitingWindowKeys = JSON.parse(String(row.limitingWindowKeysJson ?? "[]")) as string[]; return { id: String(row.key), scopeType: String(row.scopeType) as InferenceAccountingRecord["scopeType"], scopeKey: String(row.scopeKey), runId: String(row.runKey), agentId: String(row.agentKey), agentVersion: Number(row.agentVersion), provider: String(row.provider), model: String(row.model), status: String(row.status) as InferenceAccountingRecord["status"], usageStatus: String(row.usageStatus) as InferenceAccountingRecord["usageStatus"], estimate: JSON.parse(String(row.estimateJson)) as InferenceReservationEstimate, charged: JSON.parse(String(row.chargedJson)) as InferenceCharge, ...(actualUsage ? { actualUsage: JSON.parse(actualUsage) as InferenceUsage } : {}), windowKeys: JSON.parse(String(row.windowKeysJson)) as string[], reservedAt: String(row.reservedAt), leaseExpiresAt: String(row.leaseExpiresAt), ...(settledAt ? { settledAt } : {}), ...(denialReason ? { denialReason } : {}), ...(limitingWindowKeys.length > 0 ? { limitingWindowKeys } : {}), updatedAt: String(row.updatedAt), };}
function asJsonValue(value: unknown): JsonValue { return JSON.parse(JSON.stringify(value)) as JsonValue;}
function eventToJazz(event: ThoughtEvent): Record<string, unknown> { return { key: event.id, sourceSequence: event.sourceSequence, type: event.type, schemaVersion: event.schemaVersion, source: event.source, sourceKind: event.sourceKind, externalId: event.externalId, idempotencyKey: event.idempotencyKey, occurredAt: event.occurredAt, observedAt: event.observedAt, actor: event.actor, rootEventKey: event.rootEventId, parentEventKey: event.parentEventId ?? "", correlationId: event.correlationId, privacy: event.privacy, payloadJson: canonicalJson(event.payload), payloadHash: event.payloadHash, traceKey: event.traceId ?? "", createdByRuntime: event.createdByRuntime, };}
function eventFromJazz(row: Record<string, unknown>): ThoughtEvent { const parentEventId = String(row.parentEventKey ?? ""); const traceId = String(row.traceKey ?? ""); return { id: String(row.key), sourceSequence: Number(row.sourceSequence ?? 0), type: String(row.type), schemaVersion: Number(row.schemaVersion), source: String(row.source), sourceKind: String(row.sourceKind) as ThoughtEvent["sourceKind"], externalId: String(row.externalId), idempotencyKey: String(row.idempotencyKey), occurredAt: String(row.occurredAt), observedAt: String(row.observedAt), actor: String(row.actor), rootEventId: String(row.rootEventKey), ...(parentEventId ? { parentEventId } : {}), correlationId: String(row.correlationId), privacy: String(row.privacy) as ThoughtEvent["privacy"], payload: parseJsonObject(String(row.payloadJson)), payloadHash: String(row.payloadHash), ...(traceId ? { traceId } : {}), createdByRuntime: String(row.createdByRuntime), ...(row.id ? { jazzRowId: String(row.id) } : {}), };}
function documentVersionFromJazz(row: Record<string, unknown>): DocumentVersion { return { id: String(row.key), source: String(row.source), documentId: String(row.documentKey), path: String(row.path), contentType: String(row.contentType), sha256: String(row.sha256), content: String(row.content), sizeBytes: Number(row.sizeBytes), mtimeMs: Number(row.mtimeMs), createdAt: String(row.createdAt), };}
function currentDocumentFromJazz(row: Record<string, unknown>): CurrentDocument { return { id: String(row.key), source: String(row.source), documentId: String(row.documentKey), path: String(row.path), versionId: String(row.versionKey), sha256: String(row.sha256), contentType: String(row.contentType), sizeBytes: Number(row.sizeBytes), mtimeMs: Number(row.mtimeMs), deleted: Boolean(row.deleted), updatedAt: String(row.updatedAt), };}
function agentFromJazz(row: Record<string, unknown>): StoredAgent { return { id: String(row.key), version: Number(row.version), enabled: Boolean(row.enabled), spec: parseJsonObject(String(row.specJson)), specHash: String(row.specHash), updatedAt: String(row.updatedAt) };}
function runToJazz(run: AgentRun): Record<string, unknown> { return { key: run.id, executionKey: run.executionKey, triggerEventKey: run.triggerEventId, agentKey: run.agentId, agentVersion: run.agentVersion, status: run.status, inputEventIdsJson: canonicalJson(run.inputEventIds), outputEventIdsJson: canonicalJson(run.outputEventIds), attempt: run.attempt, provider: run.provider, model: run.model, adapterRevision: run.adapterRevision ?? "", promptHash: run.promptHash, contextManifestJson: canonicalJson(run.contextManifest), ...(inferenceAccountingEnabled ? { accountingReservationKey: run.accountingReservationId ?? "" } : {}), resultJson: encodeRunResult(run), errorText: run.errorText ?? "", createdAt: run.createdAt, startedAt: run.startedAt ?? "", completedAt: run.completedAt ?? "", updatedAt: run.updatedAt, };}
function runFromJazz(row: Record<string, unknown>): AgentRun { const inputEventIds = JSON.parse(String(row.inputEventIdsJson)) as string[]; const outputEventIds = JSON.parse(String(row.outputEventIdsJson)) as string[]; const persistedResult = decodeRunResult(String(row.resultJson ?? "")); return { id: String(row.key), executionKey: String(row.executionKey ?? row.key), triggerEventId: String(row.triggerEventKey ?? ""), agentId: String(row.agentKey), agentVersion: Number(row.agentVersion), status: String(row.status) as AgentRun["status"], inputEventIds, outputEventIds, attempt: Number(row.attempt), provider: String(row.provider), model: String(row.model), ...(persistedResult.privacy ? { privacy: persistedResult.privacy } : {}), ...(persistedResult.checkpointRevision ? { checkpointRevision: persistedResult.checkpointRevision } : {}), ...(String(row.adapterRevision ?? "") ? { adapterRevision: String(row.adapterRevision), executionAdapterRevision: String(row.adapterRevision), } : {}), ...(persistedResult.executionAdapterRevision ? { executionAdapterRevision: persistedResult.executionAdapterRevision } : {}), ...(persistedResult.modelAdapter ? { modelAdapter: persistedResult.modelAdapter } : {}), ...(persistedResult.adapterCatalogDigest ? { adapterCatalogDigest: persistedResult.adapterCatalogDigest } : {}), ...(persistedResult.adapterCatalogGeneration ? { adapterCatalogGeneration: persistedResult.adapterCatalogGeneration } : {}), promptHash: String(row.promptHash), contextManifest: parseJsonObject(String(row.contextManifestJson)), ...(String(row.accountingReservationKey ?? "") ? { accountingReservationId: String(row.accountingReservationKey) } : {}), ...(persistedResult.result ? { result: persistedResult.result } : {}), ...(String(row.errorText ?? "") ? { errorText: String(row.errorText) } : {}), createdAt: String(row.createdAt), ...(String(row.startedAt ?? "") ? { startedAt: String(row.startedAt) } : {}), ...(String(row.completedAt ?? "") ? { completedAt: String(row.completedAt) } : {}), updatedAt: String(row.updatedAt), };}
const RUN_RESULT_ENVELOPE_KEY = "__thoughtstreamRunPersistenceEnvelopeV1";
function encodeRunResult(run: AgentRun): string { const hasPrivateEnvelope = Boolean( run.checkpointRevision || run.privacy || run.executionAdapterRevision || run.modelAdapter || run.adapterCatalogDigest || run.adapterCatalogGeneration, ); if (!hasPrivateEnvelope) return run.result ? canonicalJson(run.result) : ""; return canonicalJson({ [RUN_RESULT_ENVELOPE_KEY]: { ...(run.checkpointRevision ? { checkpointRevision: run.checkpointRevision } : {}), ...(run.privacy ? { privacy: run.privacy } : {}), ...(run.executionAdapterRevision ? { executionAdapterRevision: run.executionAdapterRevision } : {}), ...(run.modelAdapter ? { modelAdapter: run.modelAdapter as unknown as JsonObject } : {}), ...(run.adapterCatalogDigest ? { adapterCatalogDigest: run.adapterCatalogDigest } : {}), ...(run.adapterCatalogGeneration ? { adapterCatalogGeneration: run.adapterCatalogGeneration } : {}), ...(run.result ? { result: run.result } : {}), }, });}
function decodeRunResult(value: string): { result?: JsonObject; privacy?: ThoughtEvent["privacy"]; checkpointRevision?: string; executionAdapterRevision?: string; modelAdapter?: ModelAdapterIdentity; adapterCatalogDigest?: string; adapterCatalogGeneration?: number;} { if (!value) return {}; const parsed = parseJsonObject(value); const envelope = parsed[RUN_RESULT_ENVELOPE_KEY]; if (!envelope || typeof envelope !== "object" || Array.isArray(envelope)) return { result: parsed }; const checkpointRevision = typeof envelope.checkpointRevision === "string" ? envelope.checkpointRevision : undefined; const privacy = ["public-source", "private", "sensitive"].includes(String(envelope.privacy)) ? envelope.privacy as ThoughtEvent["privacy"] : undefined; const executionAdapterRevision = typeof envelope.executionAdapterRevision === "string" ? envelope.executionAdapterRevision : undefined; const modelAdapter = envelope.modelAdapter === undefined ? undefined : modelAdapterIdentitySchema.parse(envelope.modelAdapter); const adapterCatalogDigest = typeof envelope.adapterCatalogDigest === "string" ? envelope.adapterCatalogDigest : undefined; const adapterCatalogGeneration = typeof envelope.adapterCatalogGeneration === "number" && Number.isSafeInteger(envelope.adapterCatalogGeneration) && envelope.adapterCatalogGeneration > 0 ? envelope.adapterCatalogGeneration : undefined; const result = envelope.result; return { ...(result && typeof result === "object" && !Array.isArray(result) ? { result: result as JsonObject } : {}), ...(privacy ? { privacy } : {}), ...(checkpointRevision ? { checkpointRevision } : {}), ...(executionAdapterRevision ? { executionAdapterRevision } : {}), ...(modelAdapter ? { modelAdapter } : {}), ...(adapterCatalogDigest ? { adapterCatalogDigest } : {}), ...(adapterCatalogGeneration ? { adapterCatalogGeneration } : {}), };}
function traceFromJazz(row: Record<string, unknown>): TraceChunk { return { id: String(row.key), runId: String(row.runKey), sequence: Number(row.sequence), type: String(row.type), payload: parseJsonObject(String(row.payloadJson)), createdAt: String(row.createdAt) };}
function projectionFromJazz(row: Record<string, unknown>): Projection { return { id: String(row.key), payload: parseJsonObject(String(row.payloadJson)), lastEventId: String(row.lastEventKey), projectionVersion: Number(row.projectionVersion), updatedAt: String(row.updatedAt) };}
function sourceCursorFromJazz(row: Record<string, unknown>): SourceCursor { const lastSuccessAt = String(row.lastSuccessAt ?? ""); const lastFailureAt = String(row.lastFailureAt ?? ""); const lastError = String(row.lastError ?? ""); return { id: String(row.key), source: String(row.source), cursor: parseJsonObject(String(row.cursorJson)), ...(lastSuccessAt ? { lastSuccessAt } : {}), ...(lastFailureAt ? { lastFailureAt } : {}), ...(lastError ? { lastError } : {}), updatedAt: String(row.updatedAt), };}
function buildEvent( candidate: EventCandidate, id: string, sourceSequence: number, defaultObservedAt: string, runtimeRevision: string, registry: EventRegistry,): ThoughtEvent { const payload = registry.validateEvent( candidate.type, candidate.schemaVersion, candidate.privacy, candidate.payload, ); return { id, sourceSequence, type: candidate.type, schemaVersion: candidate.schemaVersion, source: candidate.source, sourceKind: candidate.sourceKind, externalId: candidate.externalId, idempotencyKey: candidate.idempotencyKey, occurredAt: candidate.occurredAt, observedAt: candidate.observedAt ?? defaultObservedAt, actor: candidate.actor, rootEventId: candidate.rootEventId ?? id, ...(candidate.parentEventId ? { parentEventId: candidate.parentEventId } : {}), correlationId: candidate.correlationId, privacy: candidate.privacy, payload, payloadHash: hashJson(payload), ...(candidate.traceId ? { traceId: candidate.traceId } : {}), createdByRuntime: candidate.createdByRuntime ?? runtimeRevision, };}
export function eventIdFor(source: string, idempotencyKey: string): string { return `evt_${sha256(`${source}\u0000${idempotencyKey}`).slice(0, 48)}`;}
function eventId(candidate: EventCandidate): string { return eventIdFor(candidate.source, candidate.idempotencyKey);}
function jazzRowId(namespace: string, key: string): string { const bytes = createHash("sha256").update(`${namespace}\u0000${key}`).digest().subarray(0, 16); bytes[6] = (bytes[6]! & 0x0f) | 0x50; bytes[8] = (bytes[8]! & 0x3f) | 0x80; const hex = bytes.toString("hex"); return `${hex.slice(0, 8)}-${hex.slice(8, 12)}-${hex.slice(12, 16)}-${hex.slice(16, 20)}-${hex.slice(20)}`;}
function assertSameEventIdentity(event: ThoughtEvent, candidate: EventCandidate): void { if (event.source !== candidate.source || event.idempotencyKey !== candidate.idempotencyKey || event.type !== candidate.type || event.schemaVersion !== candidate.schemaVersion) { throw new Error(`Deterministic event identity conflict for ${event.id}`); } if (event.payloadHash !== hashJson(candidate.payload)) { throw new Error(`Idempotency key reused with different payload for ${event.id}`); }}
function sourceStateFromJazz(row: Record<string, unknown>): SourceState { return { id: String(row.key), kind: String(row.kind), enabled: Boolean(row.enabled), config: parseJsonObject(String(row.configJson)), lastSequence: Number(row.lastSequence ?? 0), createdAt: String(row.createdAt), updatedAt: String(row.updatedAt), };}
function consumerProgressFromJazz(row: Record<string, unknown>): ConsumerProgress { return { id: String(row.key), consumerId: String(row.consumerKey), consumerVersion: Number(row.consumerVersion), source: String(row.source), lastSequence: Number(row.lastSequence), lastEventId: String(row.lastEventKey), updatedAt: String(row.updatedAt), };}
function lettaConversationBindingFromPayload(payload: JsonObject): LettaConversationBinding { const scopeType = String(payload.scopeType); if (scopeType !== "document") throw new Error(`Unsupported Letta conversation scope: ${scopeType}`); const binding: LettaConversationBinding = { id: String(payload.id), declarationId: String(payload.declarationId), declarationVersion: Number(payload.declarationVersion), agentId: String(payload.agentId), source: String(payload.source), scopeType, scopeKey: String(payload.scopeKey), conversationId: String(payload.conversationId), remoteMarker: String(payload.remoteMarker), currentPath: String(payload.currentPath), createdAt: String(payload.createdAt), updatedAt: String(payload.updatedAt), }; if (!binding.id || !binding.declarationId || !Number.isSafeInteger(binding.declarationVersion) || binding.declarationVersion <= 0 || !binding.agentId || !binding.source || !binding.scopeKey || !binding.conversationId || !binding.remoteMarker || !binding.currentPath || !Number.isFinite(Date.parse(binding.createdAt)) || !Number.isFinite(Date.parse(binding.updatedAt))) { throw new Error("Letta conversation binding payload is invalid"); } return binding;}
function lettaConversationBindingPayload(binding: LettaConversationBinding): JsonObject { return { id: binding.id, declarationId: binding.declarationId, declarationVersion: binding.declarationVersion, agentId: binding.agentId, source: binding.source, scopeType: binding.scopeType, scopeKey: binding.scopeKey, conversationId: binding.conversationId, remoteMarker: binding.remoteMarker, currentPath: binding.currentPath, createdAt: binding.createdAt, updatedAt: binding.updatedAt, };}
function lettaConversationProjectionId(bindingId: string): string { return stableKey("projection", "letta-conversation-binding", bindingId);}
function lettaConversationProjection(binding: LettaConversationBinding, lastEventId: string): Projection { return { id: lettaConversationProjectionId(binding.id), payload: lettaConversationBindingPayload(binding), lastEventId, projectionVersion: 1, updatedAt: binding.updatedAt, };}
function assertSameLettaConversationBinding( existing: LettaConversationBinding, candidate: LettaConversationBinding,): void { for (const key of [ "id", "declarationId", "agentId", "source", "scopeType", "scopeKey", "conversationId", "remoteMarker", ] as const) { if (existing[key] !== candidate[key]) { throw new Error(`Letta conversation binding conflict for ${existing.id}: ${key}`); } }}
function cursorToJazz(cursor: SourceCursor): Record<string, unknown> { return { key: cursor.id, source: cursor.source, cursorJson: canonicalJson(cursor.cursor), lastSuccessAt: cursor.lastSuccessAt ?? "", lastFailureAt: cursor.lastFailureAt ?? "", lastError: cursor.lastError ?? "", updatedAt: cursor.updatedAt, };}
async function upsertCursorInTransaction( tx: TransactionScope, cursor: SourceCursor, storageRevision = 2,): Promise<void> { const row = await tx.one(thoughtstreamApp.sourceCursors.where({ key: cursor.id }).limit(1)); const data = cursorToJazz(cursor); if (row) tx.update(thoughtstreamApp.sourceCursors, String(row.id), data); else { tx.insert(thoughtstreamApp.sourceCursors, data, { id: jazzRowId(`cursor-v${storageRevision}`, cursor.id), }); }}
/** * Local-owner discovery (alpha55): the storage-owning process records its * server URL and backend secret in a file next to the database directory so * sibling processes can attach as clients instead of failing on the exclusive * RocksDB lock. The owner removes the file on close; a stale file (owner * crashed without cleanup) is detected by the reachability probe and * reclaimed by the next opener. */interface LocalServerDiscovery { url: string; backendSecret: string;}
function serverDiscoveryFilePath(dataPath: string): string { return `${dataPath}.server.json`;}
async function readServerDiscoveryFile(dataPath: string): Promise<LocalServerDiscovery | undefined> { try { const raw = JSON.parse(await fs.readFile(serverDiscoveryFilePath(dataPath), "utf8")) as Record<string, unknown>; if (typeof raw.url === "string" && typeof raw.backendSecret === "string" && raw.url && raw.backendSecret) { return { url: raw.url, backendSecret: raw.backendSecret }; } return undefined; } catch { return undefined; }}
async function writeServerDiscoveryFile(dataPath: string, url: string, backendSecret: string): Promise<void> { const target = serverDiscoveryFilePath(dataPath); // `writeFile`'s mode only applies at creation; chmod after so a // pre-existing file (or a umask surprise) cannot leave the secret wider // than 0600. Failures are non-fatal: attach mode is an optimization, and // the owner's own client already has the in-memory credentials. await fs.writeFile(target, JSON.stringify({ url, backendSecret }), { mode: 0o600 }) .then(() => fs.chmod(target, 0o600)) .catch(() => undefined);}
async function clearServerDiscoveryFile(dataPath: string): Promise<void> { await fs.rm(serverDiscoveryFilePath(dataPath), { force: true }).catch(() => undefined);}
async function isJazzServerReachable(url: string): Promise<boolean> { try { const parsed = new URL(url); const response = await fetch(parsed, { signal: AbortSignal.timeout(500) }); // Any HTTP response (including 4xx from a non-ws GET) proves a live // listener; connection refused or timeout means the owner is gone. return response.status > 0; } catch { return false; }}
function isRecoverableStorageError(error: unknown): boolean { if (!(error instanceof Error)) return false; return error.message.includes("object already exists") || error.message.includes("database is locked");}
async function upsertProgressInTransaction( tx: TransactionScope, progress: ConsumerProgress, row?: Record<string, unknown>, storageRevision = 3,): Promise<void> { const data = { key: progress.id, consumerKey: progress.consumerId, consumerVersion: progress.consumerVersion, source: progress.source, lastSequence: progress.lastSequence, lastEventKey: progress.lastEventId, updatedAt: progress.updatedAt, }; if (row) { if (Number(row.lastSequence) > progress.lastSequence) throw new Error(`Consumer progress regression for ${progress.id}`); if (Number(row.lastSequence) === progress.lastSequence && String(row.lastEventKey) !== progress.lastEventId) { throw new Error(`Consumer progress sequence conflict for ${progress.id}`); } if (String(row.consumerKey) !== progress.consumerId || Number(row.consumerVersion) !== progress.consumerVersion || String(row.source) !== progress.source) { throw new Error(`Consumer progress identity conflict for ${progress.id}`); } tx.update(thoughtstreamApp.consumerProgress, String(row.id), data); } else { tx.insert(thoughtstreamApp.consumerProgress, data, { id: jazzRowId(`consumer-progress-v${storageRevision}`, progress.id), }); }}
/** * The graceful-drain error classes (`GracefulShutdownSyncError`, * `SharedClientShutdownError`) are `@internal` to jazz-tools and not reachable * through its public exports map, so close-time drain failures are recognized * by their stable `name` properties. A failed drain leaves the client open and * usable; the subsequent `session.close()` still performs full teardown. */function isGracefulShutdownSyncError(error: unknown): boolean { return error instanceof Error && (error.name === "GracefulShutdownSyncError" || error.name === "SharedClientShutdownError");}