Skip to main content

Queue

Struct Queue 

Source
pub struct Queue<'a, A: Actor> { /* private fields */ }
Expand description

Typed handle over the actor message queue, returned by crate::Ctx::queue.

This is a thin typed facade over the core queue API. Send helpers CBOR-encode the message body; the *_raw variants pass bytes through unchanged. Received QueueMessage bodies are raw bytes that the caller decodes against their own schema, matching how queued messages arrive through the event loop as RuntimeEvent::QueueSend.

Implementations§

Source§

impl<'a, A: Actor> Queue<'a, A>

Source

pub async fn send<T: Serialize>( &self, name: &str, body: &T, ) -> Result<CoreQueueMessage>

Enqueues a message with a CBOR-encoded body.

Source

pub async fn send_raw( &self, name: &str, body: &[u8], ) -> Result<CoreQueueMessage>

Enqueues a message with a raw byte body.

Source

pub async fn enqueue_and_wait<T: Serialize>( &self, name: &str, body: &T, opts: EnqueueAndWaitOpts, ) -> Result<Option<Vec<u8>>>

Enqueues a message with a CBOR-encoded body and waits for the consumer to complete it, returning the raw completion response if any.

Source

pub async fn enqueue_and_wait_raw( &self, name: &str, body: &[u8], opts: EnqueueAndWaitOpts, ) -> Result<Option<Vec<u8>>>

Enqueues a raw-body message and waits for its completion response.

Source

pub async fn next( &self, opts: QueueNextOpts, ) -> Result<Option<CoreQueueMessage>>

Awaits the next queued message, optionally bounded by the opts timeout.

Source

pub async fn next_typed<M: QueueMessage>( &self, opts: QueueNextOpts, ) -> Result<Option<TypedQueueMessage<M>>>

Awaits the next queued message matching M::NAME and decodes its CBOR body.

Source

pub async fn next_batch( &self, opts: QueueNextBatchOpts, ) -> Result<Vec<CoreQueueMessage>>

Awaits up to opts.count queued messages.

Source

pub async fn next_batch_typed<M: QueueMessage>( &self, opts: QueueNextBatchOpts, ) -> Result<Vec<TypedQueueMessage<M>>>

Awaits up to opts.count queued messages matching M::NAME.

Source

pub fn try_next( &self, opts: QueueTryNextOpts, ) -> Result<Option<CoreQueueMessage>>

Returns the next queued message if one is immediately available.

Source

pub fn try_next_typed<M: QueueMessage>( &self, opts: QueueTryNextOpts, ) -> Result<Option<TypedQueueMessage<M>>>

Returns the next queued message matching M::NAME if one is immediately available.

Source

pub fn try_next_batch( &self, opts: QueueTryNextBatchOpts, ) -> Result<Vec<CoreQueueMessage>>

Returns immediately-available queued messages up to opts.count.

Source

pub fn try_next_batch_typed<M: QueueMessage>( &self, opts: QueueTryNextBatchOpts, ) -> Result<Vec<TypedQueueMessage<M>>>

Returns immediately-available queued messages matching M::NAME.

Source

pub async fn inspect_messages(&self) -> Result<Vec<CoreQueueMessage>>

Lists the currently persisted queue messages without consuming them.

Source

pub fn max_size(&self) -> u32

Returns the configured maximum queue size.

Auto Trait Implementations§

§

impl<'a, A> !RefUnwindSafe for Queue<'a, A>

§

impl<'a, A> !UnwindSafe for Queue<'a, A>

§

impl<'a, A> Freeze for Queue<'a, A>

§

impl<'a, A> Send for Queue<'a, A>

§

impl<'a, A> Sync for Queue<'a, A>

§

impl<'a, A> Unpin for Queue<'a, A>

§

impl<'a, A> UnsafeUnpin for Queue<'a, A>

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<T> ErasedDestructor for T
where T: 'static,

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> 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> RuntimeFutureOutput for T
where T: Send + 'static,

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