Understanding Magnitude Event-Core: Event, Signal, Projection, and Worker Concepts

Magnitude's event-core implements a fully typed, effect-driven event-sourcing system built around four core concepts—Event, Signal, Projection, and Worker—that together enable deterministic state management with asynchronous side-effect handling.

The magnitudedev/magnitude repository provides a TypeScript-native framework for building event-sourced applications. Its @magnitudedev/event-core package centers on these four architectural primitives, each with distinct responsibilities in the event processing pipeline.

Core Concepts Overview

Concept Purpose Persistence Execution Model
Event Immutable state-change messages Persistent (unless ephemeral) Sequential queue processing
Signal Ephemeral cross-projection notifications Never persisted Synchronous to projections, async to workers
Projection Deterministic state holder with typed handlers Configurable per-projection Synchronous event/signal handlers
Worker Asynchronous effect runner N/A (stateless) Concurrent fiber-per-handler

Event: The Immutable State Driver

Events in Magnitude event-core are timestamped, immutable messages that drive all state changes. The EventBusCore class in src/core/event-bus-core.ts orchestrates their lifecycle.

Key characteristics:

  • Sequential processing: Events queue in a single-consumer Queue to maintain ordering guarantees
  • Persistence control: Mark events as ephemeral to skip persistence entirely
  • Broadcast to projections: The ProjectionBus receives every event for synchronous handling
// From event-bus-core.ts – the central publish path
class EventBusCore<TEvent> {
  publish(event: TEvent): Effect<void> {
    // Queue → Persist → Broadcast to ProjectionBus
  }
}

Events flow through the system in five stages: publish, queue, projection handling, signal propagation, and worker notification.

Signal: Ephemeral Cross-Cutting Notifications

Signals solve the problem of cross-projection communication without hydration replay. Unlike events, signals are never persisted—they exist only for runtime coordination.

The SignalDef type in src/signal/define.ts provides the type system:

// Creating a typed signal
const changed = Signal.create<number>('Counter/changed');

Signal delivery operates on two paths:

  • Synchronous: Directly to other projections via queueSignal on the ProjectionBus
  • Asynchronous: Via Pub/Sub (PubSub.publish) to any subscribed Workers

This dual-path design lets projections react immediately while allowing workers to process side effects concurrently.

Projection: Deterministic State Containers

Projections are the heart of Magnitude's read model. Defined via Projection.define in src/projection/define.ts, they combine state management with handler registration.

A projection configuration includes:

Projection.define<TEvent>()({
  name: 'Counter',
  state: Schema.Number,
  initial: { count: 0 },
  signals: { changed },           // signal definitions
  eventHandlers: { ... },         // react to events
  signalHandlers: on => [ ... ],  // react to signals with builder pattern
})

Forked vs. Non-Forked Projections

Non-forked projections maintain global state through ProjectionInstance. Forked projections use ForkedProjectionInstance to isolate state per forkId.

The fork ID extraction happens automatically via extractForkIdFromEvent and extractForkIdFromSignal in src/worker/util.ts. Workers receive a correctly scoped read function through makeWorkerReadFn.

Ambient and Addressed Dependencies

Projections can declare two special dependency types:

  • Ambient: Global read-only values accessed via AmbientReader—instant, non-tracked reads
  • Addressed collections: Lazily-loaded state subsections (sequences, maps) that the read tracker can pin during handler execution

Both integrate into the projection's read function with automatic dependency tracking for updates.

Worker: Asynchronous Effect Handlers

Workers handle side effects that must not block the synchronous event pipeline. Defined in src/worker/define.ts, they receive three capabilities:

Capability Type Purpose
publish (event: TEvent) => Effect<void> Emit new events—the only permitted side effect
read WorkerReadFn<TEvent> Typed, fork-aware projection state access
handler Effect The async body with built-in interrupt handling

Interrupt Coordination

Every worker handler wraps with withInterrupt from the InterruptCoordinator in src/core/interrupt-coordinator.ts:

const withInterrupt = <A, RH>(handler: Effect<Eff, RH>, targetForkId: string | null) =>
  Effect.gen(function* () {
    const baseline = yield* interruptCoordinator.current(targetForkId);
    return yield* Effect.raceFirst(
      handler,
      interruptCoordinator.waitForInterrupt(targetForkId, baseline)
    );
  });

This races the handler against an interrupt signal, enabling graceful cancellation when forks terminate or systems shut down.

Worker Subscription Model

Workers subscribe to events through the WorkerBus in src/core/worker-bus.ts, which exposes only publish and subscribe—deliberately minimal surface area for safety.

Complete Event Flow Through the System

