1use compact_str::CompactString;
8use serde::{Deserialize, Serialize};
9
10use crate::mm::handle::ObjectDescriptor;
11use crate::scheduler::tcb::TaskId;
12use crate::types::capability::Capability;
13
14#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
15#[serde(deny_unknown_fields)]
16pub struct OperationAddress {
17 pub operation_id: CompactString,
18 pub task_id: TaskId,
19}
20
21impl OperationAddress {
22 pub fn new(operation_id: impl Into<CompactString>, task_id: impl Into<TaskId>) -> Self {
23 Self {
24 operation_id: operation_id.into(),
25 task_id: task_id.into(),
26 }
27 }
28}
29
30#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
31#[serde(deny_unknown_fields)]
32pub struct PayloadLocator {
33 pub reference: String,
35 pub digest: String,
36 pub size: u64,
37}
38
39#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
40#[serde(rename_all = "snake_case")]
41pub enum PayloadAvailability {
42 Available,
43 Unavailable,
44}
45
46#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
47#[serde(deny_unknown_fields)]
48pub struct CrossOperationMessage {
49 pub id: CompactString,
50 pub from: OperationAddress,
51 pub to: OperationAddress,
52 pub kind: CompactString,
53 pub object: ObjectDescriptor,
54 pub payload: PayloadLocator,
55 pub sequence: u64,
57 pub sent_at_turn: u32,
58 pub ttl_turns: u32,
59 #[serde(default)]
61 pub delegated_capabilities: Vec<Capability>,
62}
63
64#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
65#[serde(deny_unknown_fields)]
66pub struct OperationRegistration {
67 pub address: OperationAddress,
68 pub capabilities: Vec<Capability>,
69 #[serde(default)]
70 pub cancelled: bool,
71}
72
73#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
74#[serde(rename_all = "snake_case")]
75pub enum DeliveryState {
76 Queued,
77 InFlight,
78 Delivered,
79 Failed,
80 Expired,
81 Cancelled,
82}
83
84impl DeliveryState {
85 fn is_pending(self) -> bool {
86 matches!(self, Self::Queued | Self::InFlight)
87 }
88
89 fn is_terminal(self) -> bool {
90 !self.is_pending()
91 }
92}
93
94#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
95#[serde(rename_all = "snake_case")]
96pub enum DeliverySettlement {
97 Delivered,
98 RetryableFailure,
99 PermanentFailure,
100}
101
102#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
103#[serde(rename_all = "snake_case")]
104pub enum DeliveryFailure {
105 UnknownOperation,
106 Cancelled,
107 SourceDoesNotOwnObject,
108 InvalidPayloadLocator,
109 PayloadUnavailable,
110 SenderPermissionDenied,
111 ReceiverPermissionDenied,
112 CapabilityAttenuationDenied,
113 ExpiredCapabilityLease,
114 ExpiredTtl,
115 OutOfOrder,
116 Backpressure,
117 ConflictingDuplicate,
118}
119
120impl DeliveryFailure {
121 pub fn is_retryable(self) -> bool {
122 matches!(self, Self::PayloadUnavailable | Self::Backpressure)
123 }
124}
125
126#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
127#[serde(tag = "outcome", rename_all = "snake_case")]
128pub enum RouteOutcome {
129 Accepted { state: DeliveryState },
130 Duplicate { state: DeliveryState },
131 Rejected { failure: DeliveryFailure },
132}
133
134#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
135#[serde(deny_unknown_fields)]
136pub struct DurableDelivery {
137 pub message: CrossOperationMessage,
138 pub state: DeliveryState,
139 pub attempts: u32,
140}
141
142#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
143#[serde(deny_unknown_fields)]
144pub struct ReceiverQueueSnapshot {
145 pub receiver: OperationAddress,
146 pub message_ids: Vec<CompactString>,
147}
148
149#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
150#[serde(deny_unknown_fields)]
151pub struct ReceiverSequenceSnapshot {
152 pub receiver: OperationAddress,
153 pub last_sequence: u64,
154}
155
156#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
157#[serde(deny_unknown_fields)]
158pub struct CrossOperationRouterSnapshot {
159 pub max_pending_per_recipient: usize,
160 pub operations: Vec<OperationRegistration>,
161 pub deliveries: Vec<DurableDelivery>,
162 pub queues: Vec<ReceiverQueueSnapshot>,
163 pub last_sequences: Vec<ReceiverSequenceSnapshot>,
164}
165
166#[derive(Debug, Clone, PartialEq, Eq)]
167pub enum SnapshotError {
168 DuplicateOperation(OperationAddress),
169 DuplicateDelivery(CompactString),
170 DuplicateQueue(OperationAddress),
171 DuplicateSequence(OperationAddress),
172 SequenceForUnknownOperation(OperationAddress),
173 DeliveryForUnknownOperation(CompactString),
174 UnknownQueuedDelivery(CompactString),
175 QueuedDeliveryForWrongReceiver(CompactString),
176 NonQueuedDelivery(CompactString),
177 DuplicateQueuedDelivery(CompactString),
178 OutOfOrderQueue(CompactString),
179}
180
181#[derive(Debug, Clone, Default, PartialEq, Eq)]
184pub struct CrossOperationRouter {
185 max_pending_per_recipient: usize,
186 operations: std::collections::BTreeMap<OperationAddress, OperationRegistration>,
187 deliveries: std::collections::BTreeMap<CompactString, DurableDelivery>,
188 queues: std::collections::BTreeMap<OperationAddress, std::collections::VecDeque<CompactString>>,
189 last_sequences: std::collections::BTreeMap<OperationAddress, u64>,
190}
191
192impl CrossOperationRouter {
193 pub fn new(max_pending_per_recipient: usize) -> Self {
194 Self {
195 max_pending_per_recipient,
196 ..Self::default()
197 }
198 }
199
200 pub fn register(&mut self, registration: OperationRegistration) -> bool {
201 self.operations
202 .insert(registration.address.clone(), registration)
203 .is_none()
204 }
205
206 pub fn operation(&self, address: &OperationAddress) -> Option<&OperationRegistration> {
207 self.operations.get(address)
208 }
209
210 pub fn delivery(&self, id: &str) -> Option<&DurableDelivery> {
211 self.deliveries.get(id)
212 }
213
214 pub fn route(
215 &mut self,
216 message: CrossOperationMessage,
217 now_turn: u32,
218 payload_availability: PayloadAvailability,
219 ) -> RouteOutcome {
220 if let Some(existing) = self.deliveries.get(message.id.as_str()) {
221 return if existing.message == message {
222 RouteOutcome::Duplicate {
223 state: existing.state,
224 }
225 } else {
226 RouteOutcome::Rejected {
227 failure: DeliveryFailure::ConflictingDuplicate,
228 }
229 };
230 }
231
232 let Some(source) = self.operations.get(&message.from) else {
233 return rejected(DeliveryFailure::UnknownOperation);
234 };
235 let Some(receiver) = self.operations.get(&message.to) else {
236 return rejected(DeliveryFailure::UnknownOperation);
237 };
238 if source.cancelled || receiver.cancelled {
239 return rejected(DeliveryFailure::Cancelled);
240 }
241 if message.object.owner != message.from.task_id {
242 return rejected(DeliveryFailure::SourceDoesNotOwnObject);
243 }
244 if !valid_payload_locator(&message.object, &message.payload) {
245 return rejected(DeliveryFailure::InvalidPayloadLocator);
246 }
247 if payload_availability == PayloadAvailability::Unavailable {
248 return rejected(DeliveryFailure::PayloadUnavailable);
249 }
250 if expired(message.sent_at_turn, message.ttl_turns, now_turn) {
251 return rejected(DeliveryFailure::ExpiredTtl);
252 }
253 if matching_expired_lease(&source.capabilities, "share", &message.object, now_turn)
254 && !object_action_allowed(&source.capabilities, "share", &message.object, now_turn)
255 || matching_expired_lease(&receiver.capabilities, "read", &message.object, now_turn)
256 && !object_action_allowed(&receiver.capabilities, "read", &message.object, now_turn)
257 {
258 return rejected(DeliveryFailure::ExpiredCapabilityLease);
259 }
260 if !object_action_allowed(&source.capabilities, "share", &message.object, now_turn) {
261 return rejected(DeliveryFailure::SenderPermissionDenied);
262 }
263 if !object_action_allowed(&receiver.capabilities, "read", &message.object, now_turn) {
264 return rejected(DeliveryFailure::ReceiverPermissionDenied);
265 }
266 if !delegation_is_attenuated(
267 &message.delegated_capabilities,
268 &source.capabilities,
269 &receiver.capabilities,
270 now_turn,
271 ) {
272 return rejected(DeliveryFailure::CapabilityAttenuationDenied);
273 }
274 if self
275 .last_sequences
276 .get(&message.to)
277 .is_some_and(|last| message.sequence <= *last)
278 {
279 return rejected(DeliveryFailure::OutOfOrder);
280 }
281 if self.pending_for(&message.to) >= self.max_pending_per_recipient {
282 return rejected(DeliveryFailure::Backpressure);
283 }
284
285 self.last_sequences
286 .insert(message.to.clone(), message.sequence);
287 self.queues
288 .entry(message.to.clone())
289 .or_default()
290 .push_back(message.id.clone());
291 self.deliveries.insert(
292 message.id.clone(),
293 DurableDelivery {
294 message,
295 state: DeliveryState::Queued,
296 attempts: 0,
297 },
298 );
299 RouteOutcome::Accepted {
300 state: DeliveryState::Queued,
301 }
302 }
303
304 pub fn dispatch_next(
307 &mut self,
308 receiver: &OperationAddress,
309 now_turn: u32,
310 ) -> Option<CrossOperationMessage> {
311 if self
312 .operations
313 .get(receiver)
314 .is_none_or(|operation| operation.cancelled)
315 || self.has_in_flight_for(receiver)
316 {
317 return None;
318 }
319
320 loop {
321 let id = self.queues.get_mut(receiver)?.pop_front()?;
322 let delivery = self.deliveries.get_mut(id.as_str())?;
323 if delivery.state != DeliveryState::Queued {
324 continue;
325 }
326 if expired(
327 delivery.message.sent_at_turn,
328 delivery.message.ttl_turns,
329 now_turn,
330 ) {
331 delivery.state = DeliveryState::Expired;
332 continue;
333 }
334 delivery.state = DeliveryState::InFlight;
335 delivery.attempts = delivery.attempts.saturating_add(1);
336 return Some(delivery.message.clone());
337 }
338 }
339
340 pub fn settle(
341 &mut self,
342 message_id: &str,
343 settlement: DeliverySettlement,
344 now_turn: u32,
345 ) -> Option<DeliveryState> {
346 let (receiver, state) = {
347 let delivery = self.deliveries.get(message_id)?;
348 (delivery.message.to.clone(), delivery.state)
349 };
350 if state.is_terminal() || state != DeliveryState::InFlight {
351 return Some(state);
352 }
353
354 let next_state = match settlement {
355 DeliverySettlement::Delivered => DeliveryState::Delivered,
356 DeliverySettlement::PermanentFailure => DeliveryState::Failed,
357 DeliverySettlement::RetryableFailure => {
358 let delivery = self.deliveries.get(message_id)?;
359 if expired(
360 delivery.message.sent_at_turn,
361 delivery.message.ttl_turns,
362 now_turn,
363 ) || self
364 .operations
365 .get(&receiver)
366 .is_none_or(|operation| operation.cancelled)
367 {
368 if self
369 .operations
370 .get(&receiver)
371 .is_some_and(|operation| operation.cancelled)
372 {
373 DeliveryState::Cancelled
374 } else {
375 DeliveryState::Expired
376 }
377 } else {
378 self.queues
379 .entry(receiver)
380 .or_default()
381 .push_front(CompactString::from(message_id));
382 DeliveryState::Queued
383 }
384 }
385 };
386 self.deliveries.get_mut(message_id)?.state = next_state;
387 Some(next_state)
388 }
389
390 pub fn cancel_operation(&mut self, address: &OperationAddress) -> usize {
393 let Some(operation) = self.operations.get_mut(address) else {
394 return 0;
395 };
396 operation.cancelled = true;
397 let mut cancelled = 0;
398 for delivery in self.deliveries.values_mut() {
399 if (delivery.message.from == *address || delivery.message.to == *address)
400 && delivery.state.is_pending()
401 {
402 delivery.state = DeliveryState::Cancelled;
403 cancelled += 1;
404 }
405 }
406 for queue in self.queues.values_mut() {
407 queue.retain(|id| {
408 self.deliveries
409 .get(id.as_str())
410 .is_some_and(|delivery| delivery.state == DeliveryState::Queued)
411 });
412 }
413 cancelled
414 }
415
416 pub fn snapshot(&self) -> CrossOperationRouterSnapshot {
417 CrossOperationRouterSnapshot {
418 max_pending_per_recipient: self.max_pending_per_recipient,
419 operations: self.operations.values().cloned().collect(),
420 deliveries: self.deliveries.values().cloned().collect(),
421 queues: self
422 .queues
423 .iter()
424 .map(|(receiver, message_ids)| ReceiverQueueSnapshot {
425 receiver: receiver.clone(),
426 message_ids: message_ids.iter().cloned().collect(),
427 })
428 .collect(),
429 last_sequences: self
430 .last_sequences
431 .iter()
432 .map(|(receiver, last_sequence)| ReceiverSequenceSnapshot {
433 receiver: receiver.clone(),
434 last_sequence: *last_sequence,
435 })
436 .collect(),
437 }
438 }
439
440 pub fn from_snapshot(snapshot: CrossOperationRouterSnapshot) -> Result<Self, SnapshotError> {
441 let mut router = Self::new(snapshot.max_pending_per_recipient);
442 for operation in snapshot.operations {
443 if !router.register(operation.clone()) {
444 return Err(SnapshotError::DuplicateOperation(operation.address));
445 }
446 }
447 for delivery in snapshot.deliveries {
448 let mut delivery = delivery;
449 if delivery.state == DeliveryState::InFlight {
452 delivery.state = DeliveryState::Queued;
453 }
454 if router
455 .deliveries
456 .insert(delivery.message.id.clone(), delivery.clone())
457 .is_some()
458 {
459 return Err(SnapshotError::DuplicateDelivery(delivery.message.id));
460 }
461 if !router.operations.contains_key(&delivery.message.from)
462 || !router.operations.contains_key(&delivery.message.to)
463 {
464 return Err(SnapshotError::DeliveryForUnknownOperation(
465 delivery.message.id,
466 ));
467 }
468 }
469 for queue in snapshot.queues {
470 if router.queues.contains_key(&queue.receiver) {
471 return Err(SnapshotError::DuplicateQueue(queue.receiver));
472 }
473 let mut seen = std::collections::BTreeSet::new();
474 let mut last_sequence = None;
475 for id in &queue.message_ids {
476 if !seen.insert(id.clone()) {
477 return Err(SnapshotError::DuplicateQueuedDelivery(id.clone()));
478 }
479 let Some(delivery) = router.deliveries.get(id.as_str()) else {
480 return Err(SnapshotError::UnknownQueuedDelivery(id.clone()));
481 };
482 if delivery.message.to != queue.receiver {
483 return Err(SnapshotError::QueuedDeliveryForWrongReceiver(id.clone()));
484 }
485 if delivery.state != DeliveryState::Queued {
486 return Err(SnapshotError::NonQueuedDelivery(id.clone()));
487 }
488 if last_sequence.is_some_and(|last| delivery.message.sequence <= last) {
489 return Err(SnapshotError::OutOfOrderQueue(id.clone()));
490 }
491 last_sequence = Some(delivery.message.sequence);
492 }
493 router
494 .queues
495 .insert(queue.receiver, queue.message_ids.into_iter().collect());
496 }
497 let mut replay_ids: Vec<(OperationAddress, CompactString, u64)> = router
498 .deliveries
499 .values()
500 .filter(|delivery| delivery.state == DeliveryState::Queued)
501 .filter(|delivery| {
502 !router
503 .queues
504 .get(&delivery.message.to)
505 .is_some_and(|queue| queue.contains(&delivery.message.id))
506 })
507 .map(|delivery| {
508 (
509 delivery.message.to.clone(),
510 delivery.message.id.clone(),
511 delivery.message.sequence,
512 )
513 })
514 .collect();
515 replay_ids.sort_by_key(|(_, _, sequence)| *sequence);
516 for (receiver, id, _) in replay_ids.into_iter().rev() {
517 router.queues.entry(receiver).or_default().push_front(id);
518 }
519 for sequence in snapshot.last_sequences {
520 if !router.operations.contains_key(&sequence.receiver) {
521 return Err(SnapshotError::SequenceForUnknownOperation(
522 sequence.receiver,
523 ));
524 }
525 if router
526 .last_sequences
527 .insert(sequence.receiver.clone(), sequence.last_sequence)
528 .is_some()
529 {
530 return Err(SnapshotError::DuplicateSequence(sequence.receiver));
531 }
532 }
533 for delivery in router.deliveries.values() {
537 let watermark = router
538 .last_sequences
539 .entry(delivery.message.to.clone())
540 .or_default();
541 *watermark = (*watermark).max(delivery.message.sequence);
542 }
543 Ok(router)
544 }
545
546 fn pending_for(&self, receiver: &OperationAddress) -> usize {
547 self.deliveries
548 .values()
549 .filter(|delivery| delivery.message.to == *receiver && delivery.state.is_pending())
550 .count()
551 }
552
553 fn has_in_flight_for(&self, receiver: &OperationAddress) -> bool {
554 self.deliveries.values().any(|delivery| {
555 delivery.message.to == *receiver && delivery.state == DeliveryState::InFlight
556 })
557 }
558}
559
560fn rejected(failure: DeliveryFailure) -> RouteOutcome {
561 RouteOutcome::Rejected { failure }
562}
563
564fn expired(sent_at_turn: u32, ttl_turns: u32, now_turn: u32) -> bool {
565 now_turn >= sent_at_turn.saturating_add(ttl_turns)
566}
567
568fn valid_payload_locator(object: &ObjectDescriptor, payload: &PayloadLocator) -> bool {
569 !payload.reference.is_empty()
570 && !payload.digest.is_empty()
571 && object.payload_ref.as_deref() == Some(payload.reference.as_str())
572 && object.digest == payload.digest
573 && object.size == payload.size
574}
575
576fn matching_expired_lease(
577 capabilities: &[Capability],
578 action: &str,
579 object: &ObjectDescriptor,
580 now_turn: u32,
581) -> bool {
582 let resource = format!("object:{}/{}", object.owner, object.id);
583 capabilities.iter().any(|capability| {
584 capability.actions.0.contains(action)
585 && crate::types::capability::resource_matches(&capability.resource, &resource)
586 && capability
587 .lease
588 .as_ref()
589 .is_some_and(|lease| lease.is_expired(now_turn))
590 })
591}
592
593fn object_action_allowed(
594 capabilities: &[Capability],
595 action: &str,
596 object: &ObjectDescriptor,
597 now_turn: u32,
598) -> bool {
599 capabilities.iter().any(|capability| {
600 !capability
601 .lease
602 .as_ref()
603 .is_some_and(|lease| lease.is_expired(now_turn))
604 && capability.actions.0.contains(action)
605 && crate::types::capability::resource_matches(
606 &capability.resource,
607 &format!("object:{}/{}", object.owner, object.id),
608 )
609 })
610}
611
612fn delegation_is_attenuated(
613 delegated: &[Capability],
614 sender: &[Capability],
615 receiver: &[Capability],
616 now_turn: u32,
617) -> bool {
618 delegated.iter().all(|requested| {
619 requested
620 .lease
621 .as_ref()
622 .is_none_or(|lease| !lease.is_expired(now_turn))
623 && sender.iter().any(|parent| {
624 parent.delegatable
625 && !parent
626 .lease
627 .as_ref()
628 .is_some_and(|lease| lease.is_expired(now_turn))
629 && crate::types::capability::is_attenuation_of(requested, parent)
630 && lease_attenuates(requested, parent)
631 })
632 && receiver.iter().any(|parent| {
633 !parent
634 .lease
635 .as_ref()
636 .is_some_and(|lease| lease.is_expired(now_turn))
637 && crate::types::capability::is_attenuation_of(requested, parent)
638 && lease_attenuates(requested, parent)
639 })
640 })
641}
642
643fn lease_attenuates(child: &Capability, parent: &Capability) -> bool {
644 match (
645 child.lease.as_ref().and_then(|lease| lease.expires_at_turn),
646 parent
647 .lease
648 .as_ref()
649 .and_then(|lease| lease.expires_at_turn),
650 ) {
651 (_, None) => true,
652 (Some(child_expiry), Some(parent_expiry)) => child_expiry <= parent_expiry,
653 (None, Some(_)) => false,
654 }
655}
656
657#[cfg(test)]
658mod tests {
659 use super::*;
660 use crate::mm::handle::{ObjectId, ObjectKind, Residency};
661 use crate::types::capability::{
662 ActionSet, CapabilityId, CapabilityKind, ConstraintSet, Lease, Principal, ResourceSelector,
663 };
664
665 fn address(operation: &str, task: &str) -> OperationAddress {
666 OperationAddress::new(operation, task)
667 }
668
669 fn cap(actions: &[&str], delegatable: bool) -> Capability {
670 Capability {
671 id: CapabilityId("object-route".into()),
672 kind: CapabilityKind::Tool,
673 resource: ResourceSelector("object:source/7".into()),
674 actions: ActionSet(actions.iter().map(|action| (*action).into()).collect()),
675 constraints: ConstraintSet::default(),
676 lease: None,
677 delegatable,
678 issuer: Principal("kernel-test".into()),
679 }
680 }
681
682 fn external_object() -> ObjectDescriptor {
683 ObjectDescriptor::external(
684 ObjectId::from(7_u32),
685 ObjectKind::Artifact,
686 TaskId::from("source"),
687 1,
688 Residency::External {
689 payload_ref: "host://payload/7".into(),
690 digest: "sha256:abc".into(),
691 original_size: 9,
692 },
693 "preview",
694 )
695 }
696
697 fn message(id: &str, sequence: u64, ttl_turns: u32) -> CrossOperationMessage {
698 CrossOperationMessage {
699 id: id.into(),
700 from: address("operation-a", "source"),
701 to: address("operation-b", "sink"),
702 kind: "artifact".into(),
703 object: external_object(),
704 payload: PayloadLocator {
705 reference: "host://payload/7".into(),
706 digest: "sha256:abc".into(),
707 size: 9,
708 },
709 sequence,
710 sent_at_turn: 5,
711 ttl_turns,
712 delegated_capabilities: vec![cap(&["read"], false)],
713 }
714 }
715
716 fn router(max_pending: usize) -> CrossOperationRouter {
717 let mut router = CrossOperationRouter::new(max_pending);
718 assert!(router.register(OperationRegistration {
719 address: address("operation-a", "source"),
720 capabilities: vec![cap(&["share", "read"], true)],
721 cancelled: false,
722 }));
723 assert!(router.register(OperationRegistration {
724 address: address("operation-b", "sink"),
725 capabilities: vec![cap(&["read"], false)],
726 cancelled: false,
727 }));
728 router
729 }
730
731 #[test]
732 fn two_operations_route_by_address_with_payload_locator_and_receiver_attenuation() {
733 let mut router = router(8);
734 let original_receiver_caps = router
735 .operation(&address("operation-b", "sink"))
736 .unwrap()
737 .capabilities
738 .clone();
739 let msg = message("message-1", 1, 20);
740
741 assert_eq!(
742 router.route(msg.clone(), 5, PayloadAvailability::Available),
743 RouteOutcome::Accepted {
744 state: DeliveryState::Queued
745 }
746 );
747 let dispatched = router
748 .dispatch_next(&address("operation-b", "sink"), 6)
749 .expect("the host can now deliver the opaque locator");
750 assert_eq!(dispatched.payload.reference, "host://payload/7");
751 assert_eq!(
752 router.settle("message-1", DeliverySettlement::Delivered, 6),
753 Some(DeliveryState::Delivered)
754 );
755 assert_eq!(
756 router
757 .operation(&address("operation-b", "sink"))
758 .unwrap()
759 .capabilities,
760 original_receiver_caps,
761 "a sender's delegated capability must never modify receiver authority"
762 );
763
764 let mut widened = message("message-2", 2, 20);
765 widened.delegated_capabilities = vec![cap(&["write"], false)];
766 assert_eq!(
767 router.route(widened, 5, PayloadAvailability::Available),
768 RouteOutcome::Rejected {
769 failure: DeliveryFailure::CapabilityAttenuationDenied
770 }
771 );
772 }
773
774 #[test]
775 fn restart_replay_preserves_dedupe_order_and_ttl() {
776 let mut router = router(8);
777 let first = message("message-1", 1, 10);
778 let second = message("message-2", 2, 10);
779 assert!(matches!(
780 router.route(first.clone(), 5, PayloadAvailability::Available),
781 RouteOutcome::Accepted { .. }
782 ));
783 assert!(matches!(
784 router.route(second, 5, PayloadAvailability::Available),
785 RouteOutcome::Accepted { .. }
786 ));
787
788 let json = serde_json::to_string(&router.snapshot()).expect("durable snapshot serializes");
789 let snapshot: CrossOperationRouterSnapshot =
790 serde_json::from_str(&json).expect("durable snapshot restores");
791 let mut restored = CrossOperationRouter::from_snapshot(snapshot).expect("valid snapshot");
792 assert_eq!(
793 restored.route(first, 5, PayloadAvailability::Available),
794 RouteOutcome::Duplicate {
795 state: DeliveryState::Queued
796 }
797 );
798 assert_eq!(
799 restored.route(message("older", 1, 10), 5, PayloadAvailability::Available),
800 RouteOutcome::Rejected {
801 failure: DeliveryFailure::OutOfOrder
802 }
803 );
804 assert_eq!(
805 restored
806 .dispatch_next(&address("operation-b", "sink"), 6)
807 .unwrap()
808 .id,
809 CompactString::from("message-1")
810 );
811 assert_eq!(
812 restored.settle("message-1", DeliverySettlement::RetryableFailure, 6),
813 Some(DeliveryState::Queued)
814 );
815 assert_eq!(
816 restored.dispatch_next(&address("operation-b", "sink"), 15),
817 None,
818 "expired retry and later queued message are never dispatched after restart"
819 );
820 assert_eq!(
821 restored.delivery("message-1").unwrap().state,
822 DeliveryState::Expired
823 );
824 assert_eq!(
825 restored.delivery("message-2").unwrap().state,
826 DeliveryState::Expired
827 );
828 }
829
830 #[test]
831 fn restart_requeues_an_in_flight_delivery_for_at_least_once_replay() {
832 let mut router = router(8);
833 assert!(matches!(
834 router.route(
835 message("message-1", 1, 20),
836 5,
837 PayloadAvailability::Available
838 ),
839 RouteOutcome::Accepted { .. }
840 ));
841 assert!(
842 router
843 .dispatch_next(&address("operation-b", "sink"), 6)
844 .is_some()
845 );
846
847 let mut restored =
848 CrossOperationRouter::from_snapshot(router.snapshot()).expect("valid snapshot");
849 assert_eq!(
850 restored.delivery("message-1").unwrap().state,
851 DeliveryState::Queued
852 );
853 assert_eq!(
854 restored
855 .dispatch_next(&address("operation-b", "sink"), 7)
856 .unwrap()
857 .id,
858 CompactString::from("message-1")
859 );
860 }
861
862 #[test]
863 fn snapshot_rebuilds_missing_sequence_watermarks_before_accepting_new_delivery() {
864 let mut router = router(8);
865 assert!(matches!(
866 router.route(
867 message("message-2", 2, 20),
868 5,
869 PayloadAvailability::Available
870 ),
871 RouteOutcome::Accepted { .. }
872 ));
873 let mut snapshot = router.snapshot();
874 snapshot.last_sequences.clear();
875 let mut restored =
876 CrossOperationRouter::from_snapshot(snapshot).expect("watermark derives");
877
878 assert_eq!(
879 restored.route(
880 message("message-1", 1, 20),
881 5,
882 PayloadAvailability::Available
883 ),
884 RouteOutcome::Rejected {
885 failure: DeliveryFailure::OutOfOrder
886 }
887 );
888 }
889
890 #[test]
891 fn snapshot_rejects_duplicate_or_out_of_order_receiver_queue_entries() {
892 let mut router = router(8);
893 assert!(matches!(
894 router.route(
895 message("message-1", 1, 20),
896 5,
897 PayloadAvailability::Available
898 ),
899 RouteOutcome::Accepted { .. }
900 ));
901 assert!(matches!(
902 router.route(
903 message("message-2", 2, 20),
904 5,
905 PayloadAvailability::Available
906 ),
907 RouteOutcome::Accepted { .. }
908 ));
909 let mut duplicate = router.snapshot();
910 duplicate.queues[0].message_ids.push("message-1".into());
911 assert!(matches!(
912 CrossOperationRouter::from_snapshot(duplicate),
913 Err(SnapshotError::DuplicateQueuedDelivery(id)) if id == "message-1"
914 ));
915
916 let mut out_of_order = router.snapshot();
917 out_of_order.queues[0].message_ids.swap(0, 1);
918 assert!(matches!(
919 CrossOperationRouter::from_snapshot(out_of_order),
920 Err(SnapshotError::OutOfOrderQueue(id)) if id == "message-1"
921 ));
922 }
923
924 #[test]
925 fn payload_failure_backpressure_expired_lease_and_cancellation_are_explicit() {
926 let mut router = router(1);
927 let first = message("message-1", 1, 20);
928 assert_eq!(
929 router.route(first.clone(), 5, PayloadAvailability::Unavailable),
930 RouteOutcome::Rejected {
931 failure: DeliveryFailure::PayloadUnavailable
932 }
933 );
934 assert!(DeliveryFailure::PayloadUnavailable.is_retryable());
935 assert!(matches!(
936 router.route(first, 5, PayloadAvailability::Available),
937 RouteOutcome::Accepted { .. }
938 ));
939 assert_eq!(
940 router.route(
941 message("message-2", 2, 20),
942 5,
943 PayloadAvailability::Available
944 ),
945 RouteOutcome::Rejected {
946 failure: DeliveryFailure::Backpressure
947 }
948 );
949
950 let source = address("operation-a", "source");
951 router.operation(&source).unwrap();
952 router.operations.get_mut(&source).unwrap().capabilities[0].lease = Some(Lease {
953 expires_at_turn: Some(5),
954 });
955 assert_eq!(
956 router.route(
957 message("expired-lease", 2, 20),
958 5,
959 PayloadAvailability::Available
960 ),
961 RouteOutcome::Rejected {
962 failure: DeliveryFailure::ExpiredCapabilityLease
963 }
964 );
965 assert_eq!(router.cancel_operation(&address("operation-b", "sink")), 1);
966 assert_eq!(
967 router.delivery("message-1").unwrap().state,
968 DeliveryState::Cancelled
969 );
970 assert_eq!(
971 router.route(
972 message("after-cancel", 2, 20),
973 5,
974 PayloadAvailability::Available
975 ),
976 RouteOutcome::Rejected {
977 failure: DeliveryFailure::Cancelled
978 }
979 );
980 }
981
982 #[test]
983 fn unrelated_expired_capability_does_not_block_an_authorized_delivery() {
984 let mut router = router(8);
985 let source = address("operation-a", "source");
986 router
987 .operations
988 .get_mut(&source)
989 .unwrap()
990 .capabilities
991 .push(Capability {
992 id: CapabilityId("unrelated".into()),
993 kind: CapabilityKind::Tool,
994 resource: ResourceSelector("object:other/99".into()),
995 actions: ActionSet(["read".into()].into_iter().collect()),
996 constraints: ConstraintSet::default(),
997 lease: Some(Lease {
998 expires_at_turn: Some(5),
999 }),
1000 delegatable: false,
1001 issuer: Principal("kernel-test".into()),
1002 });
1003
1004 assert_eq!(
1005 router.route(
1006 message("message-1", 1, 20),
1007 5,
1008 PayloadAvailability::Available
1009 ),
1010 RouteOutcome::Accepted {
1011 state: DeliveryState::Queued
1012 }
1013 );
1014 }
1015
1016 #[test]
1017 fn an_expired_duplicate_capability_does_not_override_a_matching_live_capability() {
1018 let mut router = router(8);
1019 let source = address("operation-a", "source");
1020 router
1021 .operations
1022 .get_mut(&source)
1023 .unwrap()
1024 .capabilities
1025 .push(Capability {
1026 id: CapabilityId("expired-duplicate".into()),
1027 kind: CapabilityKind::Tool,
1028 resource: ResourceSelector("object:source/7".into()),
1029 actions: ActionSet(["share".into()].into_iter().collect()),
1030 constraints: ConstraintSet::default(),
1031 lease: Some(Lease {
1032 expires_at_turn: Some(5),
1033 }),
1034 delegatable: false,
1035 issuer: Principal("kernel-test".into()),
1036 });
1037
1038 assert!(matches!(
1039 router.route(
1040 message("message-1", 1, 20),
1041 5,
1042 PayloadAvailability::Available
1043 ),
1044 RouteOutcome::Accepted { .. }
1045 ));
1046 }
1047
1048 #[test]
1049 fn delegated_capability_cannot_outlive_its_parent_or_arrive_expired() {
1050 let mut router = router(8);
1051 let source = address("operation-a", "source");
1052 let receiver = address("operation-b", "sink");
1053 router.operations.get_mut(&source).unwrap().capabilities[0].lease = Some(Lease {
1054 expires_at_turn: Some(10),
1055 });
1056 router.operations.get_mut(&receiver).unwrap().capabilities[0].lease = Some(Lease {
1057 expires_at_turn: Some(10),
1058 });
1059
1060 let mut permanent = message("permanent", 1, 20);
1061 permanent.delegated_capabilities[0].lease = None;
1062 assert_eq!(
1063 router.route(permanent, 5, PayloadAvailability::Available),
1064 RouteOutcome::Rejected {
1065 failure: DeliveryFailure::CapabilityAttenuationDenied
1066 }
1067 );
1068
1069 let mut expired = message("expired", 1, 20);
1070 expired.delegated_capabilities[0].lease = Some(Lease {
1071 expires_at_turn: Some(5),
1072 });
1073 assert_eq!(
1074 router.route(expired, 5, PayloadAvailability::Available),
1075 RouteOutcome::Rejected {
1076 failure: DeliveryFailure::CapabilityAttenuationDenied
1077 }
1078 );
1079 }
1080}