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