services_hardening_discovery_state.rs (27194B)
1 #![forbid(unsafe_code)] 2 #![cfg(any(target_os = "linux", target_os = "macos"))] 3 4 use std::{error::Error, fs, os::unix::fs::PermissionsExt, path::Path}; 5 6 use myc::{ 7 MYC_DISCOVERY_DOCUMENT_MAX_BYTES, MYC_STATE_SCHEMA_VERSION, MycConfigProfile, 8 MycDeliveryAttemptNonce, MycDeliveryAttemptOutcome, MycDeliveryClaim, MycDeliveryJobStatus, 9 MycDeliveryRecoveryEntropy, MycDeliveryRelayId, MycDeliveryRetryJitter, MycDeliverySourceKind, 10 MycDeliveryTimeUnixMs, MycDiscoveryCommitAdmission, MycDiscoveryCommitRequest, 11 MycDiscoveryStateErrorKind, MycNip05ExportSelection, MycStateMetadata, 12 MycStateRepositoryErrorKind, RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, 13 initialize_myc_state, open_myc_state_inspection, open_myc_state_read_write, 14 parse_myc_cli_v1_from, parse_myc_config_v1, resolve_myc_runtime_context, 15 }; 16 use nostr::{Keys, SecretKey}; 17 use radroots_nostr::event::{ 18 ApplicationHandlerSpec as RadrootsNostrApplicationHandlerSpec, 19 Metadata as RadrootsNostrMetadata, Timestamp as RadrootsNostrTimestamp, 20 build_application_handler as radroots_nostr_build_application_handler_event, 21 }; 22 use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity}; 23 use radroots_storage::event::SourceGeneration; 24 use sha2::{Digest, Sha256}; 25 use sqlx::{ConnectOptions, Connection, sqlite::SqliteConnectOptions}; 26 27 const CONFIG_EXAMPLE: &[u8] = 28 include_bytes!("../contracts/services_hardening/config.v1.example.toml"); 29 const DISCOVERY_SOURCE: &str = include_str!("../src/state_discovery.rs"); 30 const CATALOG_SOURCE: &str = include_str!("../src/state_catalog.rs"); 31 const LIB_SOURCE: &str = include_str!("../src/lib.rs"); 32 const DISCOVERY_SECRET: &str = "3333333333333333333333333333333333333333333333333333333333333333"; 33 34 fn runtime(root: &Path) -> myc::MycRuntimeContext { 35 let root = root.to_str().expect("UTF-8 temporary root"); 36 let invocation = parse_myc_cli_v1_from([ 37 "myc", 38 "--profile", 39 "repo-local", 40 "--instance", 41 "primary", 42 "--repo-local-root", 43 root, 44 "run", 45 ]) 46 .expect("valid invocation"); 47 resolve_myc_runtime_context( 48 &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()), 49 &invocation, 50 ) 51 .expect("runtime context") 52 } 53 54 fn prepare_state_directory(runtime: &myc::MycRuntimeContext) { 55 let directory = runtime.context().paths().state(); 56 fs::create_dir_all(directory).expect("state directory"); 57 fs::set_permissions(directory, fs::Permissions::from_mode(0o700)).expect("state mode"); 58 } 59 60 fn discovery_keys() -> Keys { 61 Keys::new(SecretKey::parse(DISCOVERY_SECRET).expect("discovery secret")) 62 } 63 64 fn config_source(enabled: bool) -> Vec<u8> { 65 let identity = discovery_keys(); 66 let source = String::from_utf8(CONFIG_EXAMPLE.to_vec()) 67 .expect("UTF-8 configuration") 68 .replace( 69 "3333333333333333333333333333333333333333333333333333333333333333", 70 &identity.public_key().to_hex(), 71 ); 72 if enabled { 73 return source.into_bytes(); 74 } 75 let before_binding = source 76 .split_once("[identity.discovery.binding]") 77 .expect("discovery binding") 78 .0; 79 let after_binding = source.split_once("[[relays]]").expect("relay inventory").1; 80 let without_binding = format!( 81 "{}[[relays]]{}", 82 before_binding.replace( 83 "[identity.discovery]\nenabled = true", 84 "[identity.discovery]\nenabled = false" 85 ), 86 after_binding, 87 ); 88 let before_discovery = without_binding 89 .split_once("[discovery]") 90 .expect("discovery section") 91 .0; 92 format!("{before_discovery}[discovery]\nenabled = false\n").into_bytes() 93 } 94 95 fn state_metadata(runtime: &myc::MycRuntimeContext, source: &[u8]) -> MycStateMetadata { 96 let configuration = 97 parse_myc_config_v1(source, MycConfigProfile::RepoLocal).expect("configuration"); 98 MycStateMetadata::new( 99 runtime, 100 &configuration, 101 SourceGeneration::new([0x5a; 32]).expect("generation"), 102 1_725_000_000_000, 103 ) 104 .expect("metadata") 105 } 106 107 fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) { 108 let applied_at = MigrationAppliedAtUnixSeconds::new(1_725_000_000).expect("migration time"); 109 let build = MigrationBuildIdentity::new( 110 env!("CARGO_PKG_VERSION"), 111 "1111111111111111111111111111111111111111", 112 "053d0c750bf9cd683c6ea37cefe7e79617ba629f", 113 "rustc-test", 114 "test-target", 115 "service-host", 116 1, 117 MYC_STATE_SCHEMA_VERSION, 118 1, 119 1, 120 1, 121 ) 122 .expect("build identity"); 123 (applied_at, build) 124 } 125 126 fn time(value: u64) -> MycDeliveryTimeUnixMs { 127 MycDeliveryTimeUnixMs::new(value).expect("delivery time") 128 } 129 130 fn nostrconnect_url() -> String { 131 let mut query = url::form_urlencoded::Serializer::new(String::new()); 132 query.append_pair("relay", "wss://relay-primary.example.test/"); 133 query.append_pair("relay", "wss://relay-secondary.example.test/"); 134 let bunker = format!( 135 "bunker://{}?{}", 136 "2222222222222222222222222222222222222222222222222222222222222222", 137 query.finish() 138 ); 139 let encoded: String = url::form_urlencoded::byte_serialize(bunker.as_bytes()).collect(); 140 format!("https://myc.example.test/connect?uri={encoded}") 141 } 142 143 fn signed_handler_event(created_at: u64) -> Vec<u8> { 144 let metadata = RadrootsNostrMetadata { 145 name: Some("myc".to_owned()), 146 display_name: Some("Radroots Myc".to_owned()), 147 about: Some("NIP-46 signer".to_owned()), 148 website: Some("https://myc.example.test/".to_owned()), 149 picture: Some("https://myc.example.test/myc.png".to_owned()), 150 ..RadrootsNostrMetadata::default() 151 }; 152 let spec = RadrootsNostrApplicationHandlerSpec::new(vec![24_133]) 153 .with_identifier("myc") 154 .with_relays(vec![ 155 "wss://relay-primary.example.test/".to_owned(), 156 "wss://relay-secondary.example.test/".to_owned(), 157 ]) 158 .with_nostr_connect_url(nostrconnect_url()) 159 .with_metadata(metadata); 160 let event = radroots_nostr_build_application_handler_event(&spec) 161 .expect("typed handler event") 162 .custom_created_at(RadrootsNostrTimestamp::from_secs(created_at)) 163 .sign_with_keys(&discovery_keys()) 164 .expect("signed event"); 165 serde_json::to_vec(&event).expect("canonical event bytes") 166 } 167 168 async fn deliver_all_required( 169 repository: &myc::MycStateRepository<'_>, 170 job_id: myc::MycDeliveryJobId, 171 start: u64, 172 ) { 173 for (offset, relay_name) in ["primary", "secondary"].into_iter().enumerate() { 174 let relay = MycDeliveryRelayId::new(relay_name).expect("relay"); 175 let claimed = match repository 176 .claim_delivery_target( 177 job_id, 178 &relay, 179 MycDeliveryAttemptNonce::from_injected_entropy([0x70 + offset as u8; 32]), 180 time(start + offset as u64 * 10), 181 ) 182 .await 183 .expect("claim") 184 { 185 MycDeliveryClaim::Claimed(attempt) => attempt, 186 other => panic!("unexpected claim: {other:?}"), 187 }; 188 repository 189 .mark_delivery_attempt_submitted( 190 job_id, 191 &relay, 192 claimed.id(), 193 time(start + offset as u64 * 10 + 1), 194 ) 195 .await 196 .expect("submitted"); 197 repository 198 .record_delivery_attempt_outcome( 199 job_id, 200 &relay, 201 claimed.id(), 202 MycDeliveryAttemptOutcome::Delivered, 203 MycDeliveryRetryJitter::new(0).expect("zero retry jitter"), 204 time(start + offset as u64 * 10 + 2), 205 ) 206 .await 207 .expect("delivered"); 208 } 209 } 210 211 #[test] 212 fn discovery_input_is_signature_verified_bounded_and_redacted() { 213 let directory = tempfile::tempdir().expect("temporary root"); 214 let runtime = runtime(directory.path()); 215 let metadata = state_metadata(&runtime, &config_source(true)); 216 let event = signed_handler_event(1_725_000_000); 217 let request = MycDiscoveryCommitRequest::new(&metadata, &event, time(100)).expect("request"); 218 assert!( 219 !request 220 .generation_id() 221 .as_bytes() 222 .iter() 223 .all(|byte| *byte == 0) 224 ); 225 226 let mut altered = event.clone(); 227 let position = altered 228 .iter() 229 .position(|byte| *byte == b'm') 230 .expect("mutable event byte"); 231 altered[position] = b'n'; 232 assert_eq!( 233 MycDiscoveryCommitRequest::new(&metadata, &altered, time(100)) 234 .expect_err("signature/canonical drift") 235 .kind(), 236 MycDiscoveryStateErrorKind::InvalidEvent 237 ); 238 let mut whitespace = event.clone(); 239 whitespace.push(b'\n'); 240 assert_eq!( 241 MycDiscoveryCommitRequest::new(&metadata, &whitespace, time(100)) 242 .expect_err("noncanonical event") 243 .kind(), 244 MycDiscoveryStateErrorKind::InvalidEvent 245 ); 246 assert_eq!( 247 MycDiscoveryCommitRequest::new( 248 &metadata, 249 &vec![b'x'; MYC_DISCOVERY_DOCUMENT_MAX_BYTES + 1], 250 time(100), 251 ) 252 .expect_err("oversize") 253 .kind(), 254 MycDiscoveryStateErrorKind::TooLarge 255 ); 256 let disabled = state_metadata(&runtime, &config_source(false)); 257 let error = MycDiscoveryCommitRequest::new(&disabled, &event, time(100)) 258 .expect_err("disabled discovery"); 259 assert_eq!(error.kind(), MycDiscoveryStateErrorKind::Disabled); 260 assert!(Error::source(&error).is_none()); 261 let rendered = format!("{request:?} {error} {error:?}"); 262 assert!(!rendered.contains(DISCOVERY_SECRET)); 263 assert!(!rendered.contains("myc.example.test")); 264 assert!(!rendered.contains(&String::from_utf8(event).expect("event text"))); 265 } 266 267 #[tokio::test] 268 async fn discovery_desired_current_and_exact_documents_are_atomic_restart_safe_and_delivery_bound() 269 { 270 let directory = tempfile::tempdir().expect("temporary root"); 271 let runtime = runtime(directory.path()); 272 prepare_state_directory(&runtime); 273 let metadata = state_metadata(&runtime, &config_source(true)); 274 let (applied_at, build) = migration_evidence(); 275 initialize_myc_state(&runtime, &metadata, applied_at, &build) 276 .await 277 .expect("initialization"); 278 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 279 .await 280 .expect("writer"); 281 let event = signed_handler_event(1_725_000_000); 282 let request = MycDiscoveryCommitRequest::new(&metadata, &event, time(100)).expect("request"); 283 let committed = host 284 .repository() 285 .commit_discovery_desired_state(&request) 286 .await 287 .expect("commit"); 288 assert!(matches!(committed, MycDiscoveryCommitAdmission::Created(_))); 289 let record = committed.record(); 290 assert_eq!( 291 record.job().source_kind(), 292 MycDeliverySourceKind::DiscoveryHandler 293 ); 294 assert_eq!(record.job().operation_id(), None); 295 assert_eq!(record.job().status(), MycDeliveryJobStatus::Pending); 296 assert!( 297 !request 298 .desired_digest() 299 .as_bytes() 300 .iter() 301 .all(|byte| *byte == 0) 302 ); 303 assert_eq!(record.document().event_id().len(), 32); 304 assert!( 305 !record 306 .document() 307 .nip05_projection_digest() 308 .as_bytes() 309 .iter() 310 .all(|byte| *byte == 0) 311 ); 312 assert_eq!(record.document().event_bytes(), event); 313 let projection: serde_json::Value = 314 serde_json::from_slice(record.document().nip05_projection_bytes()).expect("projection"); 315 assert_eq!( 316 projection["schema"], 317 "radroots.myc.nip05-projection-input.v1" 318 ); 319 assert_eq!(projection["domain"], "myc.example.test"); 320 assert_eq!(projection["name"], "_"); 321 assert_eq!(projection["relays"].as_array().expect("relays").len(), 2); 322 assert_eq!(record.state().current_generation_id(), None); 323 assert_eq!(record.state().current_job_id(), None); 324 let job_id = record.job().id(); 325 let desired_export = host 326 .repository() 327 .render_offline_nip05(MycNip05ExportSelection::Desired) 328 .await 329 .expect("desired NIP-05 export"); 330 assert_eq!(desired_export.selection(), MycNip05ExportSelection::Desired); 331 assert_eq!(desired_export.domain(), "myc.example.test"); 332 let expected_export = format!( 333 "{{\"names\":{{\"_\":\"{}\"}},\"nip46\":{{\"relays\":[\"wss://relay-primary.example.test/\",\"wss://relay-secondary.example.test/\"],\"nostrconnect_url\":\"{}\"}}}}", 334 discovery_keys().public_key().to_hex(), 335 nostrconnect_url() 336 ); 337 assert_eq!(desired_export.bytes(), expected_export.as_bytes()); 338 assert_eq!( 339 desired_export.digest().as_bytes(), 340 &<[u8; 32]>::from(Sha256::digest(expected_export.as_bytes())) 341 ); 342 assert_eq!( 343 host.repository() 344 .render_offline_nip05(MycNip05ExportSelection::Current) 345 .await 346 .expect_err("no proven-current generation") 347 .kind(), 348 MycStateRepositoryErrorKind::Binding 349 ); 350 351 let replay = host 352 .repository() 353 .commit_discovery_desired_state(&request) 354 .await 355 .expect("exact replay"); 356 assert!(matches!( 357 replay, 358 MycDiscoveryCommitAdmission::ExactReplay(_) 359 )); 360 assert_eq!(replay.record().job().id(), job_id); 361 let same_time_event = signed_handler_event(1_725_000_002); 362 let same_time_request = MycDiscoveryCommitRequest::new(&metadata, &same_time_event, time(100)) 363 .expect("same-time request"); 364 assert_eq!( 365 host.repository() 366 .commit_discovery_desired_state(&same_time_request) 367 .await 368 .expect_err("different desired state at the same time") 369 .kind(), 370 MycStateRepositoryErrorKind::Binding 371 ); 372 assert_eq!( 373 host.repository() 374 .promote_delivered_discovery_state(job_id, time(101)) 375 .await 376 .expect_err("pending job cannot become current") 377 .kind(), 378 MycStateRepositoryErrorKind::Binding 379 ); 380 deliver_all_required(&host.repository(), job_id, 110).await; 381 let promoted = host 382 .repository() 383 .promote_delivered_discovery_state(job_id, time(140)) 384 .await 385 .expect("promoted current"); 386 assert_eq!( 387 promoted.current_generation_id(), 388 Some(request.generation_id()) 389 ); 390 assert_eq!(promoted.current_job_id(), Some(job_id)); 391 assert_eq!( 392 host.repository() 393 .render_offline_nip05(MycNip05ExportSelection::Current) 394 .await 395 .expect("current NIP-05 export") 396 .bytes(), 397 expected_export.as_bytes() 398 ); 399 assert_eq!( 400 host.repository() 401 .read_discovery_document_for_job(job_id) 402 .await 403 .expect("document read") 404 .expect("document") 405 .event_bytes(), 406 event 407 ); 408 let next_event = signed_handler_event(1_725_000_001); 409 let next_request = 410 MycDiscoveryCommitRequest::new(&metadata, &next_event, time(150)).expect("next request"); 411 let next = host 412 .repository() 413 .commit_discovery_desired_state(&next_request) 414 .await 415 .expect("next desired state"); 416 assert!(matches!(next, MycDiscoveryCommitAdmission::Created(_))); 417 assert_eq!( 418 next.record().state().desired_generation_id(), 419 next_request.generation_id() 420 ); 421 assert_eq!( 422 next.record().state().current_generation_id(), 423 Some(request.generation_id()) 424 ); 425 assert_eq!(next.record().state().current_job_id(), Some(job_id)); 426 let next_job_id = next.record().job().id(); 427 assert_eq!( 428 host.repository() 429 .promote_delivered_discovery_state(next_job_id, time(151)) 430 .await 431 .expect_err("new desired state is not yet current") 432 .kind(), 433 MycStateRepositoryErrorKind::Binding 434 ); 435 deliver_all_required(&host.repository(), next_job_id, 160).await; 436 let next_promoted = host 437 .repository() 438 .promote_delivered_discovery_state(next_job_id, time(190)) 439 .await 440 .expect("next current state"); 441 assert_eq!( 442 next_promoted.current_generation_id(), 443 Some(next_request.generation_id()) 444 ); 445 assert_eq!(next_promoted.current_job_id(), Some(next_job_id)); 446 host.close().await.expect("close"); 447 448 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 449 .await 450 .expect("reopen"); 451 let state = host 452 .repository() 453 .read_discovery_publication_state() 454 .await 455 .expect("state") 456 .expect("committed state"); 457 assert_eq!(state.current_job_id(), Some(next_job_id)); 458 host.close().await.expect("final close"); 459 460 let inspection = open_myc_state_inspection(&runtime, &metadata) 461 .await 462 .expect("inspection host"); 463 let offline = inspection 464 .repository() 465 .render_offline_nip05(MycNip05ExportSelection::Current) 466 .await 467 .expect("read-only offline export"); 468 assert_eq!(offline.bytes(), expected_export.as_bytes()); 469 let rendered = format!("{offline:?}"); 470 assert!(!rendered.contains("myc.example.test")); 471 assert!(!rendered.contains(&discovery_keys().public_key().to_hex())); 472 inspection.close().await.expect("inspection close"); 473 } 474 475 #[tokio::test] 476 async fn restart_recovery_is_bounded_jittered_and_idempotent() { 477 let directory = tempfile::tempdir().expect("temporary root"); 478 let runtime = runtime(directory.path()); 479 prepare_state_directory(&runtime); 480 let bounded_source = String::from_utf8(config_source(true)) 481 .expect("UTF-8 configuration") 482 .replace("outbox = 4096", "outbox = 1"); 483 let metadata = state_metadata(&runtime, bounded_source.as_bytes()); 484 let (applied_at, build) = migration_evidence(); 485 initialize_myc_state(&runtime, &metadata, applied_at, &build) 486 .await 487 .expect("initialization"); 488 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 489 .await 490 .expect("writer"); 491 let request = 492 MycDiscoveryCommitRequest::new(&metadata, &signed_handler_event(1_725_000_020), time(100)) 493 .expect("discovery request"); 494 let committed = host 495 .repository() 496 .commit_discovery_desired_state(&request) 497 .await 498 .expect("desired state"); 499 let job_id = committed.record().job().id(); 500 assert!(matches!( 501 host.repository() 502 .commit_discovery_desired_state(&request) 503 .await 504 .expect("exact replay at queue ceiling"), 505 MycDiscoveryCommitAdmission::ExactReplay(_) 506 )); 507 let saturated = 508 MycDiscoveryCommitRequest::new(&metadata, &signed_handler_event(1_725_000_021), time(101)) 509 .expect("second desired state"); 510 assert_eq!( 511 host.repository() 512 .commit_discovery_desired_state(&saturated) 513 .await 514 .expect_err("outbox ceiling") 515 .kind(), 516 MycStateRepositoryErrorKind::Binding 517 ); 518 let primary = MycDeliveryRelayId::new("primary").expect("primary"); 519 let attempt = match host 520 .repository() 521 .claim_delivery_target( 522 job_id, 523 &primary, 524 MycDeliveryAttemptNonce::from_injected_entropy([0xa1; 32]), 525 time(120), 526 ) 527 .await 528 .expect("claim") 529 { 530 MycDeliveryClaim::Claimed(attempt) => attempt, 531 other => panic!("unexpected claim: {other:?}"), 532 }; 533 host.repository() 534 .mark_delivery_attempt_submitted(job_id, &primary, attempt.id(), time(121)) 535 .await 536 .expect("submitted"); 537 host.close().await.expect("close before restart"); 538 539 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 540 .await 541 .expect("reopened writer"); 542 let recovered_at = time(15_121); 543 let report = host 544 .repository() 545 .recover_delivery_state( 546 recovered_at, 547 MycDeliveryRecoveryEntropy::from_injected_entropy([0xb2; 32]), 548 ) 549 .await 550 .expect("restart recovery"); 551 assert_eq!(report.examined_jobs(), 1); 552 assert_eq!(report.recovered_expired_attempts(), 1); 553 assert_eq!(report.finalized_jobs(), 0); 554 assert_eq!(report.promoted_discovery_generations(), 0); 555 assert_eq!(report.active_targets(), 0); 556 assert_eq!(report.ready_targets() + report.scheduled_targets(), 2); 557 let job = host 558 .repository() 559 .read_delivery_job(job_id) 560 .await 561 .expect("job read") 562 .expect("job"); 563 let next = job.targets()[0] 564 .next_attempt_at() 565 .expect("persisted jittered retry"); 566 assert!(next >= recovered_at); 567 assert!(next.get() <= recovered_at.get() + 250); 568 let replay = host 569 .repository() 570 .recover_delivery_state( 571 recovered_at, 572 MycDeliveryRecoveryEntropy::from_injected_entropy([0xc3; 32]), 573 ) 574 .await 575 .expect("idempotent recovery"); 576 assert_eq!(replay.recovered_expired_attempts(), 0); 577 assert_eq!(replay.examined_jobs(), 1); 578 assert_eq!( 579 host.repository() 580 .read_delivery_attempts(job_id, &primary) 581 .await 582 .expect("attempt history") 583 .len(), 584 1 585 ); 586 deliver_all_required(&host.repository(), job_id, 16_000).await; 587 assert_eq!( 588 host.repository() 589 .render_offline_nip05(MycNip05ExportSelection::Current) 590 .await 591 .expect_err("delivered desired state is not promoted implicitly") 592 .kind(), 593 MycStateRepositoryErrorKind::Binding 594 ); 595 host.close().await.expect("close after delivered job"); 596 597 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 598 .await 599 .expect("reopened writer after delivered job"); 600 let promotion = host 601 .repository() 602 .recover_delivery_state( 603 time(16_100), 604 MycDeliveryRecoveryEntropy::from_injected_entropy([0xd4; 32]), 605 ) 606 .await 607 .expect("restart promotion recovery"); 608 assert_eq!(promotion.examined_jobs(), 1); 609 assert_eq!(promotion.recovered_expired_attempts(), 0); 610 assert_eq!(promotion.finalized_jobs(), 0); 611 assert_eq!(promotion.promoted_discovery_generations(), 1); 612 assert!( 613 host.repository() 614 .render_offline_nip05(MycNip05ExportSelection::Current) 615 .await 616 .is_ok() 617 ); 618 let settled = host 619 .repository() 620 .recover_delivery_state( 621 time(16_101), 622 MycDeliveryRecoveryEntropy::from_injected_entropy([0xe5; 32]), 623 ) 624 .await 625 .expect("settled recovery"); 626 assert_eq!(settled.examined_jobs(), 0); 627 assert_eq!(settled.promoted_discovery_generations(), 0); 628 host.close().await.expect("final close"); 629 } 630 631 #[tokio::test] 632 async fn discovery_schema_guards_reject_mutation_and_corrupt_source_kinds() { 633 let directory = tempfile::tempdir().expect("temporary root"); 634 let runtime = runtime(directory.path()); 635 prepare_state_directory(&runtime); 636 let metadata = state_metadata(&runtime, &config_source(true)); 637 let (applied_at, build) = migration_evidence(); 638 initialize_myc_state(&runtime, &metadata, applied_at, &build) 639 .await 640 .expect("initialization"); 641 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 642 .await 643 .expect("writer"); 644 let request = 645 MycDiscoveryCommitRequest::new(&metadata, &signed_handler_event(1_725_000_000), time(100)) 646 .expect("request"); 647 host.repository() 648 .commit_discovery_desired_state(&request) 649 .await 650 .expect("commit"); 651 host.close().await.expect("close"); 652 653 let options = SqliteConnectOptions::new() 654 .filename(runtime.artifacts().state_database()) 655 .create_if_missing(false) 656 .disable_statement_logging(); 657 let mut connection = sqlx::SqliteConnection::connect_with(&options) 658 .await 659 .expect("inspection connection"); 660 for statement in [ 661 "DELETE FROM discovery_desired_state", 662 "UPDATE discovery_documents SET event_bytes = X'01'", 663 "DELETE FROM discovery_publication_state", 664 "UPDATE discovery_publication_state SET desired_job_id = zeroblob(32)", 665 "UPDATE delivery_jobs SET source_kind = 'invalid'", 666 ] { 667 assert!( 668 sqlx::query(statement) 669 .execute(&mut connection) 670 .await 671 .is_err(), 672 "guard accepted {statement}" 673 ); 674 } 675 for (index, source_kind) in ["signer_response", "discovery_handler"] 676 .into_iter() 677 .enumerate() 678 { 679 let job_id = [0xa0 + u8::try_from(index).expect("bounded index"); 32]; 680 let source_id = [0xf0 + u8::try_from(index).expect("bounded index"); 32]; 681 let artifact = [0xb0 + u8::try_from(index).expect("bounded index"); 32]; 682 let result = sqlx::query( 683 "INSERT INTO delivery_jobs (job_id, source_kind, source_id, artifact_sha256, \ 684 policy_mode, required_acknowledgements, max_attempts, initial_backoff_ms, \ 685 maximum_backoff_ms, attempt_deadline_ms, status, created_at_unix_ms, \ 686 updated_at_unix_ms, finalized_at_unix_ms) \ 687 VALUES (?, ?, ?, ?, 'all_required', 1, 1, 1, 1, 1, 'pending', 1, 1, NULL)", 688 ) 689 .bind(job_id.as_slice()) 690 .bind(source_kind) 691 .bind(source_id.as_slice()) 692 .bind(artifact.as_slice()) 693 .execute(&mut connection) 694 .await; 695 assert!(result.is_err(), "accepted orphaned {source_kind} source"); 696 } 697 connection.close().await.expect("connection close"); 698 } 699 700 #[test] 701 fn discovery_boundary_is_typed_sqlx_only_and_commits_before_any_publication() { 702 assert!(LIB_SOURCE.contains("mod state_discovery;")); 703 assert!(!LIB_SOURCE.contains("pub mod state_discovery;")); 704 assert!(DISCOVERY_SOURCE.contains("event.verify()")); 705 assert!(DISCOVERY_SOURCE.contains("ServiceSqliteTransaction<'_>")); 706 assert!(DISCOVERY_SOURCE.contains("source_kind = 'discovery_handler'")); 707 assert!(CATALOG_SOURCE.contains("discovery_desired_state_no_update")); 708 assert!(CATALOG_SOURCE.contains("discovery_documents_no_delete")); 709 assert!(CATALOG_SOURCE.contains("delivery_jobs_guard_insert")); 710 for forbidden in [ 711 "SqlitePool", 712 "SqliteConnection", 713 "BEGIN ", 714 "COMMIT", 715 "ROLLBACK", 716 "reqwest", 717 "send_event", 718 "publish_nip89_event", 719 "tokio::spawn", 720 "std::env", 721 "std::time", 722 "rusqlite", 723 ] { 724 assert!( 725 !DISCOVERY_SOURCE.contains(forbidden), 726 "found forbidden discovery authority `{forbidden}`" 727 ); 728 } 729 }