Skip to main content

canwu_transport/
lib.rs

1//! Domain-neutral transport execution records.
2//!
3//! The crate owns transport semantics, not the simulation scheduler or the
4//! information lifecycle. Applications persist these values through their
5//! domain-record schemas and drive transitions through canonical ingress.
6
7#![allow(clippy::missing_errors_doc)]
8
9use canwu_core::{
10    DomainRecordRef, DomainRecordVersionRef, EntityRef, EvidenceRef, KnowledgeHolderRef,
11};
12use canwu_routing::{RoutePlan, RoutingNodeRef};
13use canwu_time::SimTime;
14use serde::{Deserialize, Serialize};
15use std::cmp::Ordering;
16
17/// Semantic identity of the transport record and transition contract.
18///
19/// `v5` adds capacity pools with deterministic booking allocation, booking
20/// requests on an execution, cancellation, arrival settlement for executions
21/// without a delivery attempt, and booking confirmation or cancellation before
22/// the booking window opens.
23pub const TRANSPORT_SEMANTIC_VERSION: &str = "canwu-transport.v5";
24
25/// Hash domain of [`CapacityBookingAllocationEvidenceV1::semantic_digest`].
26pub const CAPACITY_BOOKING_ALLOCATION_DIGEST_DOMAIN: &str =
27    "canwu.transport.capacity-booking-allocation.v1";
28
29#[must_use]
30pub fn delivery_completion_operation_key(
31    execution: TransportExecutionId,
32    revision: ItineraryRevisionId,
33    attempt_version: u64,
34) -> String {
35    format!(
36        "transport/{}/revision/{}/delivery-completion/attempt-version/{}",
37        execution.0, revision.0, attempt_version
38    )
39}
40
41#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
42#[serde(transparent)]
43pub struct TransportExecutionId(pub u64);
44
45/// Stable identity for an admitted transport-domain movement intent.
46///
47/// The simulation admits movement through a host- or integration-defined command.
48/// This record adds route-plan and custody evidence when a transport domain needs
49/// to persist a richer execution than its own projected transit state.
50#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
51#[serde(transparent)]
52pub struct MovementOrderId(pub u64);
53
54/// Who initiated a movement intent. The runtime must derive this from the
55/// admitted authority and never trust an unvalidated caller-supplied label.
56#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
57#[serde(rename_all = "snake_case")]
58pub enum MovementInitiative {
59    SelfDirected,
60    Commanded,
61    Delegated,
62    Forced,
63    Automatic,
64}
65
66/// The physical role of a subject in a movement manifest.
67#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
68#[serde(rename_all = "snake_case")]
69pub enum MovementSubjectRole {
70    MovablePrincipal,
71    Cargo,
72    Carrier,
73    Passenger,
74    Attached,
75    /// An aggregate of people moving as one subject. The subject identity is
76    /// normally an application domain-record reference, and the manifest
77    /// quantity is its positive head count.
78    PersonsGroup,
79}
80
81impl MovementSubjectRole {
82    /// Whether this role is counted by a positive integer manifest quantity:
83    /// units for [`Self::Cargo`], heads for [`Self::PersonsGroup`].
84    #[must_use]
85    pub const fn requires_quantity(self) -> bool {
86        matches!(self, Self::Cargo | Self::PersonsGroup)
87    }
88}
89
90/// One typed identity in a movement manifest.
91#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
92pub struct MovementSubject {
93    pub entity: EntityRef,
94    pub role: MovementSubjectRole,
95    /// Cargo quantities are integer units and persons-group quantities are
96    /// head counts; both must be present and positive. Other roles leave this
97    /// unset because their cardinality is one identity.
98    #[serde(default, skip_serializing_if = "Option::is_none")]
99    pub quantity: Option<u64>,
100    /// Expected carrier/custodian identity at admission, when applicable.
101    #[serde(default, skip_serializing_if = "Option::is_none")]
102    pub expected_custody: Option<EntityRef>,
103}
104
105/// Immutable, admitted intent shared by transport movement domains.
106///
107/// `MovementOrder` is a contract for planning and authority evidence; it does
108/// not directly mutate a world entity. Domain handlers still own location,
109/// custody, quantity, arrival, and knowledge effects.
110#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
111pub struct MovementOrder {
112    pub id: MovementOrderId,
113    pub subjects: Vec<MovementSubject>,
114    pub origin: RoutingNodeRef,
115    pub destination: RoutingNodeRef,
116    pub plan: RoutePlan,
117    pub initiative: MovementInitiative,
118    pub ordered_at: SimTime,
119    pub expected_position_revision: u64,
120}
121
122#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
123pub enum MovementOrderError {
124    Invalid(String),
125}
126
127impl MovementOrder {
128    /// Validates the structural invariants that are domain-neutral.
129    ///
130    /// Existence, authority, capability, and position/custody matching remain
131    /// runtime or plugin responsibilities because this crate does not own the
132    /// world or domain records.
133    pub fn validate(&self) -> Result<(), MovementOrderError> {
134        if self.id.0 == 0
135            || self.expected_position_revision == 0
136            || self.origin.as_str().trim().is_empty()
137            || self.destination.as_str().trim().is_empty()
138        {
139            return Err(MovementOrderError::Invalid(
140                "movement order identity, endpoints, and expected position revision must be valid"
141                    .to_owned(),
142            ));
143        }
144        if self.subjects.is_empty()
145            || self
146                .subjects
147                .windows(2)
148                .any(|pair| pair[0].entity >= pair[1].entity)
149        {
150            return Err(MovementOrderError::Invalid(
151                "movement subjects must be non-empty, sorted, and unique by entity".to_owned(),
152            ));
153        }
154        for subject in &self.subjects {
155            if subject.quantity.is_some_and(|quantity| quantity == 0)
156                || subject.role.requires_quantity() != subject.quantity.is_some()
157            {
158                return Err(MovementOrderError::Invalid(
159                    "cargo and persons-group subjects require a positive quantity; other subjects cannot carry one"
160                        .to_owned(),
161                ));
162            }
163        }
164        if self.plan.origin != self.origin
165            || self.plan.destination != self.destination
166            || self.plan.departure_at < self.ordered_at
167            || self.plan.estimated_arrival_at < self.plan.departure_at
168            || self.plan.digest.trim().is_empty()
169        {
170            return Err(MovementOrderError::Invalid(
171                "movement plan does not match the order endpoints or time range".to_owned(),
172            ));
173        }
174        Ok(())
175    }
176}
177
178#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
179#[serde(transparent)]
180pub struct ItineraryRevisionId(pub u64);
181
182#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
183#[serde(transparent)]
184pub struct LegExecutionId(pub u64);
185
186#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
187#[serde(transparent)]
188pub struct HandoffId(pub u64);
189
190#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
191#[serde(transparent)]
192pub struct CapacityBookingId(pub u64);
193
194#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
195#[serde(rename_all = "snake_case")]
196pub enum TransportExecutionState {
197    Prepared,
198    Planning,
199    Booking,
200    Ready,
201    Executing,
202    ReplanPending,
203    ArrivalPending,
204    Settled,
205    Failed,
206    Cancelled,
207}
208
209impl TransportExecutionState {
210    #[must_use]
211    pub const fn is_terminal(self) -> bool {
212        matches!(self, Self::Settled | Self::Failed | Self::Cancelled)
213    }
214}
215
216#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
217#[serde(rename_all = "snake_case")]
218pub enum ItineraryRevisionReason {
219    Initial,
220    Disaster {
221        explanation: String,
222    },
223    CapacityUnavailable {
224        explanation: String,
225    },
226    KnowledgeUpdate {
227        explanation: String,
228    },
229    Recovery {
230        explanation: String,
231    },
232    /// A revision caused by an application-owned condition record, cited at
233    /// the exact positive record version the planner acted on. `kind` is an
234    /// application label for the condition; transport does not interpret it.
235    ExternalCondition {
236        record: DomainRecordRef,
237        version: u64,
238        kind: String,
239    },
240}
241
242impl ItineraryRevisionReason {
243    fn validate(&self) -> Result<(), TransportError> {
244        if let Self::ExternalCondition {
245            record,
246            version,
247            kind,
248        } = self
249            && (*version == 0
250                || kind.trim().is_empty()
251                || !domain_record_ref_is_well_formed(record))
252        {
253            return Err(TransportError::InvalidRevision(
254                "external-condition reason requires a well-formed record, a positive version, and a kind"
255                    .to_owned(),
256            ));
257        }
258        Ok(())
259    }
260}
261
262fn domain_record_ref_is_well_formed(record: &DomainRecordRef) -> bool {
263    !record.kind.namespace.trim().is_empty()
264        && !record.kind.name.trim().is_empty()
265        && !record.id.trim().is_empty()
266}
267
268#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
269pub struct ItineraryRevision {
270    pub id: ItineraryRevisionId,
271    pub predecessor: Option<ItineraryRevisionId>,
272    pub plan: RoutePlan,
273    pub planned_at: SimTime,
274    pub valid_from: SimTime,
275    pub reason: ItineraryRevisionReason,
276    pub superseded_at: Option<SimTime>,
277    pub evidence: Vec<EvidenceRef>,
278}
279
280#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
281#[serde(rename_all = "snake_case")]
282pub enum LegExecutionStatus {
283    Planned,
284    Booked,
285    Loaded,
286    Departed,
287    Arrived,
288    Waiting,
289    Failed,
290    Cancelled,
291}
292
293#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
294pub struct LegExecution {
295    pub id: LegExecutionId,
296    pub itinerary_revision: ItineraryRevisionId,
297    pub leg_index: usize,
298    pub status: LegExecutionStatus,
299    pub actual_departure_at: Option<SimTime>,
300    pub actual_arrival_at: Option<SimTime>,
301    pub failed_at: Option<SimTime>,
302    pub failure_reason: Option<String>,
303    pub evidence: Vec<EvidenceRef>,
304}
305
306/// How custody changed hands between two legs.
307#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
308#[serde(rename_all = "snake_case")]
309pub enum HandoffKind {
310    /// A custody transfer the itinerary planned for.
311    #[default]
312    Planned,
313    /// Custody taken by an entity outside the itinerary. Transport records the
314    /// seizing identity; the incident, hostility, and authority behind it stay
315    /// application systems. Unless it is terminal (below), a seizure follows
316    /// the same leg rules as a planned handoff: it is recorded after the
317    /// source leg arrived or failed, into a startable leg (typically the first
318    /// leg of a replacement revision), and `to_custodian` remains the
319    /// application's custody label.
320    ///
321    /// A seizure may instead end the itinerary's custody. Such a *terminal*
322    /// seizure names the failed current leg of a non-terminal execution's
323    /// active itinerary as both `from_leg` and `to_leg`. It is recorded no
324    /// earlier than that leg failed and only if no other handoff left it, and
325    /// at most once per execution. Afterwards the execution records no
326    /// handoff and cannot fail or reroute a leg, arrive, settle successfully,
327    /// or book capacity (see [`TransportExecution::custody_left_itinerary`]); its
328    /// owner closes it as failed or cancelled.
329    Seizure { by: EntityRef },
330}
331
332impl HandoffKind {
333    #[must_use]
334    pub const fn is_planned(&self) -> bool {
335        matches!(self, Self::Planned)
336    }
337}
338
339#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
340pub struct Handoff {
341    pub id: HandoffId,
342    pub from_leg: LegExecutionId,
343    pub to_leg: LegExecutionId,
344    pub from_custodian: String,
345    pub to_custodian: String,
346    pub at: SimTime,
347    pub location: String,
348    pub evidence: Vec<EvidenceRef>,
349    /// Omitted from JSON when [`HandoffKind::Planned`], so handoffs recorded
350    /// before seizure handoffs existed keep their exact shape.
351    #[serde(default, skip_serializing_if = "HandoffKind::is_planned")]
352    pub kind: HandoffKind,
353}
354
355#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
356pub struct DeliveryCompletionRequest {
357    pub operation_key: String,
358    pub execution: TransportExecutionId,
359    pub itinerary_revision: ItineraryRevisionId,
360    pub delivery_attempt: DomainRecordVersionRef,
361    pub completed_at: SimTime,
362    pub evidence: Vec<EvidenceRef>,
363}
364
365#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
366#[serde(rename_all = "snake_case")]
367pub enum CapacityBookingStatus {
368    Requested,
369    Confirmed,
370    Consumed,
371    Released,
372    Expired,
373    Cancelled,
374    Failed,
375}
376
377#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
378pub struct CapacityBooking {
379    pub id: CapacityBookingId,
380    pub execution: TransportExecutionId,
381    pub resource: String,
382    pub valid_from: SimTime,
383    pub valid_until: SimTime,
384    pub quantity: u64,
385    pub priority: i32,
386    pub status: CapacityBookingStatus,
387    pub allocation_evidence: Vec<EvidenceRef>,
388}
389
390impl CapacityBooking {
391    pub fn new(
392        id: CapacityBookingId,
393        execution: TransportExecutionId,
394        resource: String,
395        valid_from: SimTime,
396        valid_until: SimTime,
397        quantity: u64,
398        priority: i32,
399    ) -> Result<Self, TransportError> {
400        if valid_until < valid_from || quantity == 0 {
401            return Err(TransportError::InvalidBooking(
402                "capacity booking requires a positive quantity and non-inverted window".to_owned(),
403            ));
404        }
405        Ok(Self {
406            id,
407            execution,
408            resource,
409            valid_from,
410            valid_until,
411            quantity,
412            priority,
413            status: CapacityBookingStatus::Requested,
414            allocation_evidence: Vec::new(),
415        })
416    }
417
418    /// Moves the booking through its windowed state machine.
419    ///
420    /// A request may be confirmed, failed, or cancelled, and a confirmed
421    /// booking released or cancelled, before its window opens; capacity can
422    /// only be consumed inside the window, confirmed only until it ends, and
423    /// expired only after it ends.
424    pub fn transition(
425        &mut self,
426        status: CapacityBookingStatus,
427        at: SimTime,
428    ) -> Result<(), TransportError> {
429        if status == CapacityBookingStatus::Consumed
430            && (at < self.valid_from || at > self.valid_until)
431        {
432            return Err(TransportError::InvalidBooking(
433                "booking capacity can only be consumed inside its validity window".to_owned(),
434            ));
435        }
436        let allowed = matches!(
437            (self.status, status),
438            (
439                CapacityBookingStatus::Requested,
440                CapacityBookingStatus::Confirmed
441                    | CapacityBookingStatus::Failed
442                    | CapacityBookingStatus::Cancelled
443            ) | (
444                CapacityBookingStatus::Confirmed,
445                CapacityBookingStatus::Consumed
446                    | CapacityBookingStatus::Released
447                    | CapacityBookingStatus::Cancelled
448                    | CapacityBookingStatus::Expired
449            ) | (
450                CapacityBookingStatus::Consumed,
451                CapacityBookingStatus::Released
452            )
453        );
454        if !allowed {
455            return Err(TransportError::InvalidBooking(
456                "capacity booking transition is not allowed".to_owned(),
457            ));
458        }
459        if at > self.valid_until && status == CapacityBookingStatus::Confirmed {
460            return Err(TransportError::InvalidBooking(
461                "capacity booking cannot be confirmed after its validity window".to_owned(),
462            ));
463        }
464        if status == CapacityBookingStatus::Expired && at <= self.valid_until {
465            return Err(TransportError::InvalidBooking(
466                "capacity booking cannot expire before its validity window ends".to_owned(),
467            ));
468        }
469        self.status = status;
470        Ok(())
471    }
472}
473
474/// A windowed pool of interchangeable transport capacity, such as ferry
475/// crossings, carriage places, or relay mounts for one period.
476///
477/// `quantity` is the capacity offered for the whole window. `booked` is held
478/// by confirmed bookings and `consumed` by consumed ones; released, expired,
479/// cancelled, and failed bookings hold nothing. `revision` advances whenever
480/// the offer or either counter changes, so allocation evidence can cite the
481/// exact pool state it read.
482#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
483pub struct TransportCapacityPoolV1 {
484    pub id: String,
485    pub resource: String,
486    pub custodian: KnowledgeHolderRef,
487    pub window_from: SimTime,
488    pub window_until: SimTime,
489    pub quantity: u64,
490    pub booked: u64,
491    pub consumed: u64,
492    pub revision: u64,
493}
494
495impl TransportCapacityPoolV1 {
496    /// Creates an empty pool at revision one.
497    pub fn new(
498        id: String,
499        resource: String,
500        custodian: KnowledgeHolderRef,
501        window_from: SimTime,
502        window_until: SimTime,
503        quantity: u64,
504    ) -> Result<Self, TransportError> {
505        let pool = Self {
506            id,
507            resource,
508            custodian,
509            window_from,
510            window_until,
511            quantity,
512            booked: 0,
513            consumed: 0,
514            revision: 1,
515        };
516        pool.validate()?;
517        Ok(pool)
518    }
519
520    pub fn validate(&self) -> Result<(), TransportError> {
521        if !is_canonical_label(&self.id)
522            || !is_canonical_label(&self.resource)
523            || self.window_until < self.window_from
524            || self.revision == 0
525        {
526            return Err(TransportError::InvalidCapacityPool(
527                "capacity pool requires a canonical identity and resource, a non-inverted window, and a positive revision"
528                    .to_owned(),
529            ));
530        }
531        if self
532            .booked
533            .checked_add(self.consumed)
534            .is_none_or(|held| held > self.quantity)
535        {
536            return Err(TransportError::InvalidCapacityPool(
537                "capacity pool holds more than it offers".to_owned(),
538            ));
539        }
540        Ok(())
541    }
542
543    /// Capacity neither booked nor consumed.
544    #[must_use]
545    pub const fn available(&self) -> u64 {
546        self.quantity
547            .saturating_sub(self.booked)
548            .saturating_sub(self.consumed)
549    }
550
551    /// Replaces the offered window and quantity. The new quantity must still
552    /// cover everything booked or consumed.
553    pub fn revise(
554        &mut self,
555        window_from: SimTime,
556        window_until: SimTime,
557        quantity: u64,
558    ) -> Result<(), TransportError> {
559        let mut revised = self.clone();
560        revised.window_from = window_from;
561        revised.window_until = window_until;
562        revised.quantity = quantity;
563        revised.revision = self
564            .revision
565            .checked_add(1)
566            .ok_or(TransportError::Overflow)?;
567        revised.validate()?;
568        *self = revised;
569        Ok(())
570    }
571
572    /// Commits one allocation pass computed by [`allocate_capacity_bookings`]
573    /// against this exact pool revision, then advances the revision.
574    pub fn apply_allocations(
575        &mut self,
576        allocations: &[BookingAllocationV1],
577    ) -> Result<(), TransportError> {
578        if allocations.is_empty() {
579            return Ok(());
580        }
581        let mut bookings = allocations
582            .iter()
583            .map(|allocation| allocation.booking)
584            .collect::<Vec<_>>();
585        bookings.sort_unstable();
586        if bookings.windows(2).any(|pair| pair[0] == pair[1]) {
587            return Err(TransportError::InvalidCapacityPool(
588                "an allocation pass names each booking once".to_owned(),
589            ));
590        }
591        let mut confirmed = 0_u64;
592        for allocation in allocations {
593            let evidence = &allocation.evidence;
594            if evidence.pool != self.id
595                || evidence.pool_revision != self.revision
596                || evidence.booking != allocation.booking
597                || evidence.quantity != allocation.quantity
598                || evidence.status != allocation.status
599                || !evidence.digest_matches()
600            {
601                return Err(TransportError::InvalidCapacityPool(
602                    "allocation evidence does not bind this pool revision".to_owned(),
603                ));
604            }
605            match allocation.status {
606                CapacityBookingStatus::Confirmed => {
607                    confirmed = confirmed
608                        .checked_add(allocation.quantity)
609                        .ok_or(TransportError::Overflow)?;
610                }
611                CapacityBookingStatus::Failed if allocation.quantity == 0 => {}
612                _ => {
613                    return Err(TransportError::InvalidCapacityPool(
614                        "an allocation either confirms its quantity or fails with none".to_owned(),
615                    ));
616                }
617            }
618        }
619        let mut next = self.clone();
620        next.booked = self
621            .booked
622            .checked_add(confirmed)
623            .ok_or(TransportError::Overflow)?;
624        next.revision = self
625            .revision
626            .checked_add(1)
627            .ok_or(TransportError::Overflow)?;
628        next.validate()?;
629        *self = next;
630        Ok(())
631    }
632
633    /// Applies the capacity effect of one booking transition made after
634    /// allocation: consumption moves quantity from `booked` to `consumed`,
635    /// while releasing, cancelling, or expiring a confirmed booking and
636    /// releasing a consumed one return it. Cancelling or failing a request
637    /// holds nothing and leaves the pool unchanged. Confirmation happens only
638    /// through [`Self::apply_allocations`].
639    pub fn apply_booking_transition(
640        &mut self,
641        quantity: u64,
642        from: CapacityBookingStatus,
643        to: CapacityBookingStatus,
644    ) -> Result<(), TransportError> {
645        use CapacityBookingStatus::{
646            Cancelled, Confirmed, Consumed, Expired, Failed, Released, Requested,
647        };
648        let underflow =
649            || TransportError::InvalidCapacityPool("capacity pool counter underflow".to_owned());
650        let mut next = self.clone();
651        match (from, to) {
652            (Requested, Cancelled | Failed) => return Ok(()),
653            (Confirmed, Consumed) => {
654                next.booked = next.booked.checked_sub(quantity).ok_or_else(underflow)?;
655                next.consumed = next
656                    .consumed
657                    .checked_add(quantity)
658                    .ok_or(TransportError::Overflow)?;
659            }
660            (Confirmed, Released | Cancelled | Expired) => {
661                next.booked = next.booked.checked_sub(quantity).ok_or_else(underflow)?;
662            }
663            (Consumed, Released) => {
664                next.consumed = next.consumed.checked_sub(quantity).ok_or_else(underflow)?;
665            }
666            _ => {
667                return Err(TransportError::InvalidCapacityPool(
668                    "booking transition has no capacity-pool effect".to_owned(),
669                ));
670            }
671        }
672        next.revision = self
673            .revision
674            .checked_add(1)
675            .ok_or(TransportError::Overflow)?;
676        next.validate()?;
677        *self = next;
678        Ok(())
679    }
680}
681
682/// One requested booking offered to a pool allocation pass, with the caller's
683/// deterministic tie-break key and admission sequence.
684#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
685pub struct CapacityBookingRequestV1 {
686    pub booking: CapacityBooking,
687    pub tie_break: String,
688    pub admitted_sequence: u64,
689}
690
691/// Why an allocation pass failed a booking request.
692#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
693#[serde(rename_all = "snake_case")]
694pub enum CapacityAllocationFailureV1 {
695    /// The request does not fit the capacity left after earlier grants.
696    InsufficientCapacity,
697    /// The booking window is not inside the pool window.
698    OutsidePoolWindow,
699    /// The booking window ended before the allocation time.
700    WindowElapsed,
701}
702
703/// Replayable evidence of one allocation decision.
704#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
705pub struct CapacityBookingAllocationEvidenceV1 {
706    pub pool: String,
707    pub pool_revision: u64,
708    pub booking: CapacityBookingId,
709    pub execution: TransportExecutionId,
710    pub requested: u64,
711    /// Granted quantity: the full request when confirmed, zero when failed.
712    pub quantity: u64,
713    pub status: CapacityBookingStatus,
714    #[serde(default, skip_serializing_if = "Option::is_none")]
715    pub failure: Option<CapacityAllocationFailureV1>,
716    /// Pool capacity left after this decision, in allocation order.
717    pub remaining_after: u64,
718    pub allocated_at: SimTime,
719    pub operation_key: String,
720    /// Domain-separated BLAKE3 digest over every other field; see
721    /// [`CAPACITY_BOOKING_ALLOCATION_DIGEST_DOMAIN`].
722    pub semantic_digest: String,
723}
724
725impl CapacityBookingAllocationEvidenceV1 {
726    /// Recomputes the digest over every field except `semantic_digest`.
727    #[must_use]
728    pub fn expected_digest(&self) -> String {
729        let material = (
730            &self.pool,
731            self.pool_revision,
732            self.booking,
733            self.execution,
734            self.requested,
735            self.quantity,
736            self.status,
737            self.failure,
738            self.remaining_after,
739            self.allocated_at,
740            &self.operation_key,
741        );
742        domain_digest(CAPACITY_BOOKING_ALLOCATION_DIGEST_DOMAIN, &material)
743    }
744
745    #[must_use]
746    pub fn digest_matches(&self) -> bool {
747        self.semantic_digest == self.expected_digest()
748    }
749}
750
751/// One booking's allocation result, in allocation order.
752#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
753pub struct BookingAllocationV1 {
754    pub booking: CapacityBookingId,
755    /// Granted quantity: the full request when confirmed, zero when failed.
756    pub quantity: u64,
757    /// `Confirmed` or `Failed`.
758    pub status: CapacityBookingStatus,
759    pub evidence: CapacityBookingAllocationEvidenceV1,
760}
761
762/// Stable operation key of one booking's allocation at one pool revision.
763#[must_use]
764pub fn capacity_booking_allocation_operation_key(
765    pool: &str,
766    pool_revision: u64,
767    booking: CapacityBookingId,
768) -> String {
769    format!(
770        "transport/pool/{pool}/revision/{pool_revision}/booking/{}/allocation",
771        booking.0
772    )
773}
774
775/// Allocates requested bookings against one pool revision, all or nothing
776/// per booking.
777///
778/// Requests are visited by descending priority, then ascending `valid_from`,
779/// tie-break key, admission sequence, and booking identity, so the result
780/// does not depend on input order. A request is confirmed when its window
781/// lies inside the pool window, has not ended at `at`, and its full quantity
782/// fits the capacity left by earlier grants; otherwise it fails with a
783/// recorded reason, and a later, smaller request may still fit. The function
784/// reads only its arguments; commit the result with
785/// [`TransportCapacityPoolV1::apply_allocations`] and
786/// [`CapacityBooking::transition`].
787pub fn allocate_capacity_bookings(
788    pool: &TransportCapacityPoolV1,
789    requests: &[CapacityBookingRequestV1],
790    at: SimTime,
791) -> Result<Vec<BookingAllocationV1>, TransportError> {
792    pool.validate()?;
793    let mut ids = Vec::with_capacity(requests.len());
794    for request in requests {
795        let booking = &request.booking;
796        if booking.status != CapacityBookingStatus::Requested
797            || booking.resource != pool.resource
798            || booking.quantity == 0
799            || booking.valid_until < booking.valid_from
800            || !is_canonical_label(&request.tie_break)
801        {
802            return Err(TransportError::InvalidCapacityPool(
803                "allocation accepts only well-formed requested bookings for the pool resource"
804                    .to_owned(),
805            ));
806        }
807        ids.push(booking.id);
808    }
809    ids.sort_unstable();
810    if ids.windows(2).any(|pair| pair[0] == pair[1]) {
811        return Err(TransportError::InvalidCapacityPool(
812            "allocation requests must have unique booking identities".to_owned(),
813        ));
814    }
815    let mut ordered = requests.iter().collect::<Vec<_>>();
816    ordered.sort_by(|left, right| allocation_order(left, right));
817    let mut remaining = pool.available();
818    let mut allocations = Vec::with_capacity(ordered.len());
819    for request in ordered {
820        let booking = &request.booking;
821        let failure = if at > booking.valid_until {
822            Some(CapacityAllocationFailureV1::WindowElapsed)
823        } else if booking.valid_from < pool.window_from || booking.valid_until > pool.window_until {
824            Some(CapacityAllocationFailureV1::OutsidePoolWindow)
825        } else if booking.quantity > remaining {
826            Some(CapacityAllocationFailureV1::InsufficientCapacity)
827        } else {
828            None
829        };
830        let (status, quantity) = if failure.is_none() {
831            remaining -= booking.quantity;
832            (CapacityBookingStatus::Confirmed, booking.quantity)
833        } else {
834            (CapacityBookingStatus::Failed, 0)
835        };
836        let mut evidence = CapacityBookingAllocationEvidenceV1 {
837            pool: pool.id.clone(),
838            pool_revision: pool.revision,
839            booking: booking.id,
840            execution: booking.execution,
841            requested: booking.quantity,
842            quantity,
843            status,
844            failure,
845            remaining_after: remaining,
846            allocated_at: at,
847            operation_key: capacity_booking_allocation_operation_key(
848                &pool.id,
849                pool.revision,
850                booking.id,
851            ),
852            semantic_digest: String::new(),
853        };
854        evidence.semantic_digest = evidence.expected_digest();
855        allocations.push(BookingAllocationV1 {
856            booking: booking.id,
857            quantity,
858            status,
859            evidence,
860        });
861    }
862    Ok(allocations)
863}
864
865fn allocation_order(left: &CapacityBookingRequestV1, right: &CapacityBookingRequestV1) -> Ordering {
866    right
867        .booking
868        .priority
869        .cmp(&left.booking.priority)
870        .then_with(|| left.booking.valid_from.cmp(&right.booking.valid_from))
871        .then_with(|| left.tie_break.cmp(&right.tie_break))
872        .then_with(|| left.admitted_sequence.cmp(&right.admitted_sequence))
873        .then_with(|| left.booking.id.cmp(&right.booking.id))
874}
875
876fn is_canonical_label(value: &str) -> bool {
877    !value.is_empty() && value.trim() == value
878}
879
880fn domain_digest<T: Serialize>(domain: &str, value: &T) -> String {
881    let encoded = serde_json::to_vec(value).expect("transport digest material must serialize");
882    let mut hasher = blake3::Hasher::new();
883    hasher.update(domain.as_bytes());
884    hasher.update(&[0]);
885    hasher.update(&encoded);
886    hasher.finalize().to_hex().to_string()
887}
888
889#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
890#[serde(rename_all = "snake_case")]
891pub enum SagaState {
892    TransportIntent,
893    WaitingForInformation,
894    Executing,
895    ArrivalPending,
896    Settled,
897    CompensationPending,
898    Failed,
899}
900
901/// Result of reconciling the information-system delivery attempt.
902#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
903pub enum ReconciliationOutcome {
904    Success,
905    Failure { error: String },
906}
907
908#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
909pub struct DeliverySaga {
910    pub operation_key: String,
911    pub delivery_attempt: DomainRecordVersionRef,
912    pub state: SagaState,
913    pub step: u32,
914    pub expected_attempt_version: u64,
915    pub last_error: Option<String>,
916    pub evidence: Vec<EvidenceRef>,
917}
918
919#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
920pub struct TransportExecution {
921    pub id: TransportExecutionId,
922    pub delivery_attempt: Option<DomainRecordVersionRef>,
923    pub state: TransportExecutionState,
924    pub active_itinerary_revision: Option<ItineraryRevisionId>,
925    pub current_leg_index: usize,
926    pub estimated_arrival_at: Option<SimTime>,
927    pub current_endpoint: Option<String>,
928    pub revisions: Vec<ItineraryRevision>,
929    pub legs: Vec<LegExecution>,
930    pub handoffs: Vec<Handoff>,
931    pub bookings: Vec<CapacityBooking>,
932    pub saga: Option<DeliverySaga>,
933}
934
935impl TransportExecution {
936    #[must_use]
937    pub fn new(id: TransportExecutionId, delivery_attempt: Option<DomainRecordVersionRef>) -> Self {
938        Self {
939            id,
940            delivery_attempt,
941            state: TransportExecutionState::Prepared,
942            active_itinerary_revision: None,
943            current_leg_index: 0,
944            estimated_arrival_at: None,
945            current_endpoint: None,
946            revisions: Vec::new(),
947            legs: Vec::new(),
948            handoffs: Vec::new(),
949            bookings: Vec::new(),
950            saga: None,
951        }
952    }
953
954    pub fn install_initial_itinerary(
955        &mut self,
956        revision: ItineraryRevision,
957    ) -> Result<(), TransportError> {
958        if !self.revisions.is_empty() || revision.predecessor.is_some() {
959            return Err(TransportError::InvalidRevision(
960                "initial itinerary must be the first revision".to_owned(),
961            ));
962        }
963        revision.reason.validate()?;
964        self.estimated_arrival_at = Some(revision.plan.estimated_arrival_at);
965        self.active_itinerary_revision = Some(revision.id);
966        self.current_endpoint = Some(revision.plan.origin.as_str().to_owned());
967        self.legs = revision
968            .plan
969            .legs
970            .iter()
971            .enumerate()
972            .map(|(index, _)| {
973                let id = u64::try_from(index)
974                    .ok()
975                    .and_then(|index| index.checked_add(1))
976                    .ok_or(TransportError::Overflow)?;
977                Ok(LegExecution {
978                    id: LegExecutionId(id),
979                    itinerary_revision: revision.id,
980                    leg_index: index,
981                    status: LegExecutionStatus::Planned,
982                    actual_departure_at: None,
983                    actual_arrival_at: None,
984                    failed_at: None,
985                    failure_reason: None,
986                    evidence: Vec::new(),
987                })
988            })
989            .collect::<Result<Vec<_>, TransportError>>()?;
990        self.revisions.push(revision);
991        self.state = TransportExecutionState::Planning;
992        Ok(())
993    }
994
995    pub fn reroute(
996        &mut self,
997        revision: ItineraryRevision,
998        at: SimTime,
999    ) -> Result<(), TransportError> {
1000        self.ensure_custody_in_itinerary()?;
1001        let active = self
1002            .active_itinerary_revision
1003            .ok_or(TransportError::MissingItinerary)?;
1004        if revision.predecessor != Some(active) || revision.valid_from < at {
1005            return Err(TransportError::InvalidRevision("reroute must reference the active revision and start no earlier than the reroute time".to_owned()));
1006        }
1007        revision.reason.validate()?;
1008        if self.revisions.iter().any(|item| item.id == revision.id) {
1009            return Err(TransportError::InvalidRevision(
1010                "itinerary revision identity must be unique within an execution".to_owned(),
1011            ));
1012        }
1013        if let Some(previous) = self.revisions.iter_mut().find(|item| item.id == active) {
1014            previous.superseded_at = Some(at);
1015        }
1016        self.estimated_arrival_at = Some(revision.plan.estimated_arrival_at);
1017        self.active_itinerary_revision = Some(revision.id);
1018        self.current_leg_index = 0;
1019        let next_leg_id = self
1020            .legs
1021            .iter()
1022            .map(|leg| leg.id.0)
1023            .max()
1024            .unwrap_or_default()
1025            .checked_add(1)
1026            .ok_or(TransportError::Overflow)?;
1027        let mut legs = Vec::with_capacity(revision.plan.legs.len());
1028        for (index, _) in revision.plan.legs.iter().enumerate() {
1029            let id = u64::try_from(index)
1030                .ok()
1031                .and_then(|index| next_leg_id.checked_add(index))
1032                .ok_or(TransportError::Overflow)?;
1033            legs.push(LegExecution {
1034                id: LegExecutionId(id),
1035                itinerary_revision: revision.id,
1036                leg_index: index,
1037                status: LegExecutionStatus::Planned,
1038                actual_departure_at: None,
1039                actual_arrival_at: None,
1040                failed_at: None,
1041                failure_reason: None,
1042                evidence: Vec::new(),
1043            });
1044        }
1045        self.legs.extend(legs);
1046        self.revisions.push(revision);
1047        if let Some(saga) = self.saga.as_mut() {
1048            saga.operation_key = delivery_completion_operation_key(
1049                self.id,
1050                self.active_itinerary_revision
1051                    .ok_or(TransportError::MissingItinerary)?,
1052                saga.expected_attempt_version,
1053            );
1054            saga.evidence.extend(
1055                self.revisions
1056                    .last()
1057                    .map(|current| current.evidence.clone())
1058                    .unwrap_or_default(),
1059            );
1060        }
1061        self.state = TransportExecutionState::Planning;
1062        Ok(())
1063    }
1064
1065    pub fn begin_saga(
1066        &mut self,
1067        delivery_attempt: DomainRecordVersionRef,
1068        operation_key: String,
1069    ) -> Result<(), TransportError> {
1070        if self.saga.is_some() {
1071            return Err(TransportError::SagaAlreadyExists);
1072        }
1073        self.saga = Some(DeliverySaga {
1074            expected_attempt_version: delivery_attempt.version,
1075            operation_key,
1076            delivery_attempt,
1077            state: SagaState::TransportIntent,
1078            step: 0,
1079            last_error: None,
1080            evidence: Vec::new(),
1081        });
1082        self.state = TransportExecutionState::Executing;
1083        Ok(())
1084    }
1085
1086    pub fn start_current_leg(&mut self, at: SimTime) -> Result<(), TransportError> {
1087        if self.state != TransportExecutionState::Ready
1088            && self.state != TransportExecutionState::Executing
1089            && self.state != TransportExecutionState::Planning
1090        {
1091            return Err(TransportError::InvalidState(
1092                "transport execution cannot start a leg in its current state".to_owned(),
1093            ));
1094        }
1095        let active = self
1096            .active_itinerary_revision
1097            .ok_or(TransportError::MissingItinerary)?;
1098        let leg = self
1099            .legs
1100            .iter_mut()
1101            .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
1102            .ok_or(TransportError::MissingLeg)?;
1103        if !matches!(
1104            leg.status,
1105            LegExecutionStatus::Planned | LegExecutionStatus::Booked | LegExecutionStatus::Waiting
1106        ) {
1107            return Err(TransportError::InvalidState(
1108                "current leg is not startable".to_owned(),
1109            ));
1110        }
1111        leg.status = LegExecutionStatus::Departed;
1112        leg.actual_departure_at = Some(at);
1113        self.state = TransportExecutionState::Executing;
1114        Ok(())
1115    }
1116
1117    pub fn complete_current_leg(
1118        &mut self,
1119        at: SimTime,
1120        endpoint: String,
1121    ) -> Result<bool, TransportError> {
1122        let active = self
1123            .active_itinerary_revision
1124            .ok_or(TransportError::MissingItinerary)?;
1125        let active_leg_count = self
1126            .legs
1127            .iter()
1128            .filter(|leg| leg.itinerary_revision == active)
1129            .count();
1130        let next_leg_index = self
1131            .current_leg_index
1132            .checked_add(1)
1133            .ok_or(TransportError::Overflow)?;
1134        let final_leg = next_leg_index >= active_leg_count;
1135        if final_leg && self.saga.is_none() {
1136            return Err(TransportError::MissingSaga);
1137        }
1138        let leg = self
1139            .legs
1140            .iter_mut()
1141            .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
1142            .ok_or(TransportError::MissingLeg)?;
1143        if leg.status != LegExecutionStatus::Departed {
1144            return Err(TransportError::InvalidState(
1145                "current leg must be departed before arrival".to_owned(),
1146            ));
1147        }
1148        if leg
1149            .actual_departure_at
1150            .is_some_and(|departure| at < departure)
1151        {
1152            return Err(TransportError::InvalidState(
1153                "arrival precedes departure".to_owned(),
1154            ));
1155        }
1156        leg.status = LegExecutionStatus::Arrived;
1157        leg.actual_arrival_at = Some(at);
1158        self.current_endpoint = Some(endpoint);
1159        self.current_leg_index = next_leg_index;
1160        if self.current_leg_index >= active_leg_count {
1161            self.state = TransportExecutionState::ArrivalPending;
1162            self.mark_arrival_pending()?;
1163            Ok(true)
1164        } else {
1165            self.state = TransportExecutionState::Ready;
1166            Ok(false)
1167        }
1168    }
1169
1170    pub fn fail_current_leg(&mut self, reason: String, at: SimTime) -> Result<(), TransportError> {
1171        self.ensure_custody_in_itinerary()?;
1172        let active = self
1173            .active_itinerary_revision
1174            .ok_or(TransportError::MissingItinerary)?;
1175        let leg = self
1176            .legs
1177            .iter_mut()
1178            .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
1179            .ok_or(TransportError::MissingLeg)?;
1180        leg.status = LegExecutionStatus::Failed;
1181        leg.failed_at = Some(at);
1182        leg.failure_reason = Some(reason);
1183        self.state = TransportExecutionState::ReplanPending;
1184        Ok(())
1185    }
1186
1187    pub fn mark_arrival_pending(&mut self) -> Result<(), TransportError> {
1188        self.ensure_custody_in_itinerary()?;
1189        let saga = self.saga.as_mut().ok_or(TransportError::MissingSaga)?;
1190        saga.step = saga.step.checked_add(1).ok_or(TransportError::Overflow)?;
1191        saga.state = SagaState::ArrivalPending;
1192        self.state = TransportExecutionState::ArrivalPending;
1193        Ok(())
1194    }
1195
1196    pub fn record_handoff(&mut self, handoff: Handoff) -> Result<(), TransportError> {
1197        let terminal_seizure = !handoff.kind.is_planned() && handoff.from_leg == handoff.to_leg;
1198        if (handoff.from_leg == handoff.to_leg && !terminal_seizure)
1199            || handoff.from_custodian.trim().is_empty()
1200            || handoff.to_custodian.trim().is_empty()
1201            || handoff.location.trim().is_empty()
1202        {
1203            return Err(TransportError::InvalidHandoff(
1204                "handoff requires distinct legs, custodians, and location".to_owned(),
1205            ));
1206        }
1207        if let HandoffKind::Seizure {
1208            by: EntityRef::Domain(record),
1209        } = &handoff.kind
1210            && !domain_record_ref_is_well_formed(record)
1211        {
1212            return Err(TransportError::InvalidHandoff(
1213                "seizure handoff requires a well-formed seizing identity".to_owned(),
1214            ));
1215        }
1216        if self
1217            .handoffs
1218            .iter()
1219            .any(|existing| existing.id == handoff.id)
1220        {
1221            return Err(TransportError::InvalidHandoff(
1222                "handoff identity is already recorded".to_owned(),
1223            ));
1224        }
1225        if self.custody_left_itinerary() {
1226            return Err(TransportError::InvalidHandoff(
1227                "custody already left the itinerary by a terminal seizure".to_owned(),
1228            ));
1229        }
1230        let from_leg = self
1231            .legs
1232            .iter()
1233            .find(|leg| leg.id == handoff.from_leg)
1234            .ok_or(TransportError::MissingLeg)?;
1235        if terminal_seizure {
1236            if self.state.is_terminal()
1237                || Some(from_leg.itinerary_revision) != self.active_itinerary_revision
1238                || from_leg.leg_index != self.current_leg_index
1239                || from_leg.status != LegExecutionStatus::Failed
1240                || from_leg
1241                    .failed_at
1242                    .is_some_and(|failed_at| handoff.at < failed_at)
1243                || self
1244                    .handoffs
1245                    .iter()
1246                    .any(|existing| existing.from_leg == handoff.from_leg)
1247            {
1248                return Err(TransportError::InvalidHandoff(
1249                    "a terminal seizure must follow the failed current leg of a live execution, which no other handoff leaves"
1250                        .to_owned(),
1251                ));
1252            }
1253            self.handoffs.push(handoff);
1254            return Ok(());
1255        }
1256        let to_leg = self
1257            .legs
1258            .iter()
1259            .find(|leg| leg.id == handoff.to_leg)
1260            .ok_or(TransportError::MissingLeg)?;
1261        if !matches!(
1262            from_leg.status,
1263            LegExecutionStatus::Arrived | LegExecutionStatus::Failed
1264        ) || !matches!(
1265            to_leg.status,
1266            LegExecutionStatus::Planned | LegExecutionStatus::Booked | LegExecutionStatus::Waiting
1267        ) || from_leg
1268            .actual_arrival_at
1269            .is_some_and(|arrived_at| handoff.at < arrived_at)
1270            || from_leg
1271                .failed_at
1272                .is_some_and(|failed_at| handoff.at < failed_at)
1273        {
1274            return Err(TransportError::InvalidHandoff(
1275                "handoff must follow an arrived leg and precede the next leg".to_owned(),
1276            ));
1277        }
1278        self.handoffs.push(handoff);
1279        Ok(())
1280    }
1281
1282    /// Whether a terminal seizure has taken custody out of the itinerary.
1283    ///
1284    /// Such an execution cannot fail or reroute a leg, arrive, settle
1285    /// successfully, book capacity, or record another handoff. Its owner
1286    /// closes it: an execution with a delivery saga through
1287    /// [`Self::reconcile_information`] with a failure, which leaves it
1288    /// `Failed`, and one without a saga through [`Self::cancel`].
1289    #[must_use]
1290    pub fn custody_left_itinerary(&self) -> bool {
1291        self.handoffs
1292            .iter()
1293            .any(|handoff| handoff.from_leg == handoff.to_leg && !handoff.kind.is_planned())
1294    }
1295
1296    fn ensure_custody_in_itinerary(&self) -> Result<(), TransportError> {
1297        if self.custody_left_itinerary() {
1298            return Err(TransportError::InvalidState(
1299                "custody left the itinerary by a terminal seizure; the execution can only be closed"
1300                    .to_owned(),
1301            ));
1302        }
1303        Ok(())
1304    }
1305
1306    /// Adds one requested capacity booking to this execution.
1307    ///
1308    /// The booking must name this execution, still be `Requested`, carry no
1309    /// allocation evidence, and have an identity unique within the execution.
1310    /// A terminal execution accepts no further bookings.
1311    pub fn request_booking(&mut self, booking: CapacityBooking) -> Result<(), TransportError> {
1312        self.ensure_custody_in_itinerary()?;
1313        if self.state.is_terminal() {
1314            return Err(TransportError::InvalidState(
1315                "a terminal transport execution cannot request capacity".to_owned(),
1316            ));
1317        }
1318        if booking.execution != self.id
1319            || booking.status != CapacityBookingStatus::Requested
1320            || !booking.allocation_evidence.is_empty()
1321            || booking.quantity == 0
1322            || booking.valid_until < booking.valid_from
1323        {
1324            return Err(TransportError::InvalidBooking(
1325                "a booking request must name this execution, be requested, and carry a positive quantity and window"
1326                    .to_owned(),
1327            ));
1328        }
1329        if self
1330            .bookings
1331            .iter()
1332            .any(|existing| existing.id == booking.id)
1333        {
1334            return Err(TransportError::InvalidBooking(
1335                "capacity booking identity is already recorded on this execution".to_owned(),
1336            ));
1337        }
1338        self.bookings.push(booking);
1339        Ok(())
1340    }
1341
1342    /// Completes the final departed leg of an execution that carries no
1343    /// delivery attempt and settles it.
1344    ///
1345    /// A movement that does not complete an information delivery has nothing
1346    /// to reconcile, so arrival is its terminal fact. An execution with a
1347    /// delivery attempt or saga must instead use
1348    /// [`Self::complete_current_leg`], which enters `ArrivalPending`, and then
1349    /// [`Self::reconcile_information`].
1350    pub fn settle_arrival(&mut self, at: SimTime, endpoint: String) -> Result<(), TransportError> {
1351        if self.delivery_attempt.is_some() || self.saga.is_some() {
1352            return Err(TransportError::InvalidState(
1353                "an execution with a delivery attempt settles through reconciliation".to_owned(),
1354            ));
1355        }
1356        let active = self
1357            .active_itinerary_revision
1358            .ok_or(TransportError::MissingItinerary)?;
1359        let active_leg_count = self
1360            .legs
1361            .iter()
1362            .filter(|leg| leg.itinerary_revision == active)
1363            .count();
1364        let next_leg_index = self
1365            .current_leg_index
1366            .checked_add(1)
1367            .ok_or(TransportError::Overflow)?;
1368        if next_leg_index != active_leg_count {
1369            return Err(TransportError::InvalidState(
1370                "only the final leg of the active itinerary can settle an arrival".to_owned(),
1371            ));
1372        }
1373        let leg = self
1374            .legs
1375            .iter_mut()
1376            .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
1377            .ok_or(TransportError::MissingLeg)?;
1378        if leg.status != LegExecutionStatus::Departed
1379            || leg
1380                .actual_departure_at
1381                .is_some_and(|departure| at < departure)
1382        {
1383            return Err(TransportError::InvalidState(
1384                "the final leg must have departed no later than its arrival".to_owned(),
1385            ));
1386        }
1387        leg.status = LegExecutionStatus::Arrived;
1388        leg.actual_arrival_at = Some(at);
1389        self.current_endpoint = Some(endpoint);
1390        self.current_leg_index = next_leg_index;
1391        self.state = TransportExecutionState::Settled;
1392        Ok(())
1393    }
1394
1395    /// Cancels a non-terminal execution whose subject is not travelling.
1396    ///
1397    /// No leg of the active itinerary may be departed (a leg in progress must
1398    /// first arrive or fail) and an arrival-pending execution must reconcile
1399    /// instead. Unstarted legs of the active itinerary become `Cancelled`; a
1400    /// delivery saga enters `CompensationPending`. Bookings are not changed:
1401    /// the caller closes them through [`CapacityBooking::transition`].
1402    pub fn cancel(&mut self) -> Result<(), TransportError> {
1403        if self.state.is_terminal() || self.state == TransportExecutionState::ArrivalPending {
1404            return Err(TransportError::InvalidState(
1405                "only a non-terminal execution that has not arrived can be cancelled".to_owned(),
1406            ));
1407        }
1408        let active = self.active_itinerary_revision;
1409        if self.legs.iter().any(|leg| {
1410            Some(leg.itinerary_revision) == active && leg.status == LegExecutionStatus::Departed
1411        }) {
1412            return Err(TransportError::InvalidState(
1413                "a departed leg must arrive or fail before its execution is cancelled".to_owned(),
1414            ));
1415        }
1416        for leg in &mut self.legs {
1417            if Some(leg.itinerary_revision) == active
1418                && matches!(
1419                    leg.status,
1420                    LegExecutionStatus::Planned
1421                        | LegExecutionStatus::Booked
1422                        | LegExecutionStatus::Loaded
1423                        | LegExecutionStatus::Waiting
1424                )
1425            {
1426                leg.status = LegExecutionStatus::Cancelled;
1427            }
1428        }
1429        if let Some(saga) = self.saga.as_mut() {
1430            saga.step = saga.step.checked_add(1).ok_or(TransportError::Overflow)?;
1431            saga.state = SagaState::CompensationPending;
1432            saga.last_error = Some("transport execution cancelled".to_owned());
1433        }
1434        self.state = TransportExecutionState::Cancelled;
1435        Ok(())
1436    }
1437
1438    pub fn completion_request(&self) -> Result<DeliveryCompletionRequest, TransportError> {
1439        let saga = self.saga.as_ref().ok_or(TransportError::MissingSaga)?;
1440        if self.state != TransportExecutionState::ArrivalPending
1441            || saga.state != SagaState::ArrivalPending
1442        {
1443            return Err(TransportError::InvalidState(
1444                "delivery completion requires an arrival-pending execution".to_owned(),
1445            ));
1446        }
1447        let revision = self
1448            .active_itinerary_revision
1449            .ok_or(TransportError::MissingItinerary)?;
1450        let completed_at = self
1451            .legs
1452            .iter()
1453            .rev()
1454            .find(|leg| leg.itinerary_revision == revision)
1455            .and_then(|leg| leg.actual_arrival_at)
1456            .ok_or(TransportError::MissingLeg)?;
1457        let attempt = self
1458            .delivery_attempt
1459            .clone()
1460            .ok_or(TransportError::MissingDeliveryAttempt)?;
1461        if attempt != saga.delivery_attempt {
1462            return Err(TransportError::InvalidState(
1463                "saga delivery attempt does not match execution".to_owned(),
1464            ));
1465        }
1466        Ok(DeliveryCompletionRequest {
1467            operation_key: saga.operation_key.clone(),
1468            execution: self.id,
1469            itinerary_revision: revision,
1470            delivery_attempt: attempt,
1471            completed_at,
1472            evidence: saga.evidence.clone(),
1473        })
1474    }
1475
1476    pub fn reconcile_information(
1477        &mut self,
1478        outcome: ReconciliationOutcome,
1479    ) -> Result<(), TransportError> {
1480        if matches!(outcome, ReconciliationOutcome::Success) {
1481            self.ensure_custody_in_itinerary()?;
1482        }
1483        let saga = self.saga.as_mut().ok_or(TransportError::MissingSaga)?;
1484        saga.step = saga.step.checked_add(1).ok_or(TransportError::Overflow)?;
1485        match outcome {
1486            ReconciliationOutcome::Success => {
1487                saga.state = SagaState::Settled;
1488                saga.last_error = None;
1489                self.state = TransportExecutionState::Settled;
1490            }
1491            ReconciliationOutcome::Failure { error } => {
1492                saga.state = SagaState::CompensationPending;
1493                saga.last_error = Some(error);
1494                self.state = TransportExecutionState::Failed;
1495            }
1496        }
1497        Ok(())
1498    }
1499}
1500
1501#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
1502pub enum TransportError {
1503    InvalidRevision(String),
1504    MissingItinerary,
1505    MissingSaga,
1506    SagaAlreadyExists,
1507    Overflow,
1508    MissingLeg,
1509    InvalidState(String),
1510    InvalidBooking(String),
1511    InvalidHandoff(String),
1512    MissingDeliveryAttempt,
1513    /// A capacity pool, a pool revision, or an allocation input is invalid.
1514    InvalidCapacityPool(String),
1515}
1516
1517#[cfg(test)]
1518mod tests {
1519    use super::*;
1520    use canwu_core::{PersonId, TerritoryId};
1521    use canwu_routing::{
1522        ROUTING_ALGORITHM_VERSION, RouteCost, RouteLeg, RoutingConnectionRef, RoutingNodeRef,
1523        TransferMode,
1524    };
1525
1526    fn plan() -> RoutePlan {
1527        RoutePlan {
1528            algorithm_version: ROUTING_ALGORITHM_VERSION.to_owned(),
1529            policy_version: "policy.v1".to_owned(),
1530            planning_snapshot_digest: "snapshot".to_owned(),
1531            origin: RoutingNodeRef::new("a"),
1532            destination: RoutingNodeRef::new("b"),
1533            departure_at: canwu_time::SimTime::EPOCH,
1534            estimated_arrival_at: canwu_time::SimTime::from_minutes(10),
1535            cost: RouteCost {
1536                estimated_arrival_at: canwu_time::SimTime::from_minutes(10),
1537                risk_per_mille: 0,
1538                resource_cost: 0,
1539                transfers: 1,
1540            },
1541            legs: vec![RouteLeg {
1542                connection: RoutingConnectionRef::new("ab"),
1543                from: RoutingNodeRef::new("a"),
1544                to: RoutingNodeRef::new("b"),
1545                mode: TransferMode::Horse,
1546                planned_departure_at: canwu_time::SimTime::EPOCH,
1547                planned_arrival_at: canwu_time::SimTime::from_minutes(10),
1548            }],
1549            digest: "route".to_owned(),
1550        }
1551    }
1552
1553    #[test]
1554    fn movement_order_accepts_a_self_directed_person_subject() {
1555        let order = MovementOrder {
1556            id: MovementOrderId(1),
1557            subjects: vec![MovementSubject {
1558                entity: EntityRef::Person(PersonId::new(7)),
1559                role: MovementSubjectRole::MovablePrincipal,
1560                quantity: None,
1561                expected_custody: None,
1562            }],
1563            origin: RoutingNodeRef::new("a"),
1564            destination: RoutingNodeRef::new("b"),
1565            plan: plan(),
1566            initiative: MovementInitiative::SelfDirected,
1567            ordered_at: SimTime::EPOCH,
1568            expected_position_revision: 1,
1569        };
1570        order.validate().unwrap();
1571    }
1572
1573    #[test]
1574    fn movement_order_rejects_duplicate_subjects_and_missing_cargo_quantity() {
1575        let mut order = MovementOrder {
1576            id: MovementOrderId(2),
1577            subjects: vec![
1578                MovementSubject {
1579                    entity: EntityRef::Person(PersonId::new(7)),
1580                    role: MovementSubjectRole::MovablePrincipal,
1581                    quantity: None,
1582                    expected_custody: None,
1583                },
1584                MovementSubject {
1585                    entity: EntityRef::Person(PersonId::new(7)),
1586                    role: MovementSubjectRole::Cargo,
1587                    quantity: None,
1588                    expected_custody: None,
1589                },
1590            ],
1591            origin: RoutingNodeRef::new("a"),
1592            destination: RoutingNodeRef::new("b"),
1593            plan: plan(),
1594            initiative: MovementInitiative::Delegated,
1595            ordered_at: SimTime::EPOCH,
1596            expected_position_revision: 1,
1597        };
1598        assert!(matches!(
1599            order.validate(),
1600            Err(MovementOrderError::Invalid(_))
1601        ));
1602        order.subjects[1].entity = EntityRef::Territory(TerritoryId::new(8));
1603        assert!(matches!(
1604            order.validate(),
1605            Err(MovementOrderError::Invalid(_))
1606        ));
1607    }
1608
1609    #[test]
1610    fn movement_order_rejects_a_plan_that_does_not_match_the_order() {
1611        let mut order = MovementOrder {
1612            id: MovementOrderId(3),
1613            subjects: vec![MovementSubject {
1614                entity: EntityRef::Person(PersonId::new(9)),
1615                role: MovementSubjectRole::MovablePrincipal,
1616                quantity: None,
1617                expected_custody: None,
1618            }],
1619            origin: RoutingNodeRef::new("a"),
1620            destination: RoutingNodeRef::new("c"),
1621            plan: plan(),
1622            initiative: MovementInitiative::SelfDirected,
1623            ordered_at: SimTime::EPOCH,
1624            expected_position_revision: 1,
1625        };
1626        assert!(matches!(
1627            order.validate(),
1628            Err(MovementOrderError::Invalid(_))
1629        ));
1630        order.destination = RoutingNodeRef::new("b");
1631        order.plan.digest = String::new();
1632        assert!(matches!(
1633            order.validate(),
1634            Err(MovementOrderError::Invalid(_))
1635        ));
1636    }
1637
1638    #[test]
1639    fn reroute_supersedes_without_creating_a_new_delivery_attempt() {
1640        let mut execution = TransportExecution::new(TransportExecutionId(1), None);
1641        execution
1642            .install_initial_itinerary(ItineraryRevision {
1643                id: ItineraryRevisionId(1),
1644                predecessor: None,
1645                plan: plan(),
1646                planned_at: canwu_time::SimTime::EPOCH,
1647                valid_from: canwu_time::SimTime::EPOCH,
1648                reason: ItineraryRevisionReason::Initial,
1649                superseded_at: None,
1650                evidence: Vec::new(),
1651            })
1652            .unwrap();
1653        let mut replacement = plan();
1654        replacement.destination = RoutingNodeRef::new("b");
1655        execution
1656            .reroute(
1657                ItineraryRevision {
1658                    id: ItineraryRevisionId(2),
1659                    predecessor: Some(ItineraryRevisionId(1)),
1660                    plan: replacement,
1661                    planned_at: canwu_time::SimTime::from_minutes(1),
1662                    valid_from: canwu_time::SimTime::from_minutes(1),
1663                    reason: ItineraryRevisionReason::Disaster {
1664                        explanation: "bridge closed".to_owned(),
1665                    },
1666                    superseded_at: None,
1667                    evidence: Vec::new(),
1668                },
1669                canwu_time::SimTime::from_minutes(1),
1670            )
1671            .unwrap();
1672        assert_eq!(execution.revisions.len(), 2);
1673        assert_eq!(
1674            execution.active_itinerary_revision,
1675            Some(ItineraryRevisionId(2))
1676        );
1677        assert_eq!(execution.state, TransportExecutionState::Planning);
1678    }
1679
1680    #[test]
1681    fn capacity_booking_is_a_persisted_windowed_state_machine() {
1682        let mut booking = CapacityBooking::new(
1683            CapacityBookingId(1),
1684            TransportExecutionId(1),
1685            "relay-horse:wu-xi:01".to_owned(),
1686            canwu_time::SimTime::EPOCH,
1687            canwu_time::SimTime::from_minutes(60),
1688            1,
1689            10,
1690        )
1691        .unwrap();
1692        booking
1693            .transition(CapacityBookingStatus::Confirmed, canwu_time::SimTime::EPOCH)
1694            .unwrap();
1695        booking
1696            .transition(
1697                CapacityBookingStatus::Consumed,
1698                canwu_time::SimTime::from_minutes(10),
1699            )
1700            .unwrap();
1701        assert_eq!(booking.status, CapacityBookingStatus::Consumed);
1702    }
1703
1704    #[test]
1705    fn completion_operation_key_is_stable_and_revision_scoped() {
1706        let first =
1707            delivery_completion_operation_key(TransportExecutionId(4), ItineraryRevisionId(2), 7);
1708        let second =
1709            delivery_completion_operation_key(TransportExecutionId(4), ItineraryRevisionId(3), 7);
1710        assert_ne!(first, second);
1711        assert_eq!(
1712            first,
1713            delivery_completion_operation_key(TransportExecutionId(4), ItineraryRevisionId(2), 7)
1714        );
1715    }
1716
1717    #[test]
1718    fn completion_requires_saga_and_exposes_stable_bridge_request() {
1719        let attempt = DomainRecordVersionRef {
1720            record: canwu_core::DomainRecordRef::new(
1721                "fixture.information",
1722                "delivery_attempt",
1723                "delivery",
1724            ),
1725            version: 3,
1726            established_by: canwu_core::DomainRecordVersionSource::InitialScenario,
1727        };
1728        let mut execution = TransportExecution::new(TransportExecutionId(9), Some(attempt.clone()));
1729        let revision = ItineraryRevision {
1730            id: ItineraryRevisionId(1),
1731            predecessor: None,
1732            plan: plan(),
1733            planned_at: SimTime::EPOCH,
1734            valid_from: SimTime::EPOCH,
1735            reason: ItineraryRevisionReason::Initial,
1736            superseded_at: None,
1737            evidence: Vec::new(),
1738        };
1739        execution.install_initial_itinerary(revision).unwrap();
1740        assert_eq!(
1741            execution.complete_current_leg(SimTime::from_minutes(10), "b".to_owned()),
1742            Err(TransportError::MissingSaga)
1743        );
1744        execution
1745            .begin_saga(
1746                attempt,
1747                delivery_completion_operation_key(
1748                    TransportExecutionId(9),
1749                    ItineraryRevisionId(1),
1750                    3,
1751                ),
1752            )
1753            .unwrap();
1754        execution.start_current_leg(SimTime::EPOCH).unwrap();
1755        assert!(
1756            execution
1757                .complete_current_leg(SimTime::from_minutes(10), "b".to_owned())
1758                .unwrap()
1759        );
1760        let request = execution.completion_request().unwrap();
1761        assert_eq!(request.execution, TransportExecutionId(9));
1762        assert_eq!(request.itinerary_revision, ItineraryRevisionId(1));
1763        assert_eq!(request.delivery_attempt.version, 3);
1764    }
1765
1766    #[test]
1767    fn handoff_requires_arrival_and_is_idempotency_safe() {
1768        let attempt = DomainRecordVersionRef {
1769            record: canwu_core::DomainRecordRef::new(
1770                "fixture.information",
1771                "delivery_attempt",
1772                "delivery",
1773            ),
1774            version: 1,
1775            established_by: canwu_core::DomainRecordVersionSource::InitialScenario,
1776        };
1777        let mut execution =
1778            TransportExecution::new(TransportExecutionId(10), Some(attempt.clone()));
1779        let mut route = plan();
1780        route.legs.push(RouteLeg {
1781            connection: RoutingConnectionRef::new("bc"),
1782            from: RoutingNodeRef::new("b"),
1783            to: RoutingNodeRef::new("c"),
1784            mode: TransferMode::Rail,
1785            planned_departure_at: SimTime::from_minutes(10),
1786            planned_arrival_at: SimTime::from_minutes(20),
1787        });
1788        route.destination = RoutingNodeRef::new("c");
1789        route.estimated_arrival_at = SimTime::from_minutes(20);
1790        let revision = ItineraryRevision {
1791            id: ItineraryRevisionId(1),
1792            predecessor: None,
1793            plan: route,
1794            planned_at: SimTime::EPOCH,
1795            valid_from: SimTime::EPOCH,
1796            reason: ItineraryRevisionReason::Initial,
1797            superseded_at: None,
1798            evidence: Vec::new(),
1799        };
1800        execution.install_initial_itinerary(revision).unwrap();
1801        execution.start_current_leg(SimTime::EPOCH).unwrap();
1802        execution
1803            .complete_current_leg(SimTime::from_minutes(10), "b".to_owned())
1804            .unwrap();
1805        let handoff = Handoff {
1806            id: HandoffId(1),
1807            from_leg: LegExecutionId(1),
1808            to_leg: LegExecutionId(2),
1809            from_custodian: "courier/wuxi".to_owned(),
1810            to_custodian: "rail/beijing".to_owned(),
1811            at: SimTime::from_minutes(10),
1812            location: "b".to_owned(),
1813            evidence: Vec::new(),
1814            kind: HandoffKind::Planned,
1815        };
1816        execution.record_handoff(handoff.clone()).unwrap();
1817        assert_eq!(
1818            execution.record_handoff(handoff),
1819            Err(TransportError::InvalidHandoff(
1820                "handoff identity is already recorded".to_owned()
1821            ))
1822        );
1823    }
1824}