commit a1e03ec1fc8b042359d2895c4b6e8ab7d6122884 parent 85e32e4dcd77d76385184ad09283561628f686e2 Author: triesap <tyson@radroots.org> Date: Fri, 28 Aug 2026 01:55:18 +0000 storage: bound durable operation journal Diffstat:
14 files changed, 777 insertions(+), 42 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md @@ -78,9 +78,14 @@ substitute. - Product state is bound only through `radroots_runtime_paths::RuntimeContext` for service `harvestcircle` and instance `desktop`; the canonical database - and lock names are `state.sqlite` and `state.lock`. The exact schema starts - at v1, future migrations start at v2, and `radroots_service_sqlite` owns the - governed SQLite mechanics. + and lock names are `state.sqlite` and `state.lock`. Fresh state initializes + at schema v1 and is migrated through the pinned current schema v2 before host + exposure; `radroots_service_sqlite` owns the governed SQLite mechanics. +- The durable-operation journal records terminal completion time, retains + terminal receipts for exactly seven days, admits at most 1,024 unfinished + operations and 4,096 total rows, and deletes no more than 256 expired terminal + rows in one transaction. Admission reserves terminal capacity in the same row + and must never evict an in-window receipt. - `harvestcircle.sqlite3` is legacy evidence only. Never delete, rename, import, dual-read, dual-write, or otherwise treat it as current state. - SQLx is the only high-level SQLite library. HarvestCircle may use its sealed diff --git a/README.md b/README.md @@ -100,6 +100,15 @@ secret. Exact same-operation replay is idempotent only when the complete envelope matches; another operation conflicts and never overwrites the existing credential. No compatibility path reads the former plaintext credential shape. +The state database initializes at schema v1 and applies the pinned schema-v2 +operation-journal migration before the host is exposed. Terminal receipts carry +an explicit completion time and remain replayable for exactly seven days. +Admission caps unfinished operations at 1,024 and all journal rows at 4,096, +deletes at most 256 expired terminal receipts in one transaction, reserves each +accepted operation's terminal row in place, and never evicts an in-window +receipt. The migration, resulting table and guards, and both schema snapshots +are checksum-pinned; invalid legacy rows roll the migration back atomically. + ## Project documentation The consuming Radroots monorepo owns normative HarvestCircle specifications, diff --git a/core/compatibility/harvestcircle-ffi-v4.properties b/core/compatibility/harvestcircle-ffi-v4.properties @@ -2,11 +2,11 @@ schema=harvestcircle.ffi.v4 contract.id=harvestcircle-desktop-ffi-v4 contract.major=4 contract.minor=3 -contract.hash=5aebb8c68546fee90611314163db0d944705be6d88dfd1a44dccd202d0e805d2 +contract.hash=237a67e864761fded90c9a41d92da7b6a514ef22810adbdc2784b7ad5e10826a product.coordinate_digest=bf50f9ea6c2537406de255f025463e670eb6263c295f992f7e4c4db36d957064 snapshot.schema=1 storage.schema.minimum=1 -storage.schema.current=1 +storage.schema.current=2 product.version=0.1.0-alpha package.version=1.0.0 source.provenance_digest=40b9eccd486026128f92de8d55d002a9030f235a35f9b754c98c0b0d387bd8c0 diff --git a/core/compatibility/harvestcircle-storage-api-v1.txt b/core/compatibility/harvestcircle-storage-api-v1.txt @@ -67,6 +67,8 @@ pub const harvestcircle_storage::HARVESTCIRCLE_ACTOR_MAILBOX_CAPACITY: usize pub const harvestcircle_storage::HARVESTCIRCLE_APPLICATION_ID: u32 pub const harvestcircle_storage::HARVESTCIRCLE_COMMAND_DEADLINE_MAX_MS: u64 pub const harvestcircle_storage::HARVESTCIRCLE_COMMAND_DEADLINE_MIN_MS: u64 +pub const harvestcircle_storage::HARVESTCIRCLE_DURABLE_OPERATION_CAPACITY: usize +pub const harvestcircle_storage::HARVESTCIRCLE_DURABLE_OPERATION_CLEANUP_BATCH: usize pub const harvestcircle_storage::HARVESTCIRCLE_EVENTS_PER_RELAY_CAPACITY: usize pub const harvestcircle_storage::HARVESTCIRCLE_EVENTS_TOTAL_CAPACITY: usize pub const harvestcircle_storage::HARVESTCIRCLE_IDENTITY_CAPACITY: usize @@ -77,6 +79,7 @@ pub const harvestcircle_storage::HARVESTCIRCLE_RELAY_ENDPOINT_CAPACITY: usize pub const harvestcircle_storage::HARVESTCIRCLE_RELAY_URL_UTF8_BYTES: usize pub const harvestcircle_storage::HARVESTCIRCLE_SERVICE_ID: &str pub const harvestcircle_storage::HARVESTCIRCLE_STATE_SCHEMA_VERSION: u32 +pub const harvestcircle_storage::HARVESTCIRCLE_TERMINAL_RECEIPT_RETENTION_SECONDS: i64 pub const harvestcircle_storage::HARVESTCIRCLE_UNFINISHED_DURABLE_OPERATION_CAPACITY: usize pub fn harvestcircle_storage::harvestcircle_migration_catalog() -> core::result::Result<radroots_service_sqlite::migration::MigrationCatalog, harvestcircle_storage::HarvestCircleStorageContractError> pub fn harvestcircle_storage::harvestcircle_schema_catalog() -> core::result::Result<radroots_service_sqlite::integrity::catalog::SchemaCatalog, harvestcircle_storage::HarvestCircleStorageContractError> diff --git a/core/crates/harvestcircle_application/src/ports.rs b/core/crates/harvestcircle_application/src/ports.rs @@ -109,6 +109,7 @@ pub struct DurableOperationReceipt { identity: PublicKey, outcome: DurableTerminalOutcome, resulting_revision: Option<u64>, + completed_at: UnixTimestamp, } impl DurableOperationReceipt { @@ -118,12 +119,14 @@ impl DurableOperationReceipt { identity: PublicKey, outcome: DurableTerminalOutcome, resulting_revision: Option<u64>, + completed_at: UnixTimestamp, ) -> Self { Self { request_id, identity, outcome, resulting_revision, + completed_at, } } @@ -146,6 +149,11 @@ impl DurableOperationReceipt { pub const fn resulting_revision(&self) -> Option<u64> { self.resulting_revision } + + #[must_use] + pub const fn completed_at(&self) -> UnixTimestamp { + self.completed_at + } } #[derive(Clone, Debug, Eq, PartialEq)] @@ -766,9 +774,11 @@ mod tests { PublicKey::from_bytes([7; 32]).expect("valid public key"), DurableTerminalOutcome::Completed, Some(42), + UnixTimestamp::from_seconds(7).expect("time"), ); assert_eq!(receipt.request_id(), &request); assert_eq!(receipt.resulting_revision(), Some(42)); + assert_eq!(receipt.completed_at().as_seconds(), 7); for invalid in [ "", "contains space", diff --git a/core/crates/harvestcircle_application/src/recovery.rs b/core/crates/harvestcircle_application/src/recovery.rs @@ -570,7 +570,7 @@ pub(crate) mod tests { expected_phase: DurableOperationPhase, outcome: DurableTerminalOutcome, resulting_revision: Option<u64>, - _updated_at: UnixTimestamp, + updated_at: UnixTimestamp, ) -> BoxFuture<'a, Result<DurableOperationReceipt, SafeError>> { Box::pin(async move { let mut operation = self.operation(); @@ -582,6 +582,7 @@ pub(crate) mod tests { operation.identity(), outcome, resulting_revision, + updated_at, ); *operation = Self::replace( &operation, diff --git a/core/crates/harvestcircle_ffi/build.rs b/core/crates/harvestcircle_ffi/build.rs @@ -219,7 +219,7 @@ fn validate_baseline_inputs(baseline: &BTreeMap<String, String>) { assert_eq!(required(baseline, "contract.minor"), "3"); assert_eq!(required(baseline, "snapshot.schema"), "1"); assert_eq!(required(baseline, "storage.schema.minimum"), "1"); - assert_eq!(required(baseline, "storage.schema.current"), "1"); + assert_eq!(required(baseline, "storage.schema.current"), "2"); assert_eq!( required(baseline, "product.version"), env!("CARGO_PKG_VERSION") @@ -246,7 +246,7 @@ fn validate_baseline_inputs(baseline: &BTreeMap<String, String>) { required(baseline, "source.foundation_baseline") ); - assert_eq!(required(baseline, "storage.schema.current"), "1"); + assert_eq!(required(baseline, "storage.schema.current"), "2"); } fn required<'a>(baseline: &'a BTreeMap<String, String>, key: &str) -> &'a str { diff --git a/core/crates/harvestcircle_ffi/src/commands.rs b/core/crates/harvestcircle_ffi/src/commands.rs @@ -1600,7 +1600,7 @@ mod tests { .to_path_buf(); assert!(database.ends_with("data/services/harvestcircle/desktop/state.sqlite")); assert!(!database.to_string_lossy().contains("harvestcircle.sqlite3")); - assert_eq!(CURRENT_SCHEMA_VERSION, 1); + assert_eq!(CURRENT_SCHEMA_VERSION, 2); } #[test] diff --git a/core/crates/harvestcircle_storage/src/contract.rs b/core/crates/harvestcircle_storage/src/contract.rs @@ -5,17 +5,21 @@ use std::error::Error; use radroots_runtime_paths::RuntimeContext; use radroots_service_sqlite::{ - MigrationCatalog, MigrationDescriptor, SchemaCatalog, SchemaCatalogContractError, SchemaDigest, - SchemaObject, SchemaObjectKind, SchemaVersionCatalog, ServiceSqliteApplicationId, - ServiceSqlitePaths, + MigrationCatalog, MigrationChecksum, MigrationDescriptor, SchemaCatalog, + SchemaCatalogContractError, SchemaDigest, SchemaObject, SchemaObjectKind, SchemaVersionCatalog, + ServiceSqliteApplicationId, ServiceSqlitePaths, }; pub const HARVESTCIRCLE_SERVICE_ID: &str = "harvestcircle"; pub const HARVESTCIRCLE_INSTANCE_ID: &str = "desktop"; pub const HARVESTCIRCLE_APPLICATION_ID: u32 = 0x4843_5231; -pub const HARVESTCIRCLE_STATE_SCHEMA_VERSION: u32 = 1; +pub(crate) const HARVESTCIRCLE_INITIAL_STATE_SCHEMA_VERSION: u32 = 1; +pub const HARVESTCIRCLE_STATE_SCHEMA_VERSION: u32 = 2; pub const HARVESTCIRCLE_IDENTITY_CAPACITY: usize = 256; pub const HARVESTCIRCLE_UNFINISHED_DURABLE_OPERATION_CAPACITY: usize = 1_024; +pub const HARVESTCIRCLE_DURABLE_OPERATION_CAPACITY: usize = 4_096; +pub const HARVESTCIRCLE_DURABLE_OPERATION_CLEANUP_BATCH: usize = 256; +pub const HARVESTCIRCLE_TERMINAL_RECEIPT_RETENTION_SECONDS: i64 = 7 * 24 * 60 * 60; pub const HARVESTCIRCLE_PREFERENCE_VALUE_UTF8_BYTES: usize = 4_096; pub const HARVESTCIRCLE_RELAY_ENDPOINT_CAPACITY: usize = 16; pub const HARVESTCIRCLE_RELAY_URL_UTF8_BYTES: usize = 2_048; @@ -135,6 +139,140 @@ pub(crate) const CREATE_DURABLE_OPERATIONS_SQL: &str = r#"CREATE TABLE durable_o ) ) STRICT"#; +const CREATE_DURABLE_OPERATIONS_V2_SQL: &str = r#"CREATE TABLE durable_operations ( + request_id TEXT NOT NULL PRIMARY KEY CHECK ( + length(CAST(request_id AS BLOB)) = 36 + AND request_id = lower(request_id) + AND request_id NOT GLOB '*[^0-9a-f-]*' + AND substr(request_id, 9, 1) = '-' + AND substr(request_id, 14, 1) = '-' + AND substr(request_id, 15, 1) = '7' + AND substr(request_id, 19, 1) = '-' + AND substr(request_id, 20, 1) IN ('8', '9', 'a', 'b') + AND substr(request_id, 24, 1) = '-' + ), + operation_kind TEXT NOT NULL CHECK (operation_kind IN ('create', 'import', 'repair', 'remove')), + account_public_key BLOB NOT NULL CHECK (length(account_public_key) = 32), + binding_public_key BLOB NOT NULL CHECK (length(binding_public_key) = 32), + expected_revision INTEGER CHECK (expected_revision IS NULL OR expected_revision >= 0), + phase TEXT NOT NULL CHECK (phase IN ( + 'intent_recorded', 'credential_written', 'metadata_committed', 'selection_committed', + 'compensation_pending', 'credential_deleted', 'metadata_deleted', 'finalized' + )), + terminal_outcome TEXT CHECK ( + terminal_outcome IS NULL OR terminal_outcome IN ('completed', 'cancelled', 'failed') + ), + prior_selected_public_key BLOB CHECK ( + prior_selected_public_key IS NULL OR length(prior_selected_public_key) = 32 + ), + prior_binding_availability TEXT CHECK ( + prior_binding_availability IS NULL OR prior_binding_availability IN ( + 'available', 'credential_missing', 'store_unavailable' + ) + ), + resulting_revision INTEGER CHECK (resulting_revision IS NULL OR resulting_revision >= 0), + updated_at_unix_s INTEGER NOT NULL CHECK (updated_at_unix_s >= 0), + diagnostic_code TEXT CHECK (diagnostic_code IS NULL OR diagnostic_code IN ( + 'storage_unavailable', 'keyring_unavailable', 'credential_missing', + 'compensation_failed', 'conflict', 'expired' + )), completed_at_unix_s INTEGER CHECK ( + completed_at_unix_s IS NULL OR completed_at_unix_s >= 0 +), + CHECK (account_public_key = binding_public_key), + CHECK ( + (phase = 'finalized' AND terminal_outcome IS NOT NULL) + OR + (phase <> 'finalized' AND terminal_outcome IS NULL AND resulting_revision IS NULL) + ) +) STRICT"#; + +const CREATE_DURABLE_OPERATIONS_RECEIPT_INSERT_GUARD_SQL: &str = r#"CREATE TRIGGER durable_operations_receipt_insert_guard +BEFORE INSERT ON durable_operations +WHEN NOT ( + ( + NEW.phase = 'finalized' + AND NEW.terminal_outcome IS NOT NULL + AND NEW.completed_at_unix_s IS NOT NULL + ) + OR + ( + NEW.phase <> 'finalized' + AND NEW.terminal_outcome IS NULL + AND NEW.resulting_revision IS NULL + AND NEW.completed_at_unix_s IS NULL + ) +) +BEGIN + SELECT RAISE(ABORT, 'durable operation receipt invariant'); +END"#; + +const CREATE_DURABLE_OPERATIONS_RECEIPT_UPDATE_GUARD_SQL: &str = r#"CREATE TRIGGER durable_operations_receipt_update_guard +BEFORE UPDATE ON durable_operations +WHEN NOT ( + ( + NEW.phase = 'finalized' + AND NEW.terminal_outcome IS NOT NULL + AND NEW.completed_at_unix_s IS NOT NULL + ) + OR + ( + NEW.phase <> 'finalized' + AND NEW.terminal_outcome IS NULL + AND NEW.resulting_revision IS NULL + AND NEW.completed_at_unix_s IS NULL + ) +) +BEGIN + SELECT RAISE(ABORT, 'durable operation receipt invariant'); +END"#; + +pub(crate) const MIGRATE_DURABLE_OPERATIONS_V2_SQL: &str = r#"ALTER TABLE durable_operations +ADD COLUMN completed_at_unix_s INTEGER CHECK ( + completed_at_unix_s IS NULL OR completed_at_unix_s >= 0 +); +UPDATE durable_operations +SET completed_at_unix_s = updated_at_unix_s +WHERE phase = 'finalized'; +CREATE TRIGGER durable_operations_receipt_insert_guard +BEFORE INSERT ON durable_operations +WHEN NOT ( + ( + NEW.phase = 'finalized' + AND NEW.terminal_outcome IS NOT NULL + AND NEW.completed_at_unix_s IS NOT NULL + ) + OR + ( + NEW.phase <> 'finalized' + AND NEW.terminal_outcome IS NULL + AND NEW.resulting_revision IS NULL + AND NEW.completed_at_unix_s IS NULL + ) +) +BEGIN + SELECT RAISE(ABORT, 'durable operation receipt invariant'); +END; +CREATE TRIGGER durable_operations_receipt_update_guard +BEFORE UPDATE ON durable_operations +WHEN NOT ( + ( + NEW.phase = 'finalized' + AND NEW.terminal_outcome IS NOT NULL + AND NEW.completed_at_unix_s IS NOT NULL + ) + OR + ( + NEW.phase <> 'finalized' + AND NEW.terminal_outcome IS NULL + AND NEW.resulting_revision IS NULL + AND NEW.completed_at_unix_s IS NULL + ) +) +BEGIN + SELECT RAISE(ABORT, 'durable operation receipt invariant'); +END; +UPDATE durable_operations SET request_id = request_id"#; + pub(crate) const CREATE_INSTALLATION_IDENTITY_SQL: &str = r#"CREATE TABLE installation_identity ( singleton INTEGER NOT NULL PRIMARY KEY CHECK (singleton = 1), installation_id BLOB NOT NULL CHECK (length(installation_id) = 16) @@ -206,8 +344,29 @@ const VERSION_ONE_DIGEST: [u8; 32] = [ 61, 122, 56, 39, 178, 126, 179, 157, 145, 167, 19, 2, 172, 134, 213, 107, 151, 196, 212, 57, 17, 112, 163, 67, 240, 140, 61, 62, 5, 101, 14, 71, ]; +const DURABLE_OPERATIONS_V2_DIGEST: [u8; 32] = [ + 162, 24, 8, 218, 50, 32, 224, 252, 60, 229, 209, 160, 204, 46, 226, 37, 88, 102, 22, 2, 156, + 177, 28, 19, 235, 161, 26, 125, 65, 126, 59, 135, +]; +const DURABLE_OPERATIONS_RECEIPT_INSERT_GUARD_DIGEST: [u8; 32] = [ + 130, 178, 229, 235, 158, 182, 51, 151, 133, 196, 115, 85, 8, 166, 169, 136, 62, 26, 211, 92, + 215, 57, 254, 158, 103, 196, 53, 39, 106, 201, 183, 14, +]; +const DURABLE_OPERATIONS_RECEIPT_UPDATE_GUARD_DIGEST: [u8; 32] = [ + 48, 131, 220, 252, 243, 158, 221, 88, 8, 140, 207, 187, 34, 138, 145, 215, 92, 99, 32, 72, 254, + 25, 96, 241, 33, 172, 150, 186, 98, 114, 27, 142, +]; +const VERSION_TWO_DIGEST: [u8; 32] = [ + 78, 151, 73, 238, 2, 15, 71, 52, 11, 111, 100, 95, 135, 11, 170, 138, 84, 105, 106, 177, 27, 7, + 134, 156, 53, 68, 22, 220, 116, 199, 147, 5, +]; +const DURABLE_OPERATIONS_V2_MIGRATION_CHECKSUM: MigrationChecksum = + MigrationChecksum::from_bytes([ + 107, 95, 237, 250, 255, 0, 44, 110, 142, 194, 92, 163, 84, 27, 96, 31, 210, 37, 151, 186, + 210, 83, 137, 114, 251, 20, 30, 31, 11, 136, 207, 168, + ]); -/// A sealed binding between one HarvestCircle runtime context and the v1 state catalogs. +/// A sealed binding between one HarvestCircle runtime context and the governed state catalogs. /// /// External callers cannot forge alternate paths or catalogs: /// @@ -272,7 +431,7 @@ impl HarvestCircleStorageContract { #[must_use] pub const fn state_schema_version(&self) -> NonZeroU32 { NonZeroU32::new(HARVESTCIRCLE_STATE_SCHEMA_VERSION) - .expect("HarvestCircle schema v1 is nonzero") + .expect("HarvestCircle current schema version is nonzero") } } @@ -317,7 +476,14 @@ pub(crate) const fn harvestcircle_initial_schema_sql() -> &'static [&'static str pub fn harvestcircle_migration_catalog() -> Result<MigrationCatalog, HarvestCircleStorageContractError> { - MigrationCatalog::new(std::iter::empty::<MigrationDescriptor>()) + let migration = MigrationDescriptor::sql( + 2, + "bound_durable_operation_receipts", + MIGRATE_DURABLE_OPERATIONS_V2_SQL, + DURABLE_OPERATIONS_V2_MIGRATION_CHECKSUM, + ) + .map_err(|_| HarvestCircleStorageContractError::MigrationCatalog)?; + MigrationCatalog::new([migration]) .map_err(|_| HarvestCircleStorageContractError::MigrationCatalog) } @@ -329,16 +495,25 @@ pub fn harvestcircle_schema_catalog() -> Result<SchemaCatalog, HarvestCircleStor fn schema_catalog_for( migrations: &MigrationCatalog, ) -> Result<SchemaCatalog, HarvestCircleStorageContractError> { - let version = SchemaVersionCatalog::new( - HARVESTCIRCLE_STATE_SCHEMA_VERSION, - schema_objects()?, + let version_one = SchemaVersionCatalog::new( + HARVESTCIRCLE_INITIAL_STATE_SCHEMA_VERSION, + schema_objects(CREATE_DURABLE_OPERATIONS_SQL)?, SchemaDigest::from_bytes(VERSION_ONE_DIGEST), ) .map_err(schema_error)?; - SchemaCatalog::new(migrations, [version]).map_err(schema_error) + let version_two_objects = schema_objects(CREATE_DURABLE_OPERATIONS_V2_SQL)?; + let version_two = SchemaVersionCatalog::new( + HARVESTCIRCLE_STATE_SCHEMA_VERSION, + version_two_objects, + SchemaDigest::from_bytes(VERSION_TWO_DIGEST), + ) + .map_err(schema_error)?; + SchemaCatalog::new(migrations, [version_one, version_two]).map_err(schema_error) } -fn schema_objects() -> Result<Vec<SchemaObject>, HarvestCircleStorageContractError> { +fn schema_objects( + durable_operations_sql: &'static str, +) -> Result<Vec<SchemaObject>, HarvestCircleStorageContractError> { let identities = [ ( SchemaObjectKind::Table, @@ -378,15 +553,56 @@ fn schema_objects() -> Result<Vec<SchemaObject>, HarvestCircleStorageContractErr "installation_identity", ), ]; - identities + let schema_sql = [ + CREATE_ACCOUNT_IDENTITIES_SQL, + CREATE_LOCAL_SIGNER_BINDINGS_SQL, + CREATE_RUNTIME_STATE_SQL, + CREATE_PROFILE_CACHE_SQL, + CREATE_ACCOUNT_PREFERENCES_SQL, + durable_operations_sql, + CREATE_INSTALLATION_IDENTITY_SQL, + CREATE_INSTALLATION_IDENTITY_NO_UPDATE_SQL, + CREATE_INSTALLATION_IDENTITY_NO_DELETE_SQL, + ]; + let mut objects = identities .into_iter() - .zip(INITIAL_SCHEMA_SQL) + .zip(schema_sql) .zip(OBJECT_DIGESTS) .map(|(((kind, name, table), sql), digest)| { - SchemaObject::new(kind, name, table, sql, SchemaDigest::from_bytes(digest)) - .map_err(schema_error) + let digest = if sql == CREATE_DURABLE_OPERATIONS_V2_SQL { + SchemaDigest::from_bytes(DURABLE_OPERATIONS_V2_DIGEST) + } else { + SchemaDigest::from_bytes(digest) + }; + SchemaObject::new(kind, name, table, sql, digest).map_err(schema_error) }) - .collect() + .collect::<Result<Vec<_>, _>>()?; + if durable_operations_sql == CREATE_DURABLE_OPERATIONS_V2_SQL { + for (name, sql, digest) in [ + ( + "durable_operations_receipt_insert_guard", + CREATE_DURABLE_OPERATIONS_RECEIPT_INSERT_GUARD_SQL, + DURABLE_OPERATIONS_RECEIPT_INSERT_GUARD_DIGEST, + ), + ( + "durable_operations_receipt_update_guard", + CREATE_DURABLE_OPERATIONS_RECEIPT_UPDATE_GUARD_SQL, + DURABLE_OPERATIONS_RECEIPT_UPDATE_GUARD_DIGEST, + ), + ] { + objects.push( + SchemaObject::new( + SchemaObjectKind::Trigger, + name, + "durable_operations", + sql, + SchemaDigest::from_bytes(digest), + ) + .map_err(schema_error)?, + ); + } + } + Ok(objects) } const fn schema_error(_: SchemaCatalogContractError) -> HarvestCircleStorageContractError { @@ -431,10 +647,10 @@ mod tests { contract.application_id().get(), HARVESTCIRCLE_APPLICATION_ID ); - assert_eq!(contract.state_schema_version().get(), 1); - assert_eq!(contract.migrations().current_version(), 1); - assert!(contract.migrations().descriptors().is_empty()); - assert_eq!(contract.schema().versions().len(), 1); + assert_eq!(contract.state_schema_version().get(), 2); + assert_eq!(contract.migrations().current_version(), 2); + assert_eq!(contract.migrations().descriptors().len(), 1); + assert_eq!(contract.schema().versions().len(), 2); assert_eq!(harvestcircle_initial_schema_sql().len(), 9); } @@ -492,7 +708,7 @@ mod tests { assert_eq!(harvestcircle_product::STORAGE_APPLICATION_ID_TEXT, "HCR1"); assert_eq!( harvestcircle_product::STORAGE_INITIAL_SCHEMA_VERSION, - HARVESTCIRCLE_STATE_SCHEMA_VERSION.to_string() + HARVESTCIRCLE_INITIAL_STATE_SCHEMA_VERSION.to_string() ); assert_eq!( harvestcircle_product::LEGACY_DATABASE_FILENAME, @@ -566,6 +782,10 @@ mod tests { .execute(&mut connection) .await .expect("foreign keys"); + sqlx::query("PRAGMA trusted_schema = OFF") + .execute(&mut connection) + .await + .expect("trusted schema"); for statement in harvestcircle_initial_schema_sql() { sqlx::query(*statement) .execute(&mut connection) @@ -611,4 +831,166 @@ mod tests { ) })); } + + #[tokio::test] + async fn migration_v2_adds_the_exact_bounded_journal_schema() { + let mut connection = sqlx::SqliteConnection::connect(":memory:") + .await + .expect("memory database"); + sqlx::query("PRAGMA foreign_keys = ON") + .execute(&mut connection) + .await + .expect("foreign keys"); + sqlx::query("PRAGMA trusted_schema = OFF") + .execute(&mut connection) + .await + .expect("trusted schema"); + for statement in harvestcircle_initial_schema_sql() { + sqlx::query(*statement) + .execute(&mut connection) + .await + .expect("schema statement"); + } + let identity = [7_u8; 32]; + sqlx::query( + "INSERT INTO durable_operations (request_id, operation_kind, account_public_key, \ + binding_public_key, phase, terminal_outcome, updated_at_unix_s) \ + VALUES ('01890f3e-7b1c-7000-8000-000000000001', 'create', ?, ?, \ + 'finalized', 'completed', 10)", + ) + .bind(identity.as_slice()) + .bind(identity.as_slice()) + .execute(&mut connection) + .await + .expect("terminal v1 row"); + sqlx::query( + "INSERT INTO durable_operations (request_id, operation_kind, account_public_key, \ + binding_public_key, phase, updated_at_unix_s) \ + VALUES ('01890f3e-7b1c-7000-8000-000000000002', 'remove', ?, ?, \ + 'intent_recorded', 11)", + ) + .bind(identity.as_slice()) + .bind(identity.as_slice()) + .execute(&mut connection) + .await + .expect("unfinished v1 row"); + sqlx::raw_sql(MIGRATE_DURABLE_OPERATIONS_V2_SQL) + .execute(&mut connection) + .await + .expect("migration"); + let actual: String = sqlx::query_scalar( + "SELECT sql FROM sqlite_schema WHERE type = 'table' AND name = 'durable_operations'", + ) + .fetch_one(&mut connection) + .await + .expect("schema SQL"); + assert_eq!(actual, CREATE_DURABLE_OPERATIONS_V2_SQL); + let rows = sqlx::query( + "SELECT request_id, completed_at_unix_s FROM durable_operations ORDER BY request_id", + ) + .fetch_all(&mut connection) + .await + .expect("migrated rows"); + assert_eq!(rows.len(), 2); + assert_eq!( + rows[0].get::<Option<i64>, _>("completed_at_unix_s"), + Some(10) + ); + assert_eq!(rows[1].get::<Option<i64>, _>("completed_at_unix_s"), None); + assert!( + sqlx::query( + "INSERT INTO durable_operations (request_id, operation_kind, account_public_key, \ + binding_public_key, phase, terminal_outcome, updated_at_unix_s) \ + VALUES ('01890f3e-7b1c-7000-8000-000000000004', 'create', ?, ?, \ + 'finalized', 'completed', 13)", + ) + .bind(identity.as_slice()) + .bind(identity.as_slice()) + .execute(&mut connection) + .await + .is_err(), + "insert guard must require terminal completion time" + ); + assert!( + sqlx::query( + "UPDATE durable_operations SET phase = 'finalized', \ + terminal_outcome = 'completed' WHERE request_id = \ + '01890f3e-7b1c-7000-8000-000000000002'", + ) + .execute(&mut connection) + .await + .is_err(), + "update guard must require terminal completion time" + ); + + let migration = MigrationChecksum::for_sql(MIGRATE_DURABLE_OPERATIONS_V2_SQL); + let objects = schema_objects(CREATE_DURABLE_OPERATIONS_V2_SQL).expect("objects"); + let object = objects + .iter() + .find(|object| object.name() == "durable_operations") + .expect("durable operations") + .digest(); + let snapshot = SchemaVersionCatalog::computed_digest(2, objects).expect("snapshot"); + assert_eq!(migration, DURABLE_OPERATIONS_V2_MIGRATION_CHECKSUM); + assert_eq!( + object, + SchemaDigest::from_bytes(DURABLE_OPERATIONS_V2_DIGEST) + ); + assert_eq!(snapshot, SchemaDigest::from_bytes(VERSION_TWO_DIGEST)); + } + + #[tokio::test] + async fn migration_v2_rolls_back_without_partial_schema_on_invalid_v1_state() { + let mut connection = sqlx::SqliteConnection::connect(":memory:") + .await + .expect("memory database"); + for statement in harvestcircle_initial_schema_sql() { + sqlx::query(*statement) + .execute(&mut connection) + .await + .expect("schema statement"); + } + sqlx::query("PRAGMA ignore_check_constraints = ON") + .execute(&mut connection) + .await + .expect("fixture policy"); + let identity = [9_u8; 32]; + sqlx::query( + "INSERT INTO durable_operations (request_id, operation_kind, account_public_key, \ + binding_public_key, phase, updated_at_unix_s) \ + VALUES ('01890f3e-7b1c-7000-8000-000000000003', 'create', ?, ?, 'finalized', 12)", + ) + .bind(identity.as_slice()) + .bind(identity.as_slice()) + .execute(&mut connection) + .await + .expect("invalid v1 fixture"); + sqlx::query("PRAGMA ignore_check_constraints = OFF") + .execute(&mut connection) + .await + .expect("restore policy"); + + let mut transaction = connection.begin().await.expect("migration transaction"); + assert!( + sqlx::raw_sql(MIGRATE_DURABLE_OPERATIONS_V2_SQL) + .execute(&mut *transaction) + .await + .is_err() + ); + transaction.rollback().await.expect("rollback"); + + let columns: i64 = sqlx::query_scalar( + "SELECT count(*) FROM pragma_table_info('durable_operations') \ + WHERE name = 'completed_at_unix_s'", + ) + .fetch_one(&mut connection) + .await + .expect("column inventory"); + let retained: i64 = sqlx::query_scalar("SELECT count(*) FROM durable_operations") + .fetch_one(&mut connection) + .await + .expect("retained row"); + assert_eq!(columns, 0); + assert_eq!(retained, 1); + } } diff --git a/core/crates/harvestcircle_storage/src/db.rs b/core/crates/harvestcircle_storage/src/db.rs @@ -10,7 +10,9 @@ use radroots_service_sqlite::{ }; use radroots_storage::event::SourceGeneration; -use crate::contract::harvestcircle_initial_schema_sql; +use crate::contract::{ + HARVESTCIRCLE_INITIAL_STATE_SCHEMA_VERSION, harvestcircle_initial_schema_sql, +}; use crate::{HARVESTCIRCLE_STATE_SCHEMA_VERSION, HarvestCircleStorageContract}; pub const CURRENT_SCHEMA_VERSION: u32 = HARVESTCIRCLE_STATE_SCHEMA_VERSION; @@ -47,7 +49,8 @@ impl Database { let metadata = ServiceDatabaseMetadata::new( contract.paths(), generation, - NonZeroU32::new(CURRENT_SCHEMA_VERSION).expect("schema v1 is nonzero"), + NonZeroU32::new(HARVESTCIRCLE_INITIAL_STATE_SCHEMA_VERSION) + .expect("initial schema version is nonzero"), created_at_unix_ms, contract.application_id(), ) diff --git a/core/crates/harvestcircle_storage/src/journal.rs b/core/crates/harvestcircle_storage/src/journal.rs @@ -9,6 +9,10 @@ use harvestcircle_domain::{ use radroots_service_sqlite::ServiceSqliteTransaction; use sqlx::Row; +use crate::contract::{ + HARVESTCIRCLE_DURABLE_OPERATION_CAPACITY, HARVESTCIRCLE_DURABLE_OPERATION_CLEANUP_BATCH, + HARVESTCIRCLE_TERMINAL_RECEIPT_RETENTION_SECONDS, +}; use crate::db::{corrupt_storage, map_transaction_error, storage_unavailable}; use crate::{Database, HARVESTCIRCLE_UNFINISHED_DURABLE_OPERATION_CAPACITY}; @@ -26,7 +30,7 @@ const DURABLE_OPERATION_PROJECTION: &str = "SELECT \ CASE WHEN terminal_outcome IS NULL THEN NULL ELSE length(CAST(terminal_outcome AS BLOB)) END AS terminal_outcome_bytes, \ CASE WHEN prior_binding_availability IS NULL THEN NULL ELSE substr(CAST(prior_binding_availability AS BLOB), 1, 19) END AS prior_binding_availability, \ CASE WHEN prior_binding_availability IS NULL THEN NULL ELSE length(CAST(prior_binding_availability AS BLOB)) END AS prior_binding_availability_bytes, \ - resulting_revision FROM durable_operations"; + resulting_revision, completed_at_unix_s FROM durable_operations"; impl DurableOperationRepository for Database { #[allow(clippy::too_many_arguments)] @@ -56,6 +60,7 @@ impl DurableOperationRepository for Database { } return Ok(DurableOperationStart::Existing(operation)); } + cleanup_expired_terminal_receipts(transaction, updated_at).await?; let unfinished: i64 = sqlx::query_scalar( "SELECT count(*) FROM (SELECT 1 FROM durable_operations \ WHERE terminal_outcome IS NULL LIMIT 1025)", @@ -68,6 +73,18 @@ impl DurableOperationRepository for Database { }) { return Err(operation_capacity_exhausted()); } + let total: i64 = sqlx::query_scalar( + "SELECT count(*) FROM (SELECT 1 FROM durable_operations LIMIT 4097)", + ) + .fetch_one(&mut *transaction) + .await + .map_err(|_| storage_unavailable())?; + if usize::try_from(total) + .ok() + .is_none_or(|count| count >= HARVESTCIRCLE_DURABLE_OPERATION_CAPACITY) + { + return Err(operation_capacity_exhausted()); + } let result = sqlx::query( "INSERT INTO durable_operations (request_id, operation_kind, \ account_public_key, binding_public_key, expected_revision, phase, \ @@ -181,12 +198,14 @@ impl DurableOperationRepository for Database { } let result = sqlx::query( "UPDATE durable_operations SET phase = 'finalized', \ - terminal_outcome = ?, resulting_revision = ?, updated_at_unix_s = ? \ + terminal_outcome = ?, resulting_revision = ?, updated_at_unix_s = ?, \ + completed_at_unix_s = ? \ WHERE request_id = ? AND phase = ? AND terminal_outcome IS NULL", ) .bind(encode_outcome(outcome)) .bind(resulting_revision) .bind(updated_at.as_seconds()) + .bind(updated_at.as_seconds()) .bind(&request_id) .bind(encode_phase(expected_phase)) .execute(&mut *transaction) @@ -294,9 +313,22 @@ fn decode_operation(row: &sqlx::sqlite::SqliteRow) -> Result<DurableIdentityOper row.try_get("resulting_revision") .map_err(|_| corrupt_storage())?, )?; - let terminal = outcome.map(|outcome| { - DurableOperationReceipt::new(request_id.clone(), identity, outcome, resulting_revision) - }); + let completed_at = row + .try_get::<Option<i64>, _>("completed_at_unix_s") + .map_err(|_| corrupt_storage())? + .map(|value| UnixTimestamp::from_seconds(value).ok_or_else(corrupt_storage)) + .transpose()?; + let terminal = match (outcome, completed_at) { + (Some(outcome), Some(completed_at)) => Some(DurableOperationReceipt::new( + request_id.clone(), + identity, + outcome, + resulting_revision, + completed_at, + )), + (None, None) if resulting_revision.is_none() => None, + _ => return Err(corrupt_storage()), + }; Ok(DurableIdentityOperation::new( request_id, kind, @@ -310,6 +342,33 @@ fn decode_operation(row: &sqlx::sqlite::SqliteRow) -> Result<DurableIdentityOper )) } +async fn cleanup_expired_terminal_receipts( + transaction: &mut ServiceSqliteTransaction<'_>, + now: UnixTimestamp, +) -> Result<(), SafeError> { + let cutoff = now + .as_seconds() + .saturating_sub(HARVESTCIRCLE_TERMINAL_RECEIPT_RETENTION_SECONDS); + let result = sqlx::query( + "DELETE FROM durable_operations WHERE request_id IN (\ + SELECT request_id FROM durable_operations \ + WHERE completed_at_unix_s IS NOT NULL AND completed_at_unix_s < ? \ + ORDER BY completed_at_unix_s, request_id LIMIT 256\ + )", + ) + .bind(cutoff) + .execute(&mut *transaction) + .await + .map_err(|_| storage_unavailable())?; + if result.rows_affected() + > u64::try_from(HARVESTCIRCLE_DURABLE_OPERATION_CLEANUP_BATCH) + .expect("cleanup bound fits in u64") + { + return Err(corrupt_storage()); + } + Ok(()) +} + fn required_key( row: &sqlx::sqlite::SqliteRow, value: &str, diff --git a/core/crates/harvestcircle_storage/src/lib.rs b/core/crates/harvestcircle_storage/src/lib.rs @@ -17,13 +17,15 @@ pub use backup::{VerifiedHarvestCircleBackup, verify_harvestcircle_backup}; pub use contract::{ HARVESTCIRCLE_ACTOR_MAILBOX_CAPACITY, HARVESTCIRCLE_APPLICATION_ID, HARVESTCIRCLE_COMMAND_DEADLINE_MAX_MS, HARVESTCIRCLE_COMMAND_DEADLINE_MIN_MS, + HARVESTCIRCLE_DURABLE_OPERATION_CAPACITY, HARVESTCIRCLE_DURABLE_OPERATION_CLEANUP_BATCH, HARVESTCIRCLE_EVENTS_PER_RELAY_CAPACITY, HARVESTCIRCLE_EVENTS_TOTAL_CAPACITY, HARVESTCIRCLE_IDENTITY_CAPACITY, HARVESTCIRCLE_INSTANCE_ID, HARVESTCIRCLE_OBSERVER_CAPACITY, HARVESTCIRCLE_PREFERENCE_VALUE_UTF8_BYTES, HARVESTCIRCLE_RELAY_ENDPOINT_CAPACITY, HARVESTCIRCLE_RELAY_URL_UTF8_BYTES, HARVESTCIRCLE_SERVICE_ID, - HARVESTCIRCLE_STATE_SCHEMA_VERSION, HARVESTCIRCLE_UNFINISHED_DURABLE_OPERATION_CAPACITY, - HarvestCircleStorageContract, HarvestCircleStorageContractError, - harvestcircle_migration_catalog, harvestcircle_schema_catalog, + HARVESTCIRCLE_STATE_SCHEMA_VERSION, HARVESTCIRCLE_TERMINAL_RECEIPT_RETENTION_SECONDS, + HARVESTCIRCLE_UNFINISHED_DURABLE_OPERATION_CAPACITY, HarvestCircleStorageContract, + HarvestCircleStorageContractError, harvestcircle_migration_catalog, + harvestcircle_schema_catalog, }; pub use db::{CURRENT_SCHEMA_VERSION, Database}; pub use os_keyring::{CREDENTIAL_SERVICE, OsKeyringSecretStore}; diff --git a/core/crates/harvestcircle_storage/tests/package_boundary.rs b/core/crates/harvestcircle_storage/tests/package_boundary.rs @@ -12,7 +12,9 @@ fn storage_package_keeps_one_sqlite_authority_and_a_sealed_public_surface() { let workspace_manifest = read(&workspace_root.join("Cargo.toml")); let lock = read(&workspace_root.join("Cargo.lock")); let root_source = read(&crate_root.join("src/lib.rs")); + let contract_source = read(&crate_root.join("src/contract.rs")); let database_source = read(&crate_root.join("src/db.rs")); + let journal_source = read(&crate_root.join("src/journal.rs")); let keyring_source = read(&crate_root.join("src/os_keyring.rs")); let api = read(&workspace_root.join("compatibility/harvestcircle-storage-api-v1.txt")); @@ -74,6 +76,9 @@ fn storage_package_keeps_one_sqlite_authority_and_a_sealed_public_surface() { "pub fn harvestcircle_storage::OsKeyringSecretStore::delete<'a>(&'a self, &'a harvestcircle_application::ports::DurableRequestId, harvestcircle_domain::key::PublicKey) -> harvestcircle_application::ports::BoxFuture<'a", "pub fn harvestcircle_storage::harvestcircle_migration_catalog()", "pub fn harvestcircle_storage::harvestcircle_schema_catalog()", + "pub const harvestcircle_storage::HARVESTCIRCLE_DURABLE_OPERATION_CAPACITY: usize", + "pub const harvestcircle_storage::HARVESTCIRCLE_DURABLE_OPERATION_CLEANUP_BATCH: usize", + "pub const harvestcircle_storage::HARVESTCIRCLE_TERMINAL_RECEIPT_RETENTION_SECONDS: i64", ] { assert!(api.contains(required), "API baseline is missing {required}"); } @@ -96,4 +101,43 @@ fn storage_package_keeps_one_sqlite_authority_and_a_sealed_public_surface() { "API baseline exposes forbidden surface {forbidden}" ); } + + for required in [ + "pub const HARVESTCIRCLE_STATE_SCHEMA_VERSION: u32 = 2;", + "pub const HARVESTCIRCLE_UNFINISHED_DURABLE_OPERATION_CAPACITY: usize = 1_024;", + "pub const HARVESTCIRCLE_DURABLE_OPERATION_CAPACITY: usize = 4_096;", + "pub const HARVESTCIRCLE_DURABLE_OPERATION_CLEANUP_BATCH: usize = 256;", + "pub const HARVESTCIRCLE_TERMINAL_RECEIPT_RETENTION_SECONDS: i64 = 7 * 24 * 60 * 60;", + "bound_durable_operation_receipts", + "completed_at_unix_s", + "durable_operations_receipt_insert_guard", + "durable_operations_receipt_update_guard", + ] { + assert!( + contract_source.contains(required), + "storage contract is missing {required}" + ); + } + for required in [ + "LIMIT 1025", + "LIMIT 4097", + "LIMIT 256", + "completed_at_unix_s < ?", + "completed_at_unix_s = ?", + ] { + assert!( + journal_source.contains(required), + "journal enforcement is missing {required}" + ); + } + for forbidden in [ + "DELETE FROM durable_operations WHERE terminal_outcome IS NOT NULL", + "DELETE FROM durable_operations WHERE completed_at_unix_s IS NOT NULL;", + "ORDER BY completed_at_unix_s DESC", + ] { + assert!( + !journal_source.contains(forbidden), + "journal reintroduced unbounded or in-window eviction: {forbidden}" + ); + } } diff --git a/core/crates/harvestcircle_storage/tests/sqlx_storage.rs b/core/crates/harvestcircle_storage/tests/sqlx_storage.rs @@ -142,6 +142,7 @@ async fn canonical_database_preserves_legacy_state_and_enforces_identity_capacit .await .expect("database"); let generation = database.metadata().source_generation(); + assert_eq!(database.metadata().state_schema_version().get(), 2); let first_identity = identity(0); database @@ -185,6 +186,7 @@ async fn canonical_database_preserves_legacy_state_and_enforces_identity_capacit .await .expect("reopen"); assert_eq!(reopened.metadata().source_generation(), generation); + assert_eq!(reopened.metadata().state_schema_version().get(), 2); assert_eq!( reopened .list_identities() @@ -273,7 +275,7 @@ async fn unfinished_uuid_ledger_enforces_exact_capacity_and_recovers_after_final .expect_err("capacity must reject"); assert_eq!(error.code(), SafeErrorCode::InvalidApplicationState); - database + let receipt = database .finalize_durable_operation( &requests[0], DurableOperationPhase::IntentRecorded, @@ -283,6 +285,7 @@ async fn unfinished_uuid_ledger_enforces_exact_capacity_and_recovers_after_final ) .await .expect("finalize one operation"); + assert_eq!(receipt.completed_at().as_seconds(), 3); database .begin_durable_operation( &excess, @@ -296,3 +299,217 @@ async fn unfinished_uuid_ledger_enforces_exact_capacity_and_recovers_after_final .expect("capacity recovered"); database.close().await.expect("close"); } + +async fn seed_terminal_operations( + connection: &mut SqliteConnection, + count: usize, + completed_at: impl Fn(usize) -> i64, +) { + let identity = PublicKey::from_bytes([7; 32]).expect("public key"); + let mut transaction = connection.begin().await.expect("fixture transaction"); + for index in 0..count { + let request = DurableRequestId::new_v7(); + let completed_at = completed_at(index); + sqlx::query( + "INSERT INTO durable_operations (request_id, operation_kind, account_public_key, \ + binding_public_key, phase, terminal_outcome, updated_at_unix_s, \ + completed_at_unix_s) \ + VALUES (?, 'import', ?, ?, 'finalized', 'completed', ?, ?)", + ) + .bind(request.as_str()) + .bind(identity.as_bytes().as_slice()) + .bind(identity.as_bytes().as_slice()) + .bind(completed_at) + .bind(completed_at) + .execute(&mut *transaction) + .await + .expect("terminal fixture operation"); + } + transaction.commit().await.expect("fixture commit"); +} + +#[tokio::test] +async fn terminal_retention_is_bounded_batched_and_never_evicts_in_window_receipts() { + const TOTAL_CAPACITY: usize = 4_096; + const CLEANUP_BATCH: usize = 256; + const RETENTION_SECONDS: i64 = 7 * 24 * 60 * 60; + const EXPIRED: usize = 300; + const NOW: i64 = 1_000_000; + + let directory = tempdir().expect("directory"); + let context = runtime_context(&directory); + let database_path = context.paths().state().join("state.sqlite"); + let build = build_identity(); + let database = Database::open(&context, 1, 1, &build) + .await + .expect("database"); + database.close().await.expect("fixture close"); + + let options = SqliteConnectOptions::new() + .filename(&database_path) + .create_if_missing(false); + let mut connection = SqliteConnection::connect_with(&options) + .await + .expect("fixture connection"); + seed_terminal_operations(&mut connection, TOTAL_CAPACITY, |index| { + if index < EXPIRED { + NOW - RETENTION_SECONDS - 1 + } else { + NOW - RETENTION_SECONDS + } + }) + .await; + connection.close().await.expect("fixture close"); + + let database = Database::open(&context, 2, 2, &build) + .await + .expect("reopen"); + let identity = PublicKey::from_bytes([7; 32]).expect("public key"); + database + .begin_durable_operation( + &DurableRequestId::new_v7(), + DurableOperationKind::Import, + identity, + None, + OperationPriorState::new(None, None), + UnixTimestamp::from_seconds(NOW).expect("time"), + ) + .await + .expect("admission after bounded cleanup"); + database.close().await.expect("close after cleanup"); + + let mut connection = SqliteConnection::connect_with(&options) + .await + .expect("inspect connection"); + let total: i64 = sqlx::query_scalar("SELECT count(*) FROM durable_operations") + .fetch_one(&mut connection) + .await + .expect("total"); + let expired: i64 = + sqlx::query_scalar("SELECT count(*) FROM durable_operations WHERE completed_at_unix_s < ?") + .bind(NOW - RETENTION_SECONDS) + .fetch_one(&mut connection) + .await + .expect("expired"); + let in_window: i64 = + sqlx::query_scalar("SELECT count(*) FROM durable_operations WHERE completed_at_unix_s = ?") + .bind(NOW - RETENTION_SECONDS) + .fetch_one(&mut connection) + .await + .expect("in-window"); + assert_eq!( + total, + i64::try_from(TOTAL_CAPACITY - CLEANUP_BATCH + 1).unwrap() + ); + assert_eq!(expired, i64::try_from(EXPIRED - CLEANUP_BATCH).unwrap()); + assert_eq!(in_window, i64::try_from(TOTAL_CAPACITY - EXPIRED).unwrap()); + connection.close().await.expect("inspection close"); + + let database = Database::open(&context, 3, 3, &build) + .await + .expect("second reopen"); + database + .begin_durable_operation( + &DurableRequestId::new_v7(), + DurableOperationKind::Import, + identity, + None, + OperationPriorState::new(None, None), + UnixTimestamp::from_seconds(NOW).expect("time"), + ) + .await + .expect("admission after remaining cleanup"); + database.close().await.expect("second cleanup close"); + + let mut connection = SqliteConnection::connect_with(&options) + .await + .expect("final inspection"); + let expired: i64 = + sqlx::query_scalar("SELECT count(*) FROM durable_operations WHERE completed_at_unix_s < ?") + .bind(NOW - RETENTION_SECONDS) + .fetch_one(&mut connection) + .await + .expect("remaining expired"); + let in_window: i64 = + sqlx::query_scalar("SELECT count(*) FROM durable_operations WHERE completed_at_unix_s = ?") + .bind(NOW - RETENTION_SECONDS) + .fetch_one(&mut connection) + .await + .expect("remaining in-window"); + assert_eq!(expired, 0); + assert_eq!(in_window, i64::try_from(TOTAL_CAPACITY - EXPIRED).unwrap()); + connection.close().await.expect("final inspection close"); +} + +#[tokio::test] +async fn total_capacity_reserves_terminal_receipts_without_in_window_eviction() { + const TOTAL_CAPACITY: usize = 4_096; + const NOW: i64 = 1_000_000; + + let directory = tempdir().expect("directory"); + let context = runtime_context(&directory); + let database_path = context.paths().state().join("state.sqlite"); + let build = build_identity(); + let database = Database::open(&context, 1, 1, &build) + .await + .expect("database"); + database.close().await.expect("fixture close"); + + let options = SqliteConnectOptions::new() + .filename(&database_path) + .create_if_missing(false); + let mut connection = SqliteConnection::connect_with(&options) + .await + .expect("fixture connection"); + seed_terminal_operations(&mut connection, TOTAL_CAPACITY - 1, |_| NOW).await; + connection.close().await.expect("fixture close"); + + let database = Database::open(&context, 2, 2, &build) + .await + .expect("reopen"); + let reserved = DurableRequestId::new_v7(); + database + .begin_durable_operation( + &reserved, + DurableOperationKind::Create, + PublicKey::from_bytes([7; 32]).expect("public key"), + None, + OperationPriorState::new(None, None), + UnixTimestamp::from_seconds(NOW).expect("time"), + ) + .await + .expect("last row must reserve its terminal receipt"); + database + .finalize_durable_operation( + &reserved, + DurableOperationPhase::IntentRecorded, + DurableTerminalOutcome::Completed, + None, + UnixTimestamp::from_seconds(NOW + 1).expect("time"), + ) + .await + .expect("reserved row must finalize in place"); + let error = database + .begin_durable_operation( + &DurableRequestId::new_v7(), + DurableOperationKind::Create, + PublicKey::from_bytes([7; 32]).expect("public key"), + None, + OperationPriorState::new(None, None), + UnixTimestamp::from_seconds(NOW).expect("time"), + ) + .await + .expect_err("total capacity must reserve the terminal journal"); + assert_eq!(error.code(), SafeErrorCode::InvalidApplicationState); + database.close().await.expect("close"); + + let mut connection = SqliteConnection::connect_with(&options) + .await + .expect("inspection connection"); + let total: i64 = sqlx::query_scalar("SELECT count(*) FROM durable_operations") + .fetch_one(&mut connection) + .await + .expect("total"); + assert_eq!(total, i64::try_from(TOTAL_CAPACITY).unwrap()); + connection.close().await.expect("inspection close"); +}