1use vstd::prelude::*;
6
7use crate::connectives::buffer::Buffer;
8use crate::connectives::counter::Counter;
9
10verus! {
11
12pub ghost struct BoundedTransferModel<T> {
14 pub slot_capacity: nat,
16 pub retained_byte_capacity: nat,
18 pub values: Seq<T>,
20 pub registry: Seq<(u64, u64)>,
22 pub retained_bytes: nat,
24 pub head: nat,
26 pub tail: nat,
28 pub closed: bool,
30}
31
32impl<T> BoundedTransferModel<T> {
33 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 pub open spec fn ready(self) -> bool {
67 self.values.len() > 0
68 }
69
70 pub open spec fn terminal(self) -> bool {
72 self.closed && self.values.len() == 0
73 }
74}
75
76pub 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
101pub 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
124pub 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
143pub open spec fn bounded_transfer_refusal_stutters<T>(
145 pre: BoundedTransferModel<T>,
146 post: BoundedTransferModel<T>,
147) -> bool {
148 pre == post
149}
150
151pub 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
165pub 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
176pub 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
187pub 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
196pub 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
204pub 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
232pub 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
261pub 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
273pub 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
284pub 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
311pub 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
322pub open spec fn connection_initial<T>(
324 model: BoundedTransferModel<T>,
325) -> bool {
326 bounded_transfer_initial(model)
327}
328
329pub 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
338pub struct StreamGraph {
340 pub chain_length: usize,
342 pub max_inputs: usize,
344 pub record_domain_size: u64,
346 pub q1: Buffer<u64>,
348 pub q2: Buffer<u64>,
350 pub q3: Buffer<u64>,
352 pub ingested: Counter,
354 pub emitted: Counter,
356}
357
358impl StreamGraph {
359 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 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 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 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 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 pub open spec fn no_record_loss(&self) -> bool {
407 self.ingested.value_spec() == self.queue_depth() + self.emitted.value_spec()
408 }
409
410 pub open spec fn inv(&self) -> bool {
412 self.type_invariant() && self.backpressure_correct() && self.no_record_loss()
413 }
414
415 pub fn valid_config(
417 chain_length: usize,
418 capacity: usize,
419 record_domain_size: u64,
420 ) -> (valid: bool)
421 ensures valid == Self::valid_config_spec(
422 chain_length, capacity, record_domain_size),
423 {
424 (chain_length == 3 || chain_length == 4)
425 && capacity > 0
426 && record_domain_size > 0
427 }
428
429 pub fn capacity(&self) -> (capacity: usize)
431 ensures capacity == self.q1.capacity,
432 {
433 self.q1.capacity
434 }
435
436 pub fn new(
438 chain_length: usize,
439 capacity: usize,
440 max_inputs: usize,
441 record_domain_size: u64,
442 ) -> (s: StreamGraph)
443 requires Self::valid_config_spec(chain_length, capacity, record_domain_size),
444 ensures
445 s.chain_length == chain_length,
446 s.q1.capacity == capacity,
447 s.q2.capacity == capacity,
448 s.q3.capacity == capacity,
449 s.max_inputs == max_inputs,
450 s.record_domain_size == record_domain_size,
451 s.q1.values@.len() == 0,
452 s.q2.values@.len() == 0,
453 s.q3.values@.len() == 0,
454 s.ingested.value_spec() == 0,
455 s.emitted.value_spec() == 0,
456 s.inv(),
457 {
458 StreamGraph {
459 chain_length,
460 max_inputs,
461 record_domain_size,
462 q1: Buffer::new(capacity),
463 q2: Buffer::new(capacity),
464 q3: Buffer::new(capacity),
465 ingested: Counter::new(0),
466 emitted: Counter::new(0),
467 }
468 }
469
470 pub fn source_ingest(&mut self, value: u64) -> (accepted: bool)
472 requires old(self).inv(),
473 ensures
474 accepted == (value < old(self).record_domain_size
475 && old(self).ingested.value_spec() < old(self).max_inputs as nat
476 && old(self).q1.values@.len() < old(self).q1.capacity),
477 final(self).chain_length == old(self).chain_length,
478 final(self).q1.capacity == old(self).q1.capacity,
479 final(self).q2.capacity == old(self).q2.capacity,
480 final(self).q3.capacity == old(self).q3.capacity,
481 final(self).max_inputs == old(self).max_inputs,
482 final(self).record_domain_size == old(self).record_domain_size,
483 final(self).q1.values@ == if accepted {
484 old(self).q1.values@.push(value)
485 } else { old(self).q1.values@ },
486 final(self).q2.values@ == old(self).q2.values@,
487 final(self).q3.values@ == old(self).q3.values@,
488 accepted ==> final(self).ingested.value_spec()
489 == old(self).ingested.value_spec() + 1,
490 !accepted ==> final(self).ingested.value_spec()
491 == old(self).ingested.value_spec(),
492 final(self).emitted.value_spec() == old(self).emitted.value_spec(),
493 final(self).inv(),
494 {
495 if value < self.record_domain_size
496 && self.ingested.value() < self.max_inputs as u64
497 && self.q1.len() < self.q1.capacity
498 {
499 let _pushed = self.q1.push(value);
500 let _counted = self.ingested.try_increment();
501 assert(_counted);
502 assert forall|i: int| 0 <= i < self.q1.values@.len()
503 implies #[trigger] self.q1.values@[i] < self.record_domain_size by {
504 if i < old(self).q1.values@.len() {
505 assert(self.q1.values@[i] == old(self).q1.values@[i]);
506 }
507 }
508 true
509 } else {
510 false
511 }
512 }
513
514 pub fn middle2_fire(&mut self) -> (accepted: bool)
516 requires old(self).inv(),
517 ensures
518 accepted == (old(self).q1.values@.len() > 0
519 && old(self).q2.values@.len() < old(self).q2.capacity),
520 final(self).chain_length == old(self).chain_length,
521 final(self).q1.capacity == old(self).q1.capacity,
522 final(self).q2.capacity == old(self).q2.capacity,
523 final(self).q3.capacity == old(self).q3.capacity,
524 final(self).max_inputs == old(self).max_inputs,
525 final(self).record_domain_size == old(self).record_domain_size,
526 final(self).q1.values@ == if accepted {
527 old(self).q1.values@.subrange(1, old(self).q1.values@.len() as int)
528 } else { old(self).q1.values@ },
529 final(self).q2.values@ == if accepted {
530 old(self).q2.values@.push(old(self).q1.values@[0])
531 } else { old(self).q2.values@ },
532 final(self).q3.values@ == old(self).q3.values@,
533 final(self).ingested.value_spec() == old(self).ingested.value_spec(),
534 final(self).emitted.value_spec() == old(self).emitted.value_spec(),
535 final(self).inv(),
536 {
537 if self.q1.len() > 0 && self.q2.len() < self.q2.capacity {
538 let ghost old_q1 = self.q1.values@;
539 let ghost old_q2 = self.q2.values@;
540 let value = self.q1.values[0];
541 let _popped = self.q1.pop();
542 let _pushed = self.q2.push(value);
543 assert(self.q1.values@ =~= old_q1.subrange(1, old_q1.len() as int));
544 assert(self.q2.values@ =~= old_q2.push(old_q1[0]));
545 assert forall|i: int| 0 <= i < self.q2.values@.len()
546 implies #[trigger] self.q2.values@[i] < self.record_domain_size by {
547 if i < old_q2.len() {
548 assert(self.q2.values@[i] == old_q2[i]);
549 } else {
550 assert(self.q2.values@[i] == old_q1[0]);
551 }
552 }
553 true
554 } else {
555 false
556 }
557 }
558
559 pub fn middle3_fire(&mut self) -> (accepted: bool)
561 requires old(self).inv(),
562 ensures
563 accepted == (old(self).chain_length == 4
564 && old(self).q2.values@.len() > 0
565 && old(self).q3.values@.len() < old(self).q3.capacity),
566 final(self).chain_length == old(self).chain_length,
567 final(self).q1.capacity == old(self).q1.capacity,
568 final(self).q2.capacity == old(self).q2.capacity,
569 final(self).q3.capacity == old(self).q3.capacity,
570 final(self).max_inputs == old(self).max_inputs,
571 final(self).record_domain_size == old(self).record_domain_size,
572 final(self).q1.values@ == old(self).q1.values@,
573 final(self).q2.values@ == if accepted {
574 old(self).q2.values@.subrange(1, old(self).q2.values@.len() as int)
575 } else { old(self).q2.values@ },
576 final(self).q3.values@ == if accepted {
577 old(self).q3.values@.push(old(self).q2.values@[0])
578 } else { old(self).q3.values@ },
579 final(self).ingested.value_spec() == old(self).ingested.value_spec(),
580 final(self).emitted.value_spec() == old(self).emitted.value_spec(),
581 final(self).inv(),
582 {
583 if self.chain_length == 4 && self.q2.len() > 0
584 && self.q3.len() < self.q3.capacity
585 {
586 let ghost old_q2 = self.q2.values@;
587 let ghost old_q3 = self.q3.values@;
588 let value = self.q2.values[0];
589 let _popped = self.q2.pop();
590 let _pushed = self.q3.push(value);
591 assert(self.q2.values@ =~= old_q2.subrange(1, old_q2.len() as int));
592 assert(self.q3.values@ =~= old_q3.push(old_q2[0]));
593 assert forall|i: int| 0 <= i < self.q3.values@.len()
594 implies #[trigger] self.q3.values@[i] < self.record_domain_size by {
595 if i < old_q3.len() {
596 assert(self.q3.values@[i] == old_q3[i]);
597 } else {
598 assert(self.q3.values@[i] == old_q2[0]);
599 }
600 }
601 true
602 } else {
603 false
604 }
605 }
606
607 pub fn sink_consume(&mut self) -> (accepted: bool)
609 requires old(self).inv(),
610 ensures
611 accepted == if old(self).chain_length == 3 {
612 old(self).q2.values@.len() > 0
613 } else {
614 old(self).q3.values@.len() > 0
615 },
616 final(self).chain_length == old(self).chain_length,
617 final(self).q1.capacity == old(self).q1.capacity,
618 final(self).q2.capacity == old(self).q2.capacity,
619 final(self).q3.capacity == old(self).q3.capacity,
620 final(self).max_inputs == old(self).max_inputs,
621 final(self).record_domain_size == old(self).record_domain_size,
622 final(self).q1.values@ == old(self).q1.values@,
623 final(self).q2.values@ == if accepted && old(self).chain_length == 3 {
624 old(self).q2.values@.subrange(1, old(self).q2.values@.len() as int)
625 } else { old(self).q2.values@ },
626 final(self).q3.values@ == if accepted && old(self).chain_length == 4 {
627 old(self).q3.values@.subrange(1, old(self).q3.values@.len() as int)
628 } else { old(self).q3.values@ },
629 final(self).ingested.value_spec() == old(self).ingested.value_spec(),
630 accepted ==> final(self).emitted.value_spec()
631 == old(self).emitted.value_spec() + 1,
632 !accepted ==> final(self).emitted.value_spec()
633 == old(self).emitted.value_spec(),
634 final(self).inv(),
635 {
636 if self.chain_length == 3 {
637 if self.q2.len() > 0 {
638 let ghost old_q2 = self.q2.values@;
639 let _popped = self.q2.pop();
640 let _counted = self.emitted.try_increment();
641 assert(_counted);
642 assert(self.q2.values@ =~= old_q2.subrange(1, old_q2.len() as int));
643 true
644 } else {
645 false
646 }
647 } else {
648 if self.q3.len() > 0 {
649 let ghost old_q3 = self.q3.values@;
650 let _popped = self.q3.pop();
651 let _counted = self.emitted.try_increment();
652 assert(_counted);
653 assert(self.q3.values@ =~= old_q3.subrange(1, old_q3.len() as int));
654 true
655 } else {
656 false
657 }
658 }
659 }
660
661 pub fn done_stuttering(&mut self) -> (enabled: bool)
663 requires old(self).inv(),
664 ensures
665 enabled == (old(self).ingested.value_spec() == old(self).max_inputs as nat
666 && old(self).q1.values@.len() == 0
667 && old(self).q2.values@.len() == 0
668 && old(self).q3.values@.len() == 0),
669 final(self).chain_length == old(self).chain_length,
670 final(self).q1.capacity == old(self).q1.capacity,
671 final(self).q2.capacity == old(self).q2.capacity,
672 final(self).q3.capacity == old(self).q3.capacity,
673 final(self).max_inputs == old(self).max_inputs,
674 final(self).record_domain_size == old(self).record_domain_size,
675 final(self).q1.values@ == old(self).q1.values@,
676 final(self).q2.values@ == old(self).q2.values@,
677 final(self).q3.values@ == old(self).q3.values@,
678 final(self).ingested.value_spec() == old(self).ingested.value_spec(),
679 final(self).emitted.value_spec() == old(self).emitted.value_spec(),
680 final(self).inv(),
681 {
682 self.ingested.value() == self.max_inputs as u64
683 && self.q1.len() == 0
684 && self.q2.len() == 0
685 && self.q3.len() == 0
686 }
687}
688
689}