← Workstreams

Workstream: EventSourcingStreams

Status: Planned · Component: Maximize developer productivity

Goal

Make event distribution a solved primitive. Any party interested in what is happening — as it happens, or catching up after being away — subscribes to a Stream: a decentralized, subscribable event log with replay. A service publishes its domain events once; every interested consumer receives them live, replays what it missed, keeps its own copy under its own retention policy, derives filtered streams, or joins streams together — with no broker to operate and no per-service event plumbing (webhooks, pollers, bespoke fan-out) to invent.

Mental model: the event plane and the state plane

The ecosystem separates events in motion from state at rest, and EventSourcingStreams is the event half:

  • Streams are the event plane. A consumer on the event plane wants the events themselves — every one, in causal order, with full detail: an auditor, a trigger, a service feeding its own read model. Such a consumer typically never reads anyone's projection.
  • Observables are the state plane. A consumer on the state plane wants an up-to-date projection and never needs to see the underlying events — though it frequently writes to the event plane (its commands append events) without ever reading it.
  • Streams derive from streams — and subscribing is deriving. A subscription is itself a stream: by default a subscriber materializes its own copy of the upstream stream, on which it sets its own retention policy and registers its own projections. Filtered and joined streams are the same mechanism with a transform applied; the plain copy is the identity derivation, and a purely live, fire-and-forget consumer is a copy with an empty retention window — one mechanism, no special-cased "lightweight" subscription.
  • A projection is a fold of a stream — and deliberately a lossy abstraction over it, never a substitute for it. The two planes serve disjoint reader audiences and stay distinct: streams carry the changelog, projections carry the current state, and the remote-observables SUBSCRIBE → SNAPSHOT → DELTA sync (state deltas, conflatable, snapshot-not-history for late subscribers) is the state plane's counterpart to a stream's subscribe-and-replay (domain events, every one delivered, replay for late subscribers). The fold's own definition is versioned and event-sourced (see the committed shape), so the pair (stream events, definition events) deterministically reconstructs a projection at any point in time.
  • The Event Log is the same contract at rest. One core log contract — append, cursor-based replay, subscribe — with two implementations: the Event Log is the durable, single-node, retention-forever implementation (the system of record, for which "replay from the beginning" is simply the retention-forever special case), and EventSourcingStreams is the decentralized implementation of the same contract.

The committed shape

  • One contract, two implementations. The core contract lives with the existing EventLog interface, extended with a push-subscription surface alongside today's cursor replay (streamFrom / fetchBatch). The Event Log implements it as the durable local log; EventSourcingStreams implements it as a decentralized stream. Code written against the contract moves between them freely — a test runs on the in-memory log, production runs on a stream.
  • Embedded-first, no broker. Each member embeds a stream node — an embedded HierarchicalClock plus a local log buffer (which can itself be the local Event Log implementation) — and members peer over url:// via UrlResolver. A hosted, always-on member is just another member acting as a durability and retention anchor, per the deployment spectrum and the decentralized-by-default philosophy.
  • Subscribing is deriving your own stream. The default subscription materializes the subscriber's own copy of the stream — the identity case of a derived stream — with its own retention policy, independent of the origin's, and its own registered projections. This is retention sovereignty: a consumer that needs history keeps it itself instead of petitioning the origin for a longer window, and the origin's retention only bounds how far back a newly created copy can backfill. A purely live consumer is the degenerate copy with an empty retention window, not a separate subscription mechanism.
  • Retention is per-stream policy, from a bounded window to forever. A stream — original or copy — may keep a length- or age-bounded window of events, or keep everything. Bounded retention is the common case for high-volume streams; retention-forever makes a stream a distributed system of record.
  • Truncation is gated on projections, locally. A registered projection is always updated before the events it folds are truncated — a stream never drops an event one of its projections has not yet consumed. Because subscription-by-copy makes projections local to the copy they are registered on, each stream enforces this independently; there is no global registry of downstream projections to coordinate. Expired history is thereby consolidated into state rather than silently lost; consumers who need the full event detail keep a copy whose retention covers their need.
  • Derived streams: copy, filter, and join. Filtering and joining are event → event operators whose output is another stream, not state — and the plain copy a subscription creates is their identity case. Every derived stream is a full stream: its own retention policy, its own projections, its own subscribers. Events keep their identity and HierarchicalTimestamp across derivation, so cross-member and cross-stream ordering is causal — never wall clocks — and a freshness requirement satisfied by the origin is satisfiable by any sufficiently caught-up copy.
  • Projections are Observables — one projection system. Every projection, whether over the Event Log or over a stream, is an Observable projection, per the CQRS mandate. EventSourcingStreams ships a fold-a-stream-into-a-projection bridge, not a second projection framework; serving that state remotely is the Observables workstream's remote-observable path.
  • Schema evolution happens in the fold — and the fold's evolution is event-sourced. Events, once published, are immutable and are never rewritten or migrated; as event types and schemas evolve, it is the projection definition that changes. Installing a new projection version is itself recorded as an event in an accompanying, retention-forever definition stream (itself just a stream; definition events are tiny and must outlive the data stream's retention, or determinism is lost to truncation). Each installation record references the exact projection code version installed and the explicit event in the data stream as of which it takes effect — never an implied "from now on". Given the retained events and the definition stream, a projection's state at any point in time is deterministically reconstructible, replaying each historical version over the segment it governed; when the data stream has truncated, a new version builds on the prior projection's state at the recorded takeover event — which is exactly why the effective event is explicit. An installation may be backdated: the projection rewinds to a checkpoint of its own state at or before the effective event, replays up to it under the old code, then switches to the new code and continues — so a bad or not-yet-understood event type is retroactively repaired by a backdated projection version, reaching as far back as retention and prior checkpoints allow.
  • Clocks are structural. Every published event carries a HierarchicalTimestamp minted by the publishing member's embedded clock; catch-up is naturally expressed as a ClockPad ("everything I have not yet witnessed"); and a publish returns its timestamp, so a publisher can require any downstream projection to prove it reflects the publish — read-your-writes across the two planes with zero per-service plumbing. This is also the Event Log's multi-writer ordering answer: causal ordering via entangled clocks, with no stronger total-order mechanism imposed.

