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>
impl<'a, A: Actor> Queue<'a, A>
Sourcepub async fn send<T: Serialize>(
&self,
name: &str,
body: &T,
) -> Result<CoreQueueMessage>
pub async fn send<T: Serialize>( &self, name: &str, body: &T, ) -> Result<CoreQueueMessage>
Enqueues a message with a CBOR-encoded body.
Sourcepub async fn send_raw(
&self,
name: &str,
body: &[u8],
) -> Result<CoreQueueMessage>
pub async fn send_raw( &self, name: &str, body: &[u8], ) -> Result<CoreQueueMessage>
Enqueues a message with a raw byte body.
Sourcepub async fn enqueue_and_wait<T: Serialize>(
&self,
name: &str,
body: &T,
opts: EnqueueAndWaitOpts,
) -> Result<Option<Vec<u8>>>
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.
Sourcepub async fn enqueue_and_wait_raw(
&self,
name: &str,
body: &[u8],
opts: EnqueueAndWaitOpts,
) -> Result<Option<Vec<u8>>>
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.
Sourcepub async fn next(
&self,
opts: QueueNextOpts,
) -> Result<Option<CoreQueueMessage>>
pub async fn next( &self, opts: QueueNextOpts, ) -> Result<Option<CoreQueueMessage>>
Awaits the next queued message, optionally bounded by the opts timeout.
Sourcepub async fn next_typed<M: QueueMessage>(
&self,
opts: QueueNextOpts,
) -> Result<Option<TypedQueueMessage<M>>>
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.
Sourcepub async fn next_batch(
&self,
opts: QueueNextBatchOpts,
) -> Result<Vec<CoreQueueMessage>>
pub async fn next_batch( &self, opts: QueueNextBatchOpts, ) -> Result<Vec<CoreQueueMessage>>
Awaits up to opts.count queued messages.
Sourcepub async fn next_batch_typed<M: QueueMessage>(
&self,
opts: QueueNextBatchOpts,
) -> Result<Vec<TypedQueueMessage<M>>>
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.
Sourcepub fn try_next(
&self,
opts: QueueTryNextOpts,
) -> Result<Option<CoreQueueMessage>>
pub fn try_next( &self, opts: QueueTryNextOpts, ) -> Result<Option<CoreQueueMessage>>
Returns the next queued message if one is immediately available.
Sourcepub fn try_next_typed<M: QueueMessage>(
&self,
opts: QueueTryNextOpts,
) -> Result<Option<TypedQueueMessage<M>>>
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.
Sourcepub fn try_next_batch(
&self,
opts: QueueTryNextBatchOpts,
) -> Result<Vec<CoreQueueMessage>>
pub fn try_next_batch( &self, opts: QueueTryNextBatchOpts, ) -> Result<Vec<CoreQueueMessage>>
Returns immediately-available queued messages up to opts.count.
Sourcepub fn try_next_batch_typed<M: QueueMessage>(
&self,
opts: QueueTryNextBatchOpts,
) -> Result<Vec<TypedQueueMessage<M>>>
pub fn try_next_batch_typed<M: QueueMessage>( &self, opts: QueueTryNextBatchOpts, ) -> Result<Vec<TypedQueueMessage<M>>>
Returns immediately-available queued messages matching M::NAME.
Sourcepub async fn inspect_messages(&self) -> Result<Vec<CoreQueueMessage>>
pub async fn inspect_messages(&self) -> Result<Vec<CoreQueueMessage>>
Lists the currently persisted queue messages without consuming them.