Understanding Projections and Workers in Magnitude's Event Sourcing Model

Projections maintain immutable read-only state derived from event streams using Effect-Schemas, while workers handle asynchronous side-effects and external IO, with both supporting forked execution for isolated processing paths.

Magnitude's core runtime in the magnitudedev/magnitude repository implements an event sourcing model where all state changes are captured as immutable events. The framework separates read operations from side-effect execution through two primary abstractions: projections and workers. Understanding how these components interact within the source code is essential for building deterministic, replayable systems.

What Are Projections in Magnitude?

A projection derives read-only state from event streams, defining the shape of that state via an Effect-Schema. Defined in packages/event-core/src/projection/define.ts, projections declare an initial state and a set of event handlers that update that state.

When an event is published to the ProjectionBus (packages/event-core/src/core/projection-bus.ts), each registered projection receives the event in order. The projection's handler updates an internal immutable snapshot, exposing the latest state via ProjectionInstance.get() or ProjectionInstance.getAllForks().

Projections are pure functions of the event log—they contain no side effects. This makes them ideal for maintaining queryable views of system state that can be rebuilt deterministically from the event history.

What Are Workers in Magnitude?

Workers execute side-effects—such as HTTP requests, database writes, or long-running tasks—in response to events. Defined in packages/event-core/src/worker/define.ts using defineWorker, these components run in the background and may emit additional events.

When events are published, the WorkerBus (packages/event-core/src/core/worker-bus.ts) dispatches them to subscribed workers. Unlike projections, workers may:

  • Call external APIs asynchronously
  • Publish new events that other projections or workers consume
  • Run indefinitely until a completeOn event signals termination

Workers bridge the gap between the pure event-sourcing core and external systems.

Forking and Isolated Execution

Both projections and workers support forking, enabling multiple isolated views of the same event stream. This is implemented in packages/event-core/src/projection/defineForked.ts and packages/event-core/src/worker/defineForked.ts.

When you fork a projection or worker, you attach a forkId to events. This allows the engine to maintain parallel execution paths—useful for speculative execution, retries, or multi-agent coordination.

Forked projections expose getFork(forkId) to retrieve state for a specific fork. Forked workers run with isolated scopes and can be torn down independently when their completeOn condition is met.

Snapshotting and Persistence

Projections support serialization through writeProjectionSnapshot and restoration via prepareProjectionSnapshotRestore, implemented in packages/storage/src/sessions/storage.ts.

Snapshots capture the internal state of a projection at a specific point in time, enabling:

  • Fast session restoration without replaying entire event histories
  • Debugging and inspection of intermediate states
  • Fault tolerance through state persistence

Workers do not maintain snapshots—they handle side-effects ephemerally while relying on projections for durable state.

Code Examples

Defining a Counter Projection

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

export const Counter = define<IncrementEvent>()({
  name: 'Counter',
  init: () => ({ count: 0 }),
  apply: {
    increment: (state, ev) => ({ count: state.count + ev.amount })
  }
});

This projection, defined in packages/event-core/src/projection/define.ts, maintains a simple counter state that increments in response to events.

Creating a Forked Projection

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

export const WhatIfCounter = defineForked<IncrementEvent, CounterState>()({
  name: 'WhatIfCounter',
  init: () => ({ count: 0 }),
  apply: {
    increment: (state, ev) => ({ count: state.count + ev.amount })
  }
});

// Usage within a task:
const forkId = 'speculation-1';
yield* publish({ type: 'increment', amount: 5, forkId });
const forkedState = yield* getFork(forkId);

The forked variant, from packages/event-core/src/projection/defineForked.ts, allows parallel execution paths without affecting the main projection state.

Implementing a Worker for External API Calls

import { defineWorker } from '@magnitudedev/event-core/worker';
import { Effect } from 'effect';

export const HttpFetcher = defineWorker<HttpFetchEvent>()({
  name: 'HttpFetcher',
  on: {
    httpFetch: (ev) =>
      Effect.flatMap(
        Http.request(ev.url).pipe(Http.run),
        (response) => publish({ type: 'httpResponse', data: response })
      )
  },
  completeOn: ['sessionEnded']
});

Defined in packages/event-core/src/worker/define.ts, this worker handles HTTP requests asynchronously and publishes response events back into the system. The completeOn array specifies events that trigger worker termination, with behavior tested in packages/event-core/src/worker/define.interrupt-coordinator.test.ts.

Capturing and Restoring Snapshots

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

const engine = Engine.make({
  projections: [Counter],
});

yield* engine.publish({ type: 'increment', amount: 3 });
const snapshot = yield* engine.captureProjectionSnapshot(
  { index: 0, timestamp: Date.now() }, 
  'session-123'
);

// Later restoration:
yield* engine.prepareProjectionSnapshotRestore(snapshot);

This pattern utilizes functions from packages/storage/src/sessions/storage.ts to enable durable state management across sessions.

Summary

Frequently Asked Questions

How do projections and workers differ in Magnitude's event sourcing model?

Projections derive immutable read-only state from event streams using Effect-Schemas, making them deterministic and replayable. Workers handle side-effects like API calls and database writes, executing asynchronously outside the pure event-sourcing flow according to packages/event-core/src/worker/define.ts. While projections update internal state objects that can be snapshotted, workers perform actions that cannot be rolled back or replayed.

What is the purpose of forking in Magnitude projections and workers?

Forking allows multiple isolated instances of the same projection or worker to process identical event streams with different forkId identifiers. This enables "what-if" scenarios, speculative execution, and parallel agent coordination without polluting the main state. Each fork maintains independent state in projections via getFork(forkId), while forked workers run in isolated scopes until their completeOn condition triggers termination.

Can workers emit events that projections consume?

Yes. Workers can publish new events using the publish function, which flow through the event bus and trigger both projections and other workers. This creates reactive pipelines where a worker's external IO results in state updates via projections, or triggers additional side-effects through other workers according to the implementation in packages/event-core/src/core/worker-bus.ts.

How does Magnitude handle persistence of projection state?

The framework provides writeProjectionSnapshot and prepareProjectionSnapshotRestore functions in packages/storage/src/sessions/storage.ts to serialize projection state to storage and reload it later. This snapshotting capability, exposed through the Engine API, enables fast session restoration and debugging without requiring full event log replay.

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 →