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