rhi

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

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 }