Skip to main content

ruststream_zeromq/
message.rs

1//! [`ZmqMessage`]: a delivered message.
2
3use bytes::Bytes;
4use ruststream::{AckError, Headers, IncomingMessage};
5
6/// A message delivered by one of the transport's subscribers.
7///
8/// Delivery is at most once and there is no durability, so acknowledgement is reported as
9/// [`AckError::Unsupported`] rather than emulated.
10pub struct ZmqMessage {
11    pub(crate) name: String,
12    pub(crate) headers: Headers,
13    pub(crate) payload: Bytes,
14}
15
16impl std::fmt::Debug for ZmqMessage {
17    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
18        f.debug_struct("ZmqMessage")
19            .field("name", &self.name)
20            .field("payload_len", &self.payload.len())
21            .finish_non_exhaustive()
22    }
23}
24
25impl ZmqMessage {
26    /// The name frame this message carried.
27    #[must_use]
28    pub fn name(&self) -> &str {
29        &self.name
30    }
31}
32
33impl IncomingMessage for ZmqMessage {
34    fn payload(&self) -> &[u8] {
35        &self.payload
36    }
37
38    fn headers(&self) -> &Headers {
39        &self.headers
40    }
41
42    async fn ack(self) -> Result<(), AckError> {
43        Err(AckError::Unsupported)
44    }
45
46    async fn nack(self, _requeue: bool) -> Result<(), AckError> {
47        Err(AckError::Unsupported)
48    }
49}