use crate::modalities::fork_join::ForkJoin as ForkJoinCarrier;
use crate::modalities::sequential::Sequential as SequentialCarrier;
use crate::modalities::step_graph::StepGraph as StepGraphCarrier;
use crate::modalities::stream_graph::StreamGraph as StreamGraphCarrier;
use vstd::prelude::*;
verus! {
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum SequentialBuildError {
NoSteps,
EmptyValueDomain,
InitialValueOutOfRange,
}
pub struct Sequential {
inner: SequentialCarrier,
}
impl Sequential {
#[verifier::type_invariant]
closed spec fn well_formed(&self) -> bool { self.inner.inv() }
pub fn new(steps: usize, value_domain_size: u64, initial_value: u64)
-> (result: Result<Self, SequentialBuildError>) {
if steps == 0 { return Err(SequentialBuildError::NoSteps); }
if value_domain_size == 0 { return Err(SequentialBuildError::EmptyValueDomain); }
if initial_value >= value_domain_size {
return Err(SequentialBuildError::InitialValueOutOfRange);
}
Ok(Self { inner: SequentialCarrier::new(steps, value_domain_size, initial_value) })
}
pub fn steps(&self) -> usize { self.inner.steps }
pub fn completed(&self) -> usize { self.inner.pc }
pub fn value(&self) -> u64 { self.inner.value }
pub fn is_active(&self) -> bool { self.inner.active }
pub fn is_done(&self) -> bool { self.inner.pc == self.inner.steps && !self.inner.active }
#[expect(clippy::indexing_slicing, reason = "the branch proves the history index is in bounds")]
pub fn history(&self, index: usize) -> Option<u64> {
if index < self.inner.history.len() { Some(self.inner.history[index]) } else { None }
}
#[must_use]
pub fn begin_step(&mut self) -> (accepted: bool) {
proof { use_type_invariant(&*self); }
let mut carrier = sequential_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.begin_step();
core::mem::swap(&mut self.inner, &mut carrier);
accepted
}
#[must_use]
pub fn complete_step(&mut self, next_value: u64) -> (accepted: bool) {
proof { use_type_invariant(&*self); }
let mut carrier = sequential_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.complete_step(next_value);
core::mem::swap(&mut self.inner, &mut carrier);
accepted
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum ForkJoinBuildError {
EmptyValueDomain,
InitialValueOutOfRange,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum WorkerState {
Ready,
Running,
Complete,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum ForkJoinPhase {
Fork,
Join,
Done,
}
pub struct ForkJoin {
inner: ForkJoinCarrier,
}
impl ForkJoin {
#[verifier::type_invariant]
closed spec fn well_formed(&self) -> bool { self.inner.inv() }
pub fn new(workers: usize, value_domain_size: u64, initial_value: u64)
-> (result: Result<Self, ForkJoinBuildError>) {
if value_domain_size == 0 { return Err(ForkJoinBuildError::EmptyValueDomain); }
if initial_value >= value_domain_size {
return Err(ForkJoinBuildError::InitialValueOutOfRange);
}
Ok(Self { inner: ForkJoinCarrier::new(workers, value_domain_size, initial_value) })
}
pub fn len(&self) -> usize { self.inner.wstate.len() }
pub fn is_empty(&self) -> bool { self.inner.wstate.is_empty() }
pub fn phase(&self) -> ForkJoinPhase {
self.inner.phase
}
pub fn output_ready(&self) -> bool { self.inner.output_ready }
#[expect(clippy::indexing_slicing, reason = "the branch proves the worker index is in bounds")]
pub fn worker_state(&self, worker: usize) -> Option<WorkerState> {
if worker >= self.inner.wstate.len() { return None; }
Some(self.inner.wstate[worker])
}
#[expect(clippy::indexing_slicing, reason = "the branch proves the worker index is in bounds")]
pub fn worker_value(&self, worker: usize) -> Option<u64> {
if worker < self.inner.wvalue.len() { Some(self.inner.wvalue[worker]) } else { None }
}
#[expect(clippy::indexing_slicing, reason = "the branch proves the output index is in bounds")]
pub fn output(&self, worker: usize) -> Option<u64> {
if self.inner.output_ready && worker < self.inner.output_snapshot.len() {
Some(self.inner.output_snapshot[worker])
} else { None }
}
#[must_use]
pub fn start_worker(&mut self, worker: usize) -> (accepted: bool) {
proof { use_type_invariant(&*self); }
let mut carrier = fork_join_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.start_worker(worker);
core::mem::swap(&mut self.inner, &mut carrier);
accepted
}
#[must_use]
pub fn complete_worker(&mut self, worker: usize, value: u64) -> (accepted: bool) {
proof { use_type_invariant(&*self); }
let mut carrier = fork_join_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.complete_worker(worker, value);
core::mem::swap(&mut self.inner, &mut carrier);
accepted
}
#[must_use]
pub fn barrier(&mut self) -> (accepted: bool) {
proof { use_type_invariant(&*self); }
let mut carrier = fork_join_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.barrier();
core::mem::swap(&mut self.inner, &mut carrier);
accepted
}
#[must_use]
pub fn produce_output(&mut self) -> (accepted: bool) {
proof { use_type_invariant(&*self); }
let mut carrier = fork_join_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.produce_output();
core::mem::swap(&mut self.inner, &mut carrier);
accepted
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum StepGraphBuildError {
EdgeEndpointOutOfRange,
DuplicateEdge,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum StepState {
NotReady,
Ready,
Running,
Complete,
}
pub struct StepGraph {
inner: StepGraphCarrier,
}
impl StepGraph {
#[verifier::type_invariant]
closed spec fn well_formed(&self) -> bool { self.inner.inv() }
pub fn new(num_nodes: usize, edges: Vec<(usize, usize)>)
-> (result: Result<Self, StepGraphBuildError>) {
if !step_edges_valid(&edges, num_nodes) {
return Err(StepGraphBuildError::EdgeEndpointOutOfRange);
}
if !step_edges_distinct(&edges) { return Err(StepGraphBuildError::DuplicateEdge); }
Ok(Self { inner: StepGraphCarrier::new(num_nodes, edges) })
}
pub fn len(&self) -> usize { self.inner.num_nodes }
pub fn is_empty(&self) -> bool { self.inner.num_nodes == 0 }
pub fn edge_count(&self) -> usize { self.inner.edges.len() }
#[expect(clippy::indexing_slicing, reason = "the branch proves the edge index is in bounds")]
pub fn edge(&self, index: usize) -> Option<(usize, usize)> {
if index < self.inner.edges.len() { Some(self.inner.edges[index]) } else { None }
}
#[expect(clippy::indexing_slicing, reason = "the branch proves the node index is in bounds")]
pub fn state(&self, node: usize) -> Option<StepState> {
if node >= self.inner.nstate.len() { return None; }
Some(self.inner.nstate[node])
}
#[must_use]
pub fn become_ready(&mut self, node: usize) -> (accepted: bool) {
proof { use_type_invariant(&*self); }
let mut carrier = step_graph_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.become_ready(node);
core::mem::swap(&mut self.inner, &mut carrier);
accepted
}
#[must_use]
pub fn start(&mut self, node: usize) -> (accepted: bool) {
proof { use_type_invariant(&*self); }
let mut carrier = step_graph_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.start_running(node);
core::mem::swap(&mut self.inner, &mut carrier);
accepted
}
#[must_use]
pub fn complete(&mut self, node: usize) -> (accepted: bool) {
proof { use_type_invariant(&*self); }
let mut carrier = step_graph_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.complete_node(node);
core::mem::swap(&mut self.inner, &mut carrier);
accepted
}
#[expect(clippy::indexing_slicing, reason = "the loop proves the state index is in bounds")]
#[expect(clippy::arithmetic_side_effects, reason = "the loop proves the cursor remains within the vector")]
pub fn is_done(&self) -> bool {
let mut index = 0;
while index < self.inner.nstate.len()
invariant index <= self.inner.nstate.len(),
decreases self.inner.nstate.len() - index,
{
if !matches!(self.inner.nstate[index], StepState::Complete) { return false; }
index += 1;
}
true
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum StreamGraphBuildError {
UnsupportedChainLength,
ZeroCapacity,
EmptyRecordDomain,
}
pub struct StreamGraph {
inner: StreamGraphCarrier,
}
impl StreamGraph {
#[verifier::type_invariant]
closed spec fn well_formed(&self) -> bool { self.inner.inv() }
pub fn new(chain_length: usize, capacity: usize, max_inputs: usize, record_domain_size: u64)
-> (result: Result<Self, StreamGraphBuildError>) {
if chain_length != 3 && chain_length != 4 {
return Err(StreamGraphBuildError::UnsupportedChainLength);
}
if capacity == 0 { return Err(StreamGraphBuildError::ZeroCapacity); }
if record_domain_size == 0 { return Err(StreamGraphBuildError::EmptyRecordDomain); }
Ok(Self { inner: StreamGraphCarrier::new(
chain_length, capacity, max_inputs, record_domain_size,
) })
}
pub fn chain_length(&self) -> usize { self.inner.chain_length }
pub fn capacity(&self) -> usize { self.inner.capacity() }
pub fn ingested(&self) -> usize { self.inner.ingested.value() as usize }
pub fn emitted(&self) -> usize { self.inner.emitted.value() as usize }
pub fn first_queue_len(&self) -> usize { self.inner.q1.len() }
pub fn second_queue_len(&self) -> usize { self.inner.q2.len() }
pub fn third_queue_len(&self) -> usize { self.inner.q3.len() }
#[must_use]
pub fn ingest(&mut self, value: u64) -> (accepted: bool) {
proof { use_type_invariant(&*self); }
let mut carrier = stream_graph_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.source_ingest(value);
core::mem::swap(&mut self.inner, &mut carrier);
accepted
}
#[must_use]
pub fn advance_first(&mut self) -> (accepted: bool) {
proof { use_type_invariant(&*self); }
let mut carrier = stream_graph_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.middle2_fire();
core::mem::swap(&mut self.inner, &mut carrier);
accepted
}
#[must_use]
pub fn advance_second(&mut self) -> (accepted: bool) {
proof { use_type_invariant(&*self); }
let mut carrier = stream_graph_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.middle3_fire();
core::mem::swap(&mut self.inner, &mut carrier);
accepted
}
#[expect(clippy::indexing_slicing, reason = "the queue guards prove the sink head is present")]
pub fn consume(&mut self) -> (value: Option<u64>) {
proof { use_type_invariant(&*self); }
let value = if self.inner.chain_length == 3 {
if self.inner.q2.is_empty() { return None; }
self.inner.q2.values[0]
} else {
if self.inner.q3.is_empty() { return None; }
self.inner.q3.values[0]
};
let mut carrier = stream_graph_sentinel();
core::mem::swap(&mut self.inner, &mut carrier);
let accepted = carrier.sink_consume();
if !accepted {
core::mem::swap(&mut self.inner, &mut carrier);
return None;
}
core::mem::swap(&mut self.inner, &mut carrier);
Some(value)
}
pub fn is_done(&self) -> bool {
self.inner.ingested.value() == self.inner.max_inputs as u64
&& self.inner.q1.is_empty()
&& self.inner.q2.is_empty()
&& self.inner.q3.is_empty()
}
}
#[expect(clippy::indexing_slicing, reason = "the loop proves the edge index is in bounds")]
#[expect(clippy::arithmetic_side_effects, reason = "the loop proves the cursor remains within the vector")]
#[expect(clippy::ptr_arg, reason = "Verus sequence-view contracts are stated over Vec in this checked boundary")]
fn step_edges_valid(edges: &Vec<(usize, usize)>, num_nodes: usize) -> (valid: bool)
ensures valid == StepGraphCarrier::edges_valid(edges@, num_nodes),
{
let mut index = 0;
while index < edges.len()
invariant
index <= edges.len(),
forall|i: int| 0 <= i < index ==>
#[trigger] edges@[i].0 < num_nodes && edges@[i].1 < num_nodes,
decreases edges.len() - index,
{
if edges[index].0 >= num_nodes || edges[index].1 >= num_nodes {
assert(!StepGraphCarrier::edges_valid(edges@, num_nodes));
return false;
}
index = index + 1;
}
true
}
#[expect(clippy::indexing_slicing, reason = "the nested loops prove both edge indices are in bounds")]
#[expect(clippy::arithmetic_side_effects, reason = "the loops prove both cursors remain within the vector")]
#[expect(clippy::ptr_arg, reason = "Verus sequence-view contracts are stated over Vec in this checked boundary")]
fn step_edges_distinct(edges: &Vec<(usize, usize)>) -> (distinct: bool)
ensures distinct == StepGraphCarrier::edges_distinct(edges@),
{
let mut left = 0;
while left < edges.len()
invariant
left <= edges.len(),
forall|i: int, j: int| 0 <= i < left && 0 <= j < edges.len() && i != j
==> #[trigger] edges@[i] != #[trigger] edges@[j],
decreases edges.len() - left,
{
let mut right = left + 1;
while right < edges.len()
invariant
left < edges.len(),
left + 1 <= right <= edges.len(),
forall|j: int| left < j < right ==> edges@[left as int] != edges@[j],
decreases edges.len() - right,
{
if edges[left].0 == edges[right].0 && edges[left].1 == edges[right].1 {
assert(!StepGraphCarrier::edges_distinct(edges@));
return false;
}
right = right + 1;
}
left = left + 1;
}
true
}
fn sequential_sentinel() -> (carrier: SequentialCarrier)
ensures carrier.inv(),
{ SequentialCarrier::new(1, 1, 0) }
fn fork_join_sentinel() -> (carrier: ForkJoinCarrier)
ensures carrier.inv(),
{ ForkJoinCarrier::new(0, 1, 0) }
fn step_graph_sentinel() -> (carrier: StepGraphCarrier)
ensures carrier.inv(),
{
let edges: Vec<(usize, usize)> = Vec::new();
StepGraphCarrier::new(0, edges)
}
fn stream_graph_sentinel() -> (carrier: StreamGraphCarrier)
ensures carrier.inv(),
{ StreamGraphCarrier::new(3, 1, 0, 1) }
}
impl Sequential {
pub fn value_domain_size(&self) -> u64 {
self.inner.value_domain_size
}
pub fn history_values(&self) -> &[u64] {
self.inner.history.as_slice()
}
}
impl ForkJoin {
pub fn value_domain_size(&self) -> u64 {
self.inner.value_domain_size
}
pub fn worker_states(&self) -> &[WorkerState] {
self.inner.wstate.as_slice()
}
pub fn worker_values(&self) -> &[u64] {
self.inner.wvalue.as_slice()
}
pub fn outputs(&self) -> Option<&[u64]> {
self.inner
.output_ready
.then_some(self.inner.output_snapshot.as_slice())
}
}
impl StepGraph {
pub fn edges(&self) -> &[(usize, usize)] {
self.inner.edges.as_slice()
}
pub fn states(&self) -> &[StepState] {
self.inner.nstate.as_slice()
}
}
impl StreamGraph {
pub fn max_inputs(&self) -> usize {
self.inner.max_inputs
}
pub fn record_domain_size(&self) -> u64 {
self.inner.record_domain_size
}
pub fn first_queue(&self) -> &[u64] {
self.inner.q1.values.as_slice()
}
pub fn second_queue(&self) -> &[u64] {
self.inner.q2.values.as_slice()
}
pub fn third_queue(&self) -> &[u64] {
self.inner.q3.values.as_slice()
}
}
impl_observational_debug!(Sequential, "Sequential",
"steps" => steps,
"completed" => completed,
"value" => value,
"active" => is_active,
"done" => is_done,
);
impl_observational_debug!(ForkJoin, "ForkJoin",
"len" => len,
"phase" => phase,
"output_ready" => output_ready,
);
impl_observational_debug!(StepGraph, "StepGraph",
"len" => len,
"edge_count" => edge_count,
"done" => is_done,
);
impl_observational_debug!(StreamGraph, "StreamGraph",
"chain_length" => chain_length,
"capacity" => capacity,
"ingested" => ingested,
"emitted" => emitted,
"first_queue_len" => first_queue_len,
"second_queue_len" => second_queue_len,
"third_queue_len" => third_queue_len,
"done" => is_done,
);
impl_public_error!(SequentialBuildError, {
Self::NoSteps => "sequential execution requires at least one step",
Self::EmptyValueDomain => "sequential value domain is empty",
Self::InitialValueOutOfRange => "initial sequential value is outside its domain",
});
impl_public_error!(ForkJoinBuildError, {
Self::EmptyValueDomain => "fork-join value domain is empty",
Self::InitialValueOutOfRange => "initial worker value is outside its domain",
});
impl_public_error!(StepGraphBuildError, {
Self::EdgeEndpointOutOfRange => "step-graph edge endpoint is outside the node universe",
Self::DuplicateEdge => "step graph contains a duplicate edge",
});
impl_public_error!(StreamGraphBuildError, {
Self::UnsupportedChainLength => "stream graph supports only three- or four-stage chains",
Self::ZeroCapacity => "stream-graph queue capacity must be positive",
Self::EmptyRecordDomain => "stream-graph record domain is empty",
});