Understanding how these four concepts interact requires tracing a single event:

  1. Publication: WorkerBus.publish(event) or EventBusCore.publish(event) initiates
  2. Queuing: EventBusCore enqueues in its sequential Queue
  3. Projection processing: ProjectionBus runs event handlers synchronously, updating state and queuing signals
  4. Signal broadcast: Signals publish to ProjectionBus (sync) and Pub/Sub (async)
  5. Worker activation: Subscribed workers receive messages in separate fibers with fork-scoped read functions

This pipeline ensures deterministic state in projections while freeing workers for concurrent side-effect execution.

Error Handling Across All Layers

The FrameworkErrorReporter in src/core/framework-error.ts wraps every execution point. Errors become FrameworkError.* variants and report without crashing the system—critical for long-running event processors.

Key integration points:

  • Event handler errors in ProjectionBus
  • Signal handler failures
  • Worker fiber crashes
  • Event bus queue processing

Practical Examples

Defining a Projection with Signals

import { Projection, Signal } from '@magnitudedev/event-core';

const changed = Signal.create<number>('Counter/changed');

const CounterProjection = Projection.define<MyEvent>()({
  name: 'Counter',
  state: Schema.Number,
  initial: { count: 0 },
  signals: { changed },
  eventHandlers: {
    increment: ({ state, emit }) => {
      const next = state.count + 1;
      emit.changed(next);
      return { count: next };
    },
  },
  signalHandlers: on => [
    on(changed, ({ value, state }) => ({
      ...state,
    })),
  ],
});

Source: [src/projection/define.ts](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/projection/define.ts)

Consuming from a Worker

import { Worker } from '@magnitudedev/event-core';

const CounterWorker = Worker.define<MyEvent>()({
  name: 'CounterWorker',
  eventHandlers: {
    increment: async (event, publish, read) => {
      const counter = read(CounterProjection);
      console.log('Current count:', counter.count);
    },
  },
  signalHandlers: on => [
    on(CounterProjection.signals.changed, (value, publish, read) => {
      console.log('Signal received:', value);
    }),
  ],
});

Source: [src/worker/define.ts](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/worker/define.ts)

Forked Projection Pattern

const ForkedCounter = Projection.defineForked<MyEvent>()({
  name: 'ForkedCounter',
  state: Schema.Number,
  initial: { total: 0 },
  eventHandlers: {
    add: ({ state, forkId }) => ({
      total: state.total + 1,
    })
  },
});

Source: [src/projection/defineForked.ts](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/projection/defineForked.ts)

Key Source Files Reference

File Responsibility
src/core/event-bus-core.ts Central event bus, sequential queue, persistence
src/core/worker-bus.ts Worker-facing publish/subscribe API
src/core/projection-bus.ts Synchronous projection processing
src/core/interrupt-coordinator.ts Fork-aware interrupt signaling
src/core/framework-error.ts Centralized error reporting
src/projection/define.ts Projection builder with handlers and read tracking
src/worker/define.ts Worker builder with effect helpers
src/signal/define.ts Signal definitions and Pub/Sub
src/worker/util.ts Fork ID extraction utilities

Summary

  • Event: Immutable, queued, optionally persisted messages that drive all state changes through EventBusCore
  • Signal: Ephemeral, typed notifications for runtime coordination between projections and workers without persistence overhead
  • Projection: Deterministic state container with synchronous handlers, supporting both global and forked (per-ID) state isolation
  • Worker: Asynchronous effect handler with controlled side effects (publish-only), fork-scoped reads, and built-in interrupt handling

The Magnitude event-core architecture separates what happened (events), what changed now (signals), what state means (projections), and what to do about it (workers) into composable, type-safe layers.

Frequently Asked Questions

What makes Magnitude's Event type different from standard event sourcing?

Magnitude's Event type integrates directly with Effect-TS for typed error handling and supports an ephemeral flag to bypass persistence entirely—useful for high-frequency signals that don't need replay. The EventBusCore in src/core/event-bus-core.ts handles both persistent and ephemeral events through the same queue mechanism.

How do Projections and Workers communicate?

Projections emit Signals—never direct method calls. These signals propagate synchronously to other projections via the ProjectionBus and asynchronously to Workers via Pub/Sub. This design prevents tight coupling and maintains clear data flow boundaries.

What is forked state and when should I use it?

Forked state creates isolated projection instances per forkId, enabling multi-tenant or workflow-isolated scenarios. Use Projection.defineForked() when you need independent state lifecycles—common in multi-user applications or long-running workflow engines where each instance must not interfere with others.

Can Workers modify projection state directly?

No. Workers have read-only access to projections through the read function from makeWorkerReadFn. State changes only happen through event publication (publish), which the ProjectionBus processes synchronously. This enforces the event-sourcing principle that all state changes flow through events.

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 →