Skip to main content

hydro_lang/live_collections/stream/
mod.rs

1//! Definitions for the [`Stream`] live collection.
2
3use 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/// A trait implemented by valid ordering markers ([`TotalOrder`] and [`NoOrder`]).
45#[sealed::sealed]
46pub trait Ordering:
47    MinOrder<Self, Min = Self> + MinOrder<TotalOrder, Min = Self> + MinOrder<NoOrder, Min = NoOrder>
48{
49    /// The [`StreamOrder`] corresponding to this type.
50    const ORDERING_KIND: StreamOrder;
51}
52
53/// Marks the stream as being totally ordered, which means that there are
54/// no sources of non-determinism (other than intentional ones) that will
55/// affect the order of elements.
56pub enum TotalOrder {}
57
58#[sealed::sealed]
59impl Ordering for TotalOrder {
60    const ORDERING_KIND: StreamOrder = StreamOrder::TotalOrder;
61}
62
63/// Marks the stream as having no order, which means that the order of
64/// elements may be affected by non-determinism.
65///
66/// This restricts certain operators, such as `fold` and `reduce`, to only
67/// be used with commutative aggregation functions.
68pub enum NoOrder {}
69
70#[sealed::sealed]
71impl Ordering for NoOrder {
72    const ORDERING_KIND: StreamOrder = StreamOrder::NoOrder;
73}
74
75/// Marker trait for an [`Ordering`] that is available when `Self` is a weaker guarantee than
76/// `Other`, which means that a stream with `Other` guarantees can be safely converted to
77/// have `Self` guarantees instead.
78#[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/// Helper trait for determining the weakest of two orderings.
84#[sealed::sealed]
85pub trait MinOrder<Other: ?Sized> {
86    /// The weaker of the two orderings.
87    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/// A trait implemented by valid retries markers ([`ExactlyOnce`] and [`AtLeastOnce`]).
101#[sealed::sealed]
102pub trait Retries:
103    MinRetries<Self, Min = Self>
104    + MinRetries<ExactlyOnce, Min = Self>
105    + MinRetries<AtLeastOnce, Min = AtLeastOnce>
106{
107    /// The [`StreamRetry`] corresponding to this type.
108    const RETRIES_KIND: StreamRetry;
109}
110
111/// Marks the stream as having deterministic message cardinality, with no
112/// possibility of duplicates.
113pub enum ExactlyOnce {}
114
115#[sealed::sealed]
116impl Retries for ExactlyOnce {
117    const RETRIES_KIND: StreamRetry = StreamRetry::ExactlyOnce;
118}
119
120/// Marks the stream as having non-deterministic message cardinality, which
121/// means that duplicates may occur, but messages will not be dropped.
122pub enum AtLeastOnce {}
123
124#[sealed::sealed]
125impl Retries for AtLeastOnce {
126    const RETRIES_KIND: StreamRetry = StreamRetry::AtLeastOnce;
127}
128
129/// Marker trait for a [`Retries`] that is available when `Self` is a weaker guarantee than
130/// `Other`, which means that a stream with `Other` guarantees can be safely converted to
131/// have `Self` guarantees instead.
132#[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/// Helper trait for determining the weakest of two retry guarantees.
138#[sealed::sealed]
139pub trait MinRetries<Other: ?Sized> {
140    /// The weaker of the two retry guarantees.
141    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)]
160/// Marker trait that is implemented for the [`TotalOrder`] ordering guarantee.
161pub 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)]
173/// Marker trait that is implemented for the [`ExactlyOnce`] retries guarantee.
174pub trait IsExactlyOnce: Retries {}
175
176#[sealed::sealed]
177#[diagnostic::do_not_recommend]
178impl IsExactlyOnce for ExactlyOnce {}
179
180/// Streaming sequence of elements with type `Type`.
181///
182/// This live collection represents a growing sequence of elements, with new elements being
183/// asynchronously appended to the end of the sequence. This can be used to model the arrival
184/// of network input, such as API requests, or streaming ingestion.
185///
186/// By default, all streams have deterministic ordering and each element is materialized exactly
187/// once. But streams can also capture non-determinism via the `Order` and `Retries` type
188/// parameters. When the ordering / retries guarantee is relaxed, fewer APIs will be available
189/// on the stream. For example, if the stream is unordered, you cannot invoke [`Stream::first`].
190///
191/// Type Parameters:
192/// - `Type`: the type of elements in the stream
193/// - `Loc`: the location where the stream is being materialized
194/// - `Bound`: the boundedness of the stream, which is either [`Bounded`] or [`Unbounded`]
195/// - `Order`: the ordering of the stream, which is either [`TotalOrder`] or [`NoOrder`]
196///   (default is [`TotalOrder`])
197/// - `Retries`: the retry guarantee of the stream, which is either [`ExactlyOnce`] or
198///   [`AtLeastOnce`] (default is [`ExactlyOnce`])
199pub 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    /// Returns the [`Location`] where this stream is being materialized.
438    pub fn location(&self) -> &L {
439        &self.location
440    }
441
442    /// Creates a shared reference handle to this stream's handoff buffer that can be captured
443    /// inside `q!()` closures. The handle resolves to `&Vec<T>` at runtime.
444    ///
445    /// The stream must be bounded, otherwise reading it would be non-deterministic.
446    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    /// Returns a mutable reference handle to this stream's handoff buffer that can be captured
454    /// inside `q!()` closures. The handle resolves to `&mut Vec<T>` at runtime.
455    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    /// Weakens the consistency of this live collection to not guarantee any consistency across
463    /// cluster members (if this collection is on a cluster).
464    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            // already no consistency
472            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    /// Casts this live collection to have the consistency guarantees specified in the given
495    /// location type parameter. The developer must ensure that the strengthened consistency
496    /// is actually guaranteed, via the proof field (see [`crate::prelude::manual_proof`]).
497    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    /// Produces a stream based on invoking `f` on each element.
565    /// If you do not want to modify the stream and instead only want to view
566    /// each item use [`Stream::inspect`] instead.
567    ///
568    /// If the input stream is unordered **and** `f` mutably captures state (such as a
569    /// [`Singleton::by_mut`](crate::live_collections::singleton::Singleton::by_mut)
570    /// reference), `f` must be proven **commutative**: processing any two elements in
571    /// either order must leave the mutably-captured state in the same final value *and*
572    /// produce the same multiset of outputs. In particular, outputs must not expose the
573    /// processing order (e.g. emitting a running total is not commutative even though
574    /// addition is).
575    ///
576    /// # Example
577    /// ```rust
578    /// # #[cfg(feature = "deploy")] {
579    /// # use hydro_lang::prelude::*;
580    /// # use futures::StreamExt;
581    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
582    /// let words = process.source_iter(q!(vec!["hello", "world"]));
583    /// words.map(q!(|x| x.to_uppercase()))
584    /// # }, |mut stream| async move {
585    /// # for w in vec!["HELLO", "WORLD"] {
586    /// #     assert_eq!(stream.next().await.unwrap(), w);
587    /// # }
588    /// # }));
589    /// # }
590    /// ```
591    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    /// For each item `i` in the input stream, transform `i` using `f` and then treat the
624    /// result as an [`Iterator`] to produce items one by one. The implementation for [`Iterator`]
625    /// for the output type `U` must produce items in a **deterministic** order.
626    ///
627    /// For example, `U` could be a `Vec`, but not a `HashSet`. If the order of the items in `U` is
628    /// not deterministic, use [`Stream::flat_map_unordered`] instead.
629    ///
630    /// # Example
631    /// ```rust
632    /// # #[cfg(feature = "deploy")] {
633    /// # use hydro_lang::prelude::*;
634    /// # use futures::StreamExt;
635    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
636    /// process
637    ///     .source_iter(q!(vec![vec![1, 2], vec![3, 4]]))
638    ///     .flat_map_ordered(q!(|x| x))
639    /// # }, |mut stream| async move {
640    /// // 1, 2, 3, 4
641    /// # for w in (1..5) {
642    /// #     assert_eq!(stream.next().await.unwrap(), w);
643    /// # }
644    /// # }));
645    /// # }
646    /// ```
647    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    /// Like [`Stream::flat_map_ordered`], but allows the implementation of [`Iterator`]
681    /// for the output type `U` to produce items in any order.
682    ///
683    /// # Example
684    /// ```rust
685    /// # #[cfg(feature = "deploy")] {
686    /// # use hydro_lang::{prelude::*, live_collections::stream::{NoOrder, ExactlyOnce}};
687    /// # use futures::StreamExt;
688    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test::<_, _, _, NoOrder, ExactlyOnce>(|process| {
689    /// process
690    ///     .source_iter(q!(vec![
691    ///         std::collections::HashSet::<i32>::from_iter(vec![1, 2]),
692    ///         std::collections::HashSet::from_iter(vec![3, 4]),
693    ///     ]))
694    ///     .flat_map_unordered(q!(|x| x))
695    /// # }, |mut stream| async move {
696    /// // 1, 2, 3, 4, but in no particular order
697    /// # let mut results = Vec::new();
698    /// # for w in (1..5) {
699    /// #     results.push(stream.next().await.unwrap());
700    /// # }
701    /// # results.sort();
702    /// # assert_eq!(results, vec![1, 2, 3, 4]);
703    /// # }));
704    /// # }
705    /// ```
706    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    /// For each item `i` in the input stream, treat `i` as an [`Iterator`] and produce its items one by one.
740    /// The implementation for [`Iterator`] for the element type `T` must produce items in a **deterministic** order.
741    ///
742    /// For example, `T` could be a `Vec`, but not a `HashSet`. If the order of the items in `T` is
743    /// not deterministic, use [`Stream::flatten_unordered`] instead.
744    ///
745    /// ```rust
746    /// # #[cfg(feature = "deploy")] {
747    /// # use hydro_lang::prelude::*;
748    /// # use futures::StreamExt;
749    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
750    /// process
751    ///     .source_iter(q!(vec![vec![1, 2], vec![3, 4]]))
752    ///     .flatten_ordered()
753    /// # }, |mut stream| async move {
754    /// // 1, 2, 3, 4
755    /// # for w in (1..5) {
756    /// #     assert_eq!(stream.next().await.unwrap(), w);
757    /// # }
758    /// # }));
759    /// # }
760    /// ```
761    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    /// Like [`Stream::flatten_ordered`], but allows the implementation of [`Iterator`]
769    /// for the element type `T` to produce items in any order.
770    ///
771    /// # Example
772    /// ```rust
773    /// # #[cfg(feature = "deploy")] {
774    /// # use hydro_lang::{prelude::*, live_collections::stream::{NoOrder, ExactlyOnce}};
775    /// # use futures::StreamExt;
776    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test::<_, _, _, NoOrder, ExactlyOnce>(|process| {
777    /// process
778    ///     .source_iter(q!(vec![
779    ///         std::collections::HashSet::<i32>::from_iter(vec![1, 2]),
780    ///         std::collections::HashSet::from_iter(vec![3, 4]),
781    ///     ]))
782    ///     .flatten_unordered()
783    /// # }, |mut stream| async move {
784    /// // 1, 2, 3, 4, but in no particular order
785    /// # let mut results = Vec::new();
786    /// # for w in (1..5) {
787    /// #     results.push(stream.next().await.unwrap());
788    /// # }
789    /// # results.sort();
790    /// # assert_eq!(results, vec![1, 2, 3, 4]);
791    /// # }));
792    /// # }
793    /// ```
794    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    /// For each item in the input stream, apply `f` to produce a [`futures::stream::Stream`],
802    /// then emit the elements of that stream one by one. When the inner stream yields
803    /// `Pending`, this operator yields as well.
804    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    /// For each item in the input stream, treat it as a [`futures::stream::Stream`] and
838    /// emit its elements one by one. When the inner stream yields `Pending`, this operator
839    /// yields as well.
840    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    /// Creates a stream containing only the elements of the input stream that satisfy a predicate
848    /// `f`, preserving the order of the elements.
849    ///
850    /// The closure `f` receives a reference `&T` rather than an owned value `T` because filtering does
851    /// not modify or take ownership of the values. If you need to modify the values while filtering
852    /// use [`Stream::filter_map`] instead.
853    ///
854    /// If the input stream is unordered **and** `f` mutably captures state, `f` must be
855    /// proven **commutative**: processing any two elements in either order must leave
856    /// the mutably-captured state in the same final value *and* retain the same multiset
857    /// of elements. In particular, the decision for each element must not depend on the
858    /// processing order: a stateful predicate like a rate limiter is not commutative —
859    /// its budget converges either way, but *which* element passes depends on the order.
860    ///
861    /// # Example
862    /// ```rust
863    /// # #[cfg(feature = "deploy")] {
864    /// # use hydro_lang::prelude::*;
865    /// # use futures::StreamExt;
866    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
867    /// process
868    ///     .source_iter(q!(vec![1, 2, 3, 4]))
869    ///     .filter(q!(|&x| x > 2))
870    /// # }, |mut stream| async move {
871    /// // 3, 4
872    /// # for w in (3..5) {
873    /// #     assert_eq!(stream.next().await.unwrap(), w);
874    /// # }
875    /// # }));
876    /// # }
877    /// ```
878    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    /// Splits the stream into two streams based on a predicate, without cloning elements.
909    ///
910    /// Elements for which `f` returns `true` are sent to the first output stream,
911    /// and elements for which `f` returns `false` are sent to the second output stream.
912    ///
913    /// Unlike using `filter` twice, this only evaluates the predicate once per element
914    /// and does not require `T: Clone`.
915    ///
916    /// The closure `f` receives a reference `&T` rather than an owned value `T` because
917    /// the predicate is only used for routing; the element itself is moved to the
918    /// appropriate output stream.
919    ///
920    /// # Example
921    /// ```rust
922    /// # #[cfg(feature = "deploy")] {
923    /// # use hydro_lang::prelude::*;
924    /// # use hydro_lang::live_collections::stream::{NoOrder, ExactlyOnce};
925    /// # use futures::StreamExt;
926    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test::<_, _, _, NoOrder, ExactlyOnce>(|process| {
927    /// let numbers: Stream<_, _, Unbounded> = process.source_iter(q!(vec![1, 2, 3, 4, 5, 6])).into();
928    /// let (evens, odds) = numbers.partition(q!(|&x| x % 2 == 0));
929    /// // evens: 2, 4, 6 tagged with true; odds: 1, 3, 5 tagged with false
930    /// evens.map(q!(|x| (x, true)))
931    ///     .merge_unordered(odds.map(q!(|x| (x, false))))
932    /// # }, |mut stream| async move {
933    /// # let mut results = Vec::new();
934    /// # for _ in 0..6 {
935    /// #     results.push(stream.next().await.unwrap());
936    /// # }
937    /// # results.sort();
938    /// # assert_eq!(results, vec![(1, false), (2, true), (3, false), (4, true), (5, false), (6, true)]);
939    /// # }));
940    /// # }
941    /// ```
942    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    /// An operator that both filters and maps. It yields only the items for which the supplied closure `f` returns `Some(value)`.
990    ///
991    /// # Example
992    /// ```rust
993    /// # #[cfg(feature = "deploy")] {
994    /// # use hydro_lang::prelude::*;
995    /// # use futures::StreamExt;
996    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
997    /// process
998    ///     .source_iter(q!(vec!["1", "hello", "world", "2"]))
999    ///     .filter_map(q!(|s| s.parse::<usize>().ok()))
1000    /// # }, |mut stream| async move {
1001    /// // 1, 2
1002    /// # for w in (1..3) {
1003    /// #     assert_eq!(stream.next().await.unwrap(), w);
1004    /// # }
1005    /// # }));
1006    /// # }
1007    /// ```
1008    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    /// Generates a stream that maps each input element `i` to a tuple `(i, x)`,
1041    /// where `x` is the final value of `other`, a bounded [`Singleton`] or [`Optional`].
1042    /// If `other` is an empty [`Optional`], no values will be produced.
1043    ///
1044    /// # Example
1045    /// ```rust
1046    /// # #[cfg(feature = "deploy")] {
1047    /// # use hydro_lang::prelude::*;
1048    /// # use futures::StreamExt;
1049    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1050    /// let tick = process.tick();
1051    /// let batch = process
1052    ///   .source_iter(q!(vec![1, 2, 3, 4]))
1053    ///   .batch(&tick, nondet!(/** test */));
1054    /// let count = batch.clone().count(); // `count()` returns a singleton
1055    /// batch.cross_singleton(count).all_ticks()
1056    /// # }, |mut stream| async move {
1057    /// // (1, 4), (2, 4), (3, 4), (4, 4)
1058    /// # for w in vec![(1, 4), (2, 4), (3, 4), (4, 4)] {
1059    /// #     assert_eq!(stream.next().await.unwrap(), w);
1060    /// # }
1061    /// # }));
1062    /// # }
1063    /// ```
1064    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    /// Passes this stream through if the boolean signal is `true`, otherwise the output is empty.
1087    ///
1088    /// # Example
1089    /// ```rust
1090    /// # #[cfg(feature = "deploy")] {
1091    /// # use hydro_lang::prelude::*;
1092    /// # use futures::StreamExt;
1093    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1094    /// let tick = process.tick();
1095    /// // ticks are lazy by default, forces the second tick to run
1096    /// tick.spin_batch(q!(1)).all_ticks().for_each(q!(|_| {}));
1097    ///
1098    /// let signal = tick.optional_first_tick(q!(())).is_some(); // true on tick 1, false on tick 2
1099    /// let batch_first_tick = process
1100    ///   .source_iter(q!(vec![1, 2, 3, 4]))
1101    ///   .batch(&tick, nondet!(/** test */));
1102    /// let batch_second_tick = process
1103    ///   .source_iter(q!(vec![5, 6, 7, 8]))
1104    ///   .batch(&tick, nondet!(/** test */))
1105    ///   .defer_tick();
1106    /// batch_first_tick.chain(batch_second_tick)
1107    ///   .filter_if(signal)
1108    ///   .all_ticks()
1109    /// # }, |mut stream| async move {
1110    /// // [1, 2, 3, 4]
1111    /// # for w in vec![1, 2, 3, 4] {
1112    /// #     assert_eq!(stream.next().await.unwrap(), w);
1113    /// # }
1114    /// # }));
1115    /// # }
1116    /// ```
1117    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    /// Passes this stream through if the argument (a [`Bounded`] [`Optional`]) is non-null, otherwise the output is empty.
1123    ///
1124    /// Useful for gating the release of elements based on a condition, such as only processing requests if you are the
1125    /// leader of a cluster.
1126    ///
1127    /// # Example
1128    /// ```rust
1129    /// # #[cfg(feature = "deploy")] {
1130    /// # use hydro_lang::prelude::*;
1131    /// # use futures::StreamExt;
1132    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1133    /// let tick = process.tick();
1134    /// // ticks are lazy by default, forces the second tick to run
1135    /// tick.spin_batch(q!(1)).all_ticks().for_each(q!(|_| {}));
1136    ///
1137    /// let batch_first_tick = process
1138    ///   .source_iter(q!(vec![1, 2, 3, 4]))
1139    ///   .batch(&tick, nondet!(/** test */));
1140    /// let batch_second_tick = process
1141    ///   .source_iter(q!(vec![5, 6, 7, 8]))
1142    ///   .batch(&tick, nondet!(/** test */))
1143    ///   .defer_tick(); // appears on the second tick
1144    /// let some_on_first_tick = tick.optional_first_tick(q!(()));
1145    /// batch_first_tick.chain(batch_second_tick)
1146    ///   .filter_if_some(some_on_first_tick)
1147    ///   .all_ticks()
1148    /// # }, |mut stream| async move {
1149    /// // [1, 2, 3, 4]
1150    /// # for w in vec![1, 2, 3, 4] {
1151    /// #     assert_eq!(stream.next().await.unwrap(), w);
1152    /// # }
1153    /// # }));
1154    /// # }
1155    /// ```
1156    #[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    /// Passes this stream through if the argument (a [`Bounded`] [`Optional`]) is null, otherwise the output is empty.
1162    ///
1163    /// Useful for gating the release of elements based on a condition, such as triggering a protocol if you are missing
1164    /// some local state.
1165    ///
1166    /// # Example
1167    /// ```rust
1168    /// # #[cfg(feature = "deploy")] {
1169    /// # use hydro_lang::prelude::*;
1170    /// # use futures::StreamExt;
1171    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1172    /// let tick = process.tick();
1173    /// // ticks are lazy by default, forces the second tick to run
1174    /// tick.spin_batch(q!(1)).all_ticks().for_each(q!(|_| {}));
1175    ///
1176    /// let batch_first_tick = process
1177    ///   .source_iter(q!(vec![1, 2, 3, 4]))
1178    ///   .batch(&tick, nondet!(/** test */));
1179    /// let batch_second_tick = process
1180    ///   .source_iter(q!(vec![5, 6, 7, 8]))
1181    ///   .batch(&tick, nondet!(/** test */))
1182    ///   .defer_tick(); // appears on the second tick
1183    /// let some_on_first_tick = tick.optional_first_tick(q!(()));
1184    /// batch_first_tick.chain(batch_second_tick)
1185    ///   .filter_if_none(some_on_first_tick)
1186    ///   .all_ticks()
1187    /// # }, |mut stream| async move {
1188    /// // [5, 6, 7, 8]
1189    /// # for w in vec![5, 6, 7, 8] {
1190    /// #     assert_eq!(stream.next().await.unwrap(), w);
1191    /// # }
1192    /// # }));
1193    /// # }
1194    /// ```
1195    #[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    /// Forms the cross-product (Cartesian product, cross-join) of the items in the 2 input streams,
1201    /// returning all tupled pairs.
1202    ///
1203    /// When the right side is [`Bounded`], it is accumulated first and the left side streams
1204    /// through, preserving the left side's ordering. When both sides are [`Unbounded`], a
1205    /// symmetric hash join is used and ordering is [`NoOrder`].
1206    ///
1207    /// # Example
1208    /// ```rust
1209    /// # #[cfg(feature = "deploy")] {
1210    /// # use hydro_lang::prelude::*;
1211    /// # use std::collections::HashSet;
1212    /// # use futures::StreamExt;
1213    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1214    /// let tick = process.tick();
1215    /// let stream1 = process.source_iter(q!(vec![1, 2]));
1216    /// let stream2 = process.source_iter(q!(vec!['a', 'b']));
1217    /// stream1.cross_product(stream2)
1218    /// # }, |mut stream| async move {
1219    /// // (1, 'a'), (1, 'b'), (2, 'a'), (2, 'b') in any order
1220    /// # let expected = HashSet::from([(1, 'a'), (1, 'b'), (2, 'a'), (2, 'b')]);
1221    /// # stream.map(|i| assert!(expected.contains(&i)));
1222    /// # }));
1223    /// # }
1224    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    /// Takes one stream as input and filters out any duplicate occurrences. The output
1239    /// contains all unique values from the input.
1240    ///
1241    /// # Example
1242    /// ```rust
1243    /// # #[cfg(feature = "deploy")] {
1244    /// # use hydro_lang::prelude::*;
1245    /// # use futures::StreamExt;
1246    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1247    /// let tick = process.tick();
1248    /// process.source_iter(q!(vec![1, 2, 3, 2, 1, 4])).unique()
1249    /// # }, |mut stream| async move {
1250    /// # for w in vec![1, 2, 3, 4] {
1251    /// #     assert_eq!(stream.next().await.unwrap(), w);
1252    /// # }
1253    /// # }));
1254    /// # }
1255    /// ```
1256    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    /// Outputs everything in this stream that is *not* contained in the `other` stream.
1272    ///
1273    /// The `other` stream must be [`Bounded`], since this function will wait until
1274    /// all its elements are available before producing any output.
1275    /// # Example
1276    /// ```rust
1277    /// # #[cfg(feature = "deploy")] {
1278    /// # use hydro_lang::prelude::*;
1279    /// # use futures::StreamExt;
1280    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1281    /// let tick = process.tick();
1282    /// let stream = process
1283    ///   .source_iter(q!(vec![ 1, 2, 3, 4 ]))
1284    ///   .batch(&tick, nondet!(/** test */));
1285    /// let batch = process
1286    ///   .source_iter(q!(vec![1, 2]))
1287    ///   .batch(&tick, nondet!(/** test */));
1288    /// stream.filter_not_in(batch).all_ticks()
1289    /// # }, |mut stream| async move {
1290    /// # for w in vec![3, 4] {
1291    /// #     assert_eq!(stream.next().await.unwrap(), w);
1292    /// # }
1293    /// # }));
1294    /// # }
1295    /// ```
1296    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    /// An operator which allows you to "inspect" each element of a stream without
1316    /// modifying it. The closure `f` is called on a reference to each item. This is
1317    /// mainly useful for debugging, and should not be used to generate side-effects.
1318    ///
1319    /// If the input stream is unordered **and** `f` mutably captures state, `f` must be
1320    /// proven **commutative**: executing it on any two elements in either order must
1321    /// leave the mutably-captured state in the same final value. (The elements
1322    /// themselves pass through unchanged.)
1323    ///
1324    /// # Example
1325    /// ```rust
1326    /// # #[cfg(feature = "deploy")] {
1327    /// # use hydro_lang::prelude::*;
1328    /// # use futures::StreamExt;
1329    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1330    /// let nums = process.source_iter(q!(vec![1, 2]));
1331    /// // prints "1 * 10 = 10" and "2 * 10 = 20"
1332    /// nums.inspect(q!(|x| println!("{} * 10 = {}", x, x * 10)))
1333    /// # }, |mut stream| async move {
1334    /// # for w in vec![1, 2] {
1335    /// #     assert_eq!(stream.next().await.unwrap(), w);
1336    /// # }
1337    /// # }));
1338    /// # }
1339    /// ```
1340    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    /// Executes the provided closure for every element in this stream.
1374    ///
1375    /// If the stream is unordered or has retries, the closure must demonstrate commutativity
1376    /// and/or idempotence via annotations:
1377    /// ```rust,ignore
1378    /// stream.for_each(q!(
1379    ///     |x| *flag_mut |= x,
1380    ///     commutative = manual_proof!(/** boolean OR is commutative */),
1381    ///     idempotent = manual_proof!(/** boolean OR is idempotent */)
1382    /// ));
1383    /// ```
1384    ///
1385    /// **Commutative** here means that executing the closure on any two elements in
1386    /// either order must leave its side effects (e.g. mutably-captured state) in the
1387    /// same final value.
1388    ///
1389    /// On a `TotalOrder + ExactlyOnce` stream, no annotations are needed.
1390    ///
1391    /// The closure may capture singletons via `by_ref()` or `by_mut()`, as long as the
1392    /// referenced collection lives at the same location and has the same boundedness as this
1393    /// stream.
1394    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    /// Sends all elements of this stream to a provided [`futures::Sink`], such as an external
1423    /// TCP socket to some other server. You should _not_ use this API for interacting with
1424    /// external clients, instead see [`Location::bidi_external_many_bytes`] and
1425    /// [`Location::bidi_external_many_bincode`]. This should be used for custom, low-level
1426    /// interaction with asynchronous sinks.
1427    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    /// Maps each element `x` of the stream to `(i, x)`, where `i` is the index of the element.
1444    ///
1445    /// # Example
1446    /// ```rust
1447    /// # #[cfg(feature = "deploy")] {
1448    /// # use hydro_lang::{prelude::*, live_collections::stream::{TotalOrder, ExactlyOnce}};
1449    /// # use futures::StreamExt;
1450    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test::<_, _, _, TotalOrder, ExactlyOnce>(|process| {
1451    /// let tick = process.tick();
1452    /// let numbers = process.source_iter(q!(vec![1, 2, 3, 4]));
1453    /// numbers.enumerate()
1454    /// # }, |mut stream| async move {
1455    /// // (0, 1), (1, 2), (2, 3), (3, 4)
1456    /// # for w in vec![(0, 1), (1, 2), (2, 3), (3, 4)] {
1457    /// #     assert_eq!(stream.next().await.unwrap(), w);
1458    /// # }
1459    /// # }));
1460    /// # }
1461    /// ```
1462    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    /// Combines elements of the stream into a [`Singleton`], by starting with an intitial value,
1483    /// generated by the `init` closure, and then applying the `comb` closure to each element in the stream.
1484    /// Unlike iterators, `comb` takes the accumulator by `&mut` reference, so that it can be modified in place.
1485    ///
1486    /// Depending on the input stream guarantees, the closure may need to be commutative
1487    /// (for unordered streams) or idempotent (for streams with non-deterministic duplicates).
1488    ///
1489    /// # Example
1490    /// ```rust
1491    /// # #[cfg(feature = "deploy")] {
1492    /// # use hydro_lang::prelude::*;
1493    /// # use futures::StreamExt;
1494    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1495    /// let words = process.source_iter(q!(vec!["HELLO", "WORLD"]));
1496    /// words
1497    ///     .fold(q!(|| String::new()), q!(|acc, x| acc.push_str(x)))
1498    ///     .into_stream()
1499    /// # }, |mut stream| async move {
1500    /// // "HELLOWORLD"
1501    /// # assert_eq!(stream.next().await.unwrap(), "HELLOWORLD");
1502    /// # }));
1503    /// # }
1504    /// ```
1505    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        // Only assume_retries (for idempotence), not assume_ordering.
1530        // The fold hook in the simulator handles ordering non-determinism directly, so an
1531        // ordering hook on the commutativity proof binds to the fold operator itself.
1532        let nondet = nondet!(/** the combinator function is commutative and idempotent */);
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            // we do not guarantee consistency at this point because if the algebraic properties
1546            // do not hold in practice, replica consistency may fail to be maintained, so we
1547            // would like the simulator to assert consistency; in the future, this will be dynamic
1548            // based on the proof mechanism
1549        };
1550
1551        Singleton::new(retried.location.clone(), core)
1552            .assert_has_consistency_of(manual_proof!(/** algebraic properties */))
1553    }
1554
1555    /// Combines elements of the stream into an [`Optional`], by starting with the first element in the stream,
1556    /// and then applying the `comb` closure to each element in the stream. The [`Optional`] will be empty
1557    /// until the first element in the input arrives. Unlike iterators, `comb` takes the accumulator by `&mut`
1558    /// reference, so that it can be modified in place.
1559    ///
1560    /// Depending on the input stream guarantees, the closure may need to be commutative
1561    /// (for unordered streams) or idempotent (for streams with non-deterministic duplicates).
1562    ///
1563    /// # Example
1564    /// ```rust
1565    /// # #[cfg(feature = "deploy")] {
1566    /// # use hydro_lang::prelude::*;
1567    /// # use futures::StreamExt;
1568    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1569    /// let bools = process.source_iter(q!(vec![false, true, false]));
1570    /// bools.reduce(q!(|acc, x| *acc |= x)).into_stream()
1571    /// # }, |mut stream| async move {
1572    /// // true
1573    /// # assert_eq!(stream.next().await.unwrap(), true);
1574    /// # }));
1575    /// # }
1576    /// ```
1577    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!(/** the combinator function is commutative and idempotent */);
1596        let ordered_etc: Stream<T, L::DropConsistency, B> =
1597            self.assume_retries(nondet_retries).assume_ordering(nondet!(
1598                /// the combinator function is commutative; the simulator still explores
1599                /// (or scripts, via the proof's hook) the ordering
1600                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!(/** algebraic properties */))
1615    }
1616
1617    /// Computes the maximum element in the stream as an [`Optional`], which
1618    /// will be empty until the first element in the input arrives.
1619    ///
1620    /// # Example
1621    /// ```rust
1622    /// # #[cfg(feature = "deploy")] {
1623    /// # use hydro_lang::prelude::*;
1624    /// # use futures::StreamExt;
1625    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1626    /// let tick = process.tick();
1627    /// let numbers = process.source_iter(q!(vec![1, 2, 3, 4]));
1628    /// let batch = numbers.batch(&tick, nondet!(/** test */));
1629    /// batch.max().all_ticks()
1630    /// # }, |mut stream| async move {
1631    /// // 4
1632    /// # assert_eq!(stream.next().await.unwrap(), 4);
1633    /// # }));
1634    /// # }
1635    /// ```
1636    pub fn max(self) -> Optional<T, L, B::AggregatedOptional>
1637    where
1638        T: Ord,
1639    {
1640        self.assume_retries_trusted::<ExactlyOnce>(nondet!(/** max is idempotent */))
1641            .assume_ordering_trusted_bounded::<TotalOrder>(
1642                nondet!(/** max is commutative, but order affects intermediates */),
1643            )
1644            .reduce(q!(|curr, new| {
1645                if new > *curr {
1646                    *curr = new;
1647                }
1648            }))
1649    }
1650
1651    /// Computes the minimum element in the stream as an [`Optional`], which
1652    /// will be empty until the first element in the input arrives.
1653    ///
1654    /// # Example
1655    /// ```rust
1656    /// # #[cfg(feature = "deploy")] {
1657    /// # use hydro_lang::prelude::*;
1658    /// # use futures::StreamExt;
1659    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1660    /// let tick = process.tick();
1661    /// let numbers = process.source_iter(q!(vec![1, 2, 3, 4]));
1662    /// let batch = numbers.batch(&tick, nondet!(/** test */));
1663    /// batch.min().all_ticks()
1664    /// # }, |mut stream| async move {
1665    /// // 1
1666    /// # assert_eq!(stream.next().await.unwrap(), 1);
1667    /// # }));
1668    /// # }
1669    /// ```
1670    pub fn min(self) -> Optional<T, L, B::AggregatedOptional>
1671    where
1672        T: Ord,
1673    {
1674        self.assume_retries_trusted::<ExactlyOnce>(nondet!(/** min is idempotent */))
1675            .assume_ordering_trusted_bounded::<TotalOrder>(
1676                nondet!(/** max is commutative, but order affects intermediates */),
1677            )
1678            .reduce(q!(|curr, new| {
1679                if new < *curr {
1680                    *curr = new;
1681                }
1682            }))
1683    }
1684
1685    /// Computes the first element in the stream as an [`Optional`], which
1686    /// will be empty until the first element in the input arrives.
1687    ///
1688    /// This requires the stream to have a [`TotalOrder`] guarantee, otherwise
1689    /// re-ordering of elements may cause the first element to change.
1690    ///
1691    /// # Example
1692    /// ```rust
1693    /// # #[cfg(feature = "deploy")] {
1694    /// # use hydro_lang::prelude::*;
1695    /// # use futures::StreamExt;
1696    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1697    /// let tick = process.tick();
1698    /// let numbers = process.source_iter(q!(vec![1, 2, 3, 4]));
1699    /// let batch = numbers.batch(&tick, nondet!(/** test */));
1700    /// batch.first().all_ticks()
1701    /// # }, |mut stream| async move {
1702    /// // 1
1703    /// # assert_eq!(stream.next().await.unwrap(), 1);
1704    /// # }));
1705    /// # }
1706    /// ```
1707    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!(/** first is idempotent */))
1713            .generator(q!(|| ()), q!(|_, item| Generate::Return(item)))
1714            .reduce(q!(|_, _| {}))
1715    }
1716
1717    /// Computes the last element in the stream as an [`Optional`], which
1718    /// will be empty until an element in the input arrives.
1719    ///
1720    /// This requires the stream to have a [`TotalOrder`] guarantee, otherwise
1721    /// re-ordering of elements may cause the last element to change.
1722    ///
1723    /// # Example
1724    /// ```rust
1725    /// # #[cfg(feature = "deploy")] {
1726    /// # use hydro_lang::prelude::*;
1727    /// # use futures::StreamExt;
1728    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1729    /// let tick = process.tick();
1730    /// let numbers = process.source_iter(q!(vec![1, 2, 3, 4]));
1731    /// let batch = numbers.batch(&tick, nondet!(/** test */));
1732    /// batch.last().all_ticks()
1733    /// # }, |mut stream| async move {
1734    /// // 4
1735    /// # assert_eq!(stream.next().await.unwrap(), 4);
1736    /// # }));
1737    /// # }
1738    /// ```
1739    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!(/** last is idempotent */))
1745            .reduce(q!(|curr, new| *curr = new))
1746    }
1747
1748    /// Returns a stream containing at most the first `n` elements of the input stream,
1749    /// preserving the original order. Similar to `LIMIT` in SQL.
1750    ///
1751    /// This requires the stream to have a [`TotalOrder`] guarantee and [`ExactlyOnce`]
1752    /// retries, since the result depends on the order and cardinality of elements.
1753    ///
1754    /// # Example
1755    /// ```rust
1756    /// # #[cfg(feature = "deploy")] {
1757    /// # use hydro_lang::prelude::*;
1758    /// # use futures::StreamExt;
1759    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1760    /// let numbers = process.source_iter(q!(vec![10, 20, 30, 40, 50]));
1761    /// numbers.limit(q!(3))
1762    /// # }, |mut stream| async move {
1763    /// // 10, 20, 30
1764    /// # for w in vec![10, 20, 30] {
1765    /// #     assert_eq!(stream.next().await.unwrap(), w);
1766    /// # }
1767    /// # }));
1768    /// # }
1769    /// ```
1770    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    /// Collects all the elements of this stream into a single [`Vec`] element.
1796    ///
1797    /// If the input stream is [`Unbounded`], the output [`Singleton`] will be [`Unbounded`] as
1798    /// well, which means that the value of the [`Vec`] will asynchronously grow as new elements
1799    /// are added. On such a value, you can use [`Singleton::snapshot`] to grab an instance of
1800    /// the vector at an arbitrary point in time.
1801    ///
1802    /// # Example
1803    /// ```rust
1804    /// # #[cfg(feature = "deploy")] {
1805    /// # use hydro_lang::prelude::*;
1806    /// # use futures::StreamExt;
1807    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1808    /// let tick = process.tick();
1809    /// let numbers = process.source_iter(q!(vec![1, 2, 3, 4]));
1810    /// let batch = numbers.batch(&tick, nondet!(/** test */));
1811    /// batch.collect_vec().all_ticks() // emit each tick's Vec into an unbounded stream
1812    /// # }, |mut stream| async move {
1813    /// // [ vec![1, 2, 3, 4] ]
1814    /// # for w in vec![vec![1, 2, 3, 4]] {
1815    /// #     assert_eq!(stream.next().await.unwrap(), w);
1816    /// # }
1817    /// # }));
1818    /// # }
1819    /// ```
1820    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    /// Applies a function to each element of the stream, maintaining an internal state (accumulator)
1834    /// and emitting each intermediate result.
1835    ///
1836    /// Unlike `fold` which only returns the final accumulated value, `scan` produces a new stream
1837    /// containing all intermediate accumulated values. The scan operation can also terminate early
1838    /// by returning `None`.
1839    ///
1840    /// The function takes a mutable reference to the accumulator and the current element, and returns
1841    /// an `Option<U>`. If the function returns `Some(value)`, `value` is emitted to the output stream.
1842    /// If the function returns `None`, the stream is terminated and no more elements are processed.
1843    ///
1844    /// The `init` and `f` closures may capture bounded singletons, optionals, or streams by
1845    /// reference via [`by_ref()`](crate::live_collections::Singleton::by_ref), as long as the
1846    /// referenced collection lives at the same location and has the same boundedness as this
1847    /// stream.
1848    ///
1849    /// # Examples
1850    ///
1851    /// Basic usage - running sum:
1852    /// ```rust
1853    /// # #[cfg(feature = "deploy")] {
1854    /// # use hydro_lang::prelude::*;
1855    /// # use futures::StreamExt;
1856    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1857    /// process.source_iter(q!(vec![1, 2, 3, 4])).scan(
1858    ///     q!(|| 0),
1859    ///     q!(|acc, x| {
1860    ///         *acc += x;
1861    ///         Some(*acc)
1862    ///     }),
1863    /// )
1864    /// # }, |mut stream| async move {
1865    /// // Output: 1, 3, 6, 10
1866    /// # for w in vec![1, 3, 6, 10] {
1867    /// #     assert_eq!(stream.next().await.unwrap(), w);
1868    /// # }
1869    /// # }));
1870    /// # }
1871    /// ```
1872    ///
1873    /// Early termination example:
1874    /// ```rust
1875    /// # #[cfg(feature = "deploy")] {
1876    /// # use hydro_lang::prelude::*;
1877    /// # use futures::StreamExt;
1878    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1879    /// process.source_iter(q!(vec![1, 2, 3, 4])).scan(
1880    ///     q!(|| 1),
1881    ///     q!(|state, x| {
1882    ///         *state = *state * x;
1883    ///         if *state > 6 {
1884    ///             None // Terminate the stream
1885    ///         } else {
1886    ///             Some(-*state)
1887    ///         }
1888    ///     }),
1889    /// )
1890    /// # }, |mut stream| async move {
1891    /// // Output: -1, -2, -6
1892    /// # for w in vec![-1, -2, -6] {
1893    /// #     assert_eq!(stream.next().await.unwrap(), w);
1894    /// # }
1895    /// # }));
1896    /// # }
1897    /// ```
1898    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    /// Async version of [`Stream::scan`]. Applies an async function to each element of the
1932    /// stream, maintaining an internal state (accumulator) and emitting the values returned
1933    /// by the function.
1934    ///
1935    /// The closure runs synchronously (so it can mutate the accumulator), then returns a
1936    /// future. The future is polled to completion. If it resolves to `Some`, the value is
1937    /// emitted. If it resolves to `None`, the item is filtered out.
1938    ///
1939    /// The `init` and `f` closures may capture bounded singletons, optionals, or streams by
1940    /// reference via [`by_ref()`](crate::live_collections::Singleton::by_ref), as long as the
1941    /// referenced collection lives at the same location and has the same boundedness as this
1942    /// stream.
1943    ///
1944    /// # Examples
1945    ///
1946    /// ```rust
1947    /// # #[cfg(feature = "deploy")] {
1948    /// # use hydro_lang::prelude::*;
1949    /// # use futures::StreamExt;
1950    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
1951    /// process
1952    ///     .source_iter(q!(vec![1, 2, 3, 4]))
1953    ///     .scan_async_blocking(
1954    ///         q!(|| 0),
1955    ///         q!(|acc, x| {
1956    ///             *acc += x;
1957    ///             let val = *acc;
1958    ///             async move { Some(val) }
1959    ///         }),
1960    ///     )
1961    /// # }, |mut stream| async move {
1962    /// // Output: 1, 3, 6, 10
1963    /// # for w in vec![1, 3, 6, 10] {
1964    /// #     assert_eq!(stream.next().await.unwrap(), w);
1965    /// # }
1966    /// # }));
1967    /// # }
1968    /// ```
1969    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    /// Iteratively processes the elements of the stream using a state machine that can yield
2004    /// elements as it processes its inputs. This is designed to mirror the unstable generator
2005    /// syntax in Rust, without requiring special syntax.
2006    ///
2007    /// Like [`Stream::scan`], this function takes in an initializer that emits the initial
2008    /// state. The second argument defines the processing logic, taking in a mutable reference
2009    /// to the state and the value to be processed. It emits a [`Generate`] value, whose
2010    /// variants define what is emitted and whether further inputs should be processed.
2011    ///
2012    /// The `init` and `f` closures may capture bounded singletons, optionals, or streams by
2013    /// reference via [`by_ref()`](crate::live_collections::Singleton::by_ref), as long as the
2014    /// referenced collection lives at the same location and has the same boundedness as this
2015    /// stream.
2016    ///
2017    /// # Example
2018    /// ```rust
2019    /// # #[cfg(feature = "deploy")] {
2020    /// # use hydro_lang::prelude::*;
2021    /// # use futures::StreamExt;
2022    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
2023    /// process.source_iter(q!(vec![1, 3, 100, 10])).generator(
2024    ///     q!(|| 0),
2025    ///     q!(|acc, x| {
2026    ///         *acc += x;
2027    ///         if *acc > 100 {
2028    ///             hydro_lang::live_collections::keyed_stream::Generate::Return("done!".to_owned())
2029    ///         } else if *acc % 2 == 0 {
2030    ///             hydro_lang::live_collections::keyed_stream::Generate::Yield("even".to_owned())
2031    ///         } else {
2032    ///             hydro_lang::live_collections::keyed_stream::Generate::Continue
2033    ///         }
2034    ///     }),
2035    /// )
2036    /// # }, |mut stream| async move {
2037    /// // Output: "even", "done!"
2038    /// # let mut results = Vec::new();
2039    /// # for _ in 0..2 {
2040    /// #     results.push(stream.next().await.unwrap());
2041    /// # }
2042    /// # results.sort();
2043    /// # assert_eq!(results, vec!["done!".to_owned(), "even".to_owned()]);
2044    /// # }));
2045    /// # }
2046    /// ```
2047    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        // State is Option<Option<A>>:
2066        //   None = not yet initialized
2067        //   Some(Some(a)) = active with state a
2068        //   Some(None) = terminated
2069        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                        // Unlike KeyedStream, we can terminate the scan directly on
2087                        // Break/Return because there is only one state (no other keys
2088                        // that still need processing).
2089                        Generate::Break => None,
2090                        Generate::Continue => Some(None),
2091                    },
2092                    // State is Some(None) after Return; terminate the scan.
2093                    _ => 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    /// Given a time interval, returns a stream corresponding to samples taken from the
2130    /// stream roughly at that interval. The output will have elements in the same order
2131    /// as the input, but with arbitrary elements skipped between samples. There is also
2132    /// no guarantee on the exact timing of the samples.
2133    ///
2134    /// # Non-Determinism
2135    /// The output stream is non-deterministic in which elements are sampled, since this
2136    /// is controlled by a clock.
2137    ///
2138    /// In simulation tests, the internal batching of elements and of clock samples can be
2139    /// scripted through the guard's composite hook payload, e.g.
2140    /// `nondet!(/** reason */ hook = (elements_hook.into(), None))`.
2141    #[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                /// which elements are batched between samples is captured by the caller's guard
2165                hook = elements_hook
2166            ),
2167        )
2168        .filter_if(
2169            samples
2170                .batch(
2171                    &tick,
2172                    nondet!(
2173                        /// sample timing is captured by the caller's guard
2174                        hook = samples_hook
2175                    ),
2176                )
2177                .first()
2178                .is_some(),
2179        )
2180        .all_ticks()
2181        .weaken_retries()
2182    }
2183
2184    /// Given a timeout duration, returns an [`Optional`]  which will have a value if the
2185    /// stream has not emitted a value since that duration.
2186    ///
2187    /// # Non-Determinism
2188    /// Timeout relies on non-deterministic sampling of the stream, so depending on when
2189    /// samples take place, timeouts may be non-deterministically generated or missed,
2190    /// and the notification of the timeout may be delayed as well. There is also no
2191    /// guarantee on how long the [`Optional`] will have a value after the timeout is
2192    /// detected based on when the next sample is taken.
2193    #[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!(/** TODO */)
2216            ),
2217        );
2218
2219        latest_received
2220            .snapshot(
2221                &tick,
2222                nondet!(
2223                    /// sampling timing is captured by the caller's guard
2224                    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    /// Shifts this stream into an atomic context, which guarantees that any downstream logic
2242    /// will all be executed synchronously before any outputs are yielded (in [`Stream::end_atomic`]).
2243    ///
2244    /// This is useful to enforce local consistency constraints, such as ensuring that a write is
2245    /// processed before an acknowledgement is emitted.
2246    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    /// Given a tick, returns a stream corresponding to a batch of elements segmented by
2264    /// that tick. These batches are guaranteed to be contiguous across ticks and preserve
2265    /// the order of the input. The output stream will execute in the [`Tick`] that was
2266    /// used to create the atomic section.
2267    ///
2268    /// # Non-Determinism
2269    /// The batch boundaries are non-deterministic and may change across executions.
2270    ///
2271    /// In simulation tests, the batching decisions can be scripted by attaching a
2272    /// [`BatchHook`](crate::sim_hooks::BatchHook) to the guard via
2273    /// `nondet!(/** reason */ hook = my_hook)`.
2274    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    /// An operator which allows you to "name" a `HydroNode`.
2297    /// This is only used for testing, to correlate certain `HydroNode`s with IDs.
2298    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    /// Turns this [`Stream`] into a [`Optional`], under the invariant assumption that there is at
2308    /// most one element. If this invariant is broken, the program may exhibit undefined behavior,
2309    /// so uses must be carefully vetted.
2310    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    /// Explicitly "casts" the stream to a type with a different ordering
2341    /// guarantee. Useful in unsafe code where the ordering cannot be proven
2342    /// by the type-system.
2343    ///
2344    /// # Non-Determinism
2345    /// This function is used as an escape hatch, and any mistakes in the
2346    /// provided ordering guarantee will propagate into the guarantees
2347    /// for the rest of the program.
2348    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            // We can always weaken the ordering guarantee
2356            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    // like `assume_ordering_trusted`, but only if the input stream is bounded and therefore
2382    // intermediate states will not be revealed
2383    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                /// the unbounded stream exposes ordering non-determinism in intermediate states
2393                nondet
2394            ));
2395            Stream::new(self_location, inner.ir_node.replace(HydroNode::Placeholder))
2396        }
2397    }
2398
2399    // only for internal APIs that have been carefully vetted to ensure that the non-determinism
2400    // is not observable
2401    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            // We can always weaken the ordering guarantee
2409            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    /// Weakens the ordering guarantee provided by the stream to [`NoOrder`],
2434    /// which is always safe because that is the weakest possible guarantee.
2435    pub fn weakest_ordering(self) -> Stream<T, L, B, NoOrder, R> {
2436        self.weaken_ordering::<NoOrder>()
2437    }
2438
2439    /// Weakens the ordering guarantee provided by the stream to `O2`, with the type-system
2440    /// enforcing that `O2` is weaker than the input ordering guarantee.
2441    pub fn weaken_ordering<O2: WeakerOrderingThan<O>>(self) -> Stream<T, L, B, O2, R> {
2442        let nondet = nondet!(/** this is a weaker ordering guarantee, so it is safe to assume */);
2443        self.assume_ordering_trusted::<O2>(nondet)
2444    }
2445
2446    /// Strengthens the ordering guarantee to `TotalOrder`, given that `O: IsOrdered`, which
2447    /// implies that `O == TotalOrder`.
2448    pub fn make_totally_ordered(self) -> Stream<T, L, B, TotalOrder, R>
2449    where
2450        O: IsOrdered,
2451    {
2452        self.assume_ordering_trusted(nondet!(/** no-op */))
2453    }
2454
2455    /// Explicitly "casts" the stream to a type with a different retries
2456    /// guarantee. Useful in unsafe code where the lack of retries cannot
2457    /// be proven by the type-system.
2458    ///
2459    /// # Non-Determinism
2460    /// This function is used as an escape hatch, and any mistakes in the
2461    /// provided retries guarantee will propagate into the guarantees
2462    /// for the rest of the program.
2463    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            // We can always weaken the retries guarantee
2474            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    // only for internal APIs that have been carefully vetted to ensure that the non-determinism
2498    // is not observable
2499    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            // We can always weaken the retries guarantee
2507            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    /// Weakens the retries guarantee provided by the stream to [`AtLeastOnce`],
2532    /// which is always safe because that is the weakest possible guarantee.
2533    pub fn weakest_retries(self) -> Stream<T, L, B, O, AtLeastOnce> {
2534        self.weaken_retries::<AtLeastOnce>()
2535    }
2536
2537    /// Weakens the retries guarantee provided by the stream to `R2`, with the type-system
2538    /// enforcing that `R2` is weaker than the input retries guarantee.
2539    pub fn weaken_retries<R2: WeakerRetryThan<R>>(self) -> Stream<T, L, B, O, R2> {
2540        let nondet = nondet!(/** this is a weaker retry guarantee, so it is safe to assume */);
2541        self.assume_retries_trusted::<R2>(nondet)
2542    }
2543
2544    /// Strengthens the retry guarantee to `ExactlyOnce`, given that `R: IsExactlyOnce`, which
2545    /// implies that `R == ExactlyOnce`.
2546    pub fn make_exactly_once(self) -> Stream<T, L, B, O, ExactlyOnce>
2547    where
2548        R: IsExactlyOnce,
2549    {
2550        self.assume_retries_trusted(nondet!(/** no-op */))
2551    }
2552
2553    /// Strengthens the boundedness guarantee to `Bounded`, given that `B: IsBounded`, which
2554    /// implies that `B == Bounded`.
2555    pub fn make_bounded(self) -> Stream<T, L, Bounded, O, R>
2556    where
2557        B: IsBounded,
2558    {
2559        self.weaken_boundedness()
2560    }
2561
2562    /// Weakens the boundedness guarantee to an arbitrary boundedness `B2`, given that `B: IsBounded`,
2563    /// which implies that `B == Bounded`.
2564    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            // We can always weaken the boundedness
2572            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    /// Clone each element of the stream; akin to `map(q!(|d| d.clone()))`.
2590    ///
2591    /// # Example
2592    /// ```rust
2593    /// # #[cfg(feature = "deploy")] {
2594    /// # use hydro_lang::prelude::*;
2595    /// # use futures::StreamExt;
2596    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
2597    /// process.source_iter(q!(&[1, 2, 3])).cloned()
2598    /// # }, |mut stream| async move {
2599    /// // 1, 2, 3
2600    /// # for w in vec![1, 2, 3] {
2601    /// #     assert_eq!(stream.next().await.unwrap(), w);
2602    /// # }
2603    /// # }));
2604    /// # }
2605    /// ```
2606    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    /// Computes the number of elements in the stream as a [`Singleton`].
2619    ///
2620    /// # Example
2621    /// ```rust
2622    /// # #[cfg(feature = "deploy")] {
2623    /// # use hydro_lang::prelude::*;
2624    /// # use futures::StreamExt;
2625    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
2626    /// let tick = process.tick();
2627    /// let numbers = process.source_iter(q!(vec![1, 2, 3, 4]));
2628    /// let batch = numbers.batch(&tick, nondet!(/** test */));
2629    /// batch.count().all_ticks()
2630    /// # }, |mut stream| async move {
2631    /// // 4
2632    /// # assert_eq!(stream.next().await.unwrap(), 4);
2633    /// # }));
2634    /// # }
2635    /// ```
2636    pub fn count(self) -> Singleton<usize, L, B::StreamToMonotone> {
2637        self.assume_ordering_trusted::<TotalOrder>(nondet!(
2638            /// Order does not affect eventual count, and also does not affect intermediate states.
2639        ))
2640        .fold(
2641            q!(|| 0usize),
2642            q!(
2643                |count, _| *count += 1,
2644                monotone = manual_proof!(/** += 1 is monotone */)
2645            ),
2646        )
2647    }
2648}
2649
2650impl<'a, T, L: Location<'a>, O: Ordering, R: Retries> Stream<T, L, Unbounded, O, R> {
2651    /// Produces a new stream that merges the elements of the two input streams.
2652    /// The result has [`NoOrder`] because the order of merging is not guaranteed.
2653    ///
2654    /// Currently, both input streams must be [`Unbounded`]. When the streams are
2655    /// [`Bounded`], you can use [`Stream::chain`] instead.
2656    ///
2657    /// # Example
2658    /// ```rust
2659    /// # #[cfg(feature = "deploy")] {
2660    /// # use hydro_lang::prelude::*;
2661    /// # use futures::StreamExt;
2662    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
2663    /// let numbers: Stream<i32, _, Unbounded> = // 1, 2, 3, 4
2664    /// # process.source_iter(q!(vec![1, 2, 3, 4])).into();
2665    /// numbers.clone().map(q!(|x| x + 1)).merge_unordered(numbers)
2666    /// # }, |mut stream| async move {
2667    /// // 2, 3, 4, 5, and 1, 2, 3, 4 merged in unknown order
2668    /// # for w in vec![2, 3, 4, 5, 1, 2, 3, 4] {
2669    /// #     assert_eq!(stream.next().await.unwrap(), w);
2670    /// # }
2671    /// # }));
2672    /// # }
2673    /// ```
2674    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: use [`Stream::merge_unordered`] instead.
2698    #[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    /// Produces a new stream that combines the elements of the two input streams,
2712    /// preserving the relative order of elements within each input.
2713    ///
2714    /// # Non-Determinism
2715    /// The order in which elements *across* the two streams will be interleaved is
2716    /// non-deterministic, so the order of elements will vary across runs. If the output
2717    /// order is irrelevant, use [`Stream::merge_unordered`] instead, which is deterministic
2718    /// but emits an unordered stream. For deterministic first-then-second ordering on
2719    /// bounded streams, use [`Stream::chain`].
2720    ///
2721    /// In simulation tests, the interleaving decisions can be scripted by attaching a
2722    /// [`MergeOrderedHook`](crate::sim_hooks::MergeOrderedHook) to the guard via
2723    /// `nondet!(/** reason */ hook = my_hook)`.
2724    ///
2725    /// # Example
2726    /// ```rust
2727    /// # #[cfg(feature = "deploy")] {
2728    /// # use hydro_lang::prelude::*;
2729    /// # use futures::StreamExt;
2730    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
2731    /// let numbers: Stream<i32, _, Unbounded> = // 1, 3
2732    /// # process.source_iter(q!(vec![1, 3])).into();
2733    /// numbers.clone().merge_ordered(numbers.map(q!(|x| x + 1)), nondet!(/** example */))
2734    /// # }, |mut stream| async move {
2735    /// // 1, 3 and 2, 4 in some order, preserving the original local order
2736    /// # for w in vec![1, 3, 2, 4] {
2737    /// #     assert_eq!(stream.next().await.unwrap(), w);
2738    /// # }
2739    /// # }));
2740    /// # }
2741    /// ```
2742    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    /// Produces a new stream that emits the input elements in sorted order.
2775    ///
2776    /// The input stream can have any ordering guarantee, but the output stream
2777    /// will have a [`TotalOrder`] guarantee. This operator will block until all
2778    /// elements in the input stream are available, so it requires the input stream
2779    /// to be [`Bounded`].
2780    ///
2781    /// # Example
2782    /// ```rust
2783    /// # #[cfg(feature = "deploy")] {
2784    /// # use hydro_lang::prelude::*;
2785    /// # use futures::StreamExt;
2786    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
2787    /// let tick = process.tick();
2788    /// let numbers = process.source_iter(q!(vec![4, 2, 3, 1]));
2789    /// let batch = numbers.batch(&tick, nondet!(/** test */));
2790    /// batch.sort().all_ticks()
2791    /// # }, |mut stream| async move {
2792    /// // 1, 2, 3, 4
2793    /// # for w in (1..5) {
2794    /// #     assert_eq!(stream.next().await.unwrap(), w);
2795    /// # }
2796    /// # }));
2797    /// # }
2798    /// ```
2799    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    /// Produces a new stream that first emits the elements of the `self` stream,
2817    /// and then emits the elements of the `other` stream. The output stream has
2818    /// a [`TotalOrder`] guarantee if and only if both input streams have a
2819    /// [`TotalOrder`] guarantee.
2820    ///
2821    /// Currently, both input streams must be [`Bounded`]. This operator will block
2822    /// on the first stream until all its elements are available. In a future version,
2823    /// we will relax the requirement on the `other` stream.
2824    ///
2825    /// # Example
2826    /// ```rust
2827    /// # #[cfg(feature = "deploy")] {
2828    /// # use hydro_lang::prelude::*;
2829    /// # use futures::StreamExt;
2830    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
2831    /// let tick = process.tick();
2832    /// let numbers = process.source_iter(q!(vec![1, 2, 3, 4]));
2833    /// let batch = numbers.batch(&tick, nondet!(/** test */));
2834    /// batch.clone().map(q!(|x| x + 1)).chain(batch).all_ticks()
2835    /// # }, |mut stream| async move {
2836    /// // 2, 3, 4, 5, 1, 2, 3, 4
2837    /// # for w in vec![2, 3, 4, 5, 1, 2, 3, 4] {
2838    /// #     assert_eq!(stream.next().await.unwrap(), w);
2839    /// # }
2840    /// # }));
2841    /// # }
2842    /// ```
2843    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    /// Forms the cross-product (Cartesian product, cross-join) of the items in the 2 input streams.
2871    /// Unlike [`Stream::cross_product`], the output order is totally ordered when the inputs are
2872    /// because this is compiled into a nested loop.
2873    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    /// Creates a [`KeyedStream`] with the same set of keys as `keys`, but with the elements in
2903    /// `self` used as the values for *each* key.
2904    ///
2905    /// This is helpful when "broadcasting" a set of values so that all the keys have the same
2906    /// values. For example, it can be used to send the same set of elements to several cluster
2907    /// members, if the membership information is available as a [`KeyedSingleton`].
2908    ///
2909    /// # Example
2910    /// ```rust
2911    /// # #[cfg(feature = "deploy")] {
2912    /// # use hydro_lang::prelude::*;
2913    /// # use futures::StreamExt;
2914    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
2915    /// # let tick = process.tick();
2916    /// let keyed_singleton = // { 1: (), 2: () }
2917    /// # process
2918    /// #     .source_iter(q!(vec![(1, ()), (2, ())]))
2919    /// #     .into_keyed()
2920    /// #     .batch(&tick, nondet!(/** test */))
2921    /// #     .first();
2922    /// let stream = // [ "a", "b" ]
2923    /// # process
2924    /// #     .source_iter(q!(vec!["a".to_owned(), "b".to_owned()]))
2925    /// #     .batch(&tick, nondet!(/** test */));
2926    /// stream.repeat_with_keys(keyed_singleton)
2927    /// # .entries().all_ticks()
2928    /// # }, |mut stream| async move {
2929    /// // { 1: ["a", "b" ], 2: ["a", "b"] }
2930    /// # let mut results = Vec::new();
2931    /// # for _ in 0..4 {
2932    /// #     results.push(stream.next().await.unwrap());
2933    /// # }
2934    /// # results.sort();
2935    /// # assert_eq!(results, vec![(1, "a".to_owned()), (1, "b".to_owned()), (2, "a".to_owned()), (2, "b".to_owned())]);
2936    /// # }));
2937    /// # }
2938    /// ```
2939    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!(/** keyed stream does not depend on ordering of keys */),
2951            )
2952            .cross_product_nested_loop(self.make_bounded())
2953            .into_keyed()
2954    }
2955
2956    /// Consumes a stream of `Future<T>`, resolving each future while blocking subgraph
2957    /// execution until all results are available. The output order is based on when futures
2958    /// complete, and may be different than the input order.
2959    ///
2960    /// Unlike [`Stream::resolve_futures`], which allows the subgraph to continue executing
2961    /// while futures are pending, this variant blocks until the futures resolve.
2962    ///
2963    /// # Example
2964    /// ```rust
2965    /// # #[cfg(feature = "deploy")] {
2966    /// # use std::collections::HashSet;
2967    /// # use futures::StreamExt;
2968    /// # use hydro_lang::prelude::*;
2969    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
2970    /// process
2971    ///     .source_iter(q!([2, 3, 1, 9, 6, 5, 4, 7, 8]))
2972    ///     .map(q!(|x| async move {
2973    ///         tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
2974    ///         x
2975    ///     }))
2976    ///     .resolve_futures_blocking()
2977    /// #   },
2978    /// #   |mut stream| async move {
2979    /// // 1, 2, 3, 4, 5, 6, 7, 8, 9 (in any order)
2980    /// #       let mut output = HashSet::new();
2981    /// #       for _ in 1..10 {
2982    /// #           output.insert(stream.next().await.unwrap());
2983    /// #       }
2984    /// #       assert_eq!(
2985    /// #           output,
2986    /// #           HashSet::<i32>::from_iter(1..10)
2987    /// #       );
2988    /// #   },
2989    /// # ));
2990    /// # }
2991    /// ```
2992    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    /// Returns a [`Singleton`] containing `true` if the stream has no elements, or `false` otherwise.
3008    ///
3009    /// # Example
3010    /// ```rust
3011    /// # #[cfg(feature = "deploy")] {
3012    /// # use hydro_lang::prelude::*;
3013    /// # use futures::StreamExt;
3014    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
3015    /// let tick = process.tick();
3016    /// let empty: Stream<i32, _, Bounded> = process
3017    ///   .source_iter(q!(Vec::<i32>::new()))
3018    ///   .batch(&tick, nondet!(/** test */));
3019    /// empty.is_empty().all_ticks()
3020    /// # }, |mut stream| async move {
3021    /// // true
3022    /// # assert_eq!(stream.next().await.unwrap(), true);
3023    /// # }));
3024    /// # }
3025    /// ```
3026    #[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!(/** is_empty intermediates unaffected by order */),
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    /// Given two streams of pairs `(K, V1)` and `(K, V2)`, produces a new stream of nested pairs `(K, (V1, V2))`
3045    /// by equi-joining the two streams on the key attribute `K`.
3046    ///
3047    /// When the right-hand side is [`Bounded`], the join accumulates the right side first
3048    /// and streams the left side through, preserving the left side's ordering. When both
3049    /// sides are [`Unbounded`], a symmetric hash join is used and ordering is [`NoOrder`].
3050    ///
3051    /// # Example
3052    /// ```rust
3053    /// # #[cfg(feature = "deploy")] {
3054    /// # use hydro_lang::prelude::*;
3055    /// # use std::collections::HashSet;
3056    /// # use futures::StreamExt;
3057    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
3058    /// let tick = process.tick();
3059    /// let stream1 = process.source_iter(q!(vec![(1, 'a'), (2, 'b')]));
3060    /// let stream2 = process.source_iter(q!(vec![(1, 'x'), (2, 'y')]));
3061    /// stream1.join(stream2)
3062    /// # }, |mut stream| async move {
3063    /// // (1, ('a', 'x')), (2, ('b', 'y'))
3064    /// # let expected = HashSet::from([(1, ('a', 'x')), (2, ('b', 'y'))]);
3065    /// # stream.map(|i| assert!(expected.contains(&i)));
3066    /// # }));
3067    /// # }
3068    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    /// Given a stream of pairs `(K, V1)` and a bounded stream of keys `K`,
3110    /// computes the anti-join of the items in the input -- i.e. returns
3111    /// unique items in the first input that do not have a matching key
3112    /// in the second input.
3113    ///
3114    /// # Example
3115    /// ```rust
3116    /// # #[cfg(feature = "deploy")] {
3117    /// # use hydro_lang::prelude::*;
3118    /// # use futures::StreamExt;
3119    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
3120    /// let tick = process.tick();
3121    /// let stream = process
3122    ///   .source_iter(q!(vec![ (1, 'a'), (2, 'b'), (3, 'c'), (4, 'd') ]))
3123    ///   .batch(&tick, nondet!(/** test */));
3124    /// let batch = process
3125    ///   .source_iter(q!(vec![1, 2]))
3126    ///   .batch(&tick, nondet!(/** test */));
3127    /// stream.anti_join(batch).all_ticks()
3128    /// # }, |mut stream| async move {
3129    /// # for w in vec![(3, 'c'), (4, 'd')] {
3130    /// #     assert_eq!(stream.next().await.unwrap(), w);
3131    /// # }
3132    /// # }));
3133    /// # }
3134    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    /// Transforms this stream into a [`KeyedStream`], where the first element of each tuple
3160    /// is used as the key and the second element is added to the entries associated with that key.
3161    ///
3162    /// Because [`KeyedStream`] lazily groups values into buckets, this operator has zero computational
3163    /// cost and _does not_ require that the key type is hashable. Keyed streams are useful for
3164    /// performing grouped aggregations, but also for more precise ordering guarantees such as
3165    /// total ordering _within_ each group but no ordering _across_ groups.
3166    ///
3167    /// # Example
3168    /// ```rust
3169    /// # #[cfg(feature = "deploy")] {
3170    /// # use hydro_lang::prelude::*;
3171    /// # use futures::StreamExt;
3172    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
3173    /// process
3174    ///     .source_iter(q!(vec![(1, 2), (1, 3), (2, 4)]))
3175    ///     .into_keyed()
3176    /// #   .entries()
3177    /// # }, |mut stream| async move {
3178    /// // { 1: [2, 3], 2: [4] }
3179    /// # for w in vec![(1, 2), (1, 3), (2, 4)] {
3180    /// #     assert_eq!(stream.next().await.unwrap(), w);
3181    /// # }
3182    /// # }));
3183    /// # }
3184    /// ```
3185    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    /// Given a stream of pairs `(K, V)`, produces a new stream of unique keys `K`.
3204    /// # Example
3205    /// ```rust
3206    /// # #[cfg(feature = "deploy")] {
3207    /// # use hydro_lang::prelude::*;
3208    /// # use futures::StreamExt;
3209    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
3210    /// let tick = process.tick();
3211    /// let numbers = process.source_iter(q!(vec![(1, 2), (2, 3), (1, 3), (2, 4)]));
3212    /// let batch = numbers.batch(&tick, nondet!(/** test */));
3213    /// batch.keys().all_ticks()
3214    /// # }, |mut stream| async move {
3215    /// // 1, 2
3216    /// # assert_eq!(stream.next().await.unwrap(), 1);
3217    /// # assert_eq!(stream.next().await.unwrap(), 2);
3218    /// # }));
3219    /// # }
3220    /// ```
3221    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!(/** values are ignored */),
3228                    idempotent = manual_proof!(/** values are ignored */)
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    /// Returns a stream corresponding to the latest batch of elements being atomically
3240    /// processed. These batches are guaranteed to be contiguous across ticks and preserve
3241    /// the order of the input.
3242    ///
3243    /// # Non-Determinism
3244    /// The batch boundaries are non-deterministic and may change across executions.
3245    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    /// Yields the elements of this stream back into a top-level, asynchronous execution context.
3269    /// See [`Stream::atomic`] for more details.
3270    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    /// Consumes a stream of `Future<T>`, produces a new stream of the resulting `T` outputs.
3291    /// Future outputs are produced as available, regardless of input arrival order.
3292    ///
3293    /// # Example
3294    /// ```rust
3295    /// # #[cfg(feature = "deploy")] {
3296    /// # use std::collections::HashSet;
3297    /// # use futures::StreamExt;
3298    /// # use hydro_lang::prelude::*;
3299    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
3300    /// process.source_iter(q!([2, 3, 1, 9, 6, 5, 4, 7, 8]))
3301    ///     .map(q!(|x| async move {
3302    ///         tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
3303    ///         x
3304    ///     }))
3305    ///     .resolve_futures()
3306    /// #   },
3307    /// #   |mut stream| async move {
3308    /// // 1, 2, 3, 4, 5, 6, 7, 8, 9 (in any order)
3309    /// #       let mut output = HashSet::new();
3310    /// #       for _ in 1..10 {
3311    /// #           output.insert(stream.next().await.unwrap());
3312    /// #       }
3313    /// #       assert_eq!(
3314    /// #           output,
3315    /// #           HashSet::<i32>::from_iter(1..10)
3316    /// #       );
3317    /// #   },
3318    /// # ));
3319    /// # }
3320    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    /// Consumes a stream of `Future<T>`, produces a new stream of the resulting `T` outputs.
3333    /// Future outputs are produced in the same order as the input stream.
3334    ///
3335    /// # Example
3336    /// ```rust
3337    /// # #[cfg(feature = "deploy")] {
3338    /// # use std::collections::HashSet;
3339    /// # use futures::StreamExt;
3340    /// # use hydro_lang::prelude::*;
3341    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
3342    /// process.source_iter(q!([2, 3, 1, 9, 6, 5, 4, 7, 8]))
3343    ///     .map(q!(|x| async move {
3344    ///         tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
3345    ///         x
3346    ///     }))
3347    ///     .resolve_futures_ordered()
3348    /// #   },
3349    /// #   |mut stream| async move {
3350    /// // 2, 3, 1, 9, 6, 5, 4, 7, 8
3351    /// #       let mut output = Vec::new();
3352    /// #       for _ in 1..10 {
3353    /// #           output.push(stream.next().await.unwrap());
3354    /// #       }
3355    /// #       assert_eq!(
3356    /// #           output,
3357    /// #           vec![2, 3, 1, 9, 6, 5, 4, 7, 8]
3358    /// #       );
3359    /// #   },
3360    /// # ));
3361    /// # }
3362    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    /// Asynchronously yields this batch of elements outside the tick as an unbounded stream,
3380    /// which will stream all the elements across _all_ tick iterations by concatenating the batches.
3381    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    /// Synchronously yields this batch of elements outside the tick as an unbounded stream,
3399    /// which will stream all the elements across _all_ tick iterations by concatenating the batches.
3400    ///
3401    /// Unlike [`Stream::all_ticks`], this preserves synchronous execution, as the output stream
3402    /// is emitted in an [`Atomic`] context that will process elements synchronously with the input
3403    /// stream's [`Tick`] context.
3404    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    /// Transforms the stream using the given closure in "stateful" mode, where stateful operators
3420    /// such as `fold` retrain their memory across ticks rather than resetting across batches of
3421    /// input.
3422    ///
3423    /// This API is particularly useful for stateful computation on batches of data, such as
3424    /// maintaining an accumulated state that is up to date with the current batch.
3425    ///
3426    /// # Example
3427    /// ```rust
3428    /// # #[cfg(feature = "deploy")] {
3429    /// # use hydro_lang::prelude::*;
3430    /// # use futures::StreamExt;
3431    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
3432    /// let tick = process.tick();
3433    /// # // ticks are lazy by default, forces the second tick to run
3434    /// # tick.spin_batch(q!(1)).all_ticks().for_each(q!(|_| {}));
3435    /// # let batch_first_tick = process
3436    /// #   .source_iter(q!(vec![1, 2, 3, 4]))
3437    /// #  .batch(&tick, nondet!(/** test */));
3438    /// # let batch_second_tick = process
3439    /// #   .source_iter(q!(vec![5, 6, 7]))
3440    /// #   .batch(&tick, nondet!(/** test */))
3441    /// #   .defer_tick(); // appears on the second tick
3442    /// let input = // [1, 2, 3, 4 (first batch), 5, 6, 7 (second batch)]
3443    /// # batch_first_tick.chain(batch_second_tick);
3444    ///
3445    /// input.across_ticks(|s| s.count()).all_ticks()
3446    /// # }, |mut stream| async move {
3447    /// // [4, 7]
3448    /// assert_eq!(stream.next().await.unwrap(), 4);
3449    /// assert_eq!(stream.next().await.unwrap(), 7);
3450    /// # }));
3451    /// # }
3452    /// ```
3453    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    /// Shifts the elements in `self` to the **next tick**, so that the returned stream at tick `T`
3461    /// always has the elements of `self` at tick `T - 1`.
3462    ///
3463    /// At tick `0`, the output stream is empty, since there is no previous tick.
3464    ///
3465    /// This operator enables stateful iterative processing with ticks, by sending data from one
3466    /// tick to the next. For example, you can use it to compare inputs across consecutive batches.
3467    ///
3468    /// # Example
3469    /// ```rust
3470    /// # #[cfg(feature = "deploy")] {
3471    /// # use hydro_lang::prelude::*;
3472    /// # use futures::StreamExt;
3473    /// # tokio_test::block_on(hydro_lang::test_util::stream_transform_test(|process| {
3474    /// let tick = process.tick();
3475    /// // ticks are lazy by default, forces the second tick to run
3476    /// tick.spin_batch(q!(1)).all_ticks().for_each(q!(|_| {}));
3477    ///
3478    /// let batch_first_tick = process
3479    ///   .source_iter(q!(vec![1, 2, 3, 4]))
3480    ///   .batch(&tick, nondet!(/** test */));
3481    /// let batch_second_tick = process
3482    ///   .source_iter(q!(vec![0, 3, 4, 5, 6]))
3483    ///   .batch(&tick, nondet!(/** test */))
3484    ///   .defer_tick(); // appears on the second tick
3485    /// let changes_across_ticks = batch_first_tick.chain(batch_second_tick);
3486    ///
3487    /// changes_across_ticks.clone().filter_not_in(
3488    ///     changes_across_ticks.defer_tick() // the elements from the previous tick
3489    /// ).all_ticks()
3490    /// # }, |mut stream| async move {
3491    /// // [1, 2, 3, 4 /* first tick */, 0, 5, 6 /* second tick */]
3492    /// # for w in vec![1, 2, 3, 4, 0, 5, 6] {
3493    /// #     assert_eq!(stream.next().await.unwrap(), w);
3494    /// # }
3495    /// # }));
3496    /// # }
3497    /// ```
3498    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!(/** test */))
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!(/** test */));
3707            let v = use::snapshot(node.source_iter(q!(vec![1, 2, 3])).reduce(q!(|acc, v| *acc += v)), nondet!(/** test */));
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!(/** test */));
3745            let v = use::snapshot(node.source_iter(q!(vec![1, 2, 3])).reduce(q!(|acc, v| *acc += v)).into_singleton(), nondet!(/** test */));
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!(/** test */))
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!(/** test */)),
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!(/** test */))
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); // fails with nondet batching
3940        });
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!(/** test */))
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!(/** test */));
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!(/** test */));
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 // ∑ (k=1 to 4) S(4,k) × k! = 75
4014        )
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!(/** test */));
4028        let out_recv = batch
4029            .assume_ordering::<TotalOrder>(nondet!(/** test */))
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; // fails with assume_ordering
4036        });
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!(/** test */));
4049        let out_recv = batch
4050            .assume_ordering::<TotalOrder>(nondet!(/** test */))
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 // 4! * 2^{4 - 1}
4062        )
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!(/** test */))
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 // 2^4, { 0, 1, 2, 3 } can be a snapshot and 4 is always included
4088        )
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!(/** test */))
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!(/** test */));
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!(/** test */));
4166        complete_cycle_back.complete(
4167            ordered
4168                .clone()
4169                .batch(&node.tick(), nondet!(/** test */))
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!(/** test */));
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!(/** test */));
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!(/** test */))
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        // The return value of .inspect() is intentionally dropped.
4325        // Before the Null-root fix, this would silently do nothing.
4326        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        // The return value of .inspect() is bound to a variable that is still alive
4367        // when the flow is finalized by `deploy` below, so its `Drop` runs too late
4368        // to register a root the usual way. The FlowBuilder must yank the IR from
4369        // still-live collections when finalizing.
4370        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!(/** test */))
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            // Assert ordering preservation: elements from each input must
4467            // appear in their original relative order.
4468            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    /// Tests that `merge_ordered` passes through elements when only one input
4493    /// has data.
4494    #[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!(/** test */))
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        // Only one possible interleaving when one input is empty
4516        assert_eq!(instances, 1);
4517    }
4518
4519    /// Tests that `merge_ordered` correctly handles feedback cycles.
4520    /// An element output from `merge_ordered` is filtered and cycled back to
4521    /// one of its inputs. The one-at-a-time release must allow the cycled-back
4522    /// element to arrive and potentially be emitted before elements still
4523    /// waiting on the other input.
4524    #[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        // Create a forward ref for the cycle back
4533        let (complete_cycle_back, cycle_back) =
4534            node.forward_ref::<super::Stream<_, _, _, TotalOrder>>();
4535
4536        // merge_ordered: input (external) with cycle_back
4537        let merged = input.merge_ordered(cycle_back, nondet!(/** test */));
4538
4539        // Cycle back: elements equal to 1 get mapped to 10 and fed back
4540        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        // Send 1 and 2. Element 1 should cycle back as 10.
4545        // Valid orderings must have 1 before 10 (since 10 depends on 1).
4546        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            // 10 must always come after 1 (causal dependency)
4554            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            // Check if we see [1, 10, 2] — the cycled element beats the second input
4559            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    /// Tests that `merge_ordered` correctly interleaves when one input has a
4575    /// delayed element. With a: [1, _delay_, 2] and b: [3, 4], the delayed
4576    /// element 2 should be able to appear after b's elements.
4577    #[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!(/** test */))
4588            .sim_output();
4589
4590        let mut saw_delayed_interleaving = false;
4591        flow.sim().exhaustive(async || {
4592            // Send 1 from a, and 3, 4 from b
4593            in_send.send(1);
4594            in_send2.send(3);
4595            in_send2.send(4);
4596
4597            // Collect what's available so far
4598            let first_batch = out_recv.collect::<Vec<_>>().await;
4599
4600            // Now send the delayed element 2 from a
4601            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            // Check if we saw [1, 3, 4, 2] — the delayed interleaving
4611            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    /// Deploy test: `merge_ordered` with a delayed element on one input.
4623    /// Sends a=1, b=3, b=4, then after receiving those, sends a=2.
4624    /// Expects to see [1, 3, 4] first, then [2] — demonstrating that
4625    /// both inputs are pulled and the delayed element arrives later.
4626    #[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!(/** test */))
4640            .merge_ordered(
4641                input_b.assume_ordering(nondet!(/** test */)),
4642                nondet!(/** test */),
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        // Send a=1, b=3, b=4
4660        ext_a.send(1).await.unwrap();
4661        ext_b.send(3).await.unwrap();
4662        ext_b.send(4).await.unwrap();
4663
4664        // Collect the first 3 elements
4665        let mut received = Vec::new();
4666        for _ in 0..3 {
4667            received.push(ext_out.next().await.unwrap());
4668        }
4669
4670        // Now send the delayed a=2
4671        ext_a.send(2).await.unwrap();
4672        received.push(ext_out.next().await.unwrap());
4673
4674        // All elements should be present
4675        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!(/** test */)
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!(/** test */)
4771            ),
4772        );
4773
4774        // map with order_preserving should preserve monotonicity
4775        let doubled = sum.map(q!(
4776            |v| v * 2,
4777            order_preserving = manual_proof!(/** doubling preserves order */)
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    // === Compile-time type tests for join/cross_product ordering ===
4799
4800    #[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        // === Compile-time type tests for cross_product ordering ===
4839
4840        #[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    } // mod join_ordering_type_tests
4864
4865    // === Runtime correctness tests for bounded join/cross_product ===
4866
4867    #[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!(/** test */))
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!(/** test */))
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        // Create a bounded stream (source_iter is bounded within a tick)
4952        let my_stream = p1.source_iter(q!(1..=5i32));
4953
4954        let stream_ref = my_stream.by_ref();
4955
4956        // Use the stream ref to get the vec's length
4957        let out_port = p1
4958            .source_iter(q!([()]))
4959            .map(q!(|_| stream_ref.len() as i32))
4960            .send_bincode_external(&external);
4961
4962        // Also consume the stream via pipe
4963        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        // stream has 5 elements
4979        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        // Create a bounded stream
4992        let my_stream = p1.source_iter(q!(1..=3i32));
4993
4994        let stream_ref = my_stream.by_ref();
4995
4996        // Sum the referenced vec's contents
4997        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        // sum of 1+2+3 = 6
5018        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        // Create a bounded stream — no pipe consumer, only ref
5031        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        // Create a bounded stream
5066        let my_stream = p1.source_iter(q!(1..=5i32));
5067
5068        let stream_mut = my_stream.by_mut();
5069
5070        // Mutably reference the buffer to retain only items > 3
5071        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        // After retain(> 3): [4, 5] => len = 2
5095        assert_eq!(result, 2);
5096    }
5097
5098    /// A map with a mut singleton ref on an unordered input should produce > 1
5099    /// simulation instance because the ordering of elements through the mut closure
5100    /// is non-deterministic.
5101    #[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!(/** test */));
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!(/** test */)
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    /// A `scan` closure that captures a bounded singleton by reference should compile,
5142    /// run correctly, and (because the input is totally ordered) explore a single
5143    /// simulation instance.
5144    #[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!(/** test */));
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            // offset = 10, running accumulator starts at 0:
5179            //   x=1: acc += 1 + 10 = 11 -> 11
5180            //   x=2: acc += 2 + 10 = 12 -> 23
5181            //   x=3: acc += 3 + 10 = 13 -> 36
5182            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    /// A map with a mut singleton ref on a top-level unordered input should produce > 1
5193    /// simulation instance. Currently panics because `observe_nondet` doesn't support
5194    /// top-level bounded inputs yet.
5195    #[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!(/** test */)
5218            ))
5219            .assume_ordering::<TotalOrder>(nondet!(/** test */))
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}