Skip to main content

ruststream_gcp_pubsub/
message.rs

1//! [`PubSubMessage`] and the mapping between `RustStream` headers and Pub/Sub attributes.
2//!
3//! Message attributes carry headers directly - no envelope format is invented - and the
4//! partition key rides the message's ordering key in both directions.
5
6use bytes::Bytes;
7use google_cloud_pubsub::model::Message as GcpMessage;
8use google_cloud_pubsub::subscriber::handler::Handler;
9use ruststream::{AckError, Headers, IncomingMessage, OutgoingMessage, Partitioned};
10
11/// Header carrying the partition key, mapped onto the message's ordering key.
12///
13/// Mirrors the in-memory broker's convention, so services can switch brokers without changing
14/// their headers.
15pub const PARTITION_KEY_HEADER: &str = "partition-key";
16
17/// Header exposing the delivery attempt count on received messages, present when the
18/// subscription has a dead-letter policy.
19pub const DELIVERY_ATTEMPT_HEADER: &str = "pubsub-delivery-attempt";
20
21/// A message delivered by a [`PubSubSubscriber`](crate::PubSubSubscriber).
22///
23/// `ack` and `nack(requeue = true)` are native. `nack(requeue = false)` acknowledges: Pub/Sub
24/// has no "drop without redelivery" beyond acknowledgement - dead-lettering is the
25/// subscription's redrive policy, driven by repeated nacks and expired deadlines, not a
26/// per-message verb. On an exactly-once subscription the confirmed forms are used, so `Ok`
27/// from `ack` means the broker accepted it.
28pub struct PubSubMessage {
29    payload: Bytes,
30    headers: Headers,
31    handler: Handler,
32}
33
34impl std::fmt::Debug for PubSubMessage {
35    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
36        f.debug_struct("PubSubMessage")
37            .field("payload_len", &self.payload.len())
38            .finish_non_exhaustive()
39    }
40}
41
42impl PubSubMessage {
43    pub(crate) fn new(message: GcpMessage, handler: Handler) -> Self {
44        let mut headers = Headers::with_capacity(message.attributes.len() + 2);
45        for (name, value) in &message.attributes {
46            headers.insert(name.clone(), value.clone());
47        }
48        if !message.ordering_key.is_empty() {
49            headers.insert(PARTITION_KEY_HEADER, message.ordering_key.clone());
50        }
51        if let Some(attempt) = handler.delivery_attempt() {
52            headers.insert(DELIVERY_ATTEMPT_HEADER, attempt.to_string());
53        }
54        Self {
55            payload: message.data,
56            headers,
57            handler,
58        }
59    }
60}
61
62impl Partitioned for PubSubMessage {
63    fn partition_key(&self) -> Option<&[u8]> {
64        self.headers.get(PARTITION_KEY_HEADER)
65    }
66}
67
68impl IncomingMessage for PubSubMessage {
69    fn payload(&self) -> &[u8] {
70        &self.payload
71    }
72
73    fn headers(&self) -> &Headers {
74        &self.headers
75    }
76
77    async fn ack(self) -> Result<(), AckError> {
78        match self.handler {
79            // Only the confirmed form guarantees no redelivery on an exactly-once
80            // subscription; the plain form is fire-and-forget.
81            Handler::ExactlyOnce(handler) => handler
82                .confirmed_ack()
83                .await
84                .map_err(|e| AckError::Broker(Box::new(e))),
85            handler => {
86                handler.ack();
87                Ok(())
88            }
89        }
90    }
91
92    async fn nack(self, requeue: bool) -> Result<(), AckError> {
93        if requeue {
94            match self.handler {
95                Handler::ExactlyOnce(handler) => handler
96                    .confirmed_nack()
97                    .await
98                    .map_err(|e| AckError::Broker(Box::new(e))),
99                handler => {
100                    handler.nack();
101                    Ok(())
102                }
103            }
104        } else {
105            // Dropping without redelivery IS an acknowledge in Pub/Sub; the dead-letter
106            // policy on the subscription owns poison-message routing.
107            self.ack().await
108        }
109    }
110
111    fn partition_key(&self) -> Option<&[u8]> {
112        Partitioned::partition_key(self)
113    }
114}
115
116/// Builds the Pub/Sub message for an outgoing publish. Returns the message and its ordering
117/// key (empty when unordered), which the publisher needs for the resume-after-error path.
118pub(crate) fn to_gcp_message(msg: &OutgoingMessage<'_>) -> (GcpMessage, String) {
119    let headers = msg.headers();
120    let mut ordering_key = String::new();
121    let mut attributes: Vec<(String, String)> = Vec::with_capacity(headers.len());
122    for (name, value) in headers.iter() {
123        let text = String::from_utf8_lossy(value).into_owned();
124        if name == PARTITION_KEY_HEADER {
125            ordering_key = text;
126        } else {
127            attributes.push((name.to_owned(), text));
128        }
129    }
130
131    let mut message = GcpMessage::new().set_data(Bytes::copy_from_slice(msg.payload()));
132    if !attributes.is_empty() {
133        message = message.set_attributes(attributes);
134    }
135    if !ordering_key.is_empty() {
136        message = message.set_ordering_key(ordering_key.clone());
137    }
138    (message, ordering_key)
139}
140
141#[cfg(test)]
142mod tests {
143    use super::*;
144
145    #[test]
146    fn partition_key_header_becomes_the_ordering_key() {
147        let mut headers = Headers::new();
148        headers.insert(PARTITION_KEY_HEADER, "user-42");
149        headers.insert("x-tenant", "acme");
150        let outgoing = OutgoingMessage::new("orders", b"{}".as_slice()).with_headers(headers);
151
152        let (message, key) = to_gcp_message(&outgoing);
153        assert_eq!(key, "user-42");
154        assert_eq!(message.ordering_key, "user-42");
155        assert_eq!(
156            message.attributes.get("x-tenant").map(String::as_str),
157            Some("acme")
158        );
159        assert!(!message.attributes.contains_key(PARTITION_KEY_HEADER));
160    }
161
162    #[test]
163    fn plain_messages_carry_no_ordering_key() {
164        let outgoing = OutgoingMessage::new("orders", b"{}".as_slice());
165        let (message, key) = to_gcp_message(&outgoing);
166        assert!(key.is_empty());
167        assert!(message.ordering_key.is_empty());
168    }
169}