Skip to main content

MasterClientPool

Struct MasterClientPool 

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

A pool of MasterClients over independent HTTP/2 channels.

§Why

A single tonic Channel multiplexes all RPCs over one HTTP/2 connection, which caps concurrency at SETTINGS_MAX_CONCURRENT_STREAMS (default 100). Under 256-way concurrency over remote RTT the surplus requests queue in tower::Buffer, which is the measured root cause of the remote GetFileStatus / OpenFile regression vs Java (Java defaults to a channel pool). Spreading requests across master_connection_pool_size channels removes the queue.

§Scheduling

Two strategies are available via GoosefsConfig::master_connection_pool_schedule:

  • RoundRobin (default): cycles through channels in order. Wait-free, zero overhead, no in-flight tracking required.
  • P2C: Power of Two Choices — uniformly samples two distinct channels with a fast PRNG (fastrand) and picks the one with fewer in-flight RPCs. Per-channel in-flight counts are tracked inside each MasterClient (incremented in with_retry, decremented on RPC completion), so the count is accurate even for MasterClients cloned out of the pool (e.g. by GoosefsFileWriter).

§HA consistency

Every pooled client is constructed with the same inquire_client, so a failover decision is shared: all channels re-discover and switch to the same new Primary, eliminating split-brain. Each channel performs its own SASL handshake and carries a unique channel-id, fully compatible with the ArcSwap<AuthedState> model.

Implementations§

Source§

impl MasterClientPool

Source

pub async fn connect_with_inquire( config: &GoosefsConfig, inquire_client: Arc<dyn MasterInquireClient>, ) -> Result<Self>

Connect a pool of config.master_connection_pool_size master clients, all sharing the supplied inquire_client.

The size is clamped to at least 1, so this is a strict superset of the previous single-channel behaviour (size = 1).

Source

pub fn pick(&self) -> Arc<MasterClient>

Pick the next client according to the configured scheduling strategy.

  • RoundRobin (default): cycle through channels in order. Wait-free, zero overhead, no in-flight tracking required.
  • P2C: Power of Two Choices — uniformly sample two distinct channels at random and select the one with fewer in-flight RPCs. The per-channel in-flight count is maintained inside MasterClient::with_retry (not by this method), so it stays accurate even for clients cloned out of the pool.

Returns Arc<MasterClient> — callers interact with it exactly as before. The in-flight counter lives inside MasterClient itself (shared via Arc<AtomicUsize> across clones), so P2C load awareness works regardless of whether the caller holds the Arc directly or clones the inner MasterClient.

Source

pub fn size(&self) -> usize

Number of pooled channels.

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
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> 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> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

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