Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 12 additions & 5 deletions SPEC.md
Original file line number Diff line number Diff line change
Expand Up @@ -179,11 +179,18 @@ never opts in still logs and drops its update on a transient SQLite busy error,
exactly as it always has. A reactor can opt into an equivalent
retry-with-backoff for transient SQLite busy/busy-snapshot errors two ways: wrap
it in `RetryOnBusy` (retries the whole `react()` call, gated by implementing the
`IdempotentReactor` marker trait, which is a declaration that `react()` performs
solely idempotent SQLite writes with no side effect that would double-fire on
retry), or call `retry_with_backoff` / `is_retryable_sqlite_busy` directly
around just the write, leaving any prior side effects outside the retry
boundary.
`IdempotentReactor` marker trait, which is a declaration that `react()` (and any
`react_committed()` override) performs solely idempotent SQLite writes with no
side effect that would double-fire on retry), or call `retry_with_backoff` /
`is_retryable_sqlite_busy` directly around just the write, leaving any prior
side effects outside the retry boundary.

After each commit the store calls `Reactor::react_committed(event, committed)`
once per persisted event. `Committed` carries the event's `sequence` in its
aggregate's stream, so `(aggregate type, aggregate id, sequence)` is a stable
event id across restarts. The default forwards to `react`, so reactors that do
not need the sequence are unchanged. Loading an entity, rebuilding snapshots,
and catching up or rebuilding projections never dispatch to reactors.

### `SchemaRegistry`

Expand Down
176 changes: 174 additions & 2 deletions crates/event-sorcery/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -123,8 +123,9 @@ use lifecycle::Lifecycle;
pub use lifecycle::{LifecycleError, Never};
pub use projection::{Column, Projection, ProjectionError, Table};
pub use reactor::{
IdempotentReactor, RETRY_MAX_ATTEMPTS, RETRY_SCHEDULE, Reactor, RetryBaseDelay, RetryMaxDelay,
RetryOnBusy, RetrySchedule, RetryScheduleError, is_retryable_sqlite_busy, retry_with_backoff,
Committed, IdempotentReactor, RETRY_MAX_ATTEMPTS, RETRY_SCHEDULE, Reactor, RetryBaseDelay,
RetryMaxDelay, RetryOnBusy, RetrySchedule, RetryScheduleError, is_retryable_sqlite_busy,
retry_with_backoff,
};
pub use schema_registry::{ReconcileError, Reconciler, SchemaReconciliation, SchemaRegistry};
use sqlite_event_repository::SqliteEventRepository;
Expand Down Expand Up @@ -4243,4 +4244,175 @@ mod tests {
.expect("a same-aggregate send after cancellation must not deadlock on an orphaned lock");
second_result.expect("send after lock release via cancellation must succeed");
}

/// Records every event the store hands over: through `react_committed`
/// with its sequence, or through plain `react` with `None`.
#[derive(Default)]
struct CommitLog {
received: AsyncMutex<Vec<(String, String, Option<usize>)>>,
}

deps!(CommitLog, [ProjectedLedger]);

#[async_trait]
impl Reactor for CommitLog {
type Error = Never;

async fn react(
&self,
event: <Self::Dependencies as EntityList>::Event,
) -> Result<(), Self::Error> {
let (id, event) = event.into_inner();
self.received
.lock()
.await
.push((id.to_string(), event.label, None));
Ok(())
}

async fn react_committed(
&self,
event: <Self::Dependencies as EntityList>::Event,
committed: Committed,
) -> Result<(), Self::Error> {
let (id, event) = event.into_inner();
self.received.lock().await.push((
id.to_string(),
event.label,
Some(committed.sequence),
));
Ok(())
}
}

async fn append(store: &Store<ProjectedLedger>, id: u64, label: &str) {
store
.send(
&NumericId(id),
LedgerCommand::Append {
label: label.to_string(),
},
)
.await
.unwrap();
}

fn committed(id: &str, label: &str, sequence: usize) -> (String, String, Option<usize>) {
(id.to_string(), label.to_string(), Some(sequence))
}

#[tokio::test]
async fn store_hands_each_commit_to_react_committed_with_its_stored_sequence() {
let pool = projected_test_pool().await;
let reactor = Arc::new(CommitLog::default());
let (store, _projection) = StoreBuilder::<ProjectedLedger>::new(pool.clone())
.with(Arc::clone(&reactor))
.build(())
.await
.unwrap();

append(&store, 1, "a").await;
append(&store, 1, "b").await;
append(&store, 2, "c").await;

let stored: Vec<(String, i64)> = sqlx::query_as(
"SELECT aggregate_id, sequence FROM events \
WHERE aggregate_type = 'ProjectedLedger' ORDER BY rowid",
)
.fetch_all(&pool)
.await
.unwrap();

assert_eq!(
stored,
vec![
("1".to_string(), 1),
("1".to_string(), 2),
("2".to_string(), 1),
]
);
assert_eq!(
*reactor.received.lock().await,
vec![
committed("1", "a", 1),
committed("1", "b", 2),
committed("2", "c", 1),
]
);
}

