commit aa6325d7ec56ba14544d934d77d8171573213bc6
parent 2481d41fe0eb1e585aa54a99dbeb9d4d0d790264
Author: triesap <tyson@radroots.org>
Date: Tue, 28 Jul 2026 12:32:07 +0000
event_store: cover source capacity seals
- isolate initialization row-count authority
- compare persisted capacity as a complete snapshot
- prove append compare-and-swap rejects ignored writes
- remove the redundant target rebound comparison
Diffstat:
2 files changed, 114 insertions(+), 30 deletions(-)
diff --git a/contracts/event_store_production_sources.toml b/contracts/event_store_production_sources.toml
@@ -63,7 +63,7 @@ sha256 = "0527f9cc7d8d0bf1a4481327f8bd9549eba1bd6bbcce87d57f8bc5fbdb990bb7"
[[sources]]
path = "crates/event_store/src/source_maintenance_v1.rs"
-sha256 = "181576a5de365cf664b8a87091c30b0389ce0be90e7d1cc16fd7170342f6c2bc"
+sha256 = "6c90be0b0997ca33e5ab5cfb6f10e9e6ddc6679bcb5c8255dbfc21ed61d769e5"
[[sources]]
path = "crates/event_store/src/store.rs"
diff --git a/crates/event_store/src/source_maintenance_v1.rs b/crates/event_store/src/source_maintenance_v1.rs
@@ -211,21 +211,12 @@ pub(crate) async fn apply_source_maintenance_hook_v1(
let raw_tag_count: i64 = row.try_get("raw_tag_count")?;
let raw_high_water_seq: i64 = row.try_get("raw_high_water_seq")?;
let retained_generation_count = generation_count(row.try_get("retained_generation_count")?)?;
- if retained_generation_count > RADROOTS_EVENT_STORE_RETAINED_SOURCE_GENERATION_LIMIT_V1 {
- return Err(
- RadrootsEventStoreError::SourceGenerationHistoryLimitReached {
- current: retained_generation_count,
- limit: RADROOTS_EVENT_STORE_RETAINED_SOURCE_GENERATION_LIMIT_V1,
- },
- );
- }
- if raw_event_count != sqlite_capacity_value(capacity.raw_events, "raw_event_count")?
- || raw_tag_count != sqlite_capacity_value(capacity.raw_tags, "raw_tag_count")?
- {
- return source_capacity_drift(
- "measured raw row counts disagree with active source state".to_owned(),
- );
- }
+ validate_source_maintenance_initialization(
+ capacity,
+ raw_event_count,
+ raw_tag_count,
+ retained_generation_count,
+ )?;
let inserted = sqlx::query(
"INSERT INTO radroots_event_store_source_capacity_v1(singleton, source_generation, raw_event_count, raw_tag_count, raw_event_bytes, raw_tag_bytes, raw_high_water_seq, retained_generation_count, retained_generation_limit) VALUES (1, ?, ?, ?, ?, ?, ?, ?, ?)",
)
@@ -256,6 +247,33 @@ pub(crate) async fn apply_source_maintenance_hook_v1(
validate_source_capacity_authority_full_v1(connection).await
}
+fn validate_source_maintenance_initialization(
+ capacity: ReconciliationCapacity,
+ raw_event_count: i64,
+ raw_tag_count: i64,
+ retained_generation_count: u32,
+) -> Result<(), RadrootsEventStoreError> {
+ if retained_generation_count > RADROOTS_EVENT_STORE_RETAINED_SOURCE_GENERATION_LIMIT_V1 {
+ return Err(
+ RadrootsEventStoreError::SourceGenerationHistoryLimitReached {
+ current: retained_generation_count,
+ limit: RADROOTS_EVENT_STORE_RETAINED_SOURCE_GENERATION_LIMIT_V1,
+ },
+ );
+ }
+ if (raw_event_count, raw_tag_count)
+ != (
+ sqlite_capacity_value(capacity.raw_events, "raw_event_count")?,
+ sqlite_capacity_value(capacity.raw_tags, "raw_tag_count")?,
+ )
+ {
+ return source_capacity_drift(
+ "measured raw row counts disagree with active source state".to_owned(),
+ );
+ }
+ Ok(())
+}
+
pub(crate) async fn validate_source_capacity_authority_fast_v1(
connection: &mut SqliteConnection,
) -> Result<RadrootsEventStoreSourceCapacityV1, RadrootsEventStoreError> {
@@ -289,14 +307,23 @@ pub(crate) async fn validate_source_capacity_authority_fast_v1(
let raw_high_water_seq: i64 = row.try_get("raw_high_water_seq")?;
let generation_ordinal = generation_count(row.try_get("generation_ordinal")?)?;
let retained_generation_count = generation_count(row.try_get("retained_generation_count")?)?;
- if active_generation != capacity.source_generation
- || raw_event_count != capacity.capacity.raw_events
- || raw_tag_count != capacity.capacity.raw_tags
- || raw_high_water_seq != capacity.raw_high_water_seq
- || generation_ordinal != retained_generation_count
- || retained_generation_count != capacity.retained_generation_count
- || retained_generation_count > capacity.retained_generation_limit
- {
+ if (
+ active_generation,
+ raw_event_count,
+ raw_tag_count,
+ raw_high_water_seq,
+ generation_ordinal,
+ retained_generation_count,
+ retained_generation_count > capacity.retained_generation_limit,
+ ) != (
+ capacity.source_generation,
+ capacity.capacity.raw_events,
+ capacity.capacity.raw_tags,
+ capacity.raw_high_water_seq,
+ retained_generation_count,
+ capacity.retained_generation_count,
+ false,
+ ) {
return source_capacity_drift(
"capacity seal does not match active source state and generation history".to_owned(),
);
@@ -409,12 +436,7 @@ pub(crate) async fn bind_source_capacity_to_generation_v1(
updated.rows_affected()
));
}
- let rebound = validate_source_capacity_authority_fast_v1(connection).await?;
- if rebound.source_generation != target_generation {
- return source_capacity_drift(
- "source rebuild capacity bind did not select its target generation".to_owned(),
- );
- }
+ validate_source_capacity_authority_fast_v1(connection).await?;
Ok(true)
}
@@ -908,6 +930,39 @@ mod tests {
Err(RadrootsEventStoreError::SourceCapacityStateDrift { reason })
if reason == "fixture drift"
));
+
+ validate_source_maintenance_initialization(
+ ReconciliationCapacity::default(),
+ 0,
+ 0,
+ RADROOTS_EVENT_STORE_RETAINED_SOURCE_GENERATION_LIMIT_V1,
+ )
+ .expect("exact initialization authority");
+ assert!(matches!(
+ validate_source_maintenance_initialization(
+ ReconciliationCapacity::default(),
+ 0,
+ 0,
+ RADROOTS_EVENT_STORE_RETAINED_SOURCE_GENERATION_LIMIT_V1 + 1,
+ ),
+ Err(
+ RadrootsEventStoreError::SourceGenerationHistoryLimitReached {
+ current: 9,
+ limit: 8,
+ }
+ )
+ ));
+ for (raw_event_count, raw_tag_count) in [(1, 0), (0, 1)] {
+ assert!(matches!(
+ validate_source_maintenance_initialization(
+ ReconciliationCapacity::default(),
+ raw_event_count,
+ raw_tag_count,
+ 1,
+ ),
+ Err(RadrootsEventStoreError::SourceCapacityStateDrift { .. })
+ ));
+ }
}
#[tokio::test]
@@ -957,6 +1012,35 @@ mod tests {
let store = RadrootsEventStore::open_memory().await.expect("store");
let mut transaction = store.begin_write_transaction().await.expect("transaction");
+ sqlx::query("DROP TRIGGER radroots_event_store_source_capacity_update_guard")
+ .execute(&mut *transaction)
+ .await
+ .expect("drop update guard");
+ sqlx::query(
+ "CREATE TEMP TRIGGER ignore_capacity_advance BEFORE UPDATE ON radroots_event_store_source_capacity_v1 BEGIN SELECT RAISE(IGNORE); END",
+ )
+ .execute(&mut *transaction)
+ .await
+ .expect("install ignored advance");
+ assert!(matches!(
+ advance_source_capacity_after_insert_v1(
+ &mut transaction,
+ RawSourceCapacityDeltaV1 {
+ raw_events: 0,
+ raw_tags: 0,
+ raw_event_bytes: 0,
+ raw_tag_bytes: 0,
+ },
+ 0,
+ )
+ .await,
+ Err(RadrootsEventStoreError::SourceCapacityStateDrift { reason })
+ if reason.contains("append authority compare-and-swap affected 0 rows")
+ ));
+ transaction.rollback().await.expect("rollback fixture");
+
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let mut transaction = store.begin_write_transaction().await.expect("transaction");
let current = read_source_capacity_v1(&mut transaction)
.await
.expect("current capacity");