services_hardening_reconciliation_jobs.rs (120402B)
1 #![forbid(unsafe_code)] 2 #![cfg(any(target_os = "linux", target_os = "macos"))] 3 4 use std::{ 5 error::Error, 6 fs, 7 os::unix::fs::PermissionsExt, 8 path::Path, 9 sync::{ 10 Arc, 11 atomic::{AtomicUsize, Ordering}, 12 }, 13 time::Duration, 14 }; 15 16 use nostr::{Keys, SecretKey}; 17 use radroots_service_host::{ 18 EntropyError, EntropySource, SystemMonotonicClock, UnixTimeSeconds, WallClock, WallClockError, 19 }; 20 use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity}; 21 use radroots_storage::event::SourceGeneration; 22 use radroots_transport::BoxFuture; 23 use rhi::{ 24 RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, RhiDecryptedIdentity, 25 RhiEncryptedIdentityProvisioningMaterial, RhiEvidenceAttestationSupersession, 26 RhiExactPublicationSink, RhiIdentityEnvelopeBinding, RhiPreparedPublicationAttempt, 27 RhiPublicationAttemptOutcome, RhiPublicationAuthority, RhiPublicationLeaseOwner, 28 RhiPublicationMode, RhiPublicationOutboxState, RhiPublicationRetryDelayMilliseconds, 29 RhiPublicationTargetState, RhiPublicationUnixMilliseconds, RhiReconciliationAttemptErrorKind, 30 RhiReconciliationAttemptPlan, RhiReconciliationAttemptResults, 31 RhiReconciliationAttestationErrorKind, RhiReconciliationCommitErrorKind, 32 RhiReconciliationFinalizationCommitErrorKind, RhiReconciliationFinalizationErrorKind, 33 RhiReconciliationJobErrorKind, RhiReconciliationJobPolicy, RhiReconciliationJobState, 34 RhiReconciliationLease, RhiReconciliationLeaseOwner, RhiReconciliationRetryDelayMilliseconds, 35 RhiReconciliationScopePrerequisites, RhiReconciliationSourceReplayPlan, 36 RhiReconciliationSourceResult, RhiReconciliationUnixMilliseconds, RhiRuntimeContext, 37 RhiStateMetadata, RhiTimeEntropyAdapters, RhiTradeMutationAdmissionLimits, 38 RhiTradeMutationAuthoredTimePolicy, RhiTradeMutationObservedAtUnixSeconds, 39 RhiTradeSourceCompletion, TradeId, admit_rhi_trade_mutation_event, 40 build_rhi_signed_evidence_attestation, initialize_rhi_state, open_rhi_state_inspection, 41 open_rhi_state_read_write, parse_rhi_cli_v1_from, parse_rhi_config_v1, 42 provision_rhi_encrypted_identity, reduce_rhi_reconciliation_manifest, 43 resolve_rhi_runtime_context, resolve_rhi_wrapping_credential, 44 }; 45 use sha2::{Digest, Sha256}; 46 use sqlx::{Connection, SqliteConnection, sqlite::SqliteConnectOptions}; 47 use tokio::sync::Notify; 48 49 const EXAMPLE: &str = include_str!("../contracts/services_hardening/config.v1.example.toml"); 50 const TRADE_VECTOR: &str = 51 include_str!("../contracts/conformance/vectors/trade_ingest_proposal.v1.json"); 52 const SIGNED_ATTESTATION_VECTOR: &str = include_str!( 53 "../contracts/conformance/vectors/reconciliation_attestation_signed_event.v1.json" 54 ); 55 56 fn runtime(root: &Path, instance: &str) -> RhiRuntimeContext { 57 let invocation = parse_rhi_cli_v1_from([ 58 "rhi", 59 "--profile", 60 "repo-local", 61 "--instance", 62 instance, 63 "--repo-local-root", 64 root.to_str().expect("UTF-8 temporary root"), 65 "run", 66 ]) 67 .expect("runtime invocation"); 68 resolve_rhi_runtime_context( 69 &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()), 70 &invocation, 71 ) 72 .expect("runtime context") 73 } 74 75 fn metadata(runtime: &RhiRuntimeContext) -> RhiStateMetadata { 76 let config = parse_rhi_config_v1(EXAMPLE.as_bytes(), rhi::RhiConfigProfile::RepoLocal) 77 .expect("configuration"); 78 metadata_from_config(runtime, &config) 79 } 80 81 fn metadata_from_config( 82 runtime: &RhiRuntimeContext, 83 config: &rhi::RhiConfigDocumentV1, 84 ) -> RhiStateMetadata { 85 RhiStateMetadata::new( 86 runtime, 87 config, 88 SourceGeneration::new([0x5a; 32]).expect("generation"), 89 1_725_000_000_000, 90 ) 91 .expect("metadata") 92 } 93 94 fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) { 95 ( 96 MigrationAppliedAtUnixSeconds::new(1_725_000_000).expect("time"), 97 MigrationBuildIdentity::new( 98 env!("CARGO_PKG_VERSION"), 99 "1111111111111111111111111111111111111111", 100 "053d0c750bf9cd683c6ea37cefe7e79617ba629f", 101 "rustc-test", 102 "test-target", 103 "service-host", 104 1, 105 rhi::RHI_STATE_SCHEMA_VERSION, 106 1, 107 1, 108 1, 109 ) 110 .expect("build"), 111 ) 112 } 113 114 async fn initialize(runtime: &RhiRuntimeContext, metadata: &RhiStateMetadata) { 115 fs::create_dir_all(runtime.context().paths().state()).expect("state directory"); 116 fs::set_permissions( 117 runtime.context().paths().state(), 118 fs::Permissions::from_mode(0o700), 119 ) 120 .expect("state permissions"); 121 let (applied_at, build) = migration_evidence(); 122 initialize_rhi_state(runtime, metadata, applied_at, &build) 123 .await 124 .expect("initialize"); 125 } 126 127 async fn open_writer( 128 runtime: &RhiRuntimeContext, 129 metadata: &RhiStateMetadata, 130 ) -> rhi::RhiStateHost { 131 let (applied_at, build) = migration_evidence(); 132 open_rhi_state_read_write(runtime, metadata, applied_at, &build) 133 .await 134 .expect("writer") 135 } 136 137 async fn write_dirty( 138 runtime: &RhiRuntimeContext, 139 trade: TradeId, 140 generation: u64, 141 policy: [u8; 32], 142 updated_at_unix_s: u64, 143 ) { 144 let options = SqliteConnectOptions::new() 145 .filename(runtime.artifacts().state_database()) 146 .create_if_missing(false) 147 .foreign_keys(true); 148 let mut connection = SqliteConnection::connect_with(&options) 149 .await 150 .expect("offline fixture connection"); 151 sqlx::query( 152 r#"INSERT INTO trade_dirty_generations ( 153 trade_id, generation, evidence_policy_sha256, updated_at_unix_s 154 ) VALUES (?, ?, ?, ?) 155 ON CONFLICT(trade_id) DO UPDATE SET 156 generation = excluded.generation, 157 evidence_policy_sha256 = excluded.evidence_policy_sha256, 158 updated_at_unix_s = excluded.updated_at_unix_s"#, 159 ) 160 .bind(trade.as_bytes().as_slice()) 161 .bind(i64::try_from(generation).expect("generation")) 162 .bind(policy.as_slice()) 163 .bind(i64::try_from(updated_at_unix_s).expect("time")) 164 .execute(&mut connection) 165 .await 166 .expect("dirty fixture"); 167 connection.close().await.expect("fixture close"); 168 } 169 170 fn policy( 171 queue_capacity: u32, 172 lease_ms: u64, 173 renewal_ms: u64, 174 max_attempts: u16, 175 initial_backoff_ms: u64, 176 maximum_backoff_ms: u64, 177 ) -> RhiReconciliationJobPolicy { 178 RhiReconciliationJobPolicy::new( 179 queue_capacity, 180 lease_ms, 181 renewal_ms, 182 max_attempts, 183 initial_backoff_ms, 184 maximum_backoff_ms, 185 ) 186 .expect("policy") 187 } 188 189 fn configured_policy(configuration: &rhi::RhiConfigDocumentV1) -> RhiReconciliationJobPolicy { 190 RhiReconciliationJobPolicy::from_configuration(configuration).expect("configured job policy") 191 } 192 193 fn now(value: u64) -> RhiReconciliationUnixMilliseconds { 194 RhiReconciliationUnixMilliseconds::new(value).expect("time") 195 } 196 197 fn owner(byte: u8) -> RhiReconciliationLeaseOwner { 198 RhiReconciliationLeaseOwner::from_bytes([byte; 16]).expect("owner") 199 } 200 201 fn delay(value: u64) -> RhiReconciliationRetryDelayMilliseconds { 202 RhiReconciliationRetryDelayMilliseconds::new(value).expect("delay") 203 } 204 205 #[tokio::test] 206 async fn schedule_is_bounded_idempotent_and_uses_the_frozen_job_identity() { 207 let root = tempfile::tempdir().expect("root"); 208 let runtime = runtime(root.path(), "primary"); 209 let metadata = metadata(&runtime); 210 initialize(&runtime, &metadata).await; 211 let first_trade = TradeId::from_bytes([0x11; 16]); 212 let second_trade = TradeId::from_bytes([0x33; 16]); 213 write_dirty(&runtime, first_trade, 1, [0x22; 32], 1_000).await; 214 write_dirty(&runtime, second_trade, 1, [0x44; 32], 1_000).await; 215 let host = open_writer(&runtime, &metadata).await; 216 let jobs = host.repositories().reconciliation_jobs(); 217 let bounded = policy(1, 1_000, 100, 3, 100, 1_000); 218 219 let first = jobs 220 .schedule_trade(first_trade, bounded, now(1_000)) 221 .await 222 .expect("first schedule"); 223 assert!(first.created()); 224 assert_eq!(first.job().state(), RhiReconciliationJobState::Ready); 225 assert_eq!(first.job().input_generation(), 1); 226 assert_eq!(first.job().revision(), 1); 227 assert_eq!(first.job().next_attempt(), Some(now(1_000))); 228 assert_eq!( 229 lower_hex(first.job().id().as_bytes()), 230 "dc7b98b36fb8e839cba83dd7274f25fc27d021d2b46c2f5ab089ddde591fad02" 231 ); 232 let replay = jobs 233 .schedule_trade(first_trade, bounded, now(1_001)) 234 .await 235 .expect("idempotent schedule"); 236 assert!(!replay.created()); 237 assert_eq!(replay.job(), first.job()); 238 239 let full = jobs 240 .schedule_trade(second_trade, bounded, now(1_001)) 241 .await 242 .expect_err("configured queue bound"); 243 assert_eq!(full.kind(), RhiReconciliationJobErrorKind::QueueFull); 244 host.close().await.expect("close"); 245 } 246 247 #[tokio::test] 248 async fn claims_renew_only_when_due_and_expired_leases_are_reclaimed() { 249 let root = tempfile::tempdir().expect("root"); 250 let runtime = runtime(root.path(), "leases"); 251 let metadata = metadata(&runtime); 252 initialize(&runtime, &metadata).await; 253 let trade = TradeId::from_bytes([0x12; 16]); 254 write_dirty(&runtime, trade, 1, [0x23; 32], 1_000).await; 255 let host = open_writer(&runtime, &metadata).await; 256 let jobs = host.repositories().reconciliation_jobs(); 257 jobs.schedule_trade(trade, policy(8, 1_000, 100, 4, 100, 1_000), now(1_000)) 258 .await 259 .expect("schedule"); 260 261 let first = jobs 262 .claim_next(owner(1), now(1_000)) 263 .await 264 .expect("claim") 265 .expect("job"); 266 assert_eq!(first.job().attempt_count(), 1); 267 assert_eq!(first.lease_expires(), now(2_000)); 268 assert_eq!(first.renewal_due(), now(1_900)); 269 assert_eq!( 270 jobs.renew(first, now(1_899)) 271 .await 272 .expect_err("early renewal") 273 .kind(), 274 RhiReconciliationJobErrorKind::NotReady 275 ); 276 let renewed = jobs.renew(first, now(1_900)).await.expect("renew"); 277 assert_eq!(renewed.job().revision(), 3); 278 assert_eq!(renewed.lease_expires(), now(2_900)); 279 assert!( 280 jobs.claim_next(owner(2), now(2_899)) 281 .await 282 .expect("no early reclaim") 283 .is_none() 284 ); 285 let reclaimed = jobs 286 .claim_next(owner(2), now(2_900)) 287 .await 288 .expect("reclaim") 289 .expect("expired lease"); 290 assert_eq!(reclaimed.job().attempt_count(), 2); 291 assert_eq!(reclaimed.job().revision(), 4); 292 assert_eq!( 293 jobs.record_failure(renewed, now(2_900), delay(0)) 294 .await 295 .expect_err("expired lease") 296 .kind(), 297 RhiReconciliationJobErrorKind::LeaseLost 298 ); 299 assert_eq!( 300 jobs.renew(first, now(1_950)) 301 .await 302 .expect_err("stale revision") 303 .kind(), 304 RhiReconciliationJobErrorKind::LeaseLost 305 ); 306 host.close().await.expect("close"); 307 } 308 309 #[tokio::test] 310 async fn failures_schedule_exact_jitter_and_the_final_attempt_exhausts() { 311 let root = tempfile::tempdir().expect("root"); 312 let runtime = runtime(root.path(), "retry"); 313 let metadata = metadata(&runtime); 314 initialize(&runtime, &metadata).await; 315 let trade = TradeId::from_bytes([0x13; 16]); 316 write_dirty(&runtime, trade, 1, [0x24; 32], 1_000).await; 317 let host = open_writer(&runtime, &metadata).await; 318 let jobs = host.repositories().reconciliation_jobs(); 319 jobs.schedule_trade(trade, policy(8, 1_000, 100, 2, 100, 1_000), now(1_000)) 320 .await 321 .expect("schedule"); 322 let first = jobs 323 .claim_next(owner(3), now(1_000)) 324 .await 325 .expect("claim") 326 .expect("job"); 327 assert_eq!(first.retry_delay_upper_bound(), 100); 328 assert_eq!( 329 jobs.record_failure(first, now(1_100), delay(101)) 330 .await 331 .expect_err("delay above governed cap") 332 .kind(), 333 RhiReconciliationJobErrorKind::InvalidInput 334 ); 335 let retry = jobs 336 .record_failure(first, now(1_100), delay(100)) 337 .await 338 .expect("retry schedule"); 339 assert_eq!(retry.state(), RhiReconciliationJobState::Ready); 340 assert_eq!(retry.failure_count(), 1); 341 assert_eq!(retry.next_attempt(), Some(now(1_200))); 342 assert!( 343 jobs.claim_next(owner(4), now(1_199)) 344 .await 345 .expect("not due") 346 .is_none() 347 ); 348 let second = jobs 349 .claim_next(owner(4), now(1_200)) 350 .await 351 .expect("second claim") 352 .expect("job"); 353 assert_eq!(second.retry_delay_upper_bound(), 200); 354 let exhausted = jobs 355 .record_failure(second, now(1_300), delay(200)) 356 .await 357 .expect("final failure"); 358 assert_eq!(exhausted.state(), RhiReconciliationJobState::Exhausted); 359 assert_eq!(exhausted.attempt_count(), 2); 360 assert_eq!(exhausted.failure_count(), 2); 361 assert_eq!(exhausted.next_attempt(), None); 362 assert!( 363 jobs.claim_next(owner(5), now(9_000)) 364 .await 365 .expect("terminal queue") 366 .is_none() 367 ); 368 host.close().await.expect("close"); 369 } 370 371 #[tokio::test] 372 async fn an_expired_final_attempt_is_exhausted_instead_of_reclaimed() { 373 let root = tempfile::tempdir().expect("root"); 374 let runtime = runtime(root.path(), "expired-final"); 375 let metadata = metadata(&runtime); 376 initialize(&runtime, &metadata).await; 377 let trade = TradeId::from_bytes([0x18; 16]); 378 write_dirty(&runtime, trade, 1, [0x29; 32], 1_000).await; 379 let host = open_writer(&runtime, &metadata).await; 380 let jobs = host.repositories().reconciliation_jobs(); 381 let job_policy = policy(8, 1_000, 100, 1, 100, 1_000); 382 jobs.schedule_trade(trade, job_policy, now(1_000)) 383 .await 384 .expect("schedule"); 385 jobs.claim_next(owner(10), now(1_000)) 386 .await 387 .expect("claim") 388 .expect("job"); 389 390 assert!( 391 jobs.claim_next(owner(11), now(2_000)) 392 .await 393 .expect("expired final attempt") 394 .is_none() 395 ); 396 let retained = jobs 397 .schedule_trade(trade, job_policy, now(2_000)) 398 .await 399 .expect("idempotent retained job"); 400 assert!(!retained.created()); 401 assert_eq!(retained.job().state(), RhiReconciliationJobState::Exhausted); 402 assert_eq!(retained.job().attempt_count(), 1); 403 assert_eq!(retained.job().failure_count(), 1); 404 host.close().await.expect("close"); 405 } 406 407 #[tokio::test] 408 async fn lease_and_schedule_state_survive_reopen_and_new_generation_supersedes() { 409 let root = tempfile::tempdir().expect("root"); 410 let runtime = runtime(root.path(), "reopen"); 411 let metadata = metadata(&runtime); 412 initialize(&runtime, &metadata).await; 413 let trade = TradeId::from_bytes([0x14; 16]); 414 write_dirty(&runtime, trade, 1, [0x25; 32], 1_000).await; 415 let job_policy = policy(8, 1_000, 100, 3, 100, 1_000); 416 let host = open_writer(&runtime, &metadata).await; 417 let jobs = host.repositories().reconciliation_jobs(); 418 let first_job = jobs 419 .schedule_trade(trade, job_policy, now(1_000)) 420 .await 421 .expect("schedule") 422 .job(); 423 let stale = jobs 424 .claim_next(owner(6), now(1_000)) 425 .await 426 .expect("claim") 427 .expect("job"); 428 host.close().await.expect("crash boundary close"); 429 430 write_dirty(&runtime, trade, 2, [0x26; 32], 1_001).await; 431 let host = open_writer(&runtime, &metadata).await; 432 let jobs = host.repositories().reconciliation_jobs(); 433 let newer = jobs 434 .schedule_trade(trade, job_policy, now(1_900)) 435 .await 436 .expect("new generation"); 437 assert!(newer.created()); 438 assert_eq!(newer.job().input_generation(), 2); 439 assert_ne!(newer.job().id(), first_job.id()); 440 assert_eq!( 441 jobs.renew(stale, now(1_900)) 442 .await 443 .expect_err("superseded lease") 444 .kind(), 445 RhiReconciliationJobErrorKind::LeaseLost 446 ); 447 let claimed = jobs 448 .claim_next(owner(7), now(1_900)) 449 .await 450 .expect("claim new") 451 .expect("new job"); 452 assert_eq!(claimed.job().id(), newer.job().id()); 453 host.close().await.expect("close"); 454 } 455 456 #[tokio::test] 457 async fn concurrent_claims_have_one_winner_and_read_only_state_cannot_mutate() { 458 let root = tempfile::tempdir().expect("root"); 459 let runtime = runtime(root.path(), "concurrent"); 460 let metadata = metadata(&runtime); 461 initialize(&runtime, &metadata).await; 462 let trade = TradeId::from_bytes([0x15; 16]); 463 write_dirty(&runtime, trade, 1, [0x27; 32], 1_000).await; 464 let host = open_writer(&runtime, &metadata).await; 465 let jobs = host.repositories().reconciliation_jobs(); 466 jobs.schedule_trade(trade, policy(8, 1_000, 100, 3, 100, 1_000), now(1_000)) 467 .await 468 .expect("schedule"); 469 let (left, right) = tokio::join!( 470 jobs.claim_next(owner(8), now(1_000)), 471 jobs.claim_next(owner(9), now(1_000)) 472 ); 473 let winners = [left, right] 474 .into_iter() 475 .filter_map(|result| result.expect("claim result")) 476 .count(); 477 assert_eq!(winners, 1); 478 host.close().await.expect("close"); 479 480 let inspection = open_rhi_state_inspection(&runtime, &metadata) 481 .await 482 .expect("inspection"); 483 let error = inspection 484 .repositories() 485 .reconciliation_jobs() 486 .schedule_trade(trade, policy(8, 1_000, 100, 3, 100, 1_000), now(2_000)) 487 .await 488 .expect_err("read-only mutation"); 489 assert_eq!(error.kind(), RhiReconciliationJobErrorKind::InvalidMode); 490 inspection.close().await.expect("inspection close"); 491 } 492 493 #[tokio::test] 494 async fn attempt_plan_binds_the_exact_claim_policy_sources_and_frozen_identities() { 495 let root = tempfile::tempdir().expect("root"); 496 let runtime = runtime(root.path(), "attempt-plan"); 497 let configuration = parse_rhi_config_v1(EXAMPLE.as_bytes(), rhi::RhiConfigProfile::RepoLocal) 498 .expect("configuration"); 499 let metadata = metadata_from_config(&runtime, &configuration); 500 initialize(&runtime, &metadata).await; 501 let trade = TradeId::from_bytes([0x11; 16]); 502 write_dirty( 503 &runtime, 504 trade, 505 1, 506 *metadata.evidence_policy_digest().as_bytes(), 507 1_000, 508 ) 509 .await; 510 let host = open_writer(&runtime, &metadata).await; 511 let jobs = host.repositories().reconciliation_jobs(); 512 jobs.schedule_trade(trade, configured_policy(&configuration), now(1_000)) 513 .await 514 .expect("schedule"); 515 let lease = jobs 516 .claim_next(owner(0x71), now(1_000)) 517 .await 518 .expect("claim") 519 .expect("job"); 520 521 let plan = RhiReconciliationAttemptPlan::from_claim(lease, &configuration, now(1_000)) 522 .expect("attempt plan"); 523 assert_eq!(plan.job_id(), lease.job().id()); 524 assert_eq!(plan.input_generation(), 1); 525 assert_eq!( 526 plan.evidence_policy_digest(), 527 metadata.evidence_policy_digest() 528 ); 529 assert_eq!(plan.attempt_started_at(), now(1_000)); 530 assert_eq!(plan.deadline(), now(31_000)); 531 assert_eq!( 532 lower_hex(plan.job_id().as_bytes()), 533 "4e5ecfbee585698c6a67202b30487249291b446909b16a95fd8ae729c7d51e85" 534 ); 535 assert_eq!( 536 lower_hex(plan.id().as_bytes()), 537 "89b61ce985d6f11ed963a0d96a05b80e8961a16749a122010d33e1aaee04fdeb" 538 ); 539 let [request] = plan.requests() else { 540 panic!("exact source inventory") 541 }; 542 assert_eq!(request.source_id(), "trade-primary"); 543 assert_eq!(request.trade_id(), trade); 544 assert!(request.required()); 545 assert_eq!(request.attempt_started_at(), now(1_000)); 546 assert_eq!(request.deadline(), now(11_000)); 547 assert_eq!(request.lookback_seconds(), 86_400); 548 assert_eq!(request.maximum_events(), 4_096); 549 assert_eq!(request.maximum_bytes(), 8_388_608); 550 assert_eq!( 551 lower_hex(request.selector_digest().as_bytes()), 552 "2c489c22515b4db784f1be9ab2c224b578c8d28e95d3921ade420b6aa78345bd" 553 ); 554 assert_eq!( 555 lower_hex(request.id().as_bytes()), 556 "ef58f9e8a61f7964734cf0c2aabe0bdb2cbdcec529b16f40ae18ec224d0c896d" 557 ); 558 let replay = 559 RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None) 560 .expect("replay plan"); 561 assert_eq!(replay.overlap_seconds(), 300); 562 assert_eq!(replay.since_unix_seconds(), 0); 563 assert_eq!( 564 lower_hex(replay.id().as_bytes()), 565 "2490a2e9a6e85051e92f6c2fc2ff7e98a1367afd6c26f8625eb421cd2ab30c68" 566 ); 567 host.close().await.expect("close"); 568 } 569 570 #[tokio::test] 571 async fn attempt_plan_rejects_policy_mismatch_and_expired_claim_time() { 572 let root = tempfile::tempdir().expect("root"); 573 let runtime = runtime(root.path(), "attempt-policy"); 574 let configuration = parse_rhi_config_v1(EXAMPLE.as_bytes(), rhi::RhiConfigProfile::RepoLocal) 575 .expect("configuration"); 576 let metadata = metadata_from_config(&runtime, &configuration); 577 initialize(&runtime, &metadata).await; 578 let host = open_writer(&runtime, &metadata).await; 579 let jobs = host.repositories().reconciliation_jobs(); 580 let governed_policy = configured_policy(&configuration); 581 582 let mismatch_trade = TradeId::from_bytes([0x61; 16]); 583 write_dirty(&runtime, mismatch_trade, 1, [0x62; 32], 1_000).await; 584 jobs.schedule_trade(mismatch_trade, governed_policy, now(1_000)) 585 .await 586 .expect("schedule mismatch"); 587 let mismatch = jobs 588 .claim_next(owner(0x63), now(1_000)) 589 .await 590 .expect("claim mismatch") 591 .expect("job"); 592 assert_eq!( 593 RhiReconciliationAttemptPlan::from_claim(mismatch, &configuration, now(1_000)) 594 .expect_err("policy mismatch") 595 .kind(), 596 RhiReconciliationAttemptErrorKind::PolicyMismatch 597 ); 598 599 let scheduling_mismatch_trade = TradeId::from_bytes([0x69; 16]); 600 write_dirty( 601 &runtime, 602 scheduling_mismatch_trade, 603 1, 604 *metadata.evidence_policy_digest().as_bytes(), 605 1_001, 606 ) 607 .await; 608 jobs.schedule_trade( 609 scheduling_mismatch_trade, 610 policy(8, 30_000, 10_000, 3, 250, 30_000), 611 now(1_001), 612 ) 613 .await 614 .expect("schedule with mismatched job policy"); 615 let scheduling_mismatch = jobs 616 .claim_next(owner(0x6a), now(1_001)) 617 .await 618 .expect("claim scheduling mismatch") 619 .expect("job"); 620 assert_eq!( 621 RhiReconciliationAttemptPlan::from_claim(scheduling_mismatch, &configuration, now(1_001),) 622 .expect_err("job policy mismatch") 623 .kind(), 624 RhiReconciliationAttemptErrorKind::PolicyMismatch 625 ); 626 627 let expired_trade = TradeId::from_bytes([0x64; 16]); 628 write_dirty( 629 &runtime, 630 expired_trade, 631 1, 632 *metadata.evidence_policy_digest().as_bytes(), 633 1_002, 634 ) 635 .await; 636 jobs.schedule_trade(expired_trade, governed_policy, now(1_002)) 637 .await 638 .expect("schedule expired"); 639 let expired = jobs 640 .claim_next(owner(0x65), now(1_002)) 641 .await 642 .expect("claim expired") 643 .expect("job"); 644 assert_eq!( 645 RhiReconciliationAttemptPlan::from_claim(expired, &configuration, expired.lease_expires(),) 646 .expect_err("expired attempt start") 647 .kind(), 648 RhiReconciliationAttemptErrorKind::LeaseExpired 649 ); 650 host.close().await.expect("close"); 651 } 652 653 #[tokio::test] 654 async fn reclaimed_attempt_changes_identity_and_caps_deadlines_to_each_lease() { 655 let short_lease = EXAMPLE 656 .replace("lease_ms = 30000", "lease_ms = 1000") 657 .replace("lease_renewal_ms = 10000", "lease_renewal_ms = 100"); 658 let root = tempfile::tempdir().expect("root"); 659 let runtime = runtime(root.path(), "attempt-reclaim"); 660 let configuration = 661 parse_rhi_config_v1(short_lease.as_bytes(), rhi::RhiConfigProfile::RepoLocal) 662 .expect("short-lease configuration"); 663 let metadata = metadata_from_config(&runtime, &configuration); 664 initialize(&runtime, &metadata).await; 665 let trade = TradeId::from_bytes([0x66; 16]); 666 write_dirty( 667 &runtime, 668 trade, 669 1, 670 *metadata.evidence_policy_digest().as_bytes(), 671 1_000, 672 ) 673 .await; 674 let host = open_writer(&runtime, &metadata).await; 675 let jobs = host.repositories().reconciliation_jobs(); 676 jobs.schedule_trade(trade, configured_policy(&configuration), now(1_000)) 677 .await 678 .expect("schedule"); 679 let first_lease = jobs 680 .claim_next(owner(0x67), now(1_000)) 681 .await 682 .expect("first claim") 683 .expect("job"); 684 let first = RhiReconciliationAttemptPlan::from_claim(first_lease, &configuration, now(1_000)) 685 .expect("first plan"); 686 assert_eq!(first.deadline(), now(2_000)); 687 assert_eq!(first.requests()[0].deadline(), now(2_000)); 688 let later_start = 689 RhiReconciliationAttemptPlan::from_claim(first_lease, &configuration, now(1_500)) 690 .expect("same claim with a later explicit start"); 691 assert_eq!(later_start.id(), first.id()); 692 assert_eq!(later_start.requests()[0].deadline(), now(2_000)); 693 assert_ne!(later_start.requests()[0].id(), first.requests()[0].id()); 694 695 let second_lease = jobs 696 .claim_next(owner(0x68), now(2_000)) 697 .await 698 .expect("reclaim") 699 .expect("job"); 700 let second = RhiReconciliationAttemptPlan::from_claim(second_lease, &configuration, now(2_000)) 701 .expect("second plan"); 702 assert_eq!(second_lease.job().attempt_count(), 2); 703 assert_eq!(second.deadline(), now(3_000)); 704 assert_eq!(second.requests()[0].deadline(), now(3_000)); 705 assert_ne!(first.id(), second.id()); 706 assert_ne!(first.requests()[0].id(), second.requests()[0].id()); 707 assert_eq!( 708 first.requests()[0].selector_digest(), 709 second.requests()[0].selector_digest() 710 ); 711 host.close().await.expect("close"); 712 } 713 714 #[tokio::test] 715 async fn source_results_enforce_deadline_outcome_and_exact_resource_bounds() { 716 let (root, runtime, metadata, configuration, host, plan) = 717 attempt_fixture("attempt-results").await; 718 let request = &plan.requests()[0]; 719 let complete = RhiReconciliationSourceResult::new( 720 request, 721 RhiTradeSourceCompletion::Complete, 722 now(1_000), 723 now(10_999), 724 4_096, 725 8_388_608, 726 ) 727 .expect("exact maximum result"); 728 assert_eq!(complete.request_id(), request.id()); 729 assert_eq!(complete.outcome().code(), "complete"); 730 assert_eq!(complete.accepted_event_count(), 4_096); 731 assert_eq!(complete.accepted_event_bytes(), 8_388_608); 732 assert_eq!(complete.started_at(), now(1_000)); 733 assert_eq!(complete.finished_at(), now(10_999)); 734 735 let timeout = RhiReconciliationSourceResult::new( 736 request, 737 RhiTradeSourceCompletion::IncompleteTimeout, 738 now(1_000), 739 now(11_000), 740 0, 741 0, 742 ) 743 .expect("deadline timeout"); 744 assert_eq!(timeout.outcome().code(), "incomplete_timeout"); 745 for outcome in [ 746 RhiTradeSourceCompletion::IncompleteUnavailable, 747 RhiTradeSourceCompletion::IncompleteResourceLimit, 748 RhiTradeSourceCompletion::IncompleteUnknown, 749 RhiTradeSourceCompletion::Unsupported, 750 ] { 751 let result = 752 RhiReconciliationSourceResult::new(request, outcome, now(1_000), now(1_001), 0, 0) 753 .expect("safe incomplete outcome"); 754 assert_eq!(result.outcome(), outcome); 755 } 756 for invalid in [ 757 RhiReconciliationSourceResult::new( 758 request, 759 RhiTradeSourceCompletion::Complete, 760 now(1_000), 761 now(11_000), 762 0, 763 0, 764 ), 765 RhiReconciliationSourceResult::new( 766 request, 767 RhiTradeSourceCompletion::IncompleteTimeout, 768 now(1_000), 769 now(10_999), 770 0, 771 0, 772 ), 773 RhiReconciliationSourceResult::new( 774 request, 775 RhiTradeSourceCompletion::IncompleteTimeout, 776 now(11_000), 777 now(11_000), 778 0, 779 0, 780 ), 781 RhiReconciliationSourceResult::new( 782 request, 783 RhiTradeSourceCompletion::Complete, 784 now(1_000), 785 now(1_001), 786 4_097, 787 8_388_608, 788 ), 789 RhiReconciliationSourceResult::new( 790 request, 791 RhiTradeSourceCompletion::Complete, 792 now(1_000), 793 now(1_001), 794 4_096, 795 8_388_609, 796 ), 797 RhiReconciliationSourceResult::new( 798 request, 799 RhiTradeSourceCompletion::Complete, 800 now(1_000), 801 now(1_001), 802 1, 803 0, 804 ), 805 RhiReconciliationSourceResult::new( 806 request, 807 RhiTradeSourceCompletion::Unsupported, 808 now(1_000), 809 now(1_001), 810 1, 811 1, 812 ), 813 ] { 814 assert_eq!( 815 invalid.expect_err("invalid result").kind(), 816 RhiReconciliationAttemptErrorKind::InvalidInput 817 ); 818 } 819 820 let exact = RhiReconciliationAttemptResults::new(&plan, [complete]).expect("inventory"); 821 assert_eq!(exact.attempt_id(), plan.id()); 822 assert_eq!(exact.results(), [complete]); 823 assert_eq!( 824 RhiReconciliationAttemptResults::new(&plan, []) 825 .expect_err("missing result") 826 .kind(), 827 RhiReconciliationAttemptErrorKind::ResultInventory 828 ); 829 assert_eq!( 830 RhiReconciliationAttemptResults::new(&plan, std::iter::repeat(complete)) 831 .expect_err("bounded infinite excess") 832 .kind(), 833 RhiReconciliationAttemptErrorKind::ResultInventory 834 ); 835 host.close().await.expect("close"); 836 drop((configuration, metadata, runtime, root)); 837 } 838 839 #[tokio::test] 840 async fn result_inventory_rejects_reordered_configured_sources() { 841 let multi_source = EXAMPLE 842 .replace( 843 "read = false\nwrite = true\nrequired = false", 844 "read = true\nwrite = true\nrequired = false", 845 ) 846 .replace( 847 "[[evidence.sources]]\nsource_id = \"trade-primary\"", 848 "[[evidence.sources]]\nsource_id = \"a-secondary\"\nkind = \"nostr_relay\"\nrelay_id = \"relay-secondary\"\nrequired = false\nselector = \"trade_mutation_lineage_v1\"\ndeadline_ms = 5000\nlookback_seconds = 3600\noverlap_seconds = 60\n\n[[evidence.sources]]\nsource_id = \"trade-primary\"", 849 ); 850 let root = tempfile::tempdir().expect("root"); 851 let runtime = runtime(root.path(), "attempt-order"); 852 let configuration = 853 parse_rhi_config_v1(multi_source.as_bytes(), rhi::RhiConfigProfile::RepoLocal) 854 .expect("multi-source configuration"); 855 let metadata = metadata_from_config(&runtime, &configuration); 856 initialize(&runtime, &metadata).await; 857 let trade = TradeId::from_bytes([0x51; 16]); 858 write_dirty( 859 &runtime, 860 trade, 861 1, 862 *metadata.evidence_policy_digest().as_bytes(), 863 1_000, 864 ) 865 .await; 866 let host = open_writer(&runtime, &metadata).await; 867 let jobs = host.repositories().reconciliation_jobs(); 868 jobs.schedule_trade(trade, configured_policy(&configuration), now(1_000)) 869 .await 870 .expect("schedule"); 871 let lease = jobs 872 .claim_next(owner(0x52), now(1_000)) 873 .await 874 .expect("claim") 875 .expect("job"); 876 let plan = 877 RhiReconciliationAttemptPlan::from_claim(lease, &configuration, now(1_000)).expect("plan"); 878 assert_eq!( 879 plan.requests() 880 .iter() 881 .map(|request| request.source_id()) 882 .collect::<Vec<_>>(), 883 ["a-secondary", "trade-primary"] 884 ); 885 let mut results = plan 886 .requests() 887 .iter() 888 .map(|request| { 889 RhiReconciliationSourceResult::new( 890 request, 891 RhiTradeSourceCompletion::Complete, 892 now(1_000), 893 now(1_001), 894 0, 895 0, 896 ) 897 .expect("result") 898 }) 899 .collect::<Vec<_>>(); 900 assert!(RhiReconciliationAttemptResults::new(&plan, results.clone()).is_ok()); 901 results.reverse(); 902 assert_eq!( 903 RhiReconciliationAttemptResults::new(&plan, results) 904 .expect_err("reordered") 905 .kind(), 906 RhiReconciliationAttemptErrorKind::ResultInventory 907 ); 908 host.close().await.expect("close"); 909 } 910 911 #[tokio::test] 912 async fn attempt_diagnostics_are_redacted_and_source_free() { 913 let (root, runtime, metadata, configuration, host, plan) = 914 attempt_fixture("attempt-debug").await; 915 let request = &plan.requests()[0]; 916 let result = RhiReconciliationSourceResult::new( 917 request, 918 RhiTradeSourceCompletion::Complete, 919 now(1_000), 920 now(1_001), 921 1, 922 16, 923 ) 924 .expect("result"); 925 let inventory = RhiReconciliationAttemptResults::new(&plan, [result]).expect("inventory"); 926 let rendered = format!("{plan:?} {request:?} {result:?} {inventory:?}"); 927 for secret in [ 928 &lower_hex(plan.id().as_bytes()), 929 &lower_hex(request.id().as_bytes()), 930 &lower_hex(request.selector_digest().as_bytes()), 931 ] { 932 assert!(!rendered.contains(secret)); 933 } 934 let error = RhiReconciliationAttemptResults::new(&plan, []).expect_err("error"); 935 assert!(Error::source(&error).is_none()); 936 assert_eq!( 937 error.code(), 938 "reconciliation_attempt_result_inventory_invalid" 939 ); 940 assert!(!format!("{error} {error:?}").contains("trade-primary")); 941 host.close().await.expect("close"); 942 drop((configuration, metadata, runtime, root)); 943 } 944 945 #[tokio::test] 946 async fn replay_plan_binds_overlap_deduplicates_and_retains_first_provenance() { 947 let started_ms = 1_784_347_200_000; 948 let (root, runtime, metadata, configuration, host, _lease, plan) = 949 replay_fixture("replay-plan", EXAMPLE, started_ms).await; 950 let request = &plan.requests()[0]; 951 let cursor_plan = 952 RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None) 953 .expect("initial cursor plan"); 954 assert_eq!(cursor_plan.overlap_seconds(), 300); 955 assert_eq!(cursor_plan.since_unix_seconds(), 1_784_260_800); 956 assert!(cursor_plan.prior_cursor().is_none()); 957 958 let wire = replay_wire(); 959 let replay = cursor_plan 960 .finish( 961 request, 962 RhiTradeSourceCompletion::Complete, 963 now(started_ms), 964 now(started_ms + 2_000), 965 [ 966 admitted_replay(&configuration, &wire, 1_784_347_201), 967 admitted_replay(&configuration, &wire, 1_784_347_200), 968 ], 969 ) 970 .expect("canonical replay"); 971 assert_eq!(replay.accepted_event_count(), 1); 972 assert_eq!( 973 replay.accepted_original_event_bytes(), 974 u64::try_from(wire.len()).expect("wire bytes") 975 ); 976 assert_eq!(replay.duplicate_observation_count(), 1); 977 assert_eq!( 978 replay.first_observed_at().expect("first provenance").get(), 979 1_784_347_200 980 ); 981 assert_eq!(replay.result().accepted_event_count(), 1); 982 assert_eq!( 983 replay.result().accepted_event_bytes(), 984 u64::try_from(wire.len()).expect("wire bytes") 985 ); 986 let cursor = replay.eligible_cursor().expect("eligible cursor"); 987 assert_eq!(cursor.created_at_unix_seconds(), 1_784_347_200); 988 989 let incomplete = 990 RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None) 991 .expect("incomplete plan") 992 .finish( 993 request, 994 RhiTradeSourceCompletion::IncompleteUnavailable, 995 now(started_ms), 996 now(started_ms + 1_000), 997 [admitted_replay(&configuration, &wire, 1_784_347_200)], 998 ) 999 .expect("incomplete replay"); 1000 assert!(incomplete.cursor_candidate().is_some()); 1001 assert!(incomplete.eligible_cursor().is_none()); 1002 assert!(!format!("{replay:?} {incomplete:?}").contains("trade-primary")); 1003 host.close().await.expect("close"); 1004 drop((configuration, metadata, runtime, root)); 1005 } 1006 1007 #[tokio::test] 1008 async fn source_replay_commit_is_atomic_idempotent_and_mints_durable_cursor_evidence() { 1009 let started_ms = 1_784_347_200_000; 1010 let (root, runtime, metadata, configuration, host, lease, plan) = 1011 replay_fixture("replay-commit", EXAMPLE, started_ms).await; 1012 let request = &plan.requests()[0]; 1013 let wire = replay_wire(); 1014 let make_replay = || { 1015 RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None) 1016 .expect("replay plan") 1017 .finish( 1018 request, 1019 RhiTradeSourceCompletion::Complete, 1020 now(started_ms), 1021 now(started_ms + 2_000), 1022 [admitted_replay(&configuration, &wire, 1_784_347_200)], 1023 ) 1024 .expect("replay") 1025 }; 1026 let first_replay = make_replay(); 1027 let retry_replay = make_replay(); 1028 let early_replay = make_replay(); 1029 let resume_plan = plan.clone(); 1030 let attempts = host.repositories().reconciliation_attempts(); 1031 let committed = attempts 1032 .commit_source_replays(lease, plan.clone(), [first_replay]) 1033 .await 1034 .expect("commit"); 1035 assert!(committed.created()); 1036 assert_eq!(committed.source_result_count(), 1); 1037 assert_eq!(committed.checkpoint_advance_count(), 1); 1038 assert!(committed.dirty_generation_advanced()); 1039 assert_eq!(committed.committed_cursors().len(), 1); 1040 assert_eq!( 1041 committed.committed_cursors()[0] 1042 .cursor() 1043 .created_at_unix_seconds(), 1044 1_784_347_200 1045 ); 1046 let committed_cursor = committed.committed_cursors()[0].clone(); 1047 let resumed = RhiReconciliationSourceReplayPlan::from_request( 1048 &resume_plan, 1049 &resume_plan.requests()[0], 1050 &configuration, 1051 Some(committed_cursor), 1052 ) 1053 .expect("committed cursor resumes exact scope"); 1054 assert_eq!(resumed.since_unix_seconds(), 1_784_346_900); 1055 let manifest = committed 1056 .into_evidence_manifest( 1057 UnixTimeSeconds::new(1_784_347_203), 1058 RhiReconciliationScopePrerequisites::Satisfied, 1059 ) 1060 .expect("manifest"); 1061 assert_eq!(manifest.contract_version(), 1); 1062 assert_eq!( 1063 manifest.shared_manifest_contract_id(), 1064 "radroots.trade.evidence-manifest.v1" 1065 ); 1066 assert_eq!(manifest.shared_manifest_contract_version(), 1); 1067 assert_eq!(manifest.trade_id(), &TradeId::from_bytes([0x11; 16])); 1068 assert_eq!(manifest.trade_generation(), 1); 1069 assert_eq!(manifest.observed_at_unix_seconds(), 1_784_347_203); 1070 assert_eq!( 1071 (manifest.source_count(), manifest.observation_count()), 1072 (1, 1) 1073 ); 1074 assert_eq!( 1075 manifest.digest(), 1076 [ 1077 0x0b, 0x19, 0x3e, 0xd2, 0x93, 0x56, 0xd6, 0x3d, 0x37, 0x31, 0x63, 0x4b, 0x37, 0xfe, 1078 0x1f, 0x5d, 0x53, 0x21, 0x48, 0x74, 0x79, 0x3d, 0xc2, 0x3f, 0xe4, 0xa3, 0xb9, 0x84, 1079 0xf8, 0xa4, 0x51, 0xc2, 1080 ] 1081 ); 1082 let canonical_manifest = manifest.canonical_bytes().to_vec(); 1083 let projection = reduce_rhi_reconciliation_manifest(manifest).expect("pure projection"); 1084 assert_eq!(projection.contract_version(), 1); 1085 assert_eq!( 1086 projection.shared_reducer_contract_id(), 1087 "radroots.trade.reducer.v1" 1088 ); 1089 assert_eq!(projection.shared_reducer_contract_version(), 1); 1090 assert_eq!(projection.trade_id(), &TradeId::from_bytes([0x11; 16])); 1091 assert_eq!(projection.manifest().canonical_bytes(), canonical_manifest); 1092 assert!(projection.root_mutation_id().is_some()); 1093 assert_eq!(projection.issue_count(), 0); 1094 assert_eq!( 1095 projection.shared_projection_digest(), 1096 Some([ 1097 0x21, 0xd5, 0xd5, 0xe6, 0x06, 0x7a, 0x13, 0x68, 0xd0, 0xd5, 0x25, 0xa3, 0xec, 0xd1, 1098 0xb5, 0xcc, 0x99, 0xcb, 0x03, 0xd7, 0xf8, 0x06, 0xe6, 0xba, 0x47, 0xd3, 0xb9, 0x29, 1099 0x99, 0xa7, 0xe9, 0x61, 1100 ]) 1101 ); 1102 assert_eq!( 1103 projection.digest(), 1104 Some([ 1105 0xd1, 0x33, 0xa7, 0x72, 0xd2, 0x87, 0xa2, 0x56, 0x4a, 0xb3, 0xb3, 0xb2, 0xca, 0xb6, 1106 0xdc, 0xa6, 0xe5, 0xc5, 0xa0, 0x7f, 0x30, 0x8f, 0x67, 0xc5, 0xee, 0x77, 0x38, 0x06, 1107 0x40, 0x7a, 0x03, 0x71, 1108 ]) 1109 ); 1110 let claim = *projection.root_mutation_id().expect("root proposal"); 1111 let evaluation = rhi::evaluate_rhi_reconciliation_claim(projection, claim); 1112 assert_eq!(evaluation.contract_version(), 1); 1113 assert_eq!( 1114 evaluation.coverage(), 1115 rhi::RhiReconciliationCoverage::ScopeSatisfied 1116 ); 1117 assert_eq!( 1118 evaluation.outcome(), 1119 rhi::RhiReconciliationOutcome::Indeterminate 1120 ); 1121 assert_eq!( 1122 evaluation.reason_codes(), 1123 [rhi::RhiReconciliationReasonCode::AgreementClaimMissing] 1124 ); 1125 assert_eq!(evaluation.claim_mutation_id(), &claim); 1126 assert_eq!( 1127 evaluation.projection().trade_id(), 1128 &TradeId::from_bytes([0x11; 16]) 1129 ); 1130 let evaluation_debug = format!("{evaluation:?}"); 1131 assert!(!evaluation_debug.contains(&format!("{claim:?}"))); 1132 assert!(!evaluation_debug.contains("d133a772")); 1133 host.close() 1134 .await 1135 .expect("close before lost-success replay"); 1136 1137 let host = open_writer(&runtime, &metadata).await; 1138 let attempts = host.repositories().reconciliation_attempts(); 1139 let reconciled = attempts 1140 .commit_source_replays(lease, plan.clone(), [retry_replay]) 1141 .await 1142 .expect("idempotent reconcile"); 1143 assert!(!reconciled.created()); 1144 assert_eq!(reconciled.source_result_count(), 1); 1145 assert_eq!(reconciled.checkpoint_advance_count(), 1); 1146 assert!(!reconciled.dirty_generation_advanced()); 1147 let reconciled_manifest = reconciled 1148 .into_evidence_manifest( 1149 UnixTimeSeconds::new(1_784_347_203), 1150 RhiReconciliationScopePrerequisites::Satisfied, 1151 ) 1152 .expect("idempotent manifest"); 1153 assert_eq!(reconciled_manifest.canonical_bytes(), canonical_manifest); 1154 let too_early = attempts 1155 .commit_source_replays(lease, plan, [early_replay]) 1156 .await 1157 .expect("second idempotent reconcile") 1158 .into_evidence_manifest( 1159 UnixTimeSeconds::new(1_784_347_201), 1160 RhiReconciliationScopePrerequisites::Unsatisfied, 1161 ) 1162 .expect_err("observation precedes source completion"); 1163 assert_eq!( 1164 too_early.kind(), 1165 rhi::RhiReconciliationManifestErrorKind::InvalidObservationTime 1166 ); 1167 host.close().await.expect("close"); 1168 1169 let options = SqliteConnectOptions::new() 1170 .filename(runtime.artifacts().state_database()) 1171 .create_if_missing(false) 1172 .foreign_keys(true); 1173 let mut connection = SqliteConnection::connect_with(&options) 1174 .await 1175 .expect("offline fixture connection"); 1176 let attempt_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM evidence_reconciliations") 1177 .fetch_one(&mut connection) 1178 .await 1179 .expect("attempt count"); 1180 let source_count: i64 = 1181 sqlx::query_scalar("SELECT COUNT(*) FROM evidence_reconciliation_sources") 1182 .fetch_one(&mut connection) 1183 .await 1184 .expect("source count"); 1185 let checkpoint_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM relay_checkpoints") 1186 .fetch_one(&mut connection) 1187 .await 1188 .expect("checkpoint count"); 1189 let generation: i64 = sqlx::query_scalar("SELECT generation FROM trade_dirty_generations") 1190 .fetch_one(&mut connection) 1191 .await 1192 .expect("generation"); 1193 assert_eq!( 1194 (attempt_count, source_count, checkpoint_count, generation), 1195 (1, 1, 1, 2) 1196 ); 1197 connection.close().await.expect("fixture close"); 1198 drop((configuration, metadata, runtime, root)); 1199 } 1200 1201 #[tokio::test] 1202 async fn concurrent_exact_commits_converge_to_one_attempt_and_one_manifest() { 1203 let started_ms = 1_784_347_400_000; 1204 let (root, runtime, metadata, configuration, host, lease, plan) = 1205 replay_fixture("replay-concurrent-commit", EXAMPLE, started_ms).await; 1206 let request = &plan.requests()[0]; 1207 let wire = replay_wire(); 1208 let make_replay = || { 1209 RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None) 1210 .expect("replay plan") 1211 .finish( 1212 request, 1213 RhiTradeSourceCompletion::Complete, 1214 now(started_ms), 1215 now(started_ms + 2_000), 1216 [admitted_replay(&configuration, &wire, 1_784_347_400)], 1217 ) 1218 .expect("replay") 1219 }; 1220 let left_replay = make_replay(); 1221 let right_replay = make_replay(); 1222 let attempts = host.repositories().reconciliation_attempts(); 1223 let (left, right) = tokio::join!( 1224 attempts.commit_source_replays(lease, plan.clone(), [left_replay]), 1225 attempts.commit_source_replays(lease, plan.clone(), [right_replay]), 1226 ); 1227 let left = left.expect("left exact commit"); 1228 let right = right.expect("right exact commit"); 1229 assert_eq!(u8::from(left.created()) + u8::from(right.created()), 1); 1230 assert_eq!( 1231 u8::from(left.dirty_generation_advanced()) + u8::from(right.dirty_generation_advanced()), 1232 1 1233 ); 1234 let left = left 1235 .into_evidence_manifest( 1236 UnixTimeSeconds::new(1_784_347_403), 1237 RhiReconciliationScopePrerequisites::Satisfied, 1238 ) 1239 .expect("left manifest"); 1240 let right = right 1241 .into_evidence_manifest( 1242 UnixTimeSeconds::new(1_784_347_403), 1243 RhiReconciliationScopePrerequisites::Satisfied, 1244 ) 1245 .expect("right manifest"); 1246 assert_eq!(left.digest(), right.digest()); 1247 assert_eq!(left.canonical_bytes(), right.canonical_bytes()); 1248 host.close().await.expect("close"); 1249 1250 let mut connection = fixture_connection(&runtime).await; 1251 let counts: (i64, i64, i64) = sqlx::query_as( 1252 r#"SELECT 1253 (SELECT COUNT(*) FROM evidence_reconciliations), 1254 (SELECT COUNT(*) FROM evidence_reconciliation_sources), 1255 (SELECT generation FROM trade_dirty_generations)"#, 1256 ) 1257 .fetch_one(&mut connection) 1258 .await 1259 .expect("durable converged counts"); 1260 assert_eq!(counts, (1, 1, 2)); 1261 connection.close().await.expect("fixture close"); 1262 drop((configuration, metadata, runtime, root)); 1263 } 1264 1265 #[tokio::test] 1266 async fn cancelled_blocked_commit_has_no_effect_and_exact_retry_succeeds() { 1267 let started_ms = 1_784_347_500_000; 1268 let (root, runtime, metadata, configuration, host, lease, plan) = 1269 replay_fixture("replay-cancelled-commit", EXAMPLE, started_ms).await; 1270 let request = &plan.requests()[0]; 1271 let wire = replay_wire(); 1272 let make_replay = || { 1273 RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None) 1274 .expect("replay plan") 1275 .finish( 1276 request, 1277 RhiTradeSourceCompletion::Complete, 1278 now(started_ms), 1279 now(started_ms + 2_000), 1280 [admitted_replay(&configuration, &wire, 1_784_347_500)], 1281 ) 1282 .expect("replay") 1283 }; 1284 let cancelled_replay = make_replay(); 1285 let retry_replay = make_replay(); 1286 let mut blocker = fixture_connection(&runtime).await; 1287 sqlx::query("BEGIN IMMEDIATE") 1288 .execute(&mut blocker) 1289 .await 1290 .expect("exclusive SQLite write blocker"); 1291 1292 let attempts = host.repositories().reconciliation_attempts(); 1293 let cancelled = tokio::time::timeout( 1294 Duration::from_millis(100), 1295 attempts.commit_source_replays(lease, plan.clone(), [cancelled_replay]), 1296 ) 1297 .await; 1298 assert!(cancelled.is_err(), "blocked commit must remain cancellable"); 1299 let attempt_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM evidence_reconciliations") 1300 .fetch_one(&mut blocker) 1301 .await 1302 .expect("no attempt before blocker release"); 1303 assert_eq!(attempt_count, 0); 1304 sqlx::query("ROLLBACK") 1305 .execute(&mut blocker) 1306 .await 1307 .expect("release blocker"); 1308 blocker.close().await.expect("blocker close"); 1309 tokio::task::yield_now().await; 1310 1311 let committed = attempts 1312 .commit_source_replays(lease, plan, [retry_replay]) 1313 .await 1314 .expect("retry after cancellation"); 1315 assert!(committed.created()); 1316 assert_eq!(committed.source_result_count(), 1); 1317 host.close().await.expect("close"); 1318 drop((configuration, metadata, runtime, root)); 1319 } 1320 1321 #[tokio::test] 1322 async fn commit_inventory_bounds_infinite_iterators_before_any_mutation() { 1323 let started_ms = 1_784_347_600_000; 1324 let (root, runtime, metadata, configuration, host, lease, plan) = 1325 replay_fixture("replay-bounded-commit", EXAMPLE, started_ms).await; 1326 let request = &plan.requests()[0]; 1327 let wire = replay_wire(); 1328 let make_replay = || { 1329 RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None) 1330 .expect("replay plan") 1331 .finish( 1332 request, 1333 RhiTradeSourceCompletion::Complete, 1334 now(started_ms), 1335 now(started_ms + 2_000), 1336 [admitted_replay(&configuration, &wire, 1_784_347_600)], 1337 ) 1338 .expect("replay") 1339 }; 1340 let exact_replay = make_replay(); 1341 let error = host 1342 .repositories() 1343 .reconciliation_attempts() 1344 .commit_source_replays(lease, plan.clone(), std::iter::repeat_with(make_replay)) 1345 .await 1346 .expect_err("infinite inventory exceeds the exact source count"); 1347 assert_eq!(error.kind(), RhiReconciliationCommitErrorKind::InvalidInput); 1348 1349 let committed = host 1350 .repositories() 1351 .reconciliation_attempts() 1352 .commit_source_replays(lease, plan, [exact_replay]) 1353 .await 1354 .expect("exact bounded retry"); 1355 assert!(committed.created()); 1356 host.close().await.expect("close"); 1357 drop((configuration, metadata, runtime, root)); 1358 } 1359 1360 #[tokio::test] 1361 async fn incomplete_results_never_advance_and_stale_leases_fail_closed() { 1362 let started_ms = 1_784_347_200_000; 1363 let (root, runtime, metadata, configuration, host, lease, plan) = 1364 replay_fixture("replay-incomplete", EXAMPLE, started_ms).await; 1365 let request = &plan.requests()[0]; 1366 let wire = replay_wire(); 1367 let incomplete = 1368 RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None) 1369 .expect("replay plan") 1370 .finish( 1371 request, 1372 RhiTradeSourceCompletion::IncompleteTimeout, 1373 now(started_ms), 1374 now(started_ms + 10_000), 1375 [admitted_replay(&configuration, &wire, 1_784_347_200)], 1376 ) 1377 .expect("incomplete replay"); 1378 let committed = host 1379 .repositories() 1380 .reconciliation_attempts() 1381 .commit_source_replays(lease, plan, [incomplete]) 1382 .await 1383 .expect("incomplete commit"); 1384 assert_eq!(committed.checkpoint_advance_count(), 0); 1385 assert!(committed.committed_cursors().is_empty()); 1386 assert!(committed.dirty_generation_advanced()); 1387 host.close().await.expect("close"); 1388 drop((configuration, metadata, runtime, root)); 1389 1390 let started_ms = 1_784_347_300_000; 1391 let (root, runtime, metadata, configuration, host, lease, plan) = 1392 replay_fixture("replay-stale-lease", EXAMPLE, started_ms).await; 1393 let request = &plan.requests()[0]; 1394 let replay = 1395 RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None) 1396 .expect("replay plan") 1397 .finish( 1398 request, 1399 RhiTradeSourceCompletion::Complete, 1400 now(started_ms), 1401 now(started_ms + 2_000), 1402 [], 1403 ) 1404 .expect("empty replay"); 1405 host.repositories() 1406 .reconciliation_jobs() 1407 .renew(lease, now(started_ms + 20_000)) 1408 .await 1409 .expect("renewed lease"); 1410 let error = host 1411 .repositories() 1412 .reconciliation_attempts() 1413 .commit_source_replays(lease, plan, [replay]) 1414 .await 1415 .expect_err("stale lease"); 1416 assert_eq!(error.kind(), RhiReconciliationCommitErrorKind::LeaseLost); 1417 assert!(Error::source(&error).is_none()); 1418 host.close().await.expect("close"); 1419 let options = SqliteConnectOptions::new() 1420 .filename(runtime.artifacts().state_database()) 1421 .create_if_missing(false) 1422 .foreign_keys(true); 1423 let mut connection = SqliteConnection::connect_with(&options) 1424 .await 1425 .expect("offline fixture connection"); 1426 let attempt_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM evidence_reconciliations") 1427 .fetch_one(&mut connection) 1428 .await 1429 .expect("attempt count"); 1430 let checkpoint_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM relay_checkpoints") 1431 .fetch_one(&mut connection) 1432 .await 1433 .expect("checkpoint count"); 1434 let generation: i64 = sqlx::query_scalar("SELECT generation FROM trade_dirty_generations") 1435 .fetch_one(&mut connection) 1436 .await 1437 .expect("generation"); 1438 assert_eq!((attempt_count, checkpoint_count, generation), (0, 0, 1)); 1439 connection.close().await.expect("fixture close"); 1440 drop((configuration, metadata, runtime, root)); 1441 } 1442 1443 #[tokio::test] 1444 async fn finalization_preflight_is_exact_nonmutating_and_redacted() { 1445 let (root, runtime, metadata, configuration, host, lease, evaluation) = 1446 finalization_fixture("finalization-exact").await; 1447 let before = finalization_snapshot(&runtime).await; 1448 let fence = host 1449 .repositories() 1450 .reconciliation_attempts() 1451 .prepare_finalization(lease, evaluation, now(1_784_347_206_000)) 1452 .await 1453 .expect("finalization preflight"); 1454 assert_eq!(fence.contract_version(), 1); 1455 assert_eq!( 1456 fence 1457 .evaluation() 1458 .projection() 1459 .manifest() 1460 .trade_generation(), 1461 2 1462 ); 1463 assert!(fence.evaluation().projection().digest().is_some()); 1464 let rendered = format!("{fence:?}"); 1465 assert!(!rendered.contains("11111111")); 1466 assert!(!rendered.contains("reconciliation_attempt")); 1467 assert_eq!(finalization_snapshot(&runtime).await, before); 1468 host.close().await.expect("close"); 1469 drop((fence, configuration, metadata, runtime, root)); 1470 } 1471 1472 #[tokio::test] 1473 async fn finalization_rejects_stale_generation_policy_and_expired_or_lost_leases() { 1474 let (root, runtime, metadata, configuration, host, lease, evaluation) = 1475 finalization_fixture("finalization-generation").await; 1476 write_dirty( 1477 &runtime, 1478 lease.job().trade_id(), 1479 3, 1480 *lease.job().evidence_policy_digest().as_bytes(), 1481 1_784_347_207, 1482 ) 1483 .await; 1484 let error = host 1485 .repositories() 1486 .reconciliation_attempts() 1487 .prepare_finalization(lease, evaluation, now(1_784_347_207_000)) 1488 .await 1489 .expect_err("stale generation"); 1490 assert_eq!( 1491 error.kind(), 1492 RhiReconciliationFinalizationErrorKind::GenerationConflict 1493 ); 1494 assert!(Error::source(&error).is_none()); 1495 assert!(!format!("{error} {error:?}").contains("11111111")); 1496 host.close().await.expect("close"); 1497 drop((configuration, metadata, runtime, root)); 1498 1499 let (root, runtime, metadata, configuration, host, lease, evaluation) = 1500 finalization_fixture("finalization-policy").await; 1501 write_dirty( 1502 &runtime, 1503 lease.job().trade_id(), 1504 3, 1505 [0xa5; 32], 1506 1_784_347_207, 1507 ) 1508 .await; 1509 let error = host 1510 .repositories() 1511 .reconciliation_attempts() 1512 .prepare_finalization(lease, evaluation, now(1_784_347_207_000)) 1513 .await 1514 .expect_err("stale policy"); 1515 assert_eq!( 1516 error.kind(), 1517 RhiReconciliationFinalizationErrorKind::GenerationConflict 1518 ); 1519 host.close().await.expect("close"); 1520 drop((configuration, metadata, runtime, root)); 1521 1522 let (root, runtime, metadata, configuration, host, lease, evaluation) = 1523 finalization_fixture("finalization-expired").await; 1524 let error = host 1525 .repositories() 1526 .reconciliation_attempts() 1527 .prepare_finalization(lease, evaluation, lease.lease_expires()) 1528 .await 1529 .expect_err("expired lease"); 1530 assert_eq!( 1531 error.kind(), 1532 RhiReconciliationFinalizationErrorKind::LeaseLost 1533 ); 1534 host.close().await.expect("close"); 1535 drop((configuration, metadata, runtime, root)); 1536 1537 let (root, runtime, metadata, configuration, host, lease, evaluation) = 1538 finalization_fixture("finalization-lost").await; 1539 host.repositories() 1540 .reconciliation_jobs() 1541 .record_failure( 1542 lease, 1543 now(1_784_347_206_000), 1544 RhiReconciliationRetryDelayMilliseconds::new(1).expect("delay"), 1545 ) 1546 .await 1547 .expect("release lease"); 1548 let error = host 1549 .repositories() 1550 .reconciliation_attempts() 1551 .prepare_finalization(lease, evaluation, now(1_784_347_206_001)) 1552 .await 1553 .expect_err("lost lease"); 1554 assert_eq!( 1555 error.kind(), 1556 RhiReconciliationFinalizationErrorKind::LeaseLost 1557 ); 1558 host.close().await.expect("close"); 1559 drop((configuration, metadata, runtime, root)); 1560 } 1561 1562 #[tokio::test] 1563 async fn finalization_rejects_cross_attempt_relabelling_and_read_only_hosts() { 1564 let (root, runtime, metadata, configuration, host, lease, evaluation) = 1565 finalization_fixture("finalization-attempt").await; 1566 let jobs = host.repositories().reconciliation_jobs(); 1567 jobs.record_failure( 1568 lease, 1569 now(1_784_347_206_000), 1570 RhiReconciliationRetryDelayMilliseconds::new(1).expect("delay"), 1571 ) 1572 .await 1573 .expect("retry schedule"); 1574 let next_lease = jobs 1575 .claim_next(owner(0x93), now(1_784_347_206_001)) 1576 .await 1577 .expect("next claim") 1578 .expect("reclaimed job"); 1579 assert_eq!(next_lease.job().attempt_count(), 2); 1580 let error = host 1581 .repositories() 1582 .reconciliation_attempts() 1583 .prepare_finalization(next_lease, evaluation, now(1_784_347_206_002)) 1584 .await 1585 .expect_err("old evaluation cannot be relabelled"); 1586 assert_eq!( 1587 error.kind(), 1588 RhiReconciliationFinalizationErrorKind::InvalidInput 1589 ); 1590 host.close().await.expect("close"); 1591 drop((configuration, metadata, runtime, root)); 1592 1593 let (root, runtime, metadata, configuration, host, lease, evaluation) = 1594 finalization_fixture("finalization-inspection").await; 1595 host.close().await.expect("close writer"); 1596 let inspection = open_rhi_state_inspection(&runtime, &metadata) 1597 .await 1598 .expect("inspection"); 1599 let error = inspection 1600 .repositories() 1601 .reconciliation_attempts() 1602 .prepare_finalization(lease, evaluation, now(1_784_347_206_000)) 1603 .await 1604 .expect_err("inspection cannot prepare finalization"); 1605 assert_eq!( 1606 error.kind(), 1607 RhiReconciliationFinalizationErrorKind::InvalidMode 1608 ); 1609 inspection.close().await.expect("inspection close"); 1610 drop((configuration, metadata, runtime, root)); 1611 } 1612 1613 struct FixedAttestationEntropy(u8); 1614 1615 impl EntropySource for FixedAttestationEntropy { 1616 fn fill_bytes(&self, destination: &mut [u8]) -> Result<(), EntropyError> { 1617 destination.fill(self.0); 1618 Ok(()) 1619 } 1620 } 1621 1622 struct FailingAttestationEntropy; 1623 1624 impl EntropySource for FailingAttestationEntropy { 1625 fn fill_bytes(&self, _destination: &mut [u8]) -> Result<(), EntropyError> { 1626 Err(EntropyError::Unavailable) 1627 } 1628 } 1629 1630 #[derive(Clone, Copy)] 1631 struct FixedPublicationWall(u64); 1632 1633 impl WallClock for FixedPublicationWall { 1634 fn now_utc(&self) -> Result<UnixTimeSeconds, WallClockError> { 1635 Ok(UnixTimeSeconds::new(self.0)) 1636 } 1637 } 1638 1639 struct RecordingExactPublicationSink { 1640 expected: Arc<Vec<u8>>, 1641 calls: Arc<AtomicUsize>, 1642 outcome: RhiPublicationAttemptOutcome, 1643 } 1644 1645 impl RhiExactPublicationSink for RecordingExactPublicationSink { 1646 fn submit_exact<'a>( 1647 &'a self, 1648 attempt: &'a RhiPreparedPublicationAttempt, 1649 ) -> BoxFuture<'a, RhiPublicationAttemptOutcome> { 1650 Box::pin(async move { 1651 assert_eq!(attempt.exact_signed_event_bytes(), self.expected.as_slice()); 1652 assert_eq!(attempt.relay_id(), "relay-primary"); 1653 assert!(attempt.deadline_at().get() > 0); 1654 self.calls.fetch_add(1, Ordering::SeqCst); 1655 self.outcome 1656 }) 1657 } 1658 } 1659 1660 struct PendingExactPublicationSink; 1661 1662 impl RhiExactPublicationSink for PendingExactPublicationSink { 1663 fn submit_exact<'a>( 1664 &'a self, 1665 _attempt: &'a RhiPreparedPublicationAttempt, 1666 ) -> BoxFuture<'a, RhiPublicationAttemptOutcome> { 1667 Box::pin(std::future::pending()) 1668 } 1669 } 1670 1671 struct CoordinatedExactPublicationSink { 1672 expected: Arc<Vec<u8>>, 1673 calls: Arc<AtomicUsize>, 1674 started: Arc<Notify>, 1675 release: Arc<Notify>, 1676 } 1677 1678 impl RhiExactPublicationSink for CoordinatedExactPublicationSink { 1679 fn submit_exact<'a>( 1680 &'a self, 1681 attempt: &'a RhiPreparedPublicationAttempt, 1682 ) -> BoxFuture<'a, RhiPublicationAttemptOutcome> { 1683 Box::pin(async move { 1684 assert_eq!(attempt.exact_signed_event_bytes(), self.expected.as_slice()); 1685 self.calls.fetch_add(1, Ordering::SeqCst); 1686 self.started.notify_one(); 1687 self.release.notified().await; 1688 RhiPublicationAttemptOutcome::Accepted 1689 }) 1690 } 1691 } 1692 1693 fn publication_now(value: u64) -> RhiPublicationUnixMilliseconds { 1694 RhiPublicationUnixMilliseconds::new(value).expect("publication time") 1695 } 1696 1697 fn publication_owner(byte: u8) -> RhiPublicationLeaseOwner { 1698 RhiPublicationLeaseOwner::from_bytes([byte; 16]).expect("publication owner") 1699 } 1700 1701 fn publication_adapters(seconds: u64) -> RhiTimeEntropyAdapters { 1702 RhiTimeEntropyAdapters::new( 1703 FixedPublicationWall(seconds), 1704 SystemMonotonicClock::new(), 1705 FixedAttestationEntropy(0xff), 1706 ) 1707 } 1708 1709 fn attestation_secret() -> [u8; 32] { 1710 [1; 32] 1711 } 1712 1713 fn attestation_configuration(runtime: &RhiRuntimeContext, expected_public_key: &str) -> String { 1714 EXAMPLE 1715 .replace( 1716 "/var/lib/radroots/services/rhi/default/secrets/service.identity.ncrypt", 1717 runtime 1718 .identity_path() 1719 .to_str() 1720 .expect("UTF-8 identity path"), 1721 ) 1722 .replace(&"2".repeat(64), expected_public_key) 1723 } 1724 1725 fn provision_attestation_identity( 1726 runtime: &RhiRuntimeContext, 1727 configuration: &rhi::RhiConfigDocumentV1, 1728 metadata: &RhiStateMetadata, 1729 ) -> RhiDecryptedIdentity { 1730 fs::create_dir_all(runtime.context().paths().secrets()).expect("secrets directory"); 1731 fs::set_permissions( 1732 runtime.context().paths().secrets(), 1733 fs::Permissions::from_mode(0o700), 1734 ) 1735 .expect("secrets directory mode"); 1736 let credential_path = runtime 1737 .context() 1738 .paths() 1739 .secrets() 1740 .join("service_wrapping_key"); 1741 fs::write(&credential_path, [0x81; 32]).expect("wrapping credential"); 1742 fs::set_permissions(&credential_path, fs::Permissions::from_mode(0o600)) 1743 .expect("wrapping credential mode"); 1744 let binding = RhiIdentityEnvelopeBinding::from_configuration(configuration, metadata) 1745 .expect("identity binding"); 1746 let credential = 1747 resolve_rhi_wrapping_credential(runtime, &binding).expect("wrapping credential open"); 1748 provision_rhi_encrypted_identity( 1749 &binding, 1750 &credential, 1751 RhiEncryptedIdentityProvisioningMaterial::new( 1752 attestation_secret(), 1753 [0x42; 32], 1754 [0x43; 24], 1755 [0x44; 24], 1756 ) 1757 .expect("provisioning material"), 1758 ) 1759 .expect("identity provisioning") 1760 } 1761 1762 #[tokio::test] 1763 async fn signed_attestation_is_canonical_exact_verified_and_nonmutating() { 1764 let expected_public_key = 1765 Keys::new(SecretKey::from_slice(&attestation_secret()).expect("identity secret")) 1766 .public_key() 1767 .to_hex(); 1768 let (root, runtime, metadata, configuration, host, lease, evaluation) = 1769 finalization_fixture_with_source("attestation-exact", |runtime| { 1770 attestation_configuration(runtime, &expected_public_key) 1771 }) 1772 .await; 1773 let identity = provision_attestation_identity(&runtime, &configuration, &metadata); 1774 let before = finalization_snapshot(&runtime).await; 1775 let fence = host 1776 .repositories() 1777 .reconciliation_attempts() 1778 .prepare_finalization(lease, evaluation, now(1_784_347_206_000)) 1779 .await 1780 .expect("finalization fence"); 1781 let signed = build_rhi_signed_evidence_attestation( 1782 fence, 1783 &identity, 1784 UnixTimeSeconds::new(1_784_347_207), 1785 &FixedAttestationEntropy(0xa5), 1786 None, 1787 ) 1788 .expect("signed attestation"); 1789 1790 assert_eq!(signed.contract_version(), 1); 1791 assert_eq!(signed.created_at_unix_seconds(), 1_784_347_207); 1792 assert!(!signed.has_supersession()); 1793 assert!(!signed.signed_event_bytes().is_empty()); 1794 assert!(signed.signed_event_bytes().len() <= 32 * 1_024); 1795 assert_eq!( 1796 signed.signed_event_sha256(), 1797 &<[u8; 32]>::from(Sha256::digest(signed.signed_event_bytes())) 1798 ); 1799 let event: serde_json::Value = 1800 serde_json::from_slice(signed.signed_event_bytes()).expect("signed event JSON"); 1801 assert_eq!( 1802 signed.signed_event_bytes(), 1803 SIGNED_ATTESTATION_VECTOR.trim_end().as_bytes() 1804 ); 1805 assert_eq!(event["id"], lower_hex(signed.event_id())); 1806 assert_eq!(event["pubkey"], expected_public_key); 1807 assert_eq!(event["created_at"], 1_784_347_207_u64); 1808 assert_eq!(event["kind"], 3_441); 1809 assert_eq!(event["tags"].as_array().expect("tags").len(), 5); 1810 assert_eq!( 1811 event["content"].as_str().expect("content").as_bytes(), 1812 signed.canonical_report_bytes() 1813 ); 1814 let report: serde_json::Value = 1815 serde_json::from_slice(signed.canonical_report_bytes()).expect("canonical report"); 1816 assert_eq!(report["issuer_pubkey"], expected_public_key); 1817 assert_eq!(report["trade_generation"], 2); 1818 assert_eq!(report["report_id"], lower_hex(&signed.statement_digest())); 1819 assert_eq!(report["statement_digest"], report["report_id"]); 1820 assert!(report["supersedes_report_id"].is_null()); 1821 assert!(report["supersedes_event_id"].is_null()); 1822 let rendered = format!("{signed:?}"); 1823 assert!(!rendered.contains(&lower_hex(signed.event_id()))); 1824 assert!(!rendered.contains(&lower_hex(&signed.statement_digest()))); 1825 assert!(!rendered.contains(&expected_public_key)); 1826 assert_eq!(finalization_snapshot(&runtime).await, before); 1827 1828 let supersession = RhiEvidenceAttestationSupersession::from_attestation(&signed); 1829 assert_eq!( 1830 format!("{supersession:?}"), 1831 "RhiEvidenceAttestationSupersession([redacted])" 1832 ); 1833 host.close().await.expect("close"); 1834 drop(( 1835 supersession, 1836 signed, 1837 identity, 1838 configuration, 1839 metadata, 1840 runtime, 1841 root, 1842 )); 1843 } 1844 1845 #[tokio::test] 1846 async fn signed_attestation_fails_closed_when_injected_entropy_is_unavailable() { 1847 let expected_public_key = 1848 Keys::new(SecretKey::from_slice(&attestation_secret()).expect("identity secret")) 1849 .public_key() 1850 .to_hex(); 1851 let (root, runtime, metadata, configuration, host, lease, evaluation) = 1852 finalization_fixture_with_source("attestation-entropy", |runtime| { 1853 attestation_configuration(runtime, &expected_public_key) 1854 }) 1855 .await; 1856 let identity = provision_attestation_identity(&runtime, &configuration, &metadata); 1857 let before = finalization_snapshot(&runtime).await; 1858 let fence = host 1859 .repositories() 1860 .reconciliation_attempts() 1861 .prepare_finalization(lease, evaluation, now(1_784_347_206_000)) 1862 .await 1863 .expect("finalization fence"); 1864 let error = build_rhi_signed_evidence_attestation( 1865 fence, 1866 &identity, 1867 UnixTimeSeconds::new(1_784_347_207), 1868 &FailingAttestationEntropy, 1869 None, 1870 ) 1871 .expect_err("entropy failure"); 1872 assert_eq!( 1873 error.kind(), 1874 RhiReconciliationAttestationErrorKind::EntropyUnavailable 1875 ); 1876 assert_eq!( 1877 error.code(), 1878 "reconciliation_attestation_entropy_unavailable" 1879 ); 1880 assert!(Error::source(&error).is_none()); 1881 assert!(!format!("{error} {error:?}").contains(&expected_public_key)); 1882 assert_eq!(finalization_snapshot(&runtime).await, before); 1883 host.close().await.expect("close"); 1884 drop((identity, configuration, metadata, runtime, root)); 1885 } 1886 1887 #[tokio::test] 1888 async fn signed_attestation_supersession_is_verified_ordered_and_explicit() { 1889 let expected_public_key = 1890 Keys::new(SecretKey::from_slice(&attestation_secret()).expect("identity secret")) 1891 .public_key() 1892 .to_hex(); 1893 let ( 1894 prior_root, 1895 prior_runtime, 1896 prior_metadata, 1897 prior_config, 1898 prior_host, 1899 prior_lease, 1900 prior_eval, 1901 ) = finalization_fixture_with_source_and_generation("attestation-prior", 1, |runtime| { 1902 attestation_configuration(runtime, &expected_public_key) 1903 }) 1904 .await; 1905 let prior_identity = 1906 provision_attestation_identity(&prior_runtime, &prior_config, &prior_metadata); 1907 let prior_fence = prior_host 1908 .repositories() 1909 .reconciliation_attempts() 1910 .prepare_finalization(prior_lease, prior_eval, now(1_784_347_206_000)) 1911 .await 1912 .expect("prior fence"); 1913 let prior = build_rhi_signed_evidence_attestation( 1914 prior_fence, 1915 &prior_identity, 1916 UnixTimeSeconds::new(1_784_347_207), 1917 &FixedAttestationEntropy(0xa6), 1918 None, 1919 ) 1920 .expect("prior attestation"); 1921 1922 let (next_root, next_runtime, next_metadata, next_config, next_host, next_lease, next_eval) = 1923 finalization_fixture_with_source_and_generation("attestation-next", 2, |runtime| { 1924 attestation_configuration(runtime, &expected_public_key) 1925 }) 1926 .await; 1927 let next_identity = provision_attestation_identity(&next_runtime, &next_config, &next_metadata); 1928 let next_fence = next_host 1929 .repositories() 1930 .reconciliation_attempts() 1931 .prepare_finalization(next_lease, next_eval, now(1_784_347_206_000)) 1932 .await 1933 .expect("next fence"); 1934 let next = build_rhi_signed_evidence_attestation( 1935 next_fence, 1936 &next_identity, 1937 UnixTimeSeconds::new(1_784_347_208), 1938 &FixedAttestationEntropy(0xa7), 1939 Some(RhiEvidenceAttestationSupersession::from_attestation(&prior)), 1940 ) 1941 .expect("superseding attestation"); 1942 assert!(next.has_supersession()); 1943 let next_event: serde_json::Value = 1944 serde_json::from_slice(next.signed_event_bytes()).expect("next event"); 1945 assert_eq!(next_event["tags"].as_array().expect("next tags").len(), 7); 1946 let next_report: serde_json::Value = 1947 serde_json::from_slice(next.canonical_report_bytes()).expect("next report"); 1948 assert_eq!( 1949 next_report["supersedes_report_id"], 1950 lower_hex(&prior.statement_digest()) 1951 ); 1952 assert_eq!( 1953 next_report["supersedes_event_id"], 1954 lower_hex(prior.event_id()) 1955 ); 1956 assert_eq!(next_report["trade_generation"], 3); 1957 1958 let ( 1959 stale_root, 1960 stale_runtime, 1961 stale_metadata, 1962 stale_config, 1963 stale_host, 1964 stale_lease, 1965 stale_eval, 1966 ) = finalization_fixture_with_source_and_generation("attestation-stale", 1, |runtime| { 1967 attestation_configuration(runtime, &expected_public_key) 1968 }) 1969 .await; 1970 let stale_identity = 1971 provision_attestation_identity(&stale_runtime, &stale_config, &stale_metadata); 1972 let stale_fence = stale_host 1973 .repositories() 1974 .reconciliation_attempts() 1975 .prepare_finalization(stale_lease, stale_eval, now(1_784_347_206_000)) 1976 .await 1977 .expect("stale fence"); 1978 let error = build_rhi_signed_evidence_attestation( 1979 stale_fence, 1980 &stale_identity, 1981 UnixTimeSeconds::new(1_784_347_209), 1982 &FixedAttestationEntropy(0xa8), 1983 Some(RhiEvidenceAttestationSupersession::from_attestation(&next)), 1984 ) 1985 .expect_err("older generation cannot supersede"); 1986 assert_eq!( 1987 error.kind(), 1988 RhiReconciliationAttestationErrorKind::SupersessionInvalid 1989 ); 1990 assert!(Error::source(&error).is_none()); 1991 1992 prior_host.close().await.expect("prior close"); 1993 next_host.close().await.expect("next close"); 1994 stale_host.close().await.expect("stale close"); 1995 drop(( 1996 prior, 1997 next, 1998 prior_identity, 1999 next_identity, 2000 stale_identity, 2001 prior_config, 2002 next_config, 2003 stale_config, 2004 prior_metadata, 2005 next_metadata, 2006 stale_metadata, 2007 prior_runtime, 2008 next_runtime, 2009 stale_runtime, 2010 prior_root, 2011 next_root, 2012 stale_root, 2013 )); 2014 } 2015 2016 #[tokio::test] 2017 async fn atomic_finalization_commits_exact_required_inventory_and_reconciles_retry() { 2018 let ( 2019 root, 2020 runtime, 2021 metadata, 2022 configuration, 2023 host, 2024 lease, 2025 signed, 2026 publication, 2027 manifest, 2028 _trade, 2029 ) = signed_finalization_fixture("finalization-commit-required", false).await; 2030 let repositories = host.repositories(); 2031 let first = repositories 2032 .reconciliation_attempts() 2033 .commit_finalization(&signed, &publication, now(1_784_347_208_000)) 2034 .await 2035 .expect("atomic finalization"); 2036 assert!(first.created()); 2037 assert_eq!(first.publication_mode(), RhiPublicationMode::Required); 2038 assert_eq!(first.target_count(), 2); 2039 let outbox_id = first.outbox_id().expect("required outbox identity"); 2040 let committed = repositories 2041 .publication_outbox() 2042 .read_committed_publication(outbox_id) 2043 .await 2044 .expect("committed exact publication"); 2045 assert_eq!(committed.outbox_id(), outbox_id); 2046 assert_eq!(committed.event_id(), signed.event_id()); 2047 assert_eq!(committed.event_sha256(), signed.signed_event_sha256()); 2048 assert_eq!( 2049 committed.exact_signed_event_bytes(), 2050 signed.signed_event_bytes() 2051 ); 2052 let committed_debug = format!("{committed:?} {outbox_id:?}"); 2053 assert!(!committed_debug.contains("relay-primary")); 2054 assert!(!committed_debug.contains("{\"id\"")); 2055 2056 let mut progressed = fixture_connection(&runtime).await; 2057 let checkpoint = sqlx::query( 2058 r#"UPDATE relay_checkpoints 2059 SET cursor_created_at_unix_s = cursor_created_at_unix_s + 1, 2060 cursor_event_id = ?, revision = revision + 1, 2061 completed_at_unix_s = completed_at_unix_s + 1"#, 2062 ) 2063 .bind([0xfe; 32].as_slice()) 2064 .execute(&mut progressed) 2065 .await 2066 .expect("later checkpoint progression"); 2067 assert_eq!(checkpoint.rows_affected(), 1); 2068 progressed.close().await.expect("progression close"); 2069 2070 let retry = repositories 2071 .reconciliation_attempts() 2072 .commit_finalization(&signed, &publication, lease.lease_expires()) 2073 .await 2074 .expect("exact retry after consumed lease"); 2075 assert!(!retry.created()); 2076 assert_eq!(retry.publication_mode(), RhiPublicationMode::Required); 2077 assert_eq!(retry.target_count(), 2); 2078 assert_eq!(retry.outbox_id(), Some(outbox_id)); 2079 2080 let mut connection = fixture_connection(&runtime).await; 2081 let counts: (i64, i64, i64, i64, i64, i64, i64) = sqlx::query_as( 2082 r#"SELECT 2083 (SELECT COUNT(*) FROM evidence_manifests), 2084 (SELECT COUNT(*) FROM trade_projections), 2085 (SELECT COUNT(*) FROM attestation_reports), 2086 (SELECT COUNT(*) FROM signed_attestation_events), 2087 (SELECT COUNT(*) FROM publication_outbox), 2088 (SELECT COUNT(*) FROM publication_targets), 2089 (SELECT COUNT(*) FROM reconciliation_jobs WHERE state = 'completed')"#, 2090 ) 2091 .fetch_one(&mut connection) 2092 .await 2093 .expect("final inventory counts"); 2094 assert_eq!(counts, (1, 1, 1, 1, 1, 2, 1)); 2095 let exact: (Vec<u8>, Vec<u8>, Vec<u8>, String, i64) = sqlx::query_as( 2096 r#"SELECT manifest.canonical_manifest, report.canonical_report, 2097 event.canonical_event_json, outbox.state, outbox.target_count 2098 FROM evidence_manifests AS manifest 2099 JOIN attestation_reports AS report 2100 ON report.manifest_sha256 = manifest.manifest_sha256 2101 JOIN signed_attestation_events AS event 2102 ON event.statement_sha256 = report.statement_sha256 2103 JOIN publication_outbox AS outbox ON outbox.event_id = event.event_id"#, 2104 ) 2105 .fetch_one(&mut connection) 2106 .await 2107 .expect("exact final inventory"); 2108 assert_eq!(exact.0.as_slice(), manifest.as_ref()); 2109 assert_eq!(exact.1.as_slice(), signed.canonical_report_bytes()); 2110 assert_eq!(exact.2.as_slice(), signed.signed_event_bytes()); 2111 assert_eq!(exact.3, "pending"); 2112 assert_eq!(exact.4, 2); 2113 connection.close().await.expect("fixture close"); 2114 2115 let exact_signed_event_bytes = signed.signed_event_bytes().to_vec(); 2116 host.close().await.expect("host close"); 2117 let inspection = open_rhi_state_inspection(&runtime, &metadata) 2118 .await 2119 .expect("reopened inspection"); 2120 let recovered = inspection 2121 .repositories() 2122 .publication_outbox() 2123 .read_committed_publication(outbox_id) 2124 .await 2125 .expect("recovered exact publication"); 2126 assert_eq!( 2127 recovered.exact_signed_event_bytes(), 2128 exact_signed_event_bytes 2129 ); 2130 assert_eq!(recovered.event_sha256(), signed.signed_event_sha256()); 2131 inspection.close().await.expect("inspection close"); 2132 drop(( 2133 signed, 2134 publication, 2135 manifest, 2136 configuration, 2137 metadata, 2138 runtime, 2139 root, 2140 )); 2141 } 2142 2143 #[tokio::test] 2144 async fn exact_byte_publication_persists_submitted_before_io_and_commits_accepted() { 2145 let ( 2146 root, 2147 runtime, 2148 metadata, 2149 configuration, 2150 host, 2151 lease, 2152 signed, 2153 publication, 2154 manifest, 2155 _trade, 2156 ) = signed_finalization_fixture("publication-execution-accepted", false).await; 2157 let repositories = host.repositories(); 2158 let finalized = repositories 2159 .reconciliation_attempts() 2160 .commit_finalization(&signed, &publication, now(1_784_347_208_000)) 2161 .await 2162 .expect("finalization"); 2163 let outbox_id = finalized.outbox_id().expect("outbox"); 2164 let expected = Arc::new(signed.signed_event_bytes().to_vec()); 2165 let calls = Arc::new(AtomicUsize::new(0)); 2166 let sink = RecordingExactPublicationSink { 2167 expected: Arc::clone(&expected), 2168 calls: Arc::clone(&calls), 2169 outcome: RhiPublicationAttemptOutcome::Accepted, 2170 }; 2171 let outcome = repositories 2172 .publication_outbox() 2173 .execute_next_publication( 2174 publication_owner(0x91), 2175 &publication_adapters(1_784_347_208), 2176 &sink, 2177 &publication, 2178 ) 2179 .await 2180 .expect("exact publication") 2181 .expect("due publication"); 2182 assert_eq!(calls.load(Ordering::SeqCst), 1); 2183 assert_eq!(outcome.outbox_id(), outbox_id); 2184 assert_eq!(outcome.target_ordinal(), 0); 2185 assert_eq!(outcome.attempt_number(), 1); 2186 assert_eq!(outcome.outcome(), RhiPublicationAttemptOutcome::Accepted); 2187 assert_eq!(outcome.target_state(), RhiPublicationTargetState::Accepted); 2188 assert_eq!(outcome.outbox_state(), RhiPublicationOutboxState::Complete); 2189 2190 let mut connection = fixture_connection(&runtime).await; 2191 let durable: (String, i64, String, i64, i64, String, String) = sqlx::query_as( 2192 r#"SELECT outbox.state, outbox.revision, 2193 target.state, target.revision, target.attempt_count, 2194 attempt.outcome, attempt.result_code 2195 FROM publication_outbox AS outbox 2196 JOIN publication_targets AS target ON target.outbox_id = outbox.outbox_id 2197 JOIN publication_attempts AS attempt 2198 ON attempt.attempt_id = target.last_attempt_id 2199 WHERE outbox.outbox_id = ? AND target.target_ordinal = 0"#, 2200 ) 2201 .bind(outbox_id.as_bytes().as_slice()) 2202 .fetch_one(&mut connection) 2203 .await 2204 .expect("durable publication outcome"); 2205 assert_eq!( 2206 durable, 2207 ( 2208 "complete".into(), 2209 3, 2210 "accepted".into(), 2211 3, 2212 1, 2213 "accepted".into(), 2214 "accepted".into() 2215 ) 2216 ); 2217 let stored: Vec<u8> = sqlx::query_scalar( 2218 r#"SELECT event.canonical_event_json 2219 FROM publication_outbox AS outbox 2220 JOIN signed_attestation_events AS event ON event.event_id = outbox.event_id 2221 WHERE outbox.outbox_id = ?"#, 2222 ) 2223 .bind(outbox_id.as_bytes().as_slice()) 2224 .fetch_one(&mut connection) 2225 .await 2226 .expect("stored exact bytes"); 2227 assert_eq!(stored, expected.as_ref().clone()); 2228 connection.close().await.expect("fixture close"); 2229 host.close().await.expect("host close"); 2230 drop(( 2231 manifest, 2232 publication, 2233 signed, 2234 lease, 2235 configuration, 2236 metadata, 2237 runtime, 2238 root, 2239 )); 2240 } 2241 2242 #[tokio::test] 2243 async fn concurrent_publication_execution_has_one_remote_submitter() { 2244 let ( 2245 root, 2246 runtime, 2247 metadata, 2248 configuration, 2249 host, 2250 lease, 2251 signed, 2252 publication, 2253 manifest, 2254 _trade, 2255 ) = signed_finalization_fixture("publication-execution-concurrency", false).await; 2256 let repositories = host.repositories(); 2257 repositories 2258 .reconciliation_attempts() 2259 .commit_finalization(&signed, &publication, now(1_784_347_208_000)) 2260 .await 2261 .expect("finalization"); 2262 let calls = Arc::new(AtomicUsize::new(0)); 2263 let started = Arc::new(Notify::new()); 2264 let release = Arc::new(Notify::new()); 2265 let sink = CoordinatedExactPublicationSink { 2266 expected: Arc::new(signed.signed_event_bytes().to_vec()), 2267 calls: Arc::clone(&calls), 2268 started: Arc::clone(&started), 2269 release: Arc::clone(&release), 2270 }; 2271 let first_adapters = publication_adapters(1_784_347_208); 2272 let second_adapters = publication_adapters(1_784_347_208); 2273 let first_outbox = repositories.publication_outbox(); 2274 let first = first_outbox.execute_next_publication( 2275 publication_owner(0xb1), 2276 &first_adapters, 2277 &sink, 2278 &publication, 2279 ); 2280 let second = async { 2281 started.notified().await; 2282 let result = repositories 2283 .publication_outbox() 2284 .execute_next_publication( 2285 publication_owner(0xb2), 2286 &second_adapters, 2287 &sink, 2288 &publication, 2289 ) 2290 .await; 2291 release.notify_one(); 2292 result 2293 }; 2294 let (first, second) = tokio::join!(first, second); 2295 let first = first 2296 .expect("first executor") 2297 .expect("first executor claimed due publication"); 2298 assert_eq!(first.outcome(), RhiPublicationAttemptOutcome::Accepted); 2299 assert!(second.expect("second executor").is_none()); 2300 assert_eq!(calls.load(Ordering::SeqCst), 1); 2301 2302 host.close().await.expect("host close"); 2303 drop(( 2304 manifest, 2305 publication, 2306 signed, 2307 lease, 2308 configuration, 2309 metadata, 2310 runtime, 2311 root, 2312 )); 2313 } 2314 2315 #[tokio::test] 2316 async fn publication_queue_capacity_is_checked_before_finalization_mutation() { 2317 let expected_public_key = 2318 Keys::new(SecretKey::from_slice(&attestation_secret()).expect("identity secret")) 2319 .public_key() 2320 .to_hex(); 2321 let (root, runtime, metadata, configuration, host, lease, evaluation) = 2322 finalization_fixture_with_source("publication-queue-capacity", |runtime| { 2323 attestation_configuration(runtime, &expected_public_key).replacen( 2324 "publication = 4096", 2325 "publication = 1", 2326 1, 2327 ) 2328 }) 2329 .await; 2330 let identity = provision_attestation_identity(&runtime, &configuration, &metadata); 2331 let fence = host 2332 .repositories() 2333 .reconciliation_attempts() 2334 .prepare_finalization(lease, evaluation, now(1_784_347_206_000)) 2335 .await 2336 .expect("finalization fence"); 2337 let signed = build_rhi_signed_evidence_attestation( 2338 fence, 2339 &identity, 2340 UnixTimeSeconds::new(1_784_347_207), 2341 &FixedAttestationEntropy(0xb3), 2342 None, 2343 ) 2344 .expect("signed attestation"); 2345 let publication = 2346 RhiPublicationAuthority::from_config(&configuration).expect("publication authority"); 2347 assert_eq!(publication.queue_capacity(), 1); 2348 2349 let options = SqliteConnectOptions::new() 2350 .filename(runtime.artifacts().state_database()) 2351 .create_if_missing(false) 2352 .foreign_keys(false); 2353 let mut connection = SqliteConnection::connect_with(&options) 2354 .await 2355 .expect("offline queue fixture connection"); 2356 sqlx::query( 2357 r#"INSERT INTO publication_outbox ( 2358 outbox_id, event_id, event_sha256, publication_authority_sha256, 2359 target_set_sha256, target_count, required_target_count, 2360 max_attempts, initial_backoff_ms, maximum_backoff_ms, 2361 attempt_deadline_ms, state, revision, next_attempt_unix_ms, 2362 lease_owner, lease_expires_unix_ms, created_at_unix_ms, 2363 updated_at_unix_ms 2364 ) VALUES (?, ?, ?, ?, ?, 1, 1, 3, 100, 1000, 5000, 2365 'pending', 1, 1, NULL, NULL, 1, 1)"#, 2366 ) 2367 .bind([0xc1_u8; 32].as_slice()) 2368 .bind([0xc2_u8; 32].as_slice()) 2369 .bind([0xc3_u8; 32].as_slice()) 2370 .bind([0xc4_u8; 32].as_slice()) 2371 .bind([0xc5_u8; 32].as_slice()) 2372 .execute(&mut connection) 2373 .await 2374 .expect("bounded active outbox fixture"); 2375 connection.close().await.expect("queue fixture close"); 2376 2377 let error = host 2378 .repositories() 2379 .reconciliation_attempts() 2380 .commit_finalization(&signed, &publication, now(1_784_347_208_000)) 2381 .await 2382 .expect_err("full publication queue"); 2383 assert_eq!( 2384 error.kind(), 2385 RhiReconciliationFinalizationCommitErrorKind::PublicationQueueFull 2386 ); 2387 assert!(Error::source(&error).is_none()); 2388 let mut connection = fixture_connection(&runtime).await; 2389 let counts: (i64, i64, i64, i64, i64, String) = sqlx::query_as( 2390 r#"SELECT 2391 (SELECT COUNT(*) FROM evidence_manifests), 2392 (SELECT COUNT(*) FROM trade_projections), 2393 (SELECT COUNT(*) FROM attestation_reports), 2394 (SELECT COUNT(*) FROM signed_attestation_events), 2395 (SELECT COUNT(*) FROM publication_outbox), 2396 (SELECT state FROM reconciliation_jobs WHERE job_id = ?)"#, 2397 ) 2398 .bind(lease.job().id().as_bytes().as_slice()) 2399 .fetch_one(&mut connection) 2400 .await 2401 .expect("no finalization mutation"); 2402 assert_eq!(counts, (0, 0, 0, 0, 1, "leased".into())); 2403 connection.close().await.expect("verification close"); 2404 host.close().await.expect("host close"); 2405 drop(( 2406 signed, 2407 identity, 2408 publication, 2409 configuration, 2410 metadata, 2411 runtime, 2412 root, 2413 )); 2414 } 2415 2416 #[tokio::test] 2417 async fn cancelled_submitted_attempt_recovers_unknown_and_retries_exact_bytes_after_reopen() { 2418 let ( 2419 root, 2420 runtime, 2421 metadata, 2422 configuration, 2423 host, 2424 lease, 2425 signed, 2426 publication, 2427 manifest, 2428 _trade, 2429 ) = signed_finalization_fixture("publication-execution-cancel", false).await; 2430 let repositories = host.repositories(); 2431 let finalized = repositories 2432 .reconciliation_attempts() 2433 .commit_finalization(&signed, &publication, now(1_784_347_208_000)) 2434 .await 2435 .expect("finalization"); 2436 let outbox_id = finalized.outbox_id().expect("outbox"); 2437 let claimed = repositories 2438 .publication_outbox() 2439 .claim_next_publication( 2440 publication_owner(0x92), 2441 publication_now(1_784_347_208_000), 2442 &publication, 2443 ) 2444 .await 2445 .expect("claim") 2446 .expect("due outbox"); 2447 let prepared = repositories 2448 .publication_outbox() 2449 .prepare_next_publication_target(claimed, publication_now(1_784_347_208_000)) 2450 .await 2451 .expect("durable submitted"); 2452 assert_eq!( 2453 prepared.exact_signed_event_bytes(), 2454 signed.signed_event_bytes() 2455 ); 2456 let pending = PendingExactPublicationSink.submit_exact(&prepared); 2457 drop(pending); 2458 drop(prepared); 2459 2460 let mut connection = fixture_connection(&runtime).await; 2461 let submitted: (String, String, i64, i64) = sqlx::query_as( 2462 r#"SELECT outbox.state, target.state, target.attempt_count, 2463 (SELECT COUNT(*) FROM publication_attempts) 2464 FROM publication_outbox AS outbox 2465 JOIN publication_targets AS target ON target.outbox_id = outbox.outbox_id 2466 WHERE outbox.outbox_id = ? AND target.target_ordinal = 0"#, 2467 ) 2468 .bind(outbox_id.as_bytes().as_slice()) 2469 .fetch_one(&mut connection) 2470 .await 2471 .expect("submitted state"); 2472 assert_eq!(submitted, ("leased".into(), "submitted".into(), 1, 0)); 2473 connection.close().await.expect("fixture close"); 2474 host.close().await.expect("host close"); 2475 2476 let reopened = open_writer(&runtime, &metadata).await; 2477 assert!( 2478 reopened 2479 .repositories() 2480 .publication_outbox() 2481 .recover_one_expired_publication( 2482 &publication_adapters(1_784_347_224), 2483 publication_now(1_784_347_224_000), 2484 &publication, 2485 ) 2486 .await 2487 .expect("expired recovery") 2488 ); 2489 let mut connection = fixture_connection(&runtime).await; 2490 let recovered: (String, String, i64, String, i64) = sqlx::query_as( 2491 r#"SELECT outbox.state, target.state, target.attempt_count, 2492 attempt.outcome, target.next_attempt_unix_ms 2493 FROM publication_outbox AS outbox 2494 JOIN publication_targets AS target ON target.outbox_id = outbox.outbox_id 2495 JOIN publication_attempts AS attempt ON attempt.attempt_id = target.last_attempt_id 2496 WHERE outbox.outbox_id = ? AND target.target_ordinal = 0"#, 2497 ) 2498 .bind(outbox_id.as_bytes().as_slice()) 2499 .fetch_one(&mut connection) 2500 .await 2501 .expect("recovered unknown"); 2502 assert_eq!(recovered.0, "pending"); 2503 assert_eq!(recovered.1, "unknown"); 2504 assert_eq!(recovered.2, 1); 2505 assert_eq!(recovered.3, "unknown"); 2506 assert!(recovered.4 >= 1_784_347_224_000); 2507 assert!(recovered.4 <= 1_784_347_224_250); 2508 connection.close().await.expect("fixture close"); 2509 2510 let expected = Arc::new(signed.signed_event_bytes().to_vec()); 2511 let calls = Arc::new(AtomicUsize::new(0)); 2512 let sink = RecordingExactPublicationSink { 2513 expected: Arc::clone(&expected), 2514 calls: Arc::clone(&calls), 2515 outcome: RhiPublicationAttemptOutcome::Accepted, 2516 }; 2517 let retried = reopened 2518 .repositories() 2519 .publication_outbox() 2520 .execute_next_publication( 2521 publication_owner(0x93), 2522 &publication_adapters(1_784_347_300), 2523 &sink, 2524 &publication, 2525 ) 2526 .await 2527 .expect("exact retry") 2528 .expect("retried target"); 2529 assert_eq!(calls.load(Ordering::SeqCst), 1); 2530 assert_eq!(retried.attempt_number(), 2); 2531 assert_eq!(retried.outbox_state(), RhiPublicationOutboxState::Complete); 2532 reopened.close().await.expect("reopened close"); 2533 drop(( 2534 manifest, 2535 publication, 2536 signed, 2537 lease, 2538 configuration, 2539 metadata, 2540 runtime, 2541 root, 2542 )); 2543 } 2544 2545 #[tokio::test] 2546 async fn publication_outcome_commit_is_idempotent_and_inspection_is_nonmutating() { 2547 let ( 2548 root, 2549 runtime, 2550 metadata, 2551 configuration, 2552 host, 2553 lease, 2554 signed, 2555 publication, 2556 manifest, 2557 _trade, 2558 ) = signed_finalization_fixture("publication-execution-reconcile", false).await; 2559 let repositories = host.repositories(); 2560 let finalized = repositories 2561 .reconciliation_attempts() 2562 .commit_finalization(&signed, &publication, now(1_784_347_208_000)) 2563 .await 2564 .expect("finalization"); 2565 let outbox_id = finalized.outbox_id().expect("outbox"); 2566 let claimed = repositories 2567 .publication_outbox() 2568 .claim_next_publication( 2569 publication_owner(0x94), 2570 publication_now(1_784_347_208_000), 2571 &publication, 2572 ) 2573 .await 2574 .expect("claim") 2575 .expect("due outbox"); 2576 let prepared = repositories 2577 .publication_outbox() 2578 .prepare_next_publication_target(claimed, publication_now(1_784_347_208_000)) 2579 .await 2580 .expect("prepare"); 2581 let delay = RhiPublicationRetryDelayMilliseconds::new(0).expect("zero delay"); 2582 let first = repositories 2583 .publication_outbox() 2584 .record_publication_outcome( 2585 &prepared, 2586 publication_now(1_784_347_208_001), 2587 RhiPublicationAttemptOutcome::Accepted, 2588 delay, 2589 ) 2590 .await 2591 .expect("first commit"); 2592 let retry = repositories 2593 .publication_outbox() 2594 .record_publication_outcome( 2595 &prepared, 2596 publication_now(1_784_347_208_001), 2597 RhiPublicationAttemptOutcome::Accepted, 2598 delay, 2599 ) 2600 .await 2601 .expect("reconciled commit"); 2602 assert_eq!(retry, first); 2603 assert_eq!(retry.outbox_id(), outbox_id); 2604 host.close().await.expect("host close"); 2605 2606 let inspection = open_rhi_state_inspection(&runtime, &metadata) 2607 .await 2608 .expect("inspection"); 2609 let error = inspection 2610 .repositories() 2611 .publication_outbox() 2612 .claim_next_publication( 2613 publication_owner(0x95), 2614 publication_now(1_784_347_300_000), 2615 &publication, 2616 ) 2617 .await 2618 .expect_err("inspection cannot claim"); 2619 assert_eq!( 2620 error.kind(), 2621 rhi::RhiPublicationExecutionErrorKind::InvalidMode 2622 ); 2623 inspection.close().await.expect("inspection close"); 2624 drop(( 2625 manifest, 2626 publication, 2627 signed, 2628 lease, 2629 configuration, 2630 metadata, 2631 runtime, 2632 root, 2633 )); 2634 } 2635 2636 #[tokio::test] 2637 async fn publication_execution_binds_live_authority_without_mutating_on_mismatch_or_disable() { 2638 let ( 2639 root, 2640 runtime, 2641 metadata, 2642 configuration, 2643 host, 2644 lease, 2645 signed, 2646 publication, 2647 manifest, 2648 _trade, 2649 ) = signed_finalization_fixture("publication-execution-authority", false).await; 2650 let repositories = host.repositories(); 2651 let finalized = repositories 2652 .reconciliation_attempts() 2653 .commit_finalization(&signed, &publication, now(1_784_347_208_000)) 2654 .await 2655 .expect("finalization"); 2656 let outbox_id = finalized.outbox_id().expect("outbox"); 2657 let mut connection = fixture_connection(&runtime).await; 2658 let before: (String, i64, String, i64, i64) = sqlx::query_as( 2659 r#"SELECT outbox.state, outbox.revision, target.state, 2660 target.revision, target.attempt_count 2661 FROM publication_outbox AS outbox 2662 JOIN publication_targets AS target ON target.outbox_id = outbox.outbox_id 2663 WHERE outbox.outbox_id = ? AND target.target_ordinal = 0"#, 2664 ) 2665 .bind(outbox_id.as_bytes().as_slice()) 2666 .fetch_one(&mut connection) 2667 .await 2668 .expect("before authority mismatch"); 2669 connection.close().await.expect("fixture close"); 2670 2671 let changed = EXAMPLE.replacen( 2672 "attempt_deadline_ms = 15000", 2673 "attempt_deadline_ms = 14999", 2674 1, 2675 ); 2676 let changed = parse_rhi_config_v1(changed.as_bytes(), rhi::RhiConfigProfile::RepoLocal) 2677 .expect("changed configuration"); 2678 let mismatched = RhiPublicationAuthority::from_config(&changed).expect("changed authority"); 2679 let error = repositories 2680 .publication_outbox() 2681 .claim_next_publication( 2682 publication_owner(0xa1), 2683 publication_now(1_784_347_208_000), 2684 &mismatched, 2685 ) 2686 .await 2687 .expect_err("mismatched authority"); 2688 assert_eq!( 2689 error.kind(), 2690 rhi::RhiPublicationExecutionErrorKind::Invariant 2691 ); 2692 2693 let publication_offset = EXAMPLE.find("[publication]").expect("publication section"); 2694 let presence_offset = EXAMPLE.find("[presence]").expect("presence section"); 2695 let disabled_source = format!( 2696 "{}[publication]\nmode = \"disabled\"\n\n{}", 2697 &EXAMPLE[..publication_offset], 2698 &EXAMPLE[presence_offset..] 2699 ); 2700 let disabled = 2701 parse_rhi_config_v1(disabled_source.as_bytes(), rhi::RhiConfigProfile::RepoLocal) 2702 .expect("disabled configuration"); 2703 let disabled = RhiPublicationAuthority::from_config(&disabled).expect("disabled authority"); 2704 assert!( 2705 repositories 2706 .publication_outbox() 2707 .claim_next_publication( 2708 publication_owner(0xa2), 2709 publication_now(1_784_347_208_000), 2710 &disabled, 2711 ) 2712 .await 2713 .expect("disabled authority") 2714 .is_none() 2715 ); 2716 2717 let mut connection = fixture_connection(&runtime).await; 2718 let after: (String, i64, String, i64, i64) = sqlx::query_as( 2719 r#"SELECT outbox.state, outbox.revision, target.state, 2720 target.revision, target.attempt_count 2721 FROM publication_outbox AS outbox 2722 JOIN publication_targets AS target ON target.outbox_id = outbox.outbox_id 2723 WHERE outbox.outbox_id = ? AND target.target_ordinal = 0"#, 2724 ) 2725 .bind(outbox_id.as_bytes().as_slice()) 2726 .fetch_one(&mut connection) 2727 .await 2728 .expect("after authority mismatch"); 2729 assert_eq!(after, before); 2730 connection.close().await.expect("fixture close"); 2731 host.close().await.expect("host close"); 2732 drop(( 2733 manifest, 2734 publication, 2735 signed, 2736 lease, 2737 configuration, 2738 metadata, 2739 runtime, 2740 root, 2741 )); 2742 } 2743 2744 #[tokio::test] 2745 async fn terminal_required_rejection_blocks_the_outbox_without_retry_schedule() { 2746 let ( 2747 root, 2748 runtime, 2749 metadata, 2750 configuration, 2751 host, 2752 lease, 2753 signed, 2754 publication, 2755 manifest, 2756 _trade, 2757 ) = signed_finalization_fixture("publication-execution-rejected", false).await; 2758 let repositories = host.repositories(); 2759 let finalized = repositories 2760 .reconciliation_attempts() 2761 .commit_finalization(&signed, &publication, now(1_784_347_208_000)) 2762 .await 2763 .expect("finalization"); 2764 let expected = Arc::new(signed.signed_event_bytes().to_vec()); 2765 let calls = Arc::new(AtomicUsize::new(0)); 2766 let sink = RecordingExactPublicationSink { 2767 expected, 2768 calls: Arc::clone(&calls), 2769 outcome: RhiPublicationAttemptOutcome::Rejected, 2770 }; 2771 let committed = repositories 2772 .publication_outbox() 2773 .execute_next_publication( 2774 publication_owner(0xa3), 2775 &publication_adapters(1_784_347_208), 2776 &sink, 2777 &publication, 2778 ) 2779 .await 2780 .expect("terminal rejection") 2781 .expect("due outbox"); 2782 assert_eq!(calls.load(Ordering::SeqCst), 1); 2783 assert_eq!(committed.outcome(), RhiPublicationAttemptOutcome::Rejected); 2784 assert_eq!( 2785 committed.target_state(), 2786 RhiPublicationTargetState::Rejected 2787 ); 2788 assert_eq!(committed.outbox_state(), RhiPublicationOutboxState::Blocked); 2789 2790 let mut connection = fixture_connection(&runtime).await; 2791 let durable: (String, Option<i64>, String, Option<i64>) = sqlx::query_as( 2792 r#"SELECT outbox.state, outbox.next_attempt_unix_ms, 2793 target.state, target.next_attempt_unix_ms 2794 FROM publication_outbox AS outbox 2795 JOIN publication_targets AS target ON target.outbox_id = outbox.outbox_id 2796 WHERE outbox.outbox_id = ? AND target.target_ordinal = 0"#, 2797 ) 2798 .bind(finalized.outbox_id().expect("outbox").as_bytes().as_slice()) 2799 .fetch_one(&mut connection) 2800 .await 2801 .expect("terminal durable state"); 2802 assert_eq!(durable, ("blocked".into(), None, "rejected".into(), None)); 2803 connection.close().await.expect("fixture close"); 2804 host.close().await.expect("host close"); 2805 drop(( 2806 manifest, 2807 publication, 2808 signed, 2809 lease, 2810 configuration, 2811 metadata, 2812 runtime, 2813 root, 2814 )); 2815 } 2816 2817 #[tokio::test] 2818 async fn atomic_finalization_disabled_mode_creates_no_publication_rows() { 2819 let ( 2820 root, 2821 runtime, 2822 metadata, 2823 configuration, 2824 host, 2825 _lease, 2826 signed, 2827 publication, 2828 manifest, 2829 _trade, 2830 ) = signed_finalization_fixture("finalization-commit-disabled", true).await; 2831 let outcome = host 2832 .repositories() 2833 .reconciliation_attempts() 2834 .commit_finalization(&signed, &publication, now(1_784_347_208_000)) 2835 .await 2836 .expect("disabled finalization"); 2837 assert!(outcome.created()); 2838 assert_eq!(outcome.publication_mode(), RhiPublicationMode::Disabled); 2839 assert_eq!(outcome.target_count(), 0); 2840 assert_eq!(outcome.outbox_id(), None); 2841 let mut connection = fixture_connection(&runtime).await; 2842 let counts: (i64, i64, i64) = sqlx::query_as( 2843 r#"SELECT 2844 (SELECT COUNT(*) FROM signed_attestation_events), 2845 (SELECT COUNT(*) FROM publication_outbox), 2846 (SELECT COUNT(*) FROM publication_targets)"#, 2847 ) 2848 .fetch_one(&mut connection) 2849 .await 2850 .expect("disabled counts"); 2851 assert_eq!(counts, (1, 0, 0)); 2852 connection.close().await.expect("fixture close"); 2853 host.close().await.expect("host close"); 2854 drop(( 2855 signed, 2856 publication, 2857 manifest, 2858 configuration, 2859 metadata, 2860 runtime, 2861 root, 2862 )); 2863 } 2864 2865 #[tokio::test] 2866 async fn atomic_finalization_rejects_configuration_and_generation_drift_without_partial_rows() { 2867 let ( 2868 root, 2869 runtime, 2870 metadata, 2871 configuration, 2872 host, 2873 _lease, 2874 signed, 2875 publication, 2876 manifest, 2877 trade, 2878 ) = signed_finalization_fixture("finalization-commit-drift", false).await; 2879 let changed = attestation_configuration( 2880 &runtime, 2881 &Keys::new(SecretKey::from_slice(&attestation_secret()).expect("identity secret")) 2882 .public_key() 2883 .to_hex(), 2884 ) 2885 .replace("samples = 512", "samples = 511"); 2886 let changed = parse_rhi_config_v1(changed.as_bytes(), rhi::RhiConfigProfile::RepoLocal) 2887 .expect("changed configuration"); 2888 let mismatched = RhiPublicationAuthority::from_config(&changed).expect("changed authority"); 2889 assert_eq!(mismatched, publication); 2890 let error = host 2891 .repositories() 2892 .reconciliation_attempts() 2893 .commit_finalization(&signed, &mismatched, now(1_784_347_208_000)) 2894 .await 2895 .expect_err("configuration mismatch"); 2896 assert_eq!( 2897 error.kind(), 2898 RhiReconciliationFinalizationCommitErrorKind::InvalidInput 2899 ); 2900 2901 write_dirty( 2902 &runtime, 2903 trade, 2904 3, 2905 *metadata.evidence_policy_digest().as_bytes(), 2906 1_784_347_208, 2907 ) 2908 .await; 2909 let error = host 2910 .repositories() 2911 .reconciliation_attempts() 2912 .commit_finalization(&signed, &publication, now(1_784_347_208_001)) 2913 .await 2914 .expect_err("generation drift"); 2915 assert_eq!( 2916 error.kind(), 2917 RhiReconciliationFinalizationCommitErrorKind::GenerationConflict 2918 ); 2919 let mut connection = fixture_connection(&runtime).await; 2920 let counts: (i64, i64, i64, i64, i64) = sqlx::query_as( 2921 r#"SELECT 2922 (SELECT COUNT(*) FROM evidence_manifests), 2923 (SELECT COUNT(*) FROM trade_projections), 2924 (SELECT COUNT(*) FROM attestation_reports), 2925 (SELECT COUNT(*) FROM signed_attestation_events), 2926 (SELECT COUNT(*) FROM publication_outbox)"#, 2927 ) 2928 .fetch_one(&mut connection) 2929 .await 2930 .expect("rolled-back inventory"); 2931 assert_eq!(counts, (0, 0, 0, 0, 0)); 2932 connection.close().await.expect("fixture close"); 2933 host.close().await.expect("host close"); 2934 drop(( 2935 signed, 2936 publication, 2937 manifest, 2938 configuration, 2939 metadata, 2940 runtime, 2941 root, 2942 )); 2943 } 2944 2945 async fn signed_finalization_fixture( 2946 instance: &str, 2947 publication_disabled: bool, 2948 ) -> ( 2949 tempfile::TempDir, 2950 RhiRuntimeContext, 2951 RhiStateMetadata, 2952 rhi::RhiConfigDocumentV1, 2953 rhi::RhiStateHost, 2954 RhiReconciliationLease, 2955 rhi::RhiSignedEvidenceAttestation, 2956 RhiPublicationAuthority, 2957 Box<[u8]>, 2958 TradeId, 2959 ) { 2960 let expected_public_key = 2961 Keys::new(SecretKey::from_slice(&attestation_secret()).expect("identity secret")) 2962 .public_key() 2963 .to_hex(); 2964 let (root, runtime, metadata, configuration, host, lease, evaluation) = 2965 finalization_fixture_with_source(instance, |runtime| { 2966 let source = attestation_configuration(runtime, &expected_public_key); 2967 if publication_disabled { 2968 let publication = source.find("[publication]").expect("publication section"); 2969 let presence = source.find("[presence]").expect("presence section"); 2970 format!( 2971 "{}[publication]\nmode = \"disabled\"\n\n{}", 2972 &source[..publication], 2973 &source[presence..] 2974 ) 2975 } else { 2976 source 2977 } 2978 }) 2979 .await; 2980 let manifest = evaluation.projection().manifest().canonical_bytes().into(); 2981 let trade = *evaluation.projection().manifest().trade_id(); 2982 let identity = provision_attestation_identity(&runtime, &configuration, &metadata); 2983 let fence = host 2984 .repositories() 2985 .reconciliation_attempts() 2986 .prepare_finalization(lease, evaluation, now(1_784_347_206_000)) 2987 .await 2988 .expect("finalization fence"); 2989 let signed = build_rhi_signed_evidence_attestation( 2990 fence, 2991 &identity, 2992 UnixTimeSeconds::new(1_784_347_207), 2993 &FixedAttestationEntropy(0xb1), 2994 None, 2995 ) 2996 .expect("signed attestation"); 2997 let publication = 2998 RhiPublicationAuthority::from_config(&configuration).expect("publication authority"); 2999 drop(identity); 3000 ( 3001 root, 3002 runtime, 3003 metadata, 3004 configuration, 3005 host, 3006 lease, 3007 signed, 3008 publication, 3009 manifest, 3010 trade, 3011 ) 3012 } 3013 3014 async fn fixture_connection(runtime: &RhiRuntimeContext) -> SqliteConnection { 3015 let options = SqliteConnectOptions::new() 3016 .filename(runtime.artifacts().state_database()) 3017 .create_if_missing(false) 3018 .foreign_keys(true); 3019 SqliteConnection::connect_with(&options) 3020 .await 3021 .expect("offline fixture connection") 3022 } 3023 3024 async fn finalization_snapshot( 3025 runtime: &RhiRuntimeContext, 3026 ) -> (i64, i64, i64, i64, i64, i64, i64, i64) { 3027 let mut connection = fixture_connection(runtime).await; 3028 let row: (i64, i64, i64, i64, i64, i64, i64, i64) = sqlx::query_as( 3029 r#"SELECT 3030 (SELECT COUNT(*) FROM reconciliation_jobs), 3031 (SELECT COUNT(*) FROM reconciliation_jobs WHERE state = 'leased'), 3032 (SELECT SUM(revision) FROM reconciliation_jobs), 3033 (SELECT COUNT(*) FROM evidence_reconciliations), 3034 (SELECT COUNT(*) FROM evidence_reconciliation_sources), 3035 (SELECT COUNT(*) FROM relay_checkpoints), 3036 (SELECT generation FROM trade_dirty_generations LIMIT 1), 3037 (SELECT COUNT(*) FROM sqlite_schema WHERE name IN ( 3038 'evidence_manifests', 'trade_projections', 'attestation_reports', 3039 'signed_attestation_events', 'publication_outbox' 3040 ))"#, 3041 ) 3042 .fetch_one(&mut connection) 3043 .await 3044 .expect("finalization snapshot"); 3045 connection.close().await.expect("snapshot close"); 3046 row 3047 } 3048 3049 async fn finalization_fixture( 3050 instance: &str, 3051 ) -> ( 3052 tempfile::TempDir, 3053 RhiRuntimeContext, 3054 RhiStateMetadata, 3055 rhi::RhiConfigDocumentV1, 3056 rhi::RhiStateHost, 3057 RhiReconciliationLease, 3058 rhi::RhiReconciliationEvaluation, 3059 ) { 3060 finalization_fixture_with_source(instance, |_| EXAMPLE.to_owned()).await 3061 } 3062 3063 async fn finalization_fixture_with_source<F>( 3064 instance: &str, 3065 source: F, 3066 ) -> ( 3067 tempfile::TempDir, 3068 RhiRuntimeContext, 3069 RhiStateMetadata, 3070 rhi::RhiConfigDocumentV1, 3071 rhi::RhiStateHost, 3072 RhiReconciliationLease, 3073 rhi::RhiReconciliationEvaluation, 3074 ) 3075 where 3076 F: FnOnce(&RhiRuntimeContext) -> String, 3077 { 3078 finalization_fixture_with_source_and_generation(instance, 1, source).await 3079 } 3080 3081 async fn finalization_fixture_with_source_and_generation<F>( 3082 instance: &str, 3083 initial_generation: u64, 3084 source: F, 3085 ) -> ( 3086 tempfile::TempDir, 3087 RhiRuntimeContext, 3088 RhiStateMetadata, 3089 rhi::RhiConfigDocumentV1, 3090 rhi::RhiStateHost, 3091 RhiReconciliationLease, 3092 rhi::RhiReconciliationEvaluation, 3093 ) 3094 where 3095 F: FnOnce(&RhiRuntimeContext) -> String, 3096 { 3097 let started_ms = 1_784_347_200_000; 3098 let (root, runtime, metadata, configuration, host, first_lease, first_plan) = 3099 replay_fixture_with_source_and_generation(instance, started_ms, initial_generation, source) 3100 .await; 3101 let wire = replay_wire(); 3102 let first_request = &first_plan.requests()[0]; 3103 let first_replay = RhiReconciliationSourceReplayPlan::from_request( 3104 &first_plan, 3105 first_request, 3106 &configuration, 3107 None, 3108 ) 3109 .expect("first replay plan") 3110 .finish( 3111 first_request, 3112 RhiTradeSourceCompletion::Complete, 3113 now(started_ms), 3114 now(started_ms + 2_000), 3115 [admitted_replay(&configuration, &wire, 1_784_347_200)], 3116 ) 3117 .expect("first replay"); 3118 let first = host 3119 .repositories() 3120 .reconciliation_attempts() 3121 .commit_source_replays(first_lease, first_plan, [first_replay]) 3122 .await 3123 .expect("first commit"); 3124 assert!(first.dirty_generation_advanced()); 3125 let cursor = first.committed_cursors()[0].clone(); 3126 drop(first); 3127 3128 let jobs = host.repositories().reconciliation_jobs(); 3129 let scheduled = jobs 3130 .schedule_trade( 3131 first_lease.job().trade_id(), 3132 configured_policy(&configuration), 3133 now(started_ms + 3_000), 3134 ) 3135 .await 3136 .expect("schedule final generation"); 3137 assert_eq!(scheduled.job().input_generation(), initial_generation + 1); 3138 let lease = jobs 3139 .claim_next(owner(0x92), now(started_ms + 3_000)) 3140 .await 3141 .expect("claim final generation") 3142 .expect("final job"); 3143 let plan = 3144 RhiReconciliationAttemptPlan::from_claim(lease, &configuration, now(started_ms + 3_000)) 3145 .expect("final plan"); 3146 let request = &plan.requests()[0]; 3147 let replay = RhiReconciliationSourceReplayPlan::from_request( 3148 &plan, 3149 request, 3150 &configuration, 3151 Some(cursor), 3152 ) 3153 .expect("final replay plan") 3154 .finish( 3155 request, 3156 RhiTradeSourceCompletion::Complete, 3157 now(started_ms + 3_000), 3158 now(started_ms + 5_000), 3159 [admitted_replay(&configuration, &wire, 1_784_347_203)], 3160 ) 3161 .expect("final replay"); 3162 let committed = host 3163 .repositories() 3164 .reconciliation_attempts() 3165 .commit_source_replays(lease, plan, [replay]) 3166 .await 3167 .expect("final commit"); 3168 assert!(!committed.dirty_generation_advanced()); 3169 let manifest = committed 3170 .into_evidence_manifest( 3171 UnixTimeSeconds::new(1_784_347_206), 3172 RhiReconciliationScopePrerequisites::Satisfied, 3173 ) 3174 .expect("final manifest"); 3175 assert_eq!(manifest.trade_generation(), initial_generation + 1); 3176 let projection = reduce_rhi_reconciliation_manifest(manifest).expect("final projection"); 3177 let claim = *projection.root_mutation_id().expect("root claim"); 3178 let evaluation = rhi::evaluate_rhi_reconciliation_claim(projection, claim); 3179 ( 3180 root, 3181 runtime, 3182 metadata, 3183 configuration, 3184 host, 3185 lease, 3186 evaluation, 3187 ) 3188 } 3189 3190 async fn attempt_fixture( 3191 instance: &str, 3192 ) -> ( 3193 tempfile::TempDir, 3194 RhiRuntimeContext, 3195 RhiStateMetadata, 3196 rhi::RhiConfigDocumentV1, 3197 rhi::RhiStateHost, 3198 RhiReconciliationAttemptPlan, 3199 ) { 3200 let root = tempfile::tempdir().expect("root"); 3201 let runtime = runtime(root.path(), instance); 3202 let configuration = parse_rhi_config_v1(EXAMPLE.as_bytes(), rhi::RhiConfigProfile::RepoLocal) 3203 .expect("configuration"); 3204 let metadata = metadata_from_config(&runtime, &configuration); 3205 initialize(&runtime, &metadata).await; 3206 let trade = TradeId::from_bytes([0x41; 16]); 3207 write_dirty( 3208 &runtime, 3209 trade, 3210 1, 3211 *metadata.evidence_policy_digest().as_bytes(), 3212 1_000, 3213 ) 3214 .await; 3215 let host = open_writer(&runtime, &metadata).await; 3216 let jobs = host.repositories().reconciliation_jobs(); 3217 jobs.schedule_trade(trade, configured_policy(&configuration), now(1_000)) 3218 .await 3219 .expect("schedule"); 3220 let lease = jobs 3221 .claim_next(owner(0x42), now(1_000)) 3222 .await 3223 .expect("claim") 3224 .expect("job"); 3225 let plan = 3226 RhiReconciliationAttemptPlan::from_claim(lease, &configuration, now(1_000)).expect("plan"); 3227 (root, runtime, metadata, configuration, host, plan) 3228 } 3229 3230 async fn replay_fixture( 3231 instance: &str, 3232 source: &str, 3233 started_ms: u64, 3234 ) -> ( 3235 tempfile::TempDir, 3236 RhiRuntimeContext, 3237 RhiStateMetadata, 3238 rhi::RhiConfigDocumentV1, 3239 rhi::RhiStateHost, 3240 RhiReconciliationLease, 3241 RhiReconciliationAttemptPlan, 3242 ) { 3243 replay_fixture_with_source(instance, started_ms, |_| source.to_owned()).await 3244 } 3245 3246 async fn replay_fixture_with_source<F>( 3247 instance: &str, 3248 started_ms: u64, 3249 source: F, 3250 ) -> ( 3251 tempfile::TempDir, 3252 RhiRuntimeContext, 3253 RhiStateMetadata, 3254 rhi::RhiConfigDocumentV1, 3255 rhi::RhiStateHost, 3256 RhiReconciliationLease, 3257 RhiReconciliationAttemptPlan, 3258 ) 3259 where 3260 F: FnOnce(&RhiRuntimeContext) -> String, 3261 { 3262 replay_fixture_with_source_and_generation(instance, started_ms, 1, source).await 3263 } 3264 3265 async fn replay_fixture_with_source_and_generation<F>( 3266 instance: &str, 3267 started_ms: u64, 3268 initial_generation: u64, 3269 source: F, 3270 ) -> ( 3271 tempfile::TempDir, 3272 RhiRuntimeContext, 3273 RhiStateMetadata, 3274 rhi::RhiConfigDocumentV1, 3275 rhi::RhiStateHost, 3276 RhiReconciliationLease, 3277 RhiReconciliationAttemptPlan, 3278 ) 3279 where 3280 F: FnOnce(&RhiRuntimeContext) -> String, 3281 { 3282 let root = tempfile::tempdir().expect("root"); 3283 let runtime = runtime(root.path(), instance); 3284 let source = source(&runtime); 3285 let configuration = parse_rhi_config_v1(source.as_bytes(), rhi::RhiConfigProfile::RepoLocal) 3286 .expect("configuration"); 3287 let metadata = metadata_from_config(&runtime, &configuration); 3288 initialize(&runtime, &metadata).await; 3289 let trade = TradeId::from_bytes([0x11; 16]); 3290 write_dirty( 3291 &runtime, 3292 trade, 3293 initial_generation, 3294 *metadata.evidence_policy_digest().as_bytes(), 3295 started_ms / 1_000, 3296 ) 3297 .await; 3298 let host = open_writer(&runtime, &metadata).await; 3299 let jobs = host.repositories().reconciliation_jobs(); 3300 jobs.schedule_trade(trade, configured_policy(&configuration), now(started_ms)) 3301 .await 3302 .expect("schedule"); 3303 let lease = jobs 3304 .claim_next(owner(0x72), now(started_ms)) 3305 .await 3306 .expect("claim") 3307 .expect("job"); 3308 let plan = RhiReconciliationAttemptPlan::from_claim(lease, &configuration, now(started_ms)) 3309 .expect("plan"); 3310 (root, runtime, metadata, configuration, host, lease, plan) 3311 } 3312 3313 fn replay_wire() -> Vec<u8> { 3314 serde_json::from_str::<serde_json::Value>(TRADE_VECTOR).expect("trade vector")["raw_json"] 3315 .as_str() 3316 .expect("raw event") 3317 .as_bytes() 3318 .to_vec() 3319 } 3320 3321 fn admitted_replay( 3322 configuration: &rhi::RhiConfigDocumentV1, 3323 wire: &[u8], 3324 observed_at: u64, 3325 ) -> rhi::RhiAdmittedTradeMutationEvent { 3326 admit_rhi_trade_mutation_event( 3327 RhiTradeMutationAdmissionLimits::from_config(configuration).expect("admission limits"), 3328 wire, 3329 RhiTradeMutationObservedAtUnixSeconds::new(observed_at).expect("observed time"), 3330 RhiTradeMutationAuthoredTimePolicy::new(0).expect("authored-time policy"), 3331 ) 3332 .expect("admitted event") 3333 } 3334 3335 #[test] 3336 fn public_inputs_have_exact_bounds_and_diagnostics_are_redacted() { 3337 let maximum = RhiReconciliationJobPolicy::new(65_536, 300_000, 150_000, 100, 60_000, 3_600_000) 3338 .expect("exact maximum policy"); 3339 assert_eq!(maximum.queue_capacity(), 65_536); 3340 assert_eq!(maximum.lease_duration_milliseconds(), 300_000); 3341 assert_eq!(maximum.lease_renewal_milliseconds(), 150_000); 3342 assert_eq!(maximum.max_attempts(), 100); 3343 assert_eq!(maximum.initial_backoff_milliseconds(), 60_000); 3344 assert_eq!(maximum.maximum_backoff_milliseconds(), 3_600_000); 3345 let configuration = parse_rhi_config_v1(EXAMPLE.as_bytes(), rhi::RhiConfigProfile::RepoLocal) 3346 .expect("configuration"); 3347 let configured = RhiReconciliationJobPolicy::from_configuration(&configuration) 3348 .expect("configured reconciliation policy"); 3349 assert_eq!(configured.queue_capacity(), 4_096); 3350 assert_eq!(configured.lease_duration_milliseconds(), 30_000); 3351 assert_eq!(configured.lease_renewal_milliseconds(), 10_000); 3352 assert_eq!(configured.max_attempts(), 10); 3353 assert_eq!(configured.initial_backoff_milliseconds(), 250); 3354 assert_eq!(configured.maximum_backoff_milliseconds(), 30_000); 3355 assert!(RhiReconciliationJobPolicy::new(0, 1_000, 100, 1, 1, 1).is_err()); 3356 assert!(RhiReconciliationJobPolicy::new(65_537, 1_000, 100, 1, 1, 1).is_err()); 3357 assert!(RhiReconciliationJobPolicy::new(1, 999, 100, 1, 1, 1).is_err()); 3358 assert!(RhiReconciliationJobPolicy::new(1, 300_001, 100, 1, 1, 1).is_err()); 3359 assert!(RhiReconciliationJobPolicy::new(1, 1_000, 1_000, 1, 1, 1).is_err()); 3360 assert!(RhiReconciliationJobPolicy::new(1, 1_000, 100, 0, 1, 1).is_err()); 3361 assert!(RhiReconciliationJobPolicy::new(1, 1_000, 100, 101, 1, 1).is_err()); 3362 assert!(RhiReconciliationJobPolicy::new(1, 1_000, 100, 1, 0, 1).is_err()); 3363 assert!(RhiReconciliationJobPolicy::new(1, 1_000, 100, 1, 2, 1).is_err()); 3364 assert!(RhiReconciliationUnixMilliseconds::new(i64::MAX as u64).is_ok()); 3365 assert!(RhiReconciliationUnixMilliseconds::new(i64::MAX as u64 + 1).is_err()); 3366 assert!(RhiReconciliationRetryDelayMilliseconds::new(3_600_000).is_ok()); 3367 assert!(RhiReconciliationRetryDelayMilliseconds::new(3_600_001).is_err()); 3368 assert!(RhiReconciliationLeaseOwner::from_bytes([0; 16]).is_err()); 3369 let secret = RhiReconciliationLeaseOwner::from_bytes([0xab; 16]).expect("owner"); 3370 assert_eq!( 3371 format!("{secret:?}"), 3372 "RhiReconciliationLeaseOwner([redacted])" 3373 ); 3374 let error = 3375 RhiReconciliationJobPolicy::new(0, 1_000, 100, 1, 1, 1).expect_err("invalid policy"); 3376 assert!(Error::source(&error).is_none()); 3377 assert_eq!(error.code(), "reconciliation_job_input_invalid"); 3378 assert!(!format!("{error} {error:?}").contains("65536")); 3379 } 3380 3381 fn lower_hex(bytes: &[u8]) -> String { 3382 const DIGITS: &[u8; 16] = b"0123456789abcdef"; 3383 let mut output = String::with_capacity(bytes.len() * 2); 3384 for byte in bytes { 3385 output.push(char::from(DIGITS[usize::from(byte >> 4)])); 3386 output.push(char::from(DIGITS[usize::from(byte & 0x0f)])); 3387 } 3388 output 3389 }