Skip to main content

automation_structures/modalities/
stream_graph.rs

1// Faithful executable carrier for the specified three- and four-node StreamGraph chains.
2// q1, q2, and q3 correspond exactly to the named TLA+ edge queues; q3 is
3// empty for a three-node chain. Rejected calls stutter.
4
5use vstd::prelude::*;
6
7use crate::connectives::buffer::Buffer;
8use crate::connectives::counter::Counter;
9
10verus! {
11
12/// Erased logical view shared by every bounded-transfer realization.
13pub ghost struct BoundedTransferModel<T> {
14    /// Maximum retained record count.
15    pub slot_capacity: nat,
16    /// Maximum retained encoded byte count.
17    pub retained_byte_capacity: nat,
18    /// Retained records in FIFO order.
19    pub values: Seq<T>,
20    /// Encoded key-size registry for retained records.
21    pub registry: Seq<(u64, u64)>,
22    /// Total encoded bytes retained.
23    pub retained_bytes: nat,
24    /// Monotone consumed-record cursor.
25    pub head: nat,
26    /// Monotone admitted-record cursor.
27    pub tail: nat,
28    /// Whether the transfer owner is closed to new records.
29    pub closed: bool,
30}
31
32impl<T> BoundedTransferModel<T> {
33    /// Primitive-rooted and connective-rooted state invariant.
34    pub open spec fn inv(self) -> bool {
35        &&& self.slot_capacity > 0
36        &&& self.retained_byte_capacity > 0
37        &&& crate::primitives::budget::budget_safety(
38            self.slot_capacity,
39            self.values.len(),
40            0,
41            0,
42        )
43        &&& crate::primitives::budget::budget_safety(
44            self.retained_byte_capacity,
45            self.retained_bytes,
46            0,
47            0,
48        )
49        &&& crate::connectives::buffer::buffer_bounded(
50            self.values,
51            self.slot_capacity,
52        )
53        &&& self.registry.len() == self.values.len()
54        &&& crate::primitives::resource_registry::unique_mapping_entries(
55            self.registry,
56        )
57        &&& crate::connectives::cursor::cursor_admitted(self.head, self.tail)
58        &&& crate::connectives::ordering_pass::fifo_sequence_order(
59            self.registry,
60            self.head,
61            self.tail,
62        )
63    }
64
65    /// Readiness is derived from Buffer occupancy. It is not parallel state.
66    pub open spec fn ready(self) -> bool {
67        self.values.len() > 0
68    }
69
70    /// The closed transfer is terminal only after its Buffer drains.
71    pub open spec fn terminal(self) -> bool {
72        self.closed && self.values.len() == 0
73    }
74}
75
76/// One successful publication through the shared bounded-transfer composition.
77pub open spec fn bounded_transfer_publish<T>(
78    pre: BoundedTransferModel<T>,
79    post: BoundedTransferModel<T>,
80    value: T,
81    retained_bytes: u64,
82    sequence: u64,
83) -> bool {
84    &&& pre.inv()
85    &&& post.inv()
86    &&& !pre.closed
87    &&& pre.values.len() < pre.slot_capacity
88    &&& pre.retained_bytes + retained_bytes as nat <= pre.retained_byte_capacity
89    &&& pre.tail < u64::MAX as nat
90    &&& sequence as nat == pre.tail
91    &&& post.slot_capacity == pre.slot_capacity
92    &&& post.retained_byte_capacity == pre.retained_byte_capacity
93    &&& post.values == pre.values.push(value)
94    &&& post.registry == pre.registry.push((sequence, retained_bytes))
95    &&& post.retained_bytes == pre.retained_bytes + retained_bytes as nat
96    &&& post.head == pre.head
97    &&& post.tail == pre.tail + 1
98    &&& post.closed == pre.closed
99}
100
101/// One successful FIFO receipt through the shared bounded-transfer composition.
102pub open spec fn bounded_transfer_receive<T>(
103    pre: BoundedTransferModel<T>,
104    post: BoundedTransferModel<T>,
105    value: T,
106    retained_bytes: u64,
107    sequence: u64,
108) -> bool {
109    &&& pre.inv()
110    &&& post.inv()
111    &&& pre.values.len() > 0
112    &&& value == pre.values[0]
113    &&& (sequence, retained_bytes) == pre.registry[0]
114    &&& post.slot_capacity == pre.slot_capacity
115    &&& post.retained_byte_capacity == pre.retained_byte_capacity
116    &&& post.values == pre.values.skip(1)
117    &&& post.registry == pre.registry.skip(1)
118    &&& post.retained_bytes + retained_bytes as nat == pre.retained_bytes
119    &&& post.head == pre.head + 1
120    &&& post.tail == pre.tail
121    &&& post.closed == pre.closed
122}
123
124/// Idempotent producer closure through the shared bounded-transfer composition.
125pub open spec fn bounded_transfer_close<T>(
126    pre: BoundedTransferModel<T>,
127    post: BoundedTransferModel<T>,
128    changed: bool,
129) -> bool {
130    &&& pre.inv()
131    &&& post.inv()
132    &&& changed == !pre.closed
133    &&& post.slot_capacity == pre.slot_capacity
134    &&& post.retained_byte_capacity == pre.retained_byte_capacity
135    &&& post.values == pre.values
136    &&& post.registry == pre.registry
137    &&& post.retained_bytes == pre.retained_bytes
138    &&& post.head == pre.head
139    &&& post.tail == pre.tail
140    &&& post.closed
141}
142
143/// Every refused action stutters at the complete logical boundary.
144pub open spec fn bounded_transfer_refusal_stutters<T>(
145    pre: BoundedTransferModel<T>,
146    post: BoundedTransferModel<T>,
147) -> bool {
148    pre == post
149}
150
151/// Empty admitted origin for every reachable bounded-transfer execution.
152///
153/// Starting both ledgers and the registry at zero, then using only `publish` and `receive`, binds
154/// retained-byte Budget changes to the exact ResourceRegistry charge appended or removed.
155pub open spec fn bounded_transfer_initial<T>(model: BoundedTransferModel<T>) -> bool {
156    &&& model.inv()
157    &&& model.values == Seq::<T>::empty()
158    &&& model.registry == Seq::<(u64, u64)>::empty()
159    &&& model.retained_bytes == 0
160    &&& model.head == 0
161    &&& model.tail == 0
162    &&& !model.closed
163}
164
165/// Publish one value through one StreamGraph connection.
166pub open spec fn connection_publish<T>(
167    pre: BoundedTransferModel<T>,
168    post: BoundedTransferModel<T>,
169    value: T,
170    retained_bytes: u64,
171    sequence: u64,
172) -> bool {
173    bounded_transfer_publish(pre, post, value, retained_bytes, sequence)
174}
175
176/// Receive one value through one StreamGraph connection.
177pub open spec fn connection_receive<T>(
178    pre: BoundedTransferModel<T>,
179    post: BoundedTransferModel<T>,
180    value: T,
181    retained_bytes: u64,
182    sequence: u64,
183) -> bool {
184    bounded_transfer_receive(pre, post, value, retained_bytes, sequence)
185}
186
187/// Close one StreamGraph connection without disturbing queued values.
188pub open spec fn connection_close<T>(
189    pre: BoundedTransferModel<T>,
190    post: BoundedTransferModel<T>,
191    changed: bool,
192) -> bool {
193    bounded_transfer_close(pre, post, changed)
194}
195
196/// A refused StreamGraph endpoint action leaves the complete connection unchanged.
197pub open spec fn connection_refusal<T>(
198    pre: BoundedTransferModel<T>,
199    post: BoundedTransferModel<T>,
200) -> bool {
201    bounded_transfer_refusal_stutters(pre, post)
202}
203
204/// Publish one value to each of two connections as one StreamGraph action.
205pub open spec fn publish_pair<T>(
206    left_pre: BoundedTransferModel<T>,
207    left_post: BoundedTransferModel<T>,
208    left_value: T,
209    left_sequence: u64,
210    right_pre: BoundedTransferModel<T>,
211    right_post: BoundedTransferModel<T>,
212    right_value: T,
213    right_sequence: u64,
214    retained_bytes: u64,
215) -> bool {
216    &&& connection_publish(
217        left_pre,
218        left_post,
219        left_value,
220        retained_bytes,
221        left_sequence,
222    )
223    &&& connection_publish(
224        right_pre,
225        right_post,
226        right_value,
227        retained_bytes,
228        right_sequence,
229    )
230}
231
232/// Receive one oldest value from each of two connections as one StreamGraph action.
233pub open spec fn receive_pair<T>(
234    left_pre: BoundedTransferModel<T>,
235    left_post: BoundedTransferModel<T>,
236    left_value: T,
237    left_retained_bytes: u64,
238    left_sequence: u64,
239    right_pre: BoundedTransferModel<T>,
240    right_post: BoundedTransferModel<T>,
241    right_value: T,
242    right_retained_bytes: u64,
243    right_sequence: u64,
244) -> bool {
245    &&& connection_receive(
246        left_pre,
247        left_post,
248        left_value,
249        left_retained_bytes,
250        left_sequence,
251    )
252    &&& connection_receive(
253        right_pre,
254        right_post,
255        right_value,
256        right_retained_bytes,
257        right_sequence,
258    )
259}
260
261/// Close two connections together with one shared change result.
262pub open spec fn close_pair<T>(
263    left_pre: BoundedTransferModel<T>,
264    left_post: BoundedTransferModel<T>,
265    right_pre: BoundedTransferModel<T>,
266    right_post: BoundedTransferModel<T>,
267    changed: bool,
268) -> bool {
269    &&& connection_close(left_pre, left_post, changed)
270    &&& connection_close(right_pre, right_post, changed)
271}
272
273/// A refused two-connection action leaves both connections unchanged.
274pub open spec fn pair_refusal<T>(
275    left_pre: BoundedTransferModel<T>,
276    left_post: BoundedTransferModel<T>,
277    right_pre: BoundedTransferModel<T>,
278    right_post: BoundedTransferModel<T>,
279) -> bool {
280    &&& connection_refusal(left_pre, left_post)
281    &&& connection_refusal(right_pre, right_post)
282}
283
284/// Move the oldest value between adjacent connections as one StreamGraph relay action.
285pub open spec fn relay<T>(
286    source_pre: BoundedTransferModel<T>,
287    source_post: BoundedTransferModel<T>,
288    destination_pre: BoundedTransferModel<T>,
289    destination_post: BoundedTransferModel<T>,
290    destination_sequence: u64,
291) -> bool {
292    let value = source_pre.values[0];
293    let source_sequence = source_pre.registry[0].0;
294    let retained_bytes = source_pre.registry[0].1;
295    &&& connection_receive(
296        source_pre,
297        source_post,
298        value,
299        retained_bytes,
300        source_sequence,
301    )
302    &&& connection_publish(
303        destination_pre,
304        destination_post,
305        value,
306        retained_bytes,
307        destination_sequence,
308    )
309}
310
311/// A refused relay leaves both adjacent connections unchanged.
312pub open spec fn relay_refusal<T>(
313    source_pre: BoundedTransferModel<T>,
314    source_post: BoundedTransferModel<T>,
315    destination_pre: BoundedTransferModel<T>,
316    destination_post: BoundedTransferModel<T>,
317) -> bool {
318    &&& connection_refusal(source_pre, source_post)
319    &&& connection_refusal(destination_pre, destination_post)
320}
321
322/// Empty admitted origin for one StreamGraph connection.
323pub open spec fn connection_initial<T>(
324    model: BoundedTransferModel<T>,
325) -> bool {
326    bounded_transfer_initial(model)
327}
328
329/// Empty admitted origin for two StreamGraph connections.
330pub open spec fn pair_initial<T>(
331    left: BoundedTransferModel<T>,
332    right: BoundedTransferModel<T>,
333) -> bool {
334    &&& connection_initial(left)
335    &&& connection_initial(right)
336}
337
338/// Bounded linear stream owner.
339pub struct StreamGraph {
340    /// Number of stages in the linear chain.
341    pub chain_length: usize,
342    /// Maximum records admitted at the source.
343    pub max_inputs: usize,
344    /// Exclusive upper bound of record values.
345    pub record_domain_size: u64,
346    /// First FIFO edge owner.
347    pub q1: Buffer<u64>,
348    /// Second FIFO edge owner.
349    pub q2: Buffer<u64>,
350    /// Optional third FIFO edge owner.
351    pub q3: Buffer<u64>,
352    /// Source-admission counter.
353    pub ingested: Counter,
354    /// Sink-emission counter.
355    pub emitted: Counter,
356}
357
358impl StreamGraph {
359    /// Whether chain length, edge capacity, and record domain form a supported configuration.
360    pub open spec fn valid_config_spec(
361        chain_length: usize,
362        capacity: usize,
363        record_domain_size: u64,
364    ) -> bool {
365        (chain_length == 3 || chain_length == 4)
366            && capacity > 0
367            && record_domain_size > 0
368    }
369
370    /// Whether every queued value lies within `domain`.
371    pub open spec fn values_valid(q: Seq<u64>, domain: u64) -> bool {
372        forall|i: int| 0 <= i < q.len() ==> #[trigger] q[i] < domain
373    }
374
375    /// Total number of records currently retained across all stream edges.
376    pub open spec fn queue_depth(&self) -> nat {
377        if self.chain_length == 3 {
378            self.q1.values@.len() + self.q2.values@.len()
379        } else {
380            self.q1.values@.len() + self.q2.values@.len() + self.q3.values@.len()
381        }
382    }
383
384    /// Whether counters, queues, and configuration values have valid shape and bounds.
385    pub open spec fn type_invariant(&self) -> bool {
386        &&& Self::valid_config_spec(
387            self.chain_length, self.q1.capacity, self.record_domain_size)
388        &&& self.q2.capacity == self.q1.capacity
389        &&& self.q3.capacity == self.q1.capacity
390        &&& Self::values_valid(self.q1.values@, self.record_domain_size)
391        &&& Self::values_valid(self.q2.values@, self.record_domain_size)
392        &&& Self::values_valid(self.q3.values@, self.record_domain_size)
393        &&& (self.chain_length == 3 ==> self.q3.values@.len() == 0)
394        &&& self.ingested.value_spec() <= self.max_inputs as nat
395        &&& self.emitted.value_spec() <= self.max_inputs as nat
396    }
397
398    /// Whether a full downstream edge prevents the corresponding transfer.
399    pub open spec fn backpressure_correct(&self) -> bool {
400        &&& self.q1.well_formed()
401        &&& self.q2.well_formed()
402        &&& self.q3.well_formed()
403    }
404
405    /// Whether admitted-record count equals emitted plus retained-record counts.
406    pub open spec fn count_conservation(&self) -> bool {
407        self.ingested.value_spec() == self.queue_depth() + self.emitted.value_spec()
408    }
409
410    /// Compatibility alias for [`Self::count_conservation`].
411    ///
412    /// Count equality alone does not establish record identity or provenance.
413    pub open spec fn no_record_loss(&self) -> bool {
414        self.count_conservation()
415    }
416
417    /// Whether at least one modeled action is enabled in this state.
418    ///
419    /// This is a state predicate. It does not require a scheduler to choose an enabled action and
420    /// therefore does not establish temporal progress or fairness.
421    pub open spec fn some_action_enabled(&self) -> bool {
422        ||| (self.ingested.value_spec() < self.max_inputs as nat
423            && self.q1.values@.len() < self.q1.capacity)
424        ||| (self.q1.values@.len() > 0
425            && self.q2.values@.len() < self.q2.capacity)
426        ||| (self.chain_length == 3 && self.q2.values@.len() > 0)
427        ||| (self.chain_length == 4
428            && self.q2.values@.len() > 0
429            && self.q3.values@.len() < self.q3.capacity)
430        ||| (self.chain_length == 4 && self.q3.values@.len() > 0)
431        ||| (self.ingested.value_spec() == self.max_inputs as nat
432            && self.q1.values@.len() == 0
433            && self.q2.values@.len() == 0
434            && self.q3.values@.len() == 0)
435    }
436
437    /// Whether all linear-stream contract clauses hold.
438    pub open spec fn inv(&self) -> bool {
439        self.type_invariant() && self.backpressure_correct() && self.count_conservation()
440    }
441
442    /// Establish state-level enabledness from the retained carrier invariant.
443    pub proof fn lemma_some_action_enabled(&self)
444        requires self.inv(),
445        ensures self.some_action_enabled(),
446    {
447        reveal(StreamGraph::inv);
448        reveal(StreamGraph::type_invariant);
449        reveal(StreamGraph::backpressure_correct);
450        reveal(StreamGraph::some_action_enabled);
451
452        if self.ingested.value_spec() < self.max_inputs as nat {
453            if self.q1.values@.len() < self.q1.capacity {
454            } else if self.q2.values@.len() < self.q2.capacity {
455            } else if self.chain_length == 3 {
456            } else if self.q3.values@.len() < self.q3.capacity {
457            } else {
458            }
459        } else if self.q1.values@.len() > 0 {
460            if self.q2.values@.len() < self.q2.capacity {
461            } else if self.chain_length == 3 {
462            } else if self.q3.values@.len() < self.q3.capacity {
463            } else {
464            }
465        } else if self.q2.values@.len() > 0 {
466            if self.chain_length == 3 {
467            } else if self.q3.values@.len() < self.q3.capacity {
468            } else {
469            }
470        } else if self.chain_length == 4 && self.q3.values@.len() > 0 {
471        } else {
472        }
473    }
474
475    /// Evaluate the state-level enabledness predicate.
476    pub fn some_action_enabled_exec(&self) -> (enabled: bool)
477        requires self.inv(),
478        ensures enabled == self.some_action_enabled(),
479    {
480        proof { self.lemma_some_action_enabled(); }
481        (self.ingested.value() < self.max_inputs as u64 && self.q1.len() < self.q1.capacity)
482            || (!self.q1.is_empty() && self.q2.len() < self.q2.capacity)
483            || (self.chain_length == 3 && !self.q2.is_empty())
484            || (self.chain_length == 4
485                && !self.q2.is_empty()
486                && self.q3.len() < self.q3.capacity)
487            || (self.chain_length == 4 && !self.q3.is_empty())
488            || (self.ingested.value() == self.max_inputs as u64
489                && self.q1.is_empty()
490                && self.q2.is_empty()
491                && self.q3.is_empty())
492    }
493
494    /// Test whether a chain configuration is represented by this carrier.
495    pub fn valid_config(
496        chain_length: usize,
497        capacity: usize,
498        record_domain_size: u64,
499    ) -> (valid: bool)
500        ensures valid == Self::valid_config_spec(
501            chain_length, capacity, record_domain_size),
502    {
503        (chain_length == 3 || chain_length == 4)
504            && capacity > 0
505            && record_domain_size > 0
506    }
507
508    /// Capacity shared by every Buffer edge owner.
509    pub fn capacity(&self) -> (capacity: usize)
510        ensures capacity == self.q1.capacity,
511    {
512        self.q1.capacity
513    }
514
515    /// Construct an empty valid stream graph.
516    pub fn new(
517        chain_length: usize,
518        capacity: usize,
519        max_inputs: usize,
520        record_domain_size: u64,
521    ) -> (s: StreamGraph)
522        requires Self::valid_config_spec(chain_length, capacity, record_domain_size),
523        ensures
524            s.chain_length == chain_length,
525            s.q1.capacity == capacity,
526            s.q2.capacity == capacity,
527            s.q3.capacity == capacity,
528            s.max_inputs == max_inputs,
529            s.record_domain_size == record_domain_size,
530            s.q1.values@.len() == 0,
531            s.q2.values@.len() == 0,
532            s.q3.values@.len() == 0,
533            s.ingested.value_spec() == 0,
534            s.emitted.value_spec() == 0,
535            s.inv(),
536    {
537        StreamGraph {
538            chain_length,
539            max_inputs,
540            record_domain_size,
541            q1: Buffer::new(capacity),
542            q2: Buffer::new(capacity),
543            q3: Buffer::new(capacity),
544            ingested: Counter::new(0),
545            emitted: Counter::new(0),
546        }
547    }
548
549    /// Admit one in-domain record when source and backpressure bounds permit it.
550    pub fn source_ingest(&mut self, value: u64) -> (accepted: bool)
551        requires old(self).inv(),
552        ensures
553            accepted == (value < old(self).record_domain_size
554                && old(self).ingested.value_spec() < old(self).max_inputs as nat
555                && old(self).q1.values@.len() < old(self).q1.capacity),
556            final(self).chain_length == old(self).chain_length,
557            final(self).q1.capacity == old(self).q1.capacity,
558            final(self).q2.capacity == old(self).q2.capacity,
559            final(self).q3.capacity == old(self).q3.capacity,
560            final(self).max_inputs == old(self).max_inputs,
561            final(self).record_domain_size == old(self).record_domain_size,
562            final(self).q1.values@ == if accepted {
563                old(self).q1.values@.push(value)
564            } else { old(self).q1.values@ },
565            final(self).q2.values@ == old(self).q2.values@,
566            final(self).q3.values@ == old(self).q3.values@,
567            accepted ==> final(self).ingested.value_spec()
568                == old(self).ingested.value_spec() + 1,
569            !accepted ==> final(self).ingested.value_spec()
570                == old(self).ingested.value_spec(),
571            final(self).emitted.value_spec() == old(self).emitted.value_spec(),
572            final(self).inv(),
573    {
574        if value < self.record_domain_size
575            && self.ingested.value() < self.max_inputs as u64
576            && self.q1.len() < self.q1.capacity
577        {
578            let _pushed = self.q1.push(value);
579            let _counted = self.ingested.try_increment();
580            assert(_counted);
581            assert forall|i: int| 0 <= i < self.q1.values@.len()
582                implies #[trigger] self.q1.values@[i] < self.record_domain_size by {
583                if i < old(self).q1.values@.len() {
584                    assert(self.q1.values@[i] == old(self).q1.values@[i]);
585                }
586            }
587            true
588        } else {
589            false
590        }
591    }
592
593    /// Transfer the oldest first-edge record to the second edge.
594    pub fn middle2_fire(&mut self) -> (accepted: bool)
595        requires old(self).inv(),
596        ensures
597            accepted == (old(self).q1.values@.len() > 0
598                && old(self).q2.values@.len() < old(self).q2.capacity),
599            final(self).chain_length == old(self).chain_length,
600            final(self).q1.capacity == old(self).q1.capacity,
601            final(self).q2.capacity == old(self).q2.capacity,
602            final(self).q3.capacity == old(self).q3.capacity,
603            final(self).max_inputs == old(self).max_inputs,
604            final(self).record_domain_size == old(self).record_domain_size,
605            final(self).q1.values@ == if accepted {
606                old(self).q1.values@.subrange(1, old(self).q1.values@.len() as int)
607            } else { old(self).q1.values@ },
608            final(self).q2.values@ == if accepted {
609                old(self).q2.values@.push(old(self).q1.values@[0])
610            } else { old(self).q2.values@ },
611            final(self).q3.values@ == old(self).q3.values@,
612            final(self).ingested.value_spec() == old(self).ingested.value_spec(),
613            final(self).emitted.value_spec() == old(self).emitted.value_spec(),
614            final(self).inv(),
615    {
616        if self.q1.len() > 0 && self.q2.len() < self.q2.capacity {
617            let ghost old_q1 = self.q1.values@;
618            let ghost old_q2 = self.q2.values@;
619            let value = self.q1.values[0];
620            let _popped = self.q1.pop();
621            let _pushed = self.q2.push(value);
622            assert(self.q1.values@ =~= old_q1.subrange(1, old_q1.len() as int));
623            assert(self.q2.values@ =~= old_q2.push(old_q1[0]));
624            assert forall|i: int| 0 <= i < self.q2.values@.len()
625                implies #[trigger] self.q2.values@[i] < self.record_domain_size by {
626                if i < old_q2.len() {
627                    assert(self.q2.values@[i] == old_q2[i]);
628                } else {
629                    assert(self.q2.values@[i] == old_q1[0]);
630                }
631            }
632            true
633        } else {
634            false
635        }
636    }
637
638    /// Transfer the oldest second-edge record through the optional fourth stage.
639    pub fn middle3_fire(&mut self) -> (accepted: bool)
640        requires old(self).inv(),
641        ensures
642            accepted == (old(self).chain_length == 4
643                && old(self).q2.values@.len() > 0
644                && old(self).q3.values@.len() < old(self).q3.capacity),
645            final(self).chain_length == old(self).chain_length,
646            final(self).q1.capacity == old(self).q1.capacity,
647            final(self).q2.capacity == old(self).q2.capacity,
648            final(self).q3.capacity == old(self).q3.capacity,
649            final(self).max_inputs == old(self).max_inputs,
650            final(self).record_domain_size == old(self).record_domain_size,
651            final(self).q1.values@ == old(self).q1.values@,
652            final(self).q2.values@ == if accepted {
653                old(self).q2.values@.subrange(1, old(self).q2.values@.len() as int)
654            } else { old(self).q2.values@ },
655            final(self).q3.values@ == if accepted {
656                old(self).q3.values@.push(old(self).q2.values@[0])
657            } else { old(self).q3.values@ },
658            final(self).ingested.value_spec() == old(self).ingested.value_spec(),
659            final(self).emitted.value_spec() == old(self).emitted.value_spec(),
660            final(self).inv(),
661    {
662        if self.chain_length == 4 && self.q2.len() > 0
663            && self.q3.len() < self.q3.capacity
664        {
665            let ghost old_q2 = self.q2.values@;
666            let ghost old_q3 = self.q3.values@;
667            let value = self.q2.values[0];
668            let _popped = self.q2.pop();
669            let _pushed = self.q3.push(value);
670            assert(self.q2.values@ =~= old_q2.subrange(1, old_q2.len() as int));
671            assert(self.q3.values@ =~= old_q3.push(old_q2[0]));
672            assert forall|i: int| 0 <= i < self.q3.values@.len()
673                implies #[trigger] self.q3.values@[i] < self.record_domain_size by {
674                if i < old_q3.len() {
675                    assert(self.q3.values@[i] == old_q3[i]);
676                } else {
677                    assert(self.q3.values@[i] == old_q2[0]);
678                }
679            }
680            true
681        } else {
682            false
683        }
684    }
685
686    /// Consume the oldest record from the final edge.
687    pub fn sink_consume(&mut self) -> (accepted: bool)
688        requires old(self).inv(),
689        ensures
690            accepted == if old(self).chain_length == 3 {
691                old(self).q2.values@.len() > 0
692            } else {
693                old(self).q3.values@.len() > 0
694            },
695            final(self).chain_length == old(self).chain_length,
696            final(self).q1.capacity == old(self).q1.capacity,
697            final(self).q2.capacity == old(self).q2.capacity,
698            final(self).q3.capacity == old(self).q3.capacity,
699            final(self).max_inputs == old(self).max_inputs,
700            final(self).record_domain_size == old(self).record_domain_size,
701            final(self).q1.values@ == old(self).q1.values@,
702            final(self).q2.values@ == if accepted && old(self).chain_length == 3 {
703                old(self).q2.values@.subrange(1, old(self).q2.values@.len() as int)
704            } else { old(self).q2.values@ },
705            final(self).q3.values@ == if accepted && old(self).chain_length == 4 {
706                old(self).q3.values@.subrange(1, old(self).q3.values@.len() as int)
707            } else { old(self).q3.values@ },
708            final(self).ingested.value_spec() == old(self).ingested.value_spec(),
709            accepted ==> final(self).emitted.value_spec()
710                == old(self).emitted.value_spec() + 1,
711            !accepted ==> final(self).emitted.value_spec()
712                == old(self).emitted.value_spec(),
713            final(self).inv(),
714    {
715        if self.chain_length == 3 {
716            if self.q2.len() > 0 {
717                let ghost old_q2 = self.q2.values@;
718                let _popped = self.q2.pop();
719                let _counted = self.emitted.try_increment();
720                assert(_counted);
721                assert(self.q2.values@ =~= old_q2.subrange(1, old_q2.len() as int));
722                true
723            } else {
724                false
725            }
726        } else {
727            if self.q3.len() > 0 {
728                let ghost old_q3 = self.q3.values@;
729                let _popped = self.q3.pop();
730                let _counted = self.emitted.try_increment();
731                assert(_counted);
732                assert(self.q3.values@ =~= old_q3.subrange(1, old_q3.len() as int));
733                true
734            } else {
735                false
736            }
737        }
738    }
739
740    /// Execute the terminal stutter after bounded input is drained.
741    pub fn done_stuttering(&mut self) -> (enabled: bool)
742        requires old(self).inv(),
743        ensures
744            enabled == (old(self).ingested.value_spec() == old(self).max_inputs as nat
745                && old(self).q1.values@.len() == 0
746                && old(self).q2.values@.len() == 0
747                && old(self).q3.values@.len() == 0),
748            final(self).chain_length == old(self).chain_length,
749            final(self).q1.capacity == old(self).q1.capacity,
750            final(self).q2.capacity == old(self).q2.capacity,
751            final(self).q3.capacity == old(self).q3.capacity,
752            final(self).max_inputs == old(self).max_inputs,
753            final(self).record_domain_size == old(self).record_domain_size,
754            final(self).q1.values@ == old(self).q1.values@,
755            final(self).q2.values@ == old(self).q2.values@,
756            final(self).q3.values@ == old(self).q3.values@,
757            final(self).ingested.value_spec() == old(self).ingested.value_spec(),
758            final(self).emitted.value_spec() == old(self).emitted.value_spec(),
759            final(self).inv(),
760    {
761        self.ingested.value() == self.max_inputs as u64
762            && self.q1.len() == 0
763            && self.q2.len() == 0
764            && self.q3.len() == 0
765    }
766}
767
768}