Skip to main content

meerkat_machine_schema/catalog/dsl/
runtime_delivery.rs

1use 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}