@bimetal/event-sourcing
v0.37.0
Published
Domain-agnostic event-sourcing primitives: store, aggregate, projection, snapshots, middleware. Foundation for @bimetal/calendar-data, @bimetal/table-data and other domain packages.
Maintainers
Readme
@bimetal/event-sourcing
Domain-agnostic event-sourcing primitives. Zero external dependencies, zero @bimetal/* dependencies. Foundation for domain packages like @bimetal/calendar-data and @bimetal/table-data.
Installation
npm install @bimetal/event-sourcingWhat's Inside
Types
DomainEvent<T, P>— immutable, append-only event with metadata (correlationId,userId,source,version)Command<T, P>— input to a dispatch pipelineEventStore— append/read/subscribe interface; storage adapters (store-sqlite,store-eventstoredb, …) implement thisEventBus— in-process pub/subAggregate<S, C, E>—handle(state, command) → events[]+apply(state, event) → stateProjection<R>— derives read models from eventsAggregateState<T>— internal state with per-entity undo stack + processed-command-id windowAggregateSnapshot<T>,SnapshotStore<T>— point-in-time captures for fast hydrationDomainStore<S, C, E, R>— generic facade interface implemented by concrete domain storesDispatchContext<Store>,DispatchMiddleware<Cmd, Evt, Store>— pipeline hooksClock—{ now(): number }time source for event timestamps (was previously in@bimetal/temporal)EventMetadataProvider— user/source enrichmentDomainError—CONCURRENCY_CONFLICT | STREAM_NOT_FOUND | VALIDATION_ERROR | ENTITY_NOT_FOUND
Errors
DataValidationError— thrown by domain aggregates on invalid commandsConcurrencyError— thrown by anEventStorewhenexpectedVersiondoesn't match
In-Memory Adapters
createInMemoryEventStore()— non-persistedEventStorefor testscreateInMemorySnapshotStore<T>()— non-persistedSnapshotStore<T>createEventBus()—EventBus
Topic-Convention (für Broker-Konsumenten)
eventToTopic(event)— Topic-Convention für@bimetal/brokerund kompatible Adapter. Liefertevent.type. Vollständig deterministisch, keine Mapping-Magie.topicForEventType(type)— Identity-Helper für Subscriptions ohne Event-Objekt.
Da DomainEvent-Types der FQN-Konvention <vendor>.<domain>.<EventName> folgen, funktionieren Pattern-Subscriptions im RabbitMQ-Style direkt: bimetal.# (alle), bimetal.calendar.* (alle Calendar-Events), bimetal.calendar.CalendarEventCreated (exakt).
Bewusst kein Import aus @bimetal/broker — die Verbindung passiert auf der Aufrufer-Seite:
import { eventToTopic } from '@bimetal/event-sourcing';
import { createInProcessBroker } from '@bimetal/broker';
import type { DomainEvent } from '@bimetal/event-sourcing';
const broker = createInProcessBroker<DomainEvent>();
await broker.publish(eventToTopic(event), event);EventStore → Broker Bridge (bridgeEventStoreToBroker)
Eng geschnittene Contract-Utility für die kritische Verbindungsstelle EventStore ↔ Broker.
bridgeEventStoreToBroker(eventStore, broker, options?)— subscribed aufeventStore, ruftbroker.publish(topic, event)pro Event, returntUnsubscribe.- Strukturelle Type-Dependency — kein Import aus
@bimetal/broker. Jeder Adapter, der ein strukturellesPublishTarget<T> = { publish(topic, message): Promise<void> }erfüllt, passt automatisch. - Verantwortlich für: post-append Reihenfolge, async
.catchaufpublish(), sichtbarer Fehlerpfad (Defaultconsole.error, überschreibbar viaonPublishError). - Bewusst nicht zuständig für Retry, Dead-Letter, Backpressure — diese gehören in den
onPublishError-Callback bzw. um die Bridge herum.
import { bridgeEventStoreToBroker } from '@bimetal/event-sourcing';
import { createInProcessBroker } from '@bimetal/broker';
const broker = createInProcessBroker<DomainEvent>();
const off = bridgeEventStoreToBroker(eventStore, broker, {
streamId: 'main', // optional, sonst alle Streams
toTopic: (event) => event.type, // optional, Default = eventToTopic
onPublishError: (err, event, topic) => {
deadLetterQueue.enqueue({ event, topic, error: err });
},
});
// später: off();Broker → EventStore Bridge (bridgeBrokerToEventStore)
Das Gegenstück: eingehende Broker-Nachrichten in einen EventStore übernehmen — mit den beiden Bausteinen, die sonst jede App selbst baut.
bridgeBrokerToEventStore(broker, eventStore, { streamId, … })— subscribed auf den Broker (Default-Pattern#, alternativ exaktetopics), appended jedes passende Event mitexpectedVersion = envelope.version - 1und returnt ein Handle mitclose()undguardPublishTarget().- Anschluss-Guard — nur was exakt an den lokalen Stream anschließt, wird übernommen. Bereits Bekanntes fällt still weg; eine echte Lücke meldet
onGap. - Echo-Bremse —
guardPublishTarget(broker)hüllt das Publish-Target so, dass von remote übernommene Events nicht zurück auf den Bus gehen. - Der DomainStore foldet selbst. Die Bridge greift am EventStore an; die externe Append-Subscription des Stores erkennt das fremde Event (es fehlt in seinem
ownEventIds-Set) und wendet es an. Deshalb ist das Rezept domänenfrei. - Harte Voraussetzung: beide Endpunkte seeden deterministisch identisch (gleicher
streamId, gleiche Seed-Events in gleicher Reihenfolge — oder beide leer). Der Guard vergleicht Versionsnummern, keine Inhalte. - Kein Catch-up. Bei
capabilities.ordering: 'none'kann ein Event außer der Reihe ankommen; der Guard verwirft es, und Broker mitreplay: falseliefern nichts nach.onGapist das Signal „dieser Endpunkt muss neu seeden", nicht Log-Rauschen. - Keine Vertrauensgrenze von sich aus. Ohne
acceptprüft die Bridge nur die Form der Nachricht, nicht dentype, nicht dieschemaVersion, nicht die Payload. Auf einem Kanal, den mehr als der eigene Code erreichen kann (broker-broadcast: jedes Script desselben Origins), ist das nicht reparabel: der Log ist append-only — ein einziges Fremd-Event vergiftet ihn dauerhaft,hydrate()wirft ab da bei jedem Start.acceptgreift vor dem Append und ist fail-closed; abgelehnte Events werden still verworfen. close()gibt einPromise<void>. Es löst die Subscriptions und verwirft noch eingereihte Übernahmen sofort; eine bereits laufende lässt sich nicht abbrechen (kein Cancel imEventStore-Vertrag), das Promise resolved, wenn sie durch ist. Ihre Echo-Bremse bleibt bis dahin bestehen.inbound.close()als Anweisung genügt,await inbound.close()ordnet den Teardown.adoptedIdCap(Default 512) ist nicht überall nur ein Sicherheitsnetz. Bei Stores, die synchron inappend()notifizieren (in-memory, SQLite), hält das Bremsen-Set nie mehr als eine Id. Der ESDB-Adapter notifiziert lokal nicht — die Subscriber sehen das Event erst über die Live-Wire-Subscription, eine Netzwerkrunde später; dort sammeln sich die Ids an, und über der Grenze verdrängt die FIFO-Kappung noch benötigte.
import { createInProcessBroker } from '@bimetal/broker';
import {
bridgeBrokerToEventStore,
bridgeEventStoreToBroker,
createInMemoryEventStore,
type DomainEvent,
} from '@bimetal/event-sourcing';
const KNOWN_TYPES = new Set(['bimetal.calendar.CalendarEventCreated']);
const eventStore = createInMemoryEventStore();
const broker = createInProcessBroker<DomainEvent>();
const inbound = bridgeBrokerToEventStore(broker, eventStore, {
streamId: 'main',
// Vertrauensgrenze VOR dem Append — der Log ist append-only.
accept: (event) => KNOWN_TYPES.has(event.type) && event.metadata.schemaVersion === 1,
onGap: () => { window.location.reload(); }, // kein Replay — nur Neu-Seeding hilft
});
const off = bridgeEventStoreToBroker(
eventStore,
inbound.guardPublishTarget(broker), // ← Echo-Bremse
{ streamId: 'main' },
);
// später: off(); await inbound.close();Use From a Domain Package
import type {
AggregateState,
Command,
DomainEvent,
DomainStore,
Aggregate,
Projection,
} from '@bimetal/event-sourcing';
interface TableRow { readonly id: string; readonly cells: Record<string, unknown>; }
type CreateRow = Command<'CreateRow', { row: TableRow }>;
type RowCreated = DomainEvent<'RowCreated', { row: TableRow }>;
type RowState = AggregateState<TableRow>;
type RowStore = DomainStore<RowState, CreateRow, RowCreated, { rows: readonly TableRow[]; version: number }>;For complete examples, see @bimetal/calendar-data and @bimetal/table-data.
Replaces
This package contains the domain-agnostic portion of the (deprecated) @bimetal/data package, which has been removed. Calendar-specific aggregate/projection/store live in @bimetal/calendar-data.
License
Apache License 2.0
