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
impl ShuffleReceiver
Sourcepub async fn bind(
local_id: ShufflePeerId,
addr: SocketAddr,
receiver_incarnation: Uuid,
) -> Result<Self>
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.
Sourcepub fn install_process_lease_deadline(
&self,
deadline: Arc<LeaseDeadline>,
) -> Result<()>
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.
Sourcepub fn install_assignment_fence(
&self,
fence: &CheckpointAssignmentFence,
owners: &[ShufflePeerId],
) -> Result<bool>
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.
Sourcepub fn suspend_assignment_fence(&self)
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.
Sourcepub fn invalidate_assignment_fence(&self)
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.
Sourcepub fn assignment_version(&self) -> u64
pub fn assignment_version(&self) -> u64
Assignment scope currently accepted for inbound streams and staged frames.
Sourcepub fn active_assignment_digest(&self) -> Option<[u8; 32]>
pub fn active_assignment_digest(&self) -> Option<[u8; 32]>
Digest of the exact assignment certificate currently active for inbound streams.
Sourcepub fn set_recovery_gen(&self, gen: u64)
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.
Sourcepub fn recovery_gen(&self) -> u64
pub fn recovery_gen(&self) -> u64
Recovery scope currently accepted for inbound streams and staged frames.
Sourcepub const fn incarnation(&self) -> Uuid
pub const fn incarnation(&self) -> Uuid
Process incarnation bound into every accepted stream handshake.
Sourcepub const fn local_id(&self) -> ShufflePeerId
pub const fn local_id(&self) -> ShufflePeerId
Node id bound into every accepted stream handshake.
Sourcepub fn delivery_loss_incidents(&self) -> Arc<AtomicU64> ⓘ
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.
Sourcepub fn recovered_delivery_loss_incidents(&self) -> Arc<AtomicU64> ⓘ
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.
Sourcepub fn has_unrecovered_delivery_loss(&self) -> bool
pub fn has_unrecovered_delivery_loss(&self) -> bool
Whether an admitted delivery loss has not yet been repaired by coordinated recovery.
Sourcepub fn complete_recovery(&self, gen: u64) -> bool
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.
Sourcepub async fn bind_with_kv(
local_id: ShufflePeerId,
addr: SocketAddr,
kv: Arc<dyn ClusterKv>,
receiver_incarnation: Uuid,
) -> Result<Self>
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.
Sourcepub fn local_addr(&self) -> SocketAddr
pub fn local_addr(&self) -> SocketAddr
Local socket address the server is bound to.
Sourcepub async fn recv(&self) -> Option<ReceivedShuffle>
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.
Sourcepub fn drain_available(&self) -> Vec<ReceivedShuffle>
pub fn drain_available(&self) -> Vec<ReceivedShuffle>
Drain every currently available admitted message without blocking;
empty if a recv() holds the receiver.
Sourcepub fn drain_checkpointed_data_for(&self, stage: &str) -> Vec<ReceivedBatch>
pub fn drain_checkpointed_data_for(&self, stage: &str) -> Vec<ReceivedBatch>
Non-blocking drain of checkpointed operator batches for stage.
Sourcepub fn drain_staged_barriers(&self) -> Vec<ReceivedShuffle>
pub fn drain_staged_barriers(&self) -> Vec<ReceivedShuffle>
Take the barriers stashed by Self::drain_checkpointed_data_for.
Sourcepub fn has_staged_checkpoint_barriers(&self) -> bool
pub fn has_staged_checkpoint_barriers(&self) -> bool
Whether an in-band barrier currently blocks normal holdover draining.
Sourcepub fn stage_checkpointed_inbound(&self) -> bool
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.
Sourcepub fn retire_checkpoint_barriers(
&self,
attempt: CheckpointAttempt,
assignment_digest: [u8; 32],
) -> Result<()>
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.
Sourcepub fn stash_barrier(&self, barrier: ReceivedShuffle)
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.
Sourcepub fn drain_checkpointed_staged(&self) -> Vec<(String, ReceivedBatch)>
pub fn drain_checkpointed_staged(&self) -> Vec<(String, ReceivedBatch)>
Drain every staged checkpointed operator batch.
Sourcepub fn drain_checkpointed_holdover(
&self,
) -> Result<Vec<(String, ReceivedBatch)>>
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.
Sourcepub fn drain_all_staged(&self) -> Vec<(String, ReceivedBatch)>
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
impl Debug for ShuffleReceiver
Source§impl Drop for ShuffleReceiver
impl Drop for ShuffleReceiver
Auto Trait Implementations§
impl !Freeze for ShuffleReceiver
impl !RefUnwindSafe for ShuffleReceiver
impl !UnwindSafe for ShuffleReceiver
impl Send for ShuffleReceiver
impl Sync for ShuffleReceiver
impl Unpin for ShuffleReceiver
impl UnsafeUnpin for ShuffleReceiver
Blanket Implementations§
§impl<T> ArchivePointee for T
impl<T> ArchivePointee for T
§type ArchivedMetadata = ()
type ArchivedMetadata = ()
§fn pointer_metadata(
_: &<T as ArchivePointee>::ArchivedMetadata,
) -> <T as Pointee>::Metadata
fn pointer_metadata( _: &<T as ArchivePointee>::ArchivedMetadata, ) -> <T as Pointee>::Metadata
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
§impl<T> Instrument for T
impl<T> Instrument for T
§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self>
fn into_either(self, into_left: bool) -> Either<Self, Self>
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 moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
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
impl<T> IntoRequest<T> for T
§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
T in a tonic::Request§impl<L> LayerExt<L> for L
impl<L> LayerExt<L> for L
§fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>where
L: Layer<S>,
fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>where
L: Layer<S>,
Layered].