pub struct ShuffleSender { /* private fields */ }Expand description
Lazy pool of ordered streams per peer.
Implementations§
Source§impl ShuffleSender
impl ShuffleSender
Sourcepub fn new(local_id: ShufflePeerId, incarnation: Uuid) -> Self
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.
Sourcepub const fn local_id(&self) -> ShufflePeerId
pub const fn local_id(&self) -> ShufflePeerId
Node id bound into every outbound stream handshake.
Sourcepub fn install_process_lease_deadline(
&self,
deadline: Arc<LeaseDeadline>,
) -> Result<()>
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.
Sourcepub fn bind_process_lease_deadline_pair(
&self,
receiver: &ShuffleReceiver,
deadline: Arc<LeaseDeadline>,
) -> Result<()>
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.
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 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.
Sourcepub fn suspend_assignment_fence(&self)
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.
Sourcepub fn invalidate_assignment_fence(&self)
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.
Sourcepub fn assignment_version(&self) -> u64
pub fn assignment_version(&self) -> u64
Assignment scope currently accepted for newly enqueued 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 outbound streams.
Sourcepub fn set_recovery_gen(&self, gen: u64)
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.
Sourcepub fn recovery_gen(&self) -> u64
pub fn recovery_gen(&self) -> u64
Recovery scope currently accepted for newly enqueued frames.
Sourcepub const fn incarnation(&self) -> Uuid
pub const fn incarnation(&self) -> Uuid
Process incarnation bound into every outbound stream handshake.
Sourcepub fn with_kv(
local_id: ShufflePeerId,
kv: Arc<dyn ClusterKv>,
incarnation: Uuid,
) -> Self
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.
Sourcepub fn register_peer(&self, peer: ShufflePeerId, addr: SocketAddr)
pub fn register_peer(&self, peer: ShufflePeerId, addr: SocketAddr)
Register (or update) a peer’s shuffle address.
Sourcepub async fn send_to(
&self,
peer: ShufflePeerId,
msg: &ShuffleMessage,
) -> Result<()>
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.
Sourcepub async fn send_to_for_assignment(
&self,
peer: ShufflePeerId,
expected_assignment_version: u64,
msg: &ShuffleMessage,
) -> Result<()>
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.
Sourcepub async fn establish_assignment_mesh(
&self,
assignment_fence: &CheckpointAssignmentFence,
) -> Result<()>
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.
Sourcepub async fn fan_out_barrier(
&self,
peers: &[ShufflePeerId],
barrier: CheckpointBarrier,
assignment_fence: &CheckpointAssignmentFence,
) -> Result<()>
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§
Auto Trait Implementations§
impl !Freeze for ShuffleSender
impl !RefUnwindSafe for ShuffleSender
impl !UnwindSafe for ShuffleSender
impl Send for ShuffleSender
impl Sync for ShuffleSender
impl Unpin for ShuffleSender
impl UnsafeUnpin for ShuffleSender
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].