A local-first event pipeline for independent agents, built on Jazz.
19 kB

Jazz-first persistence and coordination #

Governing principle #

For every persistence, query, subscription, replication, durability, and authorization requirement, ask first: can Jazz own this directly?

If yes, thought stream uses Jazz instead of recreating the facility in TypeScript. Application code defines domain meaning, adapts external systems, invokes models, and supervises processes. It does not become a second database.

The first implementation pinned jazz-tools and jazz-napi to 2.0.0-alpha.53; the current port pins 2.0.0-alpha.55. Earlier CoFeed sketches and generic storage abstractions are not binding. The installed Jazz 2 relational API is the source of truth.

Actual process topology #

thought stream is a shared Jazz database with independently runnable producers and consumers.

  • A producer process owns one stable source id, one external-source cursor, and one source-local event sequence. There is one active writer for that source id. An email producer and a Jetstream producer may write concurrently, but they occupy different source namespaces and are not competing to ingest the same email or post.
  • A consumer process owns one stable consumer id and version, one declarative Jazz query/subscription, and its own progress through matching source sequences. There is one active process for that consumer identity. The consumer is the matcher and runner; there is no central matcher, global dispatch queue, or lease-owning scheduler between it and Jazz.
  • Producers and consumers can be launched independently on different computers. Jazz supplies the common synchronized state and live query surface.
  • A consumer may emit derived or lifecycle events. It then acts as a producer under its own source namespace, so other consumers can observe its work through the same mechanism.
  • An ad hoc reader can query or subscribe without registering durable progress. Managed consumers create progress rows only for source namespaces they actually consume.

Running redundant copies under the same producer or consumer identity is outside the initial topology. If redundancy is added later, it needs an explicit ownership protocol. The initial design should not acquire a distributed queue and lease system merely to defend against a deployment shape it does not use.

Verified surface of the pinned release #

The installed 2.0.0-alpha.55 packages expose:

  • typed relational schemas, references, selected indexes, migrations, and permission policies;
  • typed database filtering, ordering, inclusion, limits, and offsets;
  • one-shot reads and live query subscriptions with complete results per update (the alpha.53 row-delta subscribeAll API is gone; consumers diff full snapshots by row id);
  • insert with optional caller-supplied row id, caller-id upsert, update, soft delete, and restore;
  • direct batches and authority-validated transactions (failed transactions roll back atomically in alpha.55, removing the alpha.53 allocated-row-id residue the storage-revision retry worked around);
  • WriteHandle.wait({ tier }) at local, edge, and global durability tiers;
  • a supported local server harness (startLocalJazzServer from jazz-tools/dev) and self-host server (JazzServer.start from jazz-napi), with backend admission via createJazzSession from jazz-tools/backend;
  • mutation rejection reporting and row-level permission evaluation.

A local probe verified that insert(..., { id }) accepts a caller-supplied UUID and rejects both identical and conflicting second inserts with object already exists, preserving the first row. That is enough for create-once identity in one local database. Synchronized behavior still receives integration coverage, but it is not a reason to redesign the ordinary single-owner producer topology.

Responsibility matrix #

Concern Jazz owns producer or consumer owns
Schema and row shape columns, types, references, indexes, migrations event registry, domain names, payload validation
Row identity globally unique row id and caller-id insert enforcement deterministic derivation from source/domain identity
Historical records persistence, synchronization, durability, permission enforcement deciding which records are historical and how corrections are represented
Queries filtering, ordering, pagination, inclusion, execution declarative source/consumer predicate and replay policy
Live changes query subscriptions and row deltas what to do when a matching row arrives
Durability local/edge/global acknowledgement which tier must settle before source or consumer progress advances
Replication local persistence, sync transport, merge machinery deployment topology and privacy posture
Permissions row-level policy evaluation and mutation rejection actor roles, privacy classes, source namespace authority
Multi-write settlement batches and transactions choosing the domain atomic unit
Producer progress durable source and cursor rows external cursor interpretation and source-local sequence assignment
Consumer progress durable per-consumer/per-source progress rows subscription version, batching, retry, and replay policy
Process lifecycle nothing sockets, watchers, model calls, timers, cancellation, backpressure, restart

Data classes #

Historical authority #

  • events: source observations, consumer-derived observations, and lifecycle evidence.
  • documentVersions: exact observed versions required for replay and diffs.
  • traceChunks: immutable raw or redacted runner/provider stream records.

Historical rows are caller-id inserts. The deployed Jazz permission policy must deny update, delete, and restore; the current local-development policy does not yet provide that guarantee. Corrections append new events that reference prior rows.

