Skip to main content

VnodeRegistry

Struct VnodeRegistry 

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

Runtime registry of vnode topology and assignment.

Implementations§

Source§

impl VnodeRegistry

Source

pub fn new(vnode_count: u32) -> Self

Create a registry sized for vnode_count vnodes, all marked as NodeId::UNASSIGNED. The assignment version starts at 1.

§Panics

Panics if vnode_count == 0.

Source

pub fn new_unassigned(vnode_count: u32) -> Self

Create a registry with every vnode unassigned at version 0, so the first stored assignment snapshot (version >= 1) adopts through the standard rotation path.

§Panics

Panics if vnode_count == 0.

Source

pub fn single_owner(vnode_count: u32, owner: NodeId) -> Self

Create a registry where every vnode is owned by the same node.

Used by single-instance / embedded deployments.

§Panics

Panics if vnode_count == 0.

Source

pub fn vnode_count(&self) -> u32

Number of vnodes.

Source

pub fn assignment_version(&self) -> u64

Current monotonic assignment version.

Source

pub fn owner(&self, vnode: u32) -> NodeId

Owner of a given vnode. Returns NodeId::UNASSIGNED if the vnode is out of range or unassigned.

Source

pub fn snapshot(&self) -> Arc<[NodeId]>

Snapshot the current assignment vector. Cheap — internally an Arc::clone.

Source

pub fn versioned_snapshot(&self) -> VnodeAssignmentSnapshot

Consistent ownership, version, and source-handoff publication.

Source

pub fn read_assignment(&self) -> VnodeAssignmentReadGuard<'_>

Pin the current publication for a short, non-blocking source capture.

Source

pub fn set_assignment(&self, new_assignment: Arc<[NodeId]>)

Replace the full assignment and bump the version.

§Panics

Panics if new_assignment.len() != self.vnode_count.

Source

pub fn set_assignment_and_version( &self, new_assignment: Arc<[NodeId]>, version: u64, )

Replace the full assignment and set the version to version atomically. For recovery paths that must restore the registry to a persisted fence generation, not a fresh bump.

§Panics

Panics on length mismatch, or if version is less than the current one (assignment versions are monotonic).

Source

pub fn set_assignment_and_version_with_source_handoff( &self, new_assignment: Arc<[NodeId]>, version: u64, source_handoff: Arc<CommittedSourceHandoff>, )

Atomically publish ownership, its monotonic version, and the sealed connector cursors used by sources acquiring partitions in that version.

§Panics

Panics on length mismatch or version regression.

Source

pub fn set_assignment_and_version_carrying_source_handoff( &self, new_assignment: Arc<[NodeId]>, version: u64, )

Publish a version that does not acquire local ownership while retaining the prior version-bound handoff until the source has reconciled it.

Source

pub fn vnode_for_key(&self, key: &[u8]) -> u32

Map a primary key to a vnode.

Source

pub fn mark_restoring(&self, vnodes: &[u32])

Mark vnodes as Restoring.

Called during a rebalance for the vnodes a node newly acquires, before their committed state has been applied. Out-of-range ids are ignored.

Source

pub fn mark_active(&self, vnodes: &[u32])

Mark vnodes as Active.

Called once a newly-acquired vnode’s state has been applied (or immediately for vnodes that had no durable state to restore). Out-of-range ids are ignored.

Source

pub fn is_restoring(&self, vnode: u32) -> bool

Whether vnode is currently Restoring. Out-of-range ids are reported as not restoring.

Source

pub fn any_restoring(&self) -> bool

Whether any vnode is currently restoring. Cheap pre-check the emission gate uses to skip per-row work in the common case.

Source

pub fn restoring_vnodes(&self) -> Vec<u32>

Vnodes currently Restoring, ascending.

Trait Implementations§

Source§

impl Debug for VnodeRegistry

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