rhi

Coordinated trade for connected markets
git clone https://radroots.dev/git/rhi.git
Log | Files | Refs | README | LICENSE

commit 137d6e6902de01025bf2a4b404a4504dbb0dc85a
parent 560f49dc185efa36c459a2a7cfb102f53486b251
Author: triesap <tyson@radroots.org>
Date:   Mon, 24 Aug 2026 13:35:08 +0000

test(rhi): qualify durable publication wave

Diffstat:
MREADME | 8++++++++
Mcontracts/api_baselines/rhi.txt | 3+++
Acontracts/services_hardening/publication_wave_qualification.v1.json | 68++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcontracts/services_hardening/reconciliation_jobs.v1.json | 10++++++++++
Msrc/lib.rs | 6++++--
Msrc/state_catalog.rs | 190++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
Mtests/package_boundary.rs | 12++++++++++++
Mtests/services_hardening_publication_contract.rs | 8+++++---
Atests/services_hardening_publication_wave_qualification_contract.rs | 119+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/services_hardening_reconciliation_job_contract.rs | 13+++++++++++++
Mtests/services_hardening_reconciliation_jobs.rs | 197+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/services_hardening_state_catalog.rs | 58+++++++++++++++++++++++++++++++++++++++++++++++++---------
Mtests/services_hardening_state_host.rs | 213+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/services_hardening_state_resilience.rs | 2+-
14 files changed, 882 insertions(+), 25 deletions(-)

