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#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
374pub struct DeliverySaga {
375    pub operation_key: String,
376    pub delivery_attempt: DomainRecordVersionRef,
377    pub state: SagaState,
378    pub step: u32,
379    pub expected_attempt_version: u64,
380    pub last_error: Option<String>,
381    pub evidence: Vec<EvidenceRef>,
382}
383
384#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
385pub struct TransportExecution {
386    pub id: TransportExecutionId,
387    pub delivery_attempt: Option<DomainRecordVersionRef>,
388    pub state: TransportExecutionState,
389    pub active_itinerary_revision: Option<ItineraryRevisionId>,
390    pub current_leg_index: usize,
391    pub estimated_arrival_at: Option<SimTime>,
392    pub current_endpoint: Option<String>,
393    pub revisions: Vec<ItineraryRevision>,
394    pub legs: Vec<LegExecution>,
395    pub handoffs: Vec<Handoff>,
396    pub bookings: Vec<CapacityBooking>,
397    pub saga: Option<DeliverySaga>,
398}
399
400impl TransportExecution {
401    #[must_use]
402    pub fn new(id: TransportExecutionId, delivery_attempt: Option<DomainRecordVersionRef>) -> Self {
403        Self {
404            id,
405            delivery_attempt,
406            state: TransportExecutionState::Prepared,
407            active_itinerary_revision: None,
408            current_leg_index: 0,
409            estimated_arrival_at: None,
410            current_endpoint: None,
411            revisions: Vec::new(),
412            legs: Vec::new(),
413            handoffs: Vec::new(),
414            bookings: Vec::new(),
415            saga: None,
416        }
417    }
418
419    pub fn install_initial_itinerary(
420        &mut self,
421        revision: ItineraryRevision,
422    ) -> Result<(), TransportError> {
423        if !self.revisions.is_empty() || revision.predecessor.is_some() {
424            return Err(TransportError::InvalidRevision(
425                "initial itinerary must be the first revision".to_owned(),
426            ));
427        }
428        self.estimated_arrival_at = Some(revision.plan.estimated_arrival_at);
429        self.active_itinerary_revision = Some(revision.id);
430        self.current_endpoint = Some(revision.plan.origin.as_str().to_owned());
431        self.legs = revision
432            .plan
433            .legs
434            .iter()
435            .enumerate()
436            .map(|(index, _)| LegExecution {
437                id: LegExecutionId(index as u64 + 1),
438                itinerary_revision: revision.id,
439                leg_index: index,
440                status: LegExecutionStatus::Planned,
441                actual_departure_at: None,
442                actual_arrival_at: None,
443                failed_at: None,
444                failure_reason: None,
445                evidence: Vec::new(),
446            })
447            .collect();
448        self.revisions.push(revision);
449        self.state = TransportExecutionState::Planning;
450        Ok(())
451    }
452
453    pub fn reroute(
454        &mut self,
455        revision: ItineraryRevision,
456        at: SimTime,
457    ) -> Result<(), TransportError> {
458        let active = self
459            .active_itinerary_revision
460            .ok_or(TransportError::MissingItinerary)?;
461        if revision.predecessor != Some(active) || revision.valid_from < at {
462            return Err(TransportError::InvalidRevision("reroute must reference the active revision and start no earlier than the reroute time".to_owned()));
463        }
464        if self.revisions.iter().any(|item| item.id == revision.id) {
465            return Err(TransportError::InvalidRevision(
466                "itinerary revision identity must be unique within an execution".to_owned(),
467            ));
468        }
469        if let Some(previous) = self.revisions.iter_mut().find(|item| item.id == active) {
470            previous.superseded_at = Some(at);
471        }
472        self.estimated_arrival_at = Some(revision.plan.estimated_arrival_at);
473        self.active_itinerary_revision = Some(revision.id);
474        self.current_leg_index = 0;
475        let next_leg_id = self
476            .legs
477            .iter()
478            .map(|leg| leg.id.0)
479            .max()
480            .unwrap_or_default()
481            .checked_add(1)
482            .ok_or(TransportError::Overflow)?;
483        let mut legs = Vec::with_capacity(revision.plan.legs.len());
484        for (index, _) in revision.plan.legs.iter().enumerate() {
485            legs.push(LegExecution {
486                id: LegExecutionId(
487                    next_leg_id
488                        .checked_add(index as u64)
489                        .ok_or(TransportError::Overflow)?,
490                ),
491                itinerary_revision: revision.id,
492                leg_index: index,
493                status: LegExecutionStatus::Planned,
494                actual_departure_at: None,
495                actual_arrival_at: None,
496                failed_at: None,
497                failure_reason: None,
498                evidence: Vec::new(),
499            });
500        }
501        self.legs.extend(legs);
502        self.revisions.push(revision);
503        if let Some(saga) = self.saga.as_mut() {
504            saga.operation_key = delivery_completion_operation_key(
505                self.id,
506                self.active_itinerary_revision
507                    .ok_or(TransportError::MissingItinerary)?,
508                saga.expected_attempt_version,
509            );
510            saga.evidence.extend(
511                self.revisions
512                    .last()
513                    .map(|current| current.evidence.clone())
514                    .unwrap_or_default(),
515            );
516        }
517        self.state = TransportExecutionState::Planning;
518        Ok(())
519    }
520
521    pub fn begin_saga(
522        &mut self,
523        delivery_attempt: DomainRecordVersionRef,
524        operation_key: String,
525    ) -> Result<(), TransportError> {
526        if self.saga.is_some() {
527            return Err(TransportError::SagaAlreadyExists);
528        }
529        self.saga = Some(DeliverySaga {
530            expected_attempt_version: delivery_attempt.version,
531            operation_key,
532            delivery_attempt,
533            state: SagaState::TransportIntent,
534            step: 0,
535            last_error: None,
536            evidence: Vec::new(),
537        });
538        self.state = TransportExecutionState::Executing;
539        Ok(())
540    }
541
542    pub fn start_current_leg(&mut self, at: SimTime) -> Result<(), TransportError> {
543        if self.state != TransportExecutionState::Ready
544            && self.state != TransportExecutionState::Executing
545            && self.state != TransportExecutionState::Planning
546        {
547            return Err(TransportError::InvalidState(
548                "transport execution cannot start a leg in its current state".to_owned(),
549            ));
550        }
551        let active = self
552            .active_itinerary_revision
553            .ok_or(TransportError::MissingItinerary)?;
554        let leg = self
555            .legs
556            .iter_mut()
557            .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
558            .ok_or(TransportError::MissingLeg)?;
559        if !matches!(
560            leg.status,
561            LegExecutionStatus::Planned | LegExecutionStatus::Booked | LegExecutionStatus::Waiting
562        ) {
563            return Err(TransportError::InvalidState(
564                "current leg is not startable".to_owned(),
565            ));
566        }
567        leg.status = LegExecutionStatus::Departed;
568        leg.actual_departure_at = Some(at);
569        self.state = TransportExecutionState::Executing;
570        Ok(())
571    }
572
573    pub fn complete_current_leg(
574        &mut self,
575        at: SimTime,
576        endpoint: String,
577    ) -> Result<bool, TransportError> {
578        let active = self
579            .active_itinerary_revision
580            .ok_or(TransportError::MissingItinerary)?;
581        let active_leg_count = self
582            .legs
583            .iter()
584            .filter(|leg| leg.itinerary_revision == active)
585            .count();
586        let final_leg = self.current_leg_index.saturating_add(1) >= active_leg_count;
587        if final_leg && self.saga.is_none() {
588            return Err(TransportError::MissingSaga);
589        }
590        let leg = self
591            .legs
592            .iter_mut()
593            .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
594            .ok_or(TransportError::MissingLeg)?;
595        if leg.status != LegExecutionStatus::Departed {
596            return Err(TransportError::InvalidState(
597                "current leg must be departed before arrival".to_owned(),
598            ));
599        }
600        if leg
601            .actual_departure_at
602            .is_some_and(|departure| at < departure)
603        {
604            return Err(TransportError::InvalidState(
605                "arrival precedes departure".to_owned(),
606            ));
607        }
608        leg.status = LegExecutionStatus::Arrived;
609        leg.actual_arrival_at = Some(at);
610        self.current_endpoint = Some(endpoint);
611        self.current_leg_index = self.current_leg_index.saturating_add(1);
612        if self.current_leg_index >= active_leg_count {
613            self.state = TransportExecutionState::ArrivalPending;
614            self.mark_arrival_pending()?;
615            Ok(true)
616        } else {
617            self.state = TransportExecutionState::Ready;
618            Ok(false)
619        }
620    }
621
622    pub fn fail_current_leg(&mut self, reason: String, at: SimTime) -> Result<(), TransportError> {
623        let active = self
624            .active_itinerary_revision
625            .ok_or(TransportError::MissingItinerary)?;
626        let leg = self
627            .legs
628            .iter_mut()
629            .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
630            .ok_or(TransportError::MissingLeg)?;
631        leg.status = LegExecutionStatus::Failed;
632        leg.failed_at = Some(at);
633        leg.failure_reason = Some(reason);
634        self.state = TransportExecutionState::ReplanPending;
635        Ok(())
636    }
637
638    pub fn mark_arrival_pending(&mut self) -> Result<(), TransportError> {
639        let saga = self.saga.as_mut().ok_or(TransportError::MissingSaga)?;
640        saga.step = saga.step.checked_add(1).ok_or(TransportError::Overflow)?;
641        saga.state = SagaState::ArrivalPending;
642        self.state = TransportExecutionState::ArrivalPending;
643        Ok(())
644    }
645
646    pub fn record_handoff(&mut self, handoff: Handoff) -> Result<(), TransportError> {
647        if handoff.from_leg == handoff.to_leg
648            || handoff.from_custodian.trim().is_empty()
649            || handoff.to_custodian.trim().is_empty()
650            || handoff.location.trim().is_empty()
651        {
652            return Err(TransportError::InvalidHandoff(
653                "handoff requires distinct legs, custodians, and location".to_owned(),
654            ));
655        }
656        if self
657            .handoffs
658            .iter()
659            .any(|existing| existing.id == handoff.id)
660        {
661            return Err(TransportError::InvalidHandoff(
662                "handoff identity is already recorded".to_owned(),
663            ));
664        }
665        let from_leg = self
666            .legs
667            .iter()
668            .find(|leg| leg.id == handoff.from_leg)
669            .ok_or(TransportError::MissingLeg)?;
670        let to_leg = self
671            .legs
672            .iter()
673            .find(|leg| leg.id == handoff.to_leg)
674            .ok_or(TransportError::MissingLeg)?;
675        if !matches!(
676            from_leg.status,
677            LegExecutionStatus::Arrived | LegExecutionStatus::Failed
678        ) || !matches!(
679            to_leg.status,
680            LegExecutionStatus::Planned | LegExecutionStatus::Booked | LegExecutionStatus::Waiting
681        ) || from_leg
682            .actual_arrival_at
683            .is_some_and(|arrived_at| handoff.at < arrived_at)
684            || from_leg
685                .failed_at
686                .is_some_and(|failed_at| handoff.at < failed_at)
687        {
688            return Err(TransportError::InvalidHandoff(
689                "handoff must follow an arrived leg and precede the next leg".to_owned(),
690            ));
691        }
692        self.handoffs.push(handoff);
693        Ok(())
694    }
695
696    pub fn completion_request(&self) -> Result<DeliveryCompletionRequest, TransportError> {
697        let saga = self.saga.as_ref().ok_or(TransportError::MissingSaga)?;
698        if self.state != TransportExecutionState::ArrivalPending
699            || saga.state != SagaState::ArrivalPending
700        {
701            return Err(TransportError::InvalidState(
702                "delivery completion requires an arrival-pending execution".to_owned(),
703            ));
704        }
705        let revision = self
706            .active_itinerary_revision
707            .ok_or(TransportError::MissingItinerary)?;
708        let completed_at = self
709            .legs
710            .iter()
711            .rev()
712            .find(|leg| leg.itinerary_revision == revision)
713            .and_then(|leg| leg.actual_arrival_at)
714            .ok_or(TransportError::MissingLeg)?;
715        let attempt = self
716            .delivery_attempt
717            .clone()
718            .ok_or(TransportError::MissingDeliveryAttempt)?;
719        if attempt != saga.delivery_attempt {
720            return Err(TransportError::InvalidState(
721                "saga delivery attempt does not match execution".to_owned(),
722            ));
723        }
724        Ok(DeliveryCompletionRequest {
725            operation_key: saga.operation_key.clone(),
726            execution: self.id,
727            itinerary_revision: revision,
728            delivery_attempt: attempt,
729            completed_at,
730            evidence: saga.evidence.clone(),
731        })
732    }
733
734    pub fn reconcile_information(
735        &mut self,
736        success: bool,
737        error: Option<String>,
738    ) -> Result<(), TransportError> {
739        let saga = self.saga.as_mut().ok_or(TransportError::MissingSaga)?;
740        saga.step = saga.step.checked_add(1).ok_or(TransportError::Overflow)?;
741        if success {
742            saga.state = SagaState::Settled;
743            saga.last_error = None;
744            self.state = TransportExecutionState::Settled;
745        } else {
746            saga.state = SagaState::CompensationPending;
747            saga.last_error = error;
748            self.state = TransportExecutionState::Failed;
749        }
750        Ok(())
751    }
752}
753
754#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
755pub enum TransportError {
756    InvalidRevision(String),
757    MissingItinerary,
758    MissingSaga,
759    SagaAlreadyExists,
760    Overflow,
761    MissingLeg,
762    InvalidState(String),
763    InvalidBooking(String),
764    InvalidHandoff(String),
765    MissingDeliveryAttempt,
766}
767
768#[cfg(test)]
769mod tests {
770    use super::*;
771    use canwu_core::{PersonId, TerritoryId};
772    use canwu_routing::{
773        ROUTING_ALGORITHM_VERSION, RouteCost, RouteLeg, RoutingConnectionRef, RoutingNodeRef,
774        TransferMode,
775    };
776
777    fn plan() -> RoutePlan {
778        RoutePlan {
779            algorithm_version: ROUTING_ALGORITHM_VERSION.to_owned(),
780            policy_version: "policy.v1".to_owned(),
781            planning_snapshot_digest: "snapshot".to_owned(),
782            origin: RoutingNodeRef::new("a"),
783            destination: RoutingNodeRef::new("b"),
784            departure_at: canwu_time::SimTime::EPOCH,
785            estimated_arrival_at: canwu_time::SimTime::from_minutes(10),
786            cost: RouteCost {
787                estimated_arrival_at: canwu_time::SimTime::from_minutes(10),
788                risk_per_mille: 0,
789                resource_cost: 0,
790                transfers: 1,
791            },
792            legs: vec![RouteLeg {
793                connection: RoutingConnectionRef::new("ab"),
794                from: RoutingNodeRef::new("a"),
795                to: RoutingNodeRef::new("b"),
796                mode: TransferMode::Horse,
797                planned_departure_at: canwu_time::SimTime::EPOCH,
798                planned_arrival_at: canwu_time::SimTime::from_minutes(10),
799            }],
800            digest: "route".to_owned(),
801        }
802    }
803
804    #[test]
805    fn movement_order_accepts_a_self_directed_person_subject() {
806        let order = MovementOrder {
807            id: MovementOrderId(1),
808            subjects: vec![MovementSubject {
809                entity: EntityRef::Person(PersonId::new(7)),
810                role: MovementSubjectRole::MovablePrincipal,
811                quantity: None,
812                expected_custody: None,
813            }],
814            origin: RoutingNodeRef::new("a"),
815            destination: RoutingNodeRef::new("b"),
816            plan: plan(),
817            initiative: MovementInitiative::SelfDirected,
818            ordered_at: SimTime::EPOCH,
819            expected_position_revision: 1,
820        };
821        order.validate().unwrap();
822    }
823
824    #[test]
825    fn movement_order_rejects_duplicate_subjects_and_missing_cargo_quantity() {
826        let mut order = MovementOrder {
827            id: MovementOrderId(2),
828            subjects: vec![
829                MovementSubject {
830                    entity: EntityRef::Person(PersonId::new(7)),
831                    role: MovementSubjectRole::MovablePrincipal,
832                    quantity: None,
833                    expected_custody: None,
834                },
835                MovementSubject {
836                    entity: EntityRef::Person(PersonId::new(7)),
837                    role: MovementSubjectRole::Cargo,
838                    quantity: None,
839                    expected_custody: None,
840                },
841            ],
842            origin: RoutingNodeRef::new("a"),
843            destination: RoutingNodeRef::new("b"),
844            plan: plan(),
845            initiative: MovementInitiative::Delegated,
846            ordered_at: SimTime::EPOCH,
847            expected_position_revision: 1,
848        };
849        assert!(matches!(
850            order.validate(),
851            Err(MovementOrderError::Invalid(_))
852        ));
853        order.subjects[1].entity = EntityRef::Territory(TerritoryId::new(8));
854        assert!(matches!(
855            order.validate(),
856            Err(MovementOrderError::Invalid(_))
857        ));
858    }
859
860    #[test]
861    fn movement_order_rejects_a_plan_that_does_not_match_the_order() {
862        let mut order = MovementOrder {
863            id: MovementOrderId(3),
864            subjects: vec![MovementSubject {
865                entity: EntityRef::Person(PersonId::new(9)),
866                role: MovementSubjectRole::MovablePrincipal,
867                quantity: None,
868                expected_custody: None,
869            }],
870            origin: RoutingNodeRef::new("a"),
871            destination: RoutingNodeRef::new("c"),
872            plan: plan(),
873            initiative: MovementInitiative::SelfDirected,
874            ordered_at: SimTime::EPOCH,
875            expected_position_revision: 1,
876        };
877        assert!(matches!(
878            order.validate(),
879            Err(MovementOrderError::Invalid(_))
880        ));
881        order.destination = RoutingNodeRef::new("b");
882        order.plan.digest = String::new();
883        assert!(matches!(
884            order.validate(),
885            Err(MovementOrderError::Invalid(_))
886        ));
887    }
888
889    #[test]
890    fn reroute_supersedes_without_creating_a_new_delivery_attempt() {
891        let mut execution = TransportExecution::new(TransportExecutionId(1), None);
892        execution
893            .install_initial_itinerary(ItineraryRevision {
894                id: ItineraryRevisionId(1),
895                predecessor: None,
896                plan: plan(),
897                planned_at: canwu_time::SimTime::EPOCH,
898                valid_from: canwu_time::SimTime::EPOCH,
899                reason: ItineraryRevisionReason::Initial,
900                superseded_at: None,
901                evidence: Vec::new(),
902            })
903            .unwrap();
904        let mut replacement = plan();
905        replacement.destination = RoutingNodeRef::new("b");
906        execution
907            .reroute(
908                ItineraryRevision {
909                    id: ItineraryRevisionId(2),
910                    predecessor: Some(ItineraryRevisionId(1)),
911                    plan: replacement,
912                    planned_at: canwu_time::SimTime::from_minutes(1),
913                    valid_from: canwu_time::SimTime::from_minutes(1),
914                    reason: ItineraryRevisionReason::Disaster {
915                        explanation: "bridge closed".to_owned(),
916                    },
917                    superseded_at: None,
918                    evidence: Vec::new(),
919                },
920                canwu_time::SimTime::from_minutes(1),
921            )
922            .unwrap();
923        assert_eq!(execution.revisions.len(), 2);
924        assert_eq!(
925            execution.active_itinerary_revision,
926            Some(ItineraryRevisionId(2))
927        );
928        assert_eq!(execution.state, TransportExecutionState::Planning);
929    }
930
931    #[test]
932    fn capacity_booking_is_a_persisted_windowed_state_machine() {
933        let mut booking = CapacityBooking::new(
934            CapacityBookingId(1),
935            TransportExecutionId(1),
936            "relay-horse:wu-xi:01".to_owned(),
937            canwu_time::SimTime::EPOCH,
938            canwu_time::SimTime::from_minutes(60),
939            1,
940            10,
941        )
942        .unwrap();
943        booking
944            .transition(CapacityBookingStatus::Confirmed, canwu_time::SimTime::EPOCH)
945            .unwrap();
946        booking
947            .transition(
948                CapacityBookingStatus::Consumed,
949                canwu_time::SimTime::from_minutes(10),
950            )
951            .unwrap();
952        assert_eq!(booking.status, CapacityBookingStatus::Consumed);
953    }
954
955    #[test]
956    fn completion_operation_key_is_stable_and_revision_scoped() {
957        let first =
958            delivery_completion_operation_key(TransportExecutionId(4), ItineraryRevisionId(2), 7);
959        let second =
960            delivery_completion_operation_key(TransportExecutionId(4), ItineraryRevisionId(3), 7);
961        assert_ne!(first, second);
962        assert_eq!(
963            first,
964            delivery_completion_operation_key(TransportExecutionId(4), ItineraryRevisionId(2), 7)
965        );
966    }
967
968    #[test]
969    fn completion_requires_saga_and_exposes_stable_bridge_request() {
970        let attempt = DomainRecordVersionRef {
971            record: canwu_core::DomainRecordRef::new(
972                "fixture.information",
973                "delivery_attempt",
974                "delivery",
975            ),
976            version: 3,
977            established_by: canwu_core::DomainRecordVersionSource::InitialScenario,
978        };
979        let mut execution = TransportExecution::new(TransportExecutionId(9), Some(attempt.clone()));
980        let revision = ItineraryRevision {
981            id: ItineraryRevisionId(1),
982            predecessor: None,
983            plan: plan(),
984            planned_at: SimTime::EPOCH,
985            valid_from: SimTime::EPOCH,
986            reason: ItineraryRevisionReason::Initial,
987            superseded_at: None,
988            evidence: Vec::new(),
989        };
990        execution.install_initial_itinerary(revision).unwrap();
991        assert_eq!(
992            execution.complete_current_leg(SimTime::from_minutes(10), "b".to_owned()),
993            Err(TransportError::MissingSaga)
994        );
995        execution
996            .begin_saga(
997                attempt,
998                delivery_completion_operation_key(
999                    TransportExecutionId(9),
1000                    ItineraryRevisionId(1),
1001                    3,
1002                ),
1003            )
1004            .unwrap();
1005        execution.start_current_leg(SimTime::EPOCH).unwrap();
1006        assert!(
1007            execution
1008                .complete_current_leg(SimTime::from_minutes(10), "b".to_owned())
1009                .unwrap()
1010        );
1011        let request = execution.completion_request().unwrap();
1012        assert_eq!(request.execution, TransportExecutionId(9));
1013        assert_eq!(request.itinerary_revision, ItineraryRevisionId(1));
1014        assert_eq!(request.delivery_attempt.version, 3);
1015    }
1016
1017    #[test]
1018    fn handoff_requires_arrival_and_is_idempotency_safe() {
1019        let attempt = DomainRecordVersionRef {
1020            record: canwu_core::DomainRecordRef::new(
1021                "fixture.information",
1022                "delivery_attempt",
1023                "delivery",
1024            ),
1025            version: 1,
1026            established_by: canwu_core::DomainRecordVersionSource::InitialScenario,
1027        };
1028        let mut execution =
1029            TransportExecution::new(TransportExecutionId(10), Some(attempt.clone()));
1030        let mut route = plan();
1031        route.legs.push(RouteLeg {
1032            connection: RoutingConnectionRef::new("bc"),
1033            from: RoutingNodeRef::new("b"),
1034            to: RoutingNodeRef::new("c"),
1035            mode: TransferMode::Rail,
1036            planned_departure_at: SimTime::from_minutes(10),
1037            planned_arrival_at: SimTime::from_minutes(20),
1038        });
1039        route.destination = RoutingNodeRef::new("c");
1040        route.estimated_arrival_at = SimTime::from_minutes(20);
1041        let revision = ItineraryRevision {
1042            id: ItineraryRevisionId(1),
1043            predecessor: None,
1044            plan: route,
1045            planned_at: SimTime::EPOCH,
1046            valid_from: SimTime::EPOCH,
1047            reason: ItineraryRevisionReason::Initial,
1048            superseded_at: None,
1049            evidence: Vec::new(),
1050        };
1051        execution.install_initial_itinerary(revision).unwrap();
1052        execution.start_current_leg(SimTime::EPOCH).unwrap();
1053        execution
1054            .complete_current_leg(SimTime::from_minutes(10), "b".to_owned())
1055            .unwrap();
1056        let handoff = Handoff {
1057            id: HandoffId(1),
1058            from_leg: LegExecutionId(1),
1059            to_leg: LegExecutionId(2),
1060            from_custodian: "courier/wuxi".to_owned(),
1061            to_custodian: "rail/beijing".to_owned(),
1062            at: SimTime::from_minutes(10),
1063            location: "b".to_owned(),
1064            evidence: Vec::new(),
1065        };
1066        execution.record_handoff(handoff.clone()).unwrap();
1067        assert_eq!(
1068            execution.record_handoff(handoff),
1069            Err(TransportError::InvalidHandoff(
1070                "handoff identity is already recorded".to_owned()
1071            ))
1072        );
1073    }
1074}