1#![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#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
38#[serde(transparent)]
39pub struct MovementOrderId(pub u64);
40
41#[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#[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 PersonsGroup,
66}
67
68impl MovementSubjectRole {
69 #[must_use]
72 pub const fn requires_quantity(self) -> bool {
73 matches!(self, Self::Cargo | Self::PersonsGroup)
74 }
75}
76
77#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
79pub struct MovementSubject {
80 pub entity: EntityRef,
81 pub role: MovementSubjectRole,
82 #[serde(default, skip_serializing_if = "Option::is_none")]
86 pub quantity: Option<u64>,
87 #[serde(default, skip_serializing_if = "Option::is_none")]
89 pub expected_custody: Option<EntityRef>,
90}
91
92#[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 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 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#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
295#[serde(rename_all = "snake_case")]
296pub enum HandoffKind {
297 #[default]
299 Planned,
300 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 #[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#[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}