Skip to main content

ClusterController

Struct ClusterController 

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

Facade composing the cluster-control primitives.

Implementations§

Source§

impl ClusterController

Source

pub fn new( instance_id: NodeId, kv: Arc<dyn ClusterKv>, snapshot: Option<Arc<AssignmentSnapshotStore>>, members_rx: Receiver<Vec<NodeInfo>>, ) -> Self

Wrap the given primitives.

Source

pub fn register_query_handler(&self, handler: Arc<dyn RemoteQueryHandler>)

Register the handler serving cross-node RemoteScan.

Source

pub fn query_client_pool(&self) -> &QueryClientPool

Access the connection pool for remote queries.

Source

pub fn cluster_min_watermark(&self) -> Option<i64>

Latest cluster-wide minimum watermark seen by this instance. None until the leader has published a Commit with a populated min_watermark_ms.

Source

pub fn publish_cluster_min_watermark(&self, wm: i64)

Leader-side monotonic publish. The leader computes the cluster-wide minimum watermark in await_prepare_quorum (its own local watermark folded with every follower’s ack) and must mirror it into the controller atomic so its own operators see the same value that followers pick up via observe_barrier on the matching Commit. Never lowers the published value — event-time progress is monotonic.

Source

pub fn instance_id(&self) -> NodeId

This instance’s ID.

Source

pub fn kv(&self) -> &Arc<dyn ClusterKv>

The cluster gossip KV, exposed so higher layers can advertise/discover per-stream state alongside the control-plane keys.

Source

pub fn current_leader(&self) -> Option<NodeId>

Current leader (lowest id among Active peers plus self).

Source

pub fn is_leader(&self) -> bool

True if this instance is currently the leader.

Source

pub fn set_active(&self, active: bool)

Mark this node’s active status.

Source

pub fn live_instances(&self) -> Vec<NodeId>

Live instance IDs: Active peers plus self.

Source

pub fn note_unresponsive(&self, peers: &[NodeId])

Record peers that failed to ack a capture quorum in time.

Source

pub fn note_responsive(&self, peers: &[NodeId])

Clear peers that acked (they are demonstrably alive).

Source

pub fn is_recently_unresponsive(&self, peer: NodeId) -> bool

Whether peer failed a capture quorum within the TTL window.

Source

pub fn begin_drain(&self)

Mark this node as draining. Idempotent.

Source

pub fn is_draining(&self) -> bool

Whether this node is draining.

Source

pub fn assignable_instances(&self) -> Vec<NodeId>

Node ids eligible to own vnodes: Active peers, plus self unless this node is draining. Mirrors how Self::live_instances folds self in, but filters non-Active peers (see assignable_node_ids) so Joining/Suspected/Draining/Left nodes never receive vnodes.

Source

pub fn set_self_locality(&self, locality: Locality)

Record this node’s own locality. Call once at startup.

Source

pub fn assignable_with_locality(&self) -> Vec<(NodeId, Locality)>

Self::assignable_instances paired with each node’s Locality (peers’ from members_rx, self’s from Self::set_self_locality).

Source

pub fn members_watch(&self) -> Receiver<Vec<NodeInfo>>

Cloneable membership watch. Background tasks subscribe to this to react to join/leave events (changed().await) without polling Self::live_instances on a timer.

Source

pub async fn announce_snapshot_version(&self, version: u64)

Write the current assignment snapshot version to gossip KV.

Source

pub async fn read_snapshot_version(&self) -> Option<u64>

Read the snapshot version from all peers in gossip KV and return the maximum version.

Source

pub async fn start_barrier_server( &self, bind_addr: SocketAddr, advertise_host: Option<String>, ) -> Result<SocketAddr, String>

Start the direct gRPC barrier sync server.

§Errors

Propagates BarrierCoordinator::start_server errors.

Source

pub async fn announce_barrier( &self, ann: &BarrierAnnouncement, ) -> Result<(), String>

Leader-side announce.

§Errors

Propagates BarrierCoordinator::announce errors.

Source

pub async fn observe_barrier( &self, ) -> Result<Option<BarrierAnnouncement>, String>

Follower-side observe; Ok(None) if no leader is visible.

As a side effect, an Aligned or Commit announcement with a populated min_watermark_ms updates the shared cluster-min-watermark atomic so operators on this instance see the cluster-wide minimum without a separate polling path (Aligned carries it so a resuming pipeline sees fresh event-time progress before the upload-gated Commit).

§Errors

Propagates BarrierCoordinator::observe errors.

Source

pub async fn ack_barrier(&self, ack: &BarrierAck) -> Result<(), String>

Follower-side ack.

§Errors

Propagates BarrierCoordinator::ack errors.

Source

pub async fn wait_for_barrier<F>( &self, pred: F, timeout: Duration, ) -> Option<BarrierAnnouncement>

Wait until Self::observe_barrier yields an announcement matching pred, or timeout expires (→ None). Push-driven off the gRPC announcement watch when available; gossip-KV-only deployments (and KV-only announcements) are covered by a fallback poll — 250ms with the watch, 25ms without.

Source

pub async fn wait_for_quorum( &self, epoch: u64, expected: &[NodeId], deadline: Duration, ) -> QuorumOutcome

Leader-side: poll until quorum or deadline.

Source

pub fn snapshot_store(&self) -> Option<&AssignmentSnapshotStore>

Assignment snapshot store, if configured.

Trait Implementations§

Source§

impl Debug for ClusterController

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