meerkat_machine_schema/catalog/dsl/
runtime_delivery.rs1use meerkat_machine_dsl::machine;
2
3use super::OptionValueExt;
4
5machine! {
6 machine RuntimeDeliveryMachine {
7 version: 1,
8 rust: "self" / "catalog::dsl::runtime_delivery",
9
10 state {
11 lifecycle_phase: RuntimeDeliveryPhase,
12 delivery_ids: Set<String>,
13 delivery_sequences: Map<String, u64>,
14 delivery_source_sequences: Map<String, u64>,
15 committed_sequences: Set<u64>,
16 next_sequence: u64,
17 applied_cursor: u64,
18 }
19
20 init(Active) {
21 delivery_ids = EmptySet,
22 delivery_sequences = EmptyMap,
23 delivery_source_sequences = EmptyMap,
24 committed_sequences = EmptySet,
25 next_sequence = 0,
26 applied_cursor = 0,
27 }
28
29 terminal []
30
31 phase RuntimeDeliveryPhase {
32 Active,
33 }
34
35 input RuntimeDeliveryInput {
36 CommitDelivery {
37 delivery_id: String,
38 source_sequence: u64,
39 },
40 MarkDeliveryApplied {
41 delivery_id: String,
42 delivery_sequence: u64,
43 },
44 }
45
46 effect RuntimeDeliveryEffect {
47 DeliveryCommitted {
48 delivery_id: String,
49 source_sequence: u64,
50 delivery_sequence: u64,
51 },
52 DeliveryReused {
53 delivery_id: String,
54 source_sequence: u64,
55 delivery_sequence: u64,
56 },
57 DeliveryApplied {
58 delivery_id: String,
59 delivery_sequence: u64,
60 },
61 }
62
63 invariant applied_cursor_does_not_pass_committed_sequence {
64 self.applied_cursor <= self.next_sequence
65 }
66
67 invariant empty_delivery_set_has_zero_sequence {
68 self.delivery_ids.len() != 0 || self.next_sequence == 0
69 }
70
71 invariant delivery_identity_and_sequence_cardinality_match {
72 self.delivery_ids.len() == self.committed_sequences.len()
73 }
74
75 invariant committed_sequence_cardinality_tracks_high_water {
76 self.committed_sequences.len() == self.next_sequence
77 }
78
79 disposition DeliveryCommitted => routed [DetachedJobMachine] seam NoOwnerRealization,
80 disposition DeliveryReused => routed [DetachedJobMachine] seam NoOwnerRealization,
81 disposition DeliveryApplied => local seam OwnerRealizationOnly,
82
83 transition CommitNewDelivery {
84 on input CommitDelivery { delivery_id, source_sequence }
85 guard {
86 self.lifecycle_phase == Phase::Active
87 && self.delivery_ids.contains(delivery_id) == false
88 && source_sequence > 0
89 && self.next_sequence < u64::MAX
90 }
91 update {
92 self.next_sequence += 1;
93 self.delivery_ids.insert(delivery_id);
94 self.delivery_sequences.insert(delivery_id, self.next_sequence);
95 self.delivery_source_sequences.insert(delivery_id, source_sequence);
96 self.committed_sequences.insert(self.next_sequence);
97 }
98 to Active
99 emit DeliveryCommitted {
100 delivery_id: delivery_id,
101 source_sequence: source_sequence,
102 delivery_sequence: self.next_sequence
103 }
104 }
105
106 transition ReuseCommittedDelivery {
107 on input CommitDelivery { delivery_id, source_sequence }
108 guard {
109 self.lifecycle_phase == Phase::Active
110 && self.delivery_ids.contains(delivery_id)
111 && self.delivery_source_sequences.get_cloned(delivery_id).get("value") == source_sequence
112 }
113 update {}
114 to Active
115 emit DeliveryReused {
116 delivery_id: delivery_id,
117 source_sequence: source_sequence,
118 delivery_sequence: self.delivery_sequences.get_cloned(delivery_id).get("value")
119 }
120 }
121
122 transition ApplyNextDelivery {
123 on input MarkDeliveryApplied { delivery_id, delivery_sequence }
124 guard {
125 self.lifecycle_phase == Phase::Active
126 && self.delivery_ids.contains(delivery_id)
127 && self.delivery_sequences.get_cloned(delivery_id).get("value") == delivery_sequence
128 && delivery_sequence > self.applied_cursor
129 && delivery_sequence - 1 == self.applied_cursor
130 }
131 update {
132 self.applied_cursor = delivery_sequence;
133 }
134 to Active
135 emit DeliveryApplied {
136 delivery_id: delivery_id,
137 delivery_sequence: delivery_sequence
138 }
139 }
140
141 transition ObserveAlreadyAppliedDelivery {
142 on input MarkDeliveryApplied { delivery_id, delivery_sequence }
143 guard {
144 self.lifecycle_phase == Phase::Active
145 && self.delivery_ids.contains(delivery_id)
146 && self.delivery_sequences.get_cloned(delivery_id).get("value") == delivery_sequence
147 && delivery_sequence <= self.applied_cursor
148 }
149 update {}
150 to Active
151 emit DeliveryApplied {
152 delivery_id: delivery_id,
153 delivery_sequence: delivery_sequence
154 }
155 }
156 }
157}