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