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, ¤t); 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 ¤t, 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, ¤t, &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 }