push.rs (60862B)
1 //! Signing, durable enqueue, delivery, and satisfaction orchestration. 2 3 use core::num::{NonZeroU32, NonZeroU64}; 4 use radroots_event::admission::RawEvent; 5 use radroots_event_codec::{ 6 authoring::AuthoredEventPlan, 7 verify::{self, Nip01SignatureVerifier}, 8 }; 9 use radroots_protocol::runtime::v1::OperationId; 10 use radroots_signing::{ 11 Actor, AuthoredArtifactId as SigningArtifactId, SigningIntentId, SigningOperationId, 12 recovery::{RecoveryDisposition, ReplayCapability, recovery_disposition}, 13 request::{CancellationPolicy, SignPolicy}, 14 }; 15 use radroots_storage::{ 16 atomic::{AtomicCommitDigest, AtomicCommitDisposition}, 17 authored::{ 18 AdmissionState, AuthoredArtifact, AuthoredArtifactId, AuthoredOperation, FailureClass, 19 OperationSettlement, RetrySchedule, SigningState, WorkClaim, WorkFailure, WorkPhase, 20 }, 21 authored_atomic::{ 22 ApplyAdmissionResult, ApplyWorkFailure, AuthoredAtomicCommand, AuthoredAtomicOutcome, 23 AuthoredWorkTarget, CancelAuthoredTarget, CancelAuthoredWork, ClaimAuthoredTarget, 24 ClaimAuthoredWork, PrepareAuthoredOperation, RecordSignedArtifact, WorkFence, 25 }, 26 authored_delivery::{ 27 AuthoredDeliveryHistory, AuthoredDeliveryIntent, AuthoredDeliveryPlan, 28 AuthoredDeliveryPlanId, AuthoredDeliveryState, DELIVERY_PLAN_ATTEMPTS_MAX, 29 DeliveryAttemptOutcome, 30 }, 31 event::{AdmissionDisposition, EventAdmission, EventStore}, 32 journal::{IdempotencyKey, OperationInstanceId}, 33 }; 34 use radroots_transport::{ 35 DeliveryRequest, SinkFailure, Target, TransportId, 36 outcome::Retryability, 37 policy::{SatisfactionPolicy, SatisfactionState}, 38 source::{EventProvenance, ObservedEvent}, 39 target::TargetSet, 40 }; 41 use sha2::{Digest, Sha256}; 42 43 use crate::{ 44 Engine, 45 ingest::{AdmissionDecision, AdmissionPolicy, RegistryPolicy}, 46 policy::{Error, OperationKind, SyncId}, 47 }; 48 49 mod delivery; 50 51 const SIGNING_CLAIM_OWNER_EXACT: &str = "radroots-sync-signing-exact"; 52 const SIGNING_CLAIM_OWNER_LOCAL: &str = "radroots-sync-signing-local"; 53 const SIGNING_CLAIM_OWNER_NON_REPLAYABLE: &str = "radroots-sync-signing-non-replayable"; 54 const ADMISSION_CLAIM_OWNER: &str = "radroots-sync-admission"; 55 const DELIVERY_CLAIM_OWNER: &str = "radroots-sync-delivery"; 56 const WORK_RETRY_DELAY_MS: u64 = 1_000; 57 58 /// Caller-owned, replay-stable inputs for one outbound operation. 59 #[derive(Clone)] 60 pub struct PushRequest { 61 operation_id: SyncId, 62 idempotency_key: IdempotencyKey, 63 actor: Actor, 64 plan: AuthoredEventPlan, 65 targets: TargetSet, 66 satisfaction: SatisfactionPolicy, 67 delivery_deadline_unix_ms: u64, 68 cancellation: CancellationPolicy, 69 } 70 71 impl PushRequest { 72 #[allow(clippy::too_many_arguments)] 73 pub fn new( 74 operation_id: SyncId, 75 idempotency_key: IdempotencyKey, 76 actor: Actor, 77 plan: AuthoredEventPlan, 78 targets: TargetSet, 79 satisfaction: SatisfactionPolicy, 80 delivery_deadline_unix_ms: u64, 81 cancellation: CancellationPolicy, 82 ) -> Result<Self, Error> { 83 if satisfaction.validate_for(&targets).is_err() 84 || delivery_deadline_unix_ms == 0 85 || actor.public_key() != *plan.author() 86 { 87 return Err(Error::InvalidPushRequest); 88 } 89 Ok(Self { 90 operation_id, 91 idempotency_key, 92 actor, 93 plan, 94 targets, 95 satisfaction, 96 delivery_deadline_unix_ms, 97 cancellation, 98 }) 99 } 100 101 /// Builds the existing pure preparation at a caller-captured stable time. 102 /// No storage, signer, clock or network is invoked. Retain this value for 103 /// atomic submission replay, which compares captured timestamps exactly. 104 pub fn authored_preparation( 105 &self, 106 captured_at_unix_ms: u64, 107 ) -> Result<PrepareAuthoredOperation, Error> { 108 let (operation_id, artifact_id, delivery_plan_id) = authored_ids(self.operation_id)?; 109 let operation = 110 AuthoredOperation::new(operation_id, vec![artifact_id], captured_at_unix_ms) 111 .map_err(map_storage_error)?; 112 let artifact = AuthoredArtifact::planned( 113 artifact_id, 114 operation_id, 115 0, 116 self.plan(), 117 captured_at_unix_ms, 118 ) 119 .map_err(map_storage_error)?; 120 let intent = AuthoredDeliveryIntent::new( 121 delivery_request_id(self.operation_id), 122 self.targets.clone(), 123 self.satisfaction.clone(), 124 self.delivery_deadline_unix_ms, 125 ) 126 .map_err(map_storage_error)?; 127 let delivery_plan = 128 AuthoredDeliveryPlan::new(delivery_plan_id, artifact_id, intent, captured_at_unix_ms) 129 .map_err(map_storage_error)?; 130 PrepareAuthoredOperation::new( 131 operation, 132 vec![artifact], 133 vec![delivery_plan], 134 authored_push_input_digest(self)?, 135 captured_at_unix_ms, 136 ) 137 .map_err(map_storage_error) 138 } 139 140 pub const fn operation_id(&self) -> SyncId { 141 self.operation_id 142 } 143 pub const fn idempotency_key(&self) -> &IdempotencyKey { 144 &self.idempotency_key 145 } 146 pub const fn actor(&self) -> &Actor { 147 &self.actor 148 } 149 pub const fn plan(&self) -> &AuthoredEventPlan { 150 &self.plan 151 } 152 pub const fn targets(&self) -> &TargetSet { 153 &self.targets 154 } 155 pub const fn satisfaction(&self) -> &SatisfactionPolicy { 156 &self.satisfaction 157 } 158 pub const fn delivery_deadline_unix_ms(&self) -> u64 { 159 self.delivery_deadline_unix_ms 160 } 161 pub const fn cancellation(&self) -> CancellationPolicy { 162 self.cancellation 163 } 164 } 165 166 impl core::fmt::Debug for PushRequest { 167 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { 168 formatter 169 .debug_struct("PushRequest") 170 .field("operation_id", &self.operation_id) 171 .field("idempotency_key", &self.idempotency_key) 172 .field("actor", &self.actor) 173 .field("plan", &"[redacted exact authored plan]") 174 .field("targets", &self.targets) 175 .field("satisfaction", &self.satisfaction) 176 .field("delivery_deadline_unix_ms", &self.delivery_deadline_unix_ms) 177 .field("cancellation", &self.cancellation) 178 .finish() 179 } 180 } 181 182 /// Complete durable intent created before any signer or transport effect. 183 #[derive(Clone, Debug, Eq, PartialEq)] 184 pub struct PushPreparation { 185 operation: AuthoredOperation, 186 artifact: AuthoredArtifact, 187 delivery_plan: AuthoredDeliveryPlan, 188 replay: bool, 189 } 190 191 impl PushPreparation { 192 pub const fn operation(&self) -> &AuthoredOperation { 193 &self.operation 194 } 195 pub const fn artifact(&self) -> &AuthoredArtifact { 196 &self.artifact 197 } 198 pub const fn delivery_plan(&self) -> &AuthoredDeliveryPlan { 199 &self.delivery_plan 200 } 201 pub const fn is_replay(&self) -> bool { 202 self.replay 203 } 204 } 205 206 /// Current durable authored state for one push operation. 207 #[derive(Clone, Debug, Eq, PartialEq)] 208 pub struct PushStatus { 209 operation: AuthoredOperation, 210 artifact: AuthoredArtifact, 211 delivery_history: AuthoredDeliveryHistory, 212 settlement: OperationSettlement, 213 } 214 215 impl PushStatus { 216 pub const fn operation(&self) -> &AuthoredOperation { 217 &self.operation 218 } 219 pub const fn artifact(&self) -> &AuthoredArtifact { 220 &self.artifact 221 } 222 pub const fn delivery_plan(&self) -> &AuthoredDeliveryPlan { 223 self.delivery_history.plan() 224 } 225 /// Consistent original-claim history; missing provenance remains unknown. 226 pub const fn delivery_history(&self) -> &AuthoredDeliveryHistory { 227 &self.delivery_history 228 } 229 pub const fn settlement(&self) -> OperationSettlement { 230 self.settlement 231 } 232 } 233 234 /// Durable result of one bounded signing execution. 235 #[derive(Clone, Debug, Eq, PartialEq)] 236 pub struct SigningRunReceipt { 237 artifact: AuthoredArtifact, 238 replay: bool, 239 } 240 241 impl SigningRunReceipt { 242 pub const fn artifact(&self) -> &AuthoredArtifact { 243 &self.artifact 244 } 245 pub const fn is_replay(&self) -> bool { 246 self.replay 247 } 248 } 249 250 /// Durable result of one bounded local-admission execution. 251 #[derive(Clone, Debug, Eq, PartialEq)] 252 pub struct AdmissionRunReceipt { 253 artifact: AuthoredArtifact, 254 replay: bool, 255 } 256 257 /// Durable result of one bounded authored-delivery execution. 258 #[derive(Clone, Debug, Eq, PartialEq)] 259 pub struct DeliveryExecutionReceipt { 260 plan: AuthoredDeliveryPlan, 261 replay: bool, 262 } 263 264 /// Durable result of cancelling all still-cancellable phases of one push. 265 #[derive(Clone, Debug, Eq, PartialEq)] 266 pub struct PushCancellationReceipt { 267 status: PushStatus, 268 changed: bool, 269 } 270 271 impl PushCancellationReceipt { 272 pub const fn status(&self) -> &PushStatus { 273 &self.status 274 } 275 pub const fn changed(&self) -> bool { 276 self.changed 277 } 278 } 279 280 impl DeliveryExecutionReceipt { 281 pub const fn plan(&self) -> &AuthoredDeliveryPlan { 282 &self.plan 283 } 284 pub const fn is_replay(&self) -> bool { 285 self.replay 286 } 287 } 288 289 impl AdmissionRunReceipt { 290 pub const fn artifact(&self) -> &AuthoredArtifact { 291 &self.artifact 292 } 293 pub const fn is_replay(&self) -> bool { 294 self.replay 295 } 296 } 297 298 impl Engine { 299 /// Atomically persists the complete parent, exact plan, and delivery intent. 300 /// 301 /// This method invokes neither a signer nor a transport. Exact request 302 /// replay returns the original durable preparation; conflicting reuse of 303 /// the operation identity fails closed. 304 pub async fn prepare_push(&self, request: PushRequest) -> Result<PushPreparation, Error> { 305 let prepared_at = self.clock.now_unix_ms()?; 306 let (operation_id, artifact_id, delivery_plan_id) = authored_ids(request.operation_id)?; 307 let command = AuthoredAtomicCommand::Prepare(request.authored_preparation(prepared_at)?); 308 let receipt = self 309 .storage 310 .execute_authored(command) 311 .await 312 .map_err(map_storage_error)?; 313 let AuthoredAtomicOutcome::Prepared { 314 operation, 315 artifacts, 316 delivery_plans, 317 } = receipt.outcome() 318 else { 319 return Err(Error::StorageFailed); 320 }; 321 let [artifact] = artifacts.as_slice() else { 322 return Err(Error::StorageFailed); 323 }; 324 let [delivery_plan] = delivery_plans.as_slice() else { 325 return Err(Error::StorageFailed); 326 }; 327 if operation.operation_id() != operation_id 328 || artifact.artifact_id() != artifact_id 329 || delivery_plan.plan_id() != delivery_plan_id 330 { 331 return Err(Error::StorageFailed); 332 } 333 Ok(PushPreparation { 334 operation: operation.clone(), 335 artifact: artifact.clone(), 336 delivery_plan: delivery_plan.clone(), 337 replay: receipt.disposition() == AtomicCommitDisposition::Replay, 338 }) 339 } 340 341 /// Loads the complete durable state for one prepared push operation. 342 pub async fn push_status(&self, operation_id: SyncId) -> Result<Option<PushStatus>, Error> { 343 let (operation_id, expected_artifact_id, expected_plan_id) = authored_ids(operation_id)?; 344 let Some(operation) = self 345 .storage 346 .authored_operation(operation_id) 347 .await 348 .map_err(map_storage_error)? 349 else { 350 return Ok(None); 351 }; 352 if operation.artifact_ids() != [expected_artifact_id] { 353 return Err(Error::StorageFailed); 354 } 355 let artifact = self 356 .storage 357 .authored_artifact(expected_artifact_id) 358 .await 359 .map_err(map_storage_error)? 360 .ok_or(Error::StorageFailed)?; 361 let delivery_history = self 362 .storage 363 .authored_delivery_history(expected_plan_id) 364 .await 365 .map_err(map_storage_error)? 366 .ok_or(Error::StorageFailed)?; 367 delivery_history.validate().map_err(map_storage_error)?; 368 let delivery_plan = delivery_history.plan(); 369 if artifact.artifact_id() != expected_artifact_id 370 || delivery_plan.plan_id() != expected_plan_id 371 || artifact.operation_id() != operation.operation_id() 372 || delivery_plan.artifact_id() != artifact.artifact_id() 373 { 374 return Err(Error::StorageFailed); 375 } 376 let settlement = OperationSettlement::evaluate_complete( 377 &operation, 378 core::slice::from_ref(&artifact), 379 core::slice::from_ref(delivery_plan), 380 ) 381 .map_err(map_storage_error)?; 382 Ok(Some(PushStatus { 383 operation, 384 artifact, 385 delivery_history, 386 settlement, 387 })) 388 } 389 390 /// Cancels every remaining local phase without discarding signed bytes or 391 /// delivery evidence. Re-entry after a phase-boundary crash finishes any 392 /// remaining cancellation and terminal replays are no-ops. 393 pub async fn cancel_push( 394 &self, 395 operation_id: SyncId, 396 ) -> Result<PushCancellationReceipt, Error> { 397 let mut status = self 398 .push_status(operation_id) 399 .await? 400 .ok_or(Error::StorageFailed)?; 401 let mut changed = false; 402 let now = self 403 .clock 404 .now_unix_ms()? 405 .max(status.artifact.updated_at_unix_ms()) 406 .max(status.delivery_plan().updated_at_unix_ms()); 407 408 let artifact_target = match status.artifact.signing_state() { 409 SigningState::Planned | SigningState::Retryable => Some( 410 CancelAuthoredTarget::ArtifactSigning(status.artifact.artifact_id()), 411 ), 412 SigningState::Signed 413 if matches!( 414 status.artifact.admission_state(), 415 AdmissionState::Pending | AdmissionState::Retryable 416 ) => 417 { 418 Some(CancelAuthoredTarget::ArtifactAdmission( 419 status.artifact.artifact_id(), 420 )) 421 } 422 SigningState::Signed 423 | SigningState::Indeterminate 424 | SigningState::FailedTerminal 425 | SigningState::Cancelled => None, 426 }; 427 if let Some(target) = artifact_target { 428 self.storage 429 .execute_authored(AuthoredAtomicCommand::Cancel( 430 CancelAuthoredWork::new(target, status.artifact.revision(), now) 431 .map_err(map_storage_error)?, 432 )) 433 .await 434 .map_err(map_storage_error)?; 435 changed = true; 436 status = self 437 .push_status(operation_id) 438 .await? 439 .ok_or(Error::StorageFailed)?; 440 } 441 442 if status.delivery_plan().stop_requested_at_unix_ms().is_none() { 443 let command = AuthoredAtomicCommand::Cancel( 444 CancelAuthoredWork::new( 445 CancelAuthoredTarget::DeliveryPlan(status.delivery_plan().plan_id()), 446 status.delivery_plan().revision(), 447 now, 448 ) 449 .map_err(map_storage_error)?, 450 ); 451 let receipt = self 452 .storage 453 .execute_authored(command.clone()) 454 .await 455 .map_err(map_storage_error)?; 456 if !receipt.matches_command(&command) { 457 return Err(Error::StorageFailed); 458 } 459 let AuthoredAtomicOutcome::DeliveryPlan(stopped) = receipt.outcome() else { 460 return Err(Error::StorageFailed); 461 }; 462 if stopped.validate().is_err() 463 || stopped.plan_id() != status.delivery_plan().plan_id() 464 || stopped.artifact_id() != status.artifact.artifact_id() 465 || stopped.request() != status.delivery_plan().request() 466 || stopped.stop_requested_at_unix_ms() != Some(now) 467 || status.delivery_plan().revision().get().checked_add(1) 468 != Some(stopped.revision().get()) 469 { 470 return Err(Error::StorageFailed); 471 } 472 changed = true; 473 status = self 474 .push_status(operation_id) 475 .await? 476 .ok_or(Error::StorageFailed)?; 477 if status.delivery_plan().stop_requested_at_unix_ms().is_none() { 478 return Err(Error::StorageFailed); 479 } 480 } 481 482 Ok(PushCancellationReceipt { status, changed }) 483 } 484 485 /// Claims and executes one prepared signing artifact. 486 pub async fn sign_prepared(&self, request: PushRequest) -> Result<SigningRunReceipt, Error> { 487 self.prepare_push(request.clone()).await?; 488 let status = self 489 .push_status(request.operation_id) 490 .await? 491 .ok_or(Error::StorageFailed)?; 492 if status.delivery_plan().state() == AuthoredDeliveryState::Cancelled { 493 self.cancel_push(request.operation_id).await?; 494 return Err(Error::SigningCancelled); 495 } 496 let artifact = status.artifact; 497 match artifact.signing_state() { 498 SigningState::Signed => { 499 return Ok(SigningRunReceipt { 500 artifact, 501 replay: true, 502 }); 503 } 504 SigningState::Indeterminate => return Err(Error::SigningIndeterminate), 505 SigningState::FailedTerminal | SigningState::Cancelled => { 506 return Err(Error::SignerFailed); 507 } 508 SigningState::Planned | SigningState::Retryable => {} 509 } 510 let signer = self.signer.as_deref().ok_or(Error::MissingSigner)?; 511 512 let now = self.clock.now_unix_ms()?.max(artifact.updated_at_unix_ms()); 513 if let Some(existing) = artifact.signing_claim() { 514 if now < existing.expires_at_unix_ms() { 515 return Err(Error::WorkClaimConflict); 516 } 517 if existing.owner() == SIGNING_CLAIM_OWNER_NON_REPLAYABLE { 518 let (claimed, fence) = self 519 .claim_artifact( 520 artifact, 521 ClaimAuthoredTarget::ArtifactSigning, 522 SIGNING_CLAIM_OWNER_NON_REPLAYABLE, 523 now, 524 ) 525 .await?; 526 let failure = WorkFailure::new( 527 "signing_effect_unknown_after_restart", 528 WorkPhase::Signing, 529 FailureClass::Indeterminate, 530 None, 531 None, 532 ) 533 .map_err(map_storage_error)?; 534 self.apply_artifact_failure(claimed.artifact_id(), fence, failure, None, now) 535 .await?; 536 return Err(Error::SigningIndeterminate); 537 } 538 } 539 540 let signer_status = signer.status().await.map_err(|_| Error::SignerFailed)?; 541 let replay_capability = signer_replay_capability(&signer_status)?; 542 let (claimed, fence) = self 543 .claim_artifact( 544 artifact, 545 ClaimAuthoredTarget::ArtifactSigning, 546 signing_claim_owner(replay_capability), 547 now, 548 ) 549 .await?; 550 let persisted_plan = claimed 551 .plan() 552 .ok_or(Error::StorageFailed)? 553 .decode() 554 .map_err(map_storage_error)? 555 .into_plan(); 556 if persisted_plan != request.plan { 557 return Err(Error::StorageConflict); 558 } 559 let signing_claim = claimed.signing_claim().ok_or(Error::StorageFailed)?.clone(); 560 let deadline = self.deadlines.deadline_unix_ms(OperationKind::Sign, now)?; 561 let signing_operation = SigningOperationId::new(*request.operation_id.as_bytes()) 562 .map_err(|_| Error::InvalidPushRequest)?; 563 let signing_artifact = SigningArtifactId::new(*claimed.artifact_id().as_bytes()) 564 .map_err(|_| Error::InvalidPushRequest)?; 565 let sign_request = radroots_signing::SignRequest::new( 566 OperationId::SyncPush, 567 SigningIntentId::new(signing_operation, signing_artifact), 568 request.actor, 569 persisted_plan, 570 SignPolicy::new(deadline, request.cancellation) 571 .map_err(|_| Error::InvalidPushRequest)?, 572 ) 573 .map_err(|_| Error::InvalidPushRequest)?; 574 575 let expected_request = sign_request.clone(); 576 let signer_result = match signer.sign_authored_evidence(sign_request).await { 577 Ok(receipt) => { 578 // A missing observation time leaves the durable attempt unresolved; 579 // it is not evidence that signing failed without an effect. 580 let observed_at_unix_ms = self.clock.now_unix_ms()?; 581 receipt.revalidate(&expected_request, observed_at_unix_ms) 582 } 583 Err(error) => Err(error), 584 }; 585 586 match signer_result { 587 Ok(receipt) => { 588 let command = AuthoredAtomicCommand::RecordSigned( 589 RecordSignedArtifact::new( 590 claimed.operation_id(), 591 claimed.artifact_id(), 592 signing_claim, 593 receipt.signed_event().clone(), 594 receipt.observed_at_unix_ms(), 595 ) 596 .map_err(map_storage_error)?, 597 ); 598 let applied = self 599 .storage 600 .execute_authored(command.clone()) 601 .await 602 .map_err(map_storage_error)?; 603 if !applied.matches_command(&command) { 604 return Err(Error::StorageFailed); 605 } 606 let status = self 607 .push_status(request.operation_id) 608 .await? 609 .ok_or(Error::StorageFailed)?; 610 if status 611 .artifact 612 .signed() 613 .is_none_or(|signed| signed.event() != receipt.signed_event()) 614 { 615 return Err(Error::StorageFailed); 616 } 617 if expected_request.cancellation_signal().is_cancelled() 618 || status.delivery_plan().state() == AuthoredDeliveryState::Cancelled 619 { 620 let stopped = self.cancel_push(request.operation_id).await?; 621 if stopped 622 .status 623 .artifact 624 .signed() 625 .is_none_or(|signed| signed.event() != receipt.signed_event()) 626 { 627 return Err(Error::StorageFailed); 628 } 629 return Err(Error::SigningCancelled); 630 } 631 match status.artifact.signing_state() { 632 SigningState::Cancelled => return Err(Error::SigningCancelled), 633 SigningState::FailedTerminal => return Err(Error::SignerFailed), 634 SigningState::Signed => {} 635 _ => return Err(Error::StorageFailed), 636 } 637 if receipt.observed_at_unix_ms() >= deadline { 638 return Err(Error::SignerDeadlineExceeded); 639 } 640 Ok(SigningRunReceipt { 641 artifact: status.artifact, 642 replay: applied.disposition() == AtomicCommitDisposition::Replay, 643 }) 644 } 645 Err(error) => { 646 let applied_at = self.clock.now_unix_ms()?; 647 let current = self 648 .push_status(request.operation_id) 649 .await? 650 .ok_or(Error::StorageFailed)?; 651 if current.artifact.signed().is_some() 652 || current.artifact.signing_claim() != claimed.signing_claim() 653 { 654 // Another worker may have committed evidence or a stop while this one waited. 655 return Err(match error.kind() { 656 radroots_signing::error::Kind::DeadlineExceeded => { 657 Error::SignerDeadlineExceeded 658 } 659 radroots_signing::error::Kind::SignerCancelled => Error::SigningCancelled, 660 _ => Error::SignerFailed, 661 }); 662 } 663 if error.kind() == radroots_signing::error::Kind::DeadlineExceeded 664 && claimed 665 .signing_claim() 666 .is_some_and(|claim| applied_at >= claim.expires_at_unix_ms()) 667 { 668 // The signer result arrived after this worker's fence 669 // expired. It is rejected, but this stale worker must not 670 // mutate durable state; recovery will reclaim the exact 671 // request under a fresh fence. 672 return Err(Error::SignerDeadlineExceeded); 673 } 674 if error.kind() == radroots_signing::error::Kind::SignerCancelled { 675 let command = AuthoredAtomicCommand::Cancel( 676 CancelAuthoredWork::new( 677 CancelAuthoredTarget::ArtifactSigning(claimed.artifact_id()), 678 claimed.revision(), 679 applied_at, 680 ) 681 .map_err(map_storage_error)?, 682 ); 683 self.storage 684 .execute_authored(command) 685 .await 686 .map_err(map_storage_error)?; 687 return Err(Error::SigningCancelled); 688 } 689 let disposition = recovery_disposition( 690 replay_capability, 691 error.remote_effect(), 692 error.retryable(), 693 ); 694 let class = match disposition { 695 RecoveryDisposition::RetryExactRequest | RecoveryDisposition::RetryLocal => { 696 FailureClass::Retryable 697 } 698 RecoveryDisposition::Indeterminate => FailureClass::Indeterminate, 699 RecoveryDisposition::Failed => FailureClass::Terminal, 700 _ => FailureClass::Indeterminate, 701 }; 702 let retry_at = if class == FailureClass::Retryable { 703 Some( 704 applied_at 705 .checked_add(WORK_RETRY_DELAY_MS) 706 .ok_or(Error::DeadlineOverflow)?, 707 ) 708 } else { 709 None 710 }; 711 let failure = 712 WorkFailure::new(error.code(), WorkPhase::Signing, class, retry_at, None) 713 .map_err(map_storage_error)?; 714 let retry = retry_schedule(claimed.signing_retry(), &failure, retry_at)?; 715 self.apply_artifact_failure( 716 claimed.artifact_id(), 717 fence, 718 failure, 719 retry, 720 applied_at, 721 ) 722 .await?; 723 match disposition { 724 RecoveryDisposition::Indeterminate => Err(Error::SigningIndeterminate), 725 RecoveryDisposition::RetryExactRequest 726 | RecoveryDisposition::RetryLocal 727 | RecoveryDisposition::Failed => { 728 if error.kind() == radroots_signing::error::Kind::DeadlineExceeded { 729 Err(Error::SignerDeadlineExceeded) 730 } else { 731 Err(Error::SignerFailed) 732 } 733 } 734 _ => Err(Error::SigningIndeterminate), 735 } 736 } 737 } 738 } 739 740 /// Claims and executes local admission for one durably signed artifact. 741 pub async fn admit_signed(&self, operation_id: SyncId) -> Result<AdmissionRunReceipt, Error> { 742 let status = self 743 .push_status(operation_id) 744 .await? 745 .ok_or(Error::StorageFailed)?; 746 let delivery_stopped = status.delivery_plan().stop_requested_at_unix_ms().is_some(); 747 let artifact = status.artifact; 748 if artifact.signing_state() != SigningState::Signed { 749 return Err(Error::InvalidSignerOutput); 750 } 751 if artifact.admission_state().is_admitted() { 752 return Ok(AdmissionRunReceipt { 753 artifact, 754 replay: true, 755 }); 756 } 757 if delivery_stopped { 758 self.cancel_push(operation_id).await?; 759 return Err(Error::AdmissionFailed); 760 } 761 if matches!( 762 artifact.admission_state(), 763 AdmissionState::Rejected | AdmissionState::Cancelled 764 ) { 765 return Err(Error::AdmissionFailed); 766 } 767 let now = self.clock.now_unix_ms()?.max(artifact.updated_at_unix_ms()); 768 if artifact 769 .admission_claim() 770 .is_some_and(|claim| now < claim.expires_at_unix_ms()) 771 { 772 return Err(Error::WorkClaimConflict); 773 } 774 let (claimed, fence) = self 775 .claim_artifact( 776 artifact, 777 ClaimAuthoredTarget::ArtifactAdmission, 778 ADMISSION_CLAIM_OWNER, 779 now, 780 ) 781 .await?; 782 let event = claimed 783 .signed() 784 .ok_or(Error::InvalidSignerOutput)? 785 .event() 786 .clone(); 787 let persisted_plan = claimed 788 .plan() 789 .ok_or(Error::StorageFailed)? 790 .decode() 791 .map_err(map_storage_error)? 792 .into_plan(); 793 let contract_id = persisted_plan.body().contract().contract_id().as_str(); 794 let admission = match outbound_admission(&event, contract_id, now) { 795 Ok(admission) => admission, 796 Err(error) => { 797 let failure = WorkFailure::new( 798 "invalid_signed_artifact", 799 WorkPhase::Admission, 800 FailureClass::Terminal, 801 None, 802 None, 803 ) 804 .map_err(map_storage_error)?; 805 self.apply_artifact_failure(claimed.artifact_id(), fence, failure, None, now) 806 .await?; 807 return Err(error); 808 } 809 }; 810 let admission_receipt = match EventStore::admit(self.storage.as_ref(), admission).await { 811 Ok(receipt) => receipt, 812 Err(error) => { 813 let applied_at = self.clock.now_unix_ms()?.max(claimed.updated_at_unix_ms()); 814 let terminal = matches!( 815 error, 816 radroots_storage::Error::EventConflict 817 | radroots_storage::Error::AdmissionRegression 818 ); 819 let retry_at = if terminal { 820 None 821 } else { 822 Some( 823 applied_at 824 .checked_add(WORK_RETRY_DELAY_MS) 825 .ok_or(Error::DeadlineOverflow)?, 826 ) 827 }; 828 let failure = WorkFailure::new( 829 if terminal { 830 "admission_conflict" 831 } else { 832 "admission_storage_unavailable" 833 }, 834 WorkPhase::Admission, 835 if terminal { 836 FailureClass::Terminal 837 } else { 838 FailureClass::Retryable 839 }, 840 retry_at, 841 None, 842 ) 843 .map_err(map_storage_error)?; 844 let retry = retry_schedule(claimed.admission_retry(), &failure, retry_at)?; 845 self.apply_artifact_failure( 846 claimed.artifact_id(), 847 fence, 848 failure, 849 retry, 850 applied_at, 851 ) 852 .await?; 853 return Err(match error { 854 radroots_storage::Error::SpaceInsufficient => Error::StorageSpaceInsufficient, 855 _ => Error::AdmissionFailed, 856 }); 857 } 858 }; 859 let state = match admission_receipt.disposition() { 860 AdmissionDisposition::Duplicate => AdmissionState::Duplicate, 861 AdmissionDisposition::Inserted | AdmissionDisposition::Advanced => { 862 AdmissionState::Inserted 863 } 864 }; 865 let applied_at = self.clock.now_unix_ms()?.max(claimed.updated_at_unix_ms()); 866 let command = AuthoredAtomicCommand::ApplyAdmission( 867 ApplyAdmissionResult::new(claimed.artifact_id(), fence, state, None, None, applied_at) 868 .map_err(map_storage_error)?, 869 ); 870 let applied = self 871 .storage 872 .execute_authored(command) 873 .await 874 .map_err(map_storage_error)?; 875 let AuthoredAtomicOutcome::Artifact(artifact) = applied.outcome() else { 876 return Err(Error::StorageFailed); 877 }; 878 Ok(AdmissionRunReceipt { 879 artifact: artifact.clone(), 880 replay: false, 881 }) 882 } 883 884 async fn claim_artifact( 885 &self, 886 artifact: AuthoredArtifact, 887 target: fn(AuthoredArtifactId) -> ClaimAuthoredTarget, 888 owner: &'static str, 889 acquired_at: u64, 890 ) -> Result<(AuthoredArtifact, WorkFence), Error> { 891 let existing = match target(artifact.artifact_id()) { 892 ClaimAuthoredTarget::ArtifactSigning(_) => artifact.signing_claim(), 893 ClaimAuthoredTarget::ArtifactAdmission(_) => artifact.admission_claim(), 894 ClaimAuthoredTarget::DeliveryPlan(_) => return Err(Error::StorageFailed), 895 }; 896 let generation = existing.map_or(1, |claim| claim.generation().get().saturating_add(1)); 897 let generation = NonZeroU64::new(generation).ok_or(Error::StorageFailed)?; 898 let expires_at = acquired_at 899 .checked_add(self.deadlines.timeout_ms(OperationKind::Sign)) 900 .ok_or(Error::DeadlineOverflow)?; 901 let claim = WorkClaim::new( 902 *self.ids.next_id(OperationKind::Sign)?.as_bytes(), 903 owner, 904 generation, 905 acquired_at, 906 expires_at, 907 artifact.revision(), 908 ) 909 .map_err(map_storage_error)?; 910 let fence = WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()) 911 .map_err(map_storage_error)?; 912 let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 913 target(artifact.artifact_id()), 914 claim, 915 )); 916 let receipt = self 917 .storage 918 .execute_authored(command) 919 .await 920 .map_err(map_claim_error)?; 921 let AuthoredAtomicOutcome::Artifact(claimed) = receipt.outcome() else { 922 return Err(Error::StorageFailed); 923 }; 924 Ok((claimed.clone(), fence)) 925 } 926 927 async fn apply_artifact_failure( 928 &self, 929 artifact_id: AuthoredArtifactId, 930 fence: WorkFence, 931 failure: WorkFailure, 932 retry: Option<RetrySchedule>, 933 applied_at: u64, 934 ) -> Result<AuthoredArtifact, Error> { 935 let command = AuthoredAtomicCommand::ApplyFailure( 936 ApplyWorkFailure::new( 937 AuthoredWorkTarget::Artifact(artifact_id), 938 fence, 939 failure, 940 retry, 941 applied_at, 942 ) 943 .map_err(map_storage_error)?, 944 ); 945 let receipt = self 946 .storage 947 .execute_authored(command) 948 .await 949 .map_err(map_storage_error)?; 950 let AuthoredAtomicOutcome::Artifact(artifact) = receipt.outcome() else { 951 return Err(Error::StorageFailed); 952 }; 953 Ok(artifact.clone()) 954 } 955 } 956 957 fn signer_replay_capability( 958 status: &radroots_signing::SignerStatus, 959 ) -> Result<ReplayCapability, Error> { 960 let mut capabilities = status.capabilities().iter(); 961 let first = capabilities 962 .next() 963 .ok_or(Error::SignerCapabilityUnavailable)? 964 .replay(); 965 if capabilities.all(|capability| capability.replay() == first) { 966 Ok(first) 967 } else { 968 Ok(ReplayCapability::NonReplayable) 969 } 970 } 971 972 fn signing_claim_owner(replay: ReplayCapability) -> &'static str { 973 match replay { 974 ReplayCapability::ExactReplayByRequestId => SIGNING_CLAIM_OWNER_EXACT, 975 ReplayCapability::LocalReplaySafe => SIGNING_CLAIM_OWNER_LOCAL, 976 ReplayCapability::NonReplayable => SIGNING_CLAIM_OWNER_NON_REPLAYABLE, 977 _ => SIGNING_CLAIM_OWNER_NON_REPLAYABLE, 978 } 979 } 980 981 fn retry_schedule( 982 previous: Option<&RetrySchedule>, 983 failure: &WorkFailure, 984 retry_at: Option<u64>, 985 ) -> Result<Option<RetrySchedule>, Error> { 986 let Some(retry_at) = retry_at else { 987 return Ok(None); 988 }; 989 let schedule = match previous { 990 Some(previous) => previous.next_attempt(retry_at, failure.clone()), 991 None => RetrySchedule::new(NonZeroU32::MIN, retry_at, failure.clone()), 992 } 993 .map_err(map_storage_error)?; 994 Ok(Some(schedule)) 995 } 996 997 fn normalize_sink_failure( 998 request: &DeliveryRequest, 999 failure: SinkFailure, 1000 attempted_at_unix_ms: u64, 1001 ) -> Result<SinkFailure, Error> { 1002 if failure.retryability() != Retryability::Retryable { 1003 return Ok(failure); 1004 } 1005 let default_retry_at = attempted_at_unix_ms 1006 .checked_add(WORK_RETRY_DELAY_MS) 1007 .ok_or(Error::DeadlineOverflow)?; 1008 let retry_at = failure 1009 .retry_after_unix_ms() 1010 .unwrap_or(default_retry_at) 1011 .max( 1012 attempted_at_unix_ms 1013 .checked_add(1) 1014 .ok_or(Error::DeadlineOverflow)?, 1015 ); 1016 if failure.retry_after_unix_ms() == Some(retry_at) { 1017 return Ok(failure); 1018 } 1019 SinkFailure::for_request( 1020 request, 1021 failure.code(), 1022 failure.retryability(), 1023 Some(retry_at), 1024 failure.message().map(str::to_owned), 1025 failure.partial_evidence().to_vec(), 1026 ) 1027 .map_err(|_| Error::InvalidDeliveryRequest) 1028 } 1029 1030 fn delivery_retry_schedule( 1031 attempt: NonZeroU32, 1032 outcome: &DeliveryAttemptOutcome, 1033 satisfaction: SatisfactionState, 1034 attempted_at_unix_ms: u64, 1035 ) -> Result<Option<RetrySchedule>, Error> { 1036 if satisfaction != SatisfactionState::Pending || attempt.get() >= DELIVERY_PLAN_ATTEMPTS_MAX { 1037 return Ok(None); 1038 } 1039 let failure = match outcome { 1040 DeliveryAttemptOutcome::Receipt(_) => { 1041 let retry_at = attempted_at_unix_ms 1042 .checked_add(WORK_RETRY_DELAY_MS) 1043 .ok_or(Error::DeadlineOverflow)?; 1044 WorkFailure::new( 1045 "delivery_pending", 1046 WorkPhase::Delivery, 1047 FailureClass::Retryable, 1048 Some(retry_at), 1049 None, 1050 ) 1051 .map_err(map_storage_error)? 1052 } 1053 DeliveryAttemptOutcome::SinkFailure(failure) 1054 if failure.retryability() == Retryability::Retryable => 1055 { 1056 WorkFailure::new( 1057 failure.code(), 1058 WorkPhase::Delivery, 1059 FailureClass::Retryable, 1060 failure.retry_after_unix_ms(), 1061 failure.message().map(str::to_owned), 1062 ) 1063 .map_err(map_storage_error)? 1064 } 1065 DeliveryAttemptOutcome::SinkFailure(_) => return Ok(None), 1066 }; 1067 let retry_at = failure.retry_after_unix_ms().unwrap_or( 1068 attempted_at_unix_ms 1069 .checked_add(WORK_RETRY_DELAY_MS) 1070 .ok_or(Error::DeadlineOverflow)?, 1071 ); 1072 RetrySchedule::new(attempt, retry_at, failure) 1073 .map(Some) 1074 .map_err(map_storage_error) 1075 } 1076 1077 fn outbound_admission( 1078 event: &radroots_event::SignedEvent, 1079 contract_id: &str, 1080 observed_at_unix_ms: u64, 1081 ) -> Result<EventAdmission, Error> { 1082 let verified = verify::signature( 1083 verify::id(RawEvent::new(event.envelope().clone())) 1084 .map_err(|_| Error::InvalidSignerOutput)?, 1085 &Nip01SignatureVerifier, 1086 ) 1087 .map_err(|_| Error::InvalidSignerOutput)?; 1088 let validated = verified 1089 .validate_contract_for_admission(contract_id) 1090 .map_err(|_| Error::InvalidSignerOutput)?; 1091 let policy = RegistryPolicy::visible(); 1092 if policy.decide(&validated) != AdmissionDecision::Visible { 1093 return Err(Error::InvalidSignerOutput); 1094 } 1095 struct Evidence; 1096 impl radroots_event::admission::AdmissionPolicy for Evidence { 1097 type Error = core::convert::Infallible; 1098 fn policy_id(&self) -> &'static str { 1099 "radroots.registry_v7" 1100 } 1101 fn admit( 1102 &self, 1103 _: &radroots_event::admission::ContractValidatedEvent, 1104 ) -> Result<(), Self::Error> { 1105 Ok(()) 1106 } 1107 } 1108 impl radroots_event::admission::VisibilityPolicy for Evidence { 1109 type Error = core::convert::Infallible; 1110 fn policy_id(&self) -> &'static str { 1111 "radroots.registry_v7" 1112 } 1113 fn make_visible( 1114 &self, 1115 _: &radroots_event::admission::AdmittedEvent, 1116 ) -> Result<(), Self::Error> { 1117 Ok(()) 1118 } 1119 } 1120 let visible = validated 1121 .admit_with(&Evidence) 1122 .and_then(|event| event.make_visible_with(&Evidence)) 1123 .map_err(|never| match never {})?; 1124 let local = Target::new(TransportId::LOCAL, "local:sync-authored") 1125 .map_err(|_| Error::InvalidPushRequest)?; 1126 let provenance = EventProvenance::new( 1127 TransportId::LOCAL, 1128 local.fingerprint().clone(), 1129 observed_at_unix_ms, 1130 ) 1131 .map_err(|_| Error::InvalidPushRequest)?; 1132 EventAdmission::visible(ObservedEvent::new(event.clone(), provenance), visible) 1133 .map_err(map_storage_error) 1134 } 1135 1136 fn authored_push_input_digest(request: &PushRequest) -> Result<AtomicCommitDigest, Error> { 1137 let mut hasher = Sha256::new(); 1138 hash_field(&mut hasher, b"radroots.sync.authored-push.v2"); 1139 hash_field(&mut hasher, request.operation_id.as_bytes()); 1140 hash_field(&mut hasher, request.idempotency_key.as_str().as_bytes()); 1141 hash_field(&mut hasher, request.plan.digest().as_bytes()); 1142 hash_field(&mut hasher, request.actor.public_key().as_bytes()); 1143 let source = match request.actor.source() { 1144 radroots_signing::actor::ActorSource::LocalAccount(_) => 0, 1145 radroots_signing::actor::ActorSource::ExplicitPublicKey => 1, 1146 radroots_signing::actor::ActorSource::RemoteSigner(_) => 2, 1147 radroots_signing::actor::ActorSource::Service(_) => 3, 1148 _ => return Err(Error::InvalidPushRequest), 1149 }; 1150 hasher.update([source]); 1151 if let Some(account_id) = request.actor.account_id() { 1152 hash_field(&mut hasher, account_id.as_bytes()); 1153 } 1154 for role in request.actor.roles() { 1155 hash_field(&mut hasher, role.as_str().as_bytes()); 1156 } 1157 for target in request.targets.targets() { 1158 hash_field(&mut hasher, target.fingerprint().as_str().as_bytes()); 1159 } 1160 hash_satisfaction(&mut hasher, &request.satisfaction); 1161 hasher.update(request.delivery_deadline_unix_ms.to_be_bytes()); 1162 let cancellation = match request.cancellation { 1163 CancellationPolicy::PreservePublishedRequest => 0, 1164 CancellationPolicy::LocalCooperative => 1, 1165 _ => return Err(Error::InvalidPushRequest), 1166 }; 1167 hasher.update([cancellation]); 1168 Ok(AtomicCommitDigest::new(hasher.finalize().into())) 1169 } 1170 1171 fn hash_satisfaction(hasher: &mut Sha256, policy: &SatisfactionPolicy) { 1172 hasher.update([match policy.class() { 1173 radroots_transport::policy::SatisfactionClass::Accepted => 0, 1174 radroots_transport::policy::SatisfactionClass::Delivered => 1, 1175 }]); 1176 let targets = policy.targets(); 1177 if targets.is_any() { 1178 hasher.update([0]); 1179 } else if targets.is_all() { 1180 hasher.update([1]); 1181 } else if let Some(threshold) = targets.quorum_threshold() { 1182 hasher.update([2]); 1183 hasher.update(threshold.to_be_bytes()); 1184 } else if let Some(required) = targets.required_targets() { 1185 hasher.update([3]); 1186 for target in required { 1187 hash_field(hasher, target.as_str().as_bytes()); 1188 } 1189 } 1190 } 1191 1192 fn authored_ids( 1193 operation_id: SyncId, 1194 ) -> Result< 1195 ( 1196 OperationInstanceId, 1197 AuthoredArtifactId, 1198 AuthoredDeliveryPlanId, 1199 ), 1200 Error, 1201 > { 1202 let operation = 1203 OperationInstanceId::new(*operation_id.as_bytes()).map_err(map_storage_error)?; 1204 let artifact = AuthoredArtifactId::new(derive_child_id( 1205 b"radroots.sync.authored-artifact.v2", 1206 operation_id.as_bytes(), 1207 )) 1208 .map_err(map_storage_error)?; 1209 let delivery = AuthoredDeliveryPlanId::new(derive_child_id( 1210 b"radroots.sync.authored-delivery.v2", 1211 artifact.as_bytes(), 1212 )) 1213 .map_err(map_storage_error)?; 1214 Ok((operation, artifact, delivery)) 1215 } 1216 1217 fn derive_child_id(domain: &[u8], parent: &[u8; 16]) -> [u8; 16] { 1218 let mut hasher = Sha256::new(); 1219 hash_field(&mut hasher, domain); 1220 hash_field(&mut hasher, parent); 1221 let digest: [u8; 32] = hasher.finalize().into(); 1222 let mut child = [0_u8; 16]; 1223 child.copy_from_slice(&digest[..16]); 1224 child 1225 } 1226 1227 fn hash_field(hasher: &mut Sha256, value: &[u8]) { 1228 hasher.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes()); 1229 hasher.update(value); 1230 } 1231 1232 fn delivery_request_id(id: SyncId) -> String { 1233 const HEX: &[u8; 16] = b"0123456789abcdef"; 1234 let mut value = String::from("push-"); 1235 for byte in id.as_bytes() { 1236 value.push(HEX[(byte >> 4) as usize] as char); 1237 value.push(HEX[(byte & 0x0f) as usize] as char); 1238 } 1239 value 1240 } 1241 1242 fn map_storage_error(error: radroots_storage::Error) -> Error { 1243 match error { 1244 radroots_storage::Error::SpaceInsufficient => Error::StorageSpaceInsufficient, 1245 radroots_storage::Error::IdempotencyConflict 1246 | radroots_storage::Error::OperationIdentityMismatch 1247 | radroots_storage::Error::JournalRevisionConflict 1248 | radroots_storage::Error::OutboxPlanConflict 1249 | radroots_storage::Error::AtomicCommitConflict => Error::StorageConflict, 1250 _ => Error::StorageFailed, 1251 } 1252 } 1253 1254 fn map_claim_error(error: radroots_storage::Error) -> Error { 1255 match error { 1256 radroots_storage::Error::AtomicCommitConflict 1257 | radroots_storage::Error::InvalidAuthoredTransition 1258 | radroots_storage::Error::InvalidWorkClaim 1259 | radroots_storage::Error::DeliveryPlanClaimConflict => Error::WorkClaimConflict, 1260 other => map_storage_error(other), 1261 } 1262 } 1263 1264 #[cfg(test)] 1265 #[cfg_attr(coverage_nightly, coverage(off))] 1266 mod tests { 1267 use super::*; 1268 use radroots_event::{GenericEventDraft, contract::AuthorRole}; 1269 use radroots_signing::{ 1270 SignerStatus, 1271 actor::ActorSource, 1272 capability::{CancellationSupport, SignerCapability, SignerKind}, 1273 status::SignerAvailability, 1274 }; 1275 use radroots_transport::{ 1276 DeliveryReceipt, 1277 outcome::DeliveryOutcome, 1278 policy::{SatisfactionClass, TargetPolicy}, 1279 sink::{DeliveryPayload, DeliveryTargetReceipt}, 1280 }; 1281 1282 const AUTHOR: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; 1283 const OTHER_AUTHOR: &str = "e0266e3cfb0d2886f91c73f5f868f3b98273713e5fcd97c081663f5518a4b3af"; 1284 1285 fn request_with( 1286 target_policy: TargetPolicy, 1287 deadline: u64, 1288 cancellation: CancellationPolicy, 1289 ) -> PushRequest { 1290 let plan = AuthoredEventPlan::from_generic( 1291 GenericEventDraft::new( 1292 "radroots.social.geochat.v1", 1293 20_000, 1294 1_800_000_100, 1295 Vec::new(), 1296 "push helper coverage", 1297 AUTHOR, 1298 ) 1299 .unwrap(), 1300 ) 1301 .unwrap(); 1302 let actor = 1303 Actor::from_public_key_hex(AUTHOR, ActorSource::ExplicitPublicKey, [AuthorRole::Any]) 1304 .unwrap(); 1305 PushRequest::new( 1306 SyncId::new([31; 16]).unwrap(), 1307 IdempotencyKey::parse("push-helper-coverage").unwrap(), 1308 actor, 1309 plan, 1310 TargetSet::new(vec![Target::nostr_relay("wss://helper.example").unwrap()]).unwrap(), 1311 SatisfactionPolicy::new(SatisfactionClass::Accepted, target_policy), 1312 deadline, 1313 cancellation, 1314 ) 1315 .unwrap() 1316 } 1317 1318 fn delivery_plan(request: &PushRequest) -> AuthoredDeliveryPlan { 1319 let (_, artifact_id, plan_id) = authored_ids(request.operation_id()).unwrap(); 1320 let delivery_request = DeliveryRequest::new( 1321 delivery_request_id(request.operation_id()), 1322 DeliveryPayload::new( 1323 radroots_event_codec::Codec::decode_signed_event( 1324 r#"{"id":"762bee187e9e645b81ec26ade05a69b5e8398caf527be8de0d9a45311ed0c7a0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1800000100,"kind":0,"tags":[],"content":"{\"display_name\":\"Moss Street Farm\",\"bot\":false,\"website\":\"https://mossstreet.example\",\"picture\":42}","sig":"4290da0bb6422986647bc8cd5f63bd52d49f41e7b665d3b47105b8109183e8d596f322c531d4061df53e1d2b70fda12d5d1c14f3720d7a56d9d0a03746af5109"}"#, 1325 ) 1326 .unwrap(), 1327 ), 1328 request.targets().clone(), 1329 request.satisfaction().clone(), 1330 request.delivery_deadline_unix_ms(), 1331 ) 1332 .unwrap(); 1333 AuthoredDeliveryPlan::new_bound(plan_id, artifact_id, delivery_request, 10).unwrap() 1334 } 1335 1336 fn capability(replay: ReplayCapability) -> SignerCapability { 1337 SignerCapability::new( 1338 SignerKind::Local, 1339 replay, 1340 CancellationSupport::BeforePublication, 1341 false, 1342 false, 1343 ) 1344 } 1345 1346 #[test] 1347 fn push_request_validates_every_binding_and_exposes_redacted_accessors() { 1348 let request = request_with( 1349 TargetPolicy::any(), 1350 1_800_000_300_000, 1351 CancellationPolicy::LocalCooperative, 1352 ); 1353 assert_eq!(request.idempotency_key().as_str(), "push-helper-coverage"); 1354 assert_eq!(request.actor().public_key(), *request.plan().author()); 1355 assert_eq!(request.targets().len(), 1); 1356 assert!(request.satisfaction().targets().is_any()); 1357 assert_eq!(request.delivery_deadline_unix_ms(), 1_800_000_300_000); 1358 assert_eq!(request.cancellation(), CancellationPolicy::LocalCooperative); 1359 let debug = format!("{request:?}"); 1360 assert!(debug.contains("[redacted exact authored plan]")); 1361 assert!(!debug.contains("push helper coverage")); 1362 1363 let invalid_policy = TargetPolicy::required(vec![ 1364 Target::nostr_relay("wss://absent.example") 1365 .unwrap() 1366 .fingerprint() 1367 .clone(), 1368 ]) 1369 .unwrap(); 1370 let mut invalid = request.clone(); 1371 invalid.satisfaction = SatisfactionPolicy::new(SatisfactionClass::Accepted, invalid_policy); 1372 assert!(matches!( 1373 PushRequest::new( 1374 invalid.operation_id, 1375 invalid.idempotency_key, 1376 invalid.actor, 1377 invalid.plan, 1378 invalid.targets, 1379 invalid.satisfaction, 1380 invalid.delivery_deadline_unix_ms, 1381 invalid.cancellation, 1382 ), 1383 Err(Error::InvalidPushRequest) 1384 )); 1385 1386 let valid = request_with( 1387 TargetPolicy::any(), 1388 1_800_000_300_000, 1389 CancellationPolicy::PreservePublishedRequest, 1390 ); 1391 assert!(matches!( 1392 PushRequest::new( 1393 valid.operation_id, 1394 valid.idempotency_key.clone(), 1395 valid.actor.clone(), 1396 valid.plan.clone(), 1397 valid.targets.clone(), 1398 valid.satisfaction.clone(), 1399 0, 1400 valid.cancellation, 1401 ), 1402 Err(Error::InvalidPushRequest) 1403 )); 1404 let wrong_actor = Actor::from_public_key_hex( 1405 OTHER_AUTHOR, 1406 ActorSource::ExplicitPublicKey, 1407 [AuthorRole::Any], 1408 ) 1409 .unwrap(); 1410 assert!(matches!( 1411 PushRequest::new( 1412 valid.operation_id, 1413 valid.idempotency_key, 1414 wrong_actor, 1415 valid.plan, 1416 valid.targets, 1417 valid.satisfaction, 1418 valid.delivery_deadline_unix_ms, 1419 valid.cancellation, 1420 ), 1421 Err(Error::InvalidPushRequest) 1422 )); 1423 } 1424 1425 #[test] 1426 fn signer_and_retry_helpers_cover_every_stable_policy_variant() { 1427 let empty = SignerStatus::new(SignerAvailability::Ready, Vec::new(), None); 1428 assert_eq!( 1429 signer_replay_capability(&empty), 1430 Err(Error::SignerCapabilityUnavailable) 1431 ); 1432 let exact = SignerStatus::new( 1433 SignerAvailability::Ready, 1434 vec![capability(ReplayCapability::ExactReplayByRequestId)], 1435 None, 1436 ); 1437 assert_eq!( 1438 signer_replay_capability(&exact).unwrap(), 1439 ReplayCapability::ExactReplayByRequestId 1440 ); 1441 let mixed = SignerStatus::new( 1442 SignerAvailability::Ready, 1443 vec![ 1444 capability(ReplayCapability::ExactReplayByRequestId), 1445 capability(ReplayCapability::LocalReplaySafe), 1446 ], 1447 None, 1448 ); 1449 assert_eq!( 1450 signer_replay_capability(&mixed).unwrap(), 1451 ReplayCapability::NonReplayable 1452 ); 1453 assert_eq!( 1454 signing_claim_owner(ReplayCapability::ExactReplayByRequestId), 1455 SIGNING_CLAIM_OWNER_EXACT 1456 ); 1457 assert_eq!( 1458 signing_claim_owner(ReplayCapability::LocalReplaySafe), 1459 SIGNING_CLAIM_OWNER_LOCAL 1460 ); 1461 assert_eq!( 1462 signing_claim_owner(ReplayCapability::NonReplayable), 1463 SIGNING_CLAIM_OWNER_NON_REPLAYABLE 1464 ); 1465 1466 let failure = WorkFailure::new( 1467 "retry", 1468 WorkPhase::Signing, 1469 FailureClass::Retryable, 1470 Some(20), 1471 None, 1472 ) 1473 .unwrap(); 1474 assert!(retry_schedule(None, &failure, None).unwrap().is_none()); 1475 let first = retry_schedule(None, &failure, Some(20)).unwrap().unwrap(); 1476 assert_eq!(first.attempt(), NonZeroU32::MIN); 1477 let next_failure = WorkFailure::new( 1478 "retry", 1479 WorkPhase::Signing, 1480 FailureClass::Retryable, 1481 Some(21), 1482 None, 1483 ) 1484 .unwrap(); 1485 assert_eq!( 1486 retry_schedule(Some(&first), &next_failure, Some(21)) 1487 .unwrap() 1488 .unwrap() 1489 .attempt() 1490 .get(), 1491 2 1492 ); 1493 } 1494 1495 #[test] 1496 fn delivery_failure_and_retry_helpers_normalize_all_outcome_classes() { 1497 let request = request_with( 1498 TargetPolicy::any(), 1499 1_800_000_300_000, 1500 CancellationPolicy::PreservePublishedRequest, 1501 ); 1502 let plan = delivery_plan(&request); 1503 let delivery_request = plan.request().unwrap(); 1504 let terminal = SinkFailure::for_request( 1505 delivery_request, 1506 "terminal", 1507 Retryability::Terminal, 1508 None, 1509 None, 1510 Vec::new(), 1511 ) 1512 .unwrap(); 1513 assert_eq!( 1514 normalize_sink_failure(delivery_request, terminal.clone(), 100).unwrap(), 1515 terminal 1516 ); 1517 let retryable = SinkFailure::for_request( 1518 delivery_request, 1519 "retryable", 1520 Retryability::Retryable, 1521 None, 1522 Some("retry".to_owned()), 1523 Vec::new(), 1524 ) 1525 .unwrap(); 1526 let normalized = normalize_sink_failure(delivery_request, retryable, 100).unwrap(); 1527 assert_eq!(normalized.retry_after_unix_ms(), Some(1_100)); 1528 assert_eq!( 1529 normalize_sink_failure(delivery_request, normalized.clone(), 100).unwrap(), 1530 normalized 1531 ); 1532 let past = SinkFailure::for_request( 1533 delivery_request, 1534 "past", 1535 Retryability::Retryable, 1536 Some(99), 1537 None, 1538 Vec::new(), 1539 ) 1540 .unwrap(); 1541 assert_eq!( 1542 normalize_sink_failure(delivery_request, past, 100) 1543 .unwrap() 1544 .retry_after_unix_ms(), 1545 Some(101) 1546 ); 1547 1548 let receipt = DeliveryReceipt::for_request( 1549 delivery_request, 1550 delivery_request 1551 .target_set() 1552 .targets() 1553 .iter() 1554 .cloned() 1555 .map(|target| { 1556 DeliveryTargetReceipt::attempted(target, DeliveryOutcome::unavailable()) 1557 }) 1558 .collect(), 1559 ) 1560 .unwrap(); 1561 let receipt_outcome = DeliveryAttemptOutcome::Receipt(receipt); 1562 assert!( 1563 delivery_retry_schedule( 1564 NonZeroU32::new(plan.attempt_count() + 1).unwrap(), 1565 &receipt_outcome, 1566 SatisfactionState::Satisfied, 1567 100 1568 ) 1569 .unwrap() 1570 .is_none() 1571 ); 1572 let schedule = delivery_retry_schedule( 1573 NonZeroU32::new(plan.attempt_count() + 1).unwrap(), 1574 &receipt_outcome, 1575 SatisfactionState::Pending, 1576 100, 1577 ) 1578 .unwrap() 1579 .unwrap(); 1580 assert_eq!(schedule.not_before_unix_ms(), 1_100); 1581 let retry_outcome = DeliveryAttemptOutcome::SinkFailure(normalized); 1582 assert!( 1583 delivery_retry_schedule( 1584 NonZeroU32::new(plan.attempt_count() + 1).unwrap(), 1585 &retry_outcome, 1586 SatisfactionState::Pending, 1587 100 1588 ) 1589 .unwrap() 1590 .is_some() 1591 ); 1592 let terminal_outcome = DeliveryAttemptOutcome::SinkFailure(terminal); 1593 assert!( 1594 delivery_retry_schedule( 1595 NonZeroU32::new(plan.attempt_count() + 1).unwrap(), 1596 &terminal_outcome, 1597 SatisfactionState::Pending, 1598 100 1599 ) 1600 .unwrap() 1601 .is_none() 1602 ); 1603 1604 for policy in [ 1605 TargetPolicy::any(), 1606 TargetPolicy::all(), 1607 TargetPolicy::quorum(1).unwrap(), 1608 TargetPolicy::required(vec![request.targets().targets()[0].fingerprint().clone()]) 1609 .unwrap(), 1610 ] { 1611 let request = request_with( 1612 policy, 1613 1_800_000_300_000, 1614 CancellationPolicy::LocalCooperative, 1615 ); 1616 assert_ne!( 1617 authored_push_input_digest(&request).unwrap().as_bytes(), 1618 &[0; 32] 1619 ); 1620 } 1621 } 1622 1623 #[test] 1624 fn storage_error_maps_are_explicit_and_fail_closed() { 1625 assert_eq!( 1626 map_storage_error(radroots_storage::Error::SpaceInsufficient), 1627 Error::StorageSpaceInsufficient 1628 ); 1629 assert_eq!( 1630 map_claim_error(radroots_storage::Error::SpaceInsufficient), 1631 Error::StorageSpaceInsufficient 1632 ); 1633 assert_eq!( 1634 map_storage_error(radroots_storage::Error::AtomicCommitConflict), 1635 Error::StorageConflict 1636 ); 1637 assert_eq!( 1638 map_storage_error(radroots_storage::Error::BackendUnavailable), 1639 Error::StorageFailed 1640 ); 1641 assert_eq!( 1642 map_claim_error(radroots_storage::Error::DeliveryPlanClaimConflict), 1643 Error::WorkClaimConflict 1644 ); 1645 assert_eq!( 1646 map_claim_error(radroots_storage::Error::BackendUnavailable), 1647 Error::StorageFailed 1648 ); 1649 assert_eq!( 1650 delivery_request_id(SyncId::new([0xab; 16]).unwrap()).len(), 1651 37 1652 ); 1653 } 1654 }