Skip to main content

Pool

Struct Pool 

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

A multi-station connection pool — see this module’s own doc for the full design.

links is a plain append-only Vec<Arc<PooledLink>> behind an async RwLock, deliberately NOT a HashMap/HashSet — this is the one change made in direct response to the SAME class of bug found (twice, independently) porting this feature to macula-go (map[string]*link randomizing iteration order) and macula-dotnet (migrating to ConcurrentDictionary for concurrent-add safety silently broke FirstSuccess’s reliance on insertion order). A Vec, appended to under the same lock that’s ALSO taken to read it, has no separate enumeration-order concept to accidentally break — insertion order IS the order, by construction, with nothing to track alongside it the way go’s fix (tracking iteration order separately) or dotnet’s fix (an explicit Ordinal field) both had to.

Implementations§

Source§

impl Pool

Source

pub fn connect( seeds: Vec<Seed>, trust: Trust, identity: KeyPair, options: PoolOptions, ) -> Arc<Pool> ⓘ

Spawn a pool with one link per seed. Returns as soon as every link’s dial has STARTED, not once any is connected — handshakes complete asynchronously, matching macula_client:connect/2 and every other port of this pool shape.

Source

pub async fn call( self: &Arc<Self>, procedure: &str, realm: [u8; 32], payload: Value, deadline_ms: i128, ) -> Result<CallResponse, PoolCallError>

Send a signed CALL, choosing among currently-connected links per PoolOptions::link_selection, trying each in order until one answers (a transport-level failure marks that link disconnected and triggers its respawn, then moves to the next candidate — a BOLT#4 ERROR response is still a successful call as far as this pool is concerned, exactly like a bare Session::call).

Source

pub async fn publish( self: &Arc<Self>, spec: &PublishSpec, ) -> Result<(), PoolPublishError>

Send a signed PUBLISH, fanning out to up to PoolOptions::replication_factor currently-connected links (ordered by PoolOptions::link_selection). Partial success counts as success, matching macula-ts/macula-dotnet’s own publish-fanout contract.

Source

pub async fn status(&self) -> PoolStatus

Aggregate health snapshot.

Per-link snapshot, in seed-list/discovery order — see Pool’s own doc on why a plain Vec already guarantees this without any extra bookkeeping.

Source

pub async fn close(&self, reason: &str, detail: Option<&str>)

Sends GOODBYE on every currently-connected link and stops all background dial/discovery tasks. Waits for every background task (respawn lifecycles, the discovery loop) to actually be gone BEFORE draining/closing links — see Pool’s own field doc on tasks for why this ordering, specifically, is load-bearing. Does not wait for the GOODBYE writes themselves to finish being scheduled beyond Session::close’s own bounded drain.

Auto Trait Implementations§

§

impl !Freeze for Pool

§

impl !RefUnwindSafe for Pool

§

impl !UnwindSafe for Pool

§

impl Send for Pool

§

impl Sync for Pool

§

impl Unpin for Pool

§

impl UnsafeUnpin for Pool

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<'a, T, E> AsTaggedExplicit<'a, E> for T
where T: 'a,

Source§

fn explicit(self, class: Class, tag: u32) -> TaggedParser<'a, Explicit, Self, E>

Source§

impl<'a, T, E> AsTaggedImplicit<'a, E> for T
where T: 'a,

Source§

fn implicit( self, class: Class, constructed: bool, tag: u32, ) -> TaggedParser<'a, Implicit, Self, E>

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.

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