1#![allow(clippy::missing_errors_doc)]
8
9use canwu_core::{
10 DomainRecordRef, DomainRecordVersionRef, EntityRef, EvidenceRef, KnowledgeHolderRef,
11};
12use canwu_routing::{RoutePlan, RoutingNodeRef};
13use canwu_time::SimTime;
14use serde::{Deserialize, Serialize};
15use std::cmp::Ordering;
16
17pub const TRANSPORT_SEMANTIC_VERSION: &str = "canwu-transport.v5";
24
25pub const CAPACITY_BOOKING_ALLOCATION_DIGEST_DOMAIN: &str =
27 "canwu.transport.capacity-booking-allocation.v1";
28
29#[must_use]
30pub fn delivery_completion_operation_key(
31 execution: TransportExecutionId,
32 revision: ItineraryRevisionId,
33 attempt_version: u64,
34) -> String {
35 format!(
36 "transport/{}/revision/{}/delivery-completion/attempt-version/{}",
37 execution.0, revision.0, attempt_version
38 )
39}
40
41#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
42#[serde(transparent)]
43pub struct TransportExecutionId(pub u64);
44
45#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
51#[serde(transparent)]
52pub struct MovementOrderId(pub u64);
53
54#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
57#[serde(rename_all = "snake_case")]
58pub enum MovementInitiative {
59 SelfDirected,
60 Commanded,
61 Delegated,
62 Forced,
63 Automatic,
64}
65
66#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
68#[serde(rename_all = "snake_case")]
69pub enum MovementSubjectRole {
70 MovablePrincipal,
71 Cargo,
72 Carrier,
73 Passenger,
74 Attached,
75 PersonsGroup,
79}
80
81impl MovementSubjectRole {
82 #[must_use]
85 pub const fn requires_quantity(self) -> bool {
86 matches!(self, Self::Cargo | Self::PersonsGroup)
87 }
88}
89
90#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
92pub struct MovementSubject {
93 pub entity: EntityRef,
94 pub role: MovementSubjectRole,
95 #[serde(default, skip_serializing_if = "Option::is_none")]
99 pub quantity: Option<u64>,
100 #[serde(default, skip_serializing_if = "Option::is_none")]
102 pub expected_custody: Option<EntityRef>,
103}
104
105#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
111pub struct MovementOrder {
112 pub id: MovementOrderId,
113 pub subjects: Vec<MovementSubject>,
114 pub origin: RoutingNodeRef,
115 pub destination: RoutingNodeRef,
116 pub plan: RoutePlan,
117 pub initiative: MovementInitiative,
118 pub ordered_at: SimTime,
119 pub expected_position_revision: u64,
120}
121
122#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
123pub enum MovementOrderError {
124 Invalid(String),
125}
126
127impl MovementOrder {
128 pub fn validate(&self) -> Result<(), MovementOrderError> {
134 if self.id.0 == 0
135 || self.expected_position_revision == 0
136 || self.origin.as_str().trim().is_empty()
137 || self.destination.as_str().trim().is_empty()
138 {
139 return Err(MovementOrderError::Invalid(
140 "movement order identity, endpoints, and expected position revision must be valid"
141 .to_owned(),
142 ));
143 }
144 if self.subjects.is_empty()
145 || self
146 .subjects
147 .windows(2)
148 .any(|pair| pair[0].entity >= pair[1].entity)
149 {
150 return Err(MovementOrderError::Invalid(
151 "movement subjects must be non-empty, sorted, and unique by entity".to_owned(),
152 ));
153 }
154 for subject in &self.subjects {
155 if subject.quantity.is_some_and(|quantity| quantity == 0)
156 || subject.role.requires_quantity() != subject.quantity.is_some()
157 {
158 return Err(MovementOrderError::Invalid(
159 "cargo and persons-group subjects require a positive quantity; other subjects cannot carry one"
160 .to_owned(),
161 ));
162 }
163 }
164 if self.plan.origin != self.origin
165 || self.plan.destination != self.destination
166 || self.plan.departure_at < self.ordered_at
167 || self.plan.estimated_arrival_at < self.plan.departure_at
168 || self.plan.digest.trim().is_empty()
169 {
170 return Err(MovementOrderError::Invalid(
171 "movement plan does not match the order endpoints or time range".to_owned(),
172 ));
173 }
174 Ok(())
175 }
176}
177
178#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
179#[serde(transparent)]
180pub struct ItineraryRevisionId(pub u64);
181
182#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
183#[serde(transparent)]
184pub struct LegExecutionId(pub u64);
185
186#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
187#[serde(transparent)]
188pub struct HandoffId(pub u64);
189
190#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
191#[serde(transparent)]
192pub struct CapacityBookingId(pub u64);
193
194#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
195#[serde(rename_all = "snake_case")]
196pub enum TransportExecutionState {
197 Prepared,
198 Planning,
199 Booking,
200 Ready,
201 Executing,
202 ReplanPending,
203 ArrivalPending,
204 Settled,
205 Failed,
206 Cancelled,
207}
208
209impl TransportExecutionState {
210 #[must_use]
211 pub const fn is_terminal(self) -> bool {
212 matches!(self, Self::Settled | Self::Failed | Self::Cancelled)
213 }
214}
215
216#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
217#[serde(rename_all = "snake_case")]
218pub enum ItineraryRevisionReason {
219 Initial,
220 Disaster {
221 explanation: String,
222 },
223 CapacityUnavailable {
224 explanation: String,
225 },
226 KnowledgeUpdate {
227 explanation: String,
228 },
229 Recovery {
230 explanation: String,
231 },
232 ExternalCondition {
236 record: DomainRecordRef,
237 version: u64,
238 kind: String,
239 },
240}
241
242impl ItineraryRevisionReason {
243 fn validate(&self) -> Result<(), TransportError> {
244 if let Self::ExternalCondition {
245 record,
246 version,
247 kind,
248 } = self
249 && (*version == 0
250 || kind.trim().is_empty()
251 || !domain_record_ref_is_well_formed(record))
252 {
253 return Err(TransportError::InvalidRevision(
254 "external-condition reason requires a well-formed record, a positive version, and a kind"
255 .to_owned(),
256 ));
257 }
258 Ok(())
259 }
260}
261
262fn domain_record_ref_is_well_formed(record: &DomainRecordRef) -> bool {
263 !record.kind.namespace.trim().is_empty()
264 && !record.kind.name.trim().is_empty()
265 && !record.id.trim().is_empty()
266}
267
268#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
269pub struct ItineraryRevision {
270 pub id: ItineraryRevisionId,
271 pub predecessor: Option<ItineraryRevisionId>,
272 pub plan: RoutePlan,
273 pub planned_at: SimTime,
274 pub valid_from: SimTime,
275 pub reason: ItineraryRevisionReason,
276 pub superseded_at: Option<SimTime>,
277 pub evidence: Vec<EvidenceRef>,
278}
279
280#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
281#[serde(rename_all = "snake_case")]
282pub enum LegExecutionStatus {
283 Planned,
284 Booked,
285 Loaded,
286 Departed,
287 Arrived,
288 Waiting,
289 Failed,
290 Cancelled,
291}
292
293#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
294pub struct LegExecution {
295 pub id: LegExecutionId,
296 pub itinerary_revision: ItineraryRevisionId,
297 pub leg_index: usize,
298 pub status: LegExecutionStatus,
299 pub actual_departure_at: Option<SimTime>,
300 pub actual_arrival_at: Option<SimTime>,
301 pub failed_at: Option<SimTime>,
302 pub failure_reason: Option<String>,
303 pub evidence: Vec<EvidenceRef>,
304}
305
306#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
308#[serde(rename_all = "snake_case")]
309pub enum HandoffKind {
310 #[default]
312 Planned,
313 Seizure { by: EntityRef },
330}
331
332impl HandoffKind {
333 #[must_use]
334 pub const fn is_planned(&self) -> bool {
335 matches!(self, Self::Planned)
336 }
337}
338
339#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
340pub struct Handoff {
341 pub id: HandoffId,
342 pub from_leg: LegExecutionId,
343 pub to_leg: LegExecutionId,
344 pub from_custodian: String,
345 pub to_custodian: String,
346 pub at: SimTime,
347 pub location: String,
348 pub evidence: Vec<EvidenceRef>,
349 #[serde(default, skip_serializing_if = "HandoffKind::is_planned")]
352 pub kind: HandoffKind,
353}
354
355#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
356pub struct DeliveryCompletionRequest {
357 pub operation_key: String,
358 pub execution: TransportExecutionId,
359 pub itinerary_revision: ItineraryRevisionId,
360 pub delivery_attempt: DomainRecordVersionRef,
361 pub completed_at: SimTime,
362 pub evidence: Vec<EvidenceRef>,
363}
364
365#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
366#[serde(rename_all = "snake_case")]
367pub enum CapacityBookingStatus {
368 Requested,
369 Confirmed,
370 Consumed,
371 Released,
372 Expired,
373 Cancelled,
374 Failed,
375}
376
377#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
378pub struct CapacityBooking {
379 pub id: CapacityBookingId,
380 pub execution: TransportExecutionId,
381 pub resource: String,
382 pub valid_from: SimTime,
383 pub valid_until: SimTime,
384 pub quantity: u64,
385 pub priority: i32,
386 pub status: CapacityBookingStatus,
387 pub allocation_evidence: Vec<EvidenceRef>,
388}
389
390impl CapacityBooking {
391 pub fn new(
392 id: CapacityBookingId,
393 execution: TransportExecutionId,
394 resource: String,
395 valid_from: SimTime,
396 valid_until: SimTime,
397 quantity: u64,
398 priority: i32,
399 ) -> Result<Self, TransportError> {
400 if valid_until < valid_from || quantity == 0 {
401 return Err(TransportError::InvalidBooking(
402 "capacity booking requires a positive quantity and non-inverted window".to_owned(),
403 ));
404 }
405 Ok(Self {
406 id,
407 execution,
408 resource,
409 valid_from,
410 valid_until,
411 quantity,
412 priority,
413 status: CapacityBookingStatus::Requested,
414 allocation_evidence: Vec::new(),
415 })
416 }
417
418 pub fn transition(
425 &mut self,
426 status: CapacityBookingStatus,
427 at: SimTime,
428 ) -> Result<(), TransportError> {
429 if status == CapacityBookingStatus::Consumed
430 && (at < self.valid_from || at > self.valid_until)
431 {
432 return Err(TransportError::InvalidBooking(
433 "booking capacity can only be consumed inside its validity window".to_owned(),
434 ));
435 }
436 let allowed = matches!(
437 (self.status, status),
438 (
439 CapacityBookingStatus::Requested,
440 CapacityBookingStatus::Confirmed
441 | CapacityBookingStatus::Failed
442 | CapacityBookingStatus::Cancelled
443 ) | (
444 CapacityBookingStatus::Confirmed,
445 CapacityBookingStatus::Consumed
446 | CapacityBookingStatus::Released
447 | CapacityBookingStatus::Cancelled
448 | CapacityBookingStatus::Expired
449 ) | (
450 CapacityBookingStatus::Consumed,
451 CapacityBookingStatus::Released
452 )
453 );
454 if !allowed {
455 return Err(TransportError::InvalidBooking(
456 "capacity booking transition is not allowed".to_owned(),
457 ));
458 }
459 if at > self.valid_until && status == CapacityBookingStatus::Confirmed {
460 return Err(TransportError::InvalidBooking(
461 "capacity booking cannot be confirmed after its validity window".to_owned(),
462 ));
463 }
464 if status == CapacityBookingStatus::Expired && at <= self.valid_until {
465 return Err(TransportError::InvalidBooking(
466 "capacity booking cannot expire before its validity window ends".to_owned(),
467 ));
468 }
469 self.status = status;
470 Ok(())
471 }
472}
473
474#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
483pub struct TransportCapacityPoolV1 {
484 pub id: String,
485 pub resource: String,
486 pub custodian: KnowledgeHolderRef,
487 pub window_from: SimTime,
488 pub window_until: SimTime,
489 pub quantity: u64,
490 pub booked: u64,
491 pub consumed: u64,
492 pub revision: u64,
493}
494
495impl TransportCapacityPoolV1 {
496 pub fn new(
498 id: String,
499 resource: String,
500 custodian: KnowledgeHolderRef,
501 window_from: SimTime,
502 window_until: SimTime,
503 quantity: u64,
504 ) -> Result<Self, TransportError> {
505 let pool = Self {
506 id,
507 resource,
508 custodian,
509 window_from,
510 window_until,
511 quantity,
512 booked: 0,
513 consumed: 0,
514 revision: 1,
515 };
516 pool.validate()?;
517 Ok(pool)
518 }
519
520 pub fn validate(&self) -> Result<(), TransportError> {
521 if !is_canonical_label(&self.id)
522 || !is_canonical_label(&self.resource)
523 || self.window_until < self.window_from
524 || self.revision == 0
525 {
526 return Err(TransportError::InvalidCapacityPool(
527 "capacity pool requires a canonical identity and resource, a non-inverted window, and a positive revision"
528 .to_owned(),
529 ));
530 }
531 if self
532 .booked
533 .checked_add(self.consumed)
534 .is_none_or(|held| held > self.quantity)
535 {
536 return Err(TransportError::InvalidCapacityPool(
537 "capacity pool holds more than it offers".to_owned(),
538 ));
539 }
540 Ok(())
541 }
542
543 #[must_use]
545 pub const fn available(&self) -> u64 {
546 self.quantity
547 .saturating_sub(self.booked)
548 .saturating_sub(self.consumed)
549 }
550
551 pub fn revise(
554 &mut self,
555 window_from: SimTime,
556 window_until: SimTime,
557 quantity: u64,
558 ) -> Result<(), TransportError> {
559 let mut revised = self.clone();
560 revised.window_from = window_from;
561 revised.window_until = window_until;
562 revised.quantity = quantity;
563 revised.revision = self
564 .revision
565 .checked_add(1)
566 .ok_or(TransportError::Overflow)?;
567 revised.validate()?;
568 *self = revised;
569 Ok(())
570 }
571
572 pub fn apply_allocations(
575 &mut self,
576 allocations: &[BookingAllocationV1],
577 ) -> Result<(), TransportError> {
578 if allocations.is_empty() {
579 return Ok(());
580 }
581 let mut bookings = allocations
582 .iter()
583 .map(|allocation| allocation.booking)
584 .collect::<Vec<_>>();
585 bookings.sort_unstable();
586 if bookings.windows(2).any(|pair| pair[0] == pair[1]) {
587 return Err(TransportError::InvalidCapacityPool(
588 "an allocation pass names each booking once".to_owned(),
589 ));
590 }
591 let mut confirmed = 0_u64;
592 for allocation in allocations {
593 let evidence = &allocation.evidence;
594 if evidence.pool != self.id
595 || evidence.pool_revision != self.revision
596 || evidence.booking != allocation.booking
597 || evidence.quantity != allocation.quantity
598 || evidence.status != allocation.status
599 || !evidence.digest_matches()
600 {
601 return Err(TransportError::InvalidCapacityPool(
602 "allocation evidence does not bind this pool revision".to_owned(),
603 ));
604 }
605 match allocation.status {
606 CapacityBookingStatus::Confirmed => {
607 confirmed = confirmed
608 .checked_add(allocation.quantity)
609 .ok_or(TransportError::Overflow)?;
610 }
611 CapacityBookingStatus::Failed if allocation.quantity == 0 => {}
612 _ => {
613 return Err(TransportError::InvalidCapacityPool(
614 "an allocation either confirms its quantity or fails with none".to_owned(),
615 ));
616 }
617 }
618 }
619 let mut next = self.clone();
620 next.booked = self
621 .booked
622 .checked_add(confirmed)
623 .ok_or(TransportError::Overflow)?;
624 next.revision = self
625 .revision
626 .checked_add(1)
627 .ok_or(TransportError::Overflow)?;
628 next.validate()?;
629 *self = next;
630 Ok(())
631 }
632
633 pub fn apply_booking_transition(
640 &mut self,
641 quantity: u64,
642 from: CapacityBookingStatus,
643 to: CapacityBookingStatus,
644 ) -> Result<(), TransportError> {
645 use CapacityBookingStatus::{
646 Cancelled, Confirmed, Consumed, Expired, Failed, Released, Requested,
647 };
648 let underflow =
649 || TransportError::InvalidCapacityPool("capacity pool counter underflow".to_owned());
650 let mut next = self.clone();
651 match (from, to) {
652 (Requested, Cancelled | Failed) => return Ok(()),
653 (Confirmed, Consumed) => {
654 next.booked = next.booked.checked_sub(quantity).ok_or_else(underflow)?;
655 next.consumed = next
656 .consumed
657 .checked_add(quantity)
658 .ok_or(TransportError::Overflow)?;
659 }
660 (Confirmed, Released | Cancelled | Expired) => {
661 next.booked = next.booked.checked_sub(quantity).ok_or_else(underflow)?;
662 }
663 (Consumed, Released) => {
664 next.consumed = next.consumed.checked_sub(quantity).ok_or_else(underflow)?;
665 }
666 _ => {
667 return Err(TransportError::InvalidCapacityPool(
668 "booking transition has no capacity-pool effect".to_owned(),
669 ));
670 }
671 }
672 next.revision = self
673 .revision
674 .checked_add(1)
675 .ok_or(TransportError::Overflow)?;
676 next.validate()?;
677 *self = next;
678 Ok(())
679 }
680}
681
682#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
685pub struct CapacityBookingRequestV1 {
686 pub booking: CapacityBooking,
687 pub tie_break: String,
688 pub admitted_sequence: u64,
689}
690
691#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
693#[serde(rename_all = "snake_case")]
694pub enum CapacityAllocationFailureV1 {
695 InsufficientCapacity,
697 OutsidePoolWindow,
699 WindowElapsed,
701}
702
703#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
705pub struct CapacityBookingAllocationEvidenceV1 {
706 pub pool: String,
707 pub pool_revision: u64,
708 pub booking: CapacityBookingId,
709 pub execution: TransportExecutionId,
710 pub requested: u64,
711 pub quantity: u64,
713 pub status: CapacityBookingStatus,
714 #[serde(default, skip_serializing_if = "Option::is_none")]
715 pub failure: Option<CapacityAllocationFailureV1>,
716 pub remaining_after: u64,
718 pub allocated_at: SimTime,
719 pub operation_key: String,
720 pub semantic_digest: String,
723}
724
725impl CapacityBookingAllocationEvidenceV1 {
726 #[must_use]
728 pub fn expected_digest(&self) -> String {
729 let material = (
730 &self.pool,
731 self.pool_revision,
732 self.booking,
733 self.execution,
734 self.requested,
735 self.quantity,
736 self.status,
737 self.failure,
738 self.remaining_after,
739 self.allocated_at,
740 &self.operation_key,
741 );
742 domain_digest(CAPACITY_BOOKING_ALLOCATION_DIGEST_DOMAIN, &material)
743 }
744
745 #[must_use]
746 pub fn digest_matches(&self) -> bool {
747 self.semantic_digest == self.expected_digest()
748 }
749}
750
751#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
753pub struct BookingAllocationV1 {
754 pub booking: CapacityBookingId,
755 pub quantity: u64,
757 pub status: CapacityBookingStatus,
759 pub evidence: CapacityBookingAllocationEvidenceV1,
760}
761
762#[must_use]
764pub fn capacity_booking_allocation_operation_key(
765 pool: &str,
766 pool_revision: u64,
767 booking: CapacityBookingId,
768) -> String {
769 format!(
770 "transport/pool/{pool}/revision/{pool_revision}/booking/{}/allocation",
771 booking.0
772 )
773}
774
775pub fn allocate_capacity_bookings(
788 pool: &TransportCapacityPoolV1,
789 requests: &[CapacityBookingRequestV1],
790 at: SimTime,
791) -> Result<Vec<BookingAllocationV1>, TransportError> {
792 pool.validate()?;
793 let mut ids = Vec::with_capacity(requests.len());
794 for request in requests {
795 let booking = &request.booking;
796 if booking.status != CapacityBookingStatus::Requested
797 || booking.resource != pool.resource
798 || booking.quantity == 0
799 || booking.valid_until < booking.valid_from
800 || !is_canonical_label(&request.tie_break)
801 {
802 return Err(TransportError::InvalidCapacityPool(
803 "allocation accepts only well-formed requested bookings for the pool resource"
804 .to_owned(),
805 ));
806 }
807 ids.push(booking.id);
808 }
809 ids.sort_unstable();
810 if ids.windows(2).any(|pair| pair[0] == pair[1]) {
811 return Err(TransportError::InvalidCapacityPool(
812 "allocation requests must have unique booking identities".to_owned(),
813 ));
814 }
815 let mut ordered = requests.iter().collect::<Vec<_>>();
816 ordered.sort_by(|left, right| allocation_order(left, right));
817 let mut remaining = pool.available();
818 let mut allocations = Vec::with_capacity(ordered.len());
819 for request in ordered {
820 let booking = &request.booking;
821 let failure = if at > booking.valid_until {
822 Some(CapacityAllocationFailureV1::WindowElapsed)
823 } else if booking.valid_from < pool.window_from || booking.valid_until > pool.window_until {
824 Some(CapacityAllocationFailureV1::OutsidePoolWindow)
825 } else if booking.quantity > remaining {
826 Some(CapacityAllocationFailureV1::InsufficientCapacity)
827 } else {
828 None
829 };
830 let (status, quantity) = if failure.is_none() {
831 remaining -= booking.quantity;
832 (CapacityBookingStatus::Confirmed, booking.quantity)
833 } else {
834 (CapacityBookingStatus::Failed, 0)
835 };
836 let mut evidence = CapacityBookingAllocationEvidenceV1 {
837 pool: pool.id.clone(),
838 pool_revision: pool.revision,
839 booking: booking.id,
840 execution: booking.execution,
841 requested: booking.quantity,
842 quantity,
843 status,
844 failure,
845 remaining_after: remaining,
846 allocated_at: at,
847 operation_key: capacity_booking_allocation_operation_key(
848 &pool.id,
849 pool.revision,
850 booking.id,
851 ),
852 semantic_digest: String::new(),
853 };
854 evidence.semantic_digest = evidence.expected_digest();
855 allocations.push(BookingAllocationV1 {
856 booking: booking.id,
857 quantity,
858 status,
859 evidence,
860 });
861 }
862 Ok(allocations)
863}
864
865fn allocation_order(left: &CapacityBookingRequestV1, right: &CapacityBookingRequestV1) -> Ordering {
866 right
867 .booking
868 .priority
869 .cmp(&left.booking.priority)
870 .then_with(|| left.booking.valid_from.cmp(&right.booking.valid_from))
871 .then_with(|| left.tie_break.cmp(&right.tie_break))
872 .then_with(|| left.admitted_sequence.cmp(&right.admitted_sequence))
873 .then_with(|| left.booking.id.cmp(&right.booking.id))
874}
875
876fn is_canonical_label(value: &str) -> bool {
877 !value.is_empty() && value.trim() == value
878}
879
880fn domain_digest<T: Serialize>(domain: &str, value: &T) -> String {
881 let encoded = serde_json::to_vec(value).expect("transport digest material must serialize");
882 let mut hasher = blake3::Hasher::new();
883 hasher.update(domain.as_bytes());
884 hasher.update(&[0]);
885 hasher.update(&encoded);
886 hasher.finalize().to_hex().to_string()
887}
888
889#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
890#[serde(rename_all = "snake_case")]
891pub enum SagaState {
892 TransportIntent,
893 WaitingForInformation,
894 Executing,
895 ArrivalPending,
896 Settled,
897 CompensationPending,
898 Failed,
899}
900
901#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
903pub enum ReconciliationOutcome {
904 Success,
905 Failure { error: String },
906}
907
908#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
909pub struct DeliverySaga {
910 pub operation_key: String,
911 pub delivery_attempt: DomainRecordVersionRef,
912 pub state: SagaState,
913 pub step: u32,
914 pub expected_attempt_version: u64,
915 pub last_error: Option<String>,
916 pub evidence: Vec<EvidenceRef>,
917}
918
919#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
920pub struct TransportExecution {
921 pub id: TransportExecutionId,
922 pub delivery_attempt: Option<DomainRecordVersionRef>,
923 pub state: TransportExecutionState,
924 pub active_itinerary_revision: Option<ItineraryRevisionId>,
925 pub current_leg_index: usize,
926 pub estimated_arrival_at: Option<SimTime>,
927 pub current_endpoint: Option<String>,
928 pub revisions: Vec<ItineraryRevision>,
929 pub legs: Vec<LegExecution>,
930 pub handoffs: Vec<Handoff>,
931 pub bookings: Vec<CapacityBooking>,
932 pub saga: Option<DeliverySaga>,
933}
934
935impl TransportExecution {
936 #[must_use]
937 pub fn new(id: TransportExecutionId, delivery_attempt: Option<DomainRecordVersionRef>) -> Self {
938 Self {
939 id,
940 delivery_attempt,
941 state: TransportExecutionState::Prepared,
942 active_itinerary_revision: None,
943 current_leg_index: 0,
944 estimated_arrival_at: None,
945 current_endpoint: None,
946 revisions: Vec::new(),
947 legs: Vec::new(),
948 handoffs: Vec::new(),
949 bookings: Vec::new(),
950 saga: None,
951 }
952 }
953
954 pub fn install_initial_itinerary(
955 &mut self,
956 revision: ItineraryRevision,
957 ) -> Result<(), TransportError> {
958 if !self.revisions.is_empty() || revision.predecessor.is_some() {
959 return Err(TransportError::InvalidRevision(
960 "initial itinerary must be the first revision".to_owned(),
961 ));
962 }
963 revision.reason.validate()?;
964 self.estimated_arrival_at = Some(revision.plan.estimated_arrival_at);
965 self.active_itinerary_revision = Some(revision.id);
966 self.current_endpoint = Some(revision.plan.origin.as_str().to_owned());
967 self.legs = revision
968 .plan
969 .legs
970 .iter()
971 .enumerate()
972 .map(|(index, _)| {
973 let id = u64::try_from(index)
974 .ok()
975 .and_then(|index| index.checked_add(1))
976 .ok_or(TransportError::Overflow)?;
977 Ok(LegExecution {
978 id: LegExecutionId(id),
979 itinerary_revision: revision.id,
980 leg_index: index,
981 status: LegExecutionStatus::Planned,
982 actual_departure_at: None,
983 actual_arrival_at: None,
984 failed_at: None,
985 failure_reason: None,
986 evidence: Vec::new(),
987 })
988 })
989 .collect::<Result<Vec<_>, TransportError>>()?;
990 self.revisions.push(revision);
991 self.state = TransportExecutionState::Planning;
992 Ok(())
993 }
994
995 pub fn reroute(
996 &mut self,
997 revision: ItineraryRevision,
998 at: SimTime,
999 ) -> Result<(), TransportError> {
1000 self.ensure_custody_in_itinerary()?;
1001 let active = self
1002 .active_itinerary_revision
1003 .ok_or(TransportError::MissingItinerary)?;
1004 if revision.predecessor != Some(active) || revision.valid_from < at {
1005 return Err(TransportError::InvalidRevision("reroute must reference the active revision and start no earlier than the reroute time".to_owned()));
1006 }
1007 revision.reason.validate()?;
1008 if self.revisions.iter().any(|item| item.id == revision.id) {
1009 return Err(TransportError::InvalidRevision(
1010 "itinerary revision identity must be unique within an execution".to_owned(),
1011 ));
1012 }
1013 if let Some(previous) = self.revisions.iter_mut().find(|item| item.id == active) {
1014 previous.superseded_at = Some(at);
1015 }
1016 self.estimated_arrival_at = Some(revision.plan.estimated_arrival_at);
1017 self.active_itinerary_revision = Some(revision.id);
1018 self.current_leg_index = 0;
1019 let next_leg_id = self
1020 .legs
1021 .iter()
1022 .map(|leg| leg.id.0)
1023 .max()
1024 .unwrap_or_default()
1025 .checked_add(1)
1026 .ok_or(TransportError::Overflow)?;
1027 let mut legs = Vec::with_capacity(revision.plan.legs.len());
1028 for (index, _) in revision.plan.legs.iter().enumerate() {
1029 let id = u64::try_from(index)
1030 .ok()
1031 .and_then(|index| next_leg_id.checked_add(index))
1032 .ok_or(TransportError::Overflow)?;
1033 legs.push(LegExecution {
1034 id: LegExecutionId(id),
1035 itinerary_revision: revision.id,
1036 leg_index: index,
1037 status: LegExecutionStatus::Planned,
1038 actual_departure_at: None,
1039 actual_arrival_at: None,
1040 failed_at: None,
1041 failure_reason: None,
1042 evidence: Vec::new(),
1043 });
1044 }
1045 self.legs.extend(legs);
1046 self.revisions.push(revision);
1047 if let Some(saga) = self.saga.as_mut() {
1048 saga.operation_key = delivery_completion_operation_key(
1049 self.id,
1050 self.active_itinerary_revision
1051 .ok_or(TransportError::MissingItinerary)?,
1052 saga.expected_attempt_version,
1053 );
1054 saga.evidence.extend(
1055 self.revisions
1056 .last()
1057 .map(|current| current.evidence.clone())
1058 .unwrap_or_default(),
1059 );
1060 }
1061 self.state = TransportExecutionState::Planning;
1062 Ok(())
1063 }
1064
1065 pub fn begin_saga(
1066 &mut self,
1067 delivery_attempt: DomainRecordVersionRef,
1068 operation_key: String,
1069 ) -> Result<(), TransportError> {
1070 if self.saga.is_some() {
1071 return Err(TransportError::SagaAlreadyExists);
1072 }
1073 self.saga = Some(DeliverySaga {
1074 expected_attempt_version: delivery_attempt.version,
1075 operation_key,
1076 delivery_attempt,
1077 state: SagaState::TransportIntent,
1078 step: 0,
1079 last_error: None,
1080 evidence: Vec::new(),
1081 });
1082 self.state = TransportExecutionState::Executing;
1083 Ok(())
1084 }
1085
1086 pub fn start_current_leg(&mut self, at: SimTime) -> Result<(), TransportError> {
1087 if self.state != TransportExecutionState::Ready
1088 && self.state != TransportExecutionState::Executing
1089 && self.state != TransportExecutionState::Planning
1090 {
1091 return Err(TransportError::InvalidState(
1092 "transport execution cannot start a leg in its current state".to_owned(),
1093 ));
1094 }
1095 let active = self
1096 .active_itinerary_revision
1097 .ok_or(TransportError::MissingItinerary)?;
1098 let leg = self
1099 .legs
1100 .iter_mut()
1101 .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
1102 .ok_or(TransportError::MissingLeg)?;
1103 if !matches!(
1104 leg.status,
1105 LegExecutionStatus::Planned | LegExecutionStatus::Booked | LegExecutionStatus::Waiting
1106 ) {
1107 return Err(TransportError::InvalidState(
1108 "current leg is not startable".to_owned(),
1109 ));
1110 }
1111 leg.status = LegExecutionStatus::Departed;
1112 leg.actual_departure_at = Some(at);
1113 self.state = TransportExecutionState::Executing;
1114 Ok(())
1115 }
1116
1117 pub fn complete_current_leg(
1118 &mut self,
1119 at: SimTime,
1120 endpoint: String,
1121 ) -> Result<bool, TransportError> {
1122 let active = self
1123 .active_itinerary_revision
1124 .ok_or(TransportError::MissingItinerary)?;
1125 let active_leg_count = self
1126 .legs
1127 .iter()
1128 .filter(|leg| leg.itinerary_revision == active)
1129 .count();
1130 let next_leg_index = self
1131 .current_leg_index
1132 .checked_add(1)
1133 .ok_or(TransportError::Overflow)?;
1134 let final_leg = next_leg_index >= active_leg_count;
1135 if final_leg && self.saga.is_none() {
1136 return Err(TransportError::MissingSaga);
1137 }
1138 let leg = self
1139 .legs
1140 .iter_mut()
1141 .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
1142 .ok_or(TransportError::MissingLeg)?;
1143 if leg.status != LegExecutionStatus::Departed {
1144 return Err(TransportError::InvalidState(
1145 "current leg must be departed before arrival".to_owned(),
1146 ));
1147 }
1148 if leg
1149 .actual_departure_at
1150 .is_some_and(|departure| at < departure)
1151 {
1152 return Err(TransportError::InvalidState(
1153 "arrival precedes departure".to_owned(),
1154 ));
1155 }
1156 leg.status = LegExecutionStatus::Arrived;
1157 leg.actual_arrival_at = Some(at);
1158 self.current_endpoint = Some(endpoint);
1159 self.current_leg_index = next_leg_index;
1160 if self.current_leg_index >= active_leg_count {
1161 self.state = TransportExecutionState::ArrivalPending;
1162 self.mark_arrival_pending()?;
1163 Ok(true)
1164 } else {
1165 self.state = TransportExecutionState::Ready;
1166 Ok(false)
1167 }
1168 }
1169
1170 pub fn fail_current_leg(&mut self, reason: String, at: SimTime) -> Result<(), TransportError> {
1171 self.ensure_custody_in_itinerary()?;
1172 let active = self
1173 .active_itinerary_revision
1174 .ok_or(TransportError::MissingItinerary)?;
1175 let leg = self
1176 .legs
1177 .iter_mut()
1178 .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
1179 .ok_or(TransportError::MissingLeg)?;
1180 leg.status = LegExecutionStatus::Failed;
1181 leg.failed_at = Some(at);
1182 leg.failure_reason = Some(reason);
1183 self.state = TransportExecutionState::ReplanPending;
1184 Ok(())
1185 }
1186
1187 pub fn mark_arrival_pending(&mut self) -> Result<(), TransportError> {
1188 self.ensure_custody_in_itinerary()?;
1189 let saga = self.saga.as_mut().ok_or(TransportError::MissingSaga)?;
1190 saga.step = saga.step.checked_add(1).ok_or(TransportError::Overflow)?;
1191 saga.state = SagaState::ArrivalPending;
1192 self.state = TransportExecutionState::ArrivalPending;
1193 Ok(())
1194 }
1195
1196 pub fn record_handoff(&mut self, handoff: Handoff) -> Result<(), TransportError> {
1197 let terminal_seizure = !handoff.kind.is_planned() && handoff.from_leg == handoff.to_leg;
1198 if (handoff.from_leg == handoff.to_leg && !terminal_seizure)
1199 || handoff.from_custodian.trim().is_empty()
1200 || handoff.to_custodian.trim().is_empty()
1201 || handoff.location.trim().is_empty()
1202 {
1203 return Err(TransportError::InvalidHandoff(
1204 "handoff requires distinct legs, custodians, and location".to_owned(),
1205 ));
1206 }
1207 if let HandoffKind::Seizure {
1208 by: EntityRef::Domain(record),
1209 } = &handoff.kind
1210 && !domain_record_ref_is_well_formed(record)
1211 {
1212 return Err(TransportError::InvalidHandoff(
1213 "seizure handoff requires a well-formed seizing identity".to_owned(),
1214 ));
1215 }
1216 if self
1217 .handoffs
1218 .iter()
1219 .any(|existing| existing.id == handoff.id)
1220 {
1221 return Err(TransportError::InvalidHandoff(
1222 "handoff identity is already recorded".to_owned(),
1223 ));
1224 }
1225 if self.custody_left_itinerary() {
1226 return Err(TransportError::InvalidHandoff(
1227 "custody already left the itinerary by a terminal seizure".to_owned(),
1228 ));
1229 }
1230 let from_leg = self
1231 .legs
1232 .iter()
1233 .find(|leg| leg.id == handoff.from_leg)
1234 .ok_or(TransportError::MissingLeg)?;
1235 if terminal_seizure {
1236 if self.state.is_terminal()
1237 || Some(from_leg.itinerary_revision) != self.active_itinerary_revision
1238 || from_leg.leg_index != self.current_leg_index
1239 || from_leg.status != LegExecutionStatus::Failed
1240 || from_leg
1241 .failed_at
1242 .is_some_and(|failed_at| handoff.at < failed_at)
1243 || self
1244 .handoffs
1245 .iter()
1246 .any(|existing| existing.from_leg == handoff.from_leg)
1247 {
1248 return Err(TransportError::InvalidHandoff(
1249 "a terminal seizure must follow the failed current leg of a live execution, which no other handoff leaves"
1250 .to_owned(),
1251 ));
1252 }
1253 self.handoffs.push(handoff);
1254 return Ok(());
1255 }
1256 let to_leg = self
1257 .legs
1258 .iter()
1259 .find(|leg| leg.id == handoff.to_leg)
1260 .ok_or(TransportError::MissingLeg)?;
1261 if !matches!(
1262 from_leg.status,
1263 LegExecutionStatus::Arrived | LegExecutionStatus::Failed
1264 ) || !matches!(
1265 to_leg.status,
1266 LegExecutionStatus::Planned | LegExecutionStatus::Booked | LegExecutionStatus::Waiting
1267 ) || from_leg
1268 .actual_arrival_at
1269 .is_some_and(|arrived_at| handoff.at < arrived_at)
1270 || from_leg
1271 .failed_at
1272 .is_some_and(|failed_at| handoff.at < failed_at)
1273 {
1274 return Err(TransportError::InvalidHandoff(
1275 "handoff must follow an arrived leg and precede the next leg".to_owned(),
1276 ));
1277 }
1278 self.handoffs.push(handoff);
1279 Ok(())
1280 }
1281
1282 #[must_use]
1290 pub fn custody_left_itinerary(&self) -> bool {
1291 self.handoffs
1292 .iter()
1293 .any(|handoff| handoff.from_leg == handoff.to_leg && !handoff.kind.is_planned())
1294 }
1295
1296 fn ensure_custody_in_itinerary(&self) -> Result<(), TransportError> {
1297 if self.custody_left_itinerary() {
1298 return Err(TransportError::InvalidState(
1299 "custody left the itinerary by a terminal seizure; the execution can only be closed"
1300 .to_owned(),
1301 ));
1302 }
1303 Ok(())
1304 }
1305
1306 pub fn request_booking(&mut self, booking: CapacityBooking) -> Result<(), TransportError> {
1312 self.ensure_custody_in_itinerary()?;
1313 if self.state.is_terminal() {
1314 return Err(TransportError::InvalidState(
1315 "a terminal transport execution cannot request capacity".to_owned(),
1316 ));
1317 }
1318 if booking.execution != self.id
1319 || booking.status != CapacityBookingStatus::Requested
1320 || !booking.allocation_evidence.is_empty()
1321 || booking.quantity == 0
1322 || booking.valid_until < booking.valid_from
1323 {
1324 return Err(TransportError::InvalidBooking(
1325 "a booking request must name this execution, be requested, and carry a positive quantity and window"
1326 .to_owned(),
1327 ));
1328 }
1329 if self
1330 .bookings
1331 .iter()
1332 .any(|existing| existing.id == booking.id)
1333 {
1334 return Err(TransportError::InvalidBooking(
1335 "capacity booking identity is already recorded on this execution".to_owned(),
1336 ));
1337 }
1338 self.bookings.push(booking);
1339 Ok(())
1340 }
1341
1342 pub fn settle_arrival(&mut self, at: SimTime, endpoint: String) -> Result<(), TransportError> {
1351 if self.delivery_attempt.is_some() || self.saga.is_some() {
1352 return Err(TransportError::InvalidState(
1353 "an execution with a delivery attempt settles through reconciliation".to_owned(),
1354 ));
1355 }
1356 let active = self
1357 .active_itinerary_revision
1358 .ok_or(TransportError::MissingItinerary)?;
1359 let active_leg_count = self
1360 .legs
1361 .iter()
1362 .filter(|leg| leg.itinerary_revision == active)
1363 .count();
1364 let next_leg_index = self
1365 .current_leg_index
1366 .checked_add(1)
1367 .ok_or(TransportError::Overflow)?;
1368 if next_leg_index != active_leg_count {
1369 return Err(TransportError::InvalidState(
1370 "only the final leg of the active itinerary can settle an arrival".to_owned(),
1371 ));
1372 }
1373 let leg = self
1374 .legs
1375 .iter_mut()
1376 .find(|leg| leg.itinerary_revision == active && leg.leg_index == self.current_leg_index)
1377 .ok_or(TransportError::MissingLeg)?;
1378 if leg.status != LegExecutionStatus::Departed
1379 || leg
1380 .actual_departure_at
1381 .is_some_and(|departure| at < departure)
1382 {
1383 return Err(TransportError::InvalidState(
1384 "the final leg must have departed no later than its arrival".to_owned(),
1385 ));
1386 }
1387 leg.status = LegExecutionStatus::Arrived;
1388 leg.actual_arrival_at = Some(at);
1389 self.current_endpoint = Some(endpoint);
1390 self.current_leg_index = next_leg_index;
1391 self.state = TransportExecutionState::Settled;
1392 Ok(())
1393 }
1394
1395 pub fn cancel(&mut self) -> Result<(), TransportError> {
1403 if self.state.is_terminal() || self.state == TransportExecutionState::ArrivalPending {
1404 return Err(TransportError::InvalidState(
1405 "only a non-terminal execution that has not arrived can be cancelled".to_owned(),
1406 ));
1407 }
1408 let active = self.active_itinerary_revision;
1409 if self.legs.iter().any(|leg| {
1410 Some(leg.itinerary_revision) == active && leg.status == LegExecutionStatus::Departed
1411 }) {
1412 return Err(TransportError::InvalidState(
1413 "a departed leg must arrive or fail before its execution is cancelled".to_owned(),
1414 ));
1415 }
1416 for leg in &mut self.legs {
1417 if Some(leg.itinerary_revision) == active
1418 && matches!(
1419 leg.status,
1420 LegExecutionStatus::Planned
1421 | LegExecutionStatus::Booked
1422 | LegExecutionStatus::Loaded
1423 | LegExecutionStatus::Waiting
1424 )
1425 {
1426 leg.status = LegExecutionStatus::Cancelled;
1427 }
1428 }
1429 if let Some(saga) = self.saga.as_mut() {
1430 saga.step = saga.step.checked_add(1).ok_or(TransportError::Overflow)?;
1431 saga.state = SagaState::CompensationPending;
1432 saga.last_error = Some("transport execution cancelled".to_owned());
1433 }
1434 self.state = TransportExecutionState::Cancelled;
1435 Ok(())
1436 }
1437
1438 pub fn completion_request(&self) -> Result<DeliveryCompletionRequest, TransportError> {
1439 let saga = self.saga.as_ref().ok_or(TransportError::MissingSaga)?;
1440 if self.state != TransportExecutionState::ArrivalPending
1441 || saga.state != SagaState::ArrivalPending
1442 {
1443 return Err(TransportError::InvalidState(
1444 "delivery completion requires an arrival-pending execution".to_owned(),
1445 ));
1446 }
1447 let revision = self
1448 .active_itinerary_revision
1449 .ok_or(TransportError::MissingItinerary)?;
1450 let completed_at = self
1451 .legs
1452 .iter()
1453 .rev()
1454 .find(|leg| leg.itinerary_revision == revision)
1455 .and_then(|leg| leg.actual_arrival_at)
1456 .ok_or(TransportError::MissingLeg)?;
1457 let attempt = self
1458 .delivery_attempt
1459 .clone()
1460 .ok_or(TransportError::MissingDeliveryAttempt)?;
1461 if attempt != saga.delivery_attempt {
1462 return Err(TransportError::InvalidState(
1463 "saga delivery attempt does not match execution".to_owned(),
1464 ));
1465 }
1466 Ok(DeliveryCompletionRequest {
1467 operation_key: saga.operation_key.clone(),
1468 execution: self.id,
1469 itinerary_revision: revision,
1470 delivery_attempt: attempt,
1471 completed_at,
1472 evidence: saga.evidence.clone(),
1473 })
1474 }
1475
1476 pub fn reconcile_information(
1477 &mut self,
1478 outcome: ReconciliationOutcome,
1479 ) -> Result<(), TransportError> {
1480 if matches!(outcome, ReconciliationOutcome::Success) {
1481 self.ensure_custody_in_itinerary()?;
1482 }
1483 let saga = self.saga.as_mut().ok_or(TransportError::MissingSaga)?;
1484 saga.step = saga.step.checked_add(1).ok_or(TransportError::Overflow)?;
1485 match outcome {
1486 ReconciliationOutcome::Success => {
1487 saga.state = SagaState::Settled;
1488 saga.last_error = None;
1489 self.state = TransportExecutionState::Settled;
1490 }
1491 ReconciliationOutcome::Failure { error } => {
1492 saga.state = SagaState::CompensationPending;
1493 saga.last_error = Some(error);
1494 self.state = TransportExecutionState::Failed;
1495 }
1496 }
1497 Ok(())
1498 }
1499}
1500
1501#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
1502pub enum TransportError {
1503 InvalidRevision(String),
1504 MissingItinerary,
1505 MissingSaga,
1506 SagaAlreadyExists,
1507 Overflow,
1508 MissingLeg,
1509 InvalidState(String),
1510 InvalidBooking(String),
1511 InvalidHandoff(String),
1512 MissingDeliveryAttempt,
1513 InvalidCapacityPool(String),
1515}
1516
1517#[cfg(test)]
1518mod tests {
1519 use super::*;
1520 use canwu_core::{PersonId, TerritoryId};
1521 use canwu_routing::{
1522 ROUTING_ALGORITHM_VERSION, RouteCost, RouteLeg, RoutingConnectionRef, RoutingNodeRef,
1523 TransferMode,
1524 };
1525
1526 fn plan() -> RoutePlan {
1527 RoutePlan {
1528 algorithm_version: ROUTING_ALGORITHM_VERSION.to_owned(),
1529 policy_version: "policy.v1".to_owned(),
1530 planning_snapshot_digest: "snapshot".to_owned(),
1531 origin: RoutingNodeRef::new("a"),
1532 destination: RoutingNodeRef::new("b"),
1533 departure_at: canwu_time::SimTime::EPOCH,
1534 estimated_arrival_at: canwu_time::SimTime::from_minutes(10),
1535 cost: RouteCost {
1536 estimated_arrival_at: canwu_time::SimTime::from_minutes(10),
1537 risk_per_mille: 0,
1538 resource_cost: 0,
1539 transfers: 1,
1540 },
1541 legs: vec![RouteLeg {
1542 connection: RoutingConnectionRef::new("ab"),
1543 from: RoutingNodeRef::new("a"),
1544 to: RoutingNodeRef::new("b"),
1545 mode: TransferMode::Horse,
1546 planned_departure_at: canwu_time::SimTime::EPOCH,
1547 planned_arrival_at: canwu_time::SimTime::from_minutes(10),
1548 }],
1549 digest: "route".to_owned(),
1550 }
1551 }
1552
1553 #[test]
1554 fn movement_order_accepts_a_self_directed_person_subject() {
1555 let order = MovementOrder {
1556 id: MovementOrderId(1),
1557 subjects: vec![MovementSubject {
1558 entity: EntityRef::Person(PersonId::new(7)),
1559 role: MovementSubjectRole::MovablePrincipal,
1560 quantity: None,
1561 expected_custody: None,
1562 }],
1563 origin: RoutingNodeRef::new("a"),
1564 destination: RoutingNodeRef::new("b"),
1565 plan: plan(),
1566 initiative: MovementInitiative::SelfDirected,
1567 ordered_at: SimTime::EPOCH,
1568 expected_position_revision: 1,
1569 };
1570 order.validate().unwrap();
1571 }
1572
1573 #[test]
1574 fn movement_order_rejects_duplicate_subjects_and_missing_cargo_quantity() {
1575 let mut order = MovementOrder {
1576 id: MovementOrderId(2),
1577 subjects: vec![
1578 MovementSubject {
1579 entity: EntityRef::Person(PersonId::new(7)),
1580 role: MovementSubjectRole::MovablePrincipal,
1581 quantity: None,
1582 expected_custody: None,
1583 },
1584 MovementSubject {
1585 entity: EntityRef::Person(PersonId::new(7)),
1586 role: MovementSubjectRole::Cargo,
1587 quantity: None,
1588 expected_custody: None,
1589 },
1590 ],
1591 origin: RoutingNodeRef::new("a"),
1592 destination: RoutingNodeRef::new("b"),
1593 plan: plan(),
1594 initiative: MovementInitiative::Delegated,
1595 ordered_at: SimTime::EPOCH,
1596 expected_position_revision: 1,
1597 };
1598 assert!(matches!(
1599 order.validate(),
1600 Err(MovementOrderError::Invalid(_))
1601 ));
1602 order.subjects[1].entity = EntityRef::Territory(TerritoryId::new(8));
1603 assert!(matches!(
1604 order.validate(),
1605 Err(MovementOrderError::Invalid(_))
1606 ));
1607 }
1608
1609 #[test]
1610 fn movement_order_rejects_a_plan_that_does_not_match_the_order() {
1611 let mut order = MovementOrder {
1612 id: MovementOrderId(3),
1613 subjects: vec![MovementSubject {
1614 entity: EntityRef::Person(PersonId::new(9)),
1615 role: MovementSubjectRole::MovablePrincipal,
1616 quantity: None,
1617 expected_custody: None,
1618 }],
1619 origin: RoutingNodeRef::new("a"),
1620 destination: RoutingNodeRef::new("c"),
1621 plan: plan(),
1622 initiative: MovementInitiative::SelfDirected,
1623 ordered_at: SimTime::EPOCH,
1624 expected_position_revision: 1,
1625 };
1626 assert!(matches!(
1627 order.validate(),
1628 Err(MovementOrderError::Invalid(_))
1629 ));
1630 order.destination = RoutingNodeRef::new("b");
1631 order.plan.digest = String::new();
1632 assert!(matches!(
1633 order.validate(),
1634 Err(MovementOrderError::Invalid(_))
1635 ));
1636 }
1637
1638 #[test]
1639 fn reroute_supersedes_without_creating_a_new_delivery_attempt() {
1640 let mut execution = TransportExecution::new(TransportExecutionId(1), None);
1641 execution
1642 .install_initial_itinerary(ItineraryRevision {
1643 id: ItineraryRevisionId(1),
1644 predecessor: None,
1645 plan: plan(),
1646 planned_at: canwu_time::SimTime::EPOCH,
1647 valid_from: canwu_time::SimTime::EPOCH,
1648 reason: ItineraryRevisionReason::Initial,
1649 superseded_at: None,
1650 evidence: Vec::new(),
1651 })
1652 .unwrap();
1653 let mut replacement = plan();
1654 replacement.destination = RoutingNodeRef::new("b");
1655 execution
1656 .reroute(
1657 ItineraryRevision {
1658 id: ItineraryRevisionId(2),
1659 predecessor: Some(ItineraryRevisionId(1)),
1660 plan: replacement,
1661 planned_at: canwu_time::SimTime::from_minutes(1),
1662 valid_from: canwu_time::SimTime::from_minutes(1),
1663 reason: ItineraryRevisionReason::Disaster {
1664 explanation: "bridge closed".to_owned(),
1665 },
1666 superseded_at: None,
1667 evidence: Vec::new(),
1668 },
1669 canwu_time::SimTime::from_minutes(1),
1670 )
1671 .unwrap();
1672 assert_eq!(execution.revisions.len(), 2);
1673 assert_eq!(
1674 execution.active_itinerary_revision,
1675 Some(ItineraryRevisionId(2))
1676 );
1677 assert_eq!(execution.state, TransportExecutionState::Planning);
1678 }
1679
1680 #[test]
1681 fn capacity_booking_is_a_persisted_windowed_state_machine() {
1682 let mut booking = CapacityBooking::new(
1683 CapacityBookingId(1),
1684 TransportExecutionId(1),
1685 "relay-horse:wu-xi:01".to_owned(),
1686 canwu_time::SimTime::EPOCH,
1687 canwu_time::SimTime::from_minutes(60),
1688 1,
1689 10,
1690 )
1691 .unwrap();
1692 booking
1693 .transition(CapacityBookingStatus::Confirmed, canwu_time::SimTime::EPOCH)
1694 .unwrap();
1695 booking
1696 .transition(
1697 CapacityBookingStatus::Consumed,
1698 canwu_time::SimTime::from_minutes(10),
1699 )
1700 .unwrap();
1701 assert_eq!(booking.status, CapacityBookingStatus::Consumed);
1702 }
1703
1704 #[test]
1705 fn completion_operation_key_is_stable_and_revision_scoped() {
1706 let first =
1707 delivery_completion_operation_key(TransportExecutionId(4), ItineraryRevisionId(2), 7);
1708 let second =
1709 delivery_completion_operation_key(TransportExecutionId(4), ItineraryRevisionId(3), 7);
1710 assert_ne!(first, second);
1711 assert_eq!(
1712 first,
1713 delivery_completion_operation_key(TransportExecutionId(4), ItineraryRevisionId(2), 7)
1714 );
1715 }
1716
1717 #[test]
1718 fn completion_requires_saga_and_exposes_stable_bridge_request() {
1719 let attempt = DomainRecordVersionRef {
1720 record: canwu_core::DomainRecordRef::new(
1721 "fixture.information",
1722 "delivery_attempt",
1723 "delivery",
1724 ),
1725 version: 3,
1726 established_by: canwu_core::DomainRecordVersionSource::InitialScenario,
1727 };
1728 let mut execution = TransportExecution::new(TransportExecutionId(9), Some(attempt.clone()));
1729 let revision = ItineraryRevision {
1730 id: ItineraryRevisionId(1),
1731 predecessor: None,
1732 plan: plan(),
1733 planned_at: SimTime::EPOCH,
1734 valid_from: SimTime::EPOCH,
1735 reason: ItineraryRevisionReason::Initial,
1736 superseded_at: None,
1737 evidence: Vec::new(),
1738 };
1739 execution.install_initial_itinerary(revision).unwrap();
1740 assert_eq!(
1741 execution.complete_current_leg(SimTime::from_minutes(10), "b".to_owned()),
1742 Err(TransportError::MissingSaga)
1743 );
1744 execution
1745 .begin_saga(
1746 attempt,
1747 delivery_completion_operation_key(
1748 TransportExecutionId(9),
1749 ItineraryRevisionId(1),
1750 3,
1751 ),
1752 )
1753 .unwrap();
1754 execution.start_current_leg(SimTime::EPOCH).unwrap();
1755 assert!(
1756 execution
1757 .complete_current_leg(SimTime::from_minutes(10), "b".to_owned())
1758 .unwrap()
1759 );
1760 let request = execution.completion_request().unwrap();
1761 assert_eq!(request.execution, TransportExecutionId(9));
1762 assert_eq!(request.itinerary_revision, ItineraryRevisionId(1));
1763 assert_eq!(request.delivery_attempt.version, 3);
1764 }
1765
1766 #[test]
1767 fn handoff_requires_arrival_and_is_idempotency_safe() {
1768 let attempt = DomainRecordVersionRef {
1769 record: canwu_core::DomainRecordRef::new(
1770 "fixture.information",
1771 "delivery_attempt",
1772 "delivery",
1773 ),
1774 version: 1,
1775 established_by: canwu_core::DomainRecordVersionSource::InitialScenario,
1776 };
1777 let mut execution =
1778 TransportExecution::new(TransportExecutionId(10), Some(attempt.clone()));
1779 let mut route = plan();
1780 route.legs.push(RouteLeg {
1781 connection: RoutingConnectionRef::new("bc"),
1782 from: RoutingNodeRef::new("b"),
1783 to: RoutingNodeRef::new("c"),
1784 mode: TransferMode::Rail,
1785 planned_departure_at: SimTime::from_minutes(10),
1786 planned_arrival_at: SimTime::from_minutes(20),
1787 });
1788 route.destination = RoutingNodeRef::new("c");
1789 route.estimated_arrival_at = SimTime::from_minutes(20);
1790 let revision = ItineraryRevision {
1791 id: ItineraryRevisionId(1),
1792 predecessor: None,
1793 plan: route,
1794 planned_at: SimTime::EPOCH,
1795 valid_from: SimTime::EPOCH,
1796 reason: ItineraryRevisionReason::Initial,
1797 superseded_at: None,
1798 evidence: Vec::new(),
1799 };
1800 execution.install_initial_itinerary(revision).unwrap();
1801 execution.start_current_leg(SimTime::EPOCH).unwrap();
1802 execution
1803 .complete_current_leg(SimTime::from_minutes(10), "b".to_owned())
1804 .unwrap();
1805 let handoff = Handoff {
1806 id: HandoffId(1),
1807 from_leg: LegExecutionId(1),
1808 to_leg: LegExecutionId(2),
1809 from_custodian: "courier/wuxi".to_owned(),
1810 to_custodian: "rail/beijing".to_owned(),
1811 at: SimTime::from_minutes(10),
1812 location: "b".to_owned(),
1813 evidence: Vec::new(),
1814 kind: HandoffKind::Planned,
1815 };
1816 execution.record_handoff(handoff.clone()).unwrap();
1817 assert_eq!(
1818 execution.record_handoff(handoff),
1819 Err(TransportError::InvalidHandoff(
1820 "handoff identity is already recorded".to_owned()
1821 ))
1822 );
1823 }
1824}