Skip to content
Merged
Show file tree
Hide file tree
Changes from 5 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
23 changes: 20 additions & 3 deletions SPEC.md
Original file line number Diff line number Diff line change
Expand Up @@ -190,7 +190,8 @@ boundary.
Tracks `(aggregate_type, schema_version)` tuples in a `schema_registry` table.
On startup, the wiring layer compares the persisted version against the current
`SCHEMA_VERSION` constant and, on mismatch, clears stale snapshots and replays
projections from events. No manual database intervention.
projections from events. It then rebuilds missing snapshots (see
[Schema drift](#schema-drift)). No manual database intervention.

### `ViewBackend` (GAT)

Expand Down Expand Up @@ -249,8 +250,24 @@ On startup, `SchemaRegistry::reconcile()` compares the persisted
`SCHEMA_VERSION` constant:

- **Match**: nothing to do.
- **Mismatch**: snapshots are cleared (forces full event replay) and projection
tables are truncated (rebuilt from events on first read or via `catch_up`).
- **Mismatch**: snapshots are cleared and projection tables are rebuilt from
events.

On every startup, whether the version matched or not, `StoreBuilder::build()`
records the version and then writes a snapshot for each `Retain` aggregate whose
stream reaches `SNAPSHOT_SIZE` events and that has no snapshot. Without this, an
aggregate that receives no commands would replay its full stream on every load,
because commits only write a snapshot when they cross a `SNAPSHOT_SIZE`
boundary. The rebuild runs before the store accepts commands, one aggregate at a
time, and never replaces an existing snapshot. Rebuilding on every startup, not
only on a mismatch, also heals snapshots cleared by an earlier release. Because
the version is recorded first, a rebuild interrupted by a crash resumes on the
next startup instead of having its snapshots cleared again.
`CompactAfterSnapshot` aggregates are skipped, because the events behind their
snapshot may be gone. A stream that does not deserialize is skipped with a
warning, so one bad stream does not stop the store from starting; a read failure
still stops it. An aggregate that replays to a failed lifecycle gets no
snapshot, so a code fix to `evolve` can still heal it by replaying its events.

### Compaction

Expand Down
6 changes: 4 additions & 2 deletions crates/event-sorcery/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,8 @@
//! change, a stale `SCHEMA_VERSION` leaves snapshots and views
//! in the old shape. Bumping it is the operator's responsibility;
//! [`EventSourced::SCHEMA_VERSION`] plus startup reconciliation
//! clears version-mismatched snapshots and rebuilds views. On load,
//! clears version-mismatched snapshots, rebuilds missing snapshots
//! for retained aggregates, and rebuilds views. On load,
//! an incompatible snapshot for a [`CompactionPolicy::Retain`]
//! aggregate is ignored and the entity rebuilt from its
//! always-present event history; for a
Expand Down Expand Up @@ -177,7 +178,8 @@ pub enum CompactionPolicy {
/// - `SCHEMA_VERSION`: Bump when the entity's state, event, or
/// view schema changes. On startup, the wiring infrastructure
/// detects version mismatches and automatically clears stale
/// snapshots and replays views.
/// snapshots, rebuilds them for retained aggregates, and replays
/// views.
///
/// # Event-side methods
///
Expand Down
106 changes: 106 additions & 0 deletions crates/event-sorcery/src/sqlite_event_repository.rs
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,72 @@ impl SqliteEventRepository {
Ok(())
}

/// IDs of `A` aggregates whose stream reaches `min_sequence` and that
/// have no snapshot, in ID order.
///
/// Streams are grouped first, so the snapshot lookup runs once per
/// aggregate instead of once per event.
pub(crate) async fn aggregates_missing_snapshot<A: Aggregate>(
&self,
min_sequence: usize,
) -> Result<Vec<String>, PersistenceError> {
let min_sequence = i64::try_from(min_sequence).map_err(SqliteEventRepositoryError::from)?;

let aggregate_ids = sqlx::query_scalar(
"WITH streams AS ( \
SELECT aggregate_id, MAX(sequence) AS last_sequence FROM events \
WHERE aggregate_type = ?1 \
GROUP BY aggregate_id \
) \
SELECT streams.aggregate_id FROM streams \
WHERE streams.last_sequence >= ?2 \
AND NOT EXISTS ( \
SELECT 1 FROM snapshots \
WHERE snapshots.aggregate_type = ?1 \
AND snapshots.aggregate_id = streams.aggregate_id \
) \
ORDER BY streams.aggregate_id",
)
.bind(A::TYPE)
.bind(min_sequence)
.fetch_all(&self.pool)
.await
.map_err(SqliteEventRepositoryError::from)?;

Ok(aggregate_ids)
}

/// Stores a rebuilt snapshot unless the aggregate already has one, so a
/// rebuild never replaces a snapshot that a command wrote.
pub(crate) async fn insert_snapshot_if_absent<A: Aggregate>(
&self,
aggregate_id: &str,
last_sequence: usize,
snapshot_version: usize,
payload: Value,
) -> Result<(), PersistenceError> {
let last_sequence =
i64::try_from(last_sequence).map_err(SqliteEventRepositoryError::from)?;
let snapshot_version =
i64::try_from(snapshot_version).map_err(SqliteEventRepositoryError::from)?;

sqlx::query(
"INSERT OR IGNORE INTO snapshots \
(aggregate_type, aggregate_id, last_sequence, snapshot_version, payload, timestamp) \
VALUES (?1, ?2, ?3, ?4, ?5, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
)
.bind(A::TYPE)
.bind(aggregate_id)
.bind(last_sequence)
.bind(snapshot_version)
.bind(payload)
.execute(&self.pool)
.await
.map_err(SqliteEventRepositoryError::from)?;

Ok(())
}

/// Stream events from the `events` table for replay.
///
/// **Compaction caveat:** This only queries the `events` table.
Expand Down Expand Up @@ -564,4 +630,44 @@ mod tests {
Err(AggregateError::DeserializationError(_))
));
}

/// A snapshot rebuild must never replace a snapshot that already exists,
/// for example one a command wrote after the rebuild listed its targets.
#[tokio::test]
async fn insert_snapshot_if_absent_keeps_existing_snapshot() {
let pool = test_pool().await;
let repo = SqliteEventRepository::new(pool.clone(), CompactionPolicy::Retain);

repo.persist::<TestAggregate>(
&covering_events("agg-existing", 20),
Some((
"agg-existing".to_string(),
serde_json::json!({"events": ["committed"]}),
2,
)),
)
.await
.unwrap();

repo.insert_snapshot_if_absent::<TestAggregate>(
"agg-existing",
12,
1,
serde_json::json!({"events": ["rebuilt"]}),
)
.await
.unwrap();

let snapshot = repo
.get_snapshot::<TestAggregate>("agg-existing")
.await
.unwrap()
.expect("the committed snapshot stays");
assert_eq!(snapshot.current_sequence, 20);
assert_eq!(snapshot.current_snapshot, 2);
assert_eq!(
snapshot.aggregate,
serde_json::json!({"events": ["committed"]})
);
}
}
Loading
Loading