documentVersions also hold immutable bounded public ATProto object-context snapshots used by the resident retry protocol. These versions are keyed by declaration, event, and source-specific strong-reference identity; they contain only the public source packet and fetched public Markdown views, and survive projection rebuilds. Private Telegram conversation context is never stored through this snapshot path. Snapshot content is live runtime data and is excluded from Git, traces, accounting, notifications, and public projections.

Operational state #

  • sources: producer declarations, source-local sequence state, and nonsecret configuration.
  • sourceCursors: the last durable external position for each producer.
  • consumers: compiled declarative consumer specifications.
  • Adapter inventory is compiled from the immutable startup catalog and persisted on consumer/run evidence. Jazz has no mutable model-adapter registry or lifecycle authority.
  • consumerProgress: a managed consumer's durable position within each source namespace it has consumed.
  • executions: optional materialized attempt/status evidence for managed model runs; this is not a global work queue.
  • inferenceBudgetAccounts: per-scope rolling/fixed-window counters and active reservation leases.
  • inferenceAccounting: one privacy-dark reservation/settlement record per model attempt, containing only execution identity, model identity, timestamps, status, and numeric estimated/actual charges for policy-tracked dimensions. Calls and tokens are mandatory; untracked cost is omitted from estimate and charged JSON rather than serialized as zero.
  • projections: rebuildable named views and their durable progress.

Operational rows are mutable. Adapter release identity and deployment catalog identity are content-addressed outside Jazz, then copied into consumer/run evidence for historical inspection. Activation and retirement require a new catalog generation plus coordinated restart; they are not database transitions. Accounting JSON readers accept both existing cost-bearing rows and new rows with omitted cost. Settlement follows the estimate captured on each reservation, so an active legacy reservation remains cost-tracked even after a declaration adopts a token-only policy.

Stateful conversation bindings #

A file-scoped Letta conversation binding is keyed by declaration id, agent id, source id, and stable document id. Declaration version and current path are mutable operational metadata, so prompt/contract upgrades and confident renames preserve the same conversation. Every binding creation or meaningful metadata change appends sensitive stream.thought.runtime.letta-conversation.binding evidence. A rebuildable projection stores the latest binding for fast lookup; no new Jazz table or parallel database is introduced.

Conversation creation is recoverable across the remote-effect gap. Before creating a conversation, the trusted parent searches the exact agent for the deterministic summary marker. Zero matches permits one create, one match is reused, and multiple matches fail closed. Jazz then appends or verifies the binding event and materializes its projection. A missing or stale projection is rebuilt from immutable binding events. A conflicting agent, document, marker, or conversation id is an integrity failure. A path-only rename appends a new binding version without changing conversation identity.

Event identity and source sequence #

Every producer derives a deterministic UUID-compatible Jazz row id from its stable source id and source-native idempotency key. Examples include:

  • email: account and immutable message/change identity;
  • Jetstream: DID, collection, record key, revision, and operation;
  • filesystem: source root, document identity, content version, and operation;
  • consumer output: consumer id/version, execution attempt, output slot, and output type;
  • lifecycle evidence: producer/consumer id, execution identity, attempt, and transition.

Historical rows use caller-id insert, never upsert. A duplicate-id rejection triggers an exact read and comparison. Identical canonical content means replay found an already accepted observation; different content is an integrity failure.

Each producer also assigns a monotonically increasing sourceSequence within its source namespace. This is not a global database sequence and does not pretend independent sources have one total causal order. It gives consumers a durable, queryable boundary per source. The producer writes new event rows, its next source sequence, and its external cursor as one transaction before acknowledging source progress.

Producer append transactions are serialized per canonical source within one process. This prevents overlapping lifecycle or connector callbacks from reading the same source head and attempting conflicting source-sequence or caller-id writes. Independent sources remain concurrent. Cross-process writers must still have disjoint source ownership under the initial topology.

Consumer query and replay #

A consumer declaration compiles directly to a Jazz event query. Predicates may include event type, source, privacy class, address, and typed indexed payload fields. Jazz performs filtering, source-sequence ordering, bounds, and pagination. Consumers do not load the entire events table and recreate a query engine in TypeScript.

For each source namespace covered by the subscription, a managed consumer stores its last terminally handled sourceSequence. Startup is:

  1. read the consumer specification and existing per-source progress;
  2. query each included source for matching events after that source's progress, ordered by sourceSequence;
  3. process the bounded backlog;
  4. follow the same Jazz query as a live subscription;
  5. deduplicate any query/subscription overlap by deterministic execution identity.

New source declarations are themselves subscribable, so a wildcard consumer can begin tracking a source that did not exist at its previous startup. Changing a consumer's predicate or semantic version requires an explicit replay policy: continue, reset selected source progress, or start a new consumer version.

The shared UI may interleave events by observation time and id for display. Correct recovery depends on source-local sequence, not on a fictional total order across independent writers.

Atomic units #

