# How the event-core Event Sourcing System Works in Magnitude

> Discover how magnitudedev/magnitude's event-core event sourcing system achieves synchronous, consistent state updates through its three-layer Effect-TS architecture before broadcasting events.

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

---

**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`](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/core/event-bus-core.ts)) – Receives events via the public `publish()` 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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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:

1. **Publication** – The `publish(event)` method in `EventBusCore` (lines 55‑66) assigns a monotonic timestamp (or respects an existing one) and enqueues an `eventQueueItem`.
2. **Consumer Fiber** – A dedicated fiber created with `Effect.forkScoped` continuously pulls items from the queue via `Queue.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 `EventSink` unless `event.ephemeral` is true (lines 29‑35).
   - Broadcasting to a `PubSub` for 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`](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/projection/define.ts). Each projection exposes:

- **State** – A `SubscriptionRef` readable synchronously via `ProjectionBus.read` locks.
- **Signals** – Lightweight, ephemeral streams built on Effect-TS `PubSub`. Signals originate from a `SignalDef` defined in [`packages/event-core/src/signal/define.ts`](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/signal/define.ts) and attach to source projections. Emitting a signal invokes the `emit` function, which the framework wraps in an `Effect`.

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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/worker/util.ts), allowing the coordinator to target specific fibers for termination without affecting global state.

## Practical Implementation Examples

### Defining Events and Projections

```typescript
// 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

```typescript
// 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

```typescript
// 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

```typescript
// 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

```typescript
// 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.