automation_structures/modalities/
stream_graph_fanout.rs1use vstd::prelude::*;
6
7use crate::connectives::buffer::Buffer;
8use crate::connectives::counter::Counter;
9
10verus! {
11
12pub struct StreamGraphFanout {
14 pub max_inputs: usize,
16 pub record_domain_size: u64,
18 pub left_queue: Buffer<u64>,
20 pub right_queue: Buffer<u64>,
22 pub ingested: Counter,
24 pub left_emitted: Counter,
26 pub right_emitted: Counter,
28}
29
30impl StreamGraphFanout {
31 pub open spec fn values_valid(q: Seq<u64>, domain: u64) -> bool {
33 forall|i: int| 0 <= i < q.len() ==> #[trigger] q[i] < domain
34 }
35
36 pub open spec fn type_invariant(&self) -> bool {
38 &&& self.left_queue.capacity > 0
39 &&& self.right_queue.capacity == self.left_queue.capacity
40 &&& self.record_domain_size > 0
41 &&& Self::values_valid(self.left_queue.values@, self.record_domain_size)
42 &&& Self::values_valid(self.right_queue.values@, self.record_domain_size)
43 &&& self.ingested.value_spec() <= self.max_inputs as nat
44 &&& self.left_emitted.value_spec() <= self.max_inputs as nat
45 &&& self.right_emitted.value_spec() <= self.max_inputs as nat
46 }
47
48 pub open spec fn backpressure_correct(&self) -> bool {
50 self.left_queue.well_formed() && self.right_queue.well_formed()
51 }
52
53 pub open spec fn per_branch_conservation(&self) -> bool {
55 &&& self.ingested.value_spec()
56 == self.left_queue.values@.len() + self.left_emitted.value_spec()
57 &&& self.ingested.value_spec()
58 == self.right_queue.values@.len() + self.right_emitted.value_spec()
59 }
60
61 pub open spec fn inv(&self) -> bool {
63 self.type_invariant()
64 && self.backpressure_correct()
65 && self.per_branch_conservation()
66 }
67
68 pub fn valid_config(capacity: usize, record_domain_size: u64) -> (valid: bool)
70 ensures valid == (capacity > 0 && record_domain_size > 0),
71 {
72 capacity > 0 && record_domain_size > 0
73 }
74
75 pub fn new(
77 capacity: usize,
78 max_inputs: usize,
79 record_domain_size: u64,
80 ) -> (s: StreamGraphFanout)
81 requires capacity > 0, record_domain_size > 0,
82 ensures
83 s.left_queue.capacity == capacity,
84 s.right_queue.capacity == capacity,
85 s.max_inputs == max_inputs,
86 s.record_domain_size == record_domain_size,
87 s.left_queue.values@.len() == 0,
88 s.right_queue.values@.len() == 0,
89 s.ingested.value_spec() == 0,
90 s.left_emitted.value_spec() == 0,
91 s.right_emitted.value_spec() == 0,
92 s.inv(),
93 {
94 StreamGraphFanout {
95 max_inputs,
96 record_domain_size,
97 left_queue: Buffer::new(capacity),
98 right_queue: Buffer::new(capacity),
99 ingested: Counter::new(0),
100 left_emitted: Counter::new(0),
101 right_emitted: Counter::new(0),
102 }
103 }
104
105 pub fn source_ingest(&mut self, value: u64) -> (accepted: bool)
107 requires old(self).inv(),
108 ensures
109 accepted == (value < old(self).record_domain_size
110 && old(self).ingested.value_spec() < old(self).max_inputs as nat
111 && old(self).left_queue.values@.len() < old(self).left_queue.capacity
112 && old(self).right_queue.values@.len() < old(self).right_queue.capacity),
113 final(self).left_queue.capacity == old(self).left_queue.capacity,
114 final(self).right_queue.capacity == old(self).right_queue.capacity,
115 final(self).max_inputs == old(self).max_inputs,
116 final(self).record_domain_size == old(self).record_domain_size,
117 final(self).left_queue.values@ == if accepted {
118 old(self).left_queue.values@.push(value)
119 } else { old(self).left_queue.values@ },
120 final(self).right_queue.values@ == if accepted {
121 old(self).right_queue.values@.push(value)
122 } else { old(self).right_queue.values@ },
123 accepted ==> final(self).ingested.value_spec()
124 == old(self).ingested.value_spec() + 1,
125 !accepted ==> final(self).ingested.value_spec()
126 == old(self).ingested.value_spec(),
127 final(self).left_emitted.value_spec() == old(self).left_emitted.value_spec(),
128 final(self).right_emitted.value_spec() == old(self).right_emitted.value_spec(),
129 final(self).inv(),
130 {
131 if value < self.record_domain_size
132 && self.ingested.value() < self.max_inputs as u64
133 && self.left_queue.len() < self.left_queue.capacity
134 && self.right_queue.len() < self.right_queue.capacity
135 {
136 let ghost old_left = self.left_queue.values@;
137 let ghost old_right = self.right_queue.values@;
138 let _left_pushed = self.left_queue.push(value);
139 let _right_pushed = self.right_queue.push(value);
140 let _counted = self.ingested.try_increment();
141 assert(_counted);
142 assert forall|i: int| 0 <= i < self.left_queue.values@.len()
143 implies #[trigger] self.left_queue.values@[i] < self.record_domain_size by {
144 if i < old_left.len() {
145 assert(self.left_queue.values@[i] == old_left[i]);
146 }
147 }
148 assert forall|i: int| 0 <= i < self.right_queue.values@.len()
149 implies #[trigger] self.right_queue.values@[i] < self.record_domain_size by {
150 if i < old_right.len() {
151 assert(self.right_queue.values@[i] == old_right[i]);
152 }
153 }
154 true
155 } else {
156 false
157 }
158 }
159
160 pub fn consume_left(&mut self) -> (accepted: bool)
162 requires old(self).inv(),
163 ensures
164 accepted == (old(self).left_queue.values@.len() > 0),
165 final(self).left_queue.capacity == old(self).left_queue.capacity,
166 final(self).right_queue.capacity == old(self).right_queue.capacity,
167 final(self).max_inputs == old(self).max_inputs,
168 final(self).record_domain_size == old(self).record_domain_size,
169 final(self).left_queue.values@ == if accepted {
170 old(self).left_queue.values@.subrange(1, old(self).left_queue.values@.len() as int)
171 } else { old(self).left_queue.values@ },
172 final(self).right_queue.values@ == old(self).right_queue.values@,
173 final(self).ingested.value_spec() == old(self).ingested.value_spec(),
174 accepted ==> final(self).left_emitted.value_spec()
175 == old(self).left_emitted.value_spec() + 1,
176 !accepted ==> final(self).left_emitted.value_spec()
177 == old(self).left_emitted.value_spec(),
178 final(self).right_emitted.value_spec() == old(self).right_emitted.value_spec(),
179 final(self).inv(),
180 {
181 if self.left_queue.len() > 0 {
182 let ghost old_left = self.left_queue.values@;
183 let _popped = self.left_queue.pop();
184 let _counted = self.left_emitted.try_increment();
185 assert(_counted);
186 assert(self.left_queue.values@ =~=
187 old_left.subrange(1, old_left.len() as int));
188 true
189 } else {
190 false
191 }
192 }
193
194 pub fn consume_right(&mut self) -> (accepted: bool)
196 requires old(self).inv(),
197 ensures
198 accepted == (old(self).right_queue.values@.len() > 0),
199 final(self).left_queue.capacity == old(self).left_queue.capacity,
200 final(self).right_queue.capacity == old(self).right_queue.capacity,
201 final(self).max_inputs == old(self).max_inputs,
202 final(self).record_domain_size == old(self).record_domain_size,
203 final(self).left_queue.values@ == old(self).left_queue.values@,
204 final(self).right_queue.values@ == if accepted {
205 old(self).right_queue.values@.subrange(1, old(self).right_queue.values@.len() as int)
206 } else { old(self).right_queue.values@ },
207 final(self).ingested.value_spec() == old(self).ingested.value_spec(),
208 final(self).left_emitted.value_spec() == old(self).left_emitted.value_spec(),
209 accepted ==> final(self).right_emitted.value_spec()
210 == old(self).right_emitted.value_spec() + 1,
211 !accepted ==> final(self).right_emitted.value_spec()
212 == old(self).right_emitted.value_spec(),
213 final(self).inv(),
214 {
215 if self.right_queue.len() > 0 {
216 let ghost old_right = self.right_queue.values@;
217 let _popped = self.right_queue.pop();
218 let _counted = self.right_emitted.try_increment();
219 assert(_counted);
220 assert(self.right_queue.values@ =~=
221 old_right.subrange(1, old_right.len() as int));
222 true
223 } else {
224 false
225 }
226 }
227
228 pub fn terminal(&self) -> (terminal: bool)
230 requires self.inv(),
231 ensures terminal == (self.ingested.value_spec() == self.max_inputs as nat
232 && self.left_queue.values@.len() == 0
233 && self.right_queue.values@.len() == 0),
234 {
235 self.ingested.value() == self.max_inputs as u64
236 && self.left_queue.len() == 0
237 && self.right_queue.len() == 0
238 }
239}
240
241}