/// A restart builds a new store over the same events. With its snapshot
/// deleted, building the store rebuilds the snapshot from the events, and
/// building it also catches the projection up. Rebuilding the view and
/// loading the entity from the snapshot replay again. None of that is a
/// commit, so the reactor sees only the new event, and its sequence
/// continues the stored stream.
#[tokio::test]
async fn restart_replay_and_hydration_do_not_reach_react_committed() {
let pool = projected_test_pool().await;
let first_run = Arc::new(CommitLog::default());
let (store, _projection) = StoreBuilder::<ProjectedLedger>::new(pool.clone())
.with(Arc::clone(&first_run))
.build(())
.await
.unwrap();

let labels: Vec<String> = (1..=SNAPSHOT_EVENTS)
.map(|n| format!("event-{n}"))
.collect();
for label in &labels {
append(&store, 1, label).await;
}
drop(store);

sqlx::query("DELETE FROM snapshots WHERE aggregate_type = 'ProjectedLedger'")
.execute(&pool)
.await
.unwrap();

let second_run = Arc::new(CommitLog::default());
let (store, projection) = StoreBuilder::<ProjectedLedger>::new(pool.clone())
.with(Arc::clone(&second_run))
.build(())
.await
.unwrap();

let rebuilt_snapshot_sequence: i64 = sqlx::query_scalar(
"SELECT last_sequence FROM snapshots \
WHERE aggregate_type = 'ProjectedLedger' AND aggregate_id = '1'",
)
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(rebuilt_snapshot_sequence, 11);

projection.rebuild_all().await.unwrap();
projection.catch_up().await.unwrap();
let loaded = store.load(&NumericId(1)).await.unwrap();

assert_eq!(loaded, Some(ProjectedLedger { order: labels }));
assert_eq!(*second_run.received.lock().await, vec![]);

append(&store, 1, "after-restart").await;

let first_run_sequences: Vec<Option<usize>> = first_run
.received
.lock()
.await
.iter()
.map(|(_id, _label, sequence)| *sequence)
.collect();
assert_eq!(
first_run_sequences,
(1..=SNAPSHOT_EVENTS).map(Some).collect::<Vec<_>>()
);
assert_eq!(
*second_run.received.lock().await,
vec![committed("1", "after-restart", 12)]
);
}

/// One more event than `ProjectedLedger`'s default `SNAPSHOT_SIZE` of 10,
/// so its stream is long enough for a snapshot rebuild.
const SNAPSHOT_EVENTS: usize = 11;
}
78 changes: 76 additions & 2 deletions crates/event-sorcery/src/lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ use tracing::{error, warn};

use crate::EventSourced;
use crate::dependency::HasEntity;
use crate::reactor::Reactor;
use crate::reactor::{Committed, Reactor};

/// Adapter that bridges [`EventSourced`] to cqrs-es `Aggregate`.
///
Expand Down Expand Up @@ -344,12 +344,15 @@ where
// No retry here by design -- wrap the reactor in `RetryOnBusy`
// (gated by `IdempotentReactor`) to retry on transient SQLite
// busy errors. See `docs/cqrs.md`'s "Reactors" section.
if let Err(error) = self.reactor.react(injected).await {
let committed = Committed::new(envelope.sequence);

if let Err(error) = self.reactor.react_committed(injected, committed).await {
error!(
target: "cqrs",
?error,
aggregate_id = aggregate_id,
aggregate_type = Entity::AGGREGATE_TYPE,
sequence = envelope.sequence,
"Reactor failed to handle event"
);
}
Expand All @@ -364,6 +367,7 @@ mod tests {
use cqrs_es::{Aggregate, DomainEvent, EventEnvelope, View};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};

use super::*;
Expand Down Expand Up @@ -850,4 +854,74 @@ mod tests {
RETRY_MAX_ATTEMPTS + 1
);
}

/// Records the sequence of every event dispatched to it, or `None` when
/// plain `react` receives it.
#[derive(Default)]
struct SequenceLog {
received: Mutex<Vec<(CounterEvent, Option<usize>)>>,
}

impl Dependent for SequenceLog {
type Dependencies = Cons<Counter, Nil>;
}

#[async_trait]
impl Reactor for SequenceLog {
type Error = Never;

async fn react(
&self,
event: <Self::Dependencies as EntityList>::Event,
) -> Result<(), Self::Error> {
let (_id, event) = event.into_inner();
self.received.lock().unwrap().push((event, None));
Ok(())
}

async fn react_committed(
&self,
event: <Self::Dependencies as EntityList>::Event,
committed: Committed,
) -> Result<(), Self::Error> {
let (_id, event) = event.into_inner();
self.received
.lock()
.unwrap()
.push((event, Some(committed.sequence)));
Ok(())
}
}

#[tokio::test]
async fn dispatch_hands_each_envelope_sequence_to_react_committed() {
let reactor = Arc::new(SequenceLog::default());
let bridge = ReactorBridge {
reactor: Arc::clone(&reactor),
};
let envelopes: Vec<EventEnvelope<Lifecycle<Counter>>> = vec![
EventEnvelope {
aggregate_id: "counter-1".to_string(),
sequence: 4,
payload: CounterEvent::Created { initial: 1 },
metadata: HashMap::new(),
},
EventEnvelope {
aggregate_id: "counter-1".to_string(),
sequence: 5,
payload: CounterEvent::Incremented,
metadata: HashMap::new(),
},
];

bridge.dispatch("counter-1", &envelopes).await;

assert_eq!(
*reactor.received.lock().unwrap(),
vec![
(CounterEvent::Created { initial: 1 }, Some(4)),
(CounterEvent::Incremented, Some(5)),
]
);
}
}
Loading
Loading