Skip to main content

ShuffleSender

Struct ShuffleSender 

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

Lazy pool of ordered streams per peer.

Implementations§

Source§

impl ShuffleSender

Source

pub fn new(local_id: ShufflePeerId, incarnation: Uuid) -> Self

Empty sender; peers arrive via Self::register_peer or KV discovery.

§Panics

Panics when the node ID is zero or the process incarnation is nil.

Source

pub const fn local_id(&self) -> ShufflePeerId

Node id bound into every outbound stream handshake.

Source

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

Bind outbound 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 bind_process_lease_deadline_pair( &self, receiver: &ShuffleReceiver, deadline: Arc<LeaseDeadline>, ) -> Result<()>

Bind both directions of one local shuffle fabric to the same process lease.

Validation and installation are serialized across both handles, so an incompatible receiver cannot leave only the sender bound.

§Errors

Returns an error without changing either handle when the deadline is expired, either handle is already bound to a different deadline, or an active assignment prevents a new binding.

Source

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

Install the exact assignment certificate accepted by outbound shuffle streams. Changing scope closes all pooled streams and resets their scoped sequences.

§Errors

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

Source

pub fn suspend_assignment_fence(&self)

Temporarily close outbound admission because the durable assignment authority could not be read. Unlike invalidation, this preserves the exact certificate and sequence counters so the same durable head can resume without manufacturing a delivery gap.

Source

pub fn invalidate_assignment_fence(&self)

Reject every stream while a newer assignment is being adopted but has not yet earned an owner-complete process certificate. The last certificate is retained only as a monotonic conflict floor; it is inactive while Self::assignment_version is zero.

Source

pub fn assignment_version(&self) -> u64

Assignment scope currently accepted for newly enqueued frames.

Source

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

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

Source

pub fn set_recovery_gen(&self, gen: u64)

Advance the generation stamped onto outbound data frames. Called after a coordinated rewind so peers can discard anything this node produced before it.

Source

pub fn recovery_gen(&self) -> u64

Recovery scope currently accepted for newly enqueued frames.

Source

pub const fn incarnation(&self) -> Uuid

Process incarnation bound into every outbound stream handshake.

Source

pub fn with_kv( local_id: ShufflePeerId, kv: Arc<dyn ClusterKv>, incarnation: Uuid, ) -> Self

Sender that falls back to kv discovery for peers not previously registered.

Source

pub fn register_peer(&self, peer: ShufflePeerId, addr: SocketAddr)

Register (or update) a peer’s shuffle address.

Source

pub async fn send_to( &self, peer: ShufflePeerId, msg: &ShuffleMessage, ) -> Result<()>

Send msg to peer, opening a client-streaming call if necessary.

§Errors

Returns io::Error when the peer is unregistered/undiscoverable, the endpoint cannot be built, or the per-peer stream has shut down.

Source

pub async fn send_to_for_assignment( &self, peer: ShufflePeerId, expected_assignment_version: u64, msg: &ShuffleMessage, ) -> Result<()>

Send only while the sender remains in expected_assignment_version.

Routing paths use this to prevent data sliced with an old ownership publication from entering a stream opened under a newer assignment scope.

§Errors

Returns io::Error when the expected assignment is no longer current, in addition to the errors returned by Self::send_to.

Source

pub async fn establish_assignment_mesh( &self, assignment_fence: &CheckpointAssignmentFence, ) -> Result<()>

Establish an exact-scope stream to every remote assignment participant without sending a frame or consuming a delivery sequence.

§Errors

Returns the first peer connection or handshake error after attempting the full roster.

Source

pub async fn fan_out_barrier( &self, peers: &[ShufflePeerId], barrier: CheckpointBarrier, assignment_fence: &CheckpointAssignmentFence, ) -> Result<()>

Ship barrier to every required peer. All peers are attempted, but any failure rejects the cut so the coordinator can abort promptly.

§Errors

Returns the first peer error after attempting the full fan-out.

Trait Implementations§

Source§

impl Debug for ShuffleSender

Source§

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

Formats the value using the given formatter. 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