1use std::cell::RefCell;
4use std::future::Future;
5use std::hash::Hash;
6use std::marker::PhantomData;
7use std::ops::Deref;
8use std::rc::Rc;
9
10use stageleft::{IntoQuotedMut, QuotedWithContext, QuotedWithContextWithProps, q, quote_type};
11#[cfg(feature = "tokio")]
12use tokio::time::Instant;
13
14use super::OperatorContext;
15use super::boundedness::{Bounded, Boundedness, IsBounded, Unbounded};
16use super::keyed_singleton::KeyedSingleton;
17use super::keyed_stream::{Generate, KeyedStream};
18use super::optional::Optional;
19use super::singleton::Singleton;
20use crate::compile::builder::{CycleId, FlowState};
21use crate::compile::ir::{
22 CollectionKind, HydroIrOpMetadata, HydroNode, HydroRoot, SharedNode, StreamOrder, StreamRetry,
23};
24#[cfg(stageleft_runtime)]
25use crate::forward_handle::{CycleCollection, CycleCollectionWithInitial, ReceiverComplete};
26use crate::forward_handle::{ForwardRef, TickCycle};
27use crate::live_collections::batch_atomic::BatchAtomic;
28use crate::live_collections::singleton::SingletonBound;
29#[cfg(stageleft_runtime)]
30use crate::location::dynamic::{DynLocation, LocationId};
31use crate::location::tick::{Atomic, DeferTick};
32use crate::location::{Location, Tick, TopLevel, check_matching_location};
33use crate::manual_expr::ManualExpr;
34use crate::nondet::{NonDet, nondet};
35use crate::prelude::manual_proof;
36use crate::properties::{
37 AggFuncAlgebra, ApplyMonotoneStream, NotProved, StreamMapFuncAlgebra, ValidCommutativityFor,
38 ValidIdempotenceFor, ValidMutBorrowCommutativityFor, ValidMutBorrowIdempotenceFor,
39 ValidMutCommutativityFor, ValidMutIdempotenceFor,
40};
41
42pub mod networking;
43
44#[sealed::sealed]
46pub trait Ordering:
47 MinOrder<Self, Min = Self> + MinOrder<TotalOrder, Min = Self> + MinOrder<NoOrder, Min = NoOrder>
48{
49 const ORDERING_KIND: StreamOrder;
51}
52
53pub enum TotalOrder {}
57
58#[sealed::sealed]
59impl Ordering for TotalOrder {
60 const ORDERING_KIND: StreamOrder = StreamOrder::TotalOrder;
61}
62
63pub enum NoOrder {}
69
70#[sealed::sealed]
71impl Ordering for NoOrder {
72 const ORDERING_KIND: StreamOrder = StreamOrder::NoOrder;
73}
74
75#[sealed::sealed]
79pub trait WeakerOrderingThan<Other: ?Sized>: Ordering {}
80#[sealed::sealed]
81impl<O, O2: Ordering> WeakerOrderingThan<O2> for O where O: Ordering + MinOrder<O2, Min = O> {}
82
83#[sealed::sealed]
85pub trait MinOrder<Other: ?Sized> {
86 type Min: Ordering;
88}
89
90#[sealed::sealed]
91impl<O: Ordering> MinOrder<O> for TotalOrder {
92 type Min = O;
93}
94
95#[sealed::sealed]
96impl<O: Ordering> MinOrder<O> for NoOrder {
97 type Min = NoOrder;
98}
99
100#[sealed::sealed]
102pub trait Retries:
103 MinRetries<Self, Min = Self>
104 + MinRetries<ExactlyOnce, Min = Self>
105 + MinRetries<AtLeastOnce, Min = AtLeastOnce>
106{
107 const RETRIES_KIND: StreamRetry;
109}
110
111pub enum ExactlyOnce {}
114
115#[sealed::sealed]
116impl Retries for ExactlyOnce {
117 const RETRIES_KIND: StreamRetry = StreamRetry::ExactlyOnce;
118}
119
120pub enum AtLeastOnce {}
123
124#[sealed::sealed]
125impl Retries for AtLeastOnce {
126 const RETRIES_KIND: StreamRetry = StreamRetry::AtLeastOnce;
127}
128
129#[sealed::sealed]
133pub trait WeakerRetryThan<Other: ?Sized>: Retries {}
134#[sealed::sealed]
135impl<R, R2: Retries> WeakerRetryThan<R2> for R where R: Retries + MinRetries<R2, Min = R> {}
136
137#[sealed::sealed]
139pub trait MinRetries<Other: ?Sized> {
140 type Min: Retries + WeakerRetryThan<Self> + WeakerRetryThan<Other>;
142}
143
144#[sealed::sealed]
145impl<R: Retries> MinRetries<R> for ExactlyOnce {
146 type Min = R;
147}
148
149#[sealed::sealed]
150impl<R: Retries> MinRetries<R> for AtLeastOnce {
151 type Min = AtLeastOnce;
152}
153
154#[sealed::sealed]
155#[diagnostic::on_unimplemented(
156 message = "The input stream must be totally-ordered (`TotalOrder`), but has order `{Self}`. Strengthen the order upstream or consider a different API.",
157 label = "required here",
158 note = "To intentionally process the stream by observing a non-deterministic (shuffled) order of elements, use `.assume_ordering`. This introduces non-determinism so avoid unless necessary."
159)]
160pub trait IsOrdered: Ordering {}
162
163#[sealed::sealed]
164#[diagnostic::do_not_recommend]
165impl IsOrdered for TotalOrder {}
166
167#[sealed::sealed]
168#[diagnostic::on_unimplemented(
169 message = "The input stream must be exactly-once (`ExactlyOnce`), but has retries `{Self}`. Strengthen the retries guarantee upstream or consider a different API.",
170 label = "required here",
171 note = "To intentionally process the stream by observing non-deterministic (randomly duplicated) retries, use `.assume_retries`. This introduces non-determinism so avoid unless necessary."
172)]
173pub trait IsExactlyOnce: Retries {}
175
176#[sealed::sealed]
177#[diagnostic::do_not_recommend]
178impl IsExactlyOnce for ExactlyOnce {}
179
180pub struct Stream<
200 Type,
201 Loc,
202 Bound: Boundedness = Unbounded,
203 Order: Ordering = TotalOrder,
204 Retry: Retries = ExactlyOnce,
205> {
206 pub(crate) location: Loc,
207 pub(crate) ir_node: Rc<RefCell<HydroNode>>,
208 pub(crate) flow_state: FlowState,
209
210 _phantom: PhantomData<(Type, Loc, Bound, Order, Retry)>,
211}
212
213impl<T, L, B: Boundedness, O: Ordering, R: Retries> Drop for Stream<T, L, B, O, R> {
214 fn drop(&mut self) {
215 let ir_node = self.ir_node.replace(HydroNode::Placeholder);
216 if !matches!(ir_node, HydroNode::Placeholder) && !ir_node.is_shared_with_others() {
217 self.flow_state.borrow_mut().try_push_root(HydroRoot::Null {
218 input: Box::new(ir_node),
219 op_metadata: HydroIrOpMetadata::new(),
220 });
221 }
222 }
223}
224
225impl<'a, T, L, O: Ordering, R: Retries> From<Stream<T, L, Bounded, O, R>>
226 for Stream<T, L, Unbounded, O, R>
227where
228 L: Location<'a>,
229{
230 fn from(stream: Stream<T, L, Bounded, O, R>) -> Stream<T, L, Unbounded, O, R> {
231 let new_meta = stream
232 .location
233 .new_node_metadata(Stream::<T, L, Unbounded, O, R>::collection_kind());
234
235 let flow_state = stream.flow_state.clone();
236 Stream {
237 location: stream.location.clone(),
238 ir_node: super::tracked_ir_node(
239 &flow_state,
240 HydroNode::Cast {
241 inner: Box::new(stream.ir_node.replace(HydroNode::Placeholder)),
242 metadata: new_meta,
243 },
244 ),
245 flow_state,
246 _phantom: PhantomData,
247 }
248 }
249}
250
251impl<'a, T, L, B: Boundedness, R: Retries> From<Stream<T, L, B, TotalOrder, R>>
252 for Stream<T, L, B, NoOrder, R>
253where
254 L: Location<'a>,
255{
256 fn from(stream: Stream<T, L, B, TotalOrder, R>) -> Stream<T, L, B, NoOrder, R> {
257 stream.weaken_ordering()
258 }
259}
260
261impl<'a, T, L, B: Boundedness, O: Ordering> From<Stream<T, L, B, O, ExactlyOnce>>
262 for Stream<T, L, B, O, AtLeastOnce>
263where
264 L: Location<'a>,
265{
266 fn from(stream: Stream<T, L, B, O, ExactlyOnce>) -> Stream<T, L, B, O, AtLeastOnce> {
267 stream.weaken_retries()
268 }
269}
270
271impl<'a, T, L, O: Ordering, R: Retries> DeferTick for Stream<T, Tick<L>, Bounded, O, R>
272where
273 L: Location<'a>,
274{
275 fn defer_tick(self) -> Self {
276 Stream::defer_tick(self)
277 }
278}
279
280impl<'a, T, L, O: Ordering, R: Retries> CycleCollection<'a, TickCycle>
281 for Stream<T, Tick<L>, Bounded, O, R>
282where
283 L: Location<'a>,
284{
285 type Location = Tick<L>;
286
287 fn create_source(cycle_id: CycleId, location: Tick<L>) -> Self {
288 Stream::new(
289 location.clone(),
290 HydroNode::CycleSource {
291 cycle_id,
292 metadata: location.new_node_metadata(Self::collection_kind()),
293 },
294 )
295 }
296}
297
298impl<'a, T, L, O: Ordering, R: Retries> CycleCollectionWithInitial<'a, TickCycle>
299 for Stream<T, Tick<L>, Bounded, O, R>
300where
301 L: Location<'a>,
302{
303 type Location = Tick<L>;
304
305 fn location(&self) -> &Self::Location {
306 self.location()
307 }
308
309 fn create_source_with_initial(cycle_id: CycleId, initial: Self, location: Tick<L>) -> Self {
310 let from_previous_tick: Stream<T, Tick<L>, Bounded, O, R> = Stream::new(
311 location.clone(),
312 HydroNode::DeferTick {
313 input: Box::new(HydroNode::CycleSource {
314 cycle_id,
315 metadata: location.new_node_metadata(Self::collection_kind()),
316 }),
317 metadata: location.new_node_metadata(Self::collection_kind()),
318 },
319 );
320
321 from_previous_tick.chain(initial.filter_if(location.optional_first_tick(q!(())).is_some()))
322 }
323}
324
325impl<'a, T, L, O: Ordering, R: Retries> ReceiverComplete<'a, TickCycle>
326 for Stream<T, Tick<L>, Bounded, O, R>
327where
328 L: Location<'a>,
329{
330 fn complete(self, cycle_id: CycleId, expected_location: LocationId) {
331 assert_eq!(
332 Location::id(&self.location),
333 expected_location,
334 "locations do not match"
335 );
336 self.location
337 .flow_state()
338 .borrow_mut()
339 .push_root(HydroRoot::CycleSink {
340 cycle_id,
341 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
342 op_metadata: HydroIrOpMetadata::new(),
343 });
344 }
345}
346
347impl<'a, T, L, B: Boundedness, O: Ordering, R: Retries> CycleCollection<'a, ForwardRef>
348 for Stream<T, L, B, O, R>
349where
350 L: Location<'a>,
351{
352 type Location = L;
353
354 fn create_source(cycle_id: CycleId, location: L) -> Self {
355 Stream::new(
356 location.clone(),
357 HydroNode::CycleSource {
358 cycle_id,
359 metadata: location.new_node_metadata(Self::collection_kind()),
360 },
361 )
362 }
363}
364
365impl<'a, T, L, B: Boundedness, O: Ordering, R: Retries> ReceiverComplete<'a, ForwardRef>
366 for Stream<T, L, B, O, R>
367where
368 L: Location<'a>,
369{
370 fn complete(self, cycle_id: CycleId, expected_location: LocationId) {
371 assert_eq!(
372 Location::id(&self.location),
373 expected_location,
374 "locations do not match"
375 );
376 self.location
377 .flow_state()
378 .borrow_mut()
379 .push_root(HydroRoot::CycleSink {
380 cycle_id,
381 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
382 op_metadata: HydroIrOpMetadata::new(),
383 });
384 }
385}
386
387impl<'a, T, L, B: Boundedness, O: Ordering, R: Retries> Clone for Stream<T, L, B, O, R>
388where
389 T: Clone,
390 L: Location<'a>,
391{
392 fn clone(&self) -> Self {
393 if !matches!(self.ir_node.borrow().deref(), HydroNode::Tee { .. }) {
394 let orig_ir_node = self.ir_node.replace(HydroNode::Placeholder);
395 *self.ir_node.borrow_mut() = HydroNode::Tee {
396 inner: SharedNode(Rc::new(RefCell::new(orig_ir_node))),
397 metadata: self.location.new_node_metadata(Self::collection_kind()),
398 };
399 }
400
401 let HydroNode::Tee { inner, metadata } = &*self.ir_node.borrow() else {
402 unreachable!()
403 };
404 Stream {
405 location: self.location.clone(),
406 flow_state: self.flow_state.clone(),
407 ir_node: super::tracked_ir_node(
408 &self.flow_state,
409 HydroNode::Tee {
410 inner: SharedNode(inner.0.clone()),
411 metadata: metadata.clone(),
412 },
413 ),
414 _phantom: PhantomData,
415 }
416 }
417}
418
419impl<'a, T, L, B: Boundedness, O: Ordering, R: Retries> Stream<T, L, B, O, R>
420where
421 L: Location<'a>,
422{
423 pub(crate) fn new(location: L, ir_node: HydroNode) -> Self {
424 debug_assert_eq!(ir_node.metadata().location_id, Location::id(&location));
425 debug_assert_eq!(ir_node.metadata().collection_kind, Self::collection_kind());
426
427 let flow_state = location.flow_state().clone();
428 let ir_node = super::tracked_ir_node(&flow_state, ir_node);
429 Stream {
430 location,
431 flow_state,
432 ir_node,
433 _phantom: PhantomData,
434 }
435 }
436
437 pub fn location(&self) -> &L {
439 &self.location
440 }
441
442 pub fn by_ref(&self) -> crate::handoff_ref::StreamRef<'a, '_, T, L, B>
447 where
448 B: IsBounded,
449 {
450 crate::handoff_ref::StreamRef::new(&self.ir_node)
451 }
452
453 pub fn by_mut(&self) -> crate::handoff_ref::StreamMut<'a, '_, T, L, B>
456 where
457 B: IsBounded,
458 {
459 crate::handoff_ref::StreamMut::new(&self.ir_node)
460 }
461
462 pub fn weaken_consistency(self) -> Stream<T, L::DropConsistency, B, O, R>
465 where
466 L: Location<'a>,
467 {
468 if L::consistency()
469 .is_none_or(|c| c == crate::location::dynamic::ClusterConsistency::NoConsistency)
470 {
471 Stream::new(
473 self.location.drop_consistency(),
474 self.ir_node.replace(HydroNode::Placeholder),
475 )
476 } else {
477 Stream::new(
478 self.location.drop_consistency(),
479 HydroNode::Cast {
480 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
481 metadata: self.location.drop_consistency().new_node_metadata(Stream::<
482 T,
483 L::DropConsistency,
484 B,
485 O,
486 R,
487 >::collection_kind(
488 )),
489 },
490 )
491 }
492 }
493
494 pub fn assert_has_consistency_of<L2: Location<'a, DropConsistency = L::DropConsistency>>(
498 self,
499 _proof: impl crate::properties::ConsistencyProof,
500 ) -> Stream<T, L2, B, O, R>
501 where
502 L: Location<'a>,
503 {
504 if L::consistency() == L2::consistency() {
505 Stream::new(
506 self.location.with_consistency_of(),
507 self.ir_node.replace(HydroNode::Placeholder),
508 )
509 } else {
510 Stream::new(
511 self.location.with_consistency_of(),
512 HydroNode::AssertIsConsistent {
513 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
514 trusted: false,
515 metadata: self
516 .location
517 .clone()
518 .with_consistency_of::<L2>()
519 .new_node_metadata(Stream::<T, L2, B, O, R>::collection_kind()),
520 },
521 )
522 }
523 }
524
525 pub(crate) fn assert_has_consistency_of_trusted<
526 L2: Location<'a, DropConsistency = L::DropConsistency>,
527 >(
528 self,
529 _proof: impl crate::properties::ConsistencyProof,
530 ) -> Stream<T, L2, B, O, R>
531 where
532 L: Location<'a>,
533 {
534 if L::consistency() == L2::consistency() {
535 Stream::new(
536 self.location.with_consistency_of(),
537 self.ir_node.replace(HydroNode::Placeholder),
538 )
539 } else {
540 Stream::new(
541 self.location.with_consistency_of(),
542 HydroNode::AssertIsConsistent {
543 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
544 trusted: true,
545 metadata: self
546 .location
547 .clone()
548 .with_consistency_of::<L2>()
549 .new_node_metadata(Stream::<T, L2, B, O, R>::collection_kind()),
550 },
551 )
552 }
553 }
554
555 pub(crate) fn collection_kind() -> CollectionKind {
556 CollectionKind::Stream {
557 bound: B::BOUND_KIND,
558 order: O::ORDERING_KIND,
559 retry: R::RETRIES_KIND,
560 element_type: quote_type::<T>().into(),
561 }
562 }
563
564 pub fn map<U, F, C, I, const WAS_MUT: bool>(
592 self,
593 f: impl IntoQuotedMut<
594 'a,
595 F,
596 OperatorContext<L, B>,
597 StreamMapFuncAlgebra<T, B, C, I, L::SimHookScope>,
598 >,
599 ) -> Stream<U, L, B, O, R>
600 where
601 F: FnMut(T) -> U + 'a,
602 C: ValidMutCommutativityFor<F, T, U, O, WAS_MUT>,
603 I: ValidMutIdempotenceFor<F, T, U, R, WAS_MUT>,
604 {
605 let f = crate::handoff_ref::with_ref_capture(|| {
606 let (expr, proof) =
607 f.splice_fnmut1_ctx_props(&OperatorContext::<L, B>::new(&self.location));
608 proof.register_proof(&expr);
609 expr.into()
610 });
611 Stream::new(
612 self.location.clone(),
613 HydroNode::Map {
614 f,
615 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
616 metadata: self
617 .location
618 .new_node_metadata(Stream::<U, L, B, O, R>::collection_kind()),
619 },
620 )
621 }
622
623 pub fn flat_map_ordered<U, I, F, C, Idemp, const WAS_MUT: bool>(
648 self,
649 f: impl IntoQuotedMut<
650 'a,
651 F,
652 OperatorContext<L, B>,
653 StreamMapFuncAlgebra<T, B, C, Idemp, L::SimHookScope>,
654 >,
655 ) -> Stream<U, L, B, O, R>
656 where
657 I: IntoIterator<Item = U>,
658 F: FnMut(T) -> I + 'a,
659 C: ValidMutCommutativityFor<F, T, I, O, WAS_MUT>,
660 Idemp: ValidMutIdempotenceFor<F, T, I, R, WAS_MUT>,
661 {
662 let f = crate::handoff_ref::with_ref_capture(|| {
663 let (expr, proof) =
664 f.splice_fnmut1_ctx_props(&OperatorContext::<L, B>::new(&self.location));
665 proof.register_proof(&expr);
666 expr.into()
667 });
668 Stream::new(
669 self.location.clone(),
670 HydroNode::FlatMap {
671 f,
672 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
673 metadata: self
674 .location
675 .new_node_metadata(Stream::<U, L, B, O, R>::collection_kind()),
676 },
677 )
678 }
679
680 pub fn flat_map_unordered<U, I, F, C, Idemp, const WAS_MUT: bool>(
707 self,
708 f: impl IntoQuotedMut<
709 'a,
710 F,
711 OperatorContext<L, B>,
712 StreamMapFuncAlgebra<T, B, C, Idemp, L::SimHookScope>,
713 >,
714 ) -> Stream<U, L, B, NoOrder, R>
715 where
716 I: IntoIterator<Item = U>,
717 F: FnMut(T) -> I + 'a,
718 C: ValidMutCommutativityFor<F, T, I, O, WAS_MUT>,
719 Idemp: ValidMutIdempotenceFor<F, T, I, R, WAS_MUT>,
720 {
721 let f = crate::handoff_ref::with_ref_capture(|| {
722 let (expr, proof) =
723 f.splice_fnmut1_ctx_props(&OperatorContext::<L, B>::new(&self.location));
724 proof.register_proof(&expr);
725 expr.into()
726 });
727 Stream::new(
728 self.location.clone(),
729 HydroNode::FlatMap {
730 f,
731 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
732 metadata: self
733 .location
734 .new_node_metadata(Stream::<U, L, B, NoOrder, R>::collection_kind()),
735 },
736 )
737 }
738
739 pub fn flatten_ordered<U>(self) -> Stream<U, L, B, O, R>
762 where
763 T: IntoIterator<Item = U>,
764 {
765 self.flat_map_ordered(q!(|d| d))
766 }
767
768 pub fn flatten_unordered<U>(self) -> Stream<U, L, B, NoOrder, R>
795 where
796 T: IntoIterator<Item = U>,
797 {
798 self.flat_map_unordered(q!(|d| d))
799 }
800
801 pub fn flat_map_stream_blocking<U, S, F, C, Idemp, const WAS_MUT: bool>(
805 self,
806 f: impl IntoQuotedMut<
807 'a,
808 F,
809 OperatorContext<L, B>,
810 StreamMapFuncAlgebra<T, B, C, Idemp, L::SimHookScope>,
811 >,
812 ) -> Stream<U, L, B, O, R>
813 where
814 S: futures::Stream<Item = U>,
815 F: FnMut(T) -> S + 'a,
816 C: ValidMutCommutativityFor<F, T, S, O, WAS_MUT>,
817 Idemp: ValidMutIdempotenceFor<F, T, S, R, WAS_MUT>,
818 {
819 let f = crate::handoff_ref::with_ref_capture(|| {
820 let (expr, proof) =
821 f.splice_fnmut1_ctx_props(&OperatorContext::<L, B>::new(&self.location));
822 proof.register_proof(&expr);
823 expr.into()
824 });
825 Stream::new(
826 self.location.clone(),
827 HydroNode::FlatMapStreamBlocking {
828 f,
829 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
830 metadata: self
831 .location
832 .new_node_metadata(Stream::<U, L, B, O, R>::collection_kind()),
833 },
834 )
835 }
836
837 pub fn flatten_stream_blocking<U>(self) -> Stream<U, L, B, O, R>
841 where
842 T: futures::Stream<Item = U>,
843 {
844 self.flat_map_stream_blocking(q!(|d| d))
845 }
846
847 pub fn filter<F, C, Idemp, const WAS_MUT: bool>(
879 self,
880 f: impl IntoQuotedMut<
881 'a,
882 F,
883 OperatorContext<L, B>,
884 StreamMapFuncAlgebra<T, B, C, Idemp, L::SimHookScope>,
885 >,
886 ) -> Self
887 where
888 F: FnMut(&T) -> bool + 'a,
889 C: ValidMutBorrowCommutativityFor<F, T, bool, O, WAS_MUT>,
890 Idemp: ValidMutBorrowIdempotenceFor<F, T, bool, R, WAS_MUT>,
891 {
892 let f = crate::handoff_ref::with_ref_capture(|| {
893 let (expr, proof) =
894 f.splice_fnmut1_borrow_ctx_props(&OperatorContext::<L, B>::new(&self.location));
895 proof.register_proof(&expr);
896 expr.into()
897 });
898 Stream::new(
899 self.location.clone(),
900 HydroNode::Filter {
901 f,
902 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
903 metadata: self.location.new_node_metadata(Self::collection_kind()),
904 },
905 )
906 }
907
908 pub fn partition<F, C, Idemp, const WAS_MUT: bool>(
943 self,
944 f: impl IntoQuotedMut<
945 'a,
946 F,
947 OperatorContext<L, B>,
948 StreamMapFuncAlgebra<T, B, C, Idemp, L::SimHookScope>,
949 >,
950 ) -> (Stream<T, L, B, O, R>, Stream<T, L, B, O, R>)
951 where
952 F: FnMut(&T) -> bool + 'a,
953 C: ValidMutBorrowCommutativityFor<F, T, bool, O, WAS_MUT>,
954 Idemp: ValidMutBorrowIdempotenceFor<F, T, bool, R, WAS_MUT>,
955 {
956 let f = crate::handoff_ref::with_ref_capture(|| {
957 let (expr, proof) =
958 f.splice_fnmut1_borrow_ctx_props(&OperatorContext::<L, B>::new(&self.location));
959 proof.register_proof(&expr);
960 expr.into()
961 });
962 let shared = Rc::new(RefCell::new(HydroNode::PartitionShared {
963 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
964 f,
965 metadata: self.location.new_node_metadata(Self::collection_kind()),
966 }));
967
968 let true_stream = Stream::new(
969 self.location.clone(),
970 HydroNode::PartitionSide {
971 inner: SharedNode(Rc::clone(&shared)),
972 is_true: true,
973 metadata: self.location.new_node_metadata(Self::collection_kind()),
974 },
975 );
976
977 let false_stream = Stream::new(
978 self.location.clone(),
979 HydroNode::PartitionSide {
980 inner: SharedNode(shared),
981 is_true: false,
982 metadata: self.location.new_node_metadata(Self::collection_kind()),
983 },
984 );
985
986 (true_stream, false_stream)
987 }
988
989 pub fn filter_map<U, F, C, Idemp, const WAS_MUT: bool>(
1009 self,
1010 f: impl IntoQuotedMut<
1011 'a,
1012 F,
1013 OperatorContext<L, B>,
1014 StreamMapFuncAlgebra<T, B, C, Idemp, L::SimHookScope>,
1015 >,
1016 ) -> Stream<U, L, B, O, R>
1017 where
1018 F: FnMut(T) -> Option<U> + 'a,
1019 C: ValidMutCommutativityFor<F, T, Option<U>, O, WAS_MUT>,
1020 Idemp: ValidMutIdempotenceFor<F, T, Option<U>, R, WAS_MUT>,
1021 {
1022 let f = crate::handoff_ref::with_ref_capture(|| {
1023 let (expr, proof) =
1024 f.splice_fnmut1_ctx_props(&OperatorContext::<L, B>::new(&self.location));
1025 proof.register_proof(&expr);
1026 expr.into()
1027 });
1028 Stream::new(
1029 self.location.clone(),
1030 HydroNode::FilterMap {
1031 f,
1032 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
1033 metadata: self
1034 .location
1035 .new_node_metadata(Stream::<U, L, B, O, R>::collection_kind()),
1036 },
1037 )
1038 }
1039
1040 pub fn cross_singleton<O2>(
1065 self,
1066 other: impl Into<Optional<O2, L, Bounded>>,
1067 ) -> Stream<(T, O2), L, B, O, R>
1068 where
1069 O2: Clone,
1070 {
1071 let other: Optional<O2, L, Bounded> = other.into();
1072 check_matching_location(&self.location, &other.location);
1073
1074 Stream::new(
1075 self.location.clone(),
1076 HydroNode::CrossSingleton {
1077 left: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
1078 right: Box::new(other.ir_node.replace(HydroNode::Placeholder)),
1079 metadata: self
1080 .location
1081 .new_node_metadata(Stream::<(T, O2), L, B, O, R>::collection_kind()),
1082 },
1083 )
1084 }
1085
1086 pub fn filter_if(self, signal: Singleton<bool, L, Bounded>) -> Stream<T, L, B, O, R> {
1118 self.cross_singleton(signal.filter(q!(|b| *b)))
1119 .map(q!(|(d, _)| d))
1120 }
1121
1122 #[deprecated(note = "use `filter_if` with `Optional::is_some()` instead")]
1157 pub fn filter_if_some<U>(self, signal: Optional<U, L, Bounded>) -> Stream<T, L, B, O, R> {
1158 self.filter_if(signal.is_some())
1159 }
1160
1161 #[deprecated(note = "use `filter_if` with `!Optional::is_some()` instead")]
1196 pub fn filter_if_none<U>(self, other: Optional<U, L, Bounded>) -> Stream<T, L, B, O, R> {
1197 self.filter_if(other.is_none())
1198 }
1199
1200 pub fn cross_product<T2, B2: Boundedness, O2: Ordering, R2: Retries>(
1225 self,
1226 other: Stream<T2, L, B2, O2, R2>,
1227 ) -> Stream<(T, T2), L, B, B2::PreserveOrderIfBounded<O>, <R as MinRetries<R2>>::Min>
1228 where
1229 T: Clone,
1230 T2: Clone,
1231 R: MinRetries<R2>,
1232 {
1233 self.map(q!(|v| ((), v)))
1234 .join(other.map(q!(|v| ((), v))))
1235 .map(q!(|((), (v1, v2))| (v1, v2)))
1236 }
1237
1238 pub fn unique(self) -> Stream<T, L, B, O, ExactlyOnce>
1257 where
1258 T: Eq + Hash,
1259 {
1260 Stream::new(
1261 self.location.clone(),
1262 HydroNode::Unique {
1263 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
1264 metadata: self
1265 .location
1266 .new_node_metadata(Stream::<T, L, B, O, ExactlyOnce>::collection_kind()),
1267 },
1268 )
1269 }
1270
1271 pub fn filter_not_in<O2: Ordering, B2>(self, other: Stream<T, L, B2, O2, R>) -> Self
1297 where
1298 T: Eq + Hash,
1299 B2: IsBounded,
1300 {
1301 check_matching_location(&self.location, &other.location);
1302
1303 Stream::new(
1304 self.location.clone(),
1305 HydroNode::Difference {
1306 pos: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
1307 neg: Box::new(other.ir_node.replace(HydroNode::Placeholder)),
1308 metadata: self
1309 .location
1310 .new_node_metadata(Stream::<T, L, Bounded, O, R>::collection_kind()),
1311 },
1312 )
1313 }
1314
1315 pub fn inspect<F, C, Idemp, const WAS_MUT: bool>(
1341 self,
1342 f: impl IntoQuotedMut<
1343 'a,
1344 F,
1345 OperatorContext<L::DropConsistency, B>,
1346 StreamMapFuncAlgebra<T, B, C, Idemp, L::SimHookScope>,
1347 >,
1348 ) -> Self
1349 where
1350 F: FnMut(&T) + 'a,
1351 C: ValidMutBorrowCommutativityFor<F, T, (), O, WAS_MUT>,
1352 Idemp: ValidMutBorrowIdempotenceFor<F, T, (), R, WAS_MUT>,
1353 {
1354 let f = crate::handoff_ref::with_ref_capture(|| {
1355 let (expr, proof) =
1356 f.splice_fnmut1_borrow_ctx_props(&OperatorContext::<L::DropConsistency, B>::new(
1357 &self.location.drop_consistency(),
1358 ));
1359 proof.register_proof(&expr);
1360 expr.into()
1361 });
1362
1363 Stream::new(
1364 self.location.clone(),
1365 HydroNode::Inspect {
1366 f,
1367 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
1368 metadata: self.location.new_node_metadata(Self::collection_kind()),
1369 },
1370 )
1371 }
1372
1373 pub fn for_each<F: FnMut(T) + 'a, C, I>(
1395 self,
1396 f: impl IntoQuotedMut<
1397 'a,
1398 F,
1399 OperatorContext<L, B>,
1400 AggFuncAlgebra<T, B, C, I, NotProved, L::SimHookScope>,
1401 >,
1402 ) where
1403 C: ValidCommutativityFor<O>,
1404 I: ValidIdempotenceFor<R>,
1405 {
1406 let f = crate::handoff_ref::with_ref_capture(|| {
1407 let (f, proof) =
1408 f.splice_fnmut1_ctx_props(&OperatorContext::<L, B>::new(&self.location));
1409 proof.register_proof(&f);
1410 f.into()
1411 });
1412 self.location
1413 .flow_state()
1414 .borrow_mut()
1415 .push_root(HydroRoot::ForEach {
1416 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
1417 f,
1418 op_metadata: HydroIrOpMetadata::new(),
1419 });
1420 }
1421
1422 pub fn dest_sink<S>(self, sink: impl QuotedWithContext<'a, S, L>)
1428 where
1429 O: IsOrdered,
1430 R: IsExactlyOnce,
1431 S: 'a + futures::Sink<T> + Unpin,
1432 {
1433 self.location
1434 .flow_state()
1435 .borrow_mut()
1436 .push_root(HydroRoot::DestSink {
1437 sink: sink.splice_typed_ctx(&self.location).into(),
1438 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
1439 op_metadata: HydroIrOpMetadata::new(),
1440 });
1441 }
1442
1443 pub fn enumerate(self) -> Stream<(usize, T), L, B, O, R>
1463 where
1464 O: IsOrdered,
1465 R: IsExactlyOnce,
1466 {
1467 Stream::new(
1468 self.location.clone(),
1469 HydroNode::Enumerate {
1470 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
1471 metadata: self.location.new_node_metadata(Stream::<
1472 (usize, T),
1473 L,
1474 B,
1475 TotalOrder,
1476 ExactlyOnce,
1477 >::collection_kind()),
1478 },
1479 )
1480 }
1481
1482 pub fn fold<A, I, F, C, Idemp, M, B2: SingletonBound>(
1506 self,
1507 init: impl IntoQuotedMut<'a, I, OperatorContext<L, B>>,
1508 comb: impl IntoQuotedMut<
1509 'a,
1510 F,
1511 OperatorContext<L, B>,
1512 AggFuncAlgebra<T, B, C, Idemp, M, L::SimHookScope>,
1513 >,
1514 ) -> Singleton<A, L, B2>
1515 where
1516 I: Fn() -> A + 'a,
1517 F: 'a + Fn(&mut A, T),
1518 C: ValidCommutativityFor<O>,
1519 Idemp: ValidIdempotenceFor<R>,
1520 B: ApplyMonotoneStream<M, B2>,
1521 {
1522 let init = init
1523 .splice_fn0_ctx(&OperatorContext::<L, B>::new(&self.location))
1524 .into();
1525 let (comb, proof) =
1526 comb.splice_fn2_borrow_mut_ctx_props(&OperatorContext::<L, B>::new(&self.location));
1527 let ordering_hook = proof.register_proof(&comb);
1528
1529 let nondet = nondet!();
1533 let retried: Stream<T, L::DropConsistency, B, O, ExactlyOnce> = self.assume_retries(nondet);
1534
1535 let mut metadata = retried
1536 .location
1537 .new_node_metadata(Singleton::<A, L::DropConsistency, B2>::collection_kind());
1538 metadata.op.sim_hook_id = ordering_hook.map(|hook| hook.id);
1539
1540 let core = HydroNode::Fold {
1541 init,
1542 acc: comb.into(),
1543 input: Box::new(retried.ir_node.replace(HydroNode::Placeholder)),
1544 metadata,
1545 };
1550
1551 Singleton::new(retried.location.clone(), core)
1552 .assert_has_consistency_of(manual_proof!())
1553 }
1554
1555 pub fn reduce<F, C, Idemp>(
1578 self,
1579 comb: impl IntoQuotedMut<
1580 'a,
1581 F,
1582 OperatorContext<L, B>,
1583 AggFuncAlgebra<T, B, C, Idemp, NotProved, L::SimHookScope>,
1584 >,
1585 ) -> Optional<T, L, B::AggregatedOptional>
1586 where
1587 F: Fn(&mut T, T) + 'a,
1588 C: ValidCommutativityFor<O>,
1589 Idemp: ValidIdempotenceFor<R>,
1590 {
1591 let (f, proof) =
1592 comb.splice_fn2_borrow_mut_ctx_props(&OperatorContext::<L, B>::new(&self.location));
1593 let ordering_hook = proof.register_proof(&f);
1594
1595 let nondet_retries = nondet!();
1596 let ordered_etc: Stream<T, L::DropConsistency, B> =
1597 self.assume_retries(nondet_retries).assume_ordering(nondet!(
1598 hook = ordering_hook
1601 ));
1602
1603 let core = HydroNode::Reduce {
1604 f: f.into(),
1605 input: Box::new(ordered_etc.ir_node.replace(HydroNode::Placeholder)),
1606 metadata: ordered_etc.location.new_node_metadata(Optional::<
1607 T,
1608 L::DropConsistency,
1609 B::AggregatedOptional,
1610 >::collection_kind()),
1611 };
1612
1613 Optional::new(ordered_etc.location.clone(), core)
1614 .assert_has_consistency_of(manual_proof!())
1615 }
1616
1617 pub fn max(self) -> Optional<T, L, B::AggregatedOptional>
1637 where
1638 T: Ord,
1639 {
1640 self.assume_retries_trusted::<ExactlyOnce>(nondet!())
1641 .assume_ordering_trusted_bounded::<TotalOrder>(
1642 nondet!(),
1643 )
1644 .reduce(q!(|curr, new| {
1645 if new > *curr {
1646 *curr = new;
1647 }
1648 }))
1649 }
1650
1651 pub fn min(self) -> Optional<T, L, B::AggregatedOptional>
1671 where
1672 T: Ord,
1673 {
1674 self.assume_retries_trusted::<ExactlyOnce>(nondet!())
1675 .assume_ordering_trusted_bounded::<TotalOrder>(
1676 nondet!(),
1677 )
1678 .reduce(q!(|curr, new| {
1679 if new < *curr {
1680 *curr = new;
1681 }
1682 }))
1683 }
1684
1685 pub fn first(self) -> Optional<T, L, B::AggregatedOptional>
1708 where
1709 O: IsOrdered,
1710 {
1711 self.make_totally_ordered()
1712 .assume_retries_trusted::<ExactlyOnce>(nondet!())
1713 .generator(q!(|| ()), q!(|_, item| Generate::Return(item)))
1714 .reduce(q!(|_, _| {}))
1715 }
1716
1717 pub fn last(self) -> Optional<T, L, B::AggregatedOptional>
1740 where
1741 O: IsOrdered,
1742 {
1743 self.make_totally_ordered()
1744 .assume_retries_trusted::<ExactlyOnce>(nondet!())
1745 .reduce(q!(|curr, new| *curr = new))
1746 }
1747
1748 pub fn limit(
1771 self,
1772 n: impl QuotedWithContext<'a, usize, OperatorContext<L, B>> + Copy + 'a,
1773 ) -> Stream<T, L, B, TotalOrder, ExactlyOnce>
1774 where
1775 O: IsOrdered,
1776 R: IsExactlyOnce,
1777 {
1778 self.generator(
1779 q!(|| 0usize),
1780 q!(move |count, item| {
1781 if *count == n {
1782 Generate::Break
1783 } else {
1784 *count += 1;
1785 if *count == n {
1786 Generate::Return(item)
1787 } else {
1788 Generate::Yield(item)
1789 }
1790 }
1791 }),
1792 )
1793 }
1794
1795 pub fn collect_vec(self) -> Singleton<Vec<T>, L, B>
1821 where
1822 O: IsOrdered,
1823 R: IsExactlyOnce,
1824 {
1825 self.make_totally_ordered().make_exactly_once().fold(
1826 q!(|| vec![]),
1827 q!(|acc, v| {
1828 acc.push(v);
1829 }),
1830 )
1831 }
1832
1833 pub fn scan<A, U, I, F>(
1899 self,
1900 init: impl IntoQuotedMut<'a, I, OperatorContext<L, B>>,
1901 f: impl IntoQuotedMut<'a, F, OperatorContext<L, B>>,
1902 ) -> Stream<U, L, B, TotalOrder, ExactlyOnce>
1903 where
1904 O: IsOrdered,
1905 R: IsExactlyOnce,
1906 I: Fn() -> A + 'a,
1907 F: Fn(&mut A, T) -> Option<U> + 'a,
1908 {
1909 let init = crate::handoff_ref::with_ref_capture(|| {
1910 init.splice_fn0_ctx(&OperatorContext::<L, B>::new(&self.location))
1911 .into()
1912 });
1913 let f = crate::handoff_ref::with_ref_capture(|| {
1914 f.splice_fn2_borrow_mut_ctx(&OperatorContext::<L, B>::new(&self.location))
1915 .into()
1916 });
1917
1918 Stream::new(
1919 self.location.clone(),
1920 HydroNode::Scan {
1921 init,
1922 acc: f,
1923 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
1924 metadata: self.location.new_node_metadata(
1925 Stream::<U, L, B, TotalOrder, ExactlyOnce>::collection_kind(),
1926 ),
1927 },
1928 )
1929 }
1930
1931 pub fn scan_async_blocking<A, U, I, F, Fut>(
1970 self,
1971 init: impl IntoQuotedMut<'a, I, OperatorContext<L, B>>,
1972 f: impl IntoQuotedMut<'a, F, OperatorContext<L, B>>,
1973 ) -> Stream<U, L, B, TotalOrder, ExactlyOnce>
1974 where
1975 O: IsOrdered,
1976 R: IsExactlyOnce,
1977 I: Fn() -> A + 'a,
1978 F: Fn(&mut A, T) -> Fut + 'a,
1979 Fut: Future<Output = Option<U>> + 'a,
1980 {
1981 let init = crate::handoff_ref::with_ref_capture(|| {
1982 init.splice_fn0_ctx(&OperatorContext::<L, B>::new(&self.location))
1983 .into()
1984 });
1985 let f = crate::handoff_ref::with_ref_capture(|| {
1986 f.splice_fn2_borrow_mut_ctx(&OperatorContext::<L, B>::new(&self.location))
1987 .into()
1988 });
1989
1990 Stream::new(
1991 self.location.clone(),
1992 HydroNode::ScanAsyncBlocking {
1993 init,
1994 acc: f,
1995 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
1996 metadata: self.location.new_node_metadata(
1997 Stream::<U, L, B, TotalOrder, ExactlyOnce>::collection_kind(),
1998 ),
1999 },
2000 )
2001 }
2002
2003 pub fn generator<A, U, I, F>(
2048 self,
2049 init: impl IntoQuotedMut<'a, I, OperatorContext<L, B>> + Copy,
2050 f: impl IntoQuotedMut<'a, F, OperatorContext<L, B>> + Copy,
2051 ) -> Stream<U, L, B, TotalOrder, ExactlyOnce>
2052 where
2053 O: IsOrdered,
2054 R: IsExactlyOnce,
2055 I: Fn() -> A + 'a,
2056 F: Fn(&mut A, T) -> Generate<U> + 'a,
2057 {
2058 let init: ManualExpr<I, _> =
2059 ManualExpr::new(move |ctx: &OperatorContext<L, B>| init.splice_fn0_ctx(ctx));
2060 let f: ManualExpr<F, _> =
2061 ManualExpr::new(move |ctx: &OperatorContext<L, B>| f.splice_fn2_borrow_mut_ctx(ctx));
2062
2063 let this = self.make_totally_ordered().make_exactly_once();
2064
2065 let scan_init = crate::handoff_ref::with_ref_capture(|| {
2070 q!(|| None)
2071 .splice_fn0_ctx::<Option<Option<A>>>(&this.location)
2072 .into()
2073 });
2074 let scan_f = crate::handoff_ref::with_ref_capture(|| {
2075 q!(move |state: &mut Option<Option<_>>, v| {
2076 if state.is_none() {
2077 *state = Some(Some(init()));
2078 }
2079 match state {
2080 Some(Some(state_value)) => match f(state_value, v) {
2081 Generate::Yield(out) => Some(Some(out)),
2082 Generate::Return(out) => {
2083 *state = Some(None);
2084 Some(Some(out))
2085 }
2086 Generate::Break => None,
2090 Generate::Continue => Some(None),
2091 },
2092 _ => None,
2094 }
2095 })
2096 .splice_fn2_borrow_mut_ctx::<Option<Option<A>>, T, _>(&OperatorContext::<L, B>::new(
2097 &this.location,
2098 ))
2099 .into()
2100 });
2101
2102 let scan_node = HydroNode::Scan {
2103 init: scan_init,
2104 acc: scan_f,
2105 input: Box::new(this.ir_node.replace(HydroNode::Placeholder)),
2106 metadata: this.location.new_node_metadata(Stream::<
2107 Option<U>,
2108 L,
2109 B,
2110 TotalOrder,
2111 ExactlyOnce,
2112 >::collection_kind()),
2113 };
2114
2115 let flatten_f = q!(|d| d)
2116 .splice_fn1_ctx::<Option<U>, _>(&this.location)
2117 .into();
2118 let flatten_node = HydroNode::FlatMap {
2119 f: flatten_f,
2120 input: Box::new(scan_node),
2121 metadata: this
2122 .location
2123 .new_node_metadata(Stream::<U, L, B, TotalOrder, ExactlyOnce>::collection_kind()),
2124 };
2125
2126 Stream::new(this.location.clone(), flatten_node)
2127 }
2128
2129 #[cfg(feature = "tokio")]
2142 #[expect(
2143 clippy::type_complexity,
2144 reason = "composite hook payload names each internal operator's handle type"
2145 )]
2146 pub fn sample_every(
2147 self,
2148 interval: impl QuotedWithContext<'a, std::time::Duration, L> + Copy + 'a,
2149 mut nondet: NonDet<(
2150 Option<crate::sim_hooks::BatchHook<T, O, R, L::SimHookScope>>,
2151 Option<crate::sim_hooks::BatchHook<(), TotalOrder, ExactlyOnce, L::SimHookScope>>,
2152 )>,
2153 ) -> Stream<T, L::DropConsistency, Unbounded, O, AtLeastOnce>
2154 where
2155 L: TopLevel<'a>,
2156 {
2157 let samples = self.location.source_interval(interval);
2158 let (elements_hook, samples_hook) = nondet.take_hook();
2159
2160 let tick = self.location.tick();
2161 self.batch(
2162 &tick,
2163 nondet!(
2164 hook = elements_hook
2166 ),
2167 )
2168 .filter_if(
2169 samples
2170 .batch(
2171 &tick,
2172 nondet!(
2173 hook = samples_hook
2175 ),
2176 )
2177 .first()
2178 .is_some(),
2179 )
2180 .all_ticks()
2181 .weaken_retries()
2182 }
2183
2184 #[cfg(feature = "tokio")]
2194 pub fn timeout(
2195 self,
2196 duration: impl QuotedWithContext<
2197 'a,
2198 std::time::Duration,
2199 OperatorContext<Tick<L::DropConsistency>, Bounded>,
2200 > + Copy
2201 + 'a,
2202 nondet: NonDet,
2203 ) -> Optional<(), L::DropConsistency, Unbounded>
2204 where
2205 L: TopLevel<'a>,
2206 {
2207 let tick = self.location.tick();
2208
2209 let latest_received = self.assume_retries::<ExactlyOnce>(nondet).fold(
2210 q!(|| None),
2211 q!(
2212 |latest, _| {
2213 *latest = Some(Instant::now());
2214 },
2215 commutative = manual_proof!()
2216 ),
2217 );
2218
2219 latest_received
2220 .snapshot(
2221 &tick,
2222 nondet!(
2223 nondet
2225 ),
2226 )
2227 .filter_map(q!(move |latest_received| {
2228 if let Some(latest_received) = latest_received {
2229 if Instant::now().duration_since(latest_received) > duration {
2230 Some(())
2231 } else {
2232 None
2233 }
2234 } else {
2235 Some(())
2236 }
2237 }))
2238 .latest()
2239 }
2240
2241 pub fn atomic(self) -> Stream<T, Atomic<L>, B, O, R>
2247 where
2248 L: TopLevel<'a>,
2249 {
2250 let out_location = Atomic {
2251 tick: self.location.tick(),
2252 };
2253 Stream::new(
2254 out_location.clone(),
2255 HydroNode::BeginAtomic {
2256 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2257 metadata: out_location
2258 .new_node_metadata(Stream::<T, Atomic<L>, B, O, R>::collection_kind()),
2259 },
2260 )
2261 }
2262
2263 pub fn batch<L2: Location<'a, DropConsistency = L::DropConsistency>>(
2275 self,
2276 tick: &Tick<L2>,
2277 mut nondet: NonDet<Option<crate::sim_hooks::BatchHook<T, O, R, L::SimHookScope>>>,
2278 ) -> Stream<T, Tick<L::DropConsistency>, Bounded, O, R> {
2279 assert_eq!(
2280 Location::id(tick.parent_location()),
2281 Location::id(&self.location)
2282 );
2283
2284 let mut metadata =
2285 tick.new_node_metadata(Stream::<T, Tick<L>, Bounded, O, R>::collection_kind());
2286 metadata.op.sim_hook_id = nondet.take_hook().map(|h| h.id);
2287 Stream::new(
2288 tick.drop_consistency(),
2289 HydroNode::Batch {
2290 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2291 metadata,
2292 },
2293 )
2294 }
2295
2296 pub fn ir_node_named(self, name: &str) -> Stream<T, L, B, O, R> {
2299 {
2300 let mut node = self.ir_node.borrow_mut();
2301 let metadata = node.metadata_mut();
2302 metadata.tag = Some(name.to_owned());
2303 }
2304 self
2305 }
2306
2307 pub(crate) fn cast_at_most_one_element(self) -> Optional<T, L, B>
2311 where
2312 B: IsBounded,
2313 {
2314 Optional::new(
2315 self.location.clone(),
2316 HydroNode::Cast {
2317 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2318 metadata: self
2319 .location
2320 .new_node_metadata(Optional::<T, L, B>::collection_kind()),
2321 },
2322 )
2323 }
2324
2325 pub(crate) fn use_ordering_type<O2: Ordering>(self) -> Stream<T, L, B, O2, R> {
2326 if O::ORDERING_KIND == O2::ORDERING_KIND {
2327 Stream::new(
2328 self.location.clone(),
2329 self.ir_node.replace(HydroNode::Placeholder),
2330 )
2331 } else {
2332 panic!(
2333 "Runtime ordering {:?} did not match requested cast {:?}.",
2334 O::ORDERING_KIND,
2335 O2::ORDERING_KIND
2336 )
2337 }
2338 }
2339
2340 pub fn assume_ordering<O2: Ordering>(
2349 self,
2350 mut nondet: NonDet<Option<crate::sim_hooks::OrderingHook<T, B, L::SimHookScope>>>,
2351 ) -> Stream<T, L::DropConsistency, B, O2, R> {
2352 if O::ORDERING_KIND == O2::ORDERING_KIND {
2353 self.use_ordering_type().weaken_consistency()
2354 } else if O2::ORDERING_KIND == StreamOrder::NoOrder {
2355 let target_location = self.location().drop_consistency();
2357 Stream::new(
2358 target_location.clone(),
2359 HydroNode::Cast {
2360 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2361 metadata: target_location
2362 .new_node_metadata(Stream::<T, L, B, O2, R>::collection_kind()),
2363 },
2364 )
2365 } else {
2366 let target_location = self.location().drop_consistency();
2367 let mut metadata =
2368 target_location.new_node_metadata(Stream::<T, L, B, O2, R>::collection_kind());
2369 metadata.op.sim_hook_id = nondet.take_hook().map(|hook| hook.id);
2370 Stream::new(
2371 target_location,
2372 HydroNode::ObserveNonDet {
2373 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2374 trusted: false,
2375 metadata,
2376 },
2377 )
2378 }
2379 }
2380
2381 fn assume_ordering_trusted_bounded<O2: Ordering>(
2384 self,
2385 nondet: NonDet,
2386 ) -> Stream<T, L, B, O2, R> {
2387 if B::BOUNDED {
2388 self.assume_ordering_trusted(nondet)
2389 } else {
2390 let self_location = self.location.clone();
2391 let inner: Stream<T, L::DropConsistency, B, O2, R> = self.assume_ordering(nondet!(
2392 nondet
2394 ));
2395 Stream::new(self_location, inner.ir_node.replace(HydroNode::Placeholder))
2396 }
2397 }
2398
2399 pub(crate) fn assume_ordering_trusted<O2: Ordering>(
2402 self,
2403 _nondet: NonDet,
2404 ) -> Stream<T, L, B, O2, R> {
2405 if O::ORDERING_KIND == O2::ORDERING_KIND {
2406 self.use_ordering_type()
2407 } else if O2::ORDERING_KIND == StreamOrder::NoOrder {
2408 Stream::new(
2410 self.location.clone(),
2411 HydroNode::Cast {
2412 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2413 metadata: self
2414 .location
2415 .new_node_metadata(Stream::<T, L, B, O2, R>::collection_kind()),
2416 },
2417 )
2418 } else {
2419 Stream::new(
2420 self.location.clone(),
2421 HydroNode::ObserveNonDet {
2422 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2423 trusted: true,
2424 metadata: self
2425 .location
2426 .new_node_metadata(Stream::<T, L, B, O2, R>::collection_kind()),
2427 },
2428 )
2429 }
2430 }
2431
2432 #[deprecated = "use `weaken_ordering::<NoOrder>()` instead"]
2433 pub fn weakest_ordering(self) -> Stream<T, L, B, NoOrder, R> {
2436 self.weaken_ordering::<NoOrder>()
2437 }
2438
2439 pub fn weaken_ordering<O2: WeakerOrderingThan<O>>(self) -> Stream<T, L, B, O2, R> {
2442 let nondet = nondet!();
2443 self.assume_ordering_trusted::<O2>(nondet)
2444 }
2445
2446 pub fn make_totally_ordered(self) -> Stream<T, L, B, TotalOrder, R>
2449 where
2450 O: IsOrdered,
2451 {
2452 self.assume_ordering_trusted(nondet!())
2453 }
2454
2455 pub fn assume_retries<R2: Retries>(
2464 self,
2465 _nondet: NonDet,
2466 ) -> Stream<T, L::DropConsistency, B, O, R2> {
2467 if R::RETRIES_KIND == R2::RETRIES_KIND {
2468 Stream::new(
2469 self.location.drop_consistency(),
2470 self.ir_node.replace(HydroNode::Placeholder),
2471 )
2472 } else if R2::RETRIES_KIND == StreamRetry::AtLeastOnce {
2473 let target_location = self.location.drop_consistency();
2475 Stream::new(
2476 target_location.clone(),
2477 HydroNode::Cast {
2478 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2479 metadata: target_location
2480 .new_node_metadata(Stream::<T, L, B, O, R2>::collection_kind()),
2481 },
2482 )
2483 } else {
2484 let target_location = self.location.drop_consistency();
2485 Stream::new(
2486 target_location.clone(),
2487 HydroNode::ObserveNonDet {
2488 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2489 trusted: false,
2490 metadata: target_location
2491 .new_node_metadata(Stream::<T, L, B, O, R2>::collection_kind()),
2492 },
2493 )
2494 }
2495 }
2496
2497 fn assume_retries_trusted<R2: Retries>(self, _nondet: NonDet) -> Stream<T, L, B, O, R2> {
2500 if R::RETRIES_KIND == R2::RETRIES_KIND {
2501 Stream::new(
2502 self.location.clone(),
2503 self.ir_node.replace(HydroNode::Placeholder),
2504 )
2505 } else if R2::RETRIES_KIND == StreamRetry::AtLeastOnce {
2506 Stream::new(
2508 self.location.clone(),
2509 HydroNode::Cast {
2510 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2511 metadata: self
2512 .location
2513 .new_node_metadata(Stream::<T, L, B, O, R2>::collection_kind()),
2514 },
2515 )
2516 } else {
2517 Stream::new(
2518 self.location.clone(),
2519 HydroNode::ObserveNonDet {
2520 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2521 trusted: true,
2522 metadata: self
2523 .location
2524 .new_node_metadata(Stream::<T, L, B, O, R2>::collection_kind()),
2525 },
2526 )
2527 }
2528 }
2529
2530 #[deprecated = "use `weaken_retries::<AtLeastOnce>()` instead"]
2531 pub fn weakest_retries(self) -> Stream<T, L, B, O, AtLeastOnce> {
2534 self.weaken_retries::<AtLeastOnce>()
2535 }
2536
2537 pub fn weaken_retries<R2: WeakerRetryThan<R>>(self) -> Stream<T, L, B, O, R2> {
2540 let nondet = nondet!();
2541 self.assume_retries_trusted::<R2>(nondet)
2542 }
2543
2544 pub fn make_exactly_once(self) -> Stream<T, L, B, O, ExactlyOnce>
2547 where
2548 R: IsExactlyOnce,
2549 {
2550 self.assume_retries_trusted(nondet!())
2551 }
2552
2553 pub fn make_bounded(self) -> Stream<T, L, Bounded, O, R>
2556 where
2557 B: IsBounded,
2558 {
2559 self.weaken_boundedness()
2560 }
2561
2562 pub fn weaken_boundedness<B2: Boundedness>(self) -> Stream<T, L, B2, O, R> {
2565 if B::BOUNDED == B2::BOUNDED {
2566 Stream::new(
2567 self.location.clone(),
2568 self.ir_node.replace(HydroNode::Placeholder),
2569 )
2570 } else {
2571 Stream::new(
2573 self.location.clone(),
2574 HydroNode::Cast {
2575 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2576 metadata: self
2577 .location
2578 .new_node_metadata(Stream::<T, L, B2, O, R>::collection_kind()),
2579 },
2580 )
2581 }
2582 }
2583}
2584
2585impl<'a, T, L, B: Boundedness, O: Ordering, R: Retries> Stream<&T, L, B, O, R>
2586where
2587 L: Location<'a>,
2588{
2589 pub fn cloned(self) -> Stream<T, L, B, O, R>
2607 where
2608 T: Clone,
2609 {
2610 self.map(q!(|d| d.clone()))
2611 }
2612}
2613
2614impl<'a, T, L, B: Boundedness, O: Ordering> Stream<T, L, B, O, ExactlyOnce>
2615where
2616 L: Location<'a>,
2617{
2618 pub fn count(self) -> Singleton<usize, L, B::StreamToMonotone> {
2637 self.assume_ordering_trusted::<TotalOrder>(nondet!(
2638 ))
2640 .fold(
2641 q!(|| 0usize),
2642 q!(
2643 |count, _| *count += 1,
2644 monotone = manual_proof!()
2645 ),
2646 )
2647 }
2648}
2649
2650impl<'a, T, L: Location<'a>, O: Ordering, R: Retries> Stream<T, L, Unbounded, O, R> {
2651 pub fn merge_unordered<O2: Ordering, R2: Retries>(
2675 self,
2676 other: Stream<T, L, Unbounded, O2, R2>,
2677 ) -> Stream<T, L, Unbounded, NoOrder, <R as MinRetries<R2>>::Min>
2678 where
2679 R: MinRetries<R2>,
2680 {
2681 Stream::new(
2682 self.location.clone(),
2683 HydroNode::Chain {
2684 first: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2685 second: Box::new(other.ir_node.replace(HydroNode::Placeholder)),
2686 metadata: self.location.new_node_metadata(Stream::<
2687 T,
2688 L,
2689 Unbounded,
2690 NoOrder,
2691 <R as MinRetries<R2>>::Min,
2692 >::collection_kind()),
2693 },
2694 )
2695 }
2696
2697 #[deprecated(note = "use `merge_unordered` instead")]
2699 pub fn interleave<O2: Ordering, R2: Retries>(
2700 self,
2701 other: Stream<T, L, Unbounded, O2, R2>,
2702 ) -> Stream<T, L, Unbounded, NoOrder, <R as MinRetries<R2>>::Min>
2703 where
2704 R: MinRetries<R2>,
2705 {
2706 self.merge_unordered(other)
2707 }
2708}
2709
2710impl<'a, T, L: Location<'a>, B: Boundedness, R: Retries> Stream<T, L, B, TotalOrder, R> {
2711 pub fn merge_ordered<R2: Retries>(
2743 self,
2744 other: Stream<T, L, B, TotalOrder, R2>,
2745 mut nondet: NonDet<Option<crate::sim_hooks::MergeOrderedHook<T, B, L::SimHookScope>>>,
2746 ) -> Stream<T, L::DropConsistency, B, TotalOrder, <R as MinRetries<R2>>::Min>
2747 where
2748 R: MinRetries<R2>,
2749 {
2750 let target_location = self.location().drop_consistency();
2751 let mut metadata = target_location.new_node_metadata(Stream::<
2752 T,
2753 L::DropConsistency,
2754 B,
2755 TotalOrder,
2756 <R as MinRetries<R2>>::Min,
2757 >::collection_kind());
2758 metadata.op.sim_hook_id = nondet.take_hook().map(|hook| hook.id);
2759 Stream::new(
2760 target_location,
2761 HydroNode::MergeOrdered {
2762 first: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2763 second: Box::new(other.ir_node.replace(HydroNode::Placeholder)),
2764 metadata,
2765 },
2766 )
2767 }
2768}
2769
2770impl<'a, T, L, B: Boundedness, O: Ordering, R: Retries> Stream<T, L, B, O, R>
2771where
2772 L: Location<'a>,
2773{
2774 pub fn sort(self) -> Stream<T, L, Bounded, TotalOrder, R>
2800 where
2801 B: IsBounded,
2802 T: Ord,
2803 {
2804 let this = self.make_bounded();
2805 Stream::new(
2806 this.location.clone(),
2807 HydroNode::Sort {
2808 input: Box::new(this.ir_node.replace(HydroNode::Placeholder)),
2809 metadata: this
2810 .location
2811 .new_node_metadata(Stream::<T, L, Bounded, TotalOrder, R>::collection_kind()),
2812 },
2813 )
2814 }
2815
2816 pub fn chain<O2: Ordering, R2: Retries, B2: Boundedness>(
2844 self,
2845 other: Stream<T, L, B2, O2, R2>,
2846 ) -> Stream<T, L, B2, <O as MinOrder<O2>>::Min, <R as MinRetries<R2>>::Min>
2847 where
2848 B: IsBounded,
2849 O: MinOrder<O2>,
2850 R: MinRetries<R2>,
2851 {
2852 check_matching_location(&self.location, &other.location);
2853
2854 Stream::new(
2855 self.location.clone(),
2856 HydroNode::Chain {
2857 first: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
2858 second: Box::new(other.ir_node.replace(HydroNode::Placeholder)),
2859 metadata: self.location.new_node_metadata(Stream::<
2860 T,
2861 L,
2862 B2,
2863 <O as MinOrder<O2>>::Min,
2864 <R as MinRetries<R2>>::Min,
2865 >::collection_kind()),
2866 },
2867 )
2868 }
2869
2870 pub fn cross_product_nested_loop<T2, O2: Ordering + MinOrder<O>, R2: Retries>(
2874 self,
2875 other: Stream<T2, L, Bounded, O2, R2>,
2876 ) -> Stream<(T, T2), L, Bounded, <O2 as MinOrder<O>>::Min, <R as MinRetries<R2>>::Min>
2877 where
2878 B: IsBounded,
2879 T: Clone,
2880 T2: Clone,
2881 R: MinRetries<R2>,
2882 {
2883 let this = self.make_bounded();
2884 check_matching_location(&this.location, &other.location);
2885
2886 Stream::new(
2887 this.location.clone(),
2888 HydroNode::CrossProduct {
2889 left: Box::new(this.ir_node.replace(HydroNode::Placeholder)),
2890 right: Box::new(other.ir_node.replace(HydroNode::Placeholder)),
2891 metadata: this.location.new_node_metadata(Stream::<
2892 (T, T2),
2893 L,
2894 Bounded,
2895 <O2 as MinOrder<O>>::Min,
2896 <R as MinRetries<R2>>::Min,
2897 >::collection_kind()),
2898 },
2899 )
2900 }
2901
2902 pub fn repeat_with_keys<K, V2>(
2940 self,
2941 keys: KeyedSingleton<K, V2, L, Bounded>,
2942 ) -> KeyedStream<K, T, L, Bounded, O, R>
2943 where
2944 B: IsBounded,
2945 K: Clone,
2946 T: Clone,
2947 {
2948 keys.keys()
2949 .assume_ordering_trusted::<TotalOrder>(
2950 nondet!(),
2951 )
2952 .cross_product_nested_loop(self.make_bounded())
2953 .into_keyed()
2954 }
2955
2956 pub fn resolve_futures_blocking(self) -> Stream<T::Output, L, B, NoOrder, R>
2993 where
2994 T: Future,
2995 {
2996 Stream::new(
2997 self.location.clone(),
2998 HydroNode::ResolveFuturesBlocking {
2999 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
3000 metadata: self
3001 .location
3002 .new_node_metadata(Stream::<T::Output, L, B, NoOrder, R>::collection_kind()),
3003 },
3004 )
3005 }
3006
3007 #[expect(clippy::wrong_self_convention, reason = "stream function naming")]
3027 pub fn is_empty(self) -> Singleton<bool, L, Bounded>
3028 where
3029 B: IsBounded,
3030 {
3031 self.make_bounded()
3032 .assume_ordering_trusted::<TotalOrder>(
3033 nondet!(),
3034 )
3035 .first()
3036 .is_none()
3037 }
3038}
3039
3040impl<'a, K, V1, L, B: Boundedness, O: Ordering, R: Retries> Stream<(K, V1), L, B, O, R>
3041where
3042 L: Location<'a>,
3043{
3044 pub fn join<V2, B2: Boundedness, O2: Ordering, R2: Retries>(
3069 self,
3070 n: Stream<(K, V2), L, B2, O2, R2>,
3071 ) -> Stream<(K, (V1, V2)), L, B, B2::PreserveOrderIfBounded<O>, <R as MinRetries<R2>>::Min>
3072 where
3073 K: Eq + Hash + Clone,
3074 R: MinRetries<R2>,
3075 V1: Clone,
3076 V2: Clone,
3077 {
3078 check_matching_location(&self.location, &n.location);
3079
3080 let ir_node = if B2::BOUNDED {
3081 HydroNode::JoinHalf {
3082 left: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
3083 right: Box::new(n.ir_node.replace(HydroNode::Placeholder)),
3084 metadata: self.location.new_node_metadata(Stream::<
3085 (K, (V1, V2)),
3086 L,
3087 B,
3088 B2::PreserveOrderIfBounded<O>,
3089 <R as MinRetries<R2>>::Min,
3090 >::collection_kind()),
3091 }
3092 } else {
3093 HydroNode::Join {
3094 left: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
3095 right: Box::new(n.ir_node.replace(HydroNode::Placeholder)),
3096 metadata: self.location.new_node_metadata(Stream::<
3097 (K, (V1, V2)),
3098 L,
3099 B,
3100 B2::PreserveOrderIfBounded<O>,
3101 <R as MinRetries<R2>>::Min,
3102 >::collection_kind()),
3103 }
3104 };
3105
3106 Stream::new(self.location.clone(), ir_node)
3107 }
3108
3109 pub fn anti_join<O2: Ordering, R2: Retries>(
3135 self,
3136 n: Stream<K, L, Bounded, O2, R2>,
3137 ) -> Stream<(K, V1), L, B, O, R>
3138 where
3139 K: Eq + Hash,
3140 {
3141 check_matching_location(&self.location, &n.location);
3142
3143 Stream::new(
3144 self.location.clone(),
3145 HydroNode::AntiJoin {
3146 pos: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
3147 neg: Box::new(n.ir_node.replace(HydroNode::Placeholder)),
3148 metadata: self
3149 .location
3150 .new_node_metadata(Stream::<(K, V1), L, B, O, R>::collection_kind()),
3151 },
3152 )
3153 }
3154}
3155
3156impl<'a, K, V, L: Location<'a>, B: Boundedness, O: Ordering, R: Retries>
3157 Stream<(K, V), L, B, O, R>
3158{
3159 pub fn into_keyed(self) -> KeyedStream<K, V, L, B, O, R> {
3186 KeyedStream::new(
3187 self.location.clone(),
3188 HydroNode::Cast {
3189 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
3190 metadata: self
3191 .location
3192 .new_node_metadata(KeyedStream::<K, V, L, B, O, R>::collection_kind()),
3193 },
3194 )
3195 }
3196}
3197
3198impl<'a, K, V, L, O: Ordering, R: Retries> Stream<(K, V), Tick<L>, Bounded, O, R>
3199where
3200 K: Eq + Hash,
3201 L: Location<'a>,
3202{
3203 pub fn keys(self) -> Stream<K, Tick<L>, Bounded, NoOrder, ExactlyOnce> {
3222 self.into_keyed()
3223 .fold(
3224 q!(|| ()),
3225 q!(
3226 |_, _| {},
3227 commutative = manual_proof!(),
3228 idempotent = manual_proof!()
3229 ),
3230 )
3231 .keys()
3232 }
3233}
3234
3235impl<'a, T, L, B: Boundedness, O: Ordering, R: Retries> Stream<T, Atomic<L>, B, O, R>
3236where
3237 L: Location<'a>,
3238{
3239 pub fn batch_atomic<L2: Location<'a, DropConsistency = L::DropConsistency>>(
3246 self,
3247 tick: &Tick<L2>,
3248 mut nondet: NonDet<Option<crate::sim_hooks::BatchHook<T, O, R, L::SimHookScope>>>,
3249 ) -> Stream<T, Tick<L::DropConsistency>, Bounded, O, R> {
3250 assert_eq!(
3251 Location::id(tick.parent_location()),
3252 Location::id(self.location.tick.parent_location())
3253 );
3254
3255 let mut metadata =
3256 tick.new_node_metadata(Stream::<T, Tick<L>, Bounded, O, R>::collection_kind());
3257
3258 metadata.op.sim_hook_id = nondet.take_hook().map(|h| h.id);
3259 Stream::new(
3260 tick.drop_consistency(),
3261 HydroNode::Batch {
3262 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
3263 metadata,
3264 },
3265 )
3266 }
3267
3268 pub fn end_atomic(self) -> Stream<T, L, B, O, R> {
3271 Stream::new(
3272 self.location.tick.l.clone(),
3273 HydroNode::EndAtomic {
3274 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
3275 metadata: self
3276 .location
3277 .tick
3278 .l
3279 .new_node_metadata(Stream::<T, L, B, O, R>::collection_kind()),
3280 },
3281 )
3282 }
3283}
3284
3285impl<'a, F, T, L, B: Boundedness, O: Ordering, R: Retries> Stream<F, L, B, O, R>
3286where
3287 L: TopLevel<'a>,
3288 F: Future<Output = T>,
3289{
3290 pub fn resolve_futures(self) -> Stream<T, L, Unbounded, NoOrder, R> {
3321 Stream::new(
3322 self.location.clone(),
3323 HydroNode::ResolveFutures {
3324 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
3325 metadata: self
3326 .location
3327 .new_node_metadata(Stream::<T, L, Unbounded, NoOrder, R>::collection_kind()),
3328 },
3329 )
3330 }
3331
3332 pub fn resolve_futures_ordered(self) -> Stream<T, L, Unbounded, O, R> {
3363 Stream::new(
3364 self.location.clone(),
3365 HydroNode::ResolveFuturesOrdered {
3366 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
3367 metadata: self
3368 .location
3369 .new_node_metadata(Stream::<T, L, Unbounded, O, R>::collection_kind()),
3370 },
3371 )
3372 }
3373}
3374
3375impl<'a, T, L, O: Ordering, R: Retries> Stream<T, Tick<L>, Bounded, O, R>
3376where
3377 L: Location<'a>,
3378{
3379 pub fn all_ticks(self) -> Stream<T, L, Unbounded, O, R> {
3382 Stream::new(
3383 self.location.parent_location().clone(),
3384 HydroNode::YieldConcat {
3385 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
3386 metadata: self.location.parent_location().new_node_metadata(Stream::<
3387 T,
3388 L,
3389 Unbounded,
3390 O,
3391 R,
3392 >::collection_kind(
3393 )),
3394 },
3395 )
3396 }
3397
3398 pub fn all_ticks_atomic(self) -> Stream<T, Atomic<L>, Unbounded, O, R> {
3405 let out_location = Atomic {
3406 tick: self.location.clone(),
3407 };
3408
3409 Stream::new(
3410 out_location.clone(),
3411 HydroNode::YieldConcat {
3412 inner: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
3413 metadata: out_location
3414 .new_node_metadata(Stream::<T, Atomic<L>, Unbounded, O, R>::collection_kind()),
3415 },
3416 )
3417 }
3418
3419 pub fn across_ticks<Out: BatchAtomic<'a>>(
3454 self,
3455 thunk: impl FnOnce(Stream<T, Atomic<L>, Unbounded, O, R>) -> Out,
3456 ) -> Out::Batched {
3457 thunk(self.all_ticks_atomic()).batched_atomic()
3458 }
3459
3460 pub fn defer_tick(self) -> Stream<T, Tick<L>, Bounded, O, R> {
3499 Stream::new(
3500 self.location.clone(),
3501 HydroNode::DeferTick {
3502 input: Box::new(self.ir_node.replace(HydroNode::Placeholder)),
3503 metadata: self
3504 .location
3505 .new_node_metadata(Stream::<T, Tick<L>, Bounded, O, R>::collection_kind()),
3506 },
3507 )
3508 }
3509}
3510
3511#[cfg(test)]
3512mod tests {
3513 #[cfg(feature = "deploy")]
3514 use futures::{SinkExt, StreamExt};
3515 #[cfg(feature = "deploy")]
3516 use hydro_deploy::Deployment;
3517 #[cfg(feature = "deploy")]
3518 use serde::{Deserialize, Serialize};
3519 #[cfg(any(feature = "deploy", feature = "sim"))]
3520 use stageleft::q;
3521
3522 #[cfg(any(feature = "deploy", feature = "sim"))]
3523 use crate::compile::builder::FlowBuilder;
3524 #[cfg(feature = "deploy")]
3525 use crate::live_collections::sliced::sliced;
3526 #[cfg(feature = "deploy")]
3527 use crate::live_collections::stream::ExactlyOnce;
3528 #[cfg(feature = "sim")]
3529 use crate::live_collections::stream::NoOrder;
3530 #[cfg(any(feature = "deploy", feature = "sim"))]
3531 use crate::live_collections::stream::TotalOrder;
3532 #[cfg(any(feature = "deploy", feature = "sim"))]
3533 use crate::location::Location;
3534 #[cfg(feature = "sim")]
3535 use crate::networking::TCP;
3536 #[cfg(any(feature = "deploy", feature = "sim"))]
3537 use crate::nondet::nondet;
3538
3539 mod backtrace_chained_ops;
3540
3541 #[cfg(feature = "deploy")]
3542 struct P1 {}
3543 #[cfg(feature = "deploy")]
3544 struct P2 {}
3545
3546 #[cfg(feature = "deploy")]
3547 #[derive(Serialize, Deserialize, Debug)]
3548 struct SendOverNetwork {
3549 n: u32,
3550 }
3551
3552 #[cfg(feature = "deploy")]
3553 #[tokio::test]
3554 async fn first_ten_distributed() {
3555 use crate::networking::TCP;
3556
3557 let mut deployment = Deployment::new();
3558
3559 let mut flow = FlowBuilder::new();
3560 let first_node = flow.process::<P1>();
3561 let second_node = flow.process::<P2>();
3562 let external = flow.external::<P2>();
3563
3564 let numbers = first_node.source_iter(q!(0..10));
3565 let out_port = numbers
3566 .map(q!(|n| SendOverNetwork { n }))
3567 .send(&second_node, TCP.fail_stop().bincode())
3568 .send_bincode_external(&external);
3569
3570 let nodes = flow
3571 .with_process(&first_node, deployment.Localhost())
3572 .with_process(&second_node, deployment.Localhost())
3573 .with_external(&external, deployment.Localhost())
3574 .deploy(&mut deployment);
3575
3576 deployment.deploy().await.unwrap();
3577
3578 let mut external_out = nodes.connect(out_port).await;
3579
3580 deployment.start().await.unwrap();
3581
3582 for i in 0..10 {
3583 assert_eq!(external_out.next().await.unwrap().n, i);
3584 }
3585 }
3586
3587 #[cfg(feature = "deploy")]
3588 #[tokio::test]
3589 async fn first_cardinality() {
3590 let mut deployment = Deployment::new();
3591
3592 let mut flow = FlowBuilder::new();
3593 let node = flow.process::<()>();
3594 let external = flow.external::<()>();
3595
3596 let node_tick = node.tick();
3597 let count = node_tick
3598 .singleton(q!([1, 2, 3]))
3599 .into_stream()
3600 .flatten_ordered()
3601 .first()
3602 .into_stream()
3603 .count()
3604 .all_ticks()
3605 .send_bincode_external(&external);
3606
3607 let nodes = flow
3608 .with_process(&node, deployment.Localhost())
3609 .with_external(&external, deployment.Localhost())
3610 .deploy(&mut deployment);
3611
3612 deployment.deploy().await.unwrap();
3613
3614 let mut external_out = nodes.connect(count).await;
3615
3616 deployment.start().await.unwrap();
3617
3618 assert_eq!(external_out.next().await.unwrap(), 1);
3619 }
3620
3621 #[cfg(feature = "deploy")]
3622 #[tokio::test]
3623 async fn unbounded_reduce_remembers_state() {
3624 let mut deployment = Deployment::new();
3625
3626 let mut flow = FlowBuilder::new();
3627 let node = flow.process::<()>();
3628 let external = flow.external::<()>();
3629
3630 let (input_port, input) = node.source_external_bincode(&external);
3631 let out = input
3632 .reduce(q!(|acc, v| *acc += v))
3633 .sample_eager(nondet!())
3634 .send_bincode_external(&external);
3635
3636 let nodes = flow
3637 .with_process(&node, deployment.Localhost())
3638 .with_external(&external, deployment.Localhost())
3639 .deploy(&mut deployment);
3640
3641 deployment.deploy().await.unwrap();
3642
3643 let mut external_in = nodes.connect(input_port).await;
3644 let mut external_out = nodes.connect(out).await;
3645
3646 deployment.start().await.unwrap();
3647
3648 external_in.send(1).await.unwrap();
3649 assert_eq!(external_out.next().await.unwrap(), 1);
3650
3651 external_in.send(2).await.unwrap();
3652 assert_eq!(external_out.next().await.unwrap(), 3);
3653 }
3654
3655 #[cfg(feature = "deploy")]
3656 #[tokio::test]
3657 async fn top_level_bounded_cross_singleton() {
3658 let mut deployment = Deployment::new();
3659
3660 let mut flow = FlowBuilder::new();
3661 let node = flow.process::<()>();
3662 let external = flow.external::<()>();
3663
3664 let (input_port, input) =
3665 node.source_external_bincode::<_, _, TotalOrder, ExactlyOnce>(&external);
3666
3667 let out = input
3668 .cross_singleton(
3669 node.source_iter(q!(vec![1, 2, 3]))
3670 .fold(q!(|| 0), q!(|acc, v| *acc += v)),
3671 )
3672 .send_bincode_external(&external);
3673
3674 let nodes = flow
3675 .with_process(&node, deployment.Localhost())
3676 .with_external(&external, deployment.Localhost())
3677 .deploy(&mut deployment);
3678
3679 deployment.deploy().await.unwrap();
3680
3681 let mut external_in = nodes.connect(input_port).await;
3682 let mut external_out = nodes.connect(out).await;
3683
3684 deployment.start().await.unwrap();
3685
3686 external_in.send(1).await.unwrap();
3687 assert_eq!(external_out.next().await.unwrap(), (1, 6));
3688
3689 external_in.send(2).await.unwrap();
3690 assert_eq!(external_out.next().await.unwrap(), (2, 6));
3691 }
3692
3693 #[cfg(feature = "deploy")]
3694 #[tokio::test]
3695 async fn top_level_bounded_reduce_cardinality() {
3696 let mut deployment = Deployment::new();
3697
3698 let mut flow = FlowBuilder::new();
3699 let node = flow.process::<()>();
3700 let external = flow.external::<()>();
3701
3702 let (input_port, input) =
3703 node.source_external_bincode::<_, _, TotalOrder, ExactlyOnce>(&external);
3704
3705 let out = sliced! {
3706 let input = use::batch(input, nondet!());
3707 let v = use::snapshot(node.source_iter(q!(vec![1, 2, 3])).reduce(q!(|acc, v| *acc += v)), nondet!());
3708 input.cross_singleton(v.into_stream().count())
3709 }
3710 .send_bincode_external(&external);
3711
3712 let nodes = flow
3713 .with_process(&node, deployment.Localhost())
3714 .with_external(&external, deployment.Localhost())
3715 .deploy(&mut deployment);
3716
3717 deployment.deploy().await.unwrap();
3718
3719 let mut external_in = nodes.connect(input_port).await;
3720 let mut external_out = nodes.connect(out).await;
3721
3722 deployment.start().await.unwrap();
3723
3724 external_in.send(1).await.unwrap();
3725 assert_eq!(external_out.next().await.unwrap(), (1, 1));
3726
3727 external_in.send(2).await.unwrap();
3728 assert_eq!(external_out.next().await.unwrap(), (2, 1));
3729 }
3730
3731 #[cfg(feature = "deploy")]
3732 #[tokio::test]
3733 async fn top_level_bounded_into_singleton_cardinality() {
3734 let mut deployment = Deployment::new();
3735
3736 let mut flow = FlowBuilder::new();
3737 let node = flow.process::<()>();
3738 let external = flow.external::<()>();
3739
3740 let (input_port, input) =
3741 node.source_external_bincode::<_, _, TotalOrder, ExactlyOnce>(&external);
3742
3743 let out = sliced! {
3744 let input = use::batch(input, nondet!());
3745 let v = use::snapshot(node.source_iter(q!(vec![1, 2, 3])).reduce(q!(|acc, v| *acc += v)).into_singleton(), nondet!());
3746 input.cross_singleton(v.into_stream().count())
3747 }
3748 .send_bincode_external(&external);
3749
3750 let nodes = flow
3751 .with_process(&node, deployment.Localhost())
3752 .with_external(&external, deployment.Localhost())
3753 .deploy(&mut deployment);
3754
3755 deployment.deploy().await.unwrap();
3756
3757 let mut external_in = nodes.connect(input_port).await;
3758 let mut external_out = nodes.connect(out).await;
3759
3760 deployment.start().await.unwrap();
3761
3762 external_in.send(1).await.unwrap();
3763 assert_eq!(external_out.next().await.unwrap(), (1, 1));
3764
3765 external_in.send(2).await.unwrap();
3766 assert_eq!(external_out.next().await.unwrap(), (2, 1));
3767 }
3768
3769 #[cfg(feature = "deploy")]
3770 #[tokio::test]
3771 async fn atomic_fold_replays_each_tick() {
3772 let mut deployment = Deployment::new();
3773
3774 let mut flow = FlowBuilder::new();
3775 let node = flow.process::<()>();
3776 let external = flow.external::<()>();
3777
3778 let (input_port, input) =
3779 node.source_external_bincode::<_, _, TotalOrder, ExactlyOnce>(&external);
3780 let tick = node.tick();
3781
3782 let out = input
3783 .batch(&tick, nondet!())
3784 .cross_singleton(
3785 node.source_iter(q!(vec![1, 2, 3]))
3786 .atomic()
3787 .fold(q!(|| 0), q!(|acc, v| *acc += v))
3788 .snapshot_atomic(&tick, nondet!()),
3789 )
3790 .all_ticks()
3791 .send_bincode_external(&external);
3792
3793 let nodes = flow
3794 .with_process(&node, deployment.Localhost())
3795 .with_external(&external, deployment.Localhost())
3796 .deploy(&mut deployment);
3797
3798 deployment.deploy().await.unwrap();
3799
3800 let mut external_in = nodes.connect(input_port).await;
3801 let mut external_out = nodes.connect(out).await;
3802
3803 deployment.start().await.unwrap();
3804
3805 external_in.send(1).await.unwrap();
3806 assert_eq!(external_out.next().await.unwrap(), (1, 6));
3807
3808 external_in.send(2).await.unwrap();
3809 assert_eq!(external_out.next().await.unwrap(), (2, 6));
3810 }
3811
3812 #[cfg(feature = "deploy")]
3813 #[tokio::test]
3814 async fn unbounded_scan_remembers_state() {
3815 let mut deployment = Deployment::new();
3816
3817 let mut flow = FlowBuilder::new();
3818 let node = flow.process::<()>();
3819 let external = flow.external::<()>();
3820
3821 let (input_port, input) = node.source_external_bincode(&external);
3822 let out = input
3823 .scan(
3824 q!(|| 0),
3825 q!(|acc, v| {
3826 *acc += v;
3827 Some(*acc)
3828 }),
3829 )
3830 .send_bincode_external(&external);
3831
3832 let nodes = flow
3833 .with_process(&node, deployment.Localhost())
3834 .with_external(&external, deployment.Localhost())
3835 .deploy(&mut deployment);
3836
3837 deployment.deploy().await.unwrap();
3838
3839 let mut external_in = nodes.connect(input_port).await;
3840 let mut external_out = nodes.connect(out).await;
3841
3842 deployment.start().await.unwrap();
3843
3844 external_in.send(1).await.unwrap();
3845 assert_eq!(external_out.next().await.unwrap(), 1);
3846
3847 external_in.send(2).await.unwrap();
3848 assert_eq!(external_out.next().await.unwrap(), 3);
3849 }
3850
3851 #[cfg(feature = "deploy")]
3852 #[tokio::test]
3853 async fn unbounded_enumerate_remembers_state() {
3854 let mut deployment = Deployment::new();
3855
3856 let mut flow = FlowBuilder::new();
3857 let node = flow.process::<()>();
3858 let external = flow.external::<()>();
3859
3860 let (input_port, input) = node.source_external_bincode(&external);
3861 let out = input.enumerate().send_bincode_external(&external);
3862
3863 let nodes = flow
3864 .with_process(&node, deployment.Localhost())
3865 .with_external(&external, deployment.Localhost())
3866 .deploy(&mut deployment);
3867
3868 deployment.deploy().await.unwrap();
3869
3870 let mut external_in = nodes.connect(input_port).await;
3871 let mut external_out = nodes.connect(out).await;
3872
3873 deployment.start().await.unwrap();
3874
3875 external_in.send(1).await.unwrap();
3876 assert_eq!(external_out.next().await.unwrap(), (0, 1));
3877
3878 external_in.send(2).await.unwrap();
3879 assert_eq!(external_out.next().await.unwrap(), (1, 2));
3880 }
3881
3882 #[cfg(feature = "deploy")]
3883 #[tokio::test]
3884 async fn unbounded_unique_remembers_state() {
3885 let mut deployment = Deployment::new();
3886
3887 let mut flow = FlowBuilder::new();
3888 let node = flow.process::<()>();
3889 let external = flow.external::<()>();
3890
3891 let (input_port, input) =
3892 node.source_external_bincode::<_, _, TotalOrder, ExactlyOnce>(&external);
3893 let out = input.unique().send_bincode_external(&external);
3894
3895 let nodes = flow
3896 .with_process(&node, deployment.Localhost())
3897 .with_external(&external, deployment.Localhost())
3898 .deploy(&mut deployment);
3899
3900 deployment.deploy().await.unwrap();
3901
3902 let mut external_in = nodes.connect(input_port).await;
3903 let mut external_out = nodes.connect(out).await;
3904
3905 deployment.start().await.unwrap();
3906
3907 external_in.send(1).await.unwrap();
3908 assert_eq!(external_out.next().await.unwrap(), 1);
3909
3910 external_in.send(2).await.unwrap();
3911 assert_eq!(external_out.next().await.unwrap(), 2);
3912
3913 external_in.send(1).await.unwrap();
3914 external_in.send(3).await.unwrap();
3915 assert_eq!(external_out.next().await.unwrap(), 3);
3916 }
3917
3918 #[cfg(feature = "sim")]
3919 #[test]
3920 #[should_panic]
3921 fn sim_batch_nondet_size() {
3922 let mut flow = FlowBuilder::new();
3923 let node = flow.process::<()>();
3924
3925 let (in_send, input) = node.sim_input::<_, TotalOrder, _>();
3926
3927 let tick = node.tick();
3928 let out_recv = input
3929 .batch(&tick, nondet!())
3930 .count()
3931 .all_ticks()
3932 .sim_output();
3933
3934 flow.sim().exhaustive(async || {
3935 in_send.send(());
3936 in_send.send(());
3937 in_send.send(());
3938
3939 assert_eq!(out_recv.next().await, 3); });
3941 }
3942
3943 #[cfg(feature = "sim")]
3944 #[test]
3945 fn sim_batch_preserves_order() {
3946 let mut flow = FlowBuilder::new();
3947 let node = flow.process::<()>();
3948
3949 let (in_send, input) = node.sim_input();
3950
3951 let tick = node.tick();
3952 let out_recv = input
3953 .batch(&tick, nondet!())
3954 .all_ticks()
3955 .sim_output();
3956
3957 flow.sim().exhaustive(async || {
3958 in_send.send(1);
3959 in_send.send(2);
3960 in_send.send(3);
3961
3962 out_recv.assert_yields_only([1, 2, 3]).await;
3963 });
3964 }
3965
3966 #[cfg(feature = "sim")]
3967 #[test]
3968 #[should_panic]
3969 fn sim_batch_unordered_shuffles() {
3970 let mut flow = FlowBuilder::new();
3971 let node = flow.process::<()>();
3972
3973 let (in_send, input) = node.sim_input::<_, NoOrder, _>();
3974
3975 let tick = node.tick();
3976 let batch = input.batch(&tick, nondet!());
3977 let out_recv = batch
3978 .clone()
3979 .min()
3980 .zip(batch.max())
3981 .all_ticks()
3982 .sim_output();
3983
3984 flow.sim().exhaustive(async || {
3985 in_send.send_many_unordered([1, 2, 3]);
3986
3987 assert!(
3988 out_recv.collect::<Vec<_>>().await != vec![(1, 3), (2, 2)],
3989 "saw both (1, 3) and (2, 2), so batching must have shuffled the order"
3990 )
3991 });
3992 }
3993
3994 #[cfg(feature = "sim")]
3995 #[test]
3996 fn sim_batch_unordered_shuffles_count() {
3997 let mut flow = FlowBuilder::new();
3998 let node = flow.process::<()>();
3999
4000 let (in_send, input) = node.sim_input::<_, NoOrder, _>();
4001
4002 let tick = node.tick();
4003 let batch = input.batch(&tick, nondet!());
4004 let out_recv = batch.all_ticks().sim_output();
4005
4006 let instance_count = flow.sim().exhaustive(async || {
4007 in_send.send_many_unordered([1, 2, 3, 4]);
4008 out_recv.assert_yields_only_unordered([1, 2, 3, 4]).await;
4009 });
4010
4011 assert_eq!(
4012 instance_count,
4013 75 )
4015 }
4016
4017 #[cfg(feature = "sim")]
4018 #[test]
4019 #[should_panic]
4020 fn sim_observe_order_batched() {
4021 let mut flow = FlowBuilder::new();
4022 let node = flow.process::<()>();
4023
4024 let (in_send, input) = node.sim_input::<_, NoOrder, _>();
4025
4026 let tick = node.tick();
4027 let batch = input.batch(&tick, nondet!());
4028 let out_recv = batch
4029 .assume_ordering::<TotalOrder>(nondet!())
4030 .all_ticks()
4031 .sim_output();
4032
4033 flow.sim().exhaustive(async || {
4034 in_send.send_many_unordered([1, 2, 3, 4]);
4035 out_recv.assert_yields_only([1, 2, 3, 4]).await; });
4037 }
4038
4039 #[cfg(feature = "sim")]
4040 #[test]
4041 fn sim_observe_order_batched_count() {
4042 let mut flow = FlowBuilder::new();
4043 let node = flow.process::<()>();
4044
4045 let (in_send, input) = node.sim_input::<_, NoOrder, _>();
4046
4047 let tick = node.tick();
4048 let batch = input.batch(&tick, nondet!());
4049 let out_recv = batch
4050 .assume_ordering::<TotalOrder>(nondet!())
4051 .all_ticks()
4052 .sim_output();
4053
4054 let instance_count = flow.sim().exhaustive(async || {
4055 in_send.send_many_unordered([1, 2, 3, 4]);
4056 let _ = out_recv.collect::<Vec<_>>().await;
4057 });
4058
4059 assert_eq!(
4060 instance_count,
4061 192 )
4063 }
4064
4065 #[cfg(feature = "sim")]
4066 #[test]
4067 fn sim_unordered_count_instance_count() {
4068 let mut flow = FlowBuilder::new();
4069 let node = flow.process::<()>();
4070
4071 let (in_send, input) = node.sim_input::<_, NoOrder, _>();
4072
4073 let tick = node.tick();
4074 let out_recv = input
4075 .count()
4076 .snapshot(&tick, nondet!())
4077 .all_ticks()
4078 .sim_output();
4079
4080 let instance_count = flow.sim().exhaustive(async || {
4081 in_send.send_many_unordered([1, 2, 3, 4]);
4082 assert!(out_recv.collect::<Vec<_>>().await.last().unwrap() == &4);
4083 });
4084
4085 assert_eq!(
4086 instance_count,
4087 16 )
4089 }
4090
4091 #[cfg(feature = "sim")]
4092 #[test]
4093 fn sim_top_level_assume_ordering() {
4094 let mut flow = FlowBuilder::new();
4095 let node = flow.process::<()>();
4096
4097 let (in_send, input) = node.sim_input::<_, NoOrder, _>();
4098
4099 let out_recv = input
4100 .assume_ordering::<TotalOrder>(nondet!())
4101 .sim_output();
4102
4103 let instance_count = flow.sim().exhaustive(async || {
4104 in_send.send_many_unordered([1, 2, 3]);
4105 let mut out = out_recv.collect::<Vec<_>>().await;
4106 out.sort();
4107 assert_eq!(out, vec![1, 2, 3]);
4108 });
4109
4110 assert_eq!(instance_count, 6)
4111 }
4112
4113 #[cfg(feature = "sim")]
4114 #[test]
4115 fn sim_top_level_assume_ordering_cycle_back() {
4116 let mut flow = FlowBuilder::new();
4117 let node = flow.process::<()>();
4118 let node2 = flow.process::<()>();
4119
4120 let (in_send, input) = node.sim_input::<_, NoOrder, _>();
4121
4122 let (complete_cycle_back, cycle_back) =
4123 node.forward_ref::<super::Stream<_, _, _, NoOrder>>();
4124 let ordered = input
4125 .merge_unordered(cycle_back)
4126 .assume_ordering::<TotalOrder>(nondet!());
4127 complete_cycle_back.complete(
4128 ordered
4129 .clone()
4130 .map(q!(|v| v + 1))
4131 .filter(q!(|v| v % 2 == 1))
4132 .send(&node2, TCP.fail_stop().bincode())
4133 .send(&node, TCP.fail_stop().bincode()),
4134 );
4135
4136 let out_recv = ordered.sim_output();
4137
4138 let mut saw = false;
4139 let instance_count = flow.sim().exhaustive(async || {
4140 in_send.send_many_unordered([0, 2]);
4141 let out = out_recv.collect::<Vec<_>>().await;
4142
4143 if out.starts_with(&[0, 1, 2]) {
4144 saw = true;
4145 }
4146 });
4147
4148 assert!(saw, "did not see an instance with 0, 1, 2 in order");
4149 assert_eq!(instance_count, 6);
4150 }
4151
4152 #[cfg(feature = "sim")]
4153 #[test]
4154 fn sim_top_level_assume_ordering_cycle_back_tick() {
4155 let mut flow = FlowBuilder::new();
4156 let node = flow.process::<()>();
4157 let node2 = flow.process::<()>();
4158
4159 let (in_send, input) = node.sim_input::<_, NoOrder, _>();
4160
4161 let (complete_cycle_back, cycle_back) =
4162 node.forward_ref::<super::Stream<_, _, _, NoOrder>>();
4163 let ordered = input
4164 .merge_unordered(cycle_back)
4165 .assume_ordering::<TotalOrder>(nondet!());
4166 complete_cycle_back.complete(
4167 ordered
4168 .clone()
4169 .batch(&node.tick(), nondet!())
4170 .all_ticks()
4171 .map(q!(|v| v + 1))
4172 .filter(q!(|v| v % 2 == 1))
4173 .send(&node2, TCP.fail_stop().bincode())
4174 .send(&node, TCP.fail_stop().bincode()),
4175 );
4176
4177 let out_recv = ordered.sim_output();
4178
4179 let mut saw = false;
4180 let instance_count = flow.sim().exhaustive(async || {
4181 in_send.send_many_unordered([0, 2]);
4182 let out = out_recv.collect::<Vec<_>>().await;
4183
4184 if out.starts_with(&[0, 1, 2]) {
4185 saw = true;
4186 }
4187 });
4188
4189 assert!(saw, "did not see an instance with 0, 1, 2 in order");
4190 assert_eq!(instance_count, 58);
4191 }
4192
4193 #[cfg(feature = "sim")]
4194 #[test]
4195 fn sim_top_level_assume_ordering_multiple() {
4196 let mut flow = FlowBuilder::new();
4197 let node = flow.process::<()>();
4198 let node2 = flow.process::<()>();
4199
4200 let (in_send, input) = node.sim_input::<_, NoOrder, _>();
4201 let (_, input2) = node.sim_input::<_, NoOrder, _>();
4202
4203 let (complete_cycle_back, cycle_back) =
4204 node.forward_ref::<super::Stream<_, _, _, NoOrder>>();
4205 let input1_ordered = input
4206 .clone()
4207 .merge_unordered(cycle_back)
4208 .assume_ordering::<TotalOrder>(nondet!());
4209 let foo = input1_ordered
4210 .clone()
4211 .map(q!(|v| v + 3))
4212 .weaken_ordering::<NoOrder>()
4213 .merge_unordered(input2)
4214 .assume_ordering::<TotalOrder>(nondet!());
4215
4216 complete_cycle_back.complete(
4217 foo.filter(q!(|v| *v == 3))
4218 .send(&node2, TCP.fail_stop().bincode())
4219 .send(&node, TCP.fail_stop().bincode()),
4220 );
4221
4222 let out_recv = input1_ordered.sim_output();
4223
4224 let mut saw = false;
4225 let instance_count = flow.sim().exhaustive(async || {
4226 in_send.send_many_unordered([0, 1]);
4227 let out = out_recv.collect::<Vec<_>>().await;
4228
4229 if out.starts_with(&[0, 3, 1]) {
4230 saw = true;
4231 }
4232 });
4233
4234 assert!(saw, "did not see an instance with 0, 3, 1 in order");
4235 assert_eq!(instance_count, 15);
4236 }
4237
4238 #[cfg(feature = "sim")]
4239 #[test]
4240 fn sim_atomic_assume_ordering_cycle_back() {
4241 let mut flow = FlowBuilder::new();
4242 let node = flow.process::<()>();
4243 let node2 = flow.process::<()>();
4244
4245 let (in_send, input) = node.sim_input::<_, NoOrder, _>();
4246
4247 let (complete_cycle_back, cycle_back) =
4248 node.forward_ref::<super::Stream<_, _, _, NoOrder>>();
4249 let ordered = input
4250 .merge_unordered(cycle_back)
4251 .atomic()
4252 .assume_ordering::<TotalOrder>(nondet!())
4253 .end_atomic();
4254 complete_cycle_back.complete(
4255 ordered
4256 .clone()
4257 .map(q!(|v| v + 1))
4258 .filter(q!(|v| v % 2 == 1))
4259 .send(&node2, TCP.fail_stop().bincode())
4260 .send(&node, TCP.fail_stop().bincode()),
4261 );
4262
4263 let out_recv = ordered.sim_output();
4264
4265 let instance_count = flow.sim().exhaustive(async || {
4266 in_send.send_many_unordered([0, 2]);
4267 let out = out_recv.collect::<Vec<_>>().await;
4268 assert_eq!(out.len(), 4);
4269 });
4270 assert_eq!(instance_count, 22);
4271 }
4272
4273 #[cfg(feature = "deploy")]
4274 #[tokio::test]
4275 async fn partition_evens_odds() {
4276 let mut deployment = Deployment::new();
4277
4278 let mut flow = FlowBuilder::new();
4279 let node = flow.process::<()>();
4280 let external = flow.external::<()>();
4281
4282 let numbers = node.source_iter(q!(vec![1i32, 2, 3, 4, 5, 6]));
4283 let (evens, odds) = numbers.partition(q!(|x: &i32| x % 2 == 0));
4284 let evens_port = evens.send_bincode_external(&external);
4285 let odds_port = odds.send_bincode_external(&external);
4286
4287 let nodes = flow
4288 .with_process(&node, deployment.Localhost())
4289 .with_external(&external, deployment.Localhost())
4290 .deploy(&mut deployment);
4291
4292 deployment.deploy().await.unwrap();
4293
4294 let mut evens_out = nodes.connect(evens_port).await;
4295 let mut odds_out = nodes.connect(odds_port).await;
4296
4297 deployment.start().await.unwrap();
4298
4299 let mut even_results = Vec::new();
4300 for _ in 0..3 {
4301 even_results.push(evens_out.next().await.unwrap());
4302 }
4303 even_results.sort();
4304 assert_eq!(even_results, vec![2, 4, 6]);
4305
4306 let mut odd_results = Vec::new();
4307 for _ in 0..3 {
4308 odd_results.push(odds_out.next().await.unwrap());
4309 }
4310 odd_results.sort();
4311 assert_eq!(odd_results, vec![1, 3, 5]);
4312 }
4313
4314 #[cfg(feature = "deploy")]
4315 #[tokio::test]
4316 async fn unconsumed_inspect_still_runs() {
4317 use crate::deploy::DeployCrateWrapper;
4318
4319 let mut deployment = Deployment::new();
4320
4321 let mut flow = FlowBuilder::new();
4322 let node = flow.process::<()>();
4323
4324 node.source_iter(q!(0..5))
4327 .inspect(q!(|x| println!("inspect: {}", x)));
4328
4329 let nodes = flow
4330 .with_process(&node, deployment.Localhost())
4331 .deploy(&mut deployment);
4332
4333 deployment.deploy().await.unwrap();
4334
4335 let mut stdout = nodes.get_process(&node).stdout();
4336
4337 deployment.start().await.unwrap();
4338
4339 let mut lines = Vec::new();
4340 for _ in 0..5 {
4341 lines.push(stdout.recv().await.unwrap());
4342 }
4343 lines.sort();
4344 assert_eq!(
4345 lines,
4346 vec![
4347 "inspect: 0",
4348 "inspect: 1",
4349 "inspect: 2",
4350 "inspect: 3",
4351 "inspect: 4",
4352 ]
4353 );
4354 }
4355
4356 #[cfg(feature = "deploy")]
4357 #[tokio::test]
4358 async fn unconsumed_inspect_alive_at_deploy_still_runs() {
4359 use crate::deploy::DeployCrateWrapper;
4360
4361 let mut deployment = Deployment::new();
4362
4363 let mut flow = FlowBuilder::new();
4364 let node = flow.process::<()>();
4365
4366 let _inspected = node
4371 .source_iter(q!(0..5))
4372 .inspect(q!(|x| println!("inspect: {}", x)));
4373
4374 let nodes = flow
4375 .with_process(&node, deployment.Localhost())
4376 .deploy(&mut deployment);
4377
4378 deployment.deploy().await.unwrap();
4379
4380 let mut stdout = nodes.get_process(&node).stdout();
4381
4382 deployment.start().await.unwrap();
4383
4384 let mut lines = Vec::new();
4385 for _ in 0..5 {
4386 lines.push(stdout.recv().await.unwrap());
4387 }
4388 lines.sort();
4389 assert_eq!(
4390 lines,
4391 vec![
4392 "inspect: 0",
4393 "inspect: 1",
4394 "inspect: 2",
4395 "inspect: 3",
4396 "inspect: 4",
4397 ]
4398 );
4399 }
4400
4401 #[cfg(feature = "sim")]
4402 #[test]
4403 fn sim_limit() {
4404 let mut flow = FlowBuilder::new();
4405 let node = flow.process::<()>();
4406
4407 let (in_send, input) = node.sim_input();
4408
4409 let out_recv = input.limit(q!(3)).sim_output();
4410
4411 flow.sim().exhaustive(async || {
4412 in_send.send(1);
4413 in_send.send(2);
4414 in_send.send(3);
4415 in_send.send(4);
4416 in_send.send(5);
4417
4418 out_recv.assert_yields_only([1, 2, 3]).await;
4419 });
4420 }
4421
4422 #[cfg(feature = "sim")]
4423 #[test]
4424 fn sim_limit_zero() {
4425 let mut flow = FlowBuilder::new();
4426 let node = flow.process::<()>();
4427
4428 let (in_send, input) = node.sim_input();
4429
4430 let out_recv = input.limit(q!(0)).sim_output();
4431
4432 flow.sim().exhaustive(async || {
4433 in_send.send(1);
4434 in_send.send(2);
4435
4436 out_recv.assert_yields_only::<i32, _>([]).await;
4437 });
4438 }
4439
4440 #[cfg(feature = "sim")]
4441 #[test]
4442 fn sim_merge_ordered() {
4443 let mut flow = FlowBuilder::new();
4444 let node = flow.process::<()>();
4445
4446 let (in_send, input) = node.sim_input();
4447 let (in_send2, input2) = node.sim_input();
4448
4449 let out_recv = input
4450 .merge_ordered(input2, nondet!())
4451 .sim_output();
4452
4453 let mut saw_out_of_order = false;
4454 let instances = flow.sim().exhaustive(async || {
4455 in_send.send(1);
4456 in_send.send(2);
4457 in_send2.send(3);
4458 in_send2.send(4);
4459
4460 let out = out_recv.collect::<Vec<_>>().await;
4461
4462 if out == [1, 3, 2, 4] {
4463 saw_out_of_order = true;
4464 }
4465
4466 let mut first_elements = out.iter().filter(|v| **v <= 2).copied().collect::<Vec<_>>();
4469 let mut second_elements = out.iter().filter(|v| **v > 2).copied().collect::<Vec<_>>();
4470 assert_eq!(
4471 first_elements,
4472 vec![1, 2],
4473 "first input order violated: {:?}",
4474 out
4475 );
4476 assert_eq!(
4477 second_elements,
4478 vec![3, 4],
4479 "second input order violated: {:?}",
4480 out
4481 );
4482
4483 first_elements.append(&mut second_elements);
4484 first_elements.sort();
4485 assert_eq!(first_elements, vec![1, 2, 3, 4]);
4486 });
4487
4488 assert!(saw_out_of_order);
4489 assert_eq!(instances, 6);
4490 }
4491
4492 #[cfg(feature = "sim")]
4495 #[test]
4496 fn sim_merge_ordered_one_empty() {
4497 let mut flow = FlowBuilder::new();
4498 let node = flow.process::<()>();
4499
4500 let (in_send, input) = node.sim_input();
4501 let (_in_send2, input2) = node.sim_input();
4502
4503 let out_recv = input
4504 .merge_ordered(input2, nondet!())
4505 .sim_output();
4506
4507 let instances = flow.sim().exhaustive(async || {
4508 in_send.send(1);
4509 in_send.send(2);
4510
4511 let out = out_recv.collect::<Vec<_>>().await;
4512 assert_eq!(out, vec![1, 2]);
4513 });
4514
4515 assert_eq!(instances, 1);
4517 }
4518
4519 #[cfg(feature = "sim")]
4525 #[test]
4526 fn sim_merge_ordered_cycle_back() {
4527 let mut flow = FlowBuilder::new();
4528 let node = flow.process::<()>();
4529
4530 let (in_send, input) = node.sim_input();
4531
4532 let (complete_cycle_back, cycle_back) =
4534 node.forward_ref::<super::Stream<_, _, _, TotalOrder>>();
4535
4536 let merged = input.merge_ordered(cycle_back, nondet!());
4538
4539 complete_cycle_back.complete(merged.clone().filter(q!(|v| *v == 1)).map(q!(|v| v * 10)));
4541
4542 let out_recv = merged.sim_output();
4543
4544 let mut saw_cycle_before_second = false;
4547 flow.sim().exhaustive(async || {
4548 in_send.send(1);
4549 in_send.send(2);
4550
4551 let out = out_recv.collect::<Vec<_>>().await;
4552
4553 let pos_1 = out.iter().position(|v| *v == 1).unwrap();
4555 let pos_10 = out.iter().position(|v| *v == 10).unwrap();
4556 assert!(pos_1 < pos_10, "causal order violated: {:?}", out);
4557
4558 if out == [1, 10, 2] {
4560 saw_cycle_before_second = true;
4561 }
4562
4563 let mut sorted = out;
4564 sorted.sort();
4565 assert_eq!(sorted, vec![1, 2, 10]);
4566 });
4567
4568 assert!(
4569 saw_cycle_before_second,
4570 "never saw the cycled element arrive before the second input element"
4571 );
4572 }
4573
4574 #[cfg(feature = "sim")]
4578 #[test]
4579 fn sim_merge_ordered_delayed() {
4580 let mut flow = FlowBuilder::new();
4581 let node = flow.process::<()>();
4582
4583 let (in_send, input) = node.sim_input();
4584 let (in_send2, input2) = node.sim_input();
4585
4586 let out_recv = input
4587 .merge_ordered(input2, nondet!())
4588 .sim_output();
4589
4590 let mut saw_delayed_interleaving = false;
4591 flow.sim().exhaustive(async || {
4592 in_send.send(1);
4594 in_send2.send(3);
4595 in_send2.send(4);
4596
4597 let first_batch = out_recv.collect::<Vec<_>>().await;
4599
4600 in_send.send(2);
4602 let second_batch = out_recv.collect::<Vec<_>>().await;
4603
4604 let mut all: Vec<_> = first_batch
4605 .iter()
4606 .chain(second_batch.iter())
4607 .copied()
4608 .collect();
4609
4610 if all == [1, 3, 4, 2] {
4612 saw_delayed_interleaving = true;
4613 }
4614
4615 all.sort();
4616 assert_eq!(all, vec![1, 2, 3, 4]);
4617 });
4618
4619 assert!(saw_delayed_interleaving);
4620 }
4621
4622 #[cfg(feature = "deploy")]
4627 #[tokio::test]
4628 async fn deploy_merge_ordered_delayed() {
4629 let mut deployment = Deployment::new();
4630
4631 let mut flow = FlowBuilder::new();
4632 let node = flow.process::<()>();
4633 let external = flow.external::<()>();
4634
4635 let (input_a_port, input_a) = node.source_external_bincode(&external);
4636 let (input_b_port, input_b) = node.source_external_bincode(&external);
4637
4638 let out = input_a
4639 .assume_ordering(nondet!())
4640 .merge_ordered(
4641 input_b.assume_ordering(nondet!()),
4642 nondet!(),
4643 )
4644 .send_bincode_external(&external);
4645
4646 let nodes = flow
4647 .with_process(&node, deployment.Localhost())
4648 .with_external(&external, deployment.Localhost())
4649 .deploy(&mut deployment);
4650
4651 deployment.deploy().await.unwrap();
4652
4653 let mut ext_a = nodes.connect(input_a_port).await;
4654 let mut ext_b = nodes.connect(input_b_port).await;
4655 let mut ext_out = nodes.connect(out).await;
4656
4657 deployment.start().await.unwrap();
4658
4659 ext_a.send(1).await.unwrap();
4661 ext_b.send(3).await.unwrap();
4662 ext_b.send(4).await.unwrap();
4663
4664 let mut received = Vec::new();
4666 for _ in 0..3 {
4667 received.push(ext_out.next().await.unwrap());
4668 }
4669
4670 ext_a.send(2).await.unwrap();
4672 received.push(ext_out.next().await.unwrap());
4673
4674 received.sort();
4676 assert_eq!(received, vec![1, 2, 3, 4]);
4677 }
4678
4679 #[cfg(feature = "deploy")]
4680 #[tokio::test]
4681 async fn monotone_fold_threshold() {
4682 use crate::properties::manual_proof;
4683
4684 let mut deployment = Deployment::new();
4685
4686 let mut flow = FlowBuilder::new();
4687 let node = flow.process::<()>();
4688 let external = flow.external::<()>();
4689
4690 let in_unbounded: super::Stream<_, _> =
4691 node.source_iter(q!(vec![1i32, 2, 3, 4, 5, 6])).into();
4692 let sum = in_unbounded.fold(
4693 q!(|| 0),
4694 q!(
4695 |sum, v| {
4696 *sum += v;
4697 },
4698 monotone = manual_proof!()
4699 ),
4700 );
4701
4702 let threshold_out = sum
4703 .threshold_greater_or_equal(node.singleton(q!(7)))
4704 .send_bincode_external(&external);
4705
4706 let nodes = flow
4707 .with_process(&node, deployment.Localhost())
4708 .with_external(&external, deployment.Localhost())
4709 .deploy(&mut deployment);
4710
4711 deployment.deploy().await.unwrap();
4712
4713 let mut threshold_out = nodes.connect(threshold_out).await;
4714
4715 deployment.start().await.unwrap();
4716
4717 assert_eq!(threshold_out.next().await.unwrap(), 7);
4718 }
4719
4720 #[cfg(feature = "deploy")]
4721 #[tokio::test]
4722 async fn monotone_count_threshold() {
4723 let mut deployment = Deployment::new();
4724
4725 let mut flow = FlowBuilder::new();
4726 let node = flow.process::<()>();
4727 let external = flow.external::<()>();
4728
4729 let in_unbounded: super::Stream<_, _> =
4730 node.source_iter(q!(vec![1i32, 2, 3, 4, 5, 6])).into();
4731 let sum = in_unbounded.count();
4732
4733 let threshold_out = sum
4734 .threshold_greater_or_equal(node.singleton(q!(3)))
4735 .send_bincode_external(&external);
4736
4737 let nodes = flow
4738 .with_process(&node, deployment.Localhost())
4739 .with_external(&external, deployment.Localhost())
4740 .deploy(&mut deployment);
4741
4742 deployment.deploy().await.unwrap();
4743
4744 let mut threshold_out = nodes.connect(threshold_out).await;
4745
4746 deployment.start().await.unwrap();
4747
4748 assert_eq!(threshold_out.next().await.unwrap(), 3);
4749 }
4750
4751 #[cfg(feature = "deploy")]
4752 #[tokio::test]
4753 async fn monotone_map_order_preserving_threshold() {
4754 use crate::properties::manual_proof;
4755
4756 let mut deployment = Deployment::new();
4757
4758 let mut flow = FlowBuilder::new();
4759 let node = flow.process::<()>();
4760 let external = flow.external::<()>();
4761
4762 let in_unbounded: super::Stream<_, _> =
4763 node.source_iter(q!(vec![1i32, 2, 3, 4, 5, 6])).into();
4764 let sum = in_unbounded.fold(
4765 q!(|| 0),
4766 q!(
4767 |sum, v| {
4768 *sum += v;
4769 },
4770 monotone = manual_proof!()
4771 ),
4772 );
4773
4774 let doubled = sum.map(q!(
4776 |v| v * 2,
4777 order_preserving = manual_proof!()
4778 ));
4779
4780 let threshold_out = doubled
4781 .threshold_greater_or_equal(node.singleton(q!(14)))
4782 .send_bincode_external(&external);
4783
4784 let nodes = flow
4785 .with_process(&node, deployment.Localhost())
4786 .with_external(&external, deployment.Localhost())
4787 .deploy(&mut deployment);
4788
4789 deployment.deploy().await.unwrap();
4790
4791 let mut threshold_out = nodes.connect(threshold_out).await;
4792
4793 deployment.start().await.unwrap();
4794
4795 assert_eq!(threshold_out.next().await.unwrap(), 14);
4796 }
4797
4798 #[cfg(any(feature = "deploy", feature = "sim"))]
4801 mod join_ordering_type_tests {
4802 use crate::live_collections::boundedness::{Bounded, Unbounded};
4803 use crate::live_collections::stream::{ExactlyOnce, NoOrder, Stream, TotalOrder};
4804 use crate::location::{Location, Process};
4805
4806 #[expect(dead_code, reason = "compile-time type test")]
4807 fn join_unbounded_with_bounded_preserves_order<'a>(
4808 left: Stream<(i32, char), Process<'a>, Unbounded, TotalOrder, ExactlyOnce>,
4809 right: Stream<(i32, char), Process<'a>, Bounded, TotalOrder, ExactlyOnce>,
4810 ) -> Stream<(i32, (char, char)), Process<'a>, Unbounded, TotalOrder, ExactlyOnce> {
4811 left.join(right)
4812 }
4813
4814 #[expect(dead_code, reason = "compile-time type test")]
4815 fn join_unbounded_with_unbounded_is_no_order<'a>(
4816 left: Stream<(i32, char), Process<'a>, Unbounded, TotalOrder, ExactlyOnce>,
4817 right: Stream<(i32, char), Process<'a>, Unbounded, TotalOrder, ExactlyOnce>,
4818 ) -> Stream<(i32, (char, char)), Process<'a>, Unbounded, NoOrder, ExactlyOnce> {
4819 left.join(right)
4820 }
4821
4822 #[expect(dead_code, reason = "compile-time type test")]
4823 fn join_bounded_with_bounded_preserves_order<'a, L: Location<'a>>(
4824 left: Stream<(i32, char), L, Bounded, TotalOrder, ExactlyOnce>,
4825 right: Stream<(i32, char), L, Bounded, TotalOrder, ExactlyOnce>,
4826 ) -> Stream<(i32, (char, char)), L, Bounded, TotalOrder, ExactlyOnce> {
4827 left.join(right)
4828 }
4829
4830 #[expect(dead_code, reason = "compile-time type test")]
4831 fn join_unbounded_noorder_with_bounded<'a>(
4832 left: Stream<(i32, char), Process<'a>, Unbounded, NoOrder, ExactlyOnce>,
4833 right: Stream<(i32, char), Process<'a>, Bounded, NoOrder, ExactlyOnce>,
4834 ) -> Stream<(i32, (char, char)), Process<'a>, Unbounded, NoOrder, ExactlyOnce> {
4835 left.join(right)
4836 }
4837
4838 #[expect(dead_code, reason = "compile-time type test")]
4841 fn cross_product_unbounded_with_bounded_preserves_order<'a>(
4842 left: Stream<i32, Process<'a>, Unbounded, TotalOrder, ExactlyOnce>,
4843 right: Stream<char, Process<'a>, Bounded, TotalOrder, ExactlyOnce>,
4844 ) -> Stream<(i32, char), Process<'a>, Unbounded, TotalOrder, ExactlyOnce> {
4845 left.cross_product(right)
4846 }
4847
4848 #[expect(dead_code, reason = "compile-time type test")]
4849 fn cross_product_bounded_with_bounded_preserves_order<'a>(
4850 left: Stream<i32, Process<'a>, Bounded, TotalOrder, ExactlyOnce>,
4851 right: Stream<char, Process<'a>, Bounded, TotalOrder, ExactlyOnce>,
4852 ) -> Stream<(i32, char), Process<'a>, Bounded, TotalOrder, ExactlyOnce> {
4853 left.cross_product(right)
4854 }
4855
4856 #[expect(dead_code, reason = "compile-time type test")]
4857 fn cross_product_unbounded_with_unbounded_is_no_order<'a>(
4858 left: Stream<i32, Process<'a>, Unbounded, TotalOrder, ExactlyOnce>,
4859 right: Stream<char, Process<'a>, Unbounded, TotalOrder, ExactlyOnce>,
4860 ) -> Stream<(i32, char), Process<'a>, Unbounded, NoOrder, ExactlyOnce> {
4861 left.cross_product(right)
4862 }
4863 } #[cfg(feature = "sim")]
4868 #[test]
4869 fn cross_product_mixed_boundedness_correctness() {
4870 use stageleft::q;
4871
4872 use crate::compile::builder::FlowBuilder;
4873 use crate::nondet::nondet;
4874
4875 let mut flow = FlowBuilder::new();
4876 let process = flow.process::<()>();
4877 let tick = process.tick();
4878
4879 let left = process.source_iter(q!(vec![1, 2]));
4880 let right = process
4881 .source_iter(q!(vec!['a', 'b']))
4882 .batch(&tick, nondet!())
4883 .all_ticks();
4884
4885 let out = left.cross_product(right).sim_output();
4886
4887 flow.sim().exhaustive(async || {
4888 out.assert_yields_only_unordered(vec![(1, 'a'), (1, 'b'), (2, 'a'), (2, 'b')])
4889 .await;
4890 });
4891 }
4892
4893 #[cfg(feature = "sim")]
4894 #[test]
4895 fn join_mixed_boundedness_correctness() {
4896 use stageleft::q;
4897
4898 use crate::compile::builder::FlowBuilder;
4899 use crate::nondet::nondet;
4900
4901 let mut flow = FlowBuilder::new();
4902 let process = flow.process::<()>();
4903 let tick = process.tick();
4904
4905 let left = process.source_iter(q!(vec![(1, 'a'), (2, 'b')]));
4906 let right = process
4907 .source_iter(q!(vec![(1, 'x'), (2, 'y')]))
4908 .batch(&tick, nondet!())
4909 .all_ticks();
4910
4911 let out = left.join(right).sim_output();
4912
4913 flow.sim().exhaustive(async || {
4914 out.assert_yields_only_unordered(vec![(1, ('a', 'x')), (2, ('b', 'y'))])
4915 .await;
4916 });
4917 }
4918
4919 #[cfg(feature = "sim")]
4920 #[test]
4921 fn sim_merge_unordered_independent_atomics() {
4922 let mut flow = FlowBuilder::new();
4923 let node = flow.process::<()>();
4924
4925 let (in1_send, input1) = node.sim_input::<_, TotalOrder, _>();
4926 let (in2_send, input2) = node.sim_input::<_, TotalOrder, _>();
4927
4928 let out = input1
4929 .atomic()
4930 .merge_unordered(input2.atomic())
4931 .end_atomic()
4932 .sim_output();
4933
4934 flow.sim().exhaustive(async || {
4935 in1_send.send(1);
4936 in2_send.send(2);
4937
4938 out.assert_yields_only_unordered(vec![1, 2]).await;
4939 });
4940 }
4941
4942 #[cfg(feature = "deploy")]
4943 #[tokio::test]
4944 async fn test_stream_ref() {
4945 let mut deployment = Deployment::new();
4946
4947 let mut flow = FlowBuilder::new();
4948 let external = flow.external::<()>();
4949 let p1 = flow.process::<()>();
4950
4951 let my_stream = p1.source_iter(q!(1..=5i32));
4953
4954 let stream_ref = my_stream.by_ref();
4955
4956 let out_port = p1
4958 .source_iter(q!([()]))
4959 .map(q!(|_| stream_ref.len() as i32))
4960 .send_bincode_external(&external);
4961
4962 my_stream.for_each(q!(|_| {}));
4964
4965 let nodes = flow
4966 .with_default_optimize()
4967 .with_process(&p1, deployment.Localhost())
4968 .with_external(&external, deployment.Localhost())
4969 .deploy(&mut deployment);
4970
4971 deployment.deploy().await.unwrap();
4972
4973 let mut out_recv = nodes.connect(out_port).await;
4974
4975 deployment.start().await.unwrap();
4976
4977 let result = out_recv.next().await.unwrap();
4978 assert_eq!(result, 5);
4980 }
4981
4982 #[cfg(feature = "deploy")]
4983 #[tokio::test]
4984 async fn test_stream_ref_contents() {
4985 let mut deployment = Deployment::new();
4986
4987 let mut flow = FlowBuilder::new();
4988 let external = flow.external::<()>();
4989 let p1 = flow.process::<()>();
4990
4991 let my_stream = p1.source_iter(q!(1..=3i32));
4993
4994 let stream_ref = my_stream.by_ref();
4995
4996 let out_port = p1
4998 .source_iter(q!([()]))
4999 .map(q!(|_| stream_ref.iter().sum::<i32>()))
5000 .send_bincode_external(&external);
5001
5002 my_stream.for_each(q!(|_| {}));
5003
5004 let nodes = flow
5005 .with_default_optimize()
5006 .with_process(&p1, deployment.Localhost())
5007 .with_external(&external, deployment.Localhost())
5008 .deploy(&mut deployment);
5009
5010 deployment.deploy().await.unwrap();
5011
5012 let mut out_recv = nodes.connect(out_port).await;
5013
5014 deployment.start().await.unwrap();
5015
5016 let result = out_recv.next().await.unwrap();
5017 assert_eq!(result, 6);
5019 }
5020
5021 #[cfg(feature = "deploy")]
5022 #[tokio::test]
5023 async fn test_stream_ref_no_consumer() {
5024 let mut deployment = Deployment::new();
5025
5026 let mut flow = FlowBuilder::new();
5027 let external = flow.external::<()>();
5028 let p1 = flow.process::<()>();
5029
5030 let my_stream = p1.source_iter(q!(1..=4i32));
5032
5033 let stream_ref = my_stream.by_ref();
5034
5035 let out_port = p1
5036 .source_iter(q!([()]))
5037 .map(q!(|_| stream_ref.len() as i32))
5038 .send_bincode_external(&external);
5039
5040 let nodes = flow
5041 .with_default_optimize()
5042 .with_process(&p1, deployment.Localhost())
5043 .with_external(&external, deployment.Localhost())
5044 .deploy(&mut deployment);
5045
5046 deployment.deploy().await.unwrap();
5047
5048 let mut out_recv = nodes.connect(out_port).await;
5049
5050 deployment.start().await.unwrap();
5051
5052 let result = out_recv.next().await.unwrap();
5053 assert_eq!(result, 4);
5054 }
5055
5056 #[cfg(feature = "deploy")]
5057 #[tokio::test]
5058 async fn test_stream_mut() {
5059 let mut deployment = Deployment::new();
5060
5061 let mut flow = FlowBuilder::new();
5062 let external = flow.external::<()>();
5063 let p1 = flow.process::<()>();
5064
5065 let my_stream = p1.source_iter(q!(1..=5i32));
5067
5068 let stream_mut = my_stream.by_mut();
5069
5070 let out_port = p1
5072 .source_iter(q!([()]))
5073 .map(q!(|_| {
5074 stream_mut.retain(|x| *x > 3);
5075 stream_mut.len() as i32
5076 }))
5077 .send_bincode_external(&external);
5078
5079 my_stream.for_each(q!(|_| {}));
5080
5081 let nodes = flow
5082 .with_default_optimize()
5083 .with_process(&p1, deployment.Localhost())
5084 .with_external(&external, deployment.Localhost())
5085 .deploy(&mut deployment);
5086
5087 deployment.deploy().await.unwrap();
5088
5089 let mut out_recv = nodes.connect(out_port).await;
5090
5091 deployment.start().await.unwrap();
5092
5093 let result = out_recv.next().await.unwrap();
5094 assert_eq!(result, 2);
5096 }
5097
5098 #[cfg(feature = "sim")]
5102 #[test]
5103 fn sim_map_with_mut_on_unordered_explores_multiple_states() {
5104 use crate::live_collections::sliced::sliced;
5105 use crate::live_collections::stream::ExactlyOnce;
5106 use crate::properties::manual_proof;
5107
5108 let mut flow = FlowBuilder::new();
5109 let node = flow.process::<()>();
5110
5111 let (trigger_send, trigger) = node.sim_input::<i32, TotalOrder, ExactlyOnce>();
5112
5113 let out_recv = sliced! {
5114 let batch = use::batch(trigger, nondet!());
5115 let counter = batch.location().source_iter(q!(vec![0i32]))
5116 .fold(q!(|| 0i32), q!(|acc, v| *acc += v));
5117 let counter_mut = counter.by_mut();
5118 let items = batch.location().source_iter(q!(vec![1i32, 2])).weaken_ordering::<NoOrder>();
5119 items.map(q!(
5120 |x| {
5121 *counter_mut += x;
5122 *counter_mut
5123 },
5124 commutative = manual_proof!()
5125 ))
5126 }
5127 .sim_output();
5128
5129 let count = flow.sim().exhaustive(async || {
5130 trigger_send.send(1);
5131 let _all: Vec<i32> = out_recv.collect_sorted().await;
5132 });
5133
5134 assert_eq!(
5135 count, 2,
5136 "Expected 2 simulation instances due to mut on unordered input, got {}",
5137 count
5138 );
5139 }
5140
5141 #[cfg(feature = "sim")]
5145 #[test]
5146 fn sim_scan_with_ref_capture() {
5147 use crate::live_collections::sliced::sliced;
5148 use crate::live_collections::stream::ExactlyOnce;
5149
5150 let mut flow = FlowBuilder::new();
5151 let node = flow.process::<()>();
5152
5153 let (trigger_send, trigger) = node.sim_input::<i32, TotalOrder, ExactlyOnce>();
5154
5155 let out_recv = sliced! {
5156 let batch = use::batch(trigger, nondet!());
5157 let offset = batch
5158 .location()
5159 .source_iter(q!(vec![10i32]))
5160 .fold(q!(|| 0i32), q!(|acc, v| *acc += v));
5161 let offset_ref = offset.by_ref();
5162 batch
5163 .location()
5164 .source_iter(q!(vec![1i32, 2, 3]))
5165 .scan(
5166 q!(|| 0i32),
5167 q!(move |acc: &mut i32, x| {
5168 *acc += x + *offset_ref;
5169 Some(*acc)
5170 }),
5171 )
5172 }
5173 .sim_output();
5174
5175 let count = flow.sim().exhaustive(async || {
5176 trigger_send.send(1);
5177 let all: Vec<i32> = out_recv.collect().await;
5178 assert_eq!(all, vec![11, 23, 36]);
5183 });
5184
5185 assert_eq!(
5186 count, 1,
5187 "Expected a single simulation instance for a totally-ordered scan, got {}",
5188 count
5189 );
5190 }
5191
5192 #[cfg(feature = "sim")]
5196 #[test]
5197 #[ignore = "observe_nondet not yet supported for top-level bounded inputs (https://github.com/hydro-project/hydro/issues/2950)"]
5198 fn sim_map_with_mut_on_unordered_top_level() {
5199 use crate::properties::manual_proof;
5200
5201 let mut flow = FlowBuilder::new();
5202 let node = flow.process::<()>();
5203
5204 let counter = node
5205 .source_iter(q!(vec![0i32]))
5206 .fold(q!(|| 0i32), q!(|acc, v| *acc += v));
5207 let counter_mut = counter.by_mut();
5208
5209 let out_recv = node
5210 .source_iter(q!(vec![1i32, 2]))
5211 .weaken_ordering::<NoOrder>()
5212 .map(q!(
5213 |x| {
5214 *counter_mut += x;
5215 *counter_mut
5216 },
5217 commutative = manual_proof!()
5218 ))
5219 .assume_ordering::<TotalOrder>(nondet!())
5220 .sim_output();
5221
5222 counter.into_stream().for_each(q!(|_| {}));
5223
5224 let count = flow.sim().exhaustive(async || {
5225 let _all: Vec<i32> = out_recv.collect().await;
5226 });
5227
5228 assert_eq!(
5229 count, 2,
5230 "Expected 2 simulation instances due to mut on unordered input, got {}",
5231 count
5232 );
5233 }
5234}