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

> Discover event-core's five foundational event-sourcing primitives for a type-safe pipeline in Magnitude apps. Publish, persist, and project events efficiently.

- Repository: [Magnitude/magnitude](https://github.com/magnitudedev/magnitude)
- Tags: internals
- Published: 2026-09-05

---

**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`](https://github.com/magnitudedev/magnitude/blob/main/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

```typescript
// 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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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.

```typescript
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`](https://github.com/magnitudedev/magnitude/blob/main/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.

```typescript
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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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`.

```typescript
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`](https://github.com/magnitudedev/magnitude/blob/main/src/core/event-bus-core.ts)) orchestrates event publishing, timestamping, and ordered delivery through a single-consumer queue.
- **EventSink** ([`src/core/event-sink.ts`](https://github.com/magnitudedev/magnitude/blob/main/src/core/event-sink.ts)) accumulates events in memory and provides claim-based APIs for atomic, durable persistence.
- **ProjectionBus** ([`src/core/projection-bus.ts`](https://github.com/magnitudedev/magnitude/blob/main/src/core/projection-bus.ts)) executes projection handlers in dependency order across two phases (event handling → signal flushing).
- **InterruptCoordinator** ([`src/core/interrupt-coordinator.ts`](https://github.com/magnitudedev/magnitude/blob/main/src/core/interrupt-coordinator.ts)) tracks execution epochs and enables safe interruption of forked projection work.
- **BaseEvent** ([`src/types.ts`](https://github.com/magnitudedev/magnitude/blob/main/src/types.ts)) defines the minimal event contract with type discrimination and ephemeral flags.
- **defineForked** ([`src/fork/index.ts`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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.