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