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