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

> Explore Magnitude's event-core concepts Event Signal Projection and Worker Discover how this effect-driven system delivers deterministic state management and handles async side effects.

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

---

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

```typescript
// 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`](https://github.com/magnitudedev/magnitude/blob/main/src/signal/define.ts) provides the type system:

```typescript
// 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`](https://github.com/magnitudedev/magnitude/blob/main/src/projection/define.ts), they combine state management with handler registration.

A projection configuration includes:

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

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

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

### Consuming from a Worker

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

### Forked Projection Pattern

```typescript
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/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`](https://github.com/magnitudedev/magnitude/blob/main/src/core/event-bus-core.ts) | Central event bus, sequential queue, persistence |
| [`src/core/worker-bus.ts`](https://github.com/magnitudedev/magnitude/blob/main/src/core/worker-bus.ts) | Worker-facing publish/subscribe API |
| [`src/core/projection-bus.ts`](https://github.com/magnitudedev/magnitude/blob/main/src/core/projection-bus.ts) | Synchronous projection processing |
| [`src/core/interrupt-coordinator.ts`](https://github.com/magnitudedev/magnitude/blob/main/src/core/interrupt-coordinator.ts) | Fork-aware interrupt signaling |
| [`src/core/framework-error.ts`](https://github.com/magnitudedev/magnitude/blob/main/src/core/framework-error.ts) | Centralized error reporting |
| [`src/projection/define.ts`](https://github.com/magnitudedev/magnitude/blob/main/src/projection/define.ts) | Projection builder with handlers and read tracking |
| [`src/worker/define.ts`](https://github.com/magnitudedev/magnitude/blob/main/src/worker/define.ts) | Worker builder with effect helpers |
| [`src/signal/define.ts`](https://github.com/magnitudedev/magnitude/blob/main/src/signal/define.ts) | Signal definitions and Pub/Sub |
| [`src/worker/util.ts`](https://github.com/magnitudedev/magnitude/blob/main/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`](https://github.com/magnitudedev/magnitude/blob/main/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.