Skip to main content

ShuffleReceiver

Struct ShuffleReceiver 

Source
pub struct ShuffleReceiver { /* private fields */ }
Expand description

Inbound side of the shuffle fabric: a Tonic ShuffleTransport server that surfaces every received frame, attributed to its peer, on the bounded queue.

Implementations§

Source§

impl ShuffleReceiver

Source

pub async fn bind( local_id: ShufflePeerId, addr: SocketAddr, receiver_incarnation: Uuid, ) -> Result<Self>

Bind on addr and start serving; the resolved address is at Self::local_addr.

§Errors

Returns io::Error on bind failure.

Source

pub fn install_process_lease_deadline( &self, deadline: Arc<LeaseDeadline>, ) -> Result<()>

Bind inbound admission to this process’s renewable cluster lease.

This must be installed before an assignment certificate is activated.

§Errors

Returns an error for an expired lease, a different previously installed deadline, or an already-active assignment.

Source

pub fn install_assignment_fence( &self, fence: &CheckpointAssignmentFence, owners: &[ShufflePeerId], ) -> Result<bool>

Install the exact assignment certificate accepted by inbound streams. Existing streams are rejected on their next frame; the new delivery domain starts at zero.

§Errors

Returns an error for a malformed certificate, a same-version certificate conflict, or a certificate that does not bind this exact receiver process.

Source

pub fn suspend_assignment_fence(&self)

Temporarily close inbound admission after a transient durable-authority read failure. Pending handshakes are discarded, while committed delivery expectations remain bound to the retained certificate for exact same-version reactivation.

Source

pub fn invalidate_assignment_fence(&self)

Reject all streams while an adopted owner map is awaiting its owner-complete process certificate. Pending tokens and sequence expectations cannot cross this boundary.

Source

pub fn assignment_version(&self) -> u64

Assignment scope currently accepted for inbound streams and staged frames.

Source

pub fn active_assignment_digest(&self) -> Option<[u8; 32]>

Digest of the exact assignment certificate currently active for inbound streams.

Source

pub fn set_recovery_gen(&self, gen: u64)

Advance the inbound recovery scope. Existing streams are rejected; the next exact handshake re-baselines sequence after the coordinated rewind.

Source

pub fn recovery_gen(&self) -> u64

Recovery scope currently accepted for inbound streams and staged frames.

Source

pub const fn incarnation(&self) -> Uuid

Process incarnation bound into every accepted stream handshake.

Source

pub const fn local_id(&self) -> ShufflePeerId

Node id bound into every accepted stream handshake.

Source

pub fn delivery_loss_incidents(&self) -> Arc<AtomicU64>

Cumulative delivery-loss incidents. A value above the recovered floor means an epoch must not seal; exact missing-frame magnitude is diagnostic and never trusted for this non-wrapping correctness signal.

Source

pub fn recovered_delivery_loss_incidents(&self) -> Arc<AtomicU64>

Delivery-loss incidents covered by a completed coordinated rewind. The callback uses this durable in-process floor so repaired incidents do not fault restored state again.

Source

pub fn has_unrecovered_delivery_loss(&self) -> bool

Whether an admitted delivery loss has not yet been repaired by coordinated recovery.

Source

pub fn complete_recovery(&self, gen: u64) -> bool

Mark the loss cutoff captured when gen began as repaired. Loss detected after the generation advanced remains above this floor and still faults the next checkpoint.

Returns false when gen is not the currently prepared recovery generation.

Source

pub async fn bind_with_kv( local_id: ShufflePeerId, addr: SocketAddr, kv: Arc<dyn ClusterKv>, receiver_incarnation: Uuid, ) -> Result<Self>

Bind and publish the listener’s address into kv for peer discovery.

§Errors

Returns io::Error on bind failure.

Source

pub fn local_addr(&self) -> SocketAddr

Local socket address the server is bound to.

Source