Current state

EventSourcingStreams does not exist yet; every building block does:

Why it accelerates developers

  • Event distribution stops being invented per service. Publishing to a stream replaces bespoke webhooks, pollers, and fan-out; subscribing replaces "how do I find out when X happens?".
  • Catch-up is a primitive. A consumer that was down replays from its cursor (or its pad) instead of reconciling state by hand.
  • Retention stops being a negotiation. A consumer that needs more history than the origin keeps simply keeps it itself — its subscription copy carries its own retention policy — instead of asking the stream's owner to widen theirs.
  • One contract from local log to global stream. The same code runs against the in-memory log in a test, the durable Event Log in a single-node deployment, and a decentralized stream in production.
  • It composes the rest of component 1. HierarchicalClock orders the events and proves read-your-writes; Observables serve the folded state; the Event Log remains the durable record; W3Wallet gates who may publish to and subscribe from a stream.

Plan / roadmap

  • [ ] Shared core contract. Extend the EventLog contract with the push-subscription surface and a per-stream retention policy, keeping the Event Log its retention-forever implementation.
  • [ ] Embedded stream node. A member = embedded clock + local log buffer + url:// peering; prove publish, live subscribe, and cursor replay end to end between two in-process members, per the testing standards.
  • [ ] Copy-by-default subscription. Subscriber-side materialization of a stream copy with an independent retention policy, backfilled from the origin within the origin's retention — the empty-retention copy doubling as the live tap.
  • [ ] Retention and projection-gated truncation. Per-stream retention policy plus the invariant that a stream's registered projections update before it truncates, enforced locally by every copy.
  • [ ] Derived streams. Filtered streams first, then clock-ordered joins — both building on the copy mechanism.
  • [ ] Projection bridge. The fold-a-stream-into-an-Observable-projection adapter — shared work with the Event Log projection toolkit.
  • [ ] Versioned projection definitions. The retention-forever definition stream, effective-as-of installation records, deterministic point-in-time reconstruction across versions, and backdated installations with checkpoint rewind — proving a retroactive repair of a mis-handled event type end to end.
  • [ ] Read-your-writes across the planes. A publish returns its timestamp; a downstream projection read presents the pad and proves incorporation, per HierarchicalClock.
  • [ ] Hosted anchor member and reference adoption. One durable hosted member anchoring a production stream, with one real pair of services exchanging events through it; document the before/after.

Graduation

When two independent services exchange domain events through a stream in production — live subscription plus catch-up after a restart — and a registered projection demonstrably consolidates a truncating stream, this graduates to a first-class project: document it in the Documentation Repository and add it to ALL_PROJECTS.md.