The intended database atomic units are:

  1. producer events plus that producer's external cursor and source sequence;
  2. consumer output/lifecycle events plus terminal execution evidence and progress for the consumed source sequences;
  3. one projection update plus its per-source progress;
  4. one inference reservation plus its budget-account charge, and later one reservation settlement plus the corresponding charge adjustment. Budget-window arithmetic treats absent optional cost as zero internally, while persisted reservation and charge records preserve omission.

Adapter startup does not add another Jazz atomic unit. Each process loads one immutable release/deployment bundle before registration, refuses startup on identity or process-set mismatch, and records the resulting public catalog/release/binding identity on declarations and runs. Catalog change uses stop/install/restart, followed by PID/start-time, loaded-digest, and canary verification. Jazz is evidence for what ran, not consensus for what may run.

Model calls, sockets, and filesystem reads occur outside a transaction. The process records a started attempt, performs external work, then commits accepted outputs, terminal evidence, and consumer progress together. If the process dies during external work, the started attempt remains nonterminal; restart records it as status-unknown or abandoned according to policy and may begin a new attempt. No lease is required because one process owns the consumer identity.

Use a Jazz transaction when the pinned runtime proves that all writes in one unit settle together at the required authority. Use a direct batch only when grouped visibility, rather than authority validation, is the actual requirement. These semantics require integration tests, but not compare-and-swap or queue machinery under the single-owner topology.

Jazz commits each inference account/reservation update together, but the pinned relational transaction API has not proved serializable cross-process read-modify-write callbacks. A process-global queue serializes reservations per account across every JazzThoughtStore in that process, matching the one-active-process-per-consumer topology. The cap is strict per process and best-effort in aggregate across independent processes: concurrent or stale processes may briefly overshoot before synchronized state catches up. This bounded spillover is accepted for the current low-cost inference circuit breaker. A shared table alone is not a distributed lock; deployments that later require globally strict billing enforcement need a proven Jazz atomic primitive or an explicit external coordinator.

Durability and synchronization #

  • Default database: .thoughtstream/state/jazz.sqlite (a RocksDB directory owned by the Jazz server process under jazz-tools/jazz-napi 2.0.0-alpha.55).
  • Default app id: thoughtstream-local.
  • Storage-owner topology (alpha.55): exactly one process opens the on-disk database — the Jazz server. Without JAZZ_SERVER_URL, the store process hosts that server in-process through the supported startLocalJazzServer harness (jazz-tools/dev) and connects its own client over the loopback sync transport with backend admission; the client uses the memory driver and never touches the server's storage. With JAZZ_SERVER_URL, the remote server is the storage owner. There is no supported path where two processes open the same database file.
  • Durability mapping (alpha.55): the client never owns persistent storage, so a local-tier wait would only settle its in-memory cache. Every required write and read waits at edge durability — acknowledged as persisted by the storage-owning server. This is the alpha.55 equivalent of the old local-file durability guarantee; the previously separate no-server (local) and server (edge) cases now share one tier. global durability remains optional policy for selected evidence, not an assumed default.

Jazz owns persistence and synchronization. thought stream does not add a second outbox, replica log, or bespoke sync protocol unless a concrete external integration requires one.

Permissions #

Synchronization and authorization are separate Jazz facilities. Private data is not protected merely because the server is self-hosted or the UI is local.

The deployed policy must:

  • allow each producer or consumer runtime to insert only within its authorized source namespace;
  • deny update/delete/restore for events, document versions, and trace chunks;
  • scope operational-state mutation to the process identity that owns it;
  • prevent clients from reading privacy classes outside their session policy;
  • expose only explicitly public-source data to any future public role;
  • treat mutation rejection as a first-class failed write, including rejection surfaced after restart.

The permissive local-development policy is not deployable. A remote server may not be configured until row-level read/write/delete tests pass for producer, consumer, owner, unauthorized, and anonymous sessions.

Application boundary #

There is no backend-agnostic ThoughtStore. A generic interface that exposes ad hoc list* and upsert* methods while hiding Jazz queries, subscriptions, transactions, and write handles recreates database behavior in application code and preserves portability to no real second backend.

Thin Jazz-shaped modules are appropriate for:

  • event registry validation and canonical serialization;
  • deterministic identity and source-sequence derivation;
  • producer transaction construction;
  • consumer query compilation and terminal transaction construction;
  • policy-aware row conversion.

They do not implement a parallel query engine, durability abstraction, transaction model, central dispatch queue, lease service, or generic repository contract. Tests use temporary real Jazz databases rather than in-memory storage doubles.

Replay and projections #

Projection rows are disposable Jazz materialized views. A rebuild clears only projection state, queries historical events by source sequence, and verifies projection hashes. Ordinary projection replay does not regenerate model outputs. Deliberate model replay starts a new consumer version or explicit replay attempt linked to the same historical inputs.