diff --git a/README b/README @@ -422,6 +422,14 @@ and an unknown local commit must be reconciled by exact attempt identity. The exact machine contract is [`publication_execution.v1.json`](contracts/services_hardening/publication_execution.v1.json). +Step 203 closes this wave with executable crash/reopen, shared transactional +failpoint, concurrent-claim, queue-bound, and redaction qualification. Schema +v8 first scans every historical reconciliation job and fails closed without +repair when a `ready` schedule or `leased` owner/expiry is missing. Only after +that scan succeeds does the same governed migration install permanent INSERT +and UPDATE state-shape guards. The exact qualification inventory and bounds are +[`publication_wave_qualification.v1.json`](contracts/services_hardening/publication_wave_qualification.v1.json). + ## Existing-state runtime foundation `open_rhi_runtime_foundation` opens only an already initialized database from diff --git a/contracts/api_baselines/rhi.txt b/contracts/api_baselines/rhi.txt @@ -1622,6 +1622,9 @@ pub const rhi::RHI_STATE_SCHEMA_VERSION_6_SHA256: [u8; 32] pub const rhi::RHI_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256: [u8; 32] pub const rhi::RHI_STATE_SCHEMA_VERSION_7_OBJECT_COUNT: u32 pub const rhi::RHI_STATE_SCHEMA_VERSION_7_SHA256: [u8; 32] +pub const rhi::RHI_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256: [u8; 32] +pub const rhi::RHI_STATE_SCHEMA_VERSION_8_OBJECT_COUNT: u32 +pub const rhi::RHI_STATE_SCHEMA_VERSION_8_SHA256: [u8; 32] pub const rhi::RHI_STATUS_CONTRACT_VERSION: u32 pub const rhi::RHI_TRADE_EVENT_EXTRA_FIELD_MAX_COUNT: usize pub const rhi::RHI_TRADE_EVENT_EXTRA_JSON_MAX_BYTES: usize diff --git a/contracts/services_hardening/publication_wave_qualification.v1.json b/contracts/services_hardening/publication_wave_qualification.v1.json @@ -0,0 +1,68 @@ +{ + "schema": "radroots.rhi.publication-wave-qualification", + "schema_version": 1, + "contract_version": 1, + "step": 203, + "service": "rhi", + "state_schema_version": 8, + "component_corpus": [ + "exact_byte_publication_persists_submitted_before_io_and_commits_accepted", + "cancelled_submitted_attempt_recovers_unknown_and_retries_exact_bytes_after_reopen", + "publication_outcome_commit_is_idempotent_and_inspection_is_nonmutating", + "publication_execution_binds_live_authority_without_mutating_on_mismatch_or_disable", + "terminal_required_rejection_blocks_the_outbox_without_retry_schedule", + "concurrent_publication_execution_has_one_remote_submitter", + "publication_queue_capacity_is_checked_before_finalization_mutation", + "schema_v8_scans_historical_nullable_job_state_and_installs_permanent_guards" + ], + "source_locked_shared_sqlite_corpus": [ + "every_initialization_durability_edge_fails_once_and_rolls_back", + "transaction_durability_edges_preserve_exact_commit_semantics", + "backup_durability_edges_fail_once_clean_exact_stage_and_recover", + "close_durability_edges_are_once_only_retryable_or_terminal", + "every_marker_and_restore_durability_edge_is_wired_once", + "sigkill_restore_boundaries_recover_exact_topologies_and_preserve_permissions" + ], + "historical_reconciliation_job_guard": { + "migration_target_version": 8, + "scan": "forward_only_fail_closed_before_guard_installation", + "repair_or_delete_invalid_rows": false, + "permanent_guards": [ + "reconciliation_jobs_shape_guard_insert", + "reconciliation_jobs_shape_guard_update" + ], + "required_ready_fields": ["next_attempt_unix_ms"], + "required_leased_fields": ["lease_owner", "lease_expires_unix_ms"], + "terminal_nullable_fields": [ + "next_attempt_unix_ms", + "lease_owner", + "lease_expires_unix_ms" + ] + }, + "resource_bounds": { + "maximum_publication_queue": 65536, + "maximum_targets_per_outbox": 32, + "maximum_attempts_per_target": 100, + "maximum_signed_event_bytes": 32768 + }, + "invariants": { + "exact_committed_bytes_only": true, + "remote_io_outside_sql_transaction": true, + "submitted_before_remote_io": true, + "single_cas_winner": true, + "unknown_outcome_recovered_before_retry": true, + "invalid_historical_state_fails_without_repair": true, + "production_failpoint_surface": false, + "ambient_network": false, + "unbounded_resource": false + }, + "deferred": [ + "rcld_promotion", + "parent_pin_alignment", + "nix", + "oci", + "signing", + "publication", + "deployment" + ] +} diff --git a/contracts/services_hardening/reconciliation_jobs.v1.json b/contracts/services_hardening/reconciliation_jobs.v1.json @@ -42,6 +42,16 @@ "jobs_are_not_deleted": true, "restart_recovers_persisted_schedule_and_lease_expiry": true }, + "state_shape_guard": { + "introduced_schema_version": 8, + "historical_scan": "forward_only_fail_closed", + "invalid_rows_repaired_or_deleted": false, + "insert_trigger": "reconciliation_jobs_shape_guard_insert", + "update_trigger": "reconciliation_jobs_shape_guard_update", + "ready_requires_nonnull_next_attempt": true, + "leased_requires_nonnull_owner_and_expiry": true, + "terminal_requires_null_schedule_and_lease": true + }, "injected_authority": ["wall_time_unix_ms", "lease_owner", "retry_jitter_ms"], "forbidden": [ "ambient_clock", diff --git a/src/lib.rs b/src/lib.rs @@ -186,8 +186,10 @@ pub use state_catalog::{ RHI_STATE_SCHEMA_VERSION_5_SHA256, RHI_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_6_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_6_SHA256, RHI_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_7_OBJECT_COUNT, - RHI_STATE_SCHEMA_VERSION_7_SHA256, RhiStateCatalogError, RhiStateCatalogErrorKind, - rhi_migration_catalog, rhi_schema_catalog, validate_rhi_state_catalogs, + RHI_STATE_SCHEMA_VERSION_7_SHA256, RHI_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256, + RHI_STATE_SCHEMA_VERSION_8_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_8_SHA256, + RhiStateCatalogError, RhiStateCatalogErrorKind, rhi_migration_catalog, rhi_schema_catalog, + validate_rhi_state_catalogs, }; pub use state_config::{ RHI_CONFIG_BINDING_MAX_GENERATIONS, RhiConfigApplyError, RhiConfigApplyErrorKind, diff --git a/src/state_catalog.rs b/src/state_catalog.rs @@ -12,7 +12,7 @@ use radroots_service_sqlite::{ pub const RHI_STATE_BASE_SCHEMA_VERSION: u32 = 1; /// The newest governed RHI state schema understood by this binary. -pub const RHI_STATE_SCHEMA_VERSION: u32 = 7; +pub const RHI_STATE_SCHEMA_VERSION: u32 = 8; /// The shared metadata and migration-ledger objects present at schema v1. pub const RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT: u32 = 6; @@ -35,10 +35,13 @@ pub const RHI_STATE_SCHEMA_VERSION_6_OBJECT_COUNT: u32 = 39; /// The shared objects plus immutable report and publication workflow state. pub const RHI_STATE_SCHEMA_VERSION_7_OBJECT_COUNT: u32 = 63; +/// The shared objects plus fail-closed reconciliation-job shape guards. +pub const RHI_STATE_SCHEMA_VERSION_8_OBJECT_COUNT: u32 = 65; + /// SHA-256 identity of the ordered migration catalog rooted at schema v1. pub const RHI_MIGRATION_CATALOG_SHA256: [u8; 32] = [ - 0xbf, 0x95, 0x58, 0x84, 0xe9, 0x73, 0xde, 0x04, 0xb9, 0x8a, 0x57, 0x98, 0x65, 0x62, 0xef, 0x38, - 0x05, 0xad, 0x82, 0xce, 0x3d, 0xa6, 0xa8, 0x3a, 0xb5, 0x89, 0x0a, 0x50, 0xe0, 0x6b, 0xd5, 0xf4, + 0xe6, 0x5a, 0xd1, 0x15, 0x5c, 0x9d, 0xa2, 0x70, 0x32, 0x81, 0x99, 0x22, 0x68, 0x53, 0x0c, 0x87, + 0xe5, 0x24, 0x8a, 0x52, 0xf7, 0x66, 0x01, 0x1c, 0x23, 0x5e, 0x1f, 0xdb, 0x67, 0x1c, 0x77, 0x4f, ]; /// SHA-256 identity of the exact schema-v1 object snapshot. @@ -119,10 +122,22 @@ pub const RHI_STATE_SCHEMA_VERSION_7_SHA256: [u8; 32] = [ 0xc2, 0x70, 0xef, 0x75, 0x22, 0x67, 0x40, 0x37, 0xde, 0xf7, 0x96, 0x28, 0xa2, 0xb0, 0x39, 0x46, ]; +/// SHA-256 identity of the schema-v8 reconciliation-job shape-guard migration. +pub const RHI_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256: [u8; 32] = [ + 0x6d, 0x3a, 0xa0, 0x6e, 0x69, 0x08, 0xe4, 0x28, 0x1b, 0x50, 0x66, 0xed, 0x4a, 0x13, 0xe2, 0x28, + 0x7b, 0xda, 0xec, 0x05, 0xb0, 0xe5, 0x39, 0x21, 0x58, 0x91, 0x89, 0x2a, 0x5a, 0xfb, 0x58, 0x3b, +]; + +/// SHA-256 identity of the exact schema-v8 object snapshot. +pub const RHI_STATE_SCHEMA_VERSION_8_SHA256: [u8; 32] = [ + 0x7c, 0xef, 0x55, 0x9a, 0xe1, 0xe6, 0xef, 0xe1, 0x58, 0xc5, 0xd1, 0xde, 0x50, 0x11, 0x4e, 0xba, + 0xcc, 0xac, 0x90, 0x53, 0x8e, 0x0c, 0xd4, 0x9a, 0x4e, 0xe2, 0x14, 0xcb, 0x0a, 0x39, 0xc6, 0x18, +]; + /// SHA-256 identity of the schema catalog bound to the migration catalog. pub const RHI_STATE_SCHEMA_CATALOG_SHA256: [u8; 32] = [ - 0xec, 0x31, 0x80, 0x9d, 0x62, 0x07, 0xfd, 0x98, 0xb2, 0x04, 0xe5, 0x69, 0x39, 0x73, 0x11, 0x97, - 0xf2, 0x97, 0x87, 0xbd, 0x1f, 0xef, 0x62, 0x8d, 0x9b, 0x34, 0x5d, 0xb0, 0x20, 0x71, 0x1f, 0xec, + 0x93, 0x30, 0x35, 0x1d, 0x30, 0x0f, 0x70, 0x0f, 0x31, 0x7e, 0xd2, 0xc7, 0x2c, 0xbb, 0xb0, 0x85, + 0x22, 0xa8, 0x27, 0xb9, 0x62, 0xfe, 0x6c, 0xcb, 0x34, 0x3e, 0x3e, 0xe7, 0x1f, 0xdf, 0xd0, 0x7c, ]; macro_rules! rhi_config_bindings_table_sql { @@ -637,6 +652,104 @@ const CREATE_RECONCILIATION_JOBS_MIGRATION_SQL: &str = concat!( ";", ); +macro_rules! reconciliation_jobs_shape_guard_insert_sql { + () => { + r#"CREATE TRIGGER reconciliation_jobs_shape_guard_insert +BEFORE INSERT ON reconciliation_jobs +WHEN (NEW.state = 'ready' AND ( + NEW.attempt_count >= NEW.max_attempts + OR typeof(NEW.next_attempt_unix_ms) != 'integer' + OR NEW.next_attempt_unix_ms NOT BETWEEN 0 AND 9223372036854775807 + OR NEW.lease_owner IS NOT NULL + OR NEW.lease_expires_unix_ms IS NOT NULL + )) + OR (NEW.state = 'leased' AND ( + NEW.attempt_count NOT BETWEEN 1 AND NEW.max_attempts + OR NEW.next_attempt_unix_ms IS NOT NULL + OR typeof(NEW.lease_owner) != 'blob' + OR length(NEW.lease_owner) != 16 + OR typeof(NEW.lease_expires_unix_ms) != 'integer' + OR NEW.lease_expires_unix_ms NOT BETWEEN 1 AND 9223372036854775807 + )) + OR (NEW.state IN ('exhausted', 'superseded', 'completed') AND ( + NEW.next_attempt_unix_ms IS NOT NULL + OR NEW.lease_owner IS NOT NULL + OR NEW.lease_expires_unix_ms IS NOT NULL + )) +BEGIN + SELECT RAISE(ABORT, 'reconciliation job state shape is invalid'); +END"# + }; +} + +macro_rules! reconciliation_jobs_shape_guard_update_sql { + () => { + r#"CREATE TRIGGER reconciliation_jobs_shape_guard_update +BEFORE UPDATE ON reconciliation_jobs +WHEN (NEW.state = 'ready' AND ( + NEW.attempt_count >= NEW.max_attempts + OR typeof(NEW.next_attempt_unix_ms) != 'integer' + OR NEW.next_attempt_unix_ms NOT BETWEEN 0 AND 9223372036854775807 + OR NEW.lease_owner IS NOT NULL + OR NEW.lease_expires_unix_ms IS NOT NULL + )) + OR (NEW.state = 'leased' AND ( + NEW.attempt_count NOT BETWEEN 1 AND NEW.max_attempts + OR NEW.next_attempt_unix_ms IS NOT NULL + OR typeof(NEW.lease_owner) != 'blob' + OR length(NEW.lease_owner) != 16 + OR typeof(NEW.lease_expires_unix_ms) != 'integer' + OR NEW.lease_expires_unix_ms NOT BETWEEN 1 AND 9223372036854775807 + )) + OR (NEW.state IN ('exhausted', 'superseded', 'completed') AND ( + NEW.next_attempt_unix_ms IS NOT NULL + OR NEW.lease_owner IS NOT NULL + OR NEW.lease_expires_unix_ms IS NOT NULL + )) +BEGIN + SELECT RAISE(ABORT, 'reconciliation job state shape is invalid'); +END"# + }; +} + +const CREATE_RECONCILIATION_JOBS_SHAPE_GUARD_INSERT_SQL: &str = + reconciliation_jobs_shape_guard_insert_sql!(); +const CREATE_RECONCILIATION_JOBS_SHAPE_GUARD_UPDATE_SQL: &str = + reconciliation_jobs_shape_guard_update_sql!(); +const CREATE_RECONCILIATION_JOB_SHAPE_GUARDS_MIGRATION_SQL: &str = concat!( + "CREATE TABLE reconciliation_jobs_shape_scan_v1 (\n", + " invalid INTEGER NOT NULL CHECK (invalid = 0)\n", + ") STRICT;\n", + "INSERT INTO reconciliation_jobs_shape_scan_v1 (invalid)\n", + "SELECT 1 FROM reconciliation_jobs\n", + "WHERE (state = 'ready' AND (\n", + " attempt_count >= max_attempts\n", + " OR typeof(next_attempt_unix_ms) != 'integer'\n", + " OR next_attempt_unix_ms NOT BETWEEN 0 AND 9223372036854775807\n", + " OR lease_owner IS NOT NULL\n", + " OR lease_expires_unix_ms IS NOT NULL\n", + " ))\n", + " OR (state = 'leased' AND (\n", + " attempt_count NOT BETWEEN 1 AND max_attempts\n", + " OR next_attempt_unix_ms IS NOT NULL\n", + " OR typeof(lease_owner) != 'blob'\n", + " OR length(lease_owner) != 16\n", + " OR typeof(lease_expires_unix_ms) != 'integer'\n", + " OR lease_expires_unix_ms NOT BETWEEN 1 AND 9223372036854775807\n", + " ))\n", + " OR (state IN ('exhausted', 'superseded', 'completed') AND (\n", + " next_attempt_unix_ms IS NOT NULL\n", + " OR lease_owner IS NOT NULL\n", + " OR lease_expires_unix_ms IS NOT NULL\n", + " ))\n", + "LIMIT 1;\n", + "DROP TABLE reconciliation_jobs_shape_scan_v1;\n", + reconciliation_jobs_shape_guard_insert_sql!(), + ";\n", + reconciliation_jobs_shape_guard_update_sql!(), + ";", +); + macro_rules! evidence_reconciliations_table_sql { () => { r#"CREATE TABLE evidence_reconciliations ( @@ -1328,6 +1441,14 @@ const RECONCILIATION_JOBS_NO_DELETE_SHA256: [u8; 32] = [ 0x1a, 0x86, 0xb1, 0x70, 0x05, 0x05, 0x4c, 0x27, 0x1b, 0xa4, 0x9a, 0xf2, 0x1a, 0x44, 0x9c, 0x9d, 0xdc, 0xfb, 0xbc, 0xd9, 0x13, 0xfa, 0x2a, 0x8e, 0x64, 0x6d, 0x5a, 0x6d, 0x22, 0x69, 0x59, 0xa2, ]; +const RECONCILIATION_JOBS_SHAPE_GUARD_INSERT_SHA256: [u8; 32] = [ + 0x77, 0xf6, 0x83, 0x23, 0x27, 0xce, 0xe0, 0x16, 0xc4, 0x0c, 0x22, 0x6a, 0x61, 0xbe, 0x82, 0xe7, + 0x5b, 0xf0, 0x0b, 0xee, 0x07, 0xb7, 0x03, 0x6f, 0xf7, 0x0d, 0xd2, 0xc9, 0xe0, 0xdb, 0x26, 0xd5, +]; +const RECONCILIATION_JOBS_SHAPE_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0xf7, 0xbb, 0x8e, 0xb8, 0x1a, 0x6a, 0x1c, 0x19, 0xa9, 0x84, 0x67, 0x51, 0x20, 0x30, 0xe3, 0x9c, + 0xb8, 0x0d, 0x51, 0x6c, 0x1e, 0x33, 0xe7, 0x78, 0x8c, 0xfc, 0x37, 0x9e, 0x4e, 0x9e, 0xf4, 0x61, +]; const RHI_CONFIG_BINDINGS_TABLE_SHA256: [u8; 32] = [ 0x4d, 0x6e, 0x8f, 0xff, 0xda, 0x43, 0xe6, 0xf5, 0x3e, 0x23, 0x77, 0xd2, 0x77, 0xa4, 0x52, 0x9e, @@ -1531,6 +1652,13 @@ pub fn rhi_migration_catalog() -> Result<MigrationCatalog, RhiStateCatalogError> MigrationChecksum::from_bytes(RHI_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256), ) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?; + let reconciliation_job_shape_guards = MigrationDescriptor::sql( + 8, + "guard_reconciliation_job_state_shape", + CREATE_RECONCILIATION_JOB_SHAPE_GUARDS_MIGRATION_SQL, + MigrationChecksum::from_bytes(RHI_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256), + ) + .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?; let catalog = MigrationCatalog::new([ configuration, trade_evidence, @@ -1538,10 +1666,11 @@ pub fn rhi_migration_catalog() -> Result<MigrationCatalog, RhiStateCatalogError> reconciliation_jobs, reconciliation_results, report_publication, + reconciliation_job_shape_guards, ]) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?; if catalog.current_version() != RHI_STATE_SCHEMA_VERSION - || catalog.descriptors().len() != 6 + || catalog.descriptors().len() != 7 || catalog.digest().as_bytes() != &RHI_MIGRATION_CATALOG_SHA256 { return Err(RhiStateCatalogError::new( @@ -1591,11 +1720,17 @@ pub fn rhi_schema_catalog() -> Result<SchemaCatalog, RhiStateCatalogError> { ) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?; let version_seven = SchemaVersionCatalog::new( - RHI_STATE_SCHEMA_VERSION, + 7, rhi_schema_version_seven_objects()?, SchemaDigest::from_bytes(RHI_STATE_SCHEMA_VERSION_7_SHA256), ) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?; + let version_eight = SchemaVersionCatalog::new( + RHI_STATE_SCHEMA_VERSION, + rhi_schema_version_eight_objects()?, + SchemaDigest::from_bytes(RHI_STATE_SCHEMA_VERSION_8_SHA256), + ) + .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?; let catalog = SchemaCatalog::new( &migrations, [ @@ -1606,6 +1741,7 @@ pub fn rhi_schema_catalog() -> Result<SchemaCatalog, RhiStateCatalogError> { version_five, version_six, version_seven, + version_eight, ], ) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?; @@ -1620,10 +1756,10 @@ pub fn validate_rhi_state_catalogs( ) -> Result<(), RhiStateCatalogError> { let versions = schema.versions(); let valid = migrations.current_version() == RHI_STATE_SCHEMA_VERSION - && migrations.descriptors().len() == 6 + && migrations.descriptors().len() == 7 && migrations.digest().as_bytes() == &RHI_MIGRATION_CATALOG_SHA256 && schema.migration_catalog_digest() == migrations.digest() - && versions.len() == 7 + && versions.len() == 8 && versions[0].version() == RHI_STATE_BASE_SCHEMA_VERSION && versions[0].object_count() == RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT && versions[0].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_1_SHA256 @@ -1642,9 +1778,12 @@ pub fn validate_rhi_state_catalogs( && versions[5].version() == 6 && versions[5].object_count() == RHI_STATE_SCHEMA_VERSION_6_OBJECT_COUNT && versions[5].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_6_SHA256 - && versions[6].version() == RHI_STATE_SCHEMA_VERSION + && versions[6].version() == 7 && versions[6].object_count() == RHI_STATE_SCHEMA_VERSION_7_OBJECT_COUNT && versions[6].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_7_SHA256 + && versions[7].version() == RHI_STATE_SCHEMA_VERSION + && versions[7].object_count() == RHI_STATE_SCHEMA_VERSION_8_OBJECT_COUNT + && versions[7].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_8_SHA256 && schema.digest().as_bytes() == &RHI_STATE_SCHEMA_CATALOG_SHA256; if valid { Ok(()) @@ -1722,6 +1861,37 @@ fn rhi_schema_version_seven_objects() -> Result<Vec<SchemaObject>, RhiStateCatal Ok(objects) } +fn rhi_schema_version_eight_objects() -> Result<Vec<SchemaObject>, RhiStateCatalogError> { + let mut objects = rhi_schema_version_seven_objects()?; + objects.extend(rhi_reconciliation_job_shape_guard_objects()?); + Ok(objects) +} + +fn rhi_reconciliation_job_shape_guard_objects() -> Result<[SchemaObject; 2], RhiStateCatalogError> { + let object = |name, sql, digest| { + SchemaObject::new( + SchemaObjectKind::Trigger, + name, + "reconciliation_jobs", + sql, + SchemaDigest::from_bytes(digest), + ) + .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog)) + }; + Ok([ + object( + "reconciliation_jobs_shape_guard_insert", + CREATE_RECONCILIATION_JOBS_SHAPE_GUARD_INSERT_SQL, + RECONCILIATION_JOBS_SHAPE_GUARD_INSERT_SHA256, + )?, + object( + "reconciliation_jobs_shape_guard_update", + CREATE_RECONCILIATION_JOBS_SHAPE_GUARD_UPDATE_SQL, + RECONCILIATION_JOBS_SHAPE_GUARD_UPDATE_SHA256, + )?, + ]) +} + fn rhi_report_publication_objects() -> Result<[SchemaObject; 24], RhiStateCatalogError> { let object = |kind, name, table_name, sql, digest| { SchemaObject::new( diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs @@ -18,6 +18,8 @@ const PUBLICATION_ATTEMPT_CONTRACT: &str = const PUBLICATION_EXECUTION: &str = include_str!("../src/publication_execution.rs"); const PUBLICATION_EXECUTION_CONTRACT: &str = include_str!("../contracts/services_hardening/publication_execution.v1.json"); +const PUBLICATION_WAVE_QUALIFICATION_CONTRACT: &str = + include_str!("../contracts/services_hardening/publication_wave_qualification.v1.json"); const PUBLICATION_CONTRACT: &str = include_str!("../contracts/services_hardening/publication_outbox.v1.json"); const PUBLICATION_SUBMISSION: &str = include_str!("../src/publication_submission.rs"); @@ -499,6 +501,15 @@ fn publication_execution_is_sqlx_owned_exact_byte_and_fail_closed() { } assert!(!ROOT.contains("pub mod publication_execution")); assert!(!PUBLIC_API.contains("rhi::publication_execution::")); + let qualification: serde_json::Value = + serde_json::from_str(PUBLICATION_WAVE_QUALIFICATION_CONTRACT) + .expect("publication-wave qualification contract"); + assert_eq!(qualification["step"], 203); + assert_eq!(qualification["state_schema_version"], 8); + assert_eq!( + qualification["invariants"]["production_failpoint_surface"], + false + ); } #[test] @@ -1227,6 +1238,7 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() { "[`publication_attempt_evidence.v1.json`](contracts/services_hardening/publication_attempt_evidence.v1.json)", "## Durable exact-byte publication execution", "[`publication_execution.v1.json`](contracts/services_hardening/publication_execution.v1.json)", + "[`publication_wave_qualification.v1.json`](contracts/services_hardening/publication_wave_qualification.v1.json)", "The event body and exact kind-3441 structural tags", "without rebuilding, reserializing, or", "Coverage is exactly `Missing`, `Partial`, `ScopeSatisfied`, or `Unsupported`", diff --git a/tests/services_hardening_publication_contract.rs b/tests/services_hardening_publication_contract.rs @@ -4,8 +4,7 @@ use std::error::Error; use rhi::{ RHI_PUBLICATION_CONTRACT_VERSION, RHI_PUBLICATION_MAX_ATTEMPTS, RHI_PUBLICATION_MAX_TARGETS, - RHI_STATE_SCHEMA_VERSION, RhiConfigProfile, RhiPublicationAuthority, RhiPublicationMode, - parse_rhi_config_v1, + RhiConfigProfile, RhiPublicationAuthority, RhiPublicationMode, parse_rhi_config_v1, }; use serde_json::json; @@ -25,7 +24,10 @@ fn machine_contract_freezes_step_198_authority_and_schema() { contract["contract_version"], RHI_PUBLICATION_CONTRACT_VERSION ); - assert_eq!(contract["state_schema_version"], RHI_STATE_SCHEMA_VERSION); + // This Step198 contract records the schema version that introduced the + // outbox. Later forward-only migrations are governed by their own + // contracts and must not silently rewrite this historical evidence. + assert_eq!(contract["state_schema_version"], 7); assert_eq!( contract["publication_modes"], json!(["required", "disabled"]) diff --git a/tests/services_hardening_publication_wave_qualification_contract.rs b/tests/services_hardening_publication_wave_qualification_contract.rs @@ -0,0 +1,119 @@ +#![forbid(unsafe_code)] + +use serde_json::json; + +const CONTRACT: &str = + include_str!("../contracts/services_hardening/publication_wave_qualification.v1.json"); +const CATALOG_SOURCE: &str = include_str!("../src/state_catalog.rs"); +const EXECUTION_SOURCE: &str = include_str!("../src/publication_execution.rs"); +const README: &str = include_str!("../README"); + +#[test] +fn step203_machine_contract_freezes_the_complete_qualification_boundary() { + let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("contract"); + assert_eq!( + contract["schema"], + "radroots.rhi.publication-wave-qualification" + ); + assert_eq!(contract["schema_version"], 1); + assert_eq!(contract["contract_version"], 1); + assert_eq!(contract["step"], 203); + assert_eq!(contract["service"], "rhi"); + assert_eq!(contract["state_schema_version"], 8); + assert_eq!( + contract["component_corpus"], + json!([ + "exact_byte_publication_persists_submitted_before_io_and_commits_accepted", + "cancelled_submitted_attempt_recovers_unknown_and_retries_exact_bytes_after_reopen", + "publication_outcome_commit_is_idempotent_and_inspection_is_nonmutating", + "publication_execution_binds_live_authority_without_mutating_on_mismatch_or_disable", + "terminal_required_rejection_blocks_the_outbox_without_retry_schedule", + "concurrent_publication_execution_has_one_remote_submitter", + "publication_queue_capacity_is_checked_before_finalization_mutation", + "schema_v8_scans_historical_nullable_job_state_and_installs_permanent_guards" + ]) + ); + assert_eq!( + contract["source_locked_shared_sqlite_corpus"], + json!([ + "every_initialization_durability_edge_fails_once_and_rolls_back", + "transaction_durability_edges_preserve_exact_commit_semantics", + "backup_durability_edges_fail_once_clean_exact_stage_and_recover", + "close_durability_edges_are_once_only_retryable_or_terminal", + "every_marker_and_restore_durability_edge_is_wired_once", + "sigkill_restore_boundaries_recover_exact_topologies_and_preserve_permissions" + ]) + ); + assert_eq!( + contract["historical_reconciliation_job_guard"], + json!({ + "migration_target_version": 8, + "scan": "forward_only_fail_closed_before_guard_installation", + "repair_or_delete_invalid_rows": false, + "permanent_guards": [ + "reconciliation_jobs_shape_guard_insert", + "reconciliation_jobs_shape_guard_update" + ], + "required_ready_fields": ["next_attempt_unix_ms"], + "required_leased_fields": ["lease_owner", "lease_expires_unix_ms"], + "terminal_nullable_fields": [ + "next_attempt_unix_ms", + "lease_owner", + "lease_expires_unix_ms" + ] + }) + ); + assert_eq!( + contract["resource_bounds"], + json!({ + "maximum_publication_queue": 65_536, + "maximum_targets_per_outbox": 32, + "maximum_attempts_per_target": 100, + "maximum_signed_event_bytes": 32_768 + }) + ); + for invariant in [ + "exact_committed_bytes_only", + "remote_io_outside_sql_transaction", + "submitted_before_remote_io", + "single_cas_winner", + "unknown_outcome_recovered_before_retry", + "invalid_historical_state_fails_without_repair", + ] { + assert_eq!(contract["invariants"][invariant], true, "{invariant}"); + } + assert_eq!( + contract["invariants"]["production_failpoint_surface"], + false + ); + assert_eq!(contract["invariants"]["ambient_network"], false); + assert_eq!(contract["invariants"]["unbounded_resource"], false); +} + +#[test] +fn qualification_remains_private_sqlx_owned_and_documented() { + for required in [ + "CREATE TABLE reconciliation_jobs_shape_scan_v1", + "DROP TABLE reconciliation_jobs_shape_scan_v1", + "CREATE TRIGGER reconciliation_jobs_shape_guard_insert", + "CREATE TRIGGER reconciliation_jobs_shape_guard_update", + "typeof(next_attempt_unix_ms) != 'integer'", + "typeof(lease_owner) != 'blob'", + "typeof(lease_expires_unix_ms) != 'integer'", + ] { + assert!(CATALOG_SOURCE.contains(required), "missing `{required}`"); + } + for forbidden in [ + "RHI_FAILPOINT", + "RHI_TEST_", + "production_failpoint", + "pub fn sqlite", + "pub fn connection", + "pub fn transaction", + ] { + assert!(!EXECUTION_SOURCE.contains(forbidden), "found `{forbidden}`"); + } + assert!(README.contains( + "[`publication_wave_qualification.v1.json`](contracts/services_hardening/publication_wave_qualification.v1.json)" + )); +} diff --git a/tests/services_hardening_reconciliation_job_contract.rs b/tests/services_hardening_reconciliation_job_contract.rs @@ -43,6 +43,19 @@ fn machine_contract_freezes_the_complete_step_187_boundary() { contract["injected_authority"], json!(["wall_time_unix_ms", "lease_owner", "retry_jitter_ms"]) ); + assert_eq!( + contract["state_shape_guard"], + json!({ + "introduced_schema_version": 8, + "historical_scan": "forward_only_fail_closed", + "invalid_rows_repaired_or_deleted": false, + "insert_trigger": "reconciliation_jobs_shape_guard_insert", + "update_trigger": "reconciliation_jobs_shape_guard_update", + "ready_requires_nonnull_next_attempt": true, + "leased_requires_nonnull_owner_and_expiry": true, + "terminal_requires_null_schedule_and_lease": true + }) + ); for forbidden in [ "ambient_clock", "ambient_entropy", diff --git a/tests/services_hardening_reconciliation_jobs.rs b/tests/services_hardening_reconciliation_jobs.rs @@ -44,6 +44,7 @@ use rhi::{ }; use sha2::{Digest, Sha256}; use sqlx::{Connection, SqliteConnection, sqlite::SqliteConnectOptions}; +use tokio::sync::Notify; const EXAMPLE: &str = include_str!("../contracts/services_hardening/config.v1.example.toml"); const TRADE_VECTOR: &str = @@ -1667,6 +1668,28 @@ impl RhiExactPublicationSink for PendingExactPublicationSink { } } +struct CoordinatedExactPublicationSink { + expected: Arc<Vec<u8>>, + calls: Arc<AtomicUsize>, + started: Arc<Notify>, + release: Arc<Notify>, +} + +impl RhiExactPublicationSink for CoordinatedExactPublicationSink { + fn submit_exact<'a>( + &'a self, + attempt: &'a RhiPreparedPublicationAttempt, + ) -> BoxFuture<'a, RhiPublicationAttemptOutcome> { + Box::pin(async move { + assert_eq!(attempt.exact_signed_event_bytes(), self.expected.as_slice()); + self.calls.fetch_add(1, Ordering::SeqCst); + self.started.notify_one(); + self.release.notified().await; + RhiPublicationAttemptOutcome::Accepted + }) + } +} + fn publication_now(value: u64) -> RhiPublicationUnixMilliseconds { RhiPublicationUnixMilliseconds::new(value).expect("publication time") } @@ -2217,6 +2240,180 @@ WHERE outbox.outbox_id = ?"#, } #[tokio::test] +async fn concurrent_publication_execution_has_one_remote_submitter() { + let ( + root, + runtime, + metadata, + configuration, + host, + lease, + signed, + publication, + manifest, + _trade, + ) = signed_finalization_fixture("publication-execution-concurrency", false).await; + let repositories = host.repositories(); + repositories + .reconciliation_attempts() + .commit_finalization(&signed, &publication, now(1_784_347_208_000)) + .await + .expect("finalization"); + let calls = Arc::new(AtomicUsize::new(0)); + let started = Arc::new(Notify::new()); + let release = Arc::new(Notify::new()); + let sink = CoordinatedExactPublicationSink { + expected: Arc::new(signed.signed_event_bytes().to_vec()), + calls: Arc::clone(&calls), + started: Arc::clone(&started), + release: Arc::clone(&release), + }; + let first_adapters = publication_adapters(1_784_347_208); + let second_adapters = publication_adapters(1_784_347_208); + let first_outbox = repositories.publication_outbox(); + let first = first_outbox.execute_next_publication( + publication_owner(0xb1), + &first_adapters, + &sink, + &publication, + ); + let second = async { + started.notified().await; + let result = repositories + .publication_outbox() + .execute_next_publication( + publication_owner(0xb2), + &second_adapters, + &sink, + &publication, + ) + .await; + release.notify_one(); + result + }; + let (first, second) = tokio::join!(first, second); + let first = first + .expect("first executor") + .expect("first executor claimed due publication"); + assert_eq!(first.outcome(), RhiPublicationAttemptOutcome::Accepted); + assert!(second.expect("second executor").is_none()); + assert_eq!(calls.load(Ordering::SeqCst), 1); + + host.close().await.expect("host close"); + drop(( + manifest, + publication, + signed, + lease, + configuration, + metadata, + runtime, + root, + )); +} + +#[tokio::test] +async fn publication_queue_capacity_is_checked_before_finalization_mutation() { + let expected_public_key = + Keys::new(SecretKey::from_slice(&attestation_secret()).expect("identity secret")) + .public_key() + .to_hex(); + let (root, runtime, metadata, configuration, host, lease, evaluation) = + finalization_fixture_with_source("publication-queue-capacity", |runtime| { + attestation_configuration(runtime, &expected_public_key).replacen( + "publication = 4096", + "publication = 1", + 1, + ) + }) + .await; + let identity = provision_attestation_identity(&runtime, &configuration, &metadata); + let fence = host + .repositories() + .reconciliation_attempts() + .prepare_finalization(lease, evaluation, now(1_784_347_206_000)) + .await + .expect("finalization fence"); + let signed = build_rhi_signed_evidence_attestation( + fence, + &identity, + UnixTimeSeconds::new(1_784_347_207), + &FixedAttestationEntropy(0xb3), + None, + ) + .expect("signed attestation"); + let publication = + RhiPublicationAuthority::from_config(&configuration).expect("publication authority"); + assert_eq!(publication.queue_capacity(), 1); + + let options = SqliteConnectOptions::new() + .filename(runtime.artifacts().state_database()) + .create_if_missing(false) + .foreign_keys(false); + let mut connection = SqliteConnection::connect_with(&options) + .await + .expect("offline queue fixture connection"); + sqlx::query( + r#"INSERT INTO publication_outbox ( + outbox_id, event_id, event_sha256, publication_authority_sha256, + target_set_sha256, target_count, required_target_count, + max_attempts, initial_backoff_ms, maximum_backoff_ms, + attempt_deadline_ms, state, revision, next_attempt_unix_ms, + lease_owner, lease_expires_unix_ms, created_at_unix_ms, + updated_at_unix_ms + ) VALUES (?, ?, ?, ?, ?, 1, 1, 3, 100, 1000, 5000, + 'pending', 1, 1, NULL, NULL, 1, 1)"#, + ) + .bind([0xc1_u8; 32].as_slice()) + .bind([0xc2_u8; 32].as_slice()) + .bind([0xc3_u8; 32].as_slice()) + .bind([0xc4_u8; 32].as_slice()) + .bind([0xc5_u8; 32].as_slice()) + .execute(&mut connection) + .await + .expect("bounded active outbox fixture"); + connection.close().await.expect("queue fixture close"); + + let error = host + .repositories() + .reconciliation_attempts() + .commit_finalization(&signed, &publication, now(1_784_347_208_000)) + .await + .expect_err("full publication queue"); + assert_eq!( + error.kind(), + RhiReconciliationFinalizationCommitErrorKind::PublicationQueueFull + ); + assert!(Error::source(&error).is_none()); + let mut connection = fixture_connection(&runtime).await; + let counts: (i64, i64, i64, i64, i64, String) = sqlx::query_as( + r#"SELECT + (SELECT COUNT(*) FROM evidence_manifests), + (SELECT COUNT(*) FROM trade_projections), + (SELECT COUNT(*) FROM attestation_reports), + (SELECT COUNT(*) FROM signed_attestation_events), + (SELECT COUNT(*) FROM publication_outbox), + (SELECT state FROM reconciliation_jobs WHERE job_id = ?)"#, + ) + .bind(lease.job().id().as_bytes().as_slice()) + .fetch_one(&mut connection) + .await + .expect("no finalization mutation"); + assert_eq!(counts, (0, 0, 0, 0, 1, "leased".into())); + connection.close().await.expect("verification close"); + host.close().await.expect("host close"); + drop(( + signed, + identity, + publication, + configuration, + metadata, + runtime, + root, + )); +} + +#[tokio::test] async fn cancelled_submitted_attempt_recovers_unknown_and_retries_exact_bytes_after_reopen() { let ( root, diff --git a/tests/services_hardening_state_catalog.rs b/tests/services_hardening_state_catalog.rs @@ -18,8 +18,10 @@ use rhi::{ RHI_STATE_SCHEMA_VERSION_5_SHA256, RHI_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_6_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_6_SHA256, RHI_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_7_OBJECT_COUNT, - RHI_STATE_SCHEMA_VERSION_7_SHA256, RhiStateCatalogErrorKind, rhi_migration_catalog, - rhi_schema_catalog, validate_rhi_state_catalogs, + RHI_STATE_SCHEMA_VERSION_7_SHA256, RHI_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256, + RHI_STATE_SCHEMA_VERSION_8_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_8_SHA256, + RhiStateCatalogErrorKind, rhi_migration_catalog, rhi_schema_catalog, + validate_rhi_state_catalogs, }; const CATALOG_SOURCE: &str = include_str!("../src/state_catalog.rs"); @@ -27,14 +29,14 @@ const LIB_SOURCE: &str = include_str!("../src/lib.rs"); const MANIFEST: &str = include_str!("../Cargo.toml"); #[test] -fn schema_v1_through_v7_catalogs_have_exact_literal_identities() { +fn schema_v1_through_v8_catalogs_have_exact_literal_identities() { let migrations = rhi_migration_catalog().expect("RHI migration catalog"); let schema = rhi_schema_catalog().expect("RHI schema catalog"); assert_eq!(RHI_STATE_BASE_SCHEMA_VERSION, 1); - assert_eq!(RHI_STATE_SCHEMA_VERSION, 7); - assert_eq!(migrations.descriptors().len(), 6); - assert_eq!(migrations.current_version(), 7); + assert_eq!(RHI_STATE_SCHEMA_VERSION, 8); + assert_eq!(migrations.descriptors().len(), 7); + assert_eq!(migrations.current_version(), 8); assert_eq!(migrations.descriptors()[0].target_version(), 2); assert_eq!( migrations.descriptors()[0].name().as_str(), @@ -89,12 +91,21 @@ fn schema_v1_through_v7_catalogs_have_exact_literal_identities() { migrations.descriptors()[5].checksum().as_bytes(), &RHI_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256 ); + assert_eq!(migrations.descriptors()[6].target_version(), 8); + assert_eq!( + migrations.descriptors()[6].name().as_str(), + "guard_reconciliation_job_state_shape" + ); + assert_eq!( + migrations.descriptors()[6].checksum().as_bytes(), + &RHI_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256 + ); assert_eq!( migrations.digest().as_bytes(), &RHI_MIGRATION_CATALOG_SHA256 ); - assert_eq!(schema.versions().len(), 7); + assert_eq!(schema.versions().len(), 8); let version = schema.versions()[0]; assert_eq!(version.version(), 1); assert_eq!( @@ -172,13 +183,24 @@ fn schema_v1_through_v7_catalogs_have_exact_literal_identities() { version.digest().as_bytes(), &RHI_STATE_SCHEMA_VERSION_7_SHA256 ); + let version = schema.versions()[7]; + assert_eq!(version.version(), 8); + assert_eq!( + version.object_count(), + RHI_STATE_SCHEMA_VERSION_8_OBJECT_COUNT + ); + assert_eq!(version.object_count(), 65); + assert_eq!( + version.digest().as_bytes(), + &RHI_STATE_SCHEMA_VERSION_8_SHA256 + ); assert_eq!(schema.digest().as_bytes(), &RHI_STATE_SCHEMA_CATALOG_SHA256); assert_eq!(schema.migration_catalog_digest(), migrations.digest()); validate_rhi_state_catalogs(&migrations, &schema).expect("exact catalogs"); assert_eq!( lower_hex(&RHI_MIGRATION_CATALOG_SHA256), - "bf955884e973de04b98a57986562ef3805ad82ce3da6a83ab5890a50e06bd5f4" + "e65ad1155c9da2703281992268530c87e5248a52f766011c235e1fdb671c774f" ); assert_eq!( lower_hex(&RHI_STATE_SCHEMA_VERSION_1_SHA256), @@ -233,8 +255,16 @@ fn schema_v1_through_v7_catalogs_have_exact_literal_identities() { "840aa83c689f9df99d26c5b4ef117131c270ef7522674037def79628a2b03946" ); assert_eq!( + lower_hex(&RHI_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256), + "6d3aa06e6908e4281b5066ed4a13e2287bdaec05b0e539215891892a5afb583b" + ); + assert_eq!( + lower_hex(&RHI_STATE_SCHEMA_VERSION_8_SHA256), + "7cef559ae1e6efe158c5d1de50114ebaccac90538e0cd49a4ee214cb0a39c618" + ); + assert_eq!( lower_hex(&RHI_STATE_SCHEMA_CATALOG_SHA256), - "ec31809d6207fd98b204e56939731197f29787bd1fef628d9b345db020711fec" + "9330351d300f700f317ed2c72cbbb08522a827b962fe6ccb343e3ee71fdfd07c" ); } @@ -294,6 +324,10 @@ fn independent_validator_rejects_migration_or_schema_drift() { SchemaVersionCatalog::computed_digest(7, [version_two_object()]).expect("v7 digest"); let version_seven = SchemaVersionCatalog::new(7, [version_two_object()], snapshot_digest) .expect("version seven"); + let snapshot_digest = + SchemaVersionCatalog::computed_digest(8, [version_two_object()]).expect("v8 digest"); + let version_eight = SchemaVersionCatalog::new(8, [version_two_object()], snapshot_digest) + .expect("version eight"); let schema = SchemaCatalog::new( &exact_migrations, [ @@ -304,6 +338,7 @@ fn independent_validator_rejects_migration_or_schema_drift() { version_five, version_six, version_seven, + version_eight, ], ) .expect("drift schema catalog"); @@ -362,6 +397,10 @@ fn catalog_errors_are_stable_source_free_and_redacted() { SchemaVersionCatalog::computed_digest(7, [secret_object()]).expect("v7 digest"); let version_seven = SchemaVersionCatalog::new(7, [secret_object()], version_seven_digest) .expect("version seven"); + let version_eight_digest = + SchemaVersionCatalog::computed_digest(8, [secret_object()]).expect("v8 digest"); + let version_eight = SchemaVersionCatalog::new(8, [secret_object()], version_eight_digest) + .expect("version eight"); let schema = SchemaCatalog::new( &migrations, [ @@ -372,6 +411,7 @@ fn catalog_errors_are_stable_source_free_and_redacted() { version_five, version_six, version_seven, + version_eight, ], ) .expect("schema catalog"); diff --git a/tests/services_hardening_state_host.rs b/tests/services_hardening_state_host.rs @@ -75,6 +75,100 @@ fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentit (applied_at, build) } +async fn offline_connection(runtime: &rhi::RhiRuntimeContext) -> SqliteConnection { + let options = SqliteConnectOptions::new() + .filename(runtime.artifacts().state_database()) + .create_if_missing(false) + .foreign_keys(false); + SqliteConnection::connect_with(&options) + .await + .expect("offline fixture connection") +} + +async fn downgrade_fixture_to_schema_v7(runtime: &rhi::RhiRuntimeContext) { + let mut connection = offline_connection(runtime).await; + let metadata_guard: String = sqlx::query_scalar( + "SELECT sql FROM sqlite_schema WHERE type = 'trigger' AND name = 'radroots_service_metadata_guard_update'", + ) + .fetch_one(&mut connection) + .await + .expect("metadata guard SQL"); + let migration_no_update: String = sqlx::query_scalar( + "SELECT sql FROM sqlite_schema WHERE type = 'trigger' AND name = 'schema_migrations_no_update'", + ) + .fetch_one(&mut connection) + .await + .expect("migration update guard SQL"); + let migration_no_delete: String = sqlx::query_scalar( + "SELECT sql FROM sqlite_schema WHERE type = 'trigger' AND name = 'schema_migrations_no_delete'", + ) + .fetch_one(&mut connection) + .await + .expect("migration delete guard SQL"); + for statement in [ + "DROP TRIGGER reconciliation_jobs_shape_guard_insert", + "DROP TRIGGER reconciliation_jobs_shape_guard_update", + "DROP TRIGGER radroots_service_metadata_guard_update", + "DROP TRIGGER schema_migrations_no_update", + "DROP TRIGGER schema_migrations_no_delete", + "UPDATE radroots_service_metadata SET state_schema_version = 7 WHERE singleton = 1", + "DELETE FROM schema_migrations WHERE version = 8", + ] { + sqlx::query(statement) + .execute(&mut connection) + .await + .expect("downgrade exact v8 fixture state"); + } + for statement in [metadata_guard, migration_no_update, migration_no_delete] { + sqlx::query(sqlx::AssertSqlSafe(statement.as_str())) + .execute(&mut connection) + .await + .expect("restore shared immutable guard"); + } + connection.close().await.expect("downgrade fixture close"); +} + +async fn insert_historical_reconciliation_job( + runtime: &rhi::RhiRuntimeContext, + state: &str, + next_attempt_unix_ms: Option<i64>, + lease_owner: Option<Vec<u8>>, + lease_expires_unix_ms: Option<i64>, +) { + let mut connection = offline_connection(runtime).await; + let trade_id = [0x51_u8; 16]; + sqlx::query( + "INSERT INTO trade_dirty_generations (trade_id, generation, evidence_policy_sha256, updated_at_unix_s) VALUES (?, 1, ?, 1)", + ) + .bind(trade_id.as_slice()) + .bind([0x52_u8; 32].as_slice()) + .execute(&mut connection) + .await + .expect("historical dirty generation"); + sqlx::query( + r#"INSERT INTO reconciliation_jobs ( + job_id, trade_id, input_generation, evidence_policy_sha256, + state, revision, attempt_count, failure_count, max_attempts, + lease_duration_ms, lease_renewal_ms, initial_backoff_ms, + maximum_backoff_ms, next_attempt_unix_ms, lease_owner, + lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms + ) VALUES (?, ?, 1, ?, ?, 1, ?, 0, 3, 30000, 10000, 100, + 1000, ?, ?, ?, 1, 1)"#, + ) + .bind([0x53_u8; 32].as_slice()) + .bind(trade_id.as_slice()) + .bind([0x52_u8; 32].as_slice()) + .bind(state) + .bind(i64::from(state == "leased")) + .bind(next_attempt_unix_ms) + .bind(lease_owner) + .bind(lease_expires_unix_ms) + .execute(&mut connection) + .await + .expect("historical nullable reconciliation row admitted by schema v7"); + connection.close().await.expect("historical fixture close"); +} + #[tokio::test] async fn initialize_is_create_new_and_both_existing_open_modes_close_explicitly() { let directory = tempfile::tempdir().expect("temporary root"); @@ -305,6 +399,125 @@ async fn publication_schema_rejects_null_state_holes_and_accepted_target_mutatio } #[tokio::test] +async fn schema_v8_scans_historical_nullable_job_state_and_installs_permanent_guards() { + for (instance, state, next_attempt, owner, expiry) in [ + ("missing-ready-time", "ready", None, None, None), + ("missing-lease-owner", "leased", None, None, Some(30_001)), + ( + "missing-lease-expiry", + "leased", + None, + Some(vec![0x61; 16]), + None, + ), + ] { + let directory = tempfile::tempdir().expect("temporary root"); + let runtime = runtime(directory.path(), instance); + prepare_state_directory(&runtime); + let metadata = metadata(&runtime); + let (applied_at, build) = migration_evidence(); + initialize_rhi_state(&runtime, &metadata, applied_at, &build) + .await + .expect("schema-v8 initialization"); + downgrade_fixture_to_schema_v7(&runtime).await; + insert_historical_reconciliation_job(&runtime, state, next_attempt, owner, expiry).await; + + let error = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build) + .await + .expect_err("invalid historical row must block migration"); + assert_eq!(error.kind(), RhiStateHostErrorKind::ReadWriteOpen); + let mut connection = offline_connection(&runtime).await; + let durable: (i64, i64, i64, i64) = sqlx::query_as( + r#"SELECT + (SELECT state_schema_version FROM radroots_service_metadata WHERE singleton = 1), + (SELECT COUNT(*) FROM schema_migrations WHERE version = 8), + (SELECT COUNT(*) FROM sqlite_schema WHERE type = 'trigger' + AND name IN ('reconciliation_jobs_shape_guard_insert', + 'reconciliation_jobs_shape_guard_update')), + (SELECT COUNT(*) FROM reconciliation_jobs)"#, + ) + .fetch_one(&mut connection) + .await + .expect("failed migration state"); + assert_eq!(durable, (7, 0, 0, 1)); + connection.close().await.expect("failed fixture close"); + } + + let directory = tempfile::tempdir().expect("temporary root"); + let runtime = runtime(directory.path(), "valid-schema-v7"); + prepare_state_directory(&runtime); + let metadata = metadata(&runtime); + let (applied_at, build) = migration_evidence(); + initialize_rhi_state(&runtime, &metadata, applied_at, &build) + .await + .expect("schema-v8 initialization"); + downgrade_fixture_to_schema_v7(&runtime).await; + insert_historical_reconciliation_job( + &runtime, + "leased", + None, + Some(vec![0x62; 16]), + Some(30_001), + ) + .await; + let writer = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build) + .await + .expect("valid schema-v7 prefix migrates"); + writer.close().await.expect("migrated writer close"); + + let mut connection = offline_connection(&runtime).await; + let migrated: (i64, i64, i64, i64) = sqlx::query_as( + r#"SELECT + (SELECT state_schema_version FROM radroots_service_metadata WHERE singleton = 1), + (SELECT COUNT(*) FROM schema_migrations WHERE version = 8), + (SELECT COUNT(*) FROM sqlite_schema WHERE type = 'trigger' + AND name IN ('reconciliation_jobs_shape_guard_insert', + 'reconciliation_jobs_shape_guard_update')), + (SELECT COUNT(*) FROM sqlite_schema + WHERE name = 'reconciliation_jobs_shape_scan_v1')"#, + ) + .fetch_one(&mut connection) + .await + .expect("migrated schema state"); + assert_eq!(migrated, (8, 1, 2, 0)); + + let invalid_insert = sqlx::query( + r#"INSERT INTO reconciliation_jobs ( + job_id, trade_id, input_generation, evidence_policy_sha256, + state, revision, attempt_count, failure_count, max_attempts, + lease_duration_ms, lease_renewal_ms, initial_backoff_ms, + maximum_backoff_ms, next_attempt_unix_ms, lease_owner, + lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms + ) VALUES (?, ?, 1, ?, 'ready', 1, 0, 0, 3, 30000, 10000, + 100, 1000, NULL, NULL, NULL, 1, 1)"#, + ) + .bind([0x63_u8; 32].as_slice()) + .bind([0x64_u8; 16].as_slice()) + .bind([0x65_u8; 32].as_slice()) + .execute(&mut connection) + .await; + assert!( + invalid_insert.is_err(), + "insert guard rejects a ready NULL hole" + ); + + let invalid_update = sqlx::query( + r#"UPDATE reconciliation_jobs + SET revision = revision + 1, updated_at_unix_ms = updated_at_unix_ms + 1, + lease_owner = NULL + WHERE job_id = ?"#, + ) + .bind([0x53_u8; 32].as_slice()) + .execute(&mut connection) + .await; + assert!( + invalid_update.is_err(), + "update guard rejects a leased NULL hole" + ); + connection.close().await.expect("guard fixture close"); +} + +#[tokio::test] async fn missing_state_and_mismatched_evidence_fail_before_database_creation() { let directory = tempfile::tempdir().expect("temporary root"); let primary = runtime(directory.path(), "primary"); diff --git a/tests/services_hardening_state_resilience.rs b/tests/services_hardening_state_resilience.rs @@ -347,7 +347,7 @@ async fn exact_open_rejects_unexpected_migration_history_without_repair() { service_version, service_commit, lib_revision, rust_version, target, feature_profile, config_contract_version, state_contract_version, admin_contract_version, status_contract_version, provider_contract_version - ) VALUES (8, 'unexpected_schema', ?, 1725000000, '0.1.0', ?, ?, + ) VALUES (9, 'unexpected_schema', ?, 1725000000, '0.1.0', ?, ?, 'rustc-test', 'test-target', 'service-host', 1, 7, 1, 1, 1)", ) .bind([0x44_u8; 32].as_slice())