Foundational Event-Sourcing Primitives in Magnitude's event-core Package

The event-core package provides five foundational event-sourcing primitives—EventBusCore, EventSink, ProjectionBus, InterruptCoordinator, and BaseEvent—that form a complete type-safe pipeline for publishing, persisting, and projecting events in Magnitude applications.

The event-core package serves as the architectural foundation of Magnitude's event-sourcing system. It exports a minimal, composable set of TypeScript primitives that handle everything from event publication and durable storage to synchronous projection orchestration and fork-level interrupt coordination.

Central Event Publication with EventBusCore

The EventBusCore class functions as the central nervous system of the event-sourcing pipeline. Implemented in src/core/event-bus-core.ts, it handles publishing, timestamping, sequential processing, replay-safety, and broadcasting of events.

When you publish an event, the EventBusCore performs five critical operations in sequence:

  1. Wraps the event in a Timestamped<E> envelope to add millisecond precision timestamps
  2. Queues the event in a single-consumer Queue to guarantee ordered processing
  3. Sends the event to the ProjectionBus for synchronous projection handling (phase 1)
  4. Persists the event via EventSink.append unless marked as ephemeral
  5. Broadcasts the event on a PubSub stream for external listeners
// Define a concrete event type extending the base model
interface UserCreated extends EventCore.BaseEvent {
  readonly type: 'userCreated'
  readonly userId: string
}

// Publish an event through the central bus
const bus = EventBusCoreTag<UserCreated>().pipe(Effect.service)
bus.publish({ type: 'userCreated', userId: 'U123' })

Durable Persistence via EventSink

The EventSink primitive, defined in src/core/event-sink.ts, provides durable persistence through an in-memory pending event accumulator coupled with a claim-based API for exclusive write ownership. Events remain in memory until a durable sink flushes them to permanent storage, ensuring atomicity and durability guarantees.

The claim-based API prevents write conflicts by requiring projection handlers to acquire exclusive ownership before appending events. This design ensures that only one writer can claim a specific event sequence at any given time, eliminating race conditions during concurrent event processing.

Projection Orchestration with ProjectionBus

The ProjectionBus, located in src/core/projection-bus.ts, manages synchronous projection logic through a strict two-phase execution model. This primitive tracks dependencies between projections using a topological ordering system that respects reads declarations and signal dependencies.

Phase 1 executes registered projection event handlers in dependency order. Phase 2 buffers signals emitted by projections and flushes them iteratively. The bus provides utilities for reading projection state, handling addressed state, and registering ambient handlers that operate across all events.

const projBus = ProjectionBusTag<UserCreated>().pipe(Effect.service)

projBus.register(
  // Event handler
  (e) => Effect.succeed(console.log('User created:', e.userId)),
  // Event types this projection observes
  ['userCreated'],
  // Projection name for dependency graph tracking
  'UserProjection'
)

Fork-Level Interrupt Coordination

The InterruptCoordinator primitive, implemented in src/core/interrupt-coordinator.ts, manages execution epochs for forked projection contexts. It provides a safe "interrupt" mechanism that workers can await through the waitForInterrupt method, allowing coordinated pausing and resuming of fork-specific logic.

When a forked projection needs to signal an interruption, it calls interrupt(forkId), which increments the execution epoch. Workers monitoring that fork can then detect the epoch change and respond appropriately, enabling safe cancellation and cleanup of long-running projection tasks.

const interrupter = InterruptCoordinator.pipe(Effect.service)
interrupter.interrupt('myForkId')

The Base Event Model

All events in the system extend BaseEvent, defined in src/types.ts, which establishes the minimal shape required for event-sourcing operations. The base interface requires a type discriminator string and supports an optional ephemeral boolean flag.

Events marked as ephemeral bypass the EventSink persistence layer, making them suitable for temporary signals that should not survive restarts. The Timestamped<E> generic wraps any BaseEvent to add timestamp metadata during the EventBusCore.publish operation.

Building Fork-Aware Projections

For stateful projections that maintain isolated state per fork, the defineForked API in src/fork/index.ts provides a type-safe factory for creating fork-aware projection logic. This primitive automatically manages state isolation using a Map<string | null, State> structure keyed by forkId.

import { defineForked } from '@magnitudedev/event-core'

const counter = defineForked({
  name: 'Counter',
  initialState: { forks: new Map<string | null, number>() },
  handlers: {
    increment: (state, { forkId }) => {
      const cur = state.forks.get(forkId) ?? 0
      state.forks.set(forkId, cur + 1)
    },
  },
})

The Complete Event-Sourcing Pipeline

These five primitives form a unified pipeline that processes every event through a deterministic lifecycle:


Event (BaseEvent) ──► EventBusCore.publish ──► ProjectionBus.processEvent
                      │                     │
                      └─► EventSink.append ──► Durable storage

The EventBusCore coordinates the flow, ensuring that projections process events synchronously before persistence occurs, while the InterruptCoordinator provides safety valves for long-running or forked operations. This architecture guarantees that all projections remain consistent with the durable event log without requiring distributed transactions.

Summary

  • EventBusCore (src/core/event-bus-core.ts) orchestrates event publishing, timestamping, and ordered delivery through a single-consumer queue.
  • EventSink (src/core/event-sink.ts) accumulates events in memory and provides claim-based APIs for atomic, durable persistence.
  • ProjectionBus (src/core/projection-bus.ts) executes projection handlers in dependency order across two phases (event handling → signal flushing).
  • InterruptCoordinator (src/core/interrupt-coordinator.ts) tracks execution epochs and enables safe interruption of forked projection work.
  • BaseEvent (src/types.ts) defines the minimal event contract with type discrimination and ephemeral flags.
  • defineForked (src/fork/index.ts) creates state-isolated projections that maintain separate state per fork identifier.

Frequently Asked Questions

What is the difference between EventBusCore and ProjectionBus?

EventBusCore handles the ingestion, timestamping, and broadcasting of raw events to all subscribers, while ProjectionBus specifically manages synchronous projection logic that transforms events into read models. According to the source code in src/core/projection-bus.ts, the projection bus executes handlers in topological order and manages signal flushing, whereas the event bus in src/core/event-bus-core.ts focuses on queue management and persistence coordination.

How does EventSink ensure events are not lost during crashes?

EventSink maintains events in an in-memory accumulator until explicitly flushed to durable storage, using a claim-based API that requires exclusive write ownership. As implemented in src/core/event-sink.ts, this design ensures that events remain buffered until the underlying storage layer confirms the write, providing durability guarantees even if individual projection handlers fail during processing.

What are ephemeral events and when should I use them?

Ephemeral events are events flagged with ephemeral: true in the BaseEvent interface that bypass the EventSink persistence layer entirely. According to the type definitions in src/types.ts, these events are processed by projections but never stored durably, making them suitable for temporary signals, heartbeats, or transient state notifications that should not survive application restarts.

How does InterruptCoordinator handle concurrent fork execution?

InterruptCoordinator manages concurrent forks through execution epochs tracked per forkId in src/core/interrupt-coordinator.ts. When interrupt(forkId) is called, it increments the epoch for that specific fork; workers awaiting interruption via waitForInterrupt detect this epoch change and safely terminate or pause their current operation. This mechanism prevents race conditions during fork-specific cleanup or cancellation.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →