lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit 8dbf27459be1729709a0ea36bba7470e90479ca3
parent 087c17e1f72d865ce58d5c881fc05a97bf3ec110
Author: triesap <tyson@radroots.org>
Date:   Tue, 22 Sep 2026 17:10:12 +0000

storage: preserve capacity through sync and startup

- Preserve typed capacity through sync and startup adapters
- Report SDK capacity without exposing native diagnostics
- Retain original signed work and unresolved delivery claims
- Verify additive API and unchanged coverage thresholds

Diffstat:
Mcontracts/api_baselines/radroots_sdk.txt | 1+
Mcontracts/api_baselines/radroots_storage_sqlite.txt | 2++
Mcontracts/api_baselines/radroots_sync.txt | 2++
Acontracts/architecture/decisions/storage_capacity_propagation.v1.json | 15+++++++++++++++
Mcontracts/architecture/deviations.toml | 20++++++++++++++++++++
Mcrates/sdk/README.md | 6++++++
Mcrates/sdk/src/error.rs | 70++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/storage_sqlite/README.md | 5+++++
Mcrates/storage_sqlite/src/backend.rs | 33++++++++++++++++++++++++---------
Mcrates/storage_sqlite/src/migration.rs | 187++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------------
Mcrates/storage_sqlite/src/migration/authored_v10.rs | 47++++++++++++++++++++++++++++-------------------
Mcrates/storage_sqlite/src/open.rs | 104++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------
Mcrates/sync/README.md | 5+++++
Mcrates/sync/src/ingest.rs | 1+
Mcrates/sync/src/policy.rs | 3+++
Mcrates/sync/src/projection.rs | 5+++++
Mcrates/sync/src/push.rs | 14+++++++++++++-
Mcrates/sync/src/status.rs | 27++++++++++++++++++++-------
Mcrates/sync/tests/engine_composition.rs | 1+
Mcrates/sync/tests/push_enqueue.rs | 8++++++++
Acrates/sync/tests/push_enqueue/capacity.rs | 142+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/sync/tests/push_enqueue/delivery_evidence.rs | 68++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
22 files changed, 650 insertions(+), 116 deletions(-)

diff --git a/contracts/api_baselines/radroots_sdk.txt b/contracts/api_baselines/radroots_sdk.txt @@ -106,6 +106,7 @@ pub radroots_sdk::error::ErrorKind::StorageCloseFailed pub radroots_sdk::error::ErrorKind::StorageInspectionFailed pub radroots_sdk::error::ErrorKind::StorageOpenFailed pub radroots_sdk::error::ErrorKind::StorageSchemaTooNew +pub radroots_sdk::error::ErrorKind::StorageSpaceInsufficient pub radroots_sdk::error::ErrorKind::StorageUnsupportedSchema impl radroots_sdk::error::ErrorKind pub const radroots_sdk::error::ErrorKind::ALL: &'static [Self] diff --git a/contracts/api_baselines/radroots_storage_sqlite.txt b/contracts/api_baselines/radroots_storage_sqlite.txt @@ -318,6 +318,7 @@ pub radroots_storage_sqlite::open::Error::SchemaTooOld::minimum: u32 pub radroots_storage_sqlite::open::Error::SourceGenerationMismatch pub radroots_storage_sqlite::open::Error::SourceGenerationRequired pub radroots_storage_sqlite::open::Error::SourceGenerationUnavailable +pub radroots_storage_sqlite::open::Error::SpaceInsufficient pub radroots_storage_sqlite::open::Error::SymlinkPath(std::path::PathBuf) pub radroots_storage_sqlite::open::Error::UnexpectedFileName pub radroots_storage_sqlite::open::Error::UnexpectedFileName::expected: &'static str @@ -459,6 +460,7 @@ pub radroots_storage_sqlite::Error::SchemaTooOld::minimum: u32 pub radroots_storage_sqlite::Error::SourceGenerationMismatch pub radroots_storage_sqlite::Error::SourceGenerationRequired pub radroots_storage_sqlite::Error::SourceGenerationUnavailable +pub radroots_storage_sqlite::Error::SpaceInsufficient pub radroots_storage_sqlite::Error::SymlinkPath(std::path::PathBuf) pub radroots_storage_sqlite::Error::UnexpectedFileName pub radroots_storage_sqlite::Error::UnexpectedFileName::expected: &'static str diff --git a/contracts/api_baselines/radroots_sync.txt b/contracts/api_baselines/radroots_sync.txt @@ -68,6 +68,7 @@ pub radroots_sync::policy::Error::SigningCancelled pub radroots_sync::policy::Error::SigningIndeterminate pub radroots_sync::policy::Error::StorageConflict pub radroots_sync::policy::Error::StorageFailed +pub radroots_sync::policy::Error::StorageSpaceInsufficient pub radroots_sync::policy::Error::VerificationFailed pub radroots_sync::policy::Error::WorkClaimConflict impl core::error::Error for radroots_sync::policy::Error @@ -273,6 +274,7 @@ pub radroots_sync::Error::SigningCancelled pub radroots_sync::Error::SigningIndeterminate pub radroots_sync::Error::StorageConflict pub radroots_sync::Error::StorageFailed +pub radroots_sync::Error::StorageSpaceInsufficient pub radroots_sync::Error::VerificationFailed pub radroots_sync::Error::WorkClaimConflict impl core::error::Error for radroots_sync::policy::Error diff --git a/contracts/architecture/decisions/storage_capacity_propagation.v1.json b/contracts/architecture/decisions/storage_capacity_propagation.v1.json @@ -0,0 +1,15 @@ +{ + "schema": "radroots.storage-capacity-propagation.v1", + "status": "approved", + "scope": "Preserve the canonical capacity distinction through Sync orchestration and SQLite startup into SDK reporting", + "sync": "StorageSpaceInsufficient preserves Storage::SpaceInsufficient. Context-specific conflicts and all other fallback errors keep their existing classification.", + "startup": "SQLite SpaceInsufficient classifies numeric primary or extended SQLITE_FULL and typed I/O StorageFull or QuotaExceeded. SDK StorageSpaceInsufficient reports the existing protocol storage_space_insufficient code without exposing raw diagnostics.", + "uncertainty": "Capacity does not establish rollback, absence of signer or remote effects, or safe automatic retry. Preserve original identities, signed bytes, claims and unresolved receipts.", + "boundaries": "No schema, SQL policy, transaction order, durability setting, dependency, resource limit, retry engine or cleanup policy changes. Backup and restore algorithms and their independent capability contracts remain separate.", + "verification": [ + "real orchestration with capacity faults before and after durable effects", + "bounded actual SQLite migration exhaustion and original data retention", + "primary and extended codes, typed I/O, redaction and unchanged fallbacks", + "additive API, unchanged coverage thresholds, portable and full workspace qualification" + ] +} diff --git a/contracts/architecture/deviations.toml b/contracts/architecture/deviations.toml @@ -2,6 +2,26 @@ schema_version = 1 architecture_id = "radroots.crates.release.v1" [[deviation]] +id = "RCRV1-DEV-023" +date = "2026-09-22" +status = "closed" +approval = "Explicit user authorization covers necessary owning-repository repairs, verified checkpoints and non-force publication." +affected_steps = ["163", "173", "178"] +spec_anchors = ["contracts/crates/release_v1/radroots_crates_release_v1.toml#package.radroots_storage_sqlite", "contracts/crates/release_v1/radroots_crates_release_v1.toml#package.radroots_sync", "contracts/crates/release_v1/radroots_crates_release_v1.toml#package.radroots_sdk"] +source_evidence = ["Sync adapters collapse the canonical typed capacity error before the host can classify signing or delivery persistence failures.", "SQLite startup and migration adapters similarly discard SQL capacity before SDK report conversion."] +replacement_action = "Preserve capacity through existing Sync and startup adapters under storage_capacity_propagation.v1.json, without changing ownership or effect semantics." +verification = ["Exercise actual orchestration before and after durable effects with capacity faults and preserved original requests and receipts.", "Verify startup SQL and I/O classification, actual bounded migration failure, unchanged fallback errors and redacted SDK reports.", "Qualify additive API, unchanged coverage gates, portable profiles, full workspace and release preflight."] +unresolved_risk = "Typed capacity classification, actual owner regressions, additive API, unchanged coverage gates and complete workspace qualification passed. Consumers still reconcile prior effects; classification grants no eviction or automatic retry authority." +normative_architecture_change = false +adr_required = false +closure_evidence = [ + "Sync preserves typed capacity without changing original operation, claim, signer or delivery evidence. Actual before/after faults prove no false success or extra signer/network calls.", + "Startup and migration preserve SQL capacity and typed I/O distinctions through the redacted SDK report. Actual bounded SQLite migration failure retains original version and data and retries the pending suffix.", + "Three additive variants in non-exhaustive public enums, no removals or umbrella API changes. No schema, SQL policy, transaction, durability, dependency or limit changes.", + "All45 coverage gates pass unchanged thresholds. Three affected packages have fresh coverage;42 unchanged package-source reports retain provenance. Complete workspace, minimal profiles, portable, preflight and explicit native/WASM generators pass.", +] + +[[deviation]] id = "RCRV1-DEV-022" date = "2026-09-22" status = "closed" diff --git a/crates/sdk/README.md b/crates/sdk/README.md @@ -17,6 +17,12 @@ applications that need to compose storage, signing, transport, or sync capabilities directly. Ordinary Rust applications should use the curated `radroots` facade. +Persistent startup capacity failures use `ErrorKind::StorageSpaceInsufficient` +and the existing `storage_space_insufficient` protocol report. Native sources +remain available through the error chain; display, debug and protocol reports +exclude their paths and diagnostic text. Hosts retain existing stores and +reconcile prior effects before retrying. + ## Getting started The default `memory` feature provides an inert, in-process backend. The caller diff --git a/crates/sdk/src/error.rs b/crates/sdk/src/error.rs @@ -106,6 +106,13 @@ error_catalog! { message: "SDK persistent storage open failed", safe_detail_keys: [] }, + StorageSpaceInsufficient => { + code: StorageSpaceInsufficient, + operation: None, + capability: Some(CapabilityId::PERSISTENT_STORAGE), + message: "SDK persistent storage space is insufficient", + safe_detail_keys: [] + }, StorageBusy => { code: DatabaseBusy, operation: None, @@ -252,6 +259,17 @@ impl Error { use radroots_storage_sqlite::Error as SqliteError; let kind = match &source { + SqliteError::SpaceInsufficient => ErrorKind::StorageSpaceInsufficient, + SqliteError::Inspect { source, .. } + | SqliteError::WriterLockOpen { source, .. } + | SqliteError::WriterLockFailed { source, .. } + if matches!( + source.kind(), + std::io::ErrorKind::StorageFull | std::io::ErrorKind::QuotaExceeded + ) => + { + ErrorKind::StorageSpaceInsufficient + } SqliteError::WriterAlreadyActive { .. } => ErrorKind::StorageBusy, SqliteError::SchemaTooNew { .. } => ErrorKind::StorageSchemaTooNew, SqliteError::SchemaTooOld { .. } | SqliteError::SchemaMigrationRequired { .. } => { @@ -376,6 +394,58 @@ mod tests { use super::*; + #[cfg(feature = "sqlite")] + #[test] + fn startup_capacity_report_preserves_source_but_never_exposes_path() { + use radroots_storage_sqlite::Error as Sqlite; + let mut sources = vec![Sqlite::SpaceInsufficient]; + for kind in [ + std::io::ErrorKind::StorageFull, + std::io::ErrorKind::QuotaExceeded, + ] { + sources.extend([ + Sqlite::Inspect { + path: "/private/secret".into(), + source: std::io::Error::new(kind, "private detail"), + }, + Sqlite::WriterLockOpen { + path: "/private/secret".into(), + source: std::io::Error::new(kind, "private detail"), + }, + Sqlite::WriterLockFailed { + path: "/private/secret".into(), + source: std::io::Error::new(kind, "private detail"), + }, + ]); + } + for source in sources { + let error = Error::storage_open_failed(source); + assert!(error.source().is_some()); + assert_eq!(error.kind(), ErrorKind::StorageSpaceInsufficient); + let report = error.to_report(); + report.validate().unwrap(); + assert_eq!(report.code().as_str(), "storage_space_insufficient"); + assert_eq!( + report.message().as_str(), + "SDK persistent storage space is insufficient" + ); + assert!(!format!("{error:?}").contains("private")); + } + for kind in [ + std::io::ErrorKind::PermissionDenied, + std::io::ErrorKind::Other, + ] { + assert_eq!( + Error::storage_open_failed(Sqlite::WriterLockOpen { + path: "/private/secret".into(), + source: kind.into(), + }) + .kind(), + ErrorKind::StorageOpenFailed + ); + } + } + #[test] fn catalog_is_exhaustive_unique_and_protocol_valid() { assert_eq!(CATALOG.len(), ErrorKind::ALL.len()); diff --git a/crates/storage_sqlite/README.md b/crates/storage_sqlite/README.md @@ -2,6 +2,11 @@ SQLite storage backend for Radroots. +Startup and migration preserve typed `SpaceInsufficient` for numeric SQLite +capacity failures and typed capacity I/O failures. Existing corruption and +validation errors remain distinct. Capacity does not prove rollback: inspect +the original store before retrying, and retain pending or ambiguous work. + Projection document inventory uses the existing canonical projection table. Each page releases its read snapshot before returning a continuation. Metadata is bounded before decoding; the 16 MiB payload budget is charged before each diff --git a/crates/storage_sqlite/src/backend.rs b/crates/storage_sqlite/src/backend.rs @@ -3,21 +3,36 @@ use radroots_storage::Error; pub(crate) fn map_backend(source: sqlx::Error) -> Error { - let full = match &source { + if is_capacity(&source) { + Error::SpaceInsufficient + } else { + Error::BackendUnavailable + } +} + +pub(crate) fn is_capacity(source: &sqlx::Error) -> bool { + match source { sqlx::Error::Database(error) => error .code() .and_then(|code| code.parse::<u32>().ok()) .is_some_and(|code| code & 0xff == 13), - sqlx::Error::Io(error) => matches!( - error.kind(), - std::io::ErrorKind::StorageFull | std::io::ErrorKind::QuotaExceeded - ), + sqlx::Error::Io(error) => is_io_capacity(error), _ => false, - }; - if full { - Error::SpaceInsufficient + } +} + +pub(crate) fn is_io_capacity(source: &std::io::Error) -> bool { + matches!( + source.kind(), + std::io::ErrorKind::StorageFull | std::io::ErrorKind::QuotaExceeded + ) +} + +pub(crate) fn startup_error(source: &sqlx::Error, fallback: crate::Error) -> crate::Error { + if is_capacity(source) { + crate::Error::SpaceInsufficient } else { - Error::BackendUnavailable + fallback } } diff --git a/crates/storage_sqlite/src/migration.rs b/crates/storage_sqlite/src/migration.rs @@ -20,8 +20,13 @@ pub async fn preflight_authored_v10(paths: &crate::Paths) -> Result<AuthoredV10P sqlx::raw_sql("PRAGMA query_only = ON") .execute(&mut connection) .await - .map_err(|_| Error::SchemaMetadataUnavailable { - database: RUNTIME_DATABASE, + .map_err(|source| { + crate::backend::startup_error( + &source, + Error::SchemaMetadataUnavailable { + database: RUNTIME_DATABASE, + }, + ) })?; let current = metadata(&mut connection, RUNTIME_DATABASE).await?; if current.application_id != RUNTIME_APPLICATION_ID { @@ -53,12 +58,14 @@ pub async fn preflight_authored_v10(paths: &crate::Paths) -> Result<AuthoredV10P ) .await?; let report = authored_v10::inspect(&mut connection).await?.report; - connection - .close() - .await - .map_err(|_| Error::DatabaseCloseFailed { - database: RUNTIME_DATABASE, - })?; + connection.close().await.map_err(|source| { + crate::backend::startup_error( + &source, + Error::DatabaseCloseFailed { + database: RUNTIME_DATABASE, + }, + ) + })?; Ok(report) } @@ -220,12 +227,14 @@ pub(crate) async fn preflight_existing(options: &crate::OpenOptions) -> Result<( Ok::<(), Error>(()) } .await; - let closed = connection - .close() - .await - .map_err(|_| Error::DatabaseCloseFailed { - database: plan.database, - }); + let closed = connection.close().await.map_err(|source| { + crate::backend::startup_error( + &source, + Error::DatabaseCloseFailed { + database: plan.database, + }, + ) + }); inspected?; closed?; } @@ -293,9 +302,14 @@ async fn migrate( let mut transaction = connection .begin_with("BEGIN IMMEDIATE") .await - .map_err(|_| Error::SchemaMigrationFailed { - database: plan.database, - target_version: initial.version.saturating_add(1), + .map_err(|source| { + crate::backend::startup_error( + &source, + Error::SchemaMigrationFailed { + database: plan.database, + target_version: initial.version.saturating_add(1), + }, + ) })?; let transactional = match metadata(&mut transaction, plan.database).await { Ok(metadata) => metadata, @@ -312,15 +326,17 @@ async fn migrate( } if initial.version == 0 - && sqlx::raw_sql(plan.set_application_id_sql) + && let Err(source) = sqlx::raw_sql(plan.set_application_id_sql) .execute(&mut *transaction) .await - .is_err() { - let error = Error::SchemaMigrationFailed { - database: plan.database, - target_version: 1, - }; + let error = crate::backend::startup_error( + &source, + Error::SchemaMigrationFailed { + database: plan.database, + target_version: 1, + }, + ); let _rollback = transaction.rollback().await; return Err(error); } @@ -352,15 +368,14 @@ async fn migrate( database: plan.database, target_version: step.version, })?; - if sqlx::raw_sql(step.sql) - .execute(&mut *transaction) - .await - .is_err() - { - let error = Error::SchemaMigrationFailed { - database: plan.database, - target_version: step.version, - }; + if let Err(source) = sqlx::raw_sql(step.sql).execute(&mut *transaction).await { + let error = crate::backend::startup_error( + &source, + Error::SchemaMigrationFailed { + database: plan.database, + target_version: step.version, + }, + ); let _rollback = transaction.rollback().await; return Err(error); } @@ -370,43 +385,46 @@ async fn migrate( let _rollback = transaction.rollback().await; return Err(error); } - if sqlx::raw_sql(version_sql) - .execute(&mut *transaction) - .await - .is_err() - { - let error = Error::SchemaMigrationFailed { - database: plan.database, - target_version: step.version, - }; + if let Err(source) = sqlx::raw_sql(version_sql).execute(&mut *transaction).await { + let error = crate::backend::startup_error( + &source, + Error::SchemaMigrationFailed { + database: plan.database, + target_version: step.version, + }, + ); let _rollback = transaction.rollback().await; return Err(error); } - if validate_exact_catalog( + if let Err(source) = validate_exact_catalog( &mut transaction, plan.database, step.version, step.owned_objects, ) .await - .is_err() { - let error = Error::SchemaMigrationFailed { - database: plan.database, - target_version: step.version, + let error = match source { + Error::SpaceInsufficient => Error::SpaceInsufficient, + _ => Error::SchemaMigrationFailed { + database: plan.database, + target_version: step.version, + }, }; let _rollback = transaction.rollback().await; return Err(error); } applied = applied.saturating_add(1); } - transaction - .commit() - .await - .map_err(|_| Error::SchemaMigrationFailed { - database: plan.database, - target_version: plan.current_version, - })?; + transaction.commit().await.map_err(|source| { + crate::backend::startup_error( + &source, + Error::SchemaMigrationFailed { + database: plan.database, + target_version: plan.current_version, + }, + ) + })?; Ok(MigrationReport { initial_version: initial.version, final_version: plan.current_version, @@ -443,7 +461,7 @@ async fn metadata( fn schema_metadata_error(source: &sqlx::Error, database: &'static str) -> Error { match crate::open::map_database_open_error(source, database) { - error @ Error::DatabaseCorrupt { .. } => error, + error @ (Error::DatabaseCorrupt { .. } | Error::SpaceInsufficient) => error, _ => Error::SchemaMetadataUnavailable { database }, } } @@ -536,7 +554,9 @@ async fn validate_exact_catalog( ) .fetch_all(&mut *connection) .await - .map_err(|_| Error::SchemaCatalogMismatch { database, version })?; + .map_err(|source| { + crate::backend::startup_error(&source, Error::SchemaCatalogMismatch { database, version }) + })?; let actual = rows .iter() .map(|row| row.get::<String, _>("name")) @@ -720,6 +740,63 @@ mod tests { } #[tokio::test] + async fn capacity_migration_retains_original_generation_and_retries_exact_pending_suffix() { + let mut connection = connection().await; + establish_runtime_version(&mut connection, 1).await; + sqlx::query("INSERT INTO radroots_runtime_source_generations (generation, state, created_at_unix_ms) VALUES (?, 'active', 10)") + .bind([7_u8; 32].as_slice()).execute(&mut connection).await.unwrap(); + let page_count: i64 = sqlx::query_scalar("PRAGMA page_count") + .fetch_one(&mut connection) + .await + .unwrap(); + let previous_limit: i64 = sqlx::query_scalar("PRAGMA max_page_count") + .fetch_one(&mut connection) + .await + .unwrap(); + // PRAGMA assignments cannot bind; the interpolated value is an i64 from SQLite. + sqlx::raw_sql(sqlx::AssertSqlSafe(format!( + "PRAGMA max_page_count = {page_count}" + ))) + .execute(&mut connection) + .await + .unwrap(); + assert!(matches!( + migrate_runtime(&mut connection, OpenMode::ReadWriteExisting).await, + Err(Error::SpaceInsufficient) + )); + assert_eq!(pragma(&mut connection, "user_version").await, 1); + let generation: Vec<u8> = sqlx::query_scalar( + "SELECT generation FROM radroots_runtime_source_generations WHERE state = 'active'", + ) + .fetch_one(&mut connection) + .await + .unwrap(); + assert_eq!(generation, [7_u8; 32]); + // The restored limit is likewise an i64, never untrusted SQL text. + sqlx::raw_sql(sqlx::AssertSqlSafe(format!( + "PRAGMA max_page_count = {previous_limit}" + ))) + .execute(&mut connection) + .await + .unwrap(); + let report = migrate_runtime(&mut connection, OpenMode::ReadWriteExisting) + .await + .unwrap(); + assert_eq!(report.initial_version(), 1); + assert_eq!(report.final_version(), runtime::CURRENT_VERSION); + assert_eq!(report.applied(), runtime::CURRENT_VERSION - 1); + assert_eq!( + sqlx::query_scalar::<_, Vec<u8>>( + "SELECT generation FROM radroots_runtime_source_generations WHERE state = 'active'" + ) + .fetch_one(&mut connection) + .await + .unwrap(), + generation + ); + } + + #[tokio::test] async fn committed_migration_reopens_as_current_after_a_lost_success_response() { let directory = tempfile::tempdir().expect("database directory"); let path = directory.path().join(RUNTIME_DATABASE); diff --git a/crates/storage_sqlite/src/migration/authored_v10.rs b/crates/storage_sqlite/src/migration/authored_v10.rs @@ -118,7 +118,7 @@ pub(crate) async fn inspect(connection: &mut SqliteConnection) -> Result<Inspect sqlx::query("SELECT event_id, signed_event FROM radroots_runtime_events ORDER BY event_id") .fetch_all(&mut *connection) .await - .map_err(|_| metadata_error())?; + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; let mut event_metadata = Vec::with_capacity(event_rows.len()); let mut invalid_event_ids = BTreeSet::new(); for row in event_rows { @@ -162,7 +162,7 @@ pub(crate) async fn inspect(connection: &mut SqliteConnection) -> Result<Inspect ) .fetch_all(&mut *connection) .await - .map_err(|_| metadata_error())?; + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; for bytes in outbox_ids { source_hasher.update((bytes.len() as u64).to_be_bytes()); source_hasher.update(bytes.as_slice()); @@ -175,6 +175,9 @@ pub(crate) async fn inspect(connection: &mut SqliteConnection) -> Result<Inspect }; match outbox::load_record(connection, item_id).await { Ok(Some(record)) => outboxes.push(record), + Err(radroots_storage::Error::SpaceInsufficient) => { + return Err(Error::SpaceInsufficient); + } Ok(None) | Err(_) => { invalid_or_unsupported = invalid_or_unsupported.saturating_add(1); } @@ -185,7 +188,7 @@ pub(crate) async fn inspect(connection: &mut SqliteConnection) -> Result<Inspect sqlx::query("SELECT * FROM radroots_runtime_journal_operations ORDER BY instance_id") .fetch_all(&mut *connection) .await - .map_err(|_| metadata_error())?; + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; let mut candidates = Vec::new(); let mut matched_outboxes = BTreeSet::new(); let mut prepared_or_recoverable = 0_u64; @@ -226,10 +229,10 @@ pub(crate) async fn inspect(connection: &mut SqliteConnection) -> Result<Inspect continue; }; matched_outboxes.insert(*record_outbox.item_id().as_bytes()); - if validate_candidate(connection, &record_outbox, event_id) - .await - .is_err() - { + if let Err(error) = validate_candidate(connection, &record_outbox, event_id).await { + if matches!(error, Error::SpaceInsufficient) { + return Err(error); + } invalid_or_unsupported = invalid_or_unsupported.saturating_add(1); push_blocked(&mut blocked_operation_ids, &operation_id); continue; @@ -301,7 +304,7 @@ pub(crate) async fn apply( .bind(metadata.event_id.as_slice()) .execute(&mut **transaction) .await - .map_err(|_| metadata_error())?; + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; if result.rows_affected() != 1 { return Err(metadata_error()); } @@ -310,11 +313,17 @@ pub(crate) async fn apply( let (operation, artifact, plan) = convert_candidate(candidate)?; authored::persist_operation(transaction, &operation) .await - .map_err(|_| metadata_error())?; + .map_err(|source| match source { + radroots_storage::Error::SpaceInsufficient => Error::SpaceInsufficient, + _ => metadata_error(), + })?; persist_imported_artifact(transaction, &artifact).await?; authored::persist_plan_v11(transaction, &plan) .await - .map_err(|_| metadata_error())?; + .map_err(|source| match source { + radroots_storage::Error::SpaceInsufficient => Error::SpaceInsufficient, + _ => metadata_error(), + })?; } let operation_count = count_transaction( transaction, @@ -340,14 +349,14 @@ pub(crate) async fn apply( let foreign_keys = sqlx::query("PRAGMA foreign_key_check") .fetch_all(&mut **transaction) .await - .map_err(|_| metadata_error())?; + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; if !foreign_keys.is_empty() { return Err(metadata_error()); } let integrity = sqlx::query_scalar::<_, String>("PRAGMA integrity_check") .fetch_all(&mut **transaction) .await - .map_err(|_| metadata_error())?; + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; if integrity.as_slice() != ["ok"] { return Err(metadata_error()); } @@ -368,7 +377,7 @@ pub(crate) async fn apply( .bind(migration_timestamp(inspected)) .execute(&mut **transaction) .await - .map_err(|_| metadata_error())?; + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; Ok(()) } @@ -460,7 +469,7 @@ async fn persist_imported_artifact( .bind(authored::encode_snapshot(artifact).map_err(|_| metadata_error())?) .execute(&mut **transaction) .await - .map_err(|_| metadata_error())?; + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; Ok(()) } @@ -600,7 +609,7 @@ async fn validate_candidate( .bind(event_id.as_bytes().as_slice()) .fetch_optional(&mut *connection) .await - .map_err(|_| metadata_error())? + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))? .ok_or_else(metadata_error)?; if stored .try_get::<Vec<u8>, _>("signed_event") @@ -633,7 +642,7 @@ async fn event_admitted_at( .bind(event_id.as_slice()) .fetch_one(&mut *connection) .await - .map_err(|_| metadata_error())?; + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; u64_from_i64(value) } @@ -652,7 +661,7 @@ async fn orphan_count(connection: &mut SqliteConnection) -> Result<u64, Error> { ) .fetch_one(&mut *connection) .await - .map_err(|_| metadata_error())?; + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; u64_from_i64(value) } @@ -673,7 +682,7 @@ async fn count(connection: &mut SqliteConnection, table: &'static str) -> Result sqlx::query_scalar::<_, i64>(query) .fetch_one(&mut *connection) .await - .map_err(|_| metadata_error())?, + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?, ) } @@ -685,7 +694,7 @@ async fn count_transaction( sqlx::query_scalar::<_, i64>(query) .fetch_one(&mut **transaction) .await - .map_err(|_| metadata_error())?, + .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?, ) } diff --git a/crates/storage_sqlite/src/open.rs b/crates/storage_sqlite/src/open.rs @@ -156,6 +156,8 @@ fn validate_parent(path: &Path) -> Result<(), Error> { #[derive(Debug)] #[non_exhaustive] pub enum Error { + /// Local capacity failure; inspect existing state before retrying startup. + SpaceInsufficient, InvalidPath(PathBuf), UnexpectedFileName { path: PathBuf, @@ -316,6 +318,7 @@ pub enum Error { impl fmt::Display for Error { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { match self { + Self::SpaceInsufficient => formatter.write_str("storage space is insufficient"), Self::InvalidPath(path) => { write!(formatter, "invalid owned SQLite path: {}", path.display()) } @@ -681,18 +684,22 @@ impl SqliteStorage { options.busy_timeout(), ) .await?; - runtime_connection - .close() - .await - .map_err(|_| Error::DatabaseCloseFailed { - database: RUNTIME_DATABASE_NAME, - })?; - private_connection - .close() - .await - .map_err(|_| Error::DatabaseCloseFailed { - database: PRIVATE_DATABASE_NAME, - })?; + runtime_connection.close().await.map_err(|source| { + crate::backend::startup_error( + &source, + Error::DatabaseCloseFailed { + database: RUNTIME_DATABASE_NAME, + }, + ) + })?; + private_connection.close().await.map_err(|source| { + crate::backend::startup_error( + &source, + Error::DatabaseCloseFailed { + database: PRIVATE_DATABASE_NAME, + }, + ) + })?; let runtime_pool = pool(runtime_options, RUNTIME_DATABASE_NAME).await?; let private_pool = pool(private_options, PRIVATE_DATABASE_NAME).await?; @@ -751,7 +758,9 @@ pub(crate) fn map_database_open_error(source: &sqlx::Error, database: &'static s .and_then(|error| error.code()) .and_then(|code| code.parse::<i32>().ok()) .is_some_and(|code| matches!(code & 0xff, 11 | 26)); - if is_corrupt { + if crate::backend::is_capacity(source) { + Error::SpaceInsufficient + } else if is_corrupt { Error::DatabaseCorrupt { database } } else { Error::DatabaseOpenFailed { database } @@ -787,7 +796,9 @@ async fn active_source_generation( let mut transaction = connection .begin_with("BEGIN IMMEDIATE") .await - .map_err(|_| Error::SourceGenerationUnavailable)?; + .map_err(|source| { + crate::backend::startup_error(&source, Error::SourceGenerationUnavailable) + })?; let rows = active_generation_rows(&mut transaction).await?; let generation = match rows.as_slice() { [] => { @@ -801,16 +812,17 @@ async fn active_source_generation( .bind(i64::try_from(created_at).map_err(|_| Error::CorruptSourceGeneration)?) .execute(&mut *transaction) .await - .map_err(|_| Error::SourceGenerationUnavailable)?; + .map_err(|source| { + crate::backend::startup_error(&source, Error::SourceGenerationUnavailable) + })?; generation } [_] => existing_source_generation(rows.as_slice(), expected)?, _ => return Err(Error::CorruptSourceGeneration), }; - transaction - .commit() - .await - .map_err(|_| Error::SourceGenerationUnavailable)?; + transaction.commit().await.map_err(|source| { + crate::backend::startup_error(&source, Error::SourceGenerationUnavailable) + })?; Ok(generation) } @@ -825,7 +837,7 @@ async fn active_generation_rows( ) .fetch_all(connection) .await - .map_err(|_| Error::SourceGenerationUnavailable) + .map_err(|source| crate::backend::startup_error(&source, Error::SourceGenerationUnavailable)) } fn existing_source_generation( @@ -884,25 +896,35 @@ async fn verify_connection( let foreign_keys = sqlx::query_scalar::<_, i64>("PRAGMA foreign_keys") .fetch_one(&mut *connection) .await - .map_err(|_| Error::ConnectionPolicyMismatch { database })?; + .map_err(|source| { + crate::backend::startup_error(&source, Error::ConnectionPolicyMismatch { database }) + })?; let journal_mode = sqlx::query_scalar::<_, String>("PRAGMA journal_mode") .fetch_one(&mut *connection) .await - .map_err(|_| Error::ConnectionPolicyMismatch { database })?; + .map_err(|source| { + crate::backend::startup_error(&source, Error::ConnectionPolicyMismatch { database }) + })?; let configured_busy_timeout = sqlx::query_scalar::<_, i64>("PRAGMA busy_timeout") .fetch_one(&mut *connection) .await - .map_err(|_| Error::ConnectionPolicyMismatch { database })?; + .map_err(|source| { + crate::backend::startup_error(&source, Error::ConnectionPolicyMismatch { database }) + })?; let synchronous = sqlx::query_scalar::<_, i64>("PRAGMA synchronous") .fetch_one(&mut *connection) .await - .map_err(|_| Error::ConnectionPolicyMismatch { database })?; + .map_err(|source| { + crate::backend::startup_error(&source, Error::ConnectionPolicyMismatch { database }) + })?; let expected_busy_timeout = i64::try_from(busy_timeout.as_millis()) .map_err(|_| Error::ConnectionPolicyMismatch { database })?; let fullfsync = sqlx::query_scalar::<_, i64>("PRAGMA fullfsync") .fetch_one(&mut *connection) .await - .map_err(|_| Error::ConnectionPolicyMismatch { database })?; + .map_err(|source| { + crate::backend::startup_error(&source, Error::ConnectionPolicyMismatch { database }) + })?; if foreign_keys == 1 && journal_mode.eq_ignore_ascii_case("wal") && configured_busy_timeout == expected_busy_timeout @@ -933,6 +955,38 @@ impl StdError for Error { #[cfg(test)] mod error_mapping_tests { use super::*; + + #[test] + fn startup_capacity_is_typed_redacted_and_preserves_generic_fallbacks() { + for source in [ + coded_error(Some("13")), + coded_error(Some("269")), + sqlx::Error::Io(std::io::Error::new( + std::io::ErrorKind::StorageFull, + "private path", + )), + sqlx::Error::Io(std::io::Error::new( + std::io::ErrorKind::QuotaExceeded, + "private path", + )), + ] { + let mapped = map_database_open_error(&source, RUNTIME_DATABASE_NAME); + assert!(matches!(mapped, Error::SpaceInsufficient)); + assert_eq!(mapped.to_string(), "storage space is insufficient"); + assert_eq!(format!("{mapped:?}"), "SpaceInsufficient"); + assert!(matches!( + crate::backend::startup_error(&source, Error::SourceGenerationUnavailable), + Error::SpaceInsufficient + )); + } + assert!(matches!( + crate::backend::startup_error( + &sqlx::Error::PoolClosed, + Error::SourceGenerationUnavailable + ), + Error::SourceGenerationUnavailable + )); + } use sqlx::error::{DatabaseError, ErrorKind}; use std::borrow::Cow; diff --git a/crates/sync/README.md b/crates/sync/README.md @@ -6,6 +6,11 @@ The package owns the shared ingest, pull, projection, push, policy, and status boundaries. It does not create an executor, spawn workers, install timers, own process lifecycle, store UI state, or branch on concrete transport adapters. +`StorageSpaceInsufficient` preserves the storage owner's typed capacity failure. +It does not establish rollback or absence of signer or remote effects. Retain +the original requests, signed bytes and unresolved receipts, then reconcile +existing state before retrying. Sync never frees storage or retries implicitly. + Ingest performs real event-ID and signature verification before host policy. Contract failure defaults to rejection. A host can explicitly retain such a signed observation through `AdmissionPolicy::contract_failure` for canonical diff --git a/crates/sync/src/ingest.rs b/crates/sync/src/ingest.rs @@ -316,6 +316,7 @@ fn hash_field(hasher: &mut Sha256, value: &[u8]) { fn map_storage_error(error: StorageError) -> Error { match error { + StorageError::SpaceInsufficient => Error::StorageSpaceInsufficient, StorageError::EventConflict | StorageError::AtomicCommitConflict => Error::StorageConflict, _ => Error::StorageFailed, } diff --git a/crates/sync/src/policy.rs b/crates/sync/src/policy.rs @@ -220,6 +220,8 @@ pub enum Error { VerificationFailed, PolicyRejected, StorageConflict, + /// Capacity failure; original effects may already be durable and require reconciliation. + StorageSpaceInsufficient, StorageFailed, InvalidIngestReceipt, InvalidPullRequest, @@ -256,6 +258,7 @@ impl core::fmt::Display for Error { Self::VerificationFailed => "sync event verification failed", Self::PolicyRejected => "sync admission policy rejected the event", Self::StorageConflict => "sync input conflicts with durable storage state", + Self::StorageSpaceInsufficient => "sync storage space is insufficient", Self::StorageFailed => "sync storage operation failed", Self::InvalidIngestReceipt => "sync storage returned an invalid ingest receipt", Self::InvalidPullRequest => "sync pull request is outside its bounds", diff --git a/crates/sync/src/projection.rs b/crates/sync/src/projection.rs @@ -550,6 +550,7 @@ const fn admission_stage_byte(stage: AdmissionStage) -> u8 { fn map_storage_error(error: StorageError) -> Error { match error { + StorageError::SpaceInsufficient => Error::StorageSpaceInsufficient, StorageError::ProjectionCheckpointMismatch | StorageError::ProjectionCheckpointRegression | StorageError::ProjectionRevisionConflict @@ -565,6 +566,10 @@ mod tests { #[test] fn raw_source_identity_stage_encoding_and_error_mapping_are_exact() { + assert_eq!( + map_storage_error(StorageError::SpaceInsufficient), + Error::StorageSpaceInsufficient + ); let invalidation = ProjectionInvalidation::new( ProjectionId::parse("projection-helper").unwrap(), ProjectionGeneration::new([1; 32]).unwrap(), diff --git a/crates/sync/src/push.rs b/crates/sync/src/push.rs @@ -850,7 +850,10 @@ impl Engine { applied_at, ) .await?; - return Err(Error::AdmissionFailed); + return Err(match error { + radroots_storage::Error::SpaceInsufficient => Error::StorageSpaceInsufficient, + _ => Error::AdmissionFailed, + }); } }; let state = match admission_receipt.disposition() { @@ -1238,6 +1241,7 @@ fn delivery_request_id(id: SyncId) -> String { fn map_storage_error(error: radroots_storage::Error) -> Error { match error { + radroots_storage::Error::SpaceInsufficient => Error::StorageSpaceInsufficient, radroots_storage::Error::IdempotencyConflict | radroots_storage::Error::OperationIdentityMismatch | radroots_storage::Error::JournalRevisionConflict @@ -1619,6 +1623,14 @@ mod tests { #[test] fn storage_error_maps_are_explicit_and_fail_closed() { assert_eq!( + map_storage_error(radroots_storage::Error::SpaceInsufficient), + Error::StorageSpaceInsufficient + ); + assert_eq!( + map_claim_error(radroots_storage::Error::SpaceInsufficient), + Error::StorageSpaceInsufficient + ); + assert_eq!( map_storage_error(radroots_storage::Error::AtomicCommitConflict), Error::StorageConflict ); diff --git a/crates/sync/src/status.rs b/crates/sync/src/status.rs @@ -23,6 +23,13 @@ use crate::{Engine, policy::Error}; const STATUS_PROJECTION_LIMIT: usize = 256; +fn map_storage_error(error: radroots_storage::Error) -> Error { + match error { + radroots_storage::Error::SpaceInsufficient => Error::StorageSpaceInsufficient, + _ => Error::StorageFailed, + } +} + /// Typed report for one optional injected host capability. #[cfg_attr(feature = "serde", derive(serde::Serialize))] #[derive(Clone, Debug, Eq, PartialEq)] @@ -177,13 +184,13 @@ impl Engine { } let storage = StorageStatusProvider::storage_status(self.storage.as_ref()) .await - .map_err(|_| Error::StorageFailed)?; + .map_err(map_storage_error)?; let events = EventStore::status(self.storage.as_ref()) .await - .map_err(|_| Error::StorageFailed)?; + .map_err(map_storage_error)?; let outbox = Outbox::status(self.storage.as_ref()) .await - .map_err(|_| Error::StorageFailed)?; + .map_err(map_storage_error)?; if outbox.total().is_none() { return Err(Error::StorageFailed); } @@ -194,7 +201,7 @@ impl Engine { for projection_id in projection_ids { let status = ProjectionStore::status(self.storage.as_ref(), projection_id.clone()) .await - .map_err(|_| Error::StorageFailed)?; + .map_err(map_storage_error)?; projections.push(ProjectionReport { projection_id: projection_id.clone(), status, @@ -223,9 +230,7 @@ impl Engine { return Err(Error::ClockUnavailable); } if plan.request().is_some() - && plan - .delivery_satisfaction() - .map_err(|_| Error::StorageFailed)? + && plan.delivery_satisfaction().map_err(map_storage_error)? == radroots_transport::policy::SatisfactionState::Satisfied { return Ok(SyncRetryDecision::Satisfied); @@ -427,6 +432,14 @@ mod tests { #[test] fn health_and_protocol_classification_cover_every_state() { assert_eq!( + map_storage_error(radroots_storage::Error::SpaceInsufficient), + Error::StorageSpaceInsufficient + ); + assert_eq!( + map_storage_error(radroots_storage::Error::BackendUnavailable), + Error::StorageFailed + ); + assert_eq!( availability_state(Availability::Available), SyncCapabilityState::Available ); diff --git a/crates/sync/tests/engine_composition.rs b/crates/sync/tests/engine_composition.rs @@ -279,6 +279,7 @@ fn invalid_compositions_and_ambient_policy_inputs_fail_closed() { Error::VerificationFailed, Error::PolicyRejected, Error::StorageConflict, + Error::StorageSpaceInsufficient, Error::StorageFailed, Error::InvalidIngestReceipt, Error::InvalidPullRequest, diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs @@ -63,9 +63,13 @@ mod delivery_evidence; #[path = "push_enqueue/delivery_selection.rs"] mod delivery_selection; +#[path = "push_enqueue/capacity.rs"] +mod capacity; + struct MockSink; struct FaultStorage { + capacity: Mutex<Option<capacity::Fault>>, inner: Arc<MemoryStorage>, delivery_injection: Mutex<Option<delivery_evidence::Injection>>, receipt_mutation: Mutex<Option<delivery_evidence::ReceiptMutation>>, @@ -79,6 +83,7 @@ struct FaultStorage { impl FaultStorage { fn new(generation: u8) -> Self { Self { + capacity: Mutex::new(None), delivery_injection: Mutex::new(None), receipt_mutation: Mutex::new(None), history_mutation: Mutex::new(None), @@ -136,6 +141,7 @@ impl EventStore for FaultStorage { match self.admit_error.load(Ordering::Relaxed) { 1 => Box::pin(async { Err(radroots_storage::Error::EventConflict) }), 2 => Box::pin(async { Err(radroots_storage::Error::BackendUnavailable) }), + 3 => Box::pin(async { Err(radroots_storage::Error::SpaceInsufficient) }), _ => EventStore::admit(self.inner.as_ref(), value), } } @@ -489,9 +495,11 @@ impl AuthoredAtomicStorage for FaultStorage { > { Box::pin(async move { delivery_evidence::inject(self, &command, true).await; + capacity::inject(self, &command, false)?; let receipt = AuthoredAtomicStorage::execute_authored(self.inner.as_ref(), command.clone()) .await?; + capacity::inject(self, &command, true)?; delivery_evidence::inject(self, &command, false).await; if matches!( receipt.outcome(), diff --git a/crates/sync/tests/push_enqueue/capacity.rs b/crates/sync/tests/push_enqueue/capacity.rs @@ -0,0 +1,142 @@ +use super::*; +use radroots_storage::authored_atomic::AuthoredAtomicCommand; + +#[derive(Clone, Copy)] +pub(super) enum Phase { + Prepare, + Claim, + Signed, + DeliveryFact, +} + +pub(super) struct Fault { + pub(super) phase: Phase, + pub(super) after: bool, +} + +pub(super) fn inject( + storage: &FaultStorage, + command: &AuthoredAtomicCommand, + after: bool, +) -> Result<(), radroots_storage::Error> { + let mut fault = storage.capacity.lock().unwrap(); + if fault.as_ref().is_some_and(|fault| { + fault.after == after + && matches!( + (fault.phase, command), + (Phase::Prepare, AuthoredAtomicCommand::Prepare(_)) + | (Phase::Claim, AuthoredAtomicCommand::Claim(_)) + | (Phase::Signed, AuthoredAtomicCommand::RecordSigned(_)) + | ( + Phase::DeliveryFact, + AuthoredAtomicCommand::RecordDelivery(_) + ) + ) + }) { + *fault = None; + Err(radroots_storage::Error::SpaceInsufficient) + } else { + Ok(()) + } +} + +fn setup(byte: u8) -> (Engine, Arc<FaultStorage>, Arc<MockSigner>, PushRequest) { + let storage = Arc::new(FaultStorage::new(byte)); + let signer = Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })); + let engine = fault_engine(storage.clone(), signer.clone(), Arc::new(MockSink)); + ( + engine, + storage, + signer, + request(byte, "wss://capacity.example"), + ) +} + +#[test] +fn capacity_prepare_reports_failure_and_reconciles_the_original_operation() { + for after in [false, true] { + let (engine, storage, signer, push) = setup(181); + *storage.capacity.lock().unwrap() = Some(Fault { + phase: Phase::Prepare, + after, + }); + assert_eq!( + block_on(engine.prepare_push(push.clone())), + Err(Error::StorageSpaceInsufficient) + ); + let observed = block_on(engine.push_status(push.operation_id())).unwrap(); + assert_eq!(observed.is_some(), after); + let original = block_on(engine.prepare_push(push.clone())).unwrap(); + let replay = block_on(engine.prepare_push(push)).unwrap(); + assert_eq!(original.operation(), replay.operation()); + assert_eq!(original.artifact(), replay.artifact()); + assert_eq!(signer.calls.load(Ordering::Relaxed), 0); + } +} + +#[test] +fn capacity_after_signed_receipt_preserves_exact_bytes_without_resigning() { + let (engine, storage, signer, push) = setup(182); + *storage.capacity.lock().unwrap() = Some(Fault { + phase: Phase::Signed, + after: true, + }); + assert_eq!( + block_on(engine.sign_prepared(push.clone())), + Err(Error::StorageSpaceInsufficient) + ); + let original = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + let signed = original + .artifact() + .signed() + .expect("durable original signed bytes"); + let replay = block_on(engine.sign_prepared(push)).unwrap(); + assert_eq!(replay.artifact().signed().unwrap(), signed); + assert_eq!(signer.calls.load(Ordering::Relaxed), 1); +} + +#[test] +fn capacity_claim_failure_never_calls_the_signer_or_discards_prepared_work() { + for after in [false, true] { + let (engine, storage, signer, push) = setup(183); + block_on(engine.prepare_push(push.clone())).unwrap(); + *storage.capacity.lock().unwrap() = Some(Fault { + phase: Phase::Claim, + after, + }); + assert_eq!( + block_on(engine.sign_prepared(push.clone())), + Err(Error::StorageSpaceInsufficient) + ); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!(status.artifact().signing_claim().is_some(), after); + assert!(status.artifact().signed().is_none()); + assert_eq!(signer.calls.load(Ordering::Relaxed), 0); + } +} + +#[test] +fn capacity_local_admission_retains_signed_bytes_and_reports_no_success() { + let (engine, storage, signer, push) = setup(184); + block_on(engine.sign_prepared(push.clone())).unwrap(); + let original = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + storage.fail_admission_with(3); + assert_eq!( + block_on(engine.admit_signed(push.operation_id())), + Err(Error::StorageSpaceInsufficient) + ); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!(status.artifact().signed(), original.artifact().signed()); + assert!(!status.artifact().admission_state().is_admitted()); + assert_eq!(signer.calls.load(Ordering::Relaxed), 1); +} diff --git a/crates/sync/tests/push_enqueue/delivery_evidence.rs b/crates/sync/tests/push_enqueue/delivery_evidence.rs @@ -80,6 +80,74 @@ fn accept((request, sender): Pending) -> DeliveryReceipt { } #[test] +fn capacity_delivery_receipt_failure_retains_claim_and_known_or_unknown_effects() { + for after in [false, true] { + let storage = Arc::new(FaultStorage::new(185)); + let sink = Arc::new(HeldSink::default()); + let signer = Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })); + let engine = fault_engine(storage.clone(), signer.clone(), sink.clone()); + let push = request(185, "wss://capacity.example"); + execute_to_admitted(&engine, &push); + let before = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + let mut future = Box::pin(engine.deliver_push(push.operation_id())); + assert!( + future + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + let pending = sink.take(); + assert_eq!(&pending.0, before.delivery_plan().request().unwrap()); + *storage.capacity.lock().unwrap() = Some(capacity::Fault { + phase: capacity::Phase::DeliveryFact, + after, + }); + let expected = accept(pending); + assert_eq!(block_on(future), Err(Error::StorageSpaceInsufficient)); + let failed = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!(failed.artifact().signed(), before.artifact().signed()); + assert_eq!( + failed.delivery_plan().request(), + before.delivery_plan().request() + ); + assert!(!failed.delivery_history().proves_no_issued_attempt()); + assert_eq!(failed.delivery_history().has_unresolved_claims(), !after); + assert_eq!( + failed.delivery_plan().delivery_facts().len(), + usize::from(after) + ); + if after { + assert_eq!( + failed.delivery_plan().delivery_facts()[0].outcome(), + &DeliveryAttemptOutcome::Receipt(expected) + ); + assert_eq!( + block_on(engine.deliver_push(push.operation_id())), + Err(Error::WorkClaimConflict) + ); + let expires_at = failed + .delivery_plan() + .claim_evidence() + .unwrap() + .expires_at_unix_ms(); + let recovery = source_only( + storage.clone(), + Arc::new(TestClock(AtomicU64::new(expires_at))), + ); + let reconciled = block_on(recovery.deliver_push(push.operation_id())).unwrap(); + assert_eq!(reconciled.plan().state(), AuthoredDeliveryState::Satisfied); + } + assert_eq!(sink.calls.load(Ordering::Relaxed), 1); + assert_eq!(signer.calls.load(Ordering::Relaxed), 1); + } +} + +#[test] fn late_accepted_result_survives_stop_and_terminal_replay_without_sink() { let (engine, storage, clock, sink, push) = setup(121); let status = block_on(engine.push_status(push.operation_id()))