Skip to main content

hydro_lang/location/
tick.rs

1//! Clock domains for batching streaming data into discrete time steps.
2//!
3//! In Hydro, a [`Tick`] represents a logical clock that can be used to batch
4//! unbounded streaming data into discrete, bounded time steps. This is essential
5//! for implementing iterative algorithms, synchronizing data across multiple
6//! streams, and performing aggregations over windows of data.
7//!
8//! A tick is created from a top-level location (such as [`super::Process`] or [`super::Cluster`])
9//! using [`Location::tick`]. Once inside a tick, bounded live collections can be
10//! manipulated with operations like fold, reduce, and cross-product, and the
11//! results can be emitted back to the unbounded stream using methods like
12//! `all_ticks()`.
13//!
14//! The [`Atomic`] wrapper provides atomicity guarantees within a tick, ensuring
15//! that reads and writes within a tick are serialized.
16
17use stageleft::{QuotedWithContext, q};
18
19#[cfg(stageleft_runtime)]
20use super::dynamic::DynLocation;
21use super::{Location, LocationId};
22use crate::compile::builder::{ClockId, FlowState};
23use crate::compile::ir::{HydroNode, HydroSource};
24#[cfg(stageleft_runtime)]
25use crate::forward_handle::{CycleCollection, CycleCollectionWithInitial};
26use crate::forward_handle::{TickCycle, TickCycleHandle};
27#[cfg(feature = "tokio")]
28use crate::live_collections::Singleton;
29use crate::live_collections::boundedness::Bounded;
30use crate::live_collections::optional::Optional;
31use crate::live_collections::stream::{ExactlyOnce, Stream, TotalOrder};
32use crate::location::TopLevel;
33#[cfg(feature = "tokio")]
34use crate::nondet::NonDet;
35use crate::nondet::nondet;
36
37/// A location wrapper that provides atomicity guarantees within a [`Tick`].
38///
39/// An `Atomic` context establishes a happens-before relationship between operations:
40/// - Downstream computations from `atomic()` are associated with an internal tick
41/// - Outputs from `end_atomic()` are held until all computations in the tick complete
42/// - Snapshots via `use::atomic` are guaranteed to reflect all updates from associated `end_atomic()`
43///
44/// This ensures read-after-write consistency: if a client receives an acknowledgement
45/// from `end_atomic()`, any subsequent `use::atomic` snapshot will include the effects
46/// of that acknowledged operation.
47#[derive(Clone)]
48pub struct Atomic<Loc> {
49    pub(crate) tick: Tick<Loc>,
50}
51
52impl<L: DynLocation> DynLocation for Atomic<L> {
53    fn dyn_id(&self) -> LocationId {
54        LocationId::Atomic(Box::new(self.tick.dyn_id()))
55    }
56
57    fn flow_state(&self) -> &FlowState {
58        self.tick.flow_state()
59    }
60
61    fn is_top_level() -> bool {
62        L::is_top_level()
63    }
64
65    fn multiversioned(&self) -> bool {
66        self.tick.multiversioned()
67    }
68
69    fn cluster_consistency() -> Option<super::dynamic::ClusterConsistency> {
70        L::cluster_consistency()
71    }
72}
73
74impl<'a, L> Location<'a> for Atomic<L>
75where
76    L: Location<'a>,
77{
78    type Root = L::Root;
79
80    type SimHookScope = L::SimHookScope;
81
82    type DropConsistency = Atomic<L::DropConsistency>;
83
84    fn consistency() -> Option<super::dynamic::ClusterConsistency> {
85        L::consistency()
86    }
87
88    fn root(&self) -> Self::Root {
89        self.tick.root()
90    }
91
92    fn drop_consistency(&self) -> Self::DropConsistency {
93        Atomic {
94            tick: self.tick.drop_consistency(),
95        }
96    }
97
98    fn from_drop_consistency(l2: Self::DropConsistency) -> Self {
99        Atomic {
100            tick: Tick::from_drop_consistency(l2.tick),
101        }
102    }
103}
104
105/// Trait for live collections that can be deferred by one tick.
106///
107/// When a collection implements `DeferTick`, calling `defer_tick` delays its
108/// values by one clock cycle. This is primarily used internally to implement
109/// tick-based cycles ([`Tick::cycle`]), ensuring that feedback loops advance
110/// by one tick to avoid infinite recursion within a single tick.
111pub trait DeferTick {
112    /// Returns a new collection whose values are delayed by one tick.
113    fn defer_tick(self) -> Self;
114}
115
116/// Marks the stream as being inside the single global clock domain.
117#[derive(Clone)]
118pub struct Tick<L> {
119    /// `None` if `l` is `Atomic`.
120    pub(crate) id: Option<ClockId>,
121    /// Location.
122    pub(crate) l: L,
123}
124
125impl<L: DynLocation> DynLocation for Tick<L> {
126    fn dyn_id(&self) -> LocationId {
127        LocationId::Tick {
128            tick: self.id,
129            parent_location: Box::new(self.l.dyn_id()),
130        }
131    }
132
133    fn flow_state(&self) -> &FlowState {
134        self.l.flow_state()
135    }
136
137    fn is_top_level() -> bool {
138        false
139    }
140
141    fn multiversioned(&self) -> bool {
142        self.l.multiversioned()
143    }
144
145    fn cluster_consistency() -> Option<super::dynamic::ClusterConsistency> {
146        L::cluster_consistency()
147    }
148}
149
150impl<'a, L> Location<'a> for Tick<L>
151where
152    L: Location<'a>,
153{
154    type Root = L::Root;
155
156    type SimHookScope = L::SimHookScope;
157
158    type DropConsistency = Tick<L::DropConsistency>;
159
160    fn consistency() -> Option<super::dynamic::ClusterConsistency> {
161        L::consistency()
162    }
163
164    fn root(&self) -> Self::Root {
165        self.l.root()
166    }
167
168    fn drop_consistency(&self) -> Self::DropConsistency {
169        Tick {
170            id: self.id,
171            l: self.l.drop_consistency(),
172        }
173    }
174
175    fn from_drop_consistency(l2: Self::DropConsistency) -> Self {
176        Tick {
177            id: l2.id,
178            l: L::from_drop_consistency(l2.l),
179        }
180    }
181}
182
183impl<'a, L> Tick<L>
184where
185    L: Location<'a>,
186{
187    /// Returns a reference to the parent location that this tick is located at.
188    ///
189    /// For example, if a `Tick` was created from a `Process`, this returns a reference
190    /// to that `Process`.
191    pub fn parent_location(&self) -> &L {
192        &self.l
193    }
194
195    /// Use [`Self::parent_location`] instead.
196    #[deprecated(note = "use `.parent_location()` instead")]
197    pub fn outer(&self) -> &L {
198        self.parent_location()
199    }
200
201    /// Creates a bounded stream of `()` values inside this tick, with a fixed batch size.
202    ///
203    /// This is useful for driving computations inside a tick that need to process
204    /// a specific number of elements per tick. Each tick will produce exactly
205    /// `batch_size` unit values.
206    pub fn spin_batch(
207        &self,
208        batch_size: impl QuotedWithContext<
209            'a,
210            usize,
211            crate::live_collections::OperatorContext<
212                L,
213                crate::live_collections::boundedness::Unbounded,
214            >,
215        > + Copy
216        + 'a,
217    ) -> Stream<(), Self, Bounded, TotalOrder, ExactlyOnce>
218    where
219        L: TopLevel<'a>,
220    {
221        let out = self
222            .l
223            .spin()
224            .flat_map_ordered(q!(move |_| 0..batch_size))
225            .map(q!(|_| ()));
226
227        let inner = out.batch(self, nondet!(/** at runtime, `spin` produces a single value per tick, so each batch is guaranteed to be the same size. */));
228        Stream::new(self.clone(), inner.ir_node.replace(HydroNode::Placeholder))
229    }
230
231    /// Creates an [`Optional`] which has a null value on every tick.
232    ///
233    /// # Example
234    /// ```rust
235    /// # #[cfg(feature = "deploy")] {
236    /// # use hydro_lang::prelude::*;
237    /// # use futures::StreamExt;
238    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
239    /// let tick = process.tick();
240    /// let optional = tick.none::<i32>();
241    /// optional.unwrap_or(tick.singleton(q!(123)))
242    /// # .all_ticks()
243    /// # }, |mut stream| async move {
244    /// // 123
245    /// # assert_eq!(stream.next().await.unwrap(), 123);
246    /// # }));
247    /// # }
248    /// ```
249    pub fn none<T>(&self) -> Optional<T, Self, Bounded> {
250        let e = q!([]);
251        let e = QuotedWithContext::<'a, [(); 0], Self>::splice_typed_ctx(e, self);
252
253        let unit_optional: Optional<(), Self, Bounded> = Optional::new(
254            self.clone(),
255            HydroNode::Source {
256                source: HydroSource::Iter(e.into()),
257                metadata: self.new_node_metadata(Optional::<(), Self, Bounded>::collection_kind()),
258            },
259        );
260
261        unit_optional.map(q!(|_| unreachable!())) // always empty
262    }
263
264    /// Creates an [`Optional`] which will have the provided static value on the first tick, and be
265    /// null on all subsequent ticks.
266    ///
267    /// This is useful for bootstrapping stateful computations which need an initial value.
268    ///
269    /// # Example
270    /// ```rust
271    /// # #[cfg(feature = "deploy")] {
272    /// # use hydro_lang::prelude::*;
273    /// # use futures::StreamExt;
274    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
275    /// let tick = process.tick();
276    /// // ticks are lazy by default, forces the second tick to run
277    /// tick.spin_batch(q!(1)).all_ticks().for_each(q!(|_| {}));
278    /// let optional = tick.optional_first_tick(q!(5));
279    /// optional.unwrap_or(tick.singleton(q!(123))).all_ticks()
280    /// # }, |mut stream| async move {
281    /// // 5, 123, 123, 123, ...
282    /// # assert_eq!(stream.next().await.unwrap(), 5);
283    /// # assert_eq!(stream.next().await.unwrap(), 123);
284    /// # assert_eq!(stream.next().await.unwrap(), 123);
285    /// # assert_eq!(stream.next().await.unwrap(), 123);
286    /// # }));
287    /// # }
288    /// ```
289    pub fn optional_first_tick<T: Clone>(
290        &self,
291        e: impl QuotedWithContext<'a, T, Tick<L>>,
292    ) -> Optional<T, Self, Bounded> {
293        let e = e.splice_untyped_ctx(self);
294
295        Optional::new(
296            self.clone(),
297            HydroNode::SingletonSource {
298                value: e.into(),
299                first_tick_only: true,
300                metadata: self.new_node_metadata(Optional::<T, Self, Bounded>::collection_kind()),
301            },
302        )
303    }
304
305    /// Returns the current wall-clock time as a [`Singleton`] containing a
306    /// [`tokio::time::Instant`].
307    ///
308    /// # Non-Determinism
309    /// Reading wall-clock time is inherently non-deterministic because the
310    /// value depends on when the tick executes. A [`NonDet`] guard is required
311    /// to acknowledge this.
312    #[cfg(feature = "tokio")]
313    pub fn current_tick_instant(
314        &self,
315        _nondet: NonDet,
316    ) -> Singleton<tokio::time::Instant, Tick<L::DropConsistency>, Bounded>
317    where
318        Self: Sized,
319    {
320        // TODO(shadaj): this is a simulator hole, should be reported as unsupported until it is
321        self.singleton(q!(tokio::time::Instant::now()))
322    }
323
324    /// Creates a feedback cycle within this tick for implementing iterative computations.
325    ///
326    /// Returns a handle that must be completed with the actual collection, and a placeholder
327    /// collection that represents the output of the previous tick (deferred by one tick).
328    /// This is useful for implementing fixed-point computations where the output of one
329    /// tick feeds into the input of the next.
330    ///
331    /// The cycle automatically defers values by one tick to prevent infinite recursion.
332    #[expect(
333        private_bounds,
334        reason = "only Hydro collections can implement ReceiverComplete"
335    )]
336    pub fn cycle<S, L2: Location<'a, DropConsistency = Tick<L::DropConsistency>>>(
337        &self,
338    ) -> (TickCycleHandle<'a, S>, S)
339    where
340        S: CycleCollection<'a, TickCycle, Location = L2> + DeferTick,
341    {
342        let cycle_id = self.flow_state().borrow_mut().next_cycle_id();
343        (
344            TickCycleHandle::new(cycle_id, Location::id(self)),
345            S::create_source(cycle_id, self.clone().with_consistency_of()).defer_tick(),
346        )
347    }
348
349    /// Creates a feedback cycle with an initial value for the first tick.
350    ///
351    /// Similar to [`Tick::cycle`], but allows providing an initial collection
352    /// that will be used as the value on the first tick before any feedback
353    /// is available. This is useful for bootstrapping iterative computations
354    /// that need a starting state.
355    #[expect(
356        private_bounds,
357        reason = "only Hydro collections can implement ReceiverComplete"
358    )]
359    pub fn cycle_with_initial<S, L2: Location<'a, DropConsistency = Tick<L::DropConsistency>>>(
360        &self,
361        initial: S,
362    ) -> (TickCycleHandle<'a, S>, S)
363    where
364        S: CycleCollectionWithInitial<'a, TickCycle, Location = L2>,
365    {
366        let cycle_id = self.flow_state().borrow_mut().next_cycle_id();
367        (
368            TickCycleHandle::new(cycle_id, Location::id(self)),
369            // no need to defer_tick, create_source_with_initial does it for us
370            S::create_source_with_initial(cycle_id, initial, self.clone().with_consistency_of()),
371        )
372    }
373}
374
375#[cfg(test)]
376mod tests {
377    #[cfg(feature = "sim")]
378    use stageleft::q;
379
380    #[cfg(feature = "sim")]
381    use crate::live_collections::sliced::sliced;
382    #[cfg(feature = "sim")]
383    use crate::location::Location;
384    #[cfg(feature = "sim")]
385    use crate::nondet::nondet;
386    #[cfg(feature = "sim")]
387    use crate::prelude::FlowBuilder;
388
389    #[cfg(feature = "sim")]
390    #[test]
391    fn sim_atomic_stream() {
392        let mut flow = FlowBuilder::new();
393        let node = flow.process::<()>();
394
395        let (write_send, write_req) = node.sim_input();
396        let (read_send, read_req) = node.sim_input::<(), _, _>();
397
398        let atomic_write = write_req.atomic();
399        let current_state = atomic_write.clone().fold(
400            q!(|| 0),
401            q!(|state: &mut i32, v: i32| {
402                *state += v;
403            }),
404        );
405
406        let write_ack_recv = atomic_write.end_atomic().sim_output();
407        let read_response_recv = sliced! {
408            let batch_of_req = use::batch(read_req, nondet!(/** test */));
409            let latest_singleton = use::atomic(current_state, nondet!(/** test */));
410            batch_of_req.cross_singleton(latest_singleton)
411        }
412        .sim_output();
413
414        let sim_compiled = flow.sim().compiled();
415        let instances = sim_compiled.exhaustive(async || {
416            write_send.send(1);
417            write_ack_recv.assert_yields([1]).await;
418            read_send.send(());
419            assert!(read_response_recv.next().await.1 >= 1);
420        });
421
422        assert_eq!(instances, 1);
423
424        let instances_read_before_write = sim_compiled.exhaustive(async || {
425            write_send.send(1);
426            read_send.send(());
427            write_ack_recv.assert_yields([1]).await;
428            let _ = read_response_recv.next().await;
429        });
430
431        assert_eq!(instances_read_before_write, 3); // read before write, write before read, both in same tick
432    }
433
434    #[cfg(feature = "sim")]
435    #[test]
436    #[should_panic]
437    fn sim_non_atomic_stream() {
438        // shows that atomic is necessary
439        let mut flow = FlowBuilder::new();
440        let node = flow.process::<()>();
441
442        let (write_send, write_req) = node.sim_input();
443        let (read_send, read_req) = node.sim_input::<(), _, _>();
444
445        let current_state = write_req.clone().fold(
446            q!(|| 0),
447            q!(|state: &mut i32, v: i32| {
448                *state += v;
449            }),
450        );
451
452        let write_ack_recv = write_req.sim_output();
453
454        let read_response_recv = sliced! {
455            let batch_of_req = use::batch(read_req, nondet!(/** test */));
456            let latest_singleton = use::snapshot(current_state, nondet!(/** test */));
457            batch_of_req.cross_singleton(latest_singleton)
458        }
459        .sim_output();
460
461        flow.sim().exhaustive(async || {
462            write_send.send(1);
463            write_ack_recv.assert_yields([1]).await;
464            read_send.send(());
465
466            let ((), v) = read_response_recv.next().await;
467            assert_eq!(v, 1);
468        });
469    }
470}