Skip to main content

sim_expr_tree_calc/calc/
watch.rs

1use std::sync::{
2    Arc,
3    atomic::{AtomicU64, Ordering},
4};
5
6use sim_kernel::{Expr, Symbol};
7use sim_lib_stream_core::{
8    BufferPolicy, PushResult, StreamDirection, StreamItem, StreamMedia, StreamMetadata,
9    StreamPacket, StreamValue, stream_cancel_bang, stream_next_bang,
10};
11
12/// A bounded standard stream endpoint carrying expression-tree progress and
13/// change packets.
14#[derive(Clone)]
15pub struct CalcWatch {
16    stream: Arc<StreamValue>,
17    overflow_evidence: Arc<AtomicU64>,
18}
19
20impl CalcWatch {
21    pub(super) fn new(id: u64, policy: BufferPolicy) -> Self {
22        let metadata = StreamMetadata::new(
23            Symbol::qualified("expr-tree/watch", id.to_string()),
24            StreamMedia::Data,
25            StreamDirection::Source,
26            Symbol::qualified("clock", "control"),
27            policy,
28        );
29        Self {
30            stream: Arc::new(StreamValue::push(metadata)),
31            overflow_evidence: Arc::new(AtomicU64::new(0)),
32        }
33    }
34
35    /// Returns the existing standard stream value.
36    #[must_use]
37    pub fn stream(&self) -> &Arc<StreamValue> {
38        &self.stream
39    }
40
41    /// Consumes the next packet through the ordinary stream operation.
42    pub fn next(&self) -> sim_kernel::Result<Option<StreamItem>> {
43        stream_next_bang(&self.stream)
44    }
45
46    /// Cancels the endpoint through the ordinary stream operation.
47    pub fn cancel(&self) -> sim_kernel::Result<()> {
48        stream_cancel_bang(&self.stream)
49    }
50
51    /// Returns explicit lifetime overflow evidence for this endpoint.
52    #[must_use]
53    pub fn overflow_evidence(&self) -> u64 {
54        self.overflow_evidence.load(Ordering::Acquire)
55    }
56
57    pub(super) fn emit(&self, kind: &'static str, fields: Vec<(Expr, Expr)>) {
58        let overflow_before = self.overflow_evidence();
59        let mut payload = vec![
60            (
61                Expr::Symbol(Symbol::new("kind")),
62                Expr::Symbol(Symbol::qualified("expr-tree", kind)),
63            ),
64            (
65                Expr::Symbol(Symbol::new("prior-overflows")),
66                Expr::String(overflow_before.to_string()),
67            ),
68        ];
69        payload.extend(fields);
70        let item = StreamItem::new(StreamPacket::data(
71            Symbol::qualified("expr-tree", kind),
72            Expr::Map(payload),
73        ));
74        let overflowed = matches!(
75            self.stream.push_packet(item),
76            Ok(PushResult::DroppedNewest(_))
77                | Ok(PushResult::DroppedOldest(_))
78                | Ok(PushResult::Rejected(_))
79                | Ok(PushResult::Closed(_))
80                | Err(_)
81        );
82        if overflowed {
83            self.overflow_evidence.fetch_add(1, Ordering::AcqRel);
84        }
85    }
86}