Skip to main content

KeyedBatchHook

Struct KeyedBatchHook 

Source
pub struct KeyedBatchHook<K, V, O: Ordering = TotalOrder, R: Retries = ExactlyOnce> { /* private fields */ }
Expand description

A hook handle controlling a batch operator over a keyed stream with keys K, values V, per-key value ordering O, and retry guarantee R (mirroring the type of the keyed stream being batched).

A decision for a keyed batch hook says which buffered (key, value) entries form the next batch released into the tick. See hydro_lang::sim::hooks for the decisions offered.

Implementations§

Source§

impl<K, V, O: Ordering, R: Retries> KeyedBatchHook<K, V, O, R>

Source

pub fn pause(&self)

Available on crate feature sim only.

Declares that buffering at this operator is intended: while paused, the hook is exempt from the missing-decision error, never causes its tick to run, and — if its tick runs anyway because other hooks feed it — contributes its “nothing new” behavior each time. Scripting any decision implicitly resumes the hook.

A pause takes its place in the script like everything else: requested while a decision is still pending, the hold begins once that decision has been consumed.

Source

pub fn resume(&self)

Available on crate feature sim only.

Ends a Self::pause (and clears Self::auto_pause mode).

Source

pub fn auto_pause(&self)

Available on crate feature sim only.

Sets a standing mode where this hook only ever acts when scripted: it holds immediately, and every scripted decision leaves a fresh hold in place behind it.

This deliberately opts out of the forgotten-hook protection: if the test forgets a step, the operator silently holds its data instead of failing. The one auto_pause() line at the top of a test is the reviewer-visible marker that this hook’s timing is entirely script-driven, missed steps and all.

Source

pub async fn pause_while<Fut: Future>(&self, body: Fut) -> Fut::Output

Available on crate feature sim only.

Pauses the hook exactly for the duration of body (resuming even on panic), so a bracketed buffering phase cannot leak a paused hook.

Source

pub fn pause_until( &self, predicate: impl Fn(&BatchStatus) -> bool + Unpin, ) -> PauseUntilFuture<BatchStatus, impl Fn(&BatchStatus) -> bool + Unpin>

Available on crate feature sim only.

Pauses the hook and returns a future that resolves once the hook’s pending-input status satisfies predicate — a synchronization point for scripts where the right decision is not knowable upfront. The status is read on demand at every scheduling point. After the future resolves, the hook is unpaused; the ordinary missing-decision error applies from there on.

Source

pub fn pause_until_count( &self, n: usize, ) -> PauseUntilFuture<BatchStatus, impl Fn(&BatchStatus) -> bool + Unpin>

Available on crate feature sim only.

Pauses the hook and returns a future that resolves once at least n entries are buffered (in total, across all keys); see Self::pause_until.

Source§

impl<K, V, R: Retries> KeyedBatchHook<K, V, TotalOrder, R>

Source

pub fn release( &self, counts: impl IntoIterator<Item = (K, usize)>, ) -> DecisionFuture
where K: Hash + Eq,

Available on crate feature sim only.

Scripts the next batch to be exactly the next count buffered values of each named key. The tick fires at the first moment the decision can be honored in full.

§Panics

Panics immediately if counts names the same key more than once.

Source

pub fn release_values( &self, entries: impl IntoIterator<Item = (K, V)>, ) -> DecisionFuture

Available on crate feature sim only.

Scripts the next batch to be exactly these (key, value) entries. Each key’s values must match that key’s buffered prefix in order (the interleaving of different keys in the scripted sequence is irrelevant): a mismatching available value panics immediately, while a matching but incomplete prefix waits for the remaining values to arrive.

Source

pub fn release_all(&self) -> DecisionFuture

Available on crate feature sim only.

Scripts the next batch to be everything that has arrived by the time the tick fires. Under fuzzing, the released contents co-vary with the schedule being explored; use Self::release_values to name them exactly.

Source

