sink.rs (20860B)
1 //! Outbound event delivery SPI and bounded request models. 2 3 use crate::{ 4 Error, 5 outcome::{DeliveryOutcome, Retryability, validate_delivery_code, validate_delivery_message}, 6 policy::{SatisfactionPolicy, SatisfactionState, evaluate_satisfaction}, 7 source::BoxFuture, 8 target::{Target, TargetSet}, 9 }; 10 use alloc::{ 11 boxed::Box, 12 collections::{BTreeMap, BTreeSet}, 13 string::{String, ToString}, 14 vec::Vec, 15 }; 16 use radroots_event::SignedEvent; 17 18 pub use crate::status::SinkStatus; 19 20 /// Maximum encoded delivery request identity length. 21 pub const DELIVERY_REQUEST_ID_MAX_BYTES: usize = 256; 22 23 /// Sink-wide typed failure retaining safe retry and partial-target evidence. 24 #[cfg_attr(feature = "serde", derive(serde::Serialize))] 25 #[derive(Clone, Debug, Eq, PartialEq)] 26 pub struct SinkFailure { 27 request_id: DeliveryRequestId, 28 target_set: TargetSet, 29 code: String, 30 retryability: Retryability, 31 retry_after_unix_ms: Option<u64>, 32 message: Option<String>, 33 partial_evidence: Vec<DeliveryTargetReceipt>, 34 } 35 36 /// Validated caller identity for one delivery operation. 37 #[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 38 pub struct DeliveryRequestId(String); 39 40 impl DeliveryRequestId { 41 /// Parses a non-empty, bounded, printable request identity. 42 pub fn parse(value: impl Into<String>) -> Result<Self, Error> { 43 let value = value.into(); 44 if value.is_empty() { 45 return Err(Error::EmptyDeliveryRequestId); 46 } 47 if value.len() > DELIVERY_REQUEST_ID_MAX_BYTES 48 || value != value.trim() 49 || value.chars().any(char::is_control) 50 { 51 return Err(Error::InvalidDeliveryRequestId); 52 } 53 Ok(Self(value)) 54 } 55 56 /// Returns the validated request identity. 57 pub fn as_str(&self) -> &str { 58 self.0.as_str() 59 } 60 } 61 62 /// Transport-neutral outbound event payload. 63 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 64 #[cfg_attr(feature = "serde", serde(deny_unknown_fields))] 65 #[derive(Clone, Debug, Eq, PartialEq)] 66 pub struct DeliveryPayload { 67 event: SignedEvent, 68 } 69 70 impl DeliveryPayload { 71 /// Wraps an ID-checked signed event for delivery. 72 pub const fn new(event: SignedEvent) -> Self { 73 Self { event } 74 } 75 76 /// Returns the signed event. Signature verification remains a caller concern. 77 pub const fn event(&self) -> &SignedEvent { 78 &self.event 79 } 80 } 81 82 /// Bounded multi-target delivery request. 83 #[cfg_attr(feature = "serde", derive(serde::Serialize))] 84 #[derive(Clone, Debug, Eq, PartialEq)] 85 pub struct DeliveryRequest { 86 request_id: DeliveryRequestId, 87 payload: DeliveryPayload, 88 target_set: TargetSet, 89 satisfaction: SatisfactionPolicy, 90 deadline_unix_ms: u64, 91 } 92 93 impl DeliveryRequest { 94 /// Creates and validates one explicit delivery request. 95 pub fn new( 96 request_id: impl Into<String>, 97 payload: DeliveryPayload, 98 target_set: TargetSet, 99 satisfaction: SatisfactionPolicy, 100 deadline_unix_ms: u64, 101 ) -> Result<Self, Error> { 102 if deadline_unix_ms == 0 { 103 return Err(Error::InvalidDeliveryDeadline); 104 } 105 satisfaction.validate_for(&target_set)?; 106 Ok(Self { 107 request_id: DeliveryRequestId::parse(request_id)?, 108 payload, 109 target_set, 110 satisfaction, 111 deadline_unix_ms, 112 }) 113 } 114 115 /// Returns the request identity. 116 pub const fn request_id(&self) -> &DeliveryRequestId { 117 &self.request_id 118 } 119 120 /// Returns the signed event payload. 121 pub const fn payload(&self) -> &DeliveryPayload { 122 &self.payload 123 } 124 125 /// Returns the exact non-empty target set. 126 pub const fn target_set(&self) -> &TargetSet { 127 &self.target_set 128 } 129 130 /// Returns the requested success and target policy. 131 pub const fn satisfaction(&self) -> &SatisfactionPolicy { 132 &self.satisfaction 133 } 134 135 /// Returns the absolute Unix deadline in milliseconds. 136 pub const fn deadline_unix_ms(&self) -> u64 { 137 self.deadline_unix_ms 138 } 139 140 /// Validates a nonempty bounded subset of the exact original targets. 141 /// Selection never changes this request's payload, policy or identity. 142 pub fn validate_target_selection(&self, selected: &TargetSet) -> Result<(), Error> { 143 if selected 144 .targets() 145 .iter() 146 .all(|target| self.target_set.targets().contains(target)) 147 { 148 Ok(()) 149 } else { 150 Err(Error::InvalidDeliveryTargetSelection) 151 } 152 } 153 } 154 155 /// Normalized result for one requested target. 156 #[cfg_attr(feature = "serde", derive(serde::Serialize))] 157 #[derive(Clone, Debug, Eq, PartialEq)] 158 pub struct DeliveryTargetReceipt { 159 target: Target, 160 attempted: bool, 161 outcome: DeliveryOutcome, 162 } 163 164 impl DeliveryTargetReceipt { 165 /// Records the outcome of an attempted target. 166 pub const fn attempted(target: Target, outcome: DeliveryOutcome) -> Self { 167 Self { 168 target, 169 attempted: true, 170 outcome, 171 } 172 } 173 174 /// Records an unattempted target and its normalized failure reason. 175 pub fn skipped(target: Target, outcome: DeliveryOutcome) -> Result<Self, Error> { 176 if outcome.satisfies(crate::policy::SatisfactionClass::Accepted) { 177 return Err(Error::DeliveryTargetReceiptAttemptMismatch); 178 } 179 Ok(Self { 180 target, 181 attempted: false, 182 outcome, 183 }) 184 } 185 186 /// Returns the exact target. 187 pub const fn target(&self) -> &Target { 188 &self.target 189 } 190 191 /// Whether the adapter attempted remote publication. 192 pub const fn was_attempted(&self) -> bool { 193 self.attempted 194 } 195 196 /// Returns normalized target outcome data. 197 pub const fn outcome(&self) -> &DeliveryOutcome { 198 &self.outcome 199 } 200 201 fn validate(&self) -> Result<(), Error> { 202 self.outcome.validate()?; 203 if !self.attempted 204 && self 205 .outcome 206 .satisfies(crate::policy::SatisfactionClass::Accepted) 207 { 208 return Err(Error::DeliveryTargetReceiptAttemptMismatch); 209 } 210 Ok(()) 211 } 212 } 213 214 /// Request-bound per-target delivery receipt. 215 #[cfg_attr(feature = "serde", derive(serde::Serialize))] 216 #[derive(Clone, Debug, Eq, PartialEq)] 217 pub struct DeliveryReceipt { 218 request_id: DeliveryRequestId, 219 target_set: TargetSet, 220 target_receipts: Vec<DeliveryTargetReceipt>, 221 } 222 223 impl DeliveryReceipt { 224 /// Creates a complete result set in the original request target order. 225 pub fn for_request( 226 request: &DeliveryRequest, 227 target_receipts: Vec<DeliveryTargetReceipt>, 228 ) -> Result<Self, Error> { 229 Self::new( 230 request.request_id.clone(), 231 request.target_set.clone(), 232 target_receipts, 233 ) 234 } 235 236 fn new( 237 request_id: DeliveryRequestId, 238 target_set: TargetSet, 239 target_receipts: Vec<DeliveryTargetReceipt>, 240 ) -> Result<Self, Error> { 241 let mut by_fingerprint = BTreeMap::new(); 242 for receipt in target_receipts { 243 receipt.validate()?; 244 let fingerprint = receipt.target().fingerprint().as_str().to_string(); 245 if !target_set 246 .targets() 247 .iter() 248 .any(|target| target.fingerprint().as_str() == fingerprint) 249 { 250 return Err(Error::UnexpectedDeliveryTargetReceipt); 251 } 252 if by_fingerprint.insert(fingerprint, receipt).is_some() { 253 return Err(Error::DuplicateDeliveryTargetReceipt); 254 } 255 } 256 257 let mut ordered = Vec::with_capacity(target_set.len()); 258 for target in target_set.targets() { 259 let Some(receipt) = by_fingerprint.remove(target.fingerprint().as_str()) else { 260 return Err(Error::MissingDeliveryTargetReceipt); 261 }; 262 ordered.push(receipt); 263 } 264 Ok(Self { 265 request_id, 266 target_set, 267 target_receipts: ordered, 268 }) 269 } 270 271 /// Validates the receipt against the exact request identity and targets. 272 pub fn validate_for_request(&self, request: &DeliveryRequest) -> Result<(), Error> { 273 if &self.request_id != request.request_id() { 274 return Err(Error::DeliveryReceiptRequestIdMismatch); 275 } 276 if self.target_set != *request.target_set() { 277 return Err(Error::DeliveryReceiptTargetSetMismatch); 278 } 279 let rebuilt = Self::for_request(request, self.target_receipts.clone())?; 280 if rebuilt != *self { 281 return Err(Error::DeliveryReceiptTargetSetMismatch); 282 } 283 Ok(()) 284 } 285 286 /// Returns whether the receipt satisfies the request's exact policy. 287 pub fn is_satisfied(&self, request: &DeliveryRequest) -> Result<bool, Error> { 288 Ok(matches!( 289 self.satisfaction(request)?, 290 SatisfactionState::Satisfied 291 )) 292 } 293 294 /// Evaluates this receipt as satisfied, pending, or exhausted. 295 pub fn satisfaction(&self, request: &DeliveryRequest) -> Result<SatisfactionState, Error> { 296 self.validate_for_request(request)?; 297 evaluate_satisfaction( 298 request.satisfaction(), 299 request.target_set(), 300 self.target_receipts 301 .iter() 302 .map(|receipt| (receipt.target().fingerprint(), receipt.outcome())), 303 ) 304 } 305 306 /// Returns the request identity. 307 pub const fn request_id(&self) -> &DeliveryRequestId { 308 &self.request_id 309 } 310 311 /// Returns per-target results in request order. 312 pub fn target_receipts(&self) -> &[DeliveryTargetReceipt] { 313 self.target_receipts.as_slice() 314 } 315 } 316 317 impl SinkFailure { 318 /// Creates a request-bound sink-wide failure with validated partial evidence. 319 pub fn for_request( 320 request: &DeliveryRequest, 321 code: impl Into<String>, 322 retryability: Retryability, 323 retry_after_unix_ms: Option<u64>, 324 message: Option<String>, 325 partial_evidence: Vec<DeliveryTargetReceipt>, 326 ) -> Result<Self, Error> { 327 let failure = Self { 328 request_id: request.request_id.clone(), 329 target_set: request.target_set.clone(), 330 code: code.into(), 331 retryability, 332 retry_after_unix_ms, 333 message, 334 partial_evidence, 335 }; 336 failure.validate_for_request(request)?; 337 Ok(failure) 338 } 339 340 /// Returns a terminal adapter-contract failure for an exact request. 341 pub fn invalid_contract(request: &DeliveryRequest) -> Self { 342 Self::for_request( 343 request, 344 "invalid_transport_contract", 345 Retryability::Terminal, 346 None, 347 Some("transport adapter returned invalid evidence".to_string()), 348 Vec::new(), 349 ) 350 .expect("static sink failure is valid") 351 } 352 353 /// Validates identity, retry timing, and bounded partial evidence. 354 pub fn validate_for_request(&self, request: &DeliveryRequest) -> Result<(), Error> { 355 if self.request_id != *request.request_id() { 356 return Err(Error::DeliveryReceiptRequestIdMismatch); 357 } 358 if self.target_set != *request.target_set() { 359 return Err(Error::DeliveryReceiptTargetSetMismatch); 360 } 361 validate_delivery_code(self.code.as_str())?; 362 if matches!(self.retryability, Retryability::NotApplicable) 363 || matches!(self.retry_after_unix_ms, Some(0)) 364 || (self.retry_after_unix_ms.is_some() 365 && !matches!(self.retryability, Retryability::Retryable)) 366 { 367 return Err(Error::InvalidDeliveryOutcome); 368 } 369 if let Some(message) = &self.message { 370 validate_delivery_message(message)?; 371 } 372 let mut observed = BTreeSet::new(); 373 for receipt in &self.partial_evidence { 374 receipt.validate()?; 375 if !request 376 .target_set() 377 .contains(receipt.target().fingerprint()) 378 { 379 return Err(Error::UnexpectedDeliveryTargetReceipt); 380 } 381 if !observed.insert(receipt.target().fingerprint().as_str()) { 382 return Err(Error::DuplicateDeliveryTargetReceipt); 383 } 384 } 385 Ok(()) 386 } 387 388 /// Returns the stable normalized failure code. 389 pub fn code(&self) -> &str { 390 self.code.as_str() 391 } 392 393 /// Returns whether retrying the same request may be useful. 394 pub const fn retryability(&self) -> Retryability { 395 self.retryability 396 } 397 398 /// Returns the earliest absolute Unix millisecond retry time, when supplied. 399 pub const fn retry_after_unix_ms(&self) -> Option<u64> { 400 self.retry_after_unix_ms 401 } 402 403 /// Returns bounded caller-safe diagnostic detail. 404 pub fn message(&self) -> Option<&str> { 405 self.message.as_deref() 406 } 407 408 /// Returns safe target evidence collected before the sink-wide failure. 409 pub fn partial_evidence(&self) -> &[DeliveryTargetReceipt] { 410 self.partial_evidence.as_slice() 411 } 412 } 413 414 /// Host SPI for outbound event delivery. 415 /// 416 /// This trait supports external implementations and is dyn-compatible. Its 417 /// futures are `Send`; implementations must not borrow request data after a 418 /// future completes. `status` observes sink state and does not initiate 419 /// delivery. `deliver` performs only the attempts authorized by its request, 420 /// returns partial success per target, and owns no hidden retry loop. 421 /// 422 /// Dropping a returned future requests cancellation. If it is dropped before 423 /// a remote request is published, the implementation must leave no remote 424 /// operation behind. Once publication may have occurred, cancellation cannot 425 /// claim rollback; a later observation may report the remote outcome. An 426 /// explicit request deadline bounds work independently of future cancellation. 427 pub trait EventSink: Send + Sync { 428 /// Returns the sink's current runtime status. 429 fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>>; 430 431 /// Delivers an event according to the request's bounded target policy. 432 fn deliver( 433 &self, 434 request: DeliveryRequest, 435 ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>>; 436 437 /// Attempts only selected targets while retaining the full request binding. 438 /// 439 /// Implementations must validate the exact subset before I/O, keep every 440 /// original target in successful receipts, and report unselected targets as 441 /// unattempted. The default supports only a selection of every target; a 442 /// proper subset fails closed without calling ordinary delivery. 443 fn deliver_selected( 444 &self, 445 request: DeliveryRequest, 446 selected: TargetSet, 447 ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { 448 Box::pin(async move { 449 if request.validate_target_selection(&selected).is_err() { 450 return Err(SinkFailure::invalid_contract(&request)); 451 } 452 if selected.len() == request.target_set().len() { 453 return self.deliver(request).await; 454 } 455 Err(SinkFailure::for_request( 456 &request, 457 "target_selection_unsupported", 458 Retryability::Terminal, 459 None, 460 None, 461 Vec::new(), 462 ) 463 .expect("static unsupported selection failure is valid")) 464 }) 465 } 466 } 467 468 #[cfg(feature = "serde")] 469 mod serde_impl { 470 use super::*; 471 472 impl serde::Serialize for DeliveryRequestId { 473 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error> 474 where 475 S: serde::Serializer, 476 { 477 serializer.serialize_str(self.as_str()) 478 } 479 } 480 481 impl<'de> serde::Deserialize<'de> for DeliveryRequestId { 482 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> 483 where 484 D: serde::Deserializer<'de>, 485 { 486 let value = <String as serde::Deserialize>::deserialize(deserializer)?; 487 Self::parse(value).map_err(serde::de::Error::custom) 488 } 489 } 490 491 #[derive(serde::Deserialize)] 492 #[serde(deny_unknown_fields)] 493 struct DeliveryRequestWire { 494 request_id: String, 495 payload: DeliveryPayload, 496 target_set: TargetSet, 497 satisfaction: SatisfactionPolicy, 498 deadline_unix_ms: u64, 499 } 500 501 impl<'de> serde::Deserialize<'de> for DeliveryRequest { 502 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> 503 where 504 D: serde::Deserializer<'de>, 505 { 506 let wire = DeliveryRequestWire::deserialize(deserializer)?; 507 Self::new( 508 wire.request_id, 509 wire.payload, 510 wire.target_set, 511 wire.satisfaction, 512 wire.deadline_unix_ms, 513 ) 514 .map_err(serde::de::Error::custom) 515 } 516 } 517 518 #[derive(serde::Deserialize)] 519 #[serde(deny_unknown_fields)] 520 struct DeliveryTargetReceiptWire { 521 target: Target, 522 attempted: bool, 523 outcome: DeliveryOutcome, 524 } 525 526 impl<'de> serde::Deserialize<'de> for DeliveryTargetReceipt { 527 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> 528 where 529 D: serde::Deserializer<'de>, 530 { 531 let wire = DeliveryTargetReceiptWire::deserialize(deserializer)?; 532 let receipt = Self { 533 target: wire.target, 534 attempted: wire.attempted, 535 outcome: wire.outcome, 536 }; 537 receipt.validate().map_err(serde::de::Error::custom)?; 538 Ok(receipt) 539 } 540 } 541 542 #[derive(serde::Deserialize)] 543 #[serde(deny_unknown_fields)] 544 struct DeliveryReceiptWire { 545 request_id: DeliveryRequestId, 546 target_set: TargetSet, 547 target_receipts: Vec<DeliveryTargetReceipt>, 548 } 549 550 impl<'de> serde::Deserialize<'de> for DeliveryReceipt { 551 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> 552 where 553 D: serde::Deserializer<'de>, 554 { 555 let wire = DeliveryReceiptWire::deserialize(deserializer)?; 556 Self::new(wire.request_id, wire.target_set, wire.target_receipts) 557 .map_err(serde::de::Error::custom) 558 } 559 } 560 561 #[derive(serde::Deserialize)] 562 #[serde(deny_unknown_fields)] 563 struct SinkFailureWire { 564 request_id: DeliveryRequestId, 565 target_set: TargetSet, 566 code: String, 567 retryability: Retryability, 568 retry_after_unix_ms: Option<u64>, 569 message: Option<String>, 570 partial_evidence: Vec<DeliveryTargetReceipt>, 571 } 572 573 impl<'de> serde::Deserialize<'de> for SinkFailure { 574 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> 575 where 576 D: serde::Deserializer<'de>, 577 { 578 let wire = SinkFailureWire::deserialize(deserializer)?; 579 let failure = Self { 580 request_id: wire.request_id, 581 target_set: wire.target_set, 582 code: wire.code, 583 retryability: wire.retryability, 584 retry_after_unix_ms: wire.retry_after_unix_ms, 585 message: wire.message, 586 partial_evidence: wire.partial_evidence, 587 }; 588 validate_delivery_code(failure.code.as_str()).map_err(serde::de::Error::custom)?; 589 if matches!(failure.retryability, Retryability::NotApplicable) 590 || matches!(failure.retry_after_unix_ms, Some(0)) 591 || (failure.retry_after_unix_ms.is_some() 592 && !matches!(failure.retryability, Retryability::Retryable)) 593 { 594 return Err(serde::de::Error::custom(Error::InvalidDeliveryOutcome)); 595 } 596 if let Some(message) = &failure.message { 597 validate_delivery_message(message).map_err(serde::de::Error::custom)?; 598 } 599 let mut observed = BTreeSet::new(); 600 for receipt in &failure.partial_evidence { 601 receipt.validate().map_err(serde::de::Error::custom)?; 602 if !failure.target_set.contains(receipt.target().fingerprint()) { 603 return Err(serde::de::Error::custom( 604 Error::UnexpectedDeliveryTargetReceipt, 605 )); 606 } 607 if !observed.insert(receipt.target().fingerprint().as_str()) { 608 return Err(serde::de::Error::custom( 609 Error::DuplicateDeliveryTargetReceipt, 610 )); 611 } 612 } 613 Ok(failure) 614 } 615 } 616 }