pub struct StreamBufferConsumer { /* private fields */ }Expand description
Consumer side of the stream_buffer endpoint.
A consumer is bound to exactly one (topic, correlation_id) partition. Use
this for reading the responses for one HTTP request, LLM generation, MCP
call, or similar one-request/many-response workflow.
Implementations§
Source§impl StreamBufferConsumer
impl StreamBufferConsumer
pub fn new(config: &StreamBufferConfig) -> Result<Self>
Trait Implementations§
Source§impl Debug for StreamBufferConsumer
impl Debug for StreamBufferConsumer
Source§impl Drop for StreamBufferConsumer
impl Drop for StreamBufferConsumer
Source§impl MessageConsumer for StreamBufferConsumer
impl MessageConsumer for StreamBufferConsumer
Source§fn receive_batch<'life0, 'async_trait>(
&'life0 mut self,
max_messages: usize,
) -> Pin<Box<dyn Future<Output = Result<ReceivedBatch, ConsumerError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn receive_batch<'life0, 'async_trait>(
&'life0 mut self,
max_messages: usize,
) -> Pin<Box<dyn Future<Output = Result<ReceivedBatch, ConsumerError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Receives a batch of messages. Read more
fn status<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = EndpointStatus> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn as_any(&self) -> &dyn Any
Source§fn on_connect_hook(&self) -> Option<BoxFuture<'_, Result<()>>>
fn on_connect_hook(&self) -> Option<BoxFuture<'_, Result<()>>>
Returns an optional lifecycle hook that runs once after the consumer connection is created. Read more
Source§fn on_disconnect_hook(&self) -> Option<BoxFuture<'_, Result<()>>>
fn on_disconnect_hook(&self) -> Option<BoxFuture<'_, Result<()>>>
Returns an optional lifecycle hook that runs before the consumer is dropped. Read more
Source§fn receive<'life0, 'async_trait>(
&'life0 mut self,
) -> Pin<Box<dyn Future<Output = Result<Received, ConsumerError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn receive<'life0, 'async_trait>(
&'life0 mut self,
) -> Pin<Box<dyn Future<Output = Result<Received, ConsumerError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Receives a single message.
fn receive_batch_helper<'life0, 'async_trait>(
&'life0 mut self,
_max_messages: usize,
) -> Pin<Box<dyn Future<Output = Result<ReceivedBatch, ConsumerError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Auto Trait Implementations§
impl !Unpin for StreamBufferConsumer
impl !UnsafeUnpin for StreamBufferConsumer
impl Freeze for StreamBufferConsumer
impl RefUnwindSafe for StreamBufferConsumer
impl Send for StreamBufferConsumer
impl Sync for StreamBufferConsumer
impl UnwindSafe for StreamBufferConsumer
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more