Skip to main content

automation_structures/modalities/
stream_graph_fanout.rs

1// Executable carrier for the bounded one-source/two-sink StreamGraph fan-out
2// profile. A source commit broadcasts one value to both queues. Rejected calls
3// stutter, and conservation is stated independently for each branch.
4
5use vstd::prelude::*;
6
7use crate::connectives::buffer::Buffer;
8use crate::connectives::counter::Counter;
9
10verus! {
11
12/// Retained verification profile for a source fanning out to two FIFO sinks.
13pub struct StreamGraphFanout {
14    /// Maximum records admitted by the source.
15    pub max_inputs: usize,
16    /// Exclusive upper bound of record values.
17    pub record_domain_size: u64,
18    /// Left-branch FIFO owner.
19    pub left_queue: Buffer<u64>,
20    /// Right-branch FIFO owner.
21    pub right_queue: Buffer<u64>,
22    /// Source-admission counter.
23    pub ingested: Counter,
24    /// Left-sink emission counter.
25    pub left_emitted: Counter,
26    /// Right-sink emission counter.
27    pub right_emitted: Counter,
28}
29
30impl StreamGraphFanout {
31    /// Whether every queued value lies within `domain`.
32    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    /// Whether counters, queues, and configuration values have valid shape and bounds.
37    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    /// Whether either full branch prevents another broadcast.
49    pub open spec fn backpressure_correct(&self) -> bool {
50        self.left_queue.well_formed() && self.right_queue.well_formed()
51    }
52
53    /// Whether each branch independently conserves admitted records.
54    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    /// Whether all fan-out stream contract clauses hold.
62    pub open spec fn inv(&self) -> bool {
63        self.type_invariant()
64            && self.backpressure_correct()
65            && self.per_branch_conservation()
66    }
67
68    /// Test whether the shared queue capacity and value domain are valid.
69    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    /// Construct an empty valid fan-out execution.
76    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    /// Replicate one source record into both branch queues.
106    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    /// Consume one record from the left FIFO branch.
161    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    /// Consume one record from the right FIFO branch.
195    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    /// Whether bounded input was admitted and both branches are drained.
229    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}