eventador 4.0.0 copy "eventador: ^4.0.0" to clipboard
eventador: ^4.0.0 copied to clipboard

Persistence & Event Sourcing Extension for Dactor - Industrial-grade event sourcing platform with hybrid architecture

4.0.0 (unreleased) #

Nothing is printed to stdout any more. Every print (the [PersistentActor] recovery lines, Command processing failed for …, snapshot and projection failures) is a package:logging record under a logger named eventador.<Class> (eventador.PersistentActor, eventador.SnapshotManager, eventador.ProjectionManager, …): recovery progress at FINE, recoverable failures at WARNING, a dropped message or a failed recovery at SEVERE. A host speaking a protocol over stdio (an MCP server) no longer has to divert the library's output; one that wants the old lines attaches a listener to Logger.root.

Moves from isar 3.1.0+1 to isar_community 3.3.2. The original isar is unmaintained, and its Android native library is aligned to 4 KB pages, which Google Play rejects for apps targeting Android 15 or later.

This is a major release because Isar appears in the public API (the event store, snapshot store and projection checkpoints take an Isar instance), so the host must move to the same package:

dependencies:
  eventador: ^4.0.0
  isar_community: ^3.3.2
dev_dependencies:
  isar_community_generator: ^3.3.2

Change import 'package:isar/isar.dart' to import 'package:isar_community/isar.dart' and regenerate .g.dart files. The collection schemas and their ids are unchanged, so existing databases keep their layout. Hosts using DuraQ's Isar backend need duraq_isar 3.0.0.

3.1.0 #

Two additions, both of them things a host could not previously say. Nothing existing changes: upgrading from 3.0.0 is a version bump and no code edits.

A pending event awaiter can be cancelled #

ProjectionActor held a registration made with AwaitEventApplied until a matching event arrived, its own timeout elapsed, or the projection stopped. There was no way to abandon one. A caller that registers an awaiter before issuing a command — the standard way to close the gap between a command being accepted and the projection having applied its event — has nothing left to wait for when that command is rejected, but its registration stayed, holding a closure, a reply target and a timer for the whole window. On a path a remote peer can drive, that is a cost the peer can impose at will.

  • AwaitEventApplied(..., awaitId: 'my-id') (new, optional) names a registration. The id is the caller's own and is matched by value.
  • CancelEventAwait('my-id') (new) drops every registration carrying that id, cancels their timers, and answers each one AwaitFailed(reason: 'cancelled') — so a caller waiting on the original ask is released immediately rather than at the timeout. Answered with EventAwaitCancelled(cancelled: n); cancelling an id that is not registered is not an error and cancels nothing.

Additive and backwards compatible: a registration made without an awaitId behaves exactly as before and no cancel can touch it. 'cancelled' joins 'timeout', 'stopped' and 'error' as an AwaitFailed.reason.

A host can delete one aggregate's history #

IsarEventStore.deleteEvents drops the events of a single aggregate. Retention needed a way to say "this run's events are no longer worth keeping": deleteOldSnapshots dropped the snapshot while the events stayed forever. It lives on IsarEventStore rather than the EventStore interface, so in-memory and mock stores implementing that interface are untouched. It is not undoable, and the doc says so.

3.0.0 #

Dependency upgrade and correctness release. Eventador now tracks dactor 1.3.0 and duraq 3.0.0 (with duraq_isar 2.0.0 for the Isar backend). Both dependencies shipped major fixes; bringing them in exposed several defects in eventador's own saga and streaming code, which this release fixes.

