Skip to main content

IngressQueue

Struct IngressQueue 

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

Lock-free MPMC ingress queue + waker. Multiple external producers may push concurrently; the engine’s single consumer drains via drain_all on each poll cycle.

Implementations§

Source§

impl IngressQueue

Source

pub fn new() -> Self

Construct a fresh ingress queue with the default capacity (DEFAULT_INGRESS_CAPACITY).

Source

pub fn with_capacity(capacity: usize) -> Self

Construct a fresh ingress queue with the supplied bounded capacity. Per ENGINE.md §2.2 the canonical sizing is bus_capacity * 4; pass the host’s chosen bus_capacity multiplied by 4 to match.

Source

pub fn completion_result_cap(&self) -> usize

Per-complete() result-byte cap. Defaults to usize::MAX when not configured; apply_config_caps reseeds it from NodeConfig::max_completion_result_bytes.

Source

pub fn push(&self, event: IngressEvent) -> Result<(), IngressEvent>

Push an event. On success returns Ok(()) and wakes the engine if it’s sleeping. On a full queue the event comes back in Err(_) and the dropped_overflow counter is incremented; transport adapters decide whether to retry, drop with a metric, or escalate as back-pressure. The IngressEvent Err variant is large (carries a WireEnvelope with multihash PeerIds); transport adapters already box or re-queue, so the cost lives at the boundary.

Source

pub fn drain_all(&self) -> Vec<IngressEvent>

Drain all available events. Called by the engine on each poll cycle’s ingress drain.

Pre-reserves capacity for the bounded queue’s full length so the drain Vec grows once at construction, not in O(log n) reallocations as events pop. The queue itself caps inflight at self.capacity(); the drain is bounded by the same cap, so the upfront reservation is the exact-fit answer.

Source

pub fn register_waker(&self, waker: &Waker)

Register the engine’s waker so future pushes can wake it.

Source

pub fn is_empty(&self) -> bool

true when the queue currently holds no events.

Source

pub fn len(&self) -> usize

Approximate current queue depth. The underlying concurrent-queue returns an approximate len for the MPMC case; introspection callers should treat this as a snapshot, not a real-time invariant.

Source

pub fn capacity(&self) -> usize

Bounded capacity supplied at construction. concurrent-queue guarantees Some(cap) for bounded queues, so unwrapping is safe for the framework’s path that never builds an unbounded ingress queue.

Source

pub fn dropped_overflow(&self) -> u64

Total events dropped due to the queue being full since this queue was constructed. Telemetry hook for transport adapters

  • Node introspection.

Trait Implementations§

Source§

impl CompletionSink for IngressQueue

Source§

fn complete(&self, cmd_id: CommandId, result_bytes: &[u8])

Deliver a successful completion. Implementation copies bytes.
Source§

fn fail(&self, cmd_id: CommandId, detail: &str)

Deliver a failure with Display rendering of the error.
Source§

impl Default for IngressQueue

Source§

fn default() -> Self

Returns the “default value” for a type. 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
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> ErasedComponent for T
where T: Any + Send + Sync,

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