pub fn release_empty(&self) -> DecisionFuture

Available on crate feature sim only.

Scripts the next batch to be empty, holding everything buffered. Shorthand for Self::release_values with no entries.

Source§

impl<K, V, R: Retries> KeyedBatchHook<K, V, NoOrder, R>

Source

pub fn release_values( &self, entries: impl IntoIterator<Item = (K, V)>, ) -> DecisionFuture

Available on crate feature sim only.

Scripts the next batch to contain exactly these (key, value) entries. Values are matched per key as multisets (independently of arrival order); duplicates request the corresponding number of equal buffered items. The tick fires once every requested entry exists.

Source

pub fn release_all(&self) -> DecisionFuture

Available on crate feature sim only.

Scripts the next batch to be everything that has arrived by the time the tick fires. Under fuzzing, the released contents co-vary with the schedule being explored; use Self::release_values to name them exactly.

Source

pub fn release_empty(&self) -> DecisionFuture

Available on crate feature sim only.

Scripts the next batch to be empty, holding everything buffered. Shorthand for Self::release_values with no entries.

Trait Implementations§

Source§

impl<K, V, O: Ordering, R: Retries> Clone for KeyedBatchHook<K, V, O, R>

Source§

fn clone(&self) -> Self

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl<K, V, O: Ordering, R: Retries> Copy for KeyedBatchHook<K, V, O, R>

Source§

impl<K, V, O: Ordering, R: Retries> Debug for KeyedBatchHook<K, V, O, R>

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl<K, V, O: Ordering, R: Retries> SimHook for KeyedBatchHook<K, V, O, R>

Source§

fn create(next_id: &mut dyn FnMut() -> usize) -> Self

Creates every handle in this value, allocating fresh IDs via next_id.

Auto Trait Implementations§

§

impl<K, V, O, R> Freeze for KeyedBatchHook<K, V, O, R>
where PhantomData<fn(K, V, O, R)>: Freeze,

§

impl<K, V, O, R> RefUnwindSafe for KeyedBatchHook<K, V, O, R>
where PhantomData<fn(K, V, O, R)>: RefUnwindSafe,

§

impl<K, V, O, R> Send for KeyedBatchHook<K, V, O, R>
where PhantomData<fn(K, V, O, R)>: Send,

§

impl<K, V, O, R> Sync for KeyedBatchHook<K, V, O, R>
where PhantomData<fn(K, V, O, R)>: Sync,

§

impl<K, V, O, R> Unpin for KeyedBatchHook<K, V, O, R>
where PhantomData<fn(K, V, O, R)>: Unpin,

§

impl<K, V, O, R> UnsafeUnpin for KeyedBatchHook<K, V, O, R>
where PhantomData<fn(K, V, O, R)>: UnsafeUnpin,

§

impl<K, V, O, R> UnwindSafe for KeyedBatchHook<K, V, O, R>
where PhantomData<fn(K, V, O, R)>: UnwindSafe,

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
§

impl<T> AsOut<T> for T
where T: Copy,

§

fn as_out(&mut self) -> Out<'_, T>

Returns an out reference to self.
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> Instrument for T

§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided [Span], returning an Instrumented wrapper. Read more
§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
§

impl<Unshared, Shared> IntoShared<Shared> for Unshared
where Shared: FromUnshared<Unshared>,

§

fn into_shared(self) -> Shared

Creates a shared type from an unshared type.
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T> ToSinkBuild for T

Source§

fn iter_to_sink_build(self) -> SendIterBuild<Self>
where Self: Sized + Iterator,

Starts a SinkBuild adaptor chain to send all items from self as an Iterator.
Source§

fn stream_to_sink_build(self) -> SendStreamBuild<Self>
where Self: Sized + Stream,

Starts a SinkBuild adaptor chain to send all items from self as a [Stream].
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a [WithDispatch] wrapper. Read more
§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a [WithDispatch] wrapper. Read more