Upgrading from 2.x #

  1. Bump the dependencies together.

    dependencies:
      eventador: ^3.0.0
      dactor: ^1.3.0
      duraq: ^3.0.0
      duraq_isar: ^2.0.0   # if you use DuraQ's Isar backend
    

    IsarStorage moved out of duraq into duraq_isar; add import 'package:duraq_isar/duraq_isar.dart'; where you construct it. Read the duraq 2.0.0 and 3.0.0 changelogs: existing queue databases are migrated in place on first open and cannot be reopened by older releases.

  2. Register the commands your sagas send. DuraQ persists payloads as JSON. Saga.sendCommand now stores the command through Command.toMap() and the coordinator rebuilds it with CommandRegistry.fromMap, so each such command type needs a toMap() override and a CommandRegistry registration. Transient fields (replyTo, sender) do not survive the queue. Previously any sendCommand / scheduleTimeout against duraq ≥ 2 threw PayloadCodecException.

  3. SagaCoordinator is now started and stopped. The constructor is SagaCoordinator(queueManager, actorSystem, {retryPolicy, pollInterval, onError}) — the unused event store parameter is gone. Call await coordinator.start() and await coordinator.stop(). processSagaCommands() / processSagaTimeouts() still exist as single-entry steps and now return Future<bool>. Previously each was a one-shot call that processed a single entry and returned, despite the documentation describing a continuous background task; and an entry whose target actor was missing was acknowledged and lost. Now it is retried with the coordinator's RetryPolicy and then dead-lettered.

  4. Give saga states a toMap(). Saga.getSnapshotState() returns sagaState and lets the event store serialize it (via toMap() when present); onSagaStateRestored receives the decoded Map<String, dynamic>. Previously the state was CBOR-encoded twice and the restore path never matched (it checked for List<int> and received List<dynamic>), so saga snapshots were silently ignored on recovery.

  5. preStart is now Future<void> on PersistentActor and AggregateRoot. Subclasses that override it must declare Future<void> preStart() async and await super.preStart(). Dactor 1.3 awaits preStart before delivering any message, so recovery now completes before the first message is processed and spawn() returns a recovered actor. Messages sent during recovery wait in the mailbox. A failed recovery is reported through onRecoveryFailure / recoveryComplete and no longer surfaces as an unhandled asynchronous error.

  6. IsarEventStore.create() opens allSchemas (event store plus projection checkpoints), so the returned store's isar can back ProjectionActor. Existing databases gain the collection on open. The new downloadIsarCore flag lets production builds refuse to fetch the native library at runtime.

  7. Event.toMap(), Command.toMap() and State.toMap() are no longer @protected — serializers call them by design. Event.isValid() and Event.getValidationErrors() likewise.

  8. The placeholder Awesome class and eventador_base.dart export are removed.

  9. DateTime fields are now tagged in CBOR; untagged strings stay strings. 2.x turned every stored string that looked like an ISO-8601 timestamp into a DateTime on read. 3.0.0 writes DateTime values with CBOR tag 0 and decodes only tagged values as dates, so a String field that happens to contain a timestamp is returned as a String. Rows written by 2.x hold untagged dates and therefore decode as String; fromMap code that follows the README's "accept String or DateTime" pattern is unaffected. If yours casts as DateTime and you cannot change it, set CborSerializer.parseUntaggedDates = true at startup to keep the 2.x heuristic for old rows. Event.toMap() / Command.toMap() / State.toMap() still write their own timestamp / lastModified as ISO-8601 strings, as before.

  10. SnapshotManager.registerActor returns Future<void>. The actor is tracked as soon as it is called, but awaiting it guarantees the count of events since the last snapshot has been seeded from the store. PersistentActor.preStart awaits it; code that registers actors by hand should too.

  11. SnapshotConfig.enableCompression and compressionThreshold are removed, along with SnapshotConfig.shouldCompress and the identity CborSerializer.compress / decompress. None of them ever did anything; snapshots were always stored uncompressed and still are. Drop the arguments from your SnapshotConfig(...) calls. maxSnapshotSize, which was also never consulted, is now enforced (see Fixed); if your snapshots are larger than the default 10 MB, raise it.

Repository and tests #

  • example/libisar.dylib, a macOS Isar core binary that had been committed by accident (Isar downloads it next to whichever script calls initializeIsarCore(download: true)), is no longer tracked; *.dylib, *.so and *.dll are ignored.
  • The docs/ folder is now doc/, the pub layout convention, and generated .g.dart / .mocks.dart files are excluded from analysis.
  • The test suite no longer races on that download. Every test file runs in its own isolate, and parallel files fetched the same library to the same path at once, so dart test failed most runs with "file does not start with MH_MAGIC". Tests now call initIsarCoreForTests() (test/support/isar_core.dart), which downloads once into .dart_tool/isar_core/<version>/ under a file lock. The first run after a checkout or an Isar upgrade needs network access; later runs do not.

Deprecated #

  • ProjectionManager — use ProjectionActor, which is now the single projection runtime. The manager's stream listener was async and unserialized, so a handle() that awaited could overlap with the next event's handler and update the read model out of order. ProjectionActor processes events in mailbox order. The manager now pauses its subscription while handle() runs (so it is safe for the one release it remains), and will be removed in 4.0.0. Migration table in the class docs and README; checkpoints are shared, so a moved projection resumes where it left off.

