Skip to main content

ClusterSupervisor

Struct ClusterSupervisor 

Source
pub struct ClusterSupervisor<L: PeerLiveness, A: ShardAdopter> { /* private fields */ }
Expand description

Watches peer liveness and auto-adopts a dead peer’s shards (SS-5b).

Implementations§

Source§

impl<L: PeerLiveness, A: ShardAdopter> ClusterSupervisor<L, A>

Source

pub fn new( liveness: Arc<L>, adopter: Arc<A>, peers: Vec<WatchedPeer>, config: SupervisorConfig, ) -> Self

Build a supervisor over peers, polling liveness and calling adopter.adopt_shards on confirmed peer death. Peers with no owned shards are dropped from the watch set (nothing to adopt for them).

Source

pub fn with_publisher( self, publisher: Arc<ClusterEventPublisher>, self_node: impl Into<String>, ) -> Self

Attach the WS3 cluster-event publisher and this node’s name so tick() emits topology deltas. Pure builder addition — a supervisor without it behaves exactly as before (every existing test passes new only).

Source

pub fn watches_any(&self) -> bool

Whether this supervisor watches any peer (false when no peer declared owned shards — the loop would do nothing, so the caller can skip spawning).

Source

pub fn adopter(&self) -> &A

Borrow the adopter (the engine, in production) this supervisor drives. Lets a test inspect the engine it auto-adopts onto after the loss.

Source

pub async fn tick(&mut self) -> Vec<String>

Run ONE poll tick: observe every watched peer’s liveness, advance the debounce counters, and adopt the shards of any peer that has now been down for confirmations consecutive ticks and is not yet adopted.

Returned is the list of peer names adopted on THIS tick (empty on a quiet tick), so a test can assert exactly when adoption fires. Extracted from the loop so the debounce decision is unit-testable without real time.

Source

pub async fn run(self, shutdown: Receiver<bool>)

Drive the poll loop until shutdown flips true, ticking every poll_interval. Consumes self; spawn it as a background task.

Auto Trait Implementations§

§

impl<L, A> !RefUnwindSafe for ClusterSupervisor<L, A>

§

impl<L, A> !UnwindSafe for ClusterSupervisor<L, A>

§

impl<L, A> Freeze for ClusterSupervisor<L, A>
where Arc<L>: Freeze, Arc<A>: Freeze,

§

impl<L, A> Send for ClusterSupervisor<L, A>
where Arc<L>: Send, Arc<A>: Send,

§

impl<L, A> Sync for ClusterSupervisor<L, A>
where Arc<L>: Sync, Arc<A>: Sync,

§

impl<L, A> Unpin for ClusterSupervisor<L, A>
where Arc<L>: Unpin, Arc<A>: Unpin,

§

impl<L, A> UnsafeUnpin for ClusterSupervisor<L, A>
where Arc<L>: UnsafeUnpin, Arc<A>: UnsafeUnpin,

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
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<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

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

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

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

Source§

fn into_maybe_undefined(self) -> MaybeUndefined<T>

Converts this value into a three-state builder argument.
Source§

impl<T> IntoOption<T> for T

Source§

fn into_option(self) -> Option<T>

Converts this value into an optional builder argument.
Source§

impl<T> IntoRequest<T> for T

Source§

fn into_request(self) -> Request<T>

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

impl<L> LayerExt<L> for L

Source§

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.
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

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 = !

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.
Source§

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

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

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
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

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