Skip to main content

BoundedSpawner

Struct BoundedSpawner 

Source
pub struct BoundedSpawner<P: Send + 'static> { /* private fields */ }
Expand description

通用有界并发原语:有界 mpsc(cap=K) + N 个常驻 worker + 溢出策略。

Implementations§

Source§

impl<P: Send + 'static> BoundedSpawner<P>

Source

pub fn new<F, Fut>( concurrency: usize, queue_cap: usize, overflow: Overflow, reply_tx: UnboundedSender<Tick>, run: F, ) -> Self
where F: Fn(Option<Correlation>, P) -> Fut + Clone + Send + 'static, Fut: Future<Output = PortOutcome> + Send + 'static,

构造:concurrency=N 个 worker;queue_cap=K;run = 每条 job 的 I/O 执行体 (async 闭包,返回 PortOutcome,仅在 job.corr=Some 时回报)。

N 个 worker 共享单队列(竞争 recv);Persist 用 N=1 保写序。

Source

pub fn new_observed<F, Fut>( concurrency: usize, queue_cap: usize, overflow: Overflow, reply_tx: UnboundedSender<Tick>, sink: Arc<dyn AsyncMetricSink>, pool: &'static str, run: F, ) -> Self
where F: Fn(Option<Correlation>, P) -> Fut + Clone + Send + 'static, Fut: Future<Output = PortOutcome> + Send + 'static,

构造带指标的工作池,并在公开 seam 上保持与 new 相同的执行语义。

Source

pub async fn submit(&self, job: Job<P>)

入队一条 job。

  • Block:send().await,满则阻塞泵(背压传导回 tick_tx)。
  • DropNewest:try_send,满则丢本条(最新) + warn(泵不阻塞)。
Source

pub async fn shutdown(self)

graceful drain:drop tx → worker recv None 退出 → join 全部。

Source

pub async fn shutdown_with_timeout(self, limit: Duration) -> bool

有上限 graceful drain(fix/lifecycle-net-decouple):drop tx → join,全程封顶 limit。 限内全 join 完成 → true;超时 → abort 残留 worker(放弃在途 job)+ false。

为何需上限:必达 Http reqwest 单条可卡满网络 timeout(连不上 ~30s),无上限会让 helix_destroy graceful drain 等满。不丢保证由乐观态持久化(status=Sending) + 重连重发 对账承担,不由「drain 必等满网络」承担——abort 在途 Http 不丢数据(消息仍 status=Sending 留库,靠重连兜底)。完整语义见 engine.rs 五不变量⑤ + driver AGENTS.md。

Auto Trait Implementations§

§

impl<P> !RefUnwindSafe for BoundedSpawner<P>

§

impl<P> !UnwindSafe for BoundedSpawner<P>

§

impl<P> Freeze for BoundedSpawner<P>
where Sender<QueuedJob<P>>: Freeze,

§

impl<P> Send for BoundedSpawner<P>
where Sender<QueuedJob<P>>: Send,

§

impl<P> Sync for BoundedSpawner<P>
where Sender<QueuedJob<P>>: Sync,

§

impl<P> Unpin for BoundedSpawner<P>
where Sender<QueuedJob<P>>: Unpin,

§

impl<P> UnsafeUnpin for BoundedSpawner<P>
where Sender<QueuedJob<P>>: 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> FutureExt for T

Source§

fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
Source§

fn with_current_context(self) -> WithContext<Self> ⓘ

Attaches the current Context to this type, returning a WithContext wrapper. Read more
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> MaybeSend for T
where T: Send,

Source§

impl<T> MaybeSync for T
where T: Sync,

Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

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

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