Fixed #

  • Saga queues had no codec. SagaCommandEnvelope and SagaTimeout now ship QueueCodecs (SagaCommandEnvelope.codec, SagaTimeout.codec), used by Saga and SagaCoordinator.
  • Timeout ids collided across sagas. Timeout queue entries are keyed by '$sagaId/$timeoutId' (SagaTimeout.entryId) and re-scheduling replaces the pending entry. Added Saga.cancelTimeout.
  • Event and command ids could collide. Both were derived from the wall clock (cmd_<ms>_<ms % 10000> gave identical ids to two commands created in the same millisecond, which command deduplication then dropped; event ids collided within a microsecond and hit the unique eventId index). Both are now UUID v4.
  • Live event streams could skip events. allEventsWithSequence forwarded live events during the historical replay and advanced its high-water mark past unread history, dropping the remaining historical events; it also looked every live event up by id with an unindexed filter and an async listener, which could reorder events. Live events now carry their journal position from the store, are buffered while the replay runs, and are drained in order afterwards. allEvents, eventsByTag and eventsByPersistenceId share the same path.
  • persistEvents ignored automatic snapshots and only guarded against isRecovering, not an incomplete recovery. AggregateRoot persists exclusively through persistEvents, so aggregates never triggered count-based snapshots.
  • ProjectionActor no longer replays from 0 when its checkpoint cannot be read. A storage error at start was swallowed and the fallback position (0) replayed the whole journal into a read model that may already contain it. The actor now enters ProjectionStatus.error without subscribing, calls projection.onError, and ResumeProjection retries the read. Checkpoint write failures are still non-fatal but are now reported via onError and ProjectionInfo.lastError. loadCheckpoint / persistCheckpoint are @protected extension points. The deprecated ProjectionManager throws ProjectionException from registerProjection / start / resumeProjection in the same case.
  • ProjectionActor.RebuildProjection now resets the persisted checkpoint to 0, so a crash mid-rebuild resumes from the start against the reset read model instead of the pre-rebuild position.
  • ProjectionManager recorded one processed event per checkpoint write instead of per event.
  • CommandHandler validated generated events against ValidatableState (never true for an event); it now uses ValidatableEvent.
  • CBOR deserialization no longer guesses DateTime from a string's shape. CborSerializer turned any string matching yyyy-MM-ddTHH:mm:ss into a DateTime, so a string field holding a timestamp came back as the wrong type and as String casts in fromMap threw. Dates are now written with CBOR tag 0 (CborDateTimeString) and only tagged values decode as DateTime; the round trip is exact (UTC marker and microseconds preserved). Event.toCbor / EventRegistry.fromCbor and Command.toCbor / CommandRegistry.fromCbor previously had their own copies of the value conversion that stringified DateTime with toString() and never decoded it; all three paths now share one codec. The old heuristic is available as CborSerializer.parseUntaggedDates for journals written by 2.x (see Upgrading, step 9).
  • Snapshot tracking survives a restart, and the time threshold no longer fires on the first event. SnapshotManager.registerActor reset the events-since-snapshot count to zero on every start, so an actor restarted just short of its threshold began counting again from nothing, and a long-lived actor that restarted often could never reach it. Registration now seeds the count from the journal (highest sequence minus the snapshot's sequence) and the last snapshot time from the snapshot itself, so both eventCountThreshold and minTimeBetweenSnapshots are measured across runs. The time rule used to answer "yes" when no snapshot existed, which meant a fresh actor with a timeThreshold snapshotted on its very first event; it is now measured from registration until the first snapshot. The CallbackSnapshotManager timer no longer writes a snapshot when no event has been persisted since the last one.
  • The persisted ProjectionCheckpoint row now records status and last error. Its status and lastError columns were never written by either runtime (the row always read running, null), so storage could not tell an operator whether a projection was paused, stopped or failing. ProjectionActor now writes the row on every status transition (running on start and resume, paused, error when the checkpoint cannot be read, stopped on StopProjection and on an abnormal stop) and records lastError when handle() throws, the event stream fails, or a checkpoint write fails. lastError is cleared by a successful (re)start or rebuild. StopProjection always writes the row, not only when events were pending. The deprecated ProjectionManager mirrors status on pause, resume and stop. New @protected ProjectionActor.persistStatus extension point alongside persistCheckpoint.
  • CallbackSnapshotManager took two snapshots per threshold. When the event-count threshold was reached, the manager scheduled a snapshot through its callback and the actor's own post-persist check created a second one, because the first had not yet reset the counter. A snapshot in progress (scheduled or saving) now blocks a second one for the same actor: shouldCreateSnapshot answers false and createSnapshot returns false while one is in flight. The base SnapshotManager no longer schedules no-op triggers or runs a timer it cannot act on; its thresholds are applied when the actor asks.
  • SnapshotConfig.maxSnapshotSize is enforced and snapshot statistics measure real sizes. The manager serializes the state with CborSerializer before saving; a state larger than maxSnapshotSize is refused (createSnapshot returns false, counted in SnapshotStats.snapshotsRejected) and totalSnapshotSize / averageSnapshotSize report the serialized size instead of a constant per-entry guess.
  • ActorRefs in metadata no longer reach storage or break replay. Event.replyTo read metadata['replyTo'] as ActorRef?, but metadata is persisted with the event, so a ref placed there was stored as its toString() and the cast threw on every replay of that event. The same applied to Command, whose constructor read replyTo from metadata when a command was rebuilt from a saga queue. ActorRef values are now dropped from metadata by toMap() and by the event store's metadata column (Event.persistableMetadata / Command.persistableMetadata), and replyTo returns null instead of throwing when the stored value is not a ref. In-memory use is unchanged: a ref in metadata is still returned by replyTo until the event is persisted.
  • Type identity no longer has to be the Dart class name. Events, commands and snapshots were stored under runtimeType.toString(), so renaming a class, or building Flutter with --obfuscate, made every existing row undeserializable with no way to recover. Event, Command and State now expose String get typeName (default: the class name, so existing journals are unaffected). It is what toMap()['type'], the eventType / stateType envelope columns, the three registries and EventMigrationManager use. Override it with a name that never changes; for types already stored under a class name, EventRegistry.registerAlias (also on CommandRegistry / StateRegistry, or the aliases: argument to register) keeps the old rows readable.

Performance #

  • EventEnvelope has a composite (persistenceId, sequenceNumber) index. getEvents, getHighestSequenceNumber and eventsByPersistenceId read the requested range in index order instead of loading and sorting an actor's whole journal. Isar applies the index change on open.
  • Journal pagination is keyset-based (id > last) rather than offset-based.
  • Projection checkpoint reads and writes use the unique projectionId index (where()), not a full-collection filter().
  • Live subscribers no longer run an Isar query per event.

Added #

  • AggregateRoot accepts snapshotManager and maxProcessedCommandIds and forwards them to PersistentActor; aggregates previously had no way to opt into automatic snapshots or command deduplication.
  • sagaCommandQueueName / sagaTimeoutQueueName constants; Saga.commandQueue / Saga.timeoutQueue accessors.
  • SagaDeliveryException, thrown by the coordinator when a target actor is not registered (drives DuraQ retries / dead-lettering).
  • SnapshotState.registeredAt.
  • ProjectionActor.persistStatus (@protected).
  • SnapshotStats.snapshotsRejected and SnapshotStats.addRejection().
  • Event.persistableMetadata and Command.persistableMetadata: the metadata as written to storage, with transient ActorRef entries removed.
  • CborSerializer.parseUntaggedDates compatibility flag.
  • Event.typeName, Command.typeName, State.typeName; registerAlias / resolve and an aliases: argument to register on EventRegistry, CommandRegistry and StateRegistry; CommandRegistry.isRegistered.

2.2.0 #

Added #

  • Stream<Event> get appliedEvents on ProjectionActor — a broadcast stream of events the projection has successfully applied. Each event is emitted strictly after projection.handle(event) returns true and the checkpoint state has been updated, in projection processing order.

    Use this when a downstream consumer needs post-applied semantics — i.e., guarantees that the read model reflects the event by the time the consumer fires. The motivating use case is an event-stream-coupled adapter (e.g., a P2P bridge) whose listeners then query the read model: subscribing to the raw aggregate stream creates a race with the projection; subscribing to appliedEvents does not.

    Events for which handle() returns false (not handled) or throws are not emitted on this stream. The stream signals onDone when the actor processes StopProjection.

    late final ProjectionActor projectionActor;
    await system.spawn(
      'projection-${p.projectionId}',
      () {
        projectionActor = ProjectionActor(p, eventStream);
        return projectionActor;
      },
    );
    projectionActor.appliedEvents.listen((event) {
      // event has been applied to the read model; any downstream read will see it
    });
    

2.1.0 #

Added #

  • ProjectionActor<T> — a dactor-based actor wrapper around a Projection<T>, modelled after Apache Pekko / Akka's ProjectionBehavior. The actor owns the event-stream subscription, drives projection.handle(), manages batched checkpoint persistence to Isar, and serves query messages. Spawning is the registration — there is no separate manager to register with; the actor system itself is the registry, exactly as in Pekko.

    final ref = await system.spawn(
      'projection-${projection.projectionId}',
      () => ProjectionActor(projection, eventStore, isar: isar),
    );
    
  • AwaitEventApplied — the canonical CQRS command-to-read-model primitive. A coordinator that just dispatched a command can ask the projection actor to reply once the projection has applied a matching event:

    await projectionRef.ask<EventAppliedResponse>(
      AwaitEventApplied(
        (e) => e is UserEmailUpdated && e.userId == userId,
        timeout: Duration(seconds: 5),
        alreadySatisfied: () => projection.readModel.emailFor(userId) == newEmail,
      ),
    );
    

    Resolution happens strictly after projection.handle() returns for the matching event. A handle failure does not advance the checkpoint and does not resolve pending awaiters. Timeouts respond with AwaitFailed(reason: 'timeout'); stopping the actor while awaiters are pending responds with AwaitFailed(reason: 'stopped').

  • Lifecycle message protocol — GetProjectionInfo, PauseProjection, ResumeProjection, RebuildProjection, and StopProjection (Pekko's ProjectionBehavior.Stop equivalent — flushes the checkpoint, replies StoppedAck, then terminates the actor).

Notes #

  • Additive. ProjectionManager is unchanged and continues to work for callers that don't need request/response semantics. Existing callers do not need to migrate. ProjectionActor and ProjectionManager can coexist in the same application.
  • See doc/projection-actor-proposal.md for the full design rationale, Pekko mapping, and invariants.

2.0.0 #

Breaking changes #

  • Removed deprecated saga state persistence. Saga.saveSagaState() and Saga.loadSagaState() (deprecated in 1.0.0) are gone. Use the standard PersistentActor snapshot mechanism instead: override getSnapshotState() and onSnapshot() (or implement onSagaStateRestored() on Saga, which the default onSnapshot() calls). Call createSnapshot() to persist.
  • Removed EventStore.saveSagaState / loadSagaState from the interface and IsarEventStore. Custom EventStore implementations no longer need to provide these methods.
  • Removed SagaStateEnvelope and its Isar collection. The schema is no longer registered by IsarEventStore.requiredSchemas. Existing databases with SagaStateEnvelope rows will retain the collection on disk but Eventador will not read or write to it. Migrate any persisted saga state to snapshots before upgrading.

Internal cleanup #

  • Dropped the redundant _commandLock from AggregateRoot. Command serialization is already guaranteed by PersistentActor._handleCommand.

1.0.0 #

Initial pub.dev release.

  • Event sourcing framework: AggregateRoot, Command, Event, State patterns
  • Persistent actor support with Isar-backed event storage (CBOR serialization)
  • Snapshot system with configurable frequency and retention strategies
  • Saga coordination for distributed, multi-aggregate workflows
  • Projection system for CQRS read models with checkpoint-based resumability
  • EventRegistry for type-safe event deserialization across system restarts
  • ProjectionManager for automatic event streaming to registered projections
  • Built on Dactor actor model and DuraQ durable queuing
1
likes
160
points
320
downloads

Documentation

API reference

Publisher

verified publisherwerkswinkel.com

Weekly Downloads

Persistence & Event Sourcing Extension for Dactor - Industrial-grade event sourcing platform with hybrid architecture

Repository (GitHub)
View/report issues

Topics

#event-sourcing #cqrs #actors #persistence #saga

License

MIT (license)

Dependencies

cbor, dactor, duraq, isar_community, logging, meta, synchronized, uuid

More

Packages that depend on eventador