From 16b571f117d6a9b3c2cdaaee2cfbe00f8a79ff95 Mon Sep 17 00:00:00 2001 From: Patrick Lee Scott Date: Thu, 10 Sep 2026 17:22:42 -0500 Subject: [PATCH] fix!: reject incompatible aggregate event replay versions Fence generated handlers by the declared event schema version after upcasting, before decoding or invoking domain code. BREAKING CHANGE: aggregate! registrations for versioned handlers must declare version = N. Unsupported older and future event versions now fail replay even when their payload layout matches. Refs: incidents/aggregate-replay-event-version --- README.md | 16 +- distributed_macros/src/aggregate.rs | 46 ++++- distributed_macros/src/sourced.rs | 11 + tests/replay_event_version.rs | 305 ++++++++++++++++++++++++++++ tests/upcasting/aggregate.rs | 4 +- 5 files changed, 375 insertions(+), 7 deletions(-) create mode 100644 tests/replay_event_version.rs diff --git a/README.md b/README.md index eed9d11de..940e98238 100644 --- a/README.md +++ b/README.md @@ -1054,13 +1054,27 @@ fn upcast_initialized_v1_v2((id, task): InitV1) -> InitV2 { } aggregate!(Todo, entity { - "initialized"(id, task, priority) => initialize, + "initialized"(id, task, priority), version = 2 => initialize, "completed"() => complete(), } upcasters [ ("initialized", 1 => 2, InitV1 => InitV2, upcast_initialized_v1_v2), ]); ``` +Replay requires each event's schema version to match its registered handler +after upcasting. `#[sourced]` reads this version from `#[event(..., version = N)]`; +`aggregate!` declares it on the registration as above and must match the +corresponding `#[digest(..., version = N)]`. Omitted versions default to 1. +Older events need an explicit upcaster chain to the registered version; future +versions are rejected. Matching payload layouts do not bypass this check. +The check runs before payload decoding or handler invocation, including when +loading only the event tail after a snapshot in a native repository or cell. + +**Breaking change:** an `aggregate!` registration for a versioned handler must +now declare its current version. Histories with unsupported versions fail replay +instead of being interpreted using the current payload shape. Stored events are +unchanged; add the appropriate upcasters to read supported older versions. + ## Event Metadata Metadata lets you attach cross-cutting context — correlation IDs, causation IDs, user context, trace spans — to events without changing your domain model. diff --git a/distributed_macros/src/aggregate.rs b/distributed_macros/src/aggregate.rs index a3101b14f..e68d4b607 100644 --- a/distributed_macros/src/aggregate.rs +++ b/distributed_macros/src/aggregate.rs @@ -11,10 +11,10 @@ use syn::{ /// /// Both entry points produce a byte-identical impl: same associated /// `ReplayError = String`, same `entity`/`entity_mut`/`replay_event` bodies, and -/// the same optional `aggregate_type` and upcasters methods. Only the replay -/// match arms differ in how they are built upstream, so this helper takes them -/// (already rendered) along with the type name and entity field. Keeping one -/// emitter prevents the replay semantics of the two macros from drifting. +/// the same optional `aggregate_type` and upcasters methods. The version and +/// replay match arms are built upstream, so this helper takes them (already +/// rendered) along with the type name and entity field. Keeping the version +/// fence here prevents the replay semantics of the two macros from drifting. /// /// It emits only the `impl` block; callers still place `#upcaster_wrappers` /// (the free upcaster fns) where they already do. @@ -22,6 +22,7 @@ pub(crate) fn aggregate_impl_tokens( type_name: &Ident, entity_field: &Ident, aggregate_type_method: &Option, + version_arms: &[TokenStream2], replay_arms: &[TokenStream2], upcasters_method: &TokenStream2, ) -> TokenStream2 { @@ -43,6 +44,18 @@ pub(crate) fn aggregate_impl_tokens( &mut self, event: &distributed::EventRecord, ) -> Result<(), Self::ReplayError> { + // Hydration runs upcasters first. Payload compatibility alone + // cannot establish that an event has the handler's semantics. + let expected_version: u64 = match event.event_name.as_str() { + #(#version_arms)* + _ => return Err(format!("Unknown event: {}", event.event_name)), + }; + if event.event_version != expected_version { + return Err(format!( + "Unsupported event version for {}: expected {}, got {}", + event.event_name, expected_version, event.event_version, + )); + } match event.event_name.as_str() { #(#replay_arms)* _ => return Err(format!("Unknown event: {}", event.event_name)), @@ -178,6 +191,15 @@ pub(crate) fn expand_aggregate(input: TokenStream2) -> syn::Result let agg_name = &input.agg_name; let entity_field = &input.entity_field; + let version_arms: Vec<_> = input + .events + .iter() + .map(|event| { + let name = &event.event_name; + let version = &event.version; + quote! { #name => #version, } + }) + .collect(); // Generate replay match arms - deserialize and call method directly let replay_arms: Vec<_> = input @@ -249,6 +271,7 @@ pub(crate) fn expand_aggregate(input: TokenStream2) -> syn::Result agg_name, entity_field, &aggregate_type_method, + &version_arms, &replay_arms, &upcasters_method, ); @@ -327,6 +350,7 @@ struct AggregateInput { struct EventDef { event_name: LitStr, + version: syn::LitInt, args: Vec, method_name: Ident, method_args: Option>, // None = use event args, Some([]) = no args, Some([x,y]) = specific args @@ -382,6 +406,19 @@ impl Parse for AggregateInput { args_content.parse_terminated(Ident::parse, Token![,])?; let args: Vec = args.into_iter().collect(); + // `"renamed"(name), version = 2 => rename`; omitted versions + // have the same v1 default as #[digest] and #[event]. + let version = if content.peek(Token![,]) { + content.parse::()?; + let keyword: Ident = content.parse()?; + if keyword != "version" { + return Err(syn::Error::new(keyword.span(), "expected `version`")); + } + content.parse::()?; + content.parse::()? + } else { + syn::LitInt::new("1", event_name.span()) + }; content.parse::]>()?; let method_name: Ident = content.parse()?; @@ -398,6 +435,7 @@ impl Parse for AggregateInput { events.push(EventDef { event_name, + version, args, method_name, method_args, diff --git a/distributed_macros/src/sourced.rs b/distributed_macros/src/sourced.rs index 90caa3434..b96f63d57 100644 --- a/distributed_macros/src/sourced.rs +++ b/distributed_macros/src/sourced.rs @@ -212,6 +212,7 @@ fn find_and_remove_event_attr( struct EventMethodInfo { event_name: LitStr, + version: syn::LitInt, method_name: Ident, params: Vec<(Ident, syn::Type)>, /// Present when this recorder has `domain` and therefore a generated @@ -833,6 +834,7 @@ pub(crate) fn expand_sourced(attr: TokenStream2, item: TokenStream2) -> syn::Res }; event_methods.push(EventMethodInfo { + version: event_version(event_attr.version.as_ref()), event_name: event_attr.event_name, method_name: method.sig.ident.clone(), params, @@ -977,6 +979,14 @@ pub(crate) fn expand_sourced(attr: TokenStream2, item: TokenStream2) -> syn::Res // Generate impl Aggregate let entity_field = &args.entity_field; + let version_arms: Vec<_> = event_methods + .iter() + .map(|event| { + let name = &event.event_name; + let version = &event.version; + quote! { #name => #version, } + }) + .collect(); let replay_arms: Vec<_> = event_methods .iter() .map(|e| { @@ -1023,6 +1033,7 @@ pub(crate) fn expand_sourced(attr: TokenStream2, item: TokenStream2) -> syn::Res &struct_name, entity_field, &aggregate_type_method, + &version_arms, &replay_arms, &upcasters_method, ); diff --git a/tests/replay_event_version.rs b/tests/replay_event_version.rs new file mode 100644 index 000000000..8cf7b6452 --- /dev/null +++ b/tests/replay_event_version.rs @@ -0,0 +1,305 @@ +use distributed::{ + hydrate, hydrate_from_snapshot, Aggregate, AggregateRepository, CommitBatch, Entity, + EventRecord, RepositoryError, SnapshotRecord, SnapshotStore, Snapshottable, StreamIdentity, + StreamWrite, TransactionalCommit, +}; +use serde::{Deserialize, Serialize}; + +#[derive(Default, Serialize, Deserialize, distributed::Snapshot)] +struct Versioned { + entity: Entity, + value: String, + calls: usize, +} + +#[distributed::sourced(entity)] +impl Versioned { + #[event("renamed", version = 2)] + fn rename(&mut self, value: String) { + self.value = value; + self.calls += 1; + } + + #[event("cleared", version = 2)] + fn clear(&mut self) { + self.value.clear(); + self.calls += 1; + } +} + +#[derive(Default, Serialize, Deserialize, distributed::Snapshot)] +struct Registered { + entity: Entity, + value: String, +} + +impl Registered { + #[distributed::digest("renamed", version = 2)] + fn rename(&mut self, value: String) { + self.value = value; + } + + #[distributed::digest("cleared")] + fn clear(&mut self) { + self.value.clear(); + } +} + +distributed::aggregate!(Registered, entity { + "renamed"(value), version = 2 => rename, + "cleared"() => clear(), + "ignored"(payload), version = 2 => clear(), +}); + +#[test] +fn sourced_rejects_same_shaped_older_and_future_events_before_mutation() { + for version in [0, 1, 3, u64::MAX] { + for (name, payload) in [ + ( + "renamed", + bitcode::serialize(&("changed".to_owned(),)).unwrap(), + ), + ("cleared", vec![]), + ] { + let mut aggregate = Versioned::default(); + aggregate.rename("unchanged".into()).unwrap(); + let before = serde_json::to_value(&aggregate.entity).unwrap(); + let event = EventRecord::new_versioned(name, payload, 2, version); + assert!( + aggregate.replay_event(&event).is_err(), + "accepted {name} v{version}" + ); + assert_eq!(aggregate.value, "unchanged"); + assert_eq!(aggregate.calls, 1); + assert_eq!(serde_json::to_value(&aggregate.entity).unwrap(), before); + } + } +} + +#[test] +fn aggregate_macro_rejects_same_shaped_versions_before_mutation() { + for (name, expected) in [("renamed", 2), ("cleared", 1), ("ignored", 2)] { + for version in [expected - 1, expected + 1, u64::MAX] { + let mut aggregate = Registered::default(); + aggregate.rename("unchanged".into()).unwrap(); + let before = serde_json::to_value(&aggregate.entity).unwrap(); + let event = EventRecord::new_versioned( + name, + bitcode::serialize(&("changed".to_owned(),)).unwrap(), + 2, + version, + ); + assert!(aggregate.replay_event(&event).is_err()); + assert_eq!(aggregate.value, "unchanged"); + assert_eq!(serde_json::to_value(&aggregate.entity).unwrap(), before); + } + } +} + +fn rename_event(version: u64, sequence: u64) -> EventRecord { + EventRecord::new_versioned( + "renamed", + bitcode::serialize(&("changed".to_owned(),)).unwrap(), + sequence, + version, + ) +} + +#[test] +fn exact_version_replays_with_both_macros() { + let mut entity = Entity::new(); + entity.load_from_history(vec![rename_event(2, 1)]); + let sourced = hydrate::(entity.clone()).unwrap(); + let registered = hydrate::(entity).unwrap(); + assert_eq!(sourced.value, "changed"); + assert_eq!(sourced.calls, 1); + assert_eq!(registered.value, "changed"); + assert_eq!(sourced.entity.version(), 1); + assert!(sourced.entity.new_events().is_empty()); +} + +fn snapshot() -> SnapshotRecord { + let mut aggregate = A::new_empty(); + aggregate.entity_mut().set_id("item"); + SnapshotRecord::new( + A::aggregate_type(), + "item", + 10, + A::SNAPSHOT_VERSION, + bitcode::serialize(&aggregate.create_snapshot()).unwrap(), + ) +} + +fn check_snapshot_tail() { + for version in [1, 2, 3] { + // Production cell and native repositories load only the suffix after a + // cached snapshot; its prefix still counts toward the next write fence. + let mut entity = Entity::new(); + entity.set_id("item"); + let event = rename_event(version, 11); + entity.load_tail_from_history(vec![event.clone()], 10); + let result = hydrate_from_snapshot::(entity, snapshot::()); + if version == 2 { + let aggregate = result.unwrap(); + assert_eq!(aggregate.entity().version(), 11); + assert_eq!(aggregate.entity().events(), &[event]); + assert!(aggregate.entity().new_events().is_empty()); + } else { + assert!(matches!(result, Err(RepositoryError::Replay(message)) + if message == format!("Unsupported event version for renamed: expected 2, got {version}"))); + } + } +} + +#[test] +fn snapshot_prefix_tail_checks_versions_for_both_macros() { + check_snapshot_tail::(); + check_snapshot_tail::(); +} + +#[derive(Default, Serialize, Deserialize, distributed::Snapshot)] +struct Upcasted { + entity: Entity, + value: String, +} + +#[derive(Default)] +struct RegisteredUpcast { + entity: Entity, + value: String, +} + +impl RegisteredUpcast { + #[distributed::digest("renamed", version = 2)] + fn rename(&mut self, value: String) { + self.value = value; + } +} + +distributed::aggregate!(RegisteredUpcast, entity { + "renamed"(value), version = 2 => rename, +} upcasters [ + ("renamed", 1 => 2, (String,) => (String,), convert_semantics), +]); + +#[derive(Default)] +struct IncompleteUpcast { + entity: Entity, +} + +#[distributed::sourced(entity, upcasters( + ("renamed", 1 => 2, (String,) => (String,), convert_semantics), +))] +impl IncompleteUpcast { + #[event("renamed", version = 3)] + fn rename(&mut self, _value: String) { + self.entity.set_id("incompatible-event-applied"); + } +} + +#[test] +fn incomplete_same_shaped_upcast_chain_is_rejected() { + let mut entity = Entity::new(); + entity.load_from_history(vec![rename_event(1, 1)]); + assert!( + matches!(hydrate::(entity), Err(RepositoryError::Replay(message)) + if message == "Unsupported event version for renamed: expected 3, got 2") + ); +} + +fn convert_semantics((value,): (String,)) -> (String,) { + (format!("v2:{value}"),) +} + +#[distributed::sourced(entity, upcasters( + ("renamed", 1 => 2, (String,) => (String,), convert_semantics), +))] +impl Upcasted { + #[event("renamed", version = 2)] + fn rename(&mut self, value: String) { + self.value = value; + } +} + +#[test] +fn same_shaped_upcast_precedes_guard_and_preserves_stored_history() { + let event = rename_event(1, 1); + let mut entity = Entity::new(); + entity.load_from_history(vec![event.clone()]); + let aggregate = hydrate::(entity).unwrap(); + assert_eq!(aggregate.value, "v2:changed"); + assert_eq!(aggregate.entity.events(), &[event]); + + let mut entity = Entity::new(); + let event = rename_event(1, 1); + entity.load_from_history(vec![event.clone()]); + let aggregate = hydrate::(entity).unwrap(); + assert_eq!(aggregate.value, "v2:changed"); + assert_eq!(aggregate.entity.events(), &[event]); + + let tail = rename_event(1, 11); + let mut entity = Entity::new(); + entity.set_id("item"); + entity.load_tail_from_history(vec![tail.clone()], 10); + let aggregate = hydrate_from_snapshot::(entity, snapshot::()).unwrap(); + assert_eq!(aggregate.value, "v2:changed"); + assert_eq!(aggregate.entity.version(), 11); + assert_eq!(aggregate.entity.events(), &[tail]); +} + +#[tokio::test] +async fn cell_repository_rejects_incompatible_tail_without_changing_storage() { + use distributed::cell_host::CellStreamStore; + + for version in [1, 2, 3] { + let store = CellStreamStore::new(Versioned::aggregate_type(), "item").unwrap(); + let identity = StreamIdentity::new(Versioned::aggregate_type(), "item").unwrap(); + let mut entity = Entity::new(); + entity.set_id("item"); + for _ in 0..10 { + entity + .digest_v("renamed", 2, &("prefix".to_owned(),)) + .unwrap(); + } + entity + .digest_v("renamed", version, &("changed".to_owned(),)) + .unwrap(); + store + .commit_batch(CommitBatch::new(vec![StreamWrite::new( + identity.clone(), + &mut entity, + )])) + .await + .unwrap(); + store + .save_snapshot(&identity, snapshot::()) + .await + .unwrap(); + let before = store.durable_state().unwrap(); + let repository = AggregateRepository::<_, Versioned>::new(store.clone()).with_snapshots(10); + let result = repository.get("item").await; + if version == 2 { + let aggregate = result.unwrap().unwrap(); + assert_eq!(aggregate.value, "changed"); + assert_eq!(aggregate.calls, 1, "snapshot prefix must not replay"); + assert_eq!(aggregate.entity.version(), 11); + } else { + assert!(matches!(result, Err(RepositoryError::Replay(_)))); + } + assert_eq!(store.durable_state().unwrap(), before); + } +} + +#[test] +fn full_hydration_rejects_same_shaped_unsupported_versions() { + for version in [1, 3] { + let mut entity = Entity::new(); + entity.load_from_history(vec![EventRecord::new_versioned( + "renamed", + bitcode::serialize(&("changed".to_owned(),)).unwrap(), + 1, + version, + )]); + assert!(hydrate::(entity).is_err()); + } +} diff --git a/tests/upcasting/aggregate.rs b/tests/upcasting/aggregate.rs index aaa7f0c82..e939c897b 100644 --- a/tests/upcasting/aggregate.rs +++ b/tests/upcasting/aggregate.rs @@ -71,7 +71,7 @@ impl TodoV2 { } distributed::aggregate!(TodoV2, entity, aggregate_type = "Todo" { - "initialized"(id, user_id, task, priority) => initialize, + "initialized"(id, user_id, task, priority), version = 2 => initialize, "completed"() => complete(), } upcasters [ ("initialized", 1 => 2, InitializedV1 => InitializedV2, upcast_initialized_v1_v2), @@ -120,7 +120,7 @@ impl TodoV3 { } distributed::aggregate!(TodoV3, entity, aggregate_type = "Todo" { - "initialized"(id, user_id, task, priority, due_date) => initialize, + "initialized"(id, user_id, task, priority, due_date), version = 3 => initialize, "completed"() => complete(), } upcasters [ ("initialized", 1 => 2, InitializedV1 => InitializedV2, upcast_initialized_v1_v2),