delivery_worker.rs (21744B)
1 //! Durable delivery orchestration over exact committed event bytes. 2 3 use core::fmt; 4 use std::error::Error; 5 6 use sha2::{Digest as _, Sha256}; 7 8 use crate::transport_nostr_adapter::{ 9 MycNostrDeliveryAdapter, MycRelayAdapter, MycRelayExecutionOutcome, 10 }; 11 use crate::{ 12 MycConfigDocumentV1, MycDeliveryAttemptNonce, MycDeliveryAttemptOutcome, 13 MycDeliveryAttemptRecord, MycDeliveryClaim, MycDeliveryJobId, MycDeliveryJobRecord, 14 MycDeliveryRelayId, MycDeliveryRetryJitter, MycDeliverySourceKind, MycDeliveryTimeUnixMs, 15 MycStateRepository, MycTaskCancellation, 16 }; 17 18 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 19 pub(crate) enum MycDeliveryWorkerErrorKind { 20 Configuration, 21 State, 22 Artifact, 23 } 24 25 pub(crate) struct MycDeliveryWorkerError { 26 kind: MycDeliveryWorkerErrorKind, 27 } 28 29 impl MycDeliveryWorkerError { 30 #[cfg(test)] 31 pub(crate) const fn kind(&self) -> MycDeliveryWorkerErrorKind { 32 self.kind 33 } 34 } 35 36 impl fmt::Debug for MycDeliveryWorkerError { 37 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 38 formatter 39 .debug_struct("MycDeliveryWorkerError") 40 .field("kind", &self.kind) 41 .finish() 42 } 43 } 44 45 impl fmt::Display for MycDeliveryWorkerError { 46 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 47 formatter.write_str("Myc delivery worker failed") 48 } 49 } 50 51 impl Error for MycDeliveryWorkerError {} 52 53 const fn worker_error(kind: MycDeliveryWorkerErrorKind) -> MycDeliveryWorkerError { 54 MycDeliveryWorkerError { kind } 55 } 56 57 #[derive(Clone, PartialEq, Eq)] 58 pub(crate) enum MycDeliveryWorkerResult { 59 Completed(MycDeliveryJobRecord), 60 ExactReplay, 61 NotReady, 62 Terminal, 63 } 64 65 impl fmt::Debug for MycDeliveryWorkerResult { 66 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 67 formatter.write_str(match self { 68 Self::Completed(_) => "MycDeliveryWorkerResult::Completed([redacted])", 69 Self::ExactReplay => "MycDeliveryWorkerResult::ExactReplay", 70 Self::NotReady => "MycDeliveryWorkerResult::NotReady", 71 Self::Terminal => "MycDeliveryWorkerResult::Terminal", 72 }) 73 } 74 } 75 76 pub(crate) struct MycDeliveryExecutionEvidence { 77 pub(crate) claimed_at: MycDeliveryTimeUnixMs, 78 pub(crate) submitted_at: MycDeliveryTimeUnixMs, 79 pub(crate) observed_at: MycDeliveryTimeUnixMs, 80 pub(crate) retry_entropy: [u8; 8], 81 } 82 83 impl fmt::Debug for MycDeliveryExecutionEvidence { 84 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 85 formatter.write_str("MycDeliveryExecutionEvidence([sealed])") 86 } 87 } 88 89 pub(crate) struct MycDeliveryWorker { 90 adapter: MycNostrDeliveryAdapter, 91 } 92 93 impl MycDeliveryWorker { 94 pub(crate) fn from_configuration( 95 configuration: &MycConfigDocumentV1, 96 ) -> Result<Self, MycDeliveryWorkerError> { 97 let adapter = MycNostrDeliveryAdapter::from_configuration(configuration) 98 .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::Configuration))?; 99 Ok(Self { adapter }) 100 } 101 102 pub(crate) async fn run_one( 103 &self, 104 repository: &MycStateRepository<'_>, 105 job_id: MycDeliveryJobId, 106 relay_id: &MycDeliveryRelayId, 107 nonce: MycDeliveryAttemptNonce, 108 evidence: MycDeliveryExecutionEvidence, 109 cancellation: &MycTaskCancellation, 110 ) -> Result<MycDeliveryWorkerResult, MycDeliveryWorkerError> { 111 run_with_adapter( 112 &self.adapter, 113 repository, 114 job_id, 115 relay_id, 116 nonce, 117 evidence, 118 cancellation, 119 ) 120 .await 121 } 122 } 123 124 #[allow(clippy::too_many_arguments)] 125 async fn run_with_adapter<A: MycRelayAdapter>( 126 adapter: &A, 127 repository: &MycStateRepository<'_>, 128 job_id: MycDeliveryJobId, 129 relay_id: &MycDeliveryRelayId, 130 nonce: MycDeliveryAttemptNonce, 131 evidence: MycDeliveryExecutionEvidence, 132 cancellation: &MycTaskCancellation, 133 ) -> Result<MycDeliveryWorkerResult, MycDeliveryWorkerError> { 134 if cancellation.is_cancelled() { 135 return Ok(MycDeliveryWorkerResult::NotReady); 136 } 137 let attempt = match repository 138 .claim_delivery_target(job_id, relay_id, nonce, evidence.claimed_at) 139 .await 140 .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State))? 141 { 142 MycDeliveryClaim::Claimed(attempt) => attempt, 143 MycDeliveryClaim::ExactReplay(_) => return Ok(MycDeliveryWorkerResult::ExactReplay), 144 MycDeliveryClaim::NotReady => return Ok(MycDeliveryWorkerResult::NotReady), 145 MycDeliveryClaim::Terminal => return Ok(MycDeliveryWorkerResult::Terminal), 146 }; 147 let job = repository 148 .read_delivery_job(job_id) 149 .await 150 .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State))? 151 .ok_or_else(|| worker_error(MycDeliveryWorkerErrorKind::Artifact))?; 152 let exact_event_bytes = read_exact_event(repository, &job).await?; 153 let request_id = request_id(job_id, attempt.id()); 154 let prepared = match adapter.prepare( 155 relay_id, 156 request_id, 157 &exact_event_bytes, 158 attempt.lease_expires_at().get(), 159 ) { 160 Ok(prepared) => prepared, 161 Err(_) => { 162 return persist_outcome( 163 repository, 164 relay_id, 165 MycDeliveryAttemptOutcome::TransportFailed, 166 &job, 167 &attempt, 168 evidence, 169 ) 170 .await; 171 } 172 }; 173 if cancellation.is_cancelled() { 174 return persist_outcome( 175 repository, 176 relay_id, 177 MycDeliveryAttemptOutcome::TransportFailed, 178 &job, 179 &attempt, 180 evidence, 181 ) 182 .await; 183 } 184 repository 185 .mark_delivery_attempt_submitted(job_id, relay_id, attempt.id(), evidence.submitted_at) 186 .await 187 .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State))?; 188 let outcome = tokio::select! { 189 result = adapter.execute(prepared) => result.unwrap_or(MycRelayExecutionOutcome::UnknownAcknowledgement), 190 () = cancellation.cancelled() => MycRelayExecutionOutcome::UnknownAcknowledgement, 191 }; 192 let outcome = match outcome { 193 MycRelayExecutionOutcome::Accepted => MycDeliveryAttemptOutcome::Delivered, 194 MycRelayExecutionOutcome::Rejected => MycDeliveryAttemptOutcome::RelayRejected, 195 MycRelayExecutionOutcome::TransportFailed => MycDeliveryAttemptOutcome::TransportFailed, 196 MycRelayExecutionOutcome::UnknownAcknowledgement => { 197 MycDeliveryAttemptOutcome::UnknownAcknowledgement 198 } 199 }; 200 persist_outcome(repository, relay_id, outcome, &job, &attempt, evidence).await 201 } 202 203 async fn read_exact_event( 204 repository: &MycStateRepository<'_>, 205 job: &MycDeliveryJobRecord, 206 ) -> Result<Box<[u8]>, MycDeliveryWorkerError> { 207 let bytes: Box<[u8]> = match job.source_kind() { 208 MycDeliverySourceKind::SignerResponse => repository 209 .read_nip46_response(job.id()) 210 .await 211 .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State))? 212 .filter(|record| record.delivery_job().id() == job.id()) 213 .map(|record| Box::from(record.signed_response_bytes())) 214 .ok_or_else(|| worker_error(MycDeliveryWorkerErrorKind::Artifact))?, 215 MycDeliverySourceKind::DiscoveryHandler => repository 216 .read_discovery_document_for_job(job.id()) 217 .await 218 .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State))? 219 .map(|record| Box::from(record.event_bytes())) 220 .ok_or_else(|| worker_error(MycDeliveryWorkerErrorKind::Artifact))?, 221 }; 222 let digest: [u8; 32] = Sha256::digest(&bytes).into(); 223 if &digest != job.artifact_digest().as_bytes() { 224 return Err(worker_error(MycDeliveryWorkerErrorKind::Artifact)); 225 } 226 Ok(bytes) 227 } 228 229 async fn persist_outcome( 230 repository: &MycStateRepository<'_>, 231 relay_id: &MycDeliveryRelayId, 232 outcome: MycDeliveryAttemptOutcome, 233 job: &MycDeliveryJobRecord, 234 attempt: &MycDeliveryAttemptRecord, 235 evidence: MycDeliveryExecutionEvidence, 236 ) -> Result<MycDeliveryWorkerResult, MycDeliveryWorkerError> { 237 let retry_jitter = 238 delivery_retry_jitter(job, attempt.number(), outcome, evidence.retry_entropy)?; 239 repository 240 .record_delivery_attempt_outcome( 241 job.id(), 242 relay_id, 243 attempt.id(), 244 outcome, 245 retry_jitter, 246 evidence.observed_at, 247 ) 248 .await 249 .map(MycDeliveryWorkerResult::Completed) 250 .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State)) 251 } 252 253 fn delivery_retry_jitter( 254 job: &MycDeliveryJobRecord, 255 attempt_number: u32, 256 outcome: MycDeliveryAttemptOutcome, 257 entropy: [u8; 8], 258 ) -> Result<MycDeliveryRetryJitter, MycDeliveryWorkerError> { 259 if outcome == MycDeliveryAttemptOutcome::Delivered || attempt_number >= job.max_attempts() { 260 return MycDeliveryRetryJitter::new(0) 261 .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State)); 262 } 263 let exponent = attempt_number.saturating_sub(1).min(31); 264 let factor = 1_u64.checked_shl(exponent).unwrap_or(u64::MAX); 265 let maximum_delay = job 266 .initial_backoff_ms() 267 .saturating_mul(factor) 268 .min(job.maximum_backoff_ms()); 269 let jitter = u64::from_be_bytes(entropy) % (maximum_delay + 1); 270 MycDeliveryRetryJitter::new(jitter).map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State)) 271 } 272 273 fn request_id(job_id: MycDeliveryJobId, attempt_id: crate::MycDeliveryAttemptId) -> String { 274 let mut request = hex::encode(job_id.as_bytes()); 275 request.push(':'); 276 request.push_str(&hex::encode(attempt_id.as_bytes())); 277 request 278 } 279 280 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] 281 mod tests { 282 use std::{ 283 fs, 284 os::unix::fs::PermissionsExt as _, 285 sync::{Arc, Mutex}, 286 }; 287 288 use super::*; 289 use crate::nip46_wave_080_a::{ 290 OBSERVED_AT_SECONDS, RECEIVED_AT_MS, configuration, connection_time, metadata, 291 migration_evidence, runtime, 292 }; 293 use crate::nip46_wave_080_b::{active_connection, atomic_response_request}; 294 use crate::{ 295 MycDeliveryAttemptStatus, MycDeliveryTargetStatus, MycNip46CommitRequest, 296 MycTaskCancellation, initialize_myc_state, open_myc_state_read_write, 297 }; 298 use tokio::sync::Notify; 299 300 struct FakeAdapter { 301 entered: Arc<Notify>, 302 release: Arc<Notify>, 303 prepared_bytes: Arc<Mutex<Vec<u8>>>, 304 outcome: MycRelayExecutionOutcome, 305 } 306 307 impl MycRelayAdapter for FakeAdapter { 308 type Prepared = (); 309 310 fn prepare( 311 &self, 312 _relay_id: &MycDeliveryRelayId, 313 request_id: String, 314 exact_event_bytes: &[u8], 315 deadline_unix_ms: u64, 316 ) -> Result<Self::Prepared, crate::transport_nostr_adapter::MycRelayAdapterError> { 317 assert_eq!(request_id.len(), 129); 318 assert!(deadline_unix_ms > RECEIVED_AT_MS); 319 *self.prepared_bytes.lock().expect("prepared bytes") = exact_event_bytes.to_vec(); 320 Ok(()) 321 } 322 323 fn execute<'a>( 324 &'a self, 325 (): Self::Prepared, 326 ) -> crate::transport_nostr_adapter::RelayExecutionFuture<'a> { 327 Box::pin(async move { 328 self.entered.notify_one(); 329 self.release.notified().await; 330 Ok(self.outcome) 331 }) 332 } 333 } 334 335 struct NoIoAdapter; 336 337 impl MycRelayAdapter for NoIoAdapter { 338 type Prepared = (); 339 340 fn prepare( 341 &self, 342 _relay_id: &MycDeliveryRelayId, 343 _request_id: String, 344 _exact_event_bytes: &[u8], 345 _deadline_unix_ms: u64, 346 ) -> Result<Self::Prepared, crate::transport_nostr_adapter::MycRelayAdapterError> { 347 panic!("an exact replay must not prepare a second relay operation") 348 } 349 350 fn execute<'a>( 351 &'a self, 352 (): Self::Prepared, 353 ) -> crate::transport_nostr_adapter::RelayExecutionFuture<'a> { 354 panic!("an exact replay must not execute a second relay operation") 355 } 356 } 357 358 async fn committed_response_host( 359 root: &std::path::Path, 360 ) -> (crate::MycStateHost, MycDeliveryJobId, Vec<u8>) { 361 let runtime = runtime(root); 362 fs::create_dir_all(runtime.context().paths().state()).expect("state directory"); 363 fs::set_permissions( 364 runtime.context().paths().state(), 365 fs::Permissions::from_mode(0o700), 366 ) 367 .expect("state mode"); 368 let metadata = metadata(&runtime); 369 let config = configuration(); 370 let (applied_at, build) = migration_evidence(); 371 initialize_myc_state(&runtime, &metadata, applied_at, &build) 372 .await 373 .expect("initialization"); 374 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 375 .await 376 .expect("state host"); 377 let repository = host.repository(); 378 let (work, _active, decision) = active_connection(&repository, &config).await; 379 let completion = MycNip46CommitRequest::new( 380 &work, 381 Some(&decision), 382 None, 383 connection_time(RECEIVED_AT_MS + 2_001), 384 ) 385 .expect("completion"); 386 let (response, exact_bytes) = atomic_response_request(&config, &completion); 387 let committed = repository 388 .commit_nip46_response(&response) 389 .await 390 .expect("committed response"); 391 ( 392 host, 393 committed.record().response().delivery_job().id(), 394 exact_bytes, 395 ) 396 } 397 398 fn evidence() -> MycDeliveryExecutionEvidence { 399 MycDeliveryExecutionEvidence { 400 claimed_at: MycDeliveryTimeUnixMs::new(RECEIVED_AT_MS + 4_000).expect("claim time"), 401 submitted_at: MycDeliveryTimeUnixMs::new(RECEIVED_AT_MS + 4_001) 402 .expect("submitted time"), 403 observed_at: MycDeliveryTimeUnixMs::new(RECEIVED_AT_MS + 4_002).expect("observed time"), 404 retry_entropy: [0; 8], 405 } 406 } 407 408 #[tokio::test] 409 async fn injected_entropy_uses_the_exact_full_jitter_cap() { 410 let root = tempfile::tempdir().expect("temporary root"); 411 let (host, job_id, _) = committed_response_host(root.path()).await; 412 let job = host 413 .repository() 414 .read_delivery_job(job_id) 415 .await 416 .expect("job read") 417 .expect("job"); 418 let first = delivery_retry_jitter( 419 &job, 420 1, 421 MycDeliveryAttemptOutcome::TransportFailed, 422 u64::MAX.to_be_bytes(), 423 ) 424 .expect("first-attempt jitter"); 425 assert_eq!(first.get(), u64::MAX % (job.initial_backoff_ms() + 1)); 426 assert!(first.get() <= job.initial_backoff_ms()); 427 assert_eq!( 428 delivery_retry_jitter( 429 &job, 430 job.max_attempts(), 431 MycDeliveryAttemptOutcome::TransportFailed, 432 u64::MAX.to_be_bytes(), 433 ) 434 .expect("exhausted jitter") 435 .get(), 436 0 437 ); 438 assert_eq!( 439 delivery_retry_jitter( 440 &job, 441 1, 442 MycDeliveryAttemptOutcome::Delivered, 443 u64::MAX.to_be_bytes(), 444 ) 445 .expect("terminal jitter") 446 .get(), 447 0 448 ); 449 host.close().await.expect("close"); 450 } 451 452 #[tokio::test] 453 async fn exact_bytes_are_submitted_only_after_durable_submitted_state() { 454 let root = tempfile::tempdir().expect("temporary root"); 455 let (host, job_id, exact_bytes) = committed_response_host(root.path()).await; 456 let repository = host.repository(); 457 let relay = MycDeliveryRelayId::new("primary").expect("relay"); 458 let entered = Arc::new(Notify::new()); 459 let release = Arc::new(Notify::new()); 460 let prepared_bytes = Arc::new(Mutex::new(Vec::new())); 461 let adapter = FakeAdapter { 462 entered: Arc::clone(&entered), 463 release: Arc::clone(&release), 464 prepared_bytes: Arc::clone(&prepared_bytes), 465 outcome: MycRelayExecutionOutcome::Accepted, 466 }; 467 let (cancellation, _cancel) = MycTaskCancellation::test_pair(); 468 let run = run_with_adapter( 469 &adapter, 470 &repository, 471 job_id, 472 &relay, 473 MycDeliveryAttemptNonce::from_injected_entropy([0xa1; 32]), 474 evidence(), 475 &cancellation, 476 ); 477 let observe = async { 478 entered.notified().await; 479 let attempts = repository 480 .read_delivery_attempts(job_id, &relay) 481 .await 482 .expect("submitted attempt"); 483 assert_eq!(attempts.len(), 1); 484 assert_eq!(attempts[0].status(), MycDeliveryAttemptStatus::Submitted); 485 release.notify_one(); 486 }; 487 let (result, ()) = tokio::join!(run, observe); 488 let MycDeliveryWorkerResult::Completed(job) = result.expect("delivery") else { 489 panic!("delivery must complete"); 490 }; 491 assert_eq!( 492 job.targets()[0].status(), 493 MycDeliveryTargetStatus::Delivered 494 ); 495 assert_eq!( 496 prepared_bytes.lock().expect("prepared bytes").as_slice(), 497 exact_bytes 498 ); 499 host.close().await.expect("close"); 500 } 501 502 #[tokio::test] 503 async fn cancellation_after_submission_is_durably_unknown() { 504 let root = tempfile::tempdir().expect("temporary root"); 505 let (host, job_id, _) = committed_response_host(root.path()).await; 506 let repository = host.repository(); 507 let relay = MycDeliveryRelayId::new("primary").expect("relay"); 508 let entered = Arc::new(Notify::new()); 509 let release = Arc::new(Notify::new()); 510 let adapter = FakeAdapter { 511 entered: Arc::clone(&entered), 512 release, 513 prepared_bytes: Arc::new(Mutex::new(Vec::new())), 514 outcome: MycRelayExecutionOutcome::Accepted, 515 }; 516 let (cancellation, cancel) = MycTaskCancellation::test_pair(); 517 let run = run_with_adapter( 518 &adapter, 519 &repository, 520 job_id, 521 &relay, 522 MycDeliveryAttemptNonce::from_injected_entropy([0xa2; 32]), 523 evidence(), 524 &cancellation, 525 ); 526 let cancel_after_submit = async { 527 entered.notified().await; 528 cancel.cancel(); 529 }; 530 let (result, ()) = tokio::join!(run, cancel_after_submit); 531 let MycDeliveryWorkerResult::Completed(job) = result.expect("unknown delivery") else { 532 panic!("delivery must resolve"); 533 }; 534 assert_eq!(job.targets()[0].status(), MycDeliveryTargetStatus::Unknown); 535 let attempts = repository 536 .read_delivery_attempts(job_id, &relay) 537 .await 538 .expect("attempts"); 539 assert_eq!(attempts[0].status(), MycDeliveryAttemptStatus::Unknown); 540 host.close().await.expect("close"); 541 } 542 543 #[tokio::test] 544 async fn concurrent_exact_replay_performs_no_second_external_operation() { 545 let root = tempfile::tempdir().expect("temporary root"); 546 let (host, job_id, _) = committed_response_host(root.path()).await; 547 let repository = host.repository(); 548 let relay = MycDeliveryRelayId::new("primary").expect("relay"); 549 let entered = Arc::new(Notify::new()); 550 let release = Arc::new(Notify::new()); 551 let adapter = FakeAdapter { 552 entered: Arc::clone(&entered), 553 release: Arc::clone(&release), 554 prepared_bytes: Arc::new(Mutex::new(Vec::new())), 555 outcome: MycRelayExecutionOutcome::Accepted, 556 }; 557 let nonce = MycDeliveryAttemptNonce::from_injected_entropy([0xa3; 32]); 558 let replay_nonce = MycDeliveryAttemptNonce::from_injected_entropy([0xa3; 32]); 559 let (cancellation, _cancel) = MycTaskCancellation::test_pair(); 560 let run = run_with_adapter( 561 &adapter, 562 &repository, 563 job_id, 564 &relay, 565 nonce, 566 evidence(), 567 &cancellation, 568 ); 569 let replay = async { 570 entered.notified().await; 571 let result = run_with_adapter( 572 &NoIoAdapter, 573 &repository, 574 job_id, 575 &relay, 576 replay_nonce, 577 evidence(), 578 &cancellation, 579 ) 580 .await 581 .expect("exact replay"); 582 assert_eq!(result, MycDeliveryWorkerResult::ExactReplay); 583 release.notify_one(); 584 }; 585 let (result, ()) = tokio::join!(run, replay); 586 assert!(matches!( 587 result.expect("delivery"), 588 MycDeliveryWorkerResult::Completed(_) 589 )); 590 host.close().await.expect("close"); 591 } 592 593 #[test] 594 fn worker_errors_are_source_free_and_redacted() { 595 for kind in [ 596 MycDeliveryWorkerErrorKind::Configuration, 597 MycDeliveryWorkerErrorKind::State, 598 MycDeliveryWorkerErrorKind::Artifact, 599 ] { 600 let error = worker_error(kind); 601 assert_eq!(error.kind(), kind); 602 assert!(Error::source(&error).is_none()); 603 assert!(!format!("{error} {error:?}").contains("secret")); 604 } 605 assert_eq!(OBSERVED_AT_SECONDS, 1_725_000_000); 606 } 607 }