# Understanding Projections and Workers in Magnitude's Event Sourcing Model

> Learn how Magnitude's event sourcing model uses projections for immutable read-only state and workers for asynchronous side-effects. Discover forked execution for isolated processing.

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

---

**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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/projection/defineForked.ts) and [`packages/event-core/src/worker/defineForked.ts`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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

```typescript
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`](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/projection/define.ts), maintains a simple counter state that increments in response to events.

### Creating a Forked Projection

```typescript
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`](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/projection/defineForked.ts), allows parallel execution paths without affecting the main projection state.

### Implementing a Worker for External API Calls

```typescript
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`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/worker/define.interrupt-coordinator.test.ts).

### Capturing and Restoring Snapshots

```typescript
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`](https://github.com/magnitudedev/magnitude/blob/main/packages/storage/src/sessions/storage.ts) to enable durable state management across sessions.

## Summary

- **Projections** maintain pure, read-only state derived from events using handlers defined in [`packages/event-core/src/projection/define.ts`](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/projection/define.ts).
- **Workers** execute side-effects and IO operations asynchronously, defined in [`packages/event-core/src/worker/define.ts`](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/worker/define.ts).
- Both components receive events through distinct buses: the ProjectionBus ([`packages/event-core/src/core/projection-bus.ts`](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/core/projection-bus.ts)) and WorkerBus ([`packages/event-core/src/core/worker-bus.ts`](https://github.com/magnitudedev/magnitude/blob/main/packages/event-core/src/core/worker-bus.ts)).
- **Forking** enables isolated execution paths through `defineForked` variants in both projections and workers, supporting speculative execution and parallel processing.
- **Snapshots** allow projections to persist and restore state via `writeProjectionSnapshot` and `prepareProjectionSnapshotRestore`.
- Workers can terminate automatically when specific `completeOn` events occur, preventing resource leaks in long-running processes.

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