services_hardening_delivery_state.rs (21930B)
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_DELIVERY_RELAY_ID_MAX_BYTES, MYC_DELIVERY_RETRY_JITTER_MAX_MS, MYC_STATE_SCHEMA_VERSION, 8 MycConfigProfile, MycDeliveryArtifactDigest, MycDeliveryAttemptNonce, 9 MycDeliveryAttemptOutcome, MycDeliveryAttemptStatus, MycDeliveryClaim, MycDeliveryJobStatus, 10 MycDeliveryPolicyMode, MycDeliveryRelayId, MycDeliveryRetryJitter, MycDeliveryStateErrorKind, 11 MycDeliveryTargetStatus, MycDeliveryTimeUnixMs, MycDiscoveryCommitAdmission, 12 MycDiscoveryCommitRequest, MycStateMetadata, MycStateRepositoryErrorKind, 13 RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, initialize_myc_state, 14 open_myc_state_read_write, parse_myc_cli_v1_from, parse_myc_config_v1, 15 resolve_myc_runtime_context, 16 }; 17 use nostr::{Keys, SecretKey}; 18 use radroots_nostr::event::{ 19 ApplicationHandlerSpec as RadrootsNostrApplicationHandlerSpec, 20 Metadata as RadrootsNostrMetadata, Timestamp as RadrootsNostrTimestamp, 21 build_application_handler as radroots_nostr_build_application_handler_event, 22 }; 23 use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity}; 24 use radroots_storage::event::SourceGeneration; 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 DELIVERY_SOURCE: &str = include_str!("../src/state_delivery.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 metadata(runtime: &myc::MycRuntimeContext, source: &[u8]) -> MycStateMetadata { 61 let configuration = 62 parse_myc_config_v1(source, MycConfigProfile::RepoLocal).expect("configuration"); 63 MycStateMetadata::new( 64 runtime, 65 &configuration, 66 SourceGeneration::new([0x5a; 32]).expect("generation"), 67 1_725_000_000_000, 68 ) 69 .expect("metadata") 70 } 71 72 fn discovery_keys() -> Keys { 73 Keys::new(SecretKey::parse(DISCOVERY_SECRET).expect("discovery secret")) 74 } 75 76 fn config_source(source: &[u8]) -> Vec<u8> { 77 String::from_utf8(source.to_vec()) 78 .expect("UTF-8 configuration") 79 .replace( 80 "3333333333333333333333333333333333333333333333333333333333333333", 81 &discovery_keys().public_key().to_hex(), 82 ) 83 .into_bytes() 84 } 85 86 fn nostrconnect_url() -> String { 87 let mut query = url::form_urlencoded::Serializer::new(String::new()); 88 query.append_pair("relay", "wss://relay-primary.example.test/"); 89 query.append_pair("relay", "wss://relay-secondary.example.test/"); 90 let bunker = format!( 91 "bunker://{}?{}", 92 "2222222222222222222222222222222222222222222222222222222222222222", 93 query.finish() 94 ); 95 let encoded: String = url::form_urlencoded::byte_serialize(bunker.as_bytes()).collect(); 96 format!("https://myc.example.test/connect?uri={encoded}") 97 } 98 99 fn signed_handler_event(created_at: u64) -> Vec<u8> { 100 let metadata = RadrootsNostrMetadata { 101 name: Some("myc".to_owned()), 102 display_name: Some("Radroots Myc".to_owned()), 103 about: Some("NIP-46 signer".to_owned()), 104 website: Some("https://myc.example.test/".to_owned()), 105 picture: Some("https://myc.example.test/myc.png".to_owned()), 106 ..RadrootsNostrMetadata::default() 107 }; 108 let spec = RadrootsNostrApplicationHandlerSpec::new(vec![24_133]) 109 .with_identifier("myc") 110 .with_relays(vec![ 111 "wss://relay-primary.example.test/".to_owned(), 112 "wss://relay-secondary.example.test/".to_owned(), 113 ]) 114 .with_nostr_connect_url(nostrconnect_url()) 115 .with_metadata(metadata); 116 let event = radroots_nostr_build_application_handler_event(&spec) 117 .expect("typed handler event") 118 .custom_created_at(RadrootsNostrTimestamp::from_secs(created_at)) 119 .sign_with_keys(&discovery_keys()) 120 .expect("signed event"); 121 serde_json::to_vec(&event).expect("canonical event bytes") 122 } 123 124 fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) { 125 let applied_at = MigrationAppliedAtUnixSeconds::new(1_725_000_000).expect("migration time"); 126 let build = MigrationBuildIdentity::new( 127 env!("CARGO_PKG_VERSION"), 128 "1111111111111111111111111111111111111111", 129 "053d0c750bf9cd683c6ea37cefe7e79617ba629f", 130 "rustc-test", 131 "test-target", 132 "service-host", 133 1, 134 MYC_STATE_SCHEMA_VERSION, 135 1, 136 1, 137 1, 138 ) 139 .expect("build identity"); 140 (applied_at, build) 141 } 142 143 fn time(value: u64) -> MycDeliveryTimeUnixMs { 144 MycDeliveryTimeUnixMs::new(value).expect("delivery time") 145 } 146 147 fn jitter(value: u64) -> MycDeliveryRetryJitter { 148 MycDeliveryRetryJitter::new(value).expect("retry jitter") 149 } 150 151 #[test] 152 fn delivery_inputs_and_diagnostics_are_closed_bounded_and_redacted() { 153 let maximum = format!("a{}", "1".repeat(MYC_DELIVERY_RELAY_ID_MAX_BYTES - 1)); 154 assert!(MycDeliveryRelayId::new(&maximum).is_ok()); 155 for invalid in [ 156 "", 157 "Primary", 158 "relay-name", 159 "relay__name", 160 "relay_", 161 &format!("a{maximum}"), 162 &"x".repeat(1024 * 1024), 163 ] { 164 assert_eq!( 165 MycDeliveryRelayId::new(invalid) 166 .expect_err("invalid relay") 167 .kind(), 168 MycDeliveryStateErrorKind::InvalidRelayId 169 ); 170 } 171 for invalid in [0, i64::MAX.unsigned_abs() + 1] { 172 assert_eq!( 173 MycDeliveryTimeUnixMs::new(invalid) 174 .expect_err("invalid time") 175 .kind(), 176 MycDeliveryStateErrorKind::InvalidTime 177 ); 178 } 179 assert!(MycDeliveryTimeUnixMs::new(i64::MAX.unsigned_abs()).is_ok()); 180 assert!(MycDeliveryRetryJitter::new(MYC_DELIVERY_RETRY_JITTER_MAX_MS).is_ok()); 181 assert_eq!( 182 MycDeliveryRetryJitter::new(MYC_DELIVERY_RETRY_JITTER_MAX_MS + 1) 183 .expect_err("oversize retry jitter") 184 .kind(), 185 MycDeliveryStateErrorKind::InvalidRetryJitter 186 ); 187 assert_eq!( 188 [ 189 MycDeliveryPolicyMode::AtLeastOneRequired, 190 MycDeliveryPolicyMode::AllRequired, 191 MycDeliveryPolicyMode::RequiredQuorum, 192 ] 193 .map(MycDeliveryPolicyMode::as_str), 194 ["at_least_one_required", "all_required", "required_quorum"] 195 ); 196 197 let relay = MycDeliveryRelayId::new("relay_secret_123").expect("relay"); 198 let nonce = MycDeliveryAttemptNonce::from_injected_entropy([0x91; 32]); 199 let digest = MycDeliveryArtifactDigest::from_bytes([0x92; 32]); 200 let error = MycDeliveryRelayId::new("secret-invalid").expect_err("invalid"); 201 assert!(Error::source(&error).is_none()); 202 let rendered = format!("{relay:?} {nonce:?} {digest:?} {error} {error:?}"); 203 for secret in ["relay_secret_123", "secret-invalid", "145, 145", "146, 146"] { 204 assert!(!rendered.contains(secret)); 205 } 206 } 207 208 #[tokio::test] 209 async fn delivery_jobs_are_config_bound_idempotent_restart_safe_and_unknown_aware() { 210 let directory = tempfile::tempdir().expect("temporary root"); 211 let runtime = runtime(directory.path()); 212 prepare_state_directory(&runtime); 213 let source = config_source(CONFIG_EXAMPLE); 214 let metadata = metadata(&runtime, &source); 215 let (applied_at, build) = migration_evidence(); 216 initialize_myc_state(&runtime, &metadata, applied_at, &build) 217 .await 218 .expect("initialization"); 219 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 220 .await 221 .expect("writer"); 222 let request = 223 MycDiscoveryCommitRequest::new(&metadata, &signed_handler_event(1_725_000_000), time(110)) 224 .expect("discovery request"); 225 let job = host 226 .repository() 227 .commit_discovery_desired_state(&request) 228 .await 229 .expect("job"); 230 assert!(matches!(job, MycDiscoveryCommitAdmission::Created(_))); 231 let job_id = job.record().job().id(); 232 assert_eq!( 233 job.record().job().policy_mode(), 234 MycDeliveryPolicyMode::AllRequired 235 ); 236 assert_eq!(job.record().job().required_acknowledgements(), 2); 237 assert_eq!(job.record().job().max_attempts(), 5); 238 assert_eq!(job.record().job().initial_backoff_ms(), 250); 239 assert_eq!(job.record().job().maximum_backoff_ms(), 30_000); 240 assert_eq!(job.record().job().attempt_deadline_ms(), 15_000); 241 assert_eq!(job.record().job().targets().len(), 2); 242 assert_eq!( 243 job.record().job().targets()[0].relay_id().as_str(), 244 "primary" 245 ); 246 assert_eq!( 247 job.record().job().targets()[1].relay_id().as_str(), 248 "secondary" 249 ); 250 assert!( 251 job.record() 252 .job() 253 .targets() 254 .iter() 255 .all(|target| target.required()) 256 ); 257 258 let replay = host 259 .repository() 260 .commit_discovery_desired_state(&request) 261 .await 262 .expect("exact replay"); 263 assert!(matches!( 264 replay, 265 MycDiscoveryCommitAdmission::ExactReplay(_) 266 )); 267 let conflicting = 268 MycDiscoveryCommitRequest::new(&metadata, &signed_handler_event(1_725_000_001), time(110)) 269 .expect("conflicting discovery request"); 270 assert_eq!( 271 host.repository() 272 .commit_discovery_desired_state(&conflicting) 273 .await 274 .expect_err("conflicting job") 275 .kind(), 276 MycStateRepositoryErrorKind::Binding 277 ); 278 279 let primary = MycDeliveryRelayId::new("primary").expect("primary"); 280 let left_repository = host.repository(); 281 let right_repository = host.repository(); 282 let (left, right) = tokio::join!( 283 left_repository.claim_delivery_target( 284 job_id, 285 &primary, 286 MycDeliveryAttemptNonce::from_injected_entropy([0x51; 32]), 287 time(120), 288 ), 289 right_repository.claim_delivery_target( 290 job_id, 291 &primary, 292 MycDeliveryAttemptNonce::from_injected_entropy([0x51; 32]), 293 time(120), 294 ) 295 ); 296 let (claimed, replayed) = match (left.expect("left claim"), right.expect("right claim")) { 297 (MycDeliveryClaim::Claimed(claimed), MycDeliveryClaim::ExactReplay(replayed)) 298 | (MycDeliveryClaim::ExactReplay(replayed), MycDeliveryClaim::Claimed(claimed)) => { 299 (claimed, replayed) 300 } 301 outcome => panic!("unexpected concurrent outcome: {outcome:?}"), 302 }; 303 assert_eq!(claimed.id(), replayed.id()); 304 assert_eq!(claimed.number(), 1); 305 assert_eq!( 306 host.repository() 307 .record_delivery_attempt_outcome( 308 job_id, 309 &primary, 310 claimed.id(), 311 MycDeliveryAttemptOutcome::Delivered, 312 jitter(0), 313 time(121), 314 ) 315 .await 316 .expect_err("delivery before submission") 317 .kind(), 318 MycStateRepositoryErrorKind::Binding 319 ); 320 let submitted = host 321 .repository() 322 .mark_delivery_attempt_submitted(job_id, &primary, claimed.id(), time(121)) 323 .await 324 .expect("submitted"); 325 assert_eq!(submitted.status(), MycDeliveryAttemptStatus::Submitted); 326 let active = host 327 .repository() 328 .record_delivery_attempt_outcome( 329 job_id, 330 &primary, 331 claimed.id(), 332 MycDeliveryAttemptOutcome::Delivered, 333 jitter(0), 334 time(122), 335 ) 336 .await 337 .expect("primary delivered"); 338 assert_eq!(active.status(), MycDeliveryJobStatus::Active); 339 340 let secondary = MycDeliveryRelayId::new("secondary").expect("secondary"); 341 let first = match host 342 .repository() 343 .claim_delivery_target( 344 job_id, 345 &secondary, 346 MycDeliveryAttemptNonce::from_injected_entropy([0x61; 32]), 347 time(123), 348 ) 349 .await 350 .expect("secondary claim") 351 { 352 MycDeliveryClaim::Claimed(attempt) => attempt, 353 other => panic!("unexpected claim: {other:?}"), 354 }; 355 host.repository() 356 .mark_delivery_attempt_submitted(job_id, &secondary, first.id(), time(124)) 357 .await 358 .expect("secondary submitted"); 359 assert_eq!( 360 host.repository() 361 .record_delivery_attempt_outcome( 362 job_id, 363 &secondary, 364 first.id(), 365 MycDeliveryAttemptOutcome::UnknownAcknowledgement, 366 jitter(251), 367 time(125), 368 ) 369 .await 370 .expect_err("jitter exceeds first retry cap") 371 .kind(), 372 MycStateRepositoryErrorKind::Binding 373 ); 374 let unknown = host 375 .repository() 376 .record_delivery_attempt_outcome( 377 job_id, 378 &secondary, 379 first.id(), 380 MycDeliveryAttemptOutcome::UnknownAcknowledgement, 381 jitter(250), 382 time(125), 383 ) 384 .await 385 .expect("unknown acknowledgement"); 386 assert_eq!(unknown.status(), MycDeliveryJobStatus::Active); 387 assert_eq!( 388 unknown.targets()[1].status(), 389 MycDeliveryTargetStatus::Unknown 390 ); 391 assert_eq!(unknown.targets()[1].next_attempt_at(), Some(time(375))); 392 assert!(matches!( 393 host.repository() 394 .claim_delivery_target( 395 job_id, 396 &secondary, 397 MycDeliveryAttemptNonce::from_injected_entropy([0x62; 32]), 398 time(374), 399 ) 400 .await 401 .expect("not ready"), 402 MycDeliveryClaim::NotReady 403 )); 404 let second = match host 405 .repository() 406 .claim_delivery_target( 407 job_id, 408 &secondary, 409 MycDeliveryAttemptNonce::from_injected_entropy([0x62; 32]), 410 time(375), 411 ) 412 .await 413 .expect("retry claim") 414 { 415 MycDeliveryClaim::Claimed(attempt) => attempt, 416 other => panic!("unexpected retry: {other:?}"), 417 }; 418 let retry = host 419 .repository() 420 .recover_expired_delivery_lease(job_id, &secondary, second.id(), jitter(500), time(15_376)) 421 .await 422 .expect("expired pre-submit lease"); 423 assert_eq!( 424 retry.targets()[1].status(), 425 MycDeliveryTargetStatus::Retryable 426 ); 427 assert_eq!(retry.targets()[1].next_attempt_at(), Some(time(15_876))); 428 let third = match host 429 .repository() 430 .claim_delivery_target( 431 job_id, 432 &secondary, 433 MycDeliveryAttemptNonce::from_injected_entropy([0x63; 32]), 434 time(15_876), 435 ) 436 .await 437 .expect("third claim") 438 { 439 MycDeliveryClaim::Claimed(attempt) => attempt, 440 other => panic!("unexpected third claim: {other:?}"), 441 }; 442 host.repository() 443 .mark_delivery_attempt_submitted(job_id, &secondary, third.id(), time(15_877)) 444 .await 445 .expect("third submitted"); 446 let delivered = host 447 .repository() 448 .record_delivery_attempt_outcome( 449 job_id, 450 &secondary, 451 third.id(), 452 MycDeliveryAttemptOutcome::Delivered, 453 jitter(0), 454 time(15_878), 455 ) 456 .await 457 .expect("job delivered"); 458 assert_eq!(delivered.status(), MycDeliveryJobStatus::Delivered); 459 assert_eq!(delivered.finalized_at(), Some(time(15_878))); 460 let attempts = host 461 .repository() 462 .read_delivery_attempts(job_id, &secondary) 463 .await 464 .expect("attempt history"); 465 assert_eq!(attempts.len(), 3); 466 assert_eq!(attempts[0].status(), MycDeliveryAttemptStatus::Unknown); 467 assert_eq!(attempts[0].reason(), Some("acknowledgement_lost")); 468 assert_eq!(attempts[1].status(), MycDeliveryAttemptStatus::Failed); 469 assert_eq!(attempts[1].reason(), Some("lease_expired_before_submit")); 470 assert_eq!(attempts[2].status(), MycDeliveryAttemptStatus::Delivered); 471 host.close().await.expect("first close"); 472 473 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 474 .await 475 .expect("reopen"); 476 assert_eq!( 477 host.repository() 478 .read_delivery_job(job_id) 479 .await 480 .expect("read after reopen") 481 .expect("job") 482 .status(), 483 MycDeliveryJobStatus::Delivered 484 ); 485 host.close().await.expect("final close"); 486 } 487 488 #[tokio::test] 489 async fn terminal_unknown_is_not_relabelled_as_failure_and_sql_guards_preserve_evidence() { 490 let directory = tempfile::tempdir().expect("temporary root"); 491 let runtime = runtime(directory.path()); 492 prepare_state_directory(&runtime); 493 let source = String::from_utf8(config_source(CONFIG_EXAMPLE)) 494 .expect("UTF-8 config") 495 .replace("max_attempts = 5", "max_attempts = 1"); 496 let metadata = metadata(&runtime, source.as_bytes()); 497 let (applied_at, build) = migration_evidence(); 498 initialize_myc_state(&runtime, &metadata, applied_at, &build) 499 .await 500 .expect("initialization"); 501 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 502 .await 503 .expect("writer"); 504 let job = host 505 .repository() 506 .commit_discovery_desired_state( 507 &MycDiscoveryCommitRequest::new( 508 &metadata, 509 &signed_handler_event(1_725_000_010), 510 time(110), 511 ) 512 .expect("discovery request"), 513 ) 514 .await 515 .expect("job"); 516 let job_id = job.record().job().id(); 517 let primary = MycDeliveryRelayId::new("primary").expect("primary"); 518 let attempt = match host 519 .repository() 520 .claim_delivery_target( 521 job_id, 522 &primary, 523 MycDeliveryAttemptNonce::from_injected_entropy([0x73; 32]), 524 time(120), 525 ) 526 .await 527 .expect("claim") 528 { 529 MycDeliveryClaim::Claimed(attempt) => attempt, 530 other => panic!("unexpected claim: {other:?}"), 531 }; 532 host.repository() 533 .mark_delivery_attempt_submitted(job_id, &primary, attempt.id(), time(121)) 534 .await 535 .expect("submit"); 536 let job = host 537 .repository() 538 .record_delivery_attempt_outcome( 539 job_id, 540 &primary, 541 attempt.id(), 542 MycDeliveryAttemptOutcome::UnknownAcknowledgement, 543 jitter(0), 544 time(122), 545 ) 546 .await 547 .expect("unknown"); 548 assert_eq!(job.status(), MycDeliveryJobStatus::Unknown); 549 assert_eq!(job.targets()[0].status(), MycDeliveryTargetStatus::Unknown); 550 assert_eq!(job.targets()[0].next_attempt_at(), None); 551 assert_eq!(job.finalized_at(), Some(time(122))); 552 host.close().await.expect("close"); 553 554 let options = SqliteConnectOptions::new() 555 .filename(runtime.artifacts().state_database()) 556 .create_if_missing(false) 557 .disable_statement_logging(); 558 let mut connection = sqlx::SqliteConnection::connect_with(&options) 559 .await 560 .expect("inspection connection"); 561 for statement in [ 562 "DELETE FROM delivery_attempts", 563 "DELETE FROM delivery_targets", 564 "DELETE FROM delivery_jobs", 565 "UPDATE delivery_targets SET relay_id = 'changed'", 566 "UPDATE delivery_jobs SET artifact_sha256 = zeroblob(32)", 567 ] { 568 assert!( 569 sqlx::query(statement) 570 .execute(&mut connection) 571 .await 572 .is_err(), 573 "guard accepted {statement}" 574 ); 575 } 576 assert_eq!( 577 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM delivery_attempts") 578 .fetch_one(&mut connection) 579 .await 580 .expect("attempt count"), 581 1 582 ); 583 connection.close().await.expect("connection close"); 584 } 585 586 #[test] 587 fn delivery_boundary_is_typed_sqlx_only_and_has_no_external_wait_or_legacy_authority() { 588 assert!(LIB_SOURCE.contains("mod state_delivery;")); 589 assert!(!LIB_SOURCE.contains("pub mod state_delivery;")); 590 assert!(DELIVERY_SOURCE.contains("ServiceSqliteTransaction<'_>")); 591 assert!(DELIVERY_SOURCE.contains("delivery_attempts")); 592 assert!(DELIVERY_SOURCE.contains("UnknownAcknowledgement")); 593 assert!(CATALOG_SOURCE.contains("delivery_jobs_no_delete")); 594 assert!(CATALOG_SOURCE.contains("delivery_targets_no_delete")); 595 assert!(CATALOG_SOURCE.contains("delivery_attempts_no_delete")); 596 for forbidden in [ 597 "SqlitePool", 598 "SqliteConnection", 599 "BEGIN ", 600 "COMMIT", 601 "ROLLBACK", 602 "relay_url", 603 "reqwest", 604 "tokio::spawn", 605 "std::env", 606 "std::time", 607 "outbox_sqlite", 608 ] { 609 assert!( 610 !DELIVERY_SOURCE.contains(forbidden), 611 "found forbidden delivery authority `{forbidden}`" 612 ); 613 } 614 }