How the event-core Event Sourcing System Works in Magnitude
The event-core event sourcing system in magnitudedev/magnitude is an in-process engine built on Effect-TS that processes events through a three-layer architecture—EventBusCore, ProjectionBus, and EventSink—to guarantee synchronous, consistent state updates before asynchronous broadcasting.
This system drives the entire Magnitude platform by handling event ingestion, state projection, and persistence within a single process. Unlike distributed event sourcing systems that rely on external message brokers, event-core keeps everything in-memory and synchronous for the critical path, using durable storage only as a backing sink.
Architecture Overview: The Three Layers
The engine coordinates three specialized components to maintain data consistency:
- EventBusCore (
packages/event-core/src/core/event-bus-core.ts) – Receives events via the publicpublish()method, timestamps them, and manages the consumption queue. It forwards events to the projection system, persists them (unless marked ephemeral), and broadcasts them to workers via PubSub. - ProjectionBus (
packages/event-core/src/core/projection-bus.ts) – Executes synchronous projection handlers in strict dependency order, buffering signals and flushing them iteratively according to the dependency graph. - EventSink (
packages/event-core/src/core/event-sink.ts) – Acts as an in-memory accumulator that owns the pending-event claim consumed by the daemon to flush events to durable storage.
Event Flow: From Publication to Broadcast
When an event enters the system, it undergoes a rigorous, ordered pipeline:
- Publication – The
publish(event)method inEventBusCore(lines 55‑66) assigns a monotonic timestamp (or respects an existing one) and enqueues aneventQueueItem. - Consumer Fiber – A dedicated fiber created with
Effect.forkScopedcontinuously pulls items from the queue viaQueue.take. For standard events, it executes:ProjectionBus.processEvent(event)(lines 19‑22) to run the first phase of projection handlers.- A check against
hydration.isHydrating()(lines 21‑23) to skip persistence during replay scenarios. - Special handling for "interrupt" events that abort forked projections via the
InterruptCoordinator(lines 25‑27). - Conditional persistence to the
EventSinkunlessevent.ephemeralis true (lines 29‑35). - Broadcasting to a
PubSubfor asynchronous worker consumption (lines 36‑40).
If the queued item is a checkpoint, the supplied effect runs between two events, ensuring all projection state is consistent before the next event processes (lines 7‑14).
Two-Phase Projection Processing
The ProjectionBus implements deterministic state updates through distinct phases to prevent partial updates:
Phase 1: Event Handler Execution
The processEvent method runs every registered handler in dependency order, derived from each projection's reads and signal subscriptions. Handlers may queue signals via queueSignal, but signals remain buffered and are not emitted immediately.
Phase 2: Signal Flushing
After all handlers complete, the bus enters an iterative flush cycle (up to MAX_SIGNAL_FLUSH_ITERATIONS to prevent infinite loops). Buffered signals deliver to registered signal handlers, which can emit additional signals, creating a deterministic cascade that completes before publish() returns.
Because both phases execute synchronously before the EventBusCore broadcast occurs, downstream consumers always observe a fully consistent view of projection state.
State Management: Projections and Signals
Projections are declared using Projection.define() in packages/event-core/src/projection/define.ts. Each projection exposes:
- State – A
SubscriptionRefreadable synchronously viaProjectionBus.readlocks. - Signals – Lightweight, ephemeral streams built on Effect-TS
PubSub. Signals originate from aSignalDefdefined inpackages/event-core/src/signal/define.tsand attach to source projections. Emitting a signal invokes theemitfunction, which the framework wraps in anEffect.
The addressed subsystem adds a proxy layer that tracks which fields a projection reads, enabling fine-grained dependency tracking across projections without manual declaration.
Forked Projections and Ambient State
Forked Projections allow isolated state slices (e.g., per-window UI state). Their values store in a Map<string | null, unknown> (see ForkedProjectionState in projection-bus.ts). The same two-phase processing applies, but reads route to the appropriate map entry based on fork ID.
Ambient State represents global values (e.g., configuration) observable by any projection. Changes process via processAmbientChangeTransaction (lines 59‑71), which respects the projection-dependency graph to ensure ordered updates.
Hydration and Replay Safety
During rehydration (replaying persisted events after process restart), the HydrationContext signals active replay mode. The EventBusCore checks hydration.isHydrating() (line 21) and skips persisting events to the EventSink, preventing duplicate writes. Ephemeral events bypass the sink entirely, ensuring they exist only during the live process lifetime and never in the event log.
Interrupt Coordination
The InterruptCoordinator (invoked lines 25‑27) enables the daemon to abort long-running forked operations (e.g., user-cancelled tasks). The reserved "interrupt" event type carries a fork ID extracted via extractForkIdFromEvent from worker/util.ts, allowing the coordinator to target specific fibers for termination without affecting global state.
Practical Implementation Examples
Defining Events and Projections
// events.ts
export interface IncrementEvent extends BaseEvent {
readonly type: 'increment'
readonly amount: number
}
// projection.ts
import { Projection } from '@magnitudedev/event-core'
export const Counter = Projection.define({
name: 'Counter',
state: Schema.Number,
events: {
increment: (event) => Effect.succeed((state) => state + event.amount)
}
})
Publishing Events
// usage.ts
import { EventBusCoreTag } from '@magnitudedev/event-core'
import { IncrementEvent } from './events'
const bus = Effect.serviceTag(EventBusCoreTag<IncrementEvent>())
Effect.runPromise(
bus.publish({ type: 'increment', amount: 5 })
)
Subscribing to Event Streams
// worker.ts
const sub = Effect.serviceTag(EventBusCoreTag<IncrementEvent>()).subscribe()
Effect.runPromise(
sub.pipe(
Stream.map(e => console.log('got event', e)),
Stream.runDrain
)
)
Using Checkpoints for Consistency
// checkpoint.ts
const bus = Effect.serviceTag(EventBusCoreTag<IncrementEvent>())
Effect.runPromise(
bus.checkpoint(
Effect.succeed(console.log('All prior events processed, projection state is up‑to‑date'))
)
)
Emitting and Consuming Signals
// signal.ts
import { create } from '@magnitudedev/event-core/signal'
export const CountReached = create<number>('CountReached')
// In projection handler:
emit(CountReached, newCount)
// In a worker:
Effect.runPromise(
stream(CountReached).pipe(
Stream.map(n => console.log('counter reached', n)),
Stream.runDrain
)
)
Summary
- The event-core system uses Effect-TS to compose an in-process event sourcing pipeline that is testable and declarative.
- EventBusCore manages the event lifecycle from ingestion through broadcasting, while ProjectionBus guarantees synchronous consistency via two-phase processing.
- All projection handlers execute in dependency order before signals flush iteratively, ensuring no partial state updates reach subscribers.
- Ephemeral events bypass persistence entirely, while the HydrationContext prevents duplicate writes during event replay.
- Forked projections and ambient state provide scoped and global state management respectively, both respecting the dependency graph.
- InterruptCoordinator allows safe cancellation of specific forked operations without destabilizing the global event bus.
Frequently Asked Questions
How does event-core guarantee consistency between projections and workers?
The ProjectionBus processes all event handlers synchronously in Phase 1 before flushing signals in Phase 2. The EventBusCore only broadcasts events via PubSub after the projection bus completes both phases, ensuring workers receive events only after all projection state is stable.
What is the difference between events and signals in event-core?
Events are durable (unless marked ephemeral) facts that travel through the EventBusCore and may be persisted to the EventSink. Signals are ephemeral, in-memory notifications defined via SignalDef that propagate through the ProjectionBus signal queue and are never persisted; they exist only to coordinate state changes between projections.
How does the system prevent duplicate events during rehydration?
During rehydration, the HydrationContext sets an internal flag. When EventBusCore checks hydration.isHydrating() (line 21), it skips the sink.append(event) call, preventing the EventSink from writing replayed events to durable storage again.
Can long-running projections be safely cancelled?
Yes. The InterruptCoordinator listens for special "interrupt" events containing a fork ID. When received, it aborts the specific fiber associated with that fork ID via the ProjectionBus, allowing the daemon to cancel individual operations without restarting the entire event sourcing engine.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →