1use 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#[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
105pub trait DeferTick {
112 fn defer_tick(self) -> Self;
114}
115
116#[derive(Clone)]
118pub struct Tick<L> {
119 pub(crate) id: Option<ClockId>,
121 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 pub fn parent_location(&self) -> &L {
192 &self.l
193 }
194
195 #[deprecated(note = "use `.parent_location()` instead")]
197 pub fn outer(&self) -> &L {
198 self.parent_location()
199 }
200
201 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!());
228 Stream::new(self.clone(), inner.ir_node.replace(HydroNode::Placeholder))
229 }
230
231 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!())) }
263
264 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 #[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 self.singleton(q!(tokio::time::Instant::now()))
322 }
323
324 #[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 #[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 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!());
409 let latest_singleton = use::atomic(current_state, nondet!());
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); }
433
434 #[cfg(feature = "sim")]
435 #[test]
436 #[should_panic]
437 fn sim_non_atomic_stream() {
438 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!());
456 let latest_singleton = use::snapshot(current_state, nondet!());
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}