eventador 4.0.0
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 oneAwaitFailed(reason: 'cancelled')— so a caller waiting on the originalaskis released immediately rather than at the timeout. Answered withEventAwaitCancelled(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 #
-
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 backendIsarStoragemoved out ofduraqintoduraq_isar; addimport '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. -
Register the commands your sagas send. DuraQ persists payloads as JSON.
Saga.sendCommandnow stores the command throughCommand.toMap()and the coordinator rebuilds it withCommandRegistry.fromMap, so each such command type needs atoMap()override and aCommandRegistryregistration. Transient fields (replyTo,sender) do not survive the queue. Previously anysendCommand/scheduleTimeoutagainst duraq ≥ 2 threwPayloadCodecException. -
SagaCoordinatoris now started and stopped. The constructor isSagaCoordinator(queueManager, actorSystem, {retryPolicy, pollInterval, onError})— the unused event store parameter is gone. Callawait coordinator.start()andawait coordinator.stop().processSagaCommands()/processSagaTimeouts()still exist as single-entry steps and now returnFuture<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'sRetryPolicyand then dead-lettered. -
Give saga states a
toMap().Saga.getSnapshotState()returnssagaStateand lets the event store serialize it (viatoMap()when present);onSagaStateRestoredreceives the decodedMap<String, dynamic>. Previously the state was CBOR-encoded twice and the restore path never matched (it checked forList<int>and receivedList<dynamic>), so saga snapshots were silently ignored on recovery. -
preStartis nowFuture<void>onPersistentActorandAggregateRoot. Subclasses that override it must declareFuture<void> preStart() asyncandawait super.preStart(). Dactor 1.3 awaitspreStartbefore delivering any message, so recovery now completes before the first message is processed andspawn()returns a recovered actor. Messages sent during recovery wait in the mailbox. A failed recovery is reported throughonRecoveryFailure/recoveryCompleteand no longer surfaces as an unhandled asynchronous error. -
IsarEventStore.create()opensallSchemas(event store plus projection checkpoints), so the returned store'sisarcan backProjectionActor. Existing databases gain the collection on open. The newdownloadIsarCoreflag lets production builds refuse to fetch the native library at runtime. -
Event.toMap(),Command.toMap()andState.toMap()are no longer@protected— serializers call them by design.Event.isValid()andEvent.getValidationErrors()likewise. -
The placeholder
Awesomeclass andeventador_base.dartexport are removed. -
DateTimefields are now tagged in CBOR; untagged strings stay strings. 2.x turned every stored string that looked like an ISO-8601 timestamp into aDateTimeon read. 3.0.0 writesDateTimevalues with CBOR tag 0 and decodes only tagged values as dates, so aStringfield that happens to contain a timestamp is returned as aString. Rows written by 2.x hold untagged dates and therefore decode asString;fromMapcode that follows the README's "acceptStringorDateTime" pattern is unaffected. If yours castsas DateTimeand you cannot change it, setCborSerializer.parseUntaggedDates = trueat startup to keep the 2.x heuristic for old rows.Event.toMap()/Command.toMap()/State.toMap()still write their owntimestamp/lastModifiedas ISO-8601 strings, as before. -
SnapshotManager.registerActorreturnsFuture<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.preStartawaits it; code that registers actors by hand should too. -
SnapshotConfig.enableCompressionandcompressionThresholdare removed, along withSnapshotConfig.shouldCompressand the identityCborSerializer.compress/decompress. None of them ever did anything; snapshots were always stored uncompressed and still are. Drop the arguments from yourSnapshotConfig(...)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 callsinitializeIsarCore(download: true)), is no longer tracked;*.dylib,*.soand*.dllare ignored.- The
docs/folder is nowdoc/, the pub layout convention, and generated.g.dart/.mocks.dartfiles 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 testfailed most runs with "file does not start with MH_MAGIC". Tests now callinitIsarCoreForTests()(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— useProjectionActor, which is now the single projection runtime. The manager's stream listener wasasyncand unserialized, so ahandle()that awaited could overlap with the next event's handler and update the read model out of order.ProjectionActorprocesses events in mailbox order. The manager now pauses its subscription whilehandle()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.
SagaCommandEnvelopeandSagaTimeoutnow shipQueueCodecs (SagaCommandEnvelope.codec,SagaTimeout.codec), used bySagaandSagaCoordinator. - Timeout ids collided across sagas. Timeout queue entries are keyed by
'$sagaId/$timeoutId'(SagaTimeout.entryId) and re-scheduling replaces the pending entry. AddedSaga.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 uniqueeventIdindex). Both are now UUID v4. - Live event streams could skip events.
allEventsWithSequenceforwarded 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,eventsByTagandeventsByPersistenceIdshare the same path. persistEventsignored automatic snapshots and only guarded againstisRecovering, not an incomplete recovery.AggregateRootpersists exclusively throughpersistEvents, so aggregates never triggered count-based snapshots.ProjectionActorno 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 entersProjectionStatus.errorwithout subscribing, callsprojection.onError, andResumeProjectionretries the read. Checkpoint write failures are still non-fatal but are now reported viaonErrorandProjectionInfo.lastError.loadCheckpoint/persistCheckpointare@protectedextension points. The deprecatedProjectionManagerthrowsProjectionExceptionfromregisterProjection/start/resumeProjectionin the same case.ProjectionActor.RebuildProjectionnow 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.ProjectionManagerrecorded one processed event per checkpoint write instead of per event.CommandHandlervalidated generated events againstValidatableState(never true for an event); it now usesValidatableEvent.- CBOR deserialization no longer guesses
DateTimefrom a string's shape.CborSerializerturned any string matchingyyyy-MM-ddTHH:mm:ssinto aDateTime, so a string field holding a timestamp came back as the wrong type andas Stringcasts infromMapthrew. Dates are now written with CBOR tag 0 (CborDateTimeString) and only tagged values decode asDateTime; the round trip is exact (UTC marker and microseconds preserved).Event.toCbor/EventRegistry.fromCborandCommand.toCbor/CommandRegistry.fromCborpreviously had their own copies of the value conversion that stringifiedDateTimewithtoString()and never decoded it; all three paths now share one codec. The old heuristic is available asCborSerializer.parseUntaggedDatesfor 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.registerActorreset 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 botheventCountThresholdandminTimeBetweenSnapshotsare measured across runs. The time rule used to answer "yes" when no snapshot existed, which meant a fresh actor with atimeThresholdsnapshotted on its very first event; it is now measured from registration until the first snapshot. TheCallbackSnapshotManagertimer no longer writes a snapshot when no event has been persisted since the last one. - The persisted
ProjectionCheckpointrow now records status and last error. ItsstatusandlastErrorcolumns were never written by either runtime (the row always readrunning,null), so storage could not tell an operator whether a projection was paused, stopped or failing.ProjectionActornow writes the row on every status transition (runningon start and resume,paused,errorwhen the checkpoint cannot be read,stoppedonStopProjectionand on an abnormal stop) and recordslastErrorwhenhandle()throws, the event stream fails, or a checkpoint write fails.lastErroris cleared by a successful (re)start or rebuild.StopProjectionalways writes the row, not only when events were pending. The deprecatedProjectionManagermirrors status on pause, resume and stop. New@protectedProjectionActor.persistStatusextension point alongsidepersistCheckpoint. CallbackSnapshotManagertook 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:shouldCreateSnapshotanswersfalseandcreateSnapshotreturnsfalsewhile one is in flight. The baseSnapshotManagerno longer schedules no-op triggers or runs a timer it cannot act on; its thresholds are applied when the actor asks.SnapshotConfig.maxSnapshotSizeis enforced and snapshot statistics measure real sizes. The manager serializes the state withCborSerializerbefore saving; a state larger thanmaxSnapshotSizeis refused (createSnapshotreturnsfalse, counted inSnapshotStats.snapshotsRejected) andtotalSnapshotSize/averageSnapshotSizereport the serialized size instead of a constant per-entry guess.ActorRefs in metadata no longer reach storage or break replay.Event.replyToreadmetadata['replyTo'] as ActorRef?, but metadata is persisted with the event, so a ref placed there was stored as itstoString()and the cast threw on every replay of that event. The same applied toCommand, whose constructor readreplyTofrom metadata when a command was rebuilt from a saga queue.ActorRefvalues are now dropped from metadata bytoMap()and by the event store's metadata column (Event.persistableMetadata/Command.persistableMetadata), andreplyToreturnsnullinstead of throwing when the stored value is not a ref. In-memory use is unchanged: a ref in metadata is still returned byreplyTountil 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,CommandandStatenow exposeString get typeName(default: the class name, so existing journals are unaffected). It is whattoMap()['type'], theeventType/stateTypeenvelope columns, the three registries andEventMigrationManageruse. Override it with a name that never changes; for types already stored under a class name,EventRegistry.registerAlias(also onCommandRegistry/StateRegistry, or thealiases:argument toregister) keeps the old rows readable.
Performance #
EventEnvelopehas a composite(persistenceId, sequenceNumber)index.getEvents,getHighestSequenceNumberandeventsByPersistenceIdread 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
projectionIdindex (where()), not a full-collectionfilter(). - Live subscribers no longer run an Isar query per event.
Added #
AggregateRootacceptssnapshotManagerandmaxProcessedCommandIdsand forwards them toPersistentActor; aggregates previously had no way to opt into automatic snapshots or command deduplication.sagaCommandQueueName/sagaTimeoutQueueNameconstants;Saga.commandQueue/Saga.timeoutQueueaccessors.SagaDeliveryException, thrown by the coordinator when a target actor is not registered (drives DuraQ retries / dead-lettering).SnapshotState.registeredAt.ProjectionActor.persistStatus(@protected).SnapshotStats.snapshotsRejectedandSnapshotStats.addRejection().Event.persistableMetadataandCommand.persistableMetadata: the metadata as written to storage, with transientActorRefentries removed.CborSerializer.parseUntaggedDatescompatibility flag.Event.typeName,Command.typeName,State.typeName;registerAlias/resolveand analiases:argument toregisteronEventRegistry,CommandRegistryandStateRegistry;CommandRegistry.isRegistered.
2.2.0 #
Added #
-
Stream<Event> get appliedEventsonProjectionActor— a broadcast stream of events the projection has successfully applied. Each event is emitted strictly afterprojection.handle(event)returnstrueand 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
appliedEventsdoes not.Events for which
handle()returnsfalse(not handled) or throws are not emitted on this stream. The stream signalsonDonewhen the actor processesStopProjection.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>— adactor-based actor wrapper around aProjection<T>, modelled after Apache Pekko / Akka'sProjectionBehavior. The actor owns the event-stream subscription, drivesprojection.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 canaskthe 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 withAwaitFailed(reason: 'timeout'); stopping the actor while awaiters are pending responds withAwaitFailed(reason: 'stopped'). -
Lifecycle message protocol —
GetProjectionInfo,PauseProjection,ResumeProjection,RebuildProjection, andStopProjection(Pekko'sProjectionBehavior.Stopequivalent — flushes the checkpoint, repliesStoppedAck, then terminates the actor).
Notes #
- Additive.
ProjectionManageris unchanged and continues to work for callers that don't need request/response semantics. Existing callers do not need to migrate.ProjectionActorandProjectionManagercan coexist in the same application. - See
doc/projection-actor-proposal.mdfor the full design rationale, Pekko mapping, and invariants.
2.0.0 #
Breaking changes #
- Removed deprecated saga state persistence.
Saga.saveSagaState()andSaga.loadSagaState()(deprecated in 1.0.0) are gone. Use the standardPersistentActorsnapshot mechanism instead: overridegetSnapshotState()andonSnapshot()(or implementonSagaStateRestored()onSaga, which the defaultonSnapshot()calls). CallcreateSnapshot()to persist. - Removed
EventStore.saveSagaState/loadSagaStatefrom the interface andIsarEventStore. CustomEventStoreimplementations no longer need to provide these methods. - Removed
SagaStateEnvelopeand its Isar collection. The schema is no longer registered byIsarEventStore.requiredSchemas. Existing databases withSagaStateEnveloperows 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
_commandLockfromAggregateRoot. Command serialization is already guaranteed byPersistentActor._handleCommand.
1.0.0 #
Initial pub.dev release.
- Event sourcing framework:
AggregateRoot,Command,Event,Statepatterns - 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
EventRegistryfor type-safe event deserialization across system restartsProjectionManagerfor automatic event streaming to registered projections- Built on Dactor actor model and DuraQ durable queuing