Skip to main content

deepstrike_core/scheduler/
cross_operation.rs

1//! spc_016-05: deterministic cross-operation message routing.
2//!
3//! The kernel never opens a transport, fetches a payload, or retains a credential. Hosts submit
4//! delivery facts and carry the returned message over their own transport. This module owns only
5//! authorization, ordering, retry, and the durable state needed to replay those decisions.
6
7use 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    /// Host-owned opaque locator. It is never opened or interpreted by the kernel.
34    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    /// Strictly increasing per destination, supplied by the Host's transport plane.
56    pub sequence: u64,
57    pub sent_at_turn: u32,
58    pub ttl_turns: u32,
59    /// A requested attenuation, never an authority grant to the receiver.
60    #[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/// Deterministic kernel-side routing state. Every map is converted to an ordered vector at the
182/// checkpoint boundary, avoiding map-key JSON encoding quirks in cross-SDK snapshots.
183#[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    /// Returns the next delivery only after all earlier accepted deliveries for the receiver have
305    /// settled. Host transport can therefore retry without changing receiver-visible order.
306    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    /// Cancelling either endpoint is a durable cancellation fact for all nonterminal deliveries
391    /// touching that operation. A later Host acknowledgement cannot resurrect them.
392    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            // A checkpoint cannot include the host's acknowledgement. An in-flight delivery is
450            // therefore replayed at least once after restart, with the same identity and sequence.
451            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        // The persisted watermarks are an index, not an authority source. Recompute their lower
534        // bound from durable messages so a partial/older snapshot cannot accept a sequence the
535        // router had already seen before restart.
536        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}