pub async fn recv(&self) -> Option<ReceivedShuffle>

Await the next (peer_id, msg); None once the server stops and the queue drains. Concurrent callers serialise via rx_returned; cancellation-safe.

Source

pub fn drain_available(&self) -> Vec<ReceivedShuffle>

Drain every currently available admitted message without blocking; empty if a recv() holds the receiver.

Source

pub fn drain_checkpointed_data_for(&self, stage: &str) -> Vec<ReceivedBatch>

Non-blocking drain of checkpointed operator batches for stage.

Source

pub fn drain_staged_barriers(&self) -> Vec<ReceivedShuffle>

Take the barriers stashed by Self::drain_checkpointed_data_for.

Source

pub fn has_staged_checkpoint_barriers(&self) -> bool

Whether an in-band barrier currently blocks normal holdover draining.

Source

pub fn stage_checkpointed_inbound(&self) -> bool

Move available inbound frames into the bounded holdover through the first barrier. Data remains keyed for its owning stage.

Source

pub fn retire_checkpoint_barriers( &self, attempt: CheckpointAttempt, assignment_digest: [u8; 32], ) -> Result<()>

Retire barrier markers through a certified terminal checkpoint without dropping data.

§Errors

Rejects a noncanonical attempt and exact markers whose assignment digest differs from the durable terminal outcome.

Source

pub fn stash_barrier(&self, barrier: ReceivedShuffle)

Re-stash a peer barrier pulled while aligning a different checkpoint, so a lagging node still sees it when it reaches that checkpoint.

Source

pub fn drain_checkpointed_staged(&self) -> Vec<(String, ReceivedBatch)>

Drain every staged checkpointed operator batch.

Source

pub fn drain_checkpointed_holdover( &self, ) -> Result<Vec<(String, ReceivedBatch)>>

Drain only the checkpointed batches already in the holdover. Unlike Self::drain_checkpointed_staged, this does not consume the live receive queue, whose data/barrier order must remain visible to checkpoint alignment.

§Errors

Returns a typed cancellation while assignment consumption is suspended, or a process lease error after this receiver loses execution authority. An error observed before extraction leaves the holdover untouched; the caller revalidates authority immediately after a successful transfer.

Source

pub fn drain_all_staged(&self) -> Vec<(String, ReceivedBatch)>

Empty the per-stage holdover, returning every buffered (stage, batch).

Trait Implementations§

Source§

impl Debug for ShuffleReceiver

Source§

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

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

impl Drop for ShuffleReceiver

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more

Auto Trait Implementations§

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> ArchivePointee for T

§

type ArchivedMetadata = ()

The archived version of the pointer metadata for this type.
§

fn pointer_metadata( _: &<T as ArchivePointee>::ArchivedMetadata, ) -> <T as Pointee>::Metadata

Converts some archived metadata to the pointer metadata for itself.
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> 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<T> IntoRequest<T> for T

§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
§

impl<L> LayerExt<L> for L

§

fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>
where L: Layer<S>,

Applies the layer to a service and wraps it in [Layered].
§

impl<T> LayoutRaw for T

§

fn layout_raw(_: <T as Pointee>::Metadata) -> Result<Layout, LayoutError>

Returns the layout of the type.
§

impl<T, N1, N2> Niching<NichedOption<T, N1>> for N2
where T: SharedNiching<N1, N2>, N1: Niching<T>, N2: Niching<T>,

§

unsafe fn is_niched(niched: *const NichedOption<T, N1>) -> bool

Returns whether the given value has been niched. Read more
§

fn resolve_niched(out: Place<NichedOption<T, N1>>)

Writes data to out indicating that a T is niched.
§

impl<T> Pointee for T

§

type Metadata = ()

The metadata type for pointers and references to this type.
§

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

§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] only if self and other return Action::Follow. Read more
§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] if either self or other returns Action::Follow. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

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

Source§

type Error = Infallible

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

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

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