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::{DomainRecordVersionRef, EntityRef, EvidenceRef};
10use canwu_routing::{RoutePlan, RoutingNodeRef};
11use canwu_time::SimTime;
12use serde::{Deserialize, Serialize};
13
14pub const TRANSPORT_SEMANTIC_VERSION: &str = "canwu-transport.v3";
15
16#[must_use]
17pub fn delivery_completion_operation_key(
18    execution: TransportExecutionId,
19    revision: ItineraryRevisionId,
20    attempt_version: u64,
21) -> String {
22    format!(
23        "transport/{}/revision/{}/delivery-completion/attempt-version/{}",
24        execution.0, revision.0, attempt_version
25    )
26}
27
28#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
29#[serde(transparent)]
30pub struct TransportExecutionId(pub u64);
31
32/// Stable identity for an admitted transport-domain movement intent.
33///
34/// The simulation admits movement through a host- or integration-defined command.
35/// This record adds route-plan and custody evidence when a transport domain needs
36/// to persist a richer execution than its own projected transit state.
37#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
38#[serde(transparent)]
39pub struct MovementOrderId(pub u64);
40
41/// Who initiated a movement intent. The runtime must derive this from the
42/// admitted authority and never trust an unvalidated caller-supplied label.
43#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
44#[serde(rename_all = "snake_case")]
45pub enum MovementInitiative {
46    SelfDirected,
47    Commanded,
48    Delegated,
49    Forced,
50    Automatic,
51}
52
53/// The physical role of a subject in a movement manifest.
54#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
55#[serde(rename_all = "snake_case")]
56pub enum MovementSubjectRole {
57    MovablePrincipal,
58    Cargo,
59    Carrier,
60    Passenger,
61    Attached,
62}
63
64/// One typed identity in a movement manifest.
65#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
66pub struct MovementSubject {
67    pub entity: EntityRef,
68    pub role: MovementSubjectRole,
69    /// Cargo quantities are integer units and must be positive. Other roles
70    /// normally leave this unset because their cardinality is one identity.
71    #[serde(default, skip_serializing_if = "Option::is_none")]
72    pub quantity: Option<u64>,
73    /// Expected carrier/custodian identity at admission, when applicable.
74    #[serde(default, skip_serializing_if = "Option::is_none")]
75    pub expected_custody: Option<EntityRef>,
76}
77
78/// Immutable, admitted intent shared by transport movement domains.
79///
80/// `MovementOrder` is a contract for planning and authority evidence; it does
81/// not directly mutate a world entity. Domain handlers still own location,
82/// custody, quantity, arrival, and knowledge effects.
83#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
84pub struct MovementOrder {
85    pub id: MovementOrderId,
86    pub subjects: Vec<MovementSubject>,
87    pub origin: RoutingNodeRef,
88    pub destination: RoutingNodeRef,
89    pub plan: RoutePlan,
90    pub initiative: MovementInitiative,
91    pub ordered_at: SimTime,
92    pub expected_position_revision: u64,
93}
94
95#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
96pub enum MovementOrderError {
97    Invalid(String),
98}
99
100impl MovementOrder {
101    /// Validates the structural invariants that are domain-neutral.
102    ///
103    /// Existence, authority, capability, and position/custody matching remain
104    /// runtime or plugin responsibilities because this crate does not own the
105    /// world or domain records.
106    pub fn validate(&self) -> Result<(), MovementOrderError> {
107        if self.id.0 == 0
108            || self.expected_position_revision == 0
109            || self.origin.as_str().trim().is_empty()
110            || self.destination.as_str().trim().is_empty()
111        {
112            return Err(MovementOrderError::Invalid(
113                "movement order identity, endpoints, and expected position revision must be valid"
114                    .to_owned(),
115            ));
116        }
117        if self.subjects.is_empty()
118            || self
119                .subjects
120                .windows(2)
121                .any(|pair| pair[0].entity >= pair[1].entity)
122        {
123            return Err(MovementOrderError::Invalid(
124                "movement subjects must be non-empty, sorted, and unique by entity".to_owned(),
125            ));
126        }
127        for subject in &self.subjects {
128            if subject.quantity.is_some_and(|quantity| quantity == 0)
129                || (subject.role == MovementSubjectRole::Cargo && subject.quantity.is_none())
130                || (subject.role != MovementSubjectRole::Cargo && subject.quantity.is_some())
131            {
132                return Err(MovementOrderError::Invalid(
133                    "cargo requires a positive quantity and non-cargo subjects cannot carry one"
134                        .to_owned(),
135                ));
136            }
137        }
138        if self.plan.origin != self.origin
139            || self.plan.destination != self.destination
140            || self.plan.departure_at < self.ordered_at
141            || self.plan.estimated_arrival_at < self.plan.departure_at
142            || self.plan.digest.trim().is_empty()
143        {
144            return Err(MovementOrderError::Invalid(
145                "movement plan does not match the order endpoints or time range".to_owned(),
146            ));
147        }
148        Ok(())
149    }
150}
151
152#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
153#[serde(transparent)]
154pub struct ItineraryRevisionId(pub u64);
155
156#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
157#[serde(transparent)]
158pub struct LegExecutionId(pub u64);
159
160#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
161#[serde(transparent)]
162pub struct HandoffId(pub u64);
163
164#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
165#[serde(transparent)]
166pub struct CapacityBookingId(pub u64);
167
168#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
169#[serde(rename_all = "snake_case")]
170pub enum TransportExecutionState {
171    Prepared,
172    Planning,
173    Booking,
174    Ready,
175    Executing,
176    ReplanPending,
177    ArrivalPending,
178    Settled,
179    Failed,
180    Cancelled,
181}
182
183impl TransportExecutionState {
184    #[must_use]
185    pub const fn is_terminal(self) -> bool {
186        matches!(self, Self::Settled | Self::Failed | Self::Cancelled)
187    }
188}
189
190#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
191#[serde(rename_all = "snake_case")]
192pub enum ItineraryRevisionReason {
193    Initial,
194    Disaster { explanation: String },
195    CapacityUnavailable { explanation: String },
196    KnowledgeUpdate { explanation: String },
197    Recovery { explanation: String },
198}
199
200#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
201pub struct ItineraryRevision {
202    pub id: ItineraryRevisionId,
203    pub predecessor: Option<ItineraryRevisionId>,
204    pub plan: RoutePlan,
205    pub planned_at: SimTime,
206    pub valid_from: SimTime,
207    pub reason: ItineraryRevisionReason,
208    pub superseded_at: Option<SimTime>,
209    pub evidence: Vec<EvidenceRef>,
210}
211
212#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
213#[serde(rename_all = "snake_case")]
214pub enum LegExecutionStatus {
215    Planned,
216    Booked,
217    Loaded,
218    Departed,
219    Arrived,
220    Waiting,
221    Failed,
222    Cancelled,
223}
224
225#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
226pub struct LegExecution {
227    pub id: LegExecutionId,
228    pub itinerary_revision: ItineraryRevisionId,
229    pub leg_index: usize,
230    pub status: LegExecutionStatus,
231    pub actual_departure_at: Option<SimTime>,
232    pub actual_arrival_at: Option<SimTime>,
233    pub failed_at: Option<SimTime>,
234    pub failure_reason: Option<String>,
235    pub evidence: Vec<EvidenceRef>,
236}
237
238#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
239pub struct Handoff {
240    pub id: HandoffId,
241    pub from_leg: LegExecutionId,
242    pub to_leg: LegExecutionId,
243    pub from_custodian: String,
244    pub to_custodian: String,
245    pub at: SimTime,
246    pub location: String,
247    pub evidence: Vec<EvidenceRef>,
248}
249
250#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
251pub struct DeliveryCompletionRequest {
252    pub operation_key: String,
253    pub execution: TransportExecutionId,
254    pub itinerary_revision: ItineraryRevisionId,
255    pub delivery_attempt: DomainRecordVersionRef,
256    pub completed_at: SimTime,
257    pub evidence: Vec<EvidenceRef>,
258}
259
260#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
261#[serde(rename_all = "snake_case")]
262pub enum CapacityBookingStatus {
263    Requested,
264    Confirmed,
265    Consumed,
266    Released,
267    Expired,
268    Cancelled,
269    Failed,
270}
271
272#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
273pub struct CapacityBooking {
274    pub id: CapacityBookingId,
275    pub execution: TransportExecutionId,
276    pub resource: String,
277    pub valid_from: SimTime,
278    pub valid_until: SimTime,
279    pub quantity: u64,
280    pub priority: i32,
281    pub status: CapacityBookingStatus,
282    pub allocation_evidence: Vec<EvidenceRef>,
283}
284
285impl CapacityBooking {
286    pub fn new(
287        id: CapacityBookingId,
288        execution: TransportExecutionId,
289        resource: String,
290        valid_from: SimTime,
291        valid_until: SimTime,
292        quantity: u64,
293        priority: i32,
294    ) -> Result<Self, TransportError> {
295        if valid_until < valid_from || quantity == 0 {
296            return Err(TransportError::InvalidBooking(
297                "capacity booking requires a positive quantity and non-inverted window".to_owned(),
298            ));
299        }
300        Ok(Self {
301            id,
302            execution,
303            resource,
304            valid_from,
305            valid_until,
306            quantity,
307            priority,
308            status: CapacityBookingStatus::Requested,
309            allocation_evidence: Vec::new(),
310        })
311    }
312
313    pub fn transition(
314        &mut self,
315        status: CapacityBookingStatus,
316        at: SimTime,
317    ) -> Result<(), TransportError> {
318        if at < self.valid_from {
319            return Err(TransportError::InvalidBooking(
320                "booking cannot transition before its validity window".to_owned(),
321            ));
322        }
323        let allowed = matches!(
324            (self.status, status),
325            (
326                CapacityBookingStatus::Requested,
327                CapacityBookingStatus::Confirmed
328                    | CapacityBookingStatus::Failed
329                    | CapacityBookingStatus::Cancelled
330            ) | (
331                CapacityBookingStatus::Confirmed,
332                CapacityBookingStatus::Consumed
333                    | CapacityBookingStatus::Released
334                    | CapacityBookingStatus::Cancelled
335                    | CapacityBookingStatus::Expired
336            ) | (
337                CapacityBookingStatus::Consumed,
338                CapacityBookingStatus::Released
339            )
340        );
341        if !allowed {
342            return Err(TransportError::InvalidBooking(
343                "capacity booking transition is not allowed".to_owned(),
344            ));
345        }
346        if at > self.valid_until && status == CapacityBookingStatus::Confirmed {
347            return Err(TransportError::InvalidBooking(
348                "capacity booking cannot be confirmed after its validity window".to_owned(),
349            ));
350        }
351        if status == CapacityBookingStatus::Expired && at <= self.valid_until {
352            return Err(TransportError::InvalidBooking(
353                "capacity booking cannot expire before its validity window ends".to_owned(),
354            ));
355        }
356        self.status = status;
357        Ok(())
358    }
359}
360
361#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
362#[serde(rename_all = "snake_case")]
363pub enum SagaState {
364    TransportIntent,
365    WaitingForInformation,
366    Executing,
367    ArrivalPending,
368    Settled,
369    CompensationPending,
370    Failed,
371}
372
373/// Result of reconciling the information-system delivery attempt.
374#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
375pub enum ReconciliationOutcome {
376    Success,
377    Failure { error: String },
378}
379
380#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
381pub struct DeliverySaga {
382    pub operation_key: String,
383    pub delivery_attempt: DomainRecordVersionRef,
384    pub state: SagaState,
385    pub step: u32,
386    pub expected_attempt_version: u64,
387    pub last_error: Option<String>,
388    pub evidence: Vec<EvidenceRef>,
389}
390
391#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
392pub struct TransportExecution {
393    pub id: TransportExecutionId,
394    pub delivery_attempt: Option<DomainRecordVersionRef>,
395    pub state: TransportExecutionState,
396    pub active_itinerary_revision: Option<ItineraryRevisionId>,
397    pub current_leg_index: usize,
398    pub estimated_arrival_at: Option<SimTime>,
399    pub current_endpoint: Option<String>,
400    pub revisions: Vec<ItineraryRevision>,
401    pub legs: Vec<LegExecution>,
402    pub handoffs: Vec<Handoff>,
403    pub bookings: Vec<CapacityBooking>,
404    pub saga: Option<DeliverySaga>,
405}
406
407impl TransportExecution {
408    #[must_use]
409    pub fn new(id: TransportExecutionId, delivery_attempt: Option<DomainRecordVersionRef>) -> Self {
410        Self {
411            id,
412            delivery_attempt,
413            state: TransportExecutionState::Prepared,
414            active_itinerary_revision: None,
415            current_leg_index: 0,
416            estimated_arrival_at: None,
417            current_endpoint: None,
418            revisions: Vec::new(),
419            legs: Vec::new(),
420            handoffs: Vec::new(),
421            bookings: Vec::new(),
422            saga: None,
423        }
424    }
425
426    pub fn install_initial_itinerary(
427        &mut self,
428        revision: ItineraryRevision,
429    ) -> Result<(), TransportError> {
430        if !self.revisions.is_empty() || revision.predecessor.is_some() {
431            return Err(TransportError::InvalidRevision(
432                "initial itinerary must be the first revision".to_owned(),
433            ));
434        }
435        self.estimated_arrival_at = Some(revision.plan.estimated_arrival_at);
436        self.active_itinerary_revision = Some(revision.id);
437        self.current_endpoint = Some(revision.plan.origin.as_str().to_owned());
438        self.legs = revision
439            .plan
440            .legs
441            .iter()
442            .enumerate()
443            .map(|(index, _)| {
444                let id = u64::try_from(index)
445                    .ok()
446                    .and_then(|index| index.checked_add(1))
447                    .ok_or(TransportError::Overflow)?;
448                Ok(LegExecution {
449                    id: LegExecutionId(id),
450                    itinerary_revision: revision.id,
451                    leg_index: index,
452                    status: LegExecutionStatus::Planned,
453                    actual_departure_at: None,
454                    actual_arrival_at: None,
455                    failed_at: None,
456                    failure_reason: None,
457                    evidence: Vec::new(),
458                })
459            })
460            .collect::<Result<Vec<_>, TransportError>>()?;
461        self.revisions.push(revision);
462        self.state = TransportExecutionState::Planning;
463        Ok(())
464    }
465
466    pub fn reroute(
467        &mut self,
468        revision: ItineraryRevision,
469        at: SimTime,
470    ) -> Result<(), TransportError> {
471        let active = self
472            .active_itinerary_revision
473            .ok_or(TransportError::MissingItinerary)?;
474        if revision.predecessor != Some(active) || revision.valid_from < at {
475            return Err(TransportError::InvalidRevision("reroute must reference the active revision and start no earlier than the reroute time".to_owned()));
476        }
477        if self.revisions.iter().any(|item| item.id == revision.id) {
478            return Err(TransportError::InvalidRevision(
479                "itinerary revision identity must be unique within an execution".to_owned(),
480            ));
481        }
482        if let Some(previous) = self.revisions.iter_mut().find(|item| item.id == active) {
483            previous.superseded_at = Some(at);
484        }
485        self.estimated_arrival_at = Some(revision.plan.estimated_arrival_at);
486        self.active_itinerary_revision = Some(revision.id);
487        self.current_leg_index = 0;
488        let next_leg_id = self
489            .legs
490            .iter()
491            .map(|leg| leg.id.0)
492            .max()
493            .unwrap_or_default()
494            .checked_add(1)
495            .ok_or(TransportError::Overflow)?;
496        let mut legs = Vec::with_capacity(revision.plan.legs.len());
497        for (index, _) in revision.plan.legs.iter().enumerate() {
498            let id = u64::try_from(index)
499                .ok()
500                .and_then(|index| next_leg_id.checked_add(index))
501                .ok_or(TransportError::Overflow)?;
502            legs.push(LegExecution {
503                id: LegExecutionId(id),
504                itinerary_revision: revision.id,
505                leg_index: index,
506                status: LegExecutionStatus::Planned,
507                actual_departure_at: None,
508                actual_arrival_at: None,
509                failed_at: None,
510                failure_reason: None,
511                evidence: Vec::new(),
512            });
513        }
514        self.legs.extend(legs);
515        self.revisions.push(revision);
516        if let Some(saga) = self.saga.as_mut() {
517            saga.operation_key = delivery_completion_operation_key(
518                self.id,
519                self.active_itinerary_revision
520                    .ok_or(TransportError::MissingItinerary)?,
521                saga.expected_attempt_version,
522            );
523            saga.evidence.extend(
524                self.revisions
525                    .last()
526                    .map(|current| current.evidence.clone())
527                    .unwrap_or_default(),
528            );
529        }
530        self.state = TransportExecutionState::Planning;
531        Ok(())
532    }
533
534    pub fn begin_saga(
535        &mut self,
536        delivery_attempt: DomainRecordVersionRef,
537        operation_key: String,
538    ) -> Result<(), TransportError> {
539        if self.saga.is_some() {
540            return Err(TransportError::SagaAlreadyExists);
541        }
542        self.saga = Some(DeliverySaga {
543            expected_attempt_version: delivery_attempt.version,
544            operation_key,
545            delivery_attempt,
546            state: SagaState::TransportIntent,
547            step: 0,
548            last_error: None,
549            evidence: Vec::new(),
550        });
551        self.state = TransportExecutionState::Executing;
552        Ok(())
553    }
554
555    pub fn start_current_leg(&mut self, at: SimTime) -> Result<(), TransportError> {
556        if self.state != TransportExecutionState::Ready
557            && self.state != TransportExecutionState::Executing
558            && self.state != TransportExecutionState::Planning
559        {
560            return Err(TransportError::InvalidState(
561                "transport execution cannot start a leg in its current state".to_owned(),
562            ));
563        }
564        let active = self
565            .active_itinerary_revision
566            .ok_or(TransportError::MissingItinerary)?;
567        let leg = self
568            .legs
569            .iter_mut()
570            .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
571            .ok_or(TransportError::MissingLeg)?;
572        if !matches!(
573            leg.status,
574            LegExecutionStatus::Planned | LegExecutionStatus::Booked | LegExecutionStatus::Waiting
575        ) {
576            return Err(TransportError::InvalidState(
577                "current leg is not startable".to_owned(),
578            ));
579        }
580        leg.status = LegExecutionStatus::Departed;
581        leg.actual_departure_at = Some(at);
582        self.state = TransportExecutionState::Executing;
583        Ok(())
584    }
585
586    pub fn complete_current_leg(
587        &mut self,
588        at: SimTime,
589        endpoint: String,
590    ) -> Result<bool, TransportError> {
591        let active = self
592            .active_itinerary_revision
593            .ok_or(TransportError::MissingItinerary)?;
594        let active_leg_count = self
595            .legs
596            .iter()
597            .filter(|leg| leg.itinerary_revision == active)
598            .count();
599        let next_leg_index = self
600            .current_leg_index
601            .checked_add(1)
602            .ok_or(TransportError::Overflow)?;
603        let final_leg = next_leg_index >= active_leg_count;
604        if final_leg && self.saga.is_none() {
605            return Err(TransportError::MissingSaga);
606        }
607        let leg = self
608            .legs
609            .iter_mut()
610            .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
611            .ok_or(TransportError::MissingLeg)?;
612        if leg.status != LegExecutionStatus::Departed {
613            return Err(TransportError::InvalidState(
614                "current leg must be departed before arrival".to_owned(),
615            ));
616        }
617        if leg
618            .actual_departure_at
619            .is_some_and(|departure| at < departure)
620        {
621            return Err(TransportError::InvalidState(
622                "arrival precedes departure".to_owned(),
623            ));
624        }
625        leg.status = LegExecutionStatus::Arrived;
626        leg.actual_arrival_at = Some(at);
627        self.current_endpoint = Some(endpoint);
628        self.current_leg_index = next_leg_index;
629        if self.current_leg_index >= active_leg_count {
630            self.state = TransportExecutionState::ArrivalPending;
631            self.mark_arrival_pending()?;
632            Ok(true)
633        } else {
634            self.state = TransportExecutionState::Ready;
635            Ok(false)
636        }
637    }
638
639    pub fn fail_current_leg(&mut self, reason: String, at: SimTime) -> Result<(), TransportError> {
640        let active = self
641            .active_itinerary_revision
642            .ok_or(TransportError::MissingItinerary)?;
643        let leg = self
644            .legs
645            .iter_mut()
646            .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
647            .ok_or(TransportError::MissingLeg)?;
648        leg.status = LegExecutionStatus::Failed;
649        leg.failed_at = Some(at);
650        leg.failure_reason = Some(reason);
651        self.state = TransportExecutionState::ReplanPending;
652        Ok(())
653    }
654
655    pub fn mark_arrival_pending(&mut self) -> Result<(), TransportError> {
656        let saga = self.saga.as_mut().ok_or(TransportError::MissingSaga)?;
657        saga.step = saga.step.checked_add(1).ok_or(TransportError::Overflow)?;
658        saga.state = SagaState::ArrivalPending;
659        self.state = TransportExecutionState::ArrivalPending;
660        Ok(())
661    }
662
663    pub fn record_handoff(&mut self, handoff: Handoff) -> Result<(), TransportError> {
664        if handoff.from_leg == handoff.to_leg
665            || handoff.from_custodian.trim().is_empty()
666            || handoff.to_custodian.trim().is_empty()
667            || handoff.location.trim().is_empty()
668        {
669            return Err(TransportError::InvalidHandoff(
670                "handoff requires distinct legs, custodians, and location".to_owned(),
671            ));
672        }
673        if self
674            .handoffs
675            .iter()
676            .any(|existing| existing.id == handoff.id)
677        {
678            return Err(TransportError::InvalidHandoff(
679                "handoff identity is already recorded".to_owned(),
680            ));
681        }
682        let from_leg = self
683            .legs
684            .iter()
685            .find(|leg| leg.id == handoff.from_leg)
686            .ok_or(TransportError::MissingLeg)?;
687        let to_leg = self
688            .legs
689            .iter()
690            .find(|leg| leg.id == handoff.to_leg)
691            .ok_or(TransportError::MissingLeg)?;
692        if !matches!(
693            from_leg.status,
694            LegExecutionStatus::Arrived | LegExecutionStatus::Failed
695        ) || !matches!(
696            to_leg.status,
697            LegExecutionStatus::Planned | LegExecutionStatus::Booked | LegExecutionStatus::Waiting
698        ) || from_leg
699            .actual_arrival_at
700            .is_some_and(|arrived_at| handoff.at < arrived_at)
701            || from_leg
702                .failed_at
703                .is_some_and(|failed_at| handoff.at < failed_at)
704        {
705            return Err(TransportError::InvalidHandoff(
706                "handoff must follow an arrived leg and precede the next leg".to_owned(),
707            ));
708        }
709        self.handoffs.push(handoff);
710        Ok(())
711    }
712
713    pub fn completion_request(&self) -> Result<DeliveryCompletionRequest, TransportError> {
714        let saga = self.saga.as_ref().ok_or(TransportError::MissingSaga)?;
715        if self.state != TransportExecutionState::ArrivalPending
716            || saga.state != SagaState::ArrivalPending
717        {
718            return Err(TransportError::InvalidState(
719                "delivery completion requires an arrival-pending execution".to_owned(),
720            ));
721        }
722        let revision = self
723            .active_itinerary_revision
724            .ok_or(TransportError::MissingItinerary)?;
725        let completed_at = self
726            .legs
727            .iter()
728            .rev()
729            .find(|leg| leg.itinerary_revision == revision)
730            .and_then(|leg| leg.actual_arrival_at)
731            .ok_or(TransportError::MissingLeg)?;
732        let attempt = self
733            .delivery_attempt
734            .clone()
735            .ok_or(TransportError::MissingDeliveryAttempt)?;
736        if attempt != saga.delivery_attempt {
737            return Err(TransportError::InvalidState(
738                "saga delivery attempt does not match execution".to_owned(),
739            ));
740        }
741        Ok(DeliveryCompletionRequest {
742            operation_key: saga.operation_key.clone(),
743            execution: self.id,
744            itinerary_revision: revision,
745            delivery_attempt: attempt,
746            completed_at,
747            evidence: saga.evidence.clone(),
748        })
749    }
750
751    pub fn reconcile_information(
752        &mut self,
753        outcome: ReconciliationOutcome,
754    ) -> Result<(), TransportError> {
755        let saga = self.saga.as_mut().ok_or(TransportError::MissingSaga)?;
756        saga.step = saga.step.checked_add(1).ok_or(TransportError::Overflow)?;
757        match outcome {
758            ReconciliationOutcome::Success => {
759                saga.state = SagaState::Settled;
760                saga.last_error = None;
761                self.state = TransportExecutionState::Settled;
762            }
763            ReconciliationOutcome::Failure { error } => {
764                saga.state = SagaState::CompensationPending;
765                saga.last_error = Some(error);
766                self.state = TransportExecutionState::Failed;
767            }
768        }
769        Ok(())
770    }
771}
772
773#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
774pub enum TransportError {
775    InvalidRevision(String),
776    MissingItinerary,
777    MissingSaga,
778    SagaAlreadyExists,
779    Overflow,
780    MissingLeg,
781    InvalidState(String),
782    InvalidBooking(String),
783    InvalidHandoff(String),
784    MissingDeliveryAttempt,
785}
786
787#[cfg(test)]
788mod tests {
789    use super::*;
790    use canwu_core::{PersonId, TerritoryId};
791    use canwu_routing::{
792        ROUTING_ALGORITHM_VERSION, RouteCost, RouteLeg, RoutingConnectionRef, RoutingNodeRef,
793        TransferMode,
794    };
795
796    fn plan() -> RoutePlan {
797        RoutePlan {
798            algorithm_version: ROUTING_ALGORITHM_VERSION.to_owned(),
799            policy_version: "policy.v1".to_owned(),
800            planning_snapshot_digest: "snapshot".to_owned(),
801            origin: RoutingNodeRef::new("a"),
802            destination: RoutingNodeRef::new("b"),
803            departure_at: canwu_time::SimTime::EPOCH,
804            estimated_arrival_at: canwu_time::SimTime::from_minutes(10),
805            cost: RouteCost {
806                estimated_arrival_at: canwu_time::SimTime::from_minutes(10),
807                risk_per_mille: 0,
808                resource_cost: 0,
809                transfers: 1,
810            },
811            legs: vec![RouteLeg {
812                connection: RoutingConnectionRef::new("ab"),
813                from: RoutingNodeRef::new("a"),
814                to: RoutingNodeRef::new("b"),
815                mode: TransferMode::Horse,
816                planned_departure_at: canwu_time::SimTime::EPOCH,
817                planned_arrival_at: canwu_time::SimTime::from_minutes(10),
818            }],
819            digest: "route".to_owned(),
820        }
821    }
822
823    #[test]
824    fn movement_order_accepts_a_self_directed_person_subject() {
825        let order = MovementOrder {
826            id: MovementOrderId(1),
827            subjects: vec![MovementSubject {
828                entity: EntityRef::Person(PersonId::new(7)),
829                role: MovementSubjectRole::MovablePrincipal,
830                quantity: None,
831                expected_custody: None,
832            }],
833            origin: RoutingNodeRef::new("a"),
834            destination: RoutingNodeRef::new("b"),
835            plan: plan(),
836            initiative: MovementInitiative::SelfDirected,
837            ordered_at: SimTime::EPOCH,
838            expected_position_revision: 1,
839        };
840        order.validate().unwrap();
841    }
842
843    #[test]
844    fn movement_order_rejects_duplicate_subjects_and_missing_cargo_quantity() {
845        let mut order = MovementOrder {
846            id: MovementOrderId(2),
847            subjects: vec![
848                MovementSubject {
849                    entity: EntityRef::Person(PersonId::new(7)),
850                    role: MovementSubjectRole::MovablePrincipal,
851                    quantity: None,
852                    expected_custody: None,
853                },
854                MovementSubject {
855                    entity: EntityRef::Person(PersonId::new(7)),
856                    role: MovementSubjectRole::Cargo,
857                    quantity: None,
858                    expected_custody: None,
859                },
860            ],
861            origin: RoutingNodeRef::new("a"),
862            destination: RoutingNodeRef::new("b"),
863            plan: plan(),
864            initiative: MovementInitiative::Delegated,
865            ordered_at: SimTime::EPOCH,
866            expected_position_revision: 1,
867        };
868        assert!(matches!(
869            order.validate(),
870            Err(MovementOrderError::Invalid(_))
871        ));
872        order.subjects[1].entity = EntityRef::Territory(TerritoryId::new(8));
873        assert!(matches!(
874            order.validate(),
875            Err(MovementOrderError::Invalid(_))
876        ));
877    }
878
879    #[test]
880    fn movement_order_rejects_a_plan_that_does_not_match_the_order() {
881        let mut order = MovementOrder {
882            id: MovementOrderId(3),
883            subjects: vec![MovementSubject {
884                entity: EntityRef::Person(PersonId::new(9)),
885                role: MovementSubjectRole::MovablePrincipal,
886                quantity: None,
887                expected_custody: None,
888            }],
889            origin: RoutingNodeRef::new("a"),
890            destination: RoutingNodeRef::new("c"),
891            plan: plan(),
892            initiative: MovementInitiative::SelfDirected,
893            ordered_at: SimTime::EPOCH,
894            expected_position_revision: 1,
895        };
896        assert!(matches!(
897            order.validate(),
898            Err(MovementOrderError::Invalid(_))
899        ));
900        order.destination = RoutingNodeRef::new("b");
901        order.plan.digest = String::new();
902        assert!(matches!(
903            order.validate(),
904            Err(MovementOrderError::Invalid(_))
905        ));
906    }
907
908    #[test]
909    fn reroute_supersedes_without_creating_a_new_delivery_attempt() {
910        let mut execution = TransportExecution::new(TransportExecutionId(1), None);
911        execution
912            .install_initial_itinerary(ItineraryRevision {
913                id: ItineraryRevisionId(1),
914                predecessor: None,
915                plan: plan(),
916                planned_at: canwu_time::SimTime::EPOCH,
917                valid_from: canwu_time::SimTime::EPOCH,
918                reason: ItineraryRevisionReason::Initial,
919                superseded_at: None,
920                evidence: Vec::new(),
921            })
922            .unwrap();
923        let mut replacement = plan();
924        replacement.destination = RoutingNodeRef::new("b");
925        execution
926            .reroute(
927                ItineraryRevision {
928                    id: ItineraryRevisionId(2),
929                    predecessor: Some(ItineraryRevisionId(1)),
930                    plan: replacement,
931                    planned_at: canwu_time::SimTime::from_minutes(1),
932                    valid_from: canwu_time::SimTime::from_minutes(1),
933                    reason: ItineraryRevisionReason::Disaster {
934                        explanation: "bridge closed".to_owned(),
935                    },
936                    superseded_at: None,
937                    evidence: Vec::new(),
938                },
939                canwu_time::SimTime::from_minutes(1),
940            )
941            .unwrap();
942        assert_eq!(execution.revisions.len(), 2);
943        assert_eq!(
944            execution.active_itinerary_revision,
945            Some(ItineraryRevisionId(2))
946        );
947        assert_eq!(execution.state, TransportExecutionState::Planning);
948    }
949
950    #[test]
951    fn capacity_booking_is_a_persisted_windowed_state_machine() {
952        let mut booking = CapacityBooking::new(
953            CapacityBookingId(1),
954            TransportExecutionId(1),
955            "relay-horse:wu-xi:01".to_owned(),
956            canwu_time::SimTime::EPOCH,
957            canwu_time::SimTime::from_minutes(60),
958            1,
959            10,
960        )
961        .unwrap();
962        booking
963            .transition(CapacityBookingStatus::Confirmed, canwu_time::SimTime::EPOCH)
964            .unwrap();
965        booking
966            .transition(
967                CapacityBookingStatus::Consumed,
968                canwu_time::SimTime::from_minutes(10),
969            )
970            .unwrap();
971        assert_eq!(booking.status, CapacityBookingStatus::Consumed);
972    }
973
974    #[test]
975    fn completion_operation_key_is_stable_and_revision_scoped() {
976        let first =
977            delivery_completion_operation_key(TransportExecutionId(4), ItineraryRevisionId(2), 7);
978        let second =
979            delivery_completion_operation_key(TransportExecutionId(4), ItineraryRevisionId(3), 7);
980        assert_ne!(first, second);
981        assert_eq!(
982            first,
983            delivery_completion_operation_key(TransportExecutionId(4), ItineraryRevisionId(2), 7)
984        );
985    }
986
987    #[test]
988    fn completion_requires_saga_and_exposes_stable_bridge_request() {
989        let attempt = DomainRecordVersionRef {
990            record: canwu_core::DomainRecordRef::new(
991                "fixture.information",
992                "delivery_attempt",
993                "delivery",
994            ),
995            version: 3,
996            established_by: canwu_core::DomainRecordVersionSource::InitialScenario,
997        };
998        let mut execution = TransportExecution::new(TransportExecutionId(9), Some(attempt.clone()));
999        let revision = ItineraryRevision {
1000            id: ItineraryRevisionId(1),
1001            predecessor: None,
1002            plan: plan(),
1003            planned_at: SimTime::EPOCH,
1004            valid_from: SimTime::EPOCH,
1005            reason: ItineraryRevisionReason::Initial,
1006            superseded_at: None,
1007            evidence: Vec::new(),
1008        };
1009        execution.install_initial_itinerary(revision).unwrap();
1010        assert_eq!(
1011            execution.complete_current_leg(SimTime::from_minutes(10), "b".to_owned()),
1012            Err(TransportError::MissingSaga)
1013        );
1014        execution
1015            .begin_saga(
1016                attempt,
1017                delivery_completion_operation_key(
1018                    TransportExecutionId(9),
1019                    ItineraryRevisionId(1),
1020                    3,
1021                ),
1022            )
1023            .unwrap();
1024        execution.start_current_leg(SimTime::EPOCH).unwrap();
1025        assert!(
1026            execution
1027                .complete_current_leg(SimTime::from_minutes(10), "b".to_owned())
1028                .unwrap()
1029        );
1030        let request = execution.completion_request().unwrap();
1031        assert_eq!(request.execution, TransportExecutionId(9));
1032        assert_eq!(request.itinerary_revision, ItineraryRevisionId(1));
1033        assert_eq!(request.delivery_attempt.version, 3);
1034    }
1035
1036    #[test]
1037    fn handoff_requires_arrival_and_is_idempotency_safe() {
1038        let attempt = DomainRecordVersionRef {
1039            record: canwu_core::DomainRecordRef::new(
1040                "fixture.information",
1041                "delivery_attempt",
1042                "delivery",
1043            ),
1044            version: 1,
1045            established_by: canwu_core::DomainRecordVersionSource::InitialScenario,
1046        };
1047        let mut execution =
1048            TransportExecution::new(TransportExecutionId(10), Some(attempt.clone()));
1049        let mut route = plan();
1050        route.legs.push(RouteLeg {
1051            connection: RoutingConnectionRef::new("bc"),
1052            from: RoutingNodeRef::new("b"),
1053            to: RoutingNodeRef::new("c"),
1054            mode: TransferMode::Rail,
1055            planned_departure_at: SimTime::from_minutes(10),
1056            planned_arrival_at: SimTime::from_minutes(20),
1057        });
1058        route.destination = RoutingNodeRef::new("c");
1059        route.estimated_arrival_at = SimTime::from_minutes(20);
1060        let revision = ItineraryRevision {
1061            id: ItineraryRevisionId(1),
1062            predecessor: None,
1063            plan: route,
1064            planned_at: SimTime::EPOCH,
1065            valid_from: SimTime::EPOCH,
1066            reason: ItineraryRevisionReason::Initial,
1067            superseded_at: None,
1068            evidence: Vec::new(),
1069        };
1070        execution.install_initial_itinerary(revision).unwrap();
1071        execution.start_current_leg(SimTime::EPOCH).unwrap();
1072        execution
1073            .complete_current_leg(SimTime::from_minutes(10), "b".to_owned())
1074            .unwrap();
1075        let handoff = Handoff {
1076            id: HandoffId(1),
1077            from_leg: LegExecutionId(1),
1078            to_leg: LegExecutionId(2),
1079            from_custodian: "courier/wuxi".to_owned(),
1080            to_custodian: "rail/beijing".to_owned(),
1081            at: SimTime::from_minutes(10),
1082            location: "b".to_owned(),
1083            evidence: Vec::new(),
1084        };
1085        execution.record_handoff(handoff.clone()).unwrap();
1086        assert_eq!(
1087            execution.record_handoff(handoff),
1088            Err(TransportError::InvalidHandoff(
1089                "handoff identity is already recorded".to_owned()
1090            ))
1091        );
1092    }
1093}