rhi

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

services_hardening_source_ingest.rs (32373B)


      1 #![forbid(unsafe_code)]
      2 #![cfg(any(target_os = "linux", target_os = "macos"))]
      3 
      4 use std::{
      5     collections::VecDeque,
      6     fs,
      7     os::unix::fs::PermissionsExt,
      8     path::Path,
      9     sync::{Arc, Mutex},
     10 };
     11 
     12 use nostr::secp256k1::{Keypair, Message};
     13 use nostr::{EventId, Keys, SECP256K1};
     14 use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity};
     15 use radroots_storage::event::SourceGeneration;
     16 use radroots_transport::{
     17     BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink, EventSource, EventSubscriber,
     18     EventSubscription, FetchPage, FetchRequest, SinkFailure, SinkStatus, SourceStatus,
     19     SubscriptionRequest,
     20     outcome::{FetchTargetOutcome, FetchTargetState},
     21     source::{EventProvenance, FetchCursor, NextPage, ObservedEvent},
     22 };
     23 use rhi::{
     24     RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, RhiConfigProfile,
     25     RhiStateMetadata, RhiTradeMutationAuthoredTimePolicy, RhiTradeMutationObservedAtUnixSeconds,
     26     RhiTradeSourceAttempt, RhiTradeSourceCompletion, RhiTradeSourceIngestErrorKind,
     27     RhiTransportAdapters, TradeId, UnixTimeSeconds, apply_rhi_configuration,
     28     ingest_rhi_trade_source, initialize_rhi_state, open_rhi_state_read_write,
     29     open_rhi_state_read_write_from_config, parse_rhi_cli_v1_from, parse_rhi_config_v1,
     30     resolve_rhi_runtime_context,
     31 };
     32 use serde_json::Value;
     33 use sqlx::{ConnectOptions, Connection, Row, SqliteConnection, sqlite::SqliteConnectOptions};
     34 use tokio::sync::Barrier;
     35 
     36 const CONFIG: &str = include_str!("../contracts/services_hardening/config.v1.example.toml");
     37 const VECTOR: &str = include_str!("../contracts/conformance/vectors/trade_ingest_proposal.v1.json");
     38 
     39 #[derive(Clone)]
     40 struct PageSpec {
     41     events: Vec<String>,
     42     state: FetchTargetState,
     43     next: Option<&'static str>,
     44 }
     45 
     46 #[derive(Clone)]
     47 struct ScriptedTransport {
     48     pages: Arc<Mutex<VecDeque<PageSpec>>>,
     49     requests: Arc<Mutex<Vec<FetchRequest>>>,
     50     fetch_barrier: Option<Arc<Barrier>>,
     51 }
     52 
     53 impl ScriptedTransport {
     54     fn new(pages: impl IntoIterator<Item = PageSpec>) -> Self {
     55         Self {
     56             pages: Arc::new(Mutex::new(pages.into_iter().collect())),
     57             requests: Arc::new(Mutex::new(Vec::new())),
     58             fetch_barrier: None,
     59         }
     60     }
     61 
     62     fn with_fetch_barrier(mut self, parties: usize) -> Self {
     63         self.fetch_barrier = Some(Arc::new(Barrier::new(parties)));
     64         self
     65     }
     66 
     67     fn requests(&self) -> Vec<FetchRequest> {
     68         self.requests.lock().expect("requests").clone()
     69     }
     70 }
     71 
     72 impl EventSource for ScriptedTransport {
     73     fn status(&self) -> BoxFuture<'_, Result<SourceStatus, radroots_transport::Error>> {
     74         Box::pin(async { Err(radroots_transport::Error::UnsupportedOperation) })
     75     }
     76 
     77     fn fetch(
     78         &self,
     79         request: FetchRequest,
     80     ) -> BoxFuture<'_, Result<FetchPage, radroots_transport::Error>> {
     81         self.requests
     82             .lock()
     83             .expect("requests")
     84             .push(request.clone());
     85         let page = self.pages.lock().expect("pages").pop_front();
     86         let fetch_barrier = self.fetch_barrier.clone();
     87         Box::pin(async move {
     88             if let Some(barrier) = fetch_barrier {
     89                 barrier.wait().await;
     90             }
     91             let spec = page.ok_or(radroots_transport::Error::UnsupportedOperation)?;
     92             let target = request.target_set().targets().first().expect("one target");
     93             let events = spec
     94                 .events
     95                 .into_iter()
     96                 .map(|raw| {
     97                     let event = radroots_event_codec::decode::signed_event(raw.as_str())
     98                         .expect("signed event");
     99                     let mut provenance = EventProvenance::new(
    100                         radroots_transport::TransportId::NOSTR,
    101                         target.fingerprint().clone(),
    102                         1,
    103                     )
    104                     .expect("provenance");
    105                     if let Some(cursor) = request.cursor().cloned() {
    106                         provenance = provenance.with_cursor(cursor);
    107                     }
    108                     ObservedEvent::new(event, provenance)
    109                 })
    110                 .collect();
    111             let outcome = FetchTargetOutcome::new(target.fingerprint().clone(), spec.state);
    112             let next = match spec.next {
    113                 Some(cursor) => NextPage::Cursor(FetchCursor::parse(cursor).expect("cursor")),
    114                 None => NextPage::Complete,
    115             };
    116             FetchPage::for_request(&request, events, vec![outcome], next)
    117         })
    118     }
    119 }
    120 
    121 impl EventSubscriber for ScriptedTransport {
    122     fn subscribe(
    123         &self,
    124         _request: SubscriptionRequest,
    125     ) -> BoxFuture<'_, Result<Box<dyn EventSubscription>, radroots_transport::Error>> {
    126         Box::pin(async { Err(radroots_transport::Error::UnsupportedOperation) })
    127     }
    128 }
    129 
    130 impl EventSink for ScriptedTransport {
    131     fn status(&self) -> BoxFuture<'_, Result<SinkStatus, radroots_transport::Error>> {
    132         Box::pin(async { Err(radroots_transport::Error::UnsupportedOperation) })
    133     }
    134 
    135     fn deliver(
    136         &self,
    137         request: DeliveryRequest,
    138     ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
    139         Box::pin(async move { Err(SinkFailure::invalid_contract(&request)) })
    140     }
    141 }
    142 
    143 fn adapters(source: &ScriptedTransport) -> RhiTransportAdapters {
    144     RhiTransportAdapters::new(
    145         Arc::new(source.clone()),
    146         Arc::new(source.clone()),
    147         Arc::new(source.clone()),
    148     )
    149 }
    150 
    151 fn runtime(root: &Path) -> rhi::RhiRuntimeContext {
    152     let invocation = parse_rhi_cli_v1_from([
    153         "rhi",
    154         "--profile",
    155         "repo-local",
    156         "--instance",
    157         "primary",
    158         "--repo-local-root",
    159         root.to_str().expect("UTF-8 root"),
    160         "run",
    161     ])
    162     .expect("invocation");
    163     resolve_rhi_runtime_context(
    164         &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()),
    165         &invocation,
    166     )
    167     .expect("runtime")
    168 }
    169 
    170 fn configuration() -> rhi::RhiConfigDocumentV1 {
    171     parse_rhi_config_v1(CONFIG.as_bytes(), RhiConfigProfile::RepoLocal).expect("configuration")
    172 }
    173 
    174 fn metadata(
    175     runtime: &rhi::RhiRuntimeContext,
    176     configuration: &rhi::RhiConfigDocumentV1,
    177 ) -> RhiStateMetadata {
    178     RhiStateMetadata::new(
    179         runtime,
    180         configuration,
    181         SourceGeneration::new([0x5a; 32]).expect("generation"),
    182         1_725_000_000_000,
    183     )
    184     .expect("metadata")
    185 }
    186 
    187 fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) {
    188     (
    189         MigrationAppliedAtUnixSeconds::new(1_725_000_000).expect("migration time"),
    190         MigrationBuildIdentity::new(
    191             env!("CARGO_PKG_VERSION"),
    192             "1111111111111111111111111111111111111111",
    193             "053d0c750bf9cd683c6ea37cefe7e79617ba629f",
    194             "rustc-test",
    195             "test-target",
    196             "service-host",
    197             1,
    198             rhi::RHI_STATE_SCHEMA_VERSION,
    199             1,
    200             1,
    201             1,
    202         )
    203         .expect("build identity"),
    204     )
    205 }
    206 
    207 fn prepare_state_directory(runtime: &rhi::RhiRuntimeContext) {
    208     fs::create_dir_all(runtime.context().paths().state()).expect("state directory");
    209     fs::set_permissions(
    210         runtime.context().paths().state(),
    211         fs::Permissions::from_mode(0o700),
    212     )
    213     .expect("state mode");
    214 }
    215 
    216 fn vector_wire() -> String {
    217     let vector: Value = serde_json::from_str(VECTOR).expect("vector");
    218     vector["raw_json"].as_str().expect("raw event").to_owned()
    219 }
    220 
    221 fn resign_wire(raw: &str, auxiliary: u8) -> String {
    222     let mut event: Value = serde_json::from_str(raw).expect("event JSON");
    223     let keys = Keys::parse("10c5304d6c9ae3a1a16f7860f1cc8f5e3a76225a2663b3a989a0d775919b7df5")
    224         .expect("approved fixture keys");
    225     let event_id = EventId::from_hex(event["id"].as_str().expect("event id")).expect("event id");
    226     let message = Message::from_digest(event_id.to_bytes());
    227     let keypair = Keypair::from_secret_key(SECP256K1, keys.secret_key());
    228     event["sig"] = SECP256K1
    229         .sign_schnorr_with_aux_rand(&message, &keypair, &[auxiliary; 32])
    230         .to_string()
    231         .into();
    232     serde_json::to_string(&event).expect("event JSON")
    233 }
    234 
    235 fn attempt(request_id: &str, started_at: u64, observed_at: u64) -> RhiTradeSourceAttempt {
    236     RhiTradeSourceAttempt::new(
    237         request_id,
    238         UnixTimeSeconds::new(started_at),
    239         RhiTradeMutationObservedAtUnixSeconds::new(observed_at).expect("observed"),
    240         RhiTradeMutationAuthoredTimePolicy::new(0).expect("authored policy"),
    241     )
    242     .expect("attempt")
    243 }
    244 
    245 async fn offline_scalar(runtime: &rhi::RhiRuntimeContext, query: &'static str) -> i64 {
    246     let options = SqliteConnectOptions::new()
    247         .filename(runtime.artifacts().state_database())
    248         .create_if_missing(false)
    249         .disable_statement_logging();
    250     let mut connection = SqliteConnection::connect_with(&options)
    251         .await
    252         .expect("offline connection");
    253     let value = sqlx::query(query)
    254         .fetch_one(&mut connection)
    255         .await
    256         .expect("offline scalar")
    257         .try_get::<i64, _>(0)
    258         .expect("scalar value");
    259     connection.close().await.expect("offline close");
    260     value
    261 }
    262 
    263 async fn offline_signed_event_keys(runtime: &rhi::RhiRuntimeContext) -> Vec<(Vec<u8>, Vec<u8>)> {
    264     let options = SqliteConnectOptions::new()
    265         .filename(runtime.artifacts().state_database())
    266         .create_if_missing(false)
    267         .disable_statement_logging();
    268     let mut connection = SqliteConnection::connect_with(&options)
    269         .await
    270         .expect("offline connection");
    271     let rows = sqlx::query(
    272         "SELECT event_id, event_signature FROM nostr_events
    273          ORDER BY event_id, event_signature",
    274     )
    275     .fetch_all(&mut connection)
    276     .await
    277     .expect("signed-event inventory");
    278     let values = rows
    279         .into_iter()
    280         .map(|row| {
    281             (
    282                 row.try_get::<Vec<u8>, _>("event_id").expect("event id"),
    283                 row.try_get::<Vec<u8>, _>("event_signature")
    284                     .expect("event signature"),
    285             )
    286         })
    287         .collect();
    288     connection.close().await.expect("offline close");
    289     values
    290 }
    291 
    292 #[tokio::test]
    293 async fn exact_selector_fetch_replay_and_rejection_preserve_cursor_and_generation() {
    294     let root = tempfile::tempdir().expect("root");
    295     let runtime = runtime(root.path());
    296     prepare_state_directory(&runtime);
    297     let configuration = configuration();
    298     let metadata = metadata(&runtime, &configuration);
    299     let (applied_at, build) = migration_evidence();
    300     initialize_rhi_state(&runtime, &metadata, applied_at, &build)
    301         .await
    302         .expect("initialize");
    303     let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
    304         .await
    305         .expect("writer");
    306     let trade_id = TradeId::parse("11111111111111111111111111111111").expect("trade id");
    307     let wire = vector_wire();
    308 
    309     let source = ScriptedTransport::new([PageSpec {
    310         events: vec![wire.clone()],
    311         state: FetchTargetState::Complete,
    312         next: None,
    313     }]);
    314     let outcome = ingest_rhi_trade_source(
    315         &host.repositories(),
    316         &adapters(&source),
    317         &configuration,
    318         "trade-primary",
    319         trade_id,
    320         attempt("attempt-1", 1_784_347_200, 1_784_347_200),
    321     )
    322     .await
    323     .expect("ingest");
    324     assert_eq!(outcome.completion(), RhiTradeSourceCompletion::Complete);
    325     assert_eq!(outcome.received_events(), 1);
    326     assert_eq!(outcome.admitted_events(), 1);
    327     assert_eq!(outcome.rejected_events(), 0);
    328     assert_eq!(outcome.inserted_mutations(), 1);
    329     assert_eq!(outcome.inserted_signed_events(), 1);
    330     assert_eq!(outcome.inserted_observations(), 1);
    331     assert!(outcome.checkpoint_advanced());
    332     assert_eq!(outcome.dirty_generation().expect("generation").get(), 1);
    333     assert!(outcome.dirty_generation_advanced());
    334 
    335     let requests = source.requests();
    336     assert_eq!(requests.len(), 1);
    337     assert_eq!(
    338         requests[0].selector().kinds(),
    339         &[3470, 3471, 3472, 3473, 3474]
    340     );
    341     assert_eq!(
    342         requests[0]
    343             .selector()
    344             .exact_tag_filters()
    345             .collect::<Vec<_>>(),
    346         vec![('d', &["11111111111111111111111111111111".to_owned()][..])]
    347     );
    348     assert_eq!(
    349         requests[0].selector().since_unix_seconds(),
    350         Some(1_784_260_800)
    351     );
    352     assert_eq!(requests[0].bounds().deadline_unix_ms(), 1_784_347_210_000);
    353 
    354     let replay_source = ScriptedTransport::new([PageSpec {
    355         events: vec![wire.clone(), wire.clone()],
    356         state: FetchTargetState::Complete,
    357         next: None,
    358     }]);
    359     let replay = ingest_rhi_trade_source(
    360         &host.repositories(),
    361         &adapters(&replay_source),
    362         &configuration,
    363         "trade-primary",
    364         trade_id,
    365         attempt("attempt-2", 1_784_347_201, 1_784_347_201),
    366     )
    367     .await
    368     .expect("replay");
    369     assert_eq!(replay.admitted_events(), 1);
    370     assert_eq!(replay.duplicate_events(), 1);
    371     assert_eq!(replay.inserted_mutations(), 0);
    372     assert_eq!(replay.inserted_signed_events(), 0);
    373     assert_eq!(replay.inserted_observations(), 1);
    374     assert!(!replay.checkpoint_advanced());
    375     assert_eq!(replay.dirty_generation().expect("generation").get(), 1);
    376     assert!(!replay.dirty_generation_advanced());
    377 
    378     let rejected_source = ScriptedTransport::new([PageSpec {
    379         events: vec![resign_wire(&wire, 9)],
    380         state: FetchTargetState::Complete,
    381         next: None,
    382     }]);
    383     let rejected = ingest_rhi_trade_source(
    384         &host.repositories(),
    385         &adapters(&rejected_source),
    386         &configuration,
    387         "trade-primary",
    388         trade_id,
    389         attempt("attempt-3", 1_784_347_000, 1_784_347_199),
    390     )
    391     .await
    392     .expect("rejected attempt");
    393     assert_eq!(rejected.completion(), RhiTradeSourceCompletion::Complete);
    394     assert_eq!(rejected.rejected_events(), 1);
    395     assert_eq!(rejected.inserted_signed_events(), 0);
    396     assert!(!rejected.checkpoint_advanced());
    397     assert_eq!(rejected.dirty_generation().expect("generation").get(), 1);
    398     assert!(!rejected.dirty_generation_advanced());
    399 
    400     host.close().await.expect("close");
    401 }
    402 
    403 #[tokio::test]
    404 async fn incomplete_and_unsupported_results_never_advance_checkpoint_or_dirty_generation() {
    405     let root = tempfile::tempdir().expect("root");
    406     let runtime = runtime(root.path());
    407     prepare_state_directory(&runtime);
    408     let configuration = configuration();
    409     let metadata = metadata(&runtime, &configuration);
    410     let (applied_at, build) = migration_evidence();
    411     initialize_rhi_state(&runtime, &metadata, applied_at, &build)
    412         .await
    413         .expect("initialize");
    414     let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
    415         .await
    416         .expect("writer");
    417     let trade_id = TradeId::parse("11111111111111111111111111111111").expect("trade id");
    418     let vectors = [
    419         (
    420             FetchTargetState::Cancelled,
    421             RhiTradeSourceCompletion::IncompleteTimeout,
    422         ),
    423         (
    424             FetchTargetState::Unavailable,
    425             RhiTradeSourceCompletion::IncompleteUnavailable,
    426         ),
    427         (
    428             FetchTargetState::FailedRetryable,
    429             RhiTradeSourceCompletion::IncompleteUnavailable,
    430         ),
    431         (
    432             FetchTargetState::FailedTerminal,
    433             RhiTradeSourceCompletion::IncompleteUnknown,
    434         ),
    435         (
    436             FetchTargetState::Partial,
    437             RhiTradeSourceCompletion::IncompleteUnknown,
    438         ),
    439     ];
    440     for (index, (state, expected)) in vectors.into_iter().enumerate() {
    441         let source = ScriptedTransport::new([PageSpec {
    442             events: Vec::new(),
    443             state,
    444             next: None,
    445         }]);
    446         let outcome = ingest_rhi_trade_source(
    447             &host.repositories(),
    448             &adapters(&source),
    449             &configuration,
    450             "trade-primary",
    451             trade_id,
    452             attempt(
    453                 &format!("incomplete-{index}"),
    454                 1_784_347_300 + index as u64,
    455                 1_784_347_300 + index as u64,
    456             ),
    457         )
    458         .await
    459         .expect("classified incomplete result");
    460         assert_eq!(outcome.completion(), expected);
    461         assert_eq!(outcome.checkpoint(), None);
    462         assert_eq!(outcome.dirty_generation(), None);
    463         assert!(!outcome.checkpoint_advanced());
    464         assert!(!outcome.dirty_generation_advanced());
    465     }
    466 
    467     let unsupported = ScriptedTransport::new([]);
    468     let outcome = ingest_rhi_trade_source(
    469         &host.repositories(),
    470         &adapters(&unsupported),
    471         &configuration,
    472         "trade-primary",
    473         trade_id,
    474         attempt("unsupported", 1_784_347_400, 1_784_347_400),
    475     )
    476     .await
    477     .expect("unsupported result");
    478     assert_eq!(outcome.completion(), RhiTradeSourceCompletion::Unsupported);
    479     assert_eq!(outcome.checkpoint(), None);
    480     assert_eq!(outcome.dirty_generation(), None);
    481 
    482     host.close().await.expect("close");
    483 }
    484 
    485 #[tokio::test]
    486 async fn configured_result_bound_retains_admitted_evidence_without_checkpoint() {
    487     let root = tempfile::tempdir().expect("root");
    488     let runtime = runtime(root.path());
    489     prepare_state_directory(&runtime);
    490     let configuration = configuration();
    491     let metadata = metadata(&runtime, &configuration);
    492     let (applied_at, build) = migration_evidence();
    493     initialize_rhi_state(&runtime, &metadata, applied_at, &build)
    494         .await
    495         .expect("initialize");
    496     let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
    497         .await
    498         .expect("writer");
    499     let trade_id = TradeId::parse("11111111111111111111111111111111").expect("trade id");
    500     let wire = vector_wire();
    501     let source = ScriptedTransport::new([
    502         PageSpec {
    503             events: vec![wire.clone(); 1_000],
    504             state: FetchTargetState::Complete,
    505             next: Some("page-1"),
    506         },
    507         PageSpec {
    508             events: vec![wire.clone(); 1_000],
    509             state: FetchTargetState::Complete,
    510             next: Some("page-2"),
    511         },
    512         PageSpec {
    513             events: vec![wire.clone(); 1_000],
    514             state: FetchTargetState::Complete,
    515             next: Some("page-3"),
    516         },
    517         PageSpec {
    518             events: vec![wire.clone(); 1_000],
    519             state: FetchTargetState::Complete,
    520             next: Some("page-4"),
    521         },
    522         PageSpec {
    523             events: vec![wire; 97],
    524             state: FetchTargetState::Complete,
    525             next: None,
    526         },
    527     ]);
    528     let outcome = ingest_rhi_trade_source(
    529         &host.repositories(),
    530         &adapters(&source),
    531         &configuration,
    532         "trade-primary",
    533         trade_id,
    534         attempt("resource-limit", 1_784_347_500, 1_784_347_500),
    535     )
    536     .await
    537     .expect("resource-limited result");
    538     assert_eq!(
    539         outcome.completion(),
    540         RhiTradeSourceCompletion::IncompleteResourceLimit
    541     );
    542     assert_eq!(outcome.received_events(), 4_096);
    543     assert_eq!(outcome.admitted_events(), 1);
    544     assert_eq!(outcome.inserted_mutations(), 1);
    545     assert_eq!(outcome.inserted_signed_events(), 1);
    546     assert_eq!(outcome.inserted_observations(), 1);
    547     assert_eq!(outcome.checkpoint(), None);
    548     assert!(!outcome.checkpoint_advanced());
    549     assert_eq!(outcome.dirty_generation().expect("dirty").get(), 1);
    550     assert!(outcome.dirty_generation_advanced());
    551 
    552     host.close().await.expect("close");
    553 }
    554 
    555 #[tokio::test]
    556 async fn signed_event_identity_permutations_preserve_both_signatures_and_canonical_state() {
    557     let wire = vector_wire();
    558     let first = resign_wire(&wire, 1);
    559     let second = resign_wire(&wire, 2);
    560     let first_json: Value = serde_json::from_str(&first).expect("first event");
    561     let second_json: Value = serde_json::from_str(&second).expect("second event");
    562     assert_eq!(first_json["id"], second_json["id"]);
    563     assert_ne!(first_json["sig"], second_json["sig"]);
    564 
    565     let mut inventories = Vec::new();
    566     for (index, events) in [
    567         vec![first.clone(), second.clone()],
    568         vec![second.clone(), first.clone()],
    569     ]
    570     .into_iter()
    571     .enumerate()
    572     {
    573         let root = tempfile::tempdir().expect("root");
    574         let runtime = runtime(root.path());
    575         prepare_state_directory(&runtime);
    576         let configuration = configuration();
    577         let metadata = metadata(&runtime, &configuration);
    578         let (applied_at, build) = migration_evidence();
    579         initialize_rhi_state(&runtime, &metadata, applied_at, &build)
    580             .await
    581             .expect("initialize");
    582         let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
    583             .await
    584             .expect("writer");
    585         let trade_id = TradeId::parse("11111111111111111111111111111111").expect("trade id");
    586         let source = ScriptedTransport::new([PageSpec {
    587             events,
    588             state: FetchTargetState::Complete,
    589             next: None,
    590         }]);
    591         let outcome = ingest_rhi_trade_source(
    592             &host.repositories(),
    593             &adapters(&source),
    594             &configuration,
    595             "trade-primary",
    596             trade_id,
    597             attempt(
    598                 &format!("signature-permutation-{index}"),
    599                 1_784_347_600 + index as u64,
    600                 1_784_347_600 + index as u64,
    601             ),
    602         )
    603         .await
    604         .expect("permutation ingest");
    605         assert_eq!(outcome.completion(), RhiTradeSourceCompletion::Complete);
    606         assert_eq!(outcome.admitted_events(), 2);
    607         assert_eq!(outcome.duplicate_events(), 0);
    608         assert_eq!(outcome.inserted_mutations(), 1);
    609         assert_eq!(outcome.inserted_signed_events(), 2);
    610         assert_eq!(outcome.inserted_observations(), 2);
    611         assert_eq!(outcome.dirty_generation().expect("dirty").get(), 1);
    612         assert!(outcome.dirty_generation_advanced());
    613         host.close().await.expect("close");
    614 
    615         assert_eq!(
    616             offline_scalar(&runtime, "SELECT COUNT(*) FROM trade_mutations").await,
    617             1
    618         );
    619         assert_eq!(
    620             offline_scalar(&runtime, "SELECT COUNT(*) FROM nostr_events").await,
    621             2
    622         );
    623         assert_eq!(
    624             offline_scalar(&runtime, "SELECT COUNT(*) FROM relay_observations").await,
    625             2
    626         );
    627         assert_eq!(
    628             offline_scalar(&runtime, "SELECT COUNT(*) FROM relay_checkpoints").await,
    629             1
    630         );
    631         inventories.push(offline_signed_event_keys(&runtime).await);
    632     }
    633     assert_eq!(inventories[0], inventories[1]);
    634 }
    635 
    636 #[tokio::test]
    637 async fn checkpoint_and_dirty_generation_survive_close_reopen_and_exact_replay() {
    638     let root = tempfile::tempdir().expect("root");
    639     let runtime = runtime(root.path());
    640     prepare_state_directory(&runtime);
    641     let configuration = configuration();
    642     let metadata = metadata(&runtime, &configuration);
    643     let (applied_at, build) = migration_evidence();
    644     initialize_rhi_state(&runtime, &metadata, applied_at, &build)
    645         .await
    646         .expect("initialize");
    647     let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
    648         .await
    649         .expect("writer");
    650     let trade_id = TradeId::parse("11111111111111111111111111111111").expect("trade id");
    651     let wire = vector_wire();
    652     let first_source = ScriptedTransport::new([PageSpec {
    653         events: vec![wire.clone()],
    654         state: FetchTargetState::Complete,
    655         next: None,
    656     }]);
    657     let first = ingest_rhi_trade_source(
    658         &host.repositories(),
    659         &adapters(&first_source),
    660         &configuration,
    661         "trade-primary",
    662         trade_id,
    663         attempt("reopen-first", 1_784_347_700, 1_784_347_700),
    664     )
    665     .await
    666     .expect("first ingest");
    667     assert!(first.checkpoint_advanced());
    668     assert!(first.dirty_generation_advanced());
    669     host.close().await.expect("first close");
    670 
    671     let reopened = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
    672         .await
    673         .expect("reopen");
    674     let replay_source = ScriptedTransport::new([PageSpec {
    675         events: vec![wire],
    676         state: FetchTargetState::Complete,
    677         next: None,
    678     }]);
    679     let replay = ingest_rhi_trade_source(
    680         &reopened.repositories(),
    681         &adapters(&replay_source),
    682         &configuration,
    683         "trade-primary",
    684         trade_id,
    685         attempt("reopen-replay", 1_784_347_701, 1_784_347_701),
    686     )
    687     .await
    688     .expect("replay after reopen");
    689     assert_eq!(replay.inserted_mutations(), 0);
    690     assert_eq!(replay.inserted_signed_events(), 0);
    691     assert_eq!(replay.inserted_observations(), 1);
    692     assert!(!replay.checkpoint_advanced());
    693     assert_eq!(replay.dirty_generation().expect("dirty").get(), 1);
    694     assert!(!replay.dirty_generation_advanced());
    695     reopened.close().await.expect("second close");
    696 
    697     assert_eq!(
    698         offline_scalar(&runtime, "SELECT COUNT(*) FROM trade_mutations").await,
    699         1
    700     );
    701     assert_eq!(
    702         offline_scalar(&runtime, "SELECT COUNT(*) FROM nostr_events").await,
    703         1
    704     );
    705     assert_eq!(
    706         offline_scalar(&runtime, "SELECT COUNT(*) FROM relay_observations").await,
    707         2
    708     );
    709     assert_eq!(
    710         offline_scalar(&runtime, "SELECT generation FROM trade_dirty_generations").await,
    711         1
    712     );
    713 }
    714 
    715 #[tokio::test]
    716 async fn policy_change_dirties_existing_trade_once_and_starts_a_new_scoped_checkpoint() {
    717     let root = tempfile::tempdir().expect("root");
    718     let runtime = runtime(root.path());
    719     prepare_state_directory(&runtime);
    720     let current = configuration();
    721     let changed = parse_rhi_config_v1(
    722         CONFIG
    723             .replacen(
    724                 "policy_id = \"production-primary\"",
    725                 "policy_id = \"production-secondary\"",
    726                 1,
    727             )
    728             .as_bytes(),
    729         RhiConfigProfile::RepoLocal,
    730     )
    731     .expect("changed configuration");
    732     let metadata = metadata(&runtime, &current);
    733     let (applied_at, build) = migration_evidence();
    734     initialize_rhi_state(&runtime, &metadata, applied_at, &build)
    735         .await
    736         .expect("initialize");
    737     let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
    738         .await
    739         .expect("writer");
    740     let trade_id = TradeId::parse("11111111111111111111111111111111").expect("trade id");
    741     let wire = vector_wire();
    742     let first_source = ScriptedTransport::new([PageSpec {
    743         events: vec![wire.clone()],
    744         state: FetchTargetState::Complete,
    745         next: None,
    746     }]);
    747     let first = ingest_rhi_trade_source(
    748         &host.repositories(),
    749         &adapters(&first_source),
    750         &current,
    751         "trade-primary",
    752         trade_id,
    753         attempt("policy-first", 1_784_347_800, 1_784_347_800),
    754     )
    755     .await
    756     .expect("initial policy ingest");
    757     assert_eq!(first.dirty_generation().expect("dirty").get(), 1);
    758     host.close().await.expect("close before apply");
    759 
    760     let changed_at = MigrationAppliedAtUnixSeconds::new(1_784_347_801).expect("policy time");
    761     let (_, changed_build) = migration_evidence();
    762     let applied = apply_rhi_configuration(&runtime, &current, &changed, changed_at, &changed_build)
    763         .await
    764         .expect("policy apply");
    765     assert_eq!(applied.generation(), 2);
    766     assert!(applied.changed());
    767 
    768     let changed_host =
    769         open_rhi_state_read_write_from_config(&runtime, &changed, changed_at, &changed_build)
    770             .await
    771             .expect("changed writer");
    772     let changed_source = ScriptedTransport::new([PageSpec {
    773         events: vec![wire],
    774         state: FetchTargetState::Complete,
    775         next: None,
    776     }]);
    777     let replay = ingest_rhi_trade_source(
    778         &changed_host.repositories(),
    779         &adapters(&changed_source),
    780         &changed,
    781         "trade-primary",
    782         trade_id,
    783         attempt("policy-replay", 1_784_347_802, 1_784_347_802),
    784     )
    785     .await
    786     .expect("new-policy replay");
    787     assert_eq!(replay.inserted_mutations(), 0);
    788     assert_eq!(replay.inserted_signed_events(), 0);
    789     assert_eq!(replay.inserted_observations(), 1);
    790     assert!(replay.checkpoint_advanced());
    791     assert_eq!(replay.dirty_generation().expect("dirty").get(), 2);
    792     assert!(!replay.dirty_generation_advanced());
    793     changed_host.close().await.expect("changed close");
    794 
    795     assert_eq!(
    796         offline_scalar(&runtime, "SELECT COUNT(*) FROM relay_checkpoints").await,
    797         2
    798     );
    799     assert_eq!(
    800         offline_scalar(&runtime, "SELECT generation FROM trade_dirty_generations").await,
    801         2
    802     );
    803 }
    804 
    805 #[tokio::test]
    806 async fn concurrent_source_attempts_commit_once_and_generation_conflict_rolls_back_loser() {
    807     let root = tempfile::tempdir().expect("root");
    808     let runtime = runtime(root.path());
    809     prepare_state_directory(&runtime);
    810     let configuration = configuration();
    811     let metadata = metadata(&runtime, &configuration);
    812     let (applied_at, build) = migration_evidence();
    813     initialize_rhi_state(&runtime, &metadata, applied_at, &build)
    814         .await
    815         .expect("initialize");
    816     let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
    817         .await
    818         .expect("writer");
    819     let trade_id = TradeId::parse("11111111111111111111111111111111").expect("trade id");
    820     let wire = vector_wire();
    821     let source = ScriptedTransport::new([
    822         PageSpec {
    823             events: vec![wire.clone()],
    824             state: FetchTargetState::Complete,
    825             next: None,
    826         },
    827         PageSpec {
    828             events: vec![wire],
    829             state: FetchTargetState::Complete,
    830             next: None,
    831         },
    832     ])
    833     .with_fetch_barrier(2);
    834     let transports = adapters(&source);
    835     let repositories = host.repositories();
    836     let (first, second) = tokio::join!(
    837         ingest_rhi_trade_source(
    838             &repositories,
    839             &transports,
    840             &configuration,
    841             "trade-primary",
    842             trade_id,
    843             attempt("concurrent-first", 1_784_347_900, 1_784_347_900),
    844         ),
    845         ingest_rhi_trade_source(
    846             &repositories,
    847             &transports,
    848             &configuration,
    849             "trade-primary",
    850             trade_id,
    851             attempt("concurrent-second", 1_784_347_901, 1_784_347_901),
    852         ),
    853     );
    854     let results = [first, second];
    855     assert_eq!(results.iter().filter(|result| result.is_ok()).count(), 1);
    856     let error = results
    857         .iter()
    858         .find_map(|result| result.as_ref().err())
    859         .expect("generation conflict");
    860     assert_eq!(
    861         error.kind(),
    862         RhiTradeSourceIngestErrorKind::GenerationConflict
    863     );
    864     host.close().await.expect("close");
    865 
    866     assert_eq!(
    867         offline_scalar(&runtime, "SELECT COUNT(*) FROM trade_mutations").await,
    868         1
    869     );
    870     assert_eq!(
    871         offline_scalar(&runtime, "SELECT COUNT(*) FROM nostr_events").await,
    872         1
    873     );
    874     assert_eq!(
    875         offline_scalar(&runtime, "SELECT COUNT(*) FROM relay_observations").await,
    876         1
    877     );
    878     assert_eq!(
    879         offline_scalar(&runtime, "SELECT COUNT(*) FROM relay_checkpoints").await,
    880         1
    881     );
    882     assert_eq!(
    883         offline_scalar(&runtime, "SELECT generation FROM trade_dirty_generations").await,
    884         1
    885     );
    886 }
    887 
    888 #[test]
    889 fn attempt_and_public_diagnostics_are_bounded_and_redacted() {
    890     let observed = RhiTradeMutationObservedAtUnixSeconds::new(100).expect("observed");
    891     let policy = RhiTradeMutationAuthoredTimePolicy::new(0).expect("policy");
    892     assert!(RhiTradeSourceAttempt::new("", UnixTimeSeconds::new(1), observed, policy).is_err());
    893     assert!(
    894         RhiTradeSourceAttempt::new("x".repeat(257), UnixTimeSeconds::new(1), observed, policy,)
    895             .is_err()
    896     );
    897     let attempt = RhiTradeSourceAttempt::new(
    898         "secret-attempt-id",
    899         UnixTimeSeconds::new(99),
    900         observed,
    901         policy,
    902     )
    903     .expect("attempt");
    904     assert!(!format!("{attempt:?}").contains("secret-attempt-id"));
    905 
    906     let error = RhiTradeSourceAttempt::new(
    907         "\nsecret-request",
    908         UnixTimeSeconds::new(99),
    909         observed,
    910         policy,
    911     )
    912     .expect_err("control character");
    913     assert_eq!(error.code(), "trade_source_input_invalid");
    914     assert!(!format!("{error}").contains("secret-request"));
    915     assert!(!format!("{error:?}").contains("secret-request"));
    916     assert!(std::error::Error::source(&error).is_none());
    917 }