ruststream_gcp_pubsub/
message.rs1use 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
11pub const PARTITION_KEY_HEADER: &str = "partition-key";
16
17pub const DELIVERY_ATTEMPT_HEADER: &str = "pubsub-delivery-attempt";
20
21pub 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 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 self.ack().await
108 }
109 }
110
111 fn partition_key(&self) -> Option<&[u8]> {
112 Partitioned::partition_key(self)
113 }
114}
115
116pub(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}