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
Queueto maintain ordering guarantees - Persistence control: Mark events as
ephemeralto skip persistence entirely - Broadcast to projections: The
ProjectionBusreceives 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
queueSignalon theProjectionBus - 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:
- Publication:
WorkerBus.publish(event)orEventBusCore.publish(event)initiates - Queuing:
EventBusCoreenqueues in its sequentialQueue - Projection processing:
ProjectionBusruns event handlers synchronously, updating state and queuing signals - Signal broadcast: Signals publish to
ProjectionBus(sync) and Pub/Sub (async) - Worker activation: Subscribed workers receive messages in separate fibers with fork-scoped
readfunctions
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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →