How Magnitude Registers Projections and Manages State in event-core
Magnitude's event-core package registers projections via an Effect-TS layer-based registry and manages their state through in-memory instances, deterministic event replay, and optional snapshots backed by the Addressed storage abstraction.
The magnitudedev/magnitude repository implements an event-sourcing architecture where read models materialize from the immutable event stream through projections. Understanding how event-core handles projection registration and state management is essential for extending the framework or debugging replay behavior.
Projection Registration Architecture
The Projection Interface
Every projection implements the Projection interface defined in packages/event-core/src/core/projection.ts. This contract declares the event handlers responsible for mutating the projection's internal state in response to domain events.
The Registration API
The central registration entry point resides in packages/event-core/src/core/engine.ts. The EventEngine exposes a registerProjection method with the following signature:
registerProjection<T extends Projection>(name: string, projection: T): Layer.Tag<T>
This function accepts a unique string identifier and a concrete projection instance, returning an Effect-TS Layer.Tag for dependency injection. The engine maintains an internal registry map (Map<string, Projection>) that stores all registered projections by name.
Layer Composition and Discovery
Magnitude leverages Effect-TS layers for dependency injection. Each call to registerProjection produces a layer that must be merged into the main application layer. For example, in the agent package (packages/agent/src/projections/index.ts), all projection layers are composed:
export const AgentProjectionsLayer = Layer.mergeAll(
TurnProjectionLayer,
SessionContextLayer,
// ... other projections
);
During engine startup, EventEngine iterates over the merged layers, extracts the Projection tags, and populates the registry map.
State Management Lifecycle
In-Memory Projection State
Each projection instance maintains its own mutable state within a plain JavaScript object. The state shape is defined by the projection author. For instance, a Turn projection might track the current turn ID and accumulated messages, while a counter projection stores a simple numeric value.
Durable Storage with Addressed
For projections requiring persistence beyond a single process lifecycle, event-core provides the Addressed abstraction in packages/event-core/src/core/addressed.ts. This generic key-value store, backed by the session storage layer in packages/event-core/src/core/storage.ts, allows projections like SessionContext to read and write durable state across restarts.
Event Replay and State Reconstruction
When the engine initializes or restores from a snapshot, it replays the event log to reconstruct projection state. The replay loop, implemented in EventEngine, iterates through stored events and dispatches them to the appropriate handler:
for (const ev of eventLog) {
const proj = registry.get(ev.projectionName);
proj?.handle(ev);
}
All state updates occur within a single Effect-TS transaction, ensuring atomicity and serializing access to prevent race conditions between projections.
Snapshotting for Fast Recovery
To avoid replaying the entire event history, projections can implement a snapshot() method that returns a serializable representation of their current state. The ProjectionSnapshotRestorePlan in packages/event-core/src/core/snapshot.ts uses these snapshots to fast-load state, significantly reducing startup time for large event logs.
Practical Example: Building a Counter Projection
Below is a complete example demonstrating the registration flow and state management for a custom projection:
// 1. Define the projection class
class CounterProjection implements Projection {
private count = 0;
handle(event: { type: 'increment'; amount: number }) {
if (event.type === 'increment') {
this.count += event.amount;
}
}
getCount(): number {
return this.count;
}
snapshot(): number {
return this.count;
}
}
// 2. Register with the engine
import { EventEngine, registerProjection } from '@magnitudedev/event-core';
import { Layer } from '@effect-ts/core';
const CounterLayer = registerProjection('counter', new CounterProjection());
// 3. Merge into the application layer
const AppLayer = Layer.mergeAll(
EventEngine.layer,
CounterLayer
);
// 4. Consume the projection
import { Effect } from '@effect-ts/core';
const getCountEffect = Effect.serviceWithEffect(CounterProjection)((proj) =>
Effect.succeed(proj.getCount())
);
This pattern illustrates how state is encapsulated within the projection instance, registered via the layer system, and accessed through Effect-TS effects.
Key Source Files
| File Path | Responsibility |
|---|---|
packages/event-core/src/core/projection.ts |
Defines the Projection interface and base types |
packages/event-core/src/core/engine.ts |
Implements EventEngine, registration API, and replay loop |
packages/event-core/src/core/addressed.ts |
Provides the Addressed storage abstraction for durable state |
packages/event-core/src/core/snapshot.ts |
Handles snapshot creation and restoration logic |
packages/agent/src/projections/index.ts |
Composes concrete projection layers for the agent package |
Summary
- Registration occurs via
registerProjectioninEventEngine, which stores instances in aMap<string, Projection>registry. - Dependency injection uses Effect-TS layers, requiring projection layers to be merged into the main application layer.
- State is stored in-memory within projection instances, with optional persistence through the
Addressedabstraction. - Reconstruction happens through deterministic event replay, with snapshots available to optimize startup performance.
- All operations are atomic and serialized through Effect-TS to prevent concurrency issues.
Frequently Asked Questions
How does event-core ensure projection state consistency during replays?
The EventEngine serializes all event handling through Effect-TS transactions. During replay, events are processed sequentially in a single transaction, ensuring that projections update atomically and maintain consistency even when handling complex dependency graphs between projections.
Can projections share state or access each other directly?
Projections should not share mutable state directly. Instead, they communicate through the event log. If one projection needs data from another, it should react to the same events or query a shared read model through the Addressed storage abstraction in packages/event-core/src/core/addressed.ts.
What happens if a projection fails to handle an event during replay?
If a projection's handle method throws or returns a failing Effect, the entire transaction fails. This halts the replay process and prevents inconsistent state, forcing the system to retry or alert operators based on the configured error handling strategy in EventEngine.
Where is the projection registry stored at runtime?
The registry exists as an in-memory Map inside the EventEngine instance defined in packages/event-core/src/core/engine.ts. It is populated during layer initialization and persists for the application's lifetime, mapping string names to projection instances.
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 →