rumqttc/v5/mqttbytes/v5/
subscribe.rs

1use super::*;
2use bytes::{Buf, Bytes};
3
4/// Subscription packet
5#[derive(Clone, Debug, PartialEq, Eq, Default)]
6pub struct Subscribe {
7    pub pkid: u16,
8    pub filters: Vec<Filter>,
9    pub properties: Option<SubscribeProperties>,
10}
11
12impl Subscribe {
13    pub fn new(filter: Filter, properties: Option<SubscribeProperties>) -> Self {
14        Self {
15            filters: vec![filter],
16            properties,
17            ..Default::default()
18        }
19    }
20
21    pub fn new_many<F>(filters: F, properties: Option<SubscribeProperties>) -> Self
22    where
23        F: IntoIterator<Item = Filter>,
24    {
25        Self {
26            filters: filters.into_iter().collect(),
27            properties,
28            ..Default::default()
29        }
30    }
31
32    pub fn size(&self) -> usize {
33        let len = self.len();
34        let remaining_len_size = len_len(len);
35
36        1 + remaining_len_size + len
37    }
38
39    fn len(&self) -> usize {
40        let mut len = 2 + self.filters.iter().fold(0, |s, t| s + t.len());
41
42        if let Some(p) = &self.properties {
43            let properties_len = p.len();
44            let properties_len_len = len_len(properties_len);
45            len += properties_len_len + properties_len;
46        } else {
47            // just 1 byte representing 0 len
48            len += 1;
49        }
50
51        len
52    }
53
54    pub fn read(fixed_header: FixedHeader, mut bytes: Bytes) -> Result<Subscribe, Error> {
55        let variable_header_index = fixed_header.fixed_header_len;
56        bytes.advance(variable_header_index);
57
58        let pkid = read_u16(&mut bytes)?;
59        let properties = SubscribeProperties::read(&mut bytes)?;
60
61        // variable header size = 2 (packet identifier)
62        let filters = Filter::read(&mut bytes)?;
63
64        match filters.len() {
65            0 => Err(Error::EmptySubscription),
66            _ => Ok(Subscribe {
67                pkid,
68                filters,
69                properties,
70            }),
71        }
72    }
73
74    pub fn write(&self, buffer: &mut BytesMut) -> Result<usize, Error> {
75        // write packet type
76        buffer.put_u8(0x82);
77
78        // write remaining length
79        let remaining_len = self.len();
80        let remaining_len_bytes = write_remaining_length(buffer, remaining_len)?;
81
82        // write packet id
83        buffer.put_u16(self.pkid);
84
85        if let Some(p) = &self.properties {
86            p.write(buffer)?;
87        } else {
88            write_remaining_length(buffer, 0)?;
89        }
90
91        // write filters
92        for f in self.filters.iter() {
93            f.write(buffer);
94        }
95
96        Ok(1 + remaining_len_bytes + remaining_len)
97    }
98}
99
100///  Subscription filter
101#[derive(Clone, Debug, PartialEq, Eq, Default)]
102pub struct Filter {
103    pub path: String,
104    pub qos: QoS,
105    pub nolocal: bool,
106    pub preserve_retain: bool,
107    pub retain_forward_rule: RetainForwardRule,
108}
109
110impl Filter {
111    pub fn new<T: Into<String>>(topic: T, qos: QoS) -> Self {
112        Self {
113            path: topic.into(),
114            qos,
115            ..Default::default()
116        }
117    }
118
119    fn len(&self) -> usize {
120        // filter len + filter + options
121        2 + self.path.len() + 1
122    }
123
124    pub fn read(bytes: &mut Bytes) -> Result<Vec<Filter>, Error> {
125        // variable header size = 2 (packet identifier)
126        let mut filters = Vec::new();
127
128        while bytes.has_remaining() {
129            let path = read_mqtt_string(bytes)?;
130            let options = read_u8(bytes)?;
131            let requested_qos = options & 0b0000_0011;
132
133            let nolocal = options >> 2 & 0b0000_0001;
134            let nolocal = nolocal != 0;
135
136            let preserve_retain = options >> 3 & 0b0000_0001;
137            let preserve_retain = preserve_retain != 0;
138
139            let retain_forward_rule = (options >> 4) & 0b0000_0011;
140            let retain_forward_rule = match retain_forward_rule {
141                0 => RetainForwardRule::OnEverySubscribe,
142                1 => RetainForwardRule::OnNewSubscribe,
143                2 => RetainForwardRule::Never,
144                r => return Err(Error::InvalidRetainForwardRule(r)),
145            };
146
147            filters.push(Filter {
148                path,
149                qos: qos(requested_qos).ok_or(Error::InvalidQoS(requested_qos))?,
150                nolocal,
151                preserve_retain,
152                retain_forward_rule,
153            });
154        }
155
156        Ok(filters)
157    }
158
159    pub fn write(&self, buffer: &mut BytesMut) {
160        let mut options = 0;
161        options |= self.qos as u8;
162
163        if self.nolocal {
164            options |= 0b0000_0100;
165        }
166
167        if self.preserve_retain {
168            options |= 0b0000_1000;
169        }
170
171        options |= match self.retain_forward_rule {
172            RetainForwardRule::OnEverySubscribe => 0b0000_0000,
173            RetainForwardRule::OnNewSubscribe => 0b0001_0000,
174            RetainForwardRule::Never => 0b0010_0000,
175        };
176
177        write_mqtt_string(buffer, self.path.as_str());
178        buffer.put_u8(options);
179    }
180}
181
182#[derive(Debug, Clone, PartialEq, Eq)]
183pub enum RetainForwardRule {
184    OnEverySubscribe,
185    OnNewSubscribe,
186    Never,
187}
188
189impl Default for RetainForwardRule {
190    fn default() -> Self {
191        Self::OnEverySubscribe
192    }
193}
194
195#[derive(Debug, Clone, PartialEq, Eq)]
196pub struct SubscribeProperties {
197    pub id: Option<usize>,
198    pub user_properties: Vec<(String, String)>,
199}
200
201impl SubscribeProperties {
202    fn len(&self) -> usize {
203        let mut len = 0;
204
205        if let Some(id) = &self.id {
206            len += 1 + len_len(*id);
207        }
208
209        for (key, value) in self.user_properties.iter() {
210            len += 1 + 2 + key.len() + 2 + value.len();
211        }
212
213        len
214    }
215
216    pub fn read(bytes: &mut Bytes) -> Result<Option<SubscribeProperties>, Error> {
217        let mut id = None;
218        let mut user_properties = Vec::new();
219
220        let (properties_len_len, properties_len) = length(bytes.iter())?;
221        bytes.advance(properties_len_len);
222
223        if properties_len == 0 {
224            return Ok(None);
225        }
226
227        let mut cursor = 0;
228        // read until cursor reaches property length. properties_len = 0 will skip this loop
229        while cursor < properties_len {
230            let prop = read_u8(bytes)?;
231            cursor += 1;
232
233            match property(prop)? {
234                PropertyType::SubscriptionIdentifier => {
235                    let (id_len, sub_id) = length(bytes.iter())?;
236                    // TODO: Validate 1 +. Tests are working either way
237                    cursor += 1 + id_len;
238                    bytes.advance(id_len);
239                    id = Some(sub_id)
240                }
241                PropertyType::UserProperty => {
242                    let key = read_mqtt_string(bytes)?;
243                    let value = read_mqtt_string(bytes)?;
244                    cursor += 2 + key.len() + 2 + value.len();
245                    user_properties.push((key, value));
246                }
247                _ => return Err(Error::InvalidPropertyType(prop)),
248            }
249        }
250
251        Ok(Some(SubscribeProperties {
252            id,
253            user_properties,
254        }))
255    }
256
257    pub fn write(&self, buffer: &mut BytesMut) -> Result<(), Error> {
258        let len = self.len();
259        write_remaining_length(buffer, len)?;
260
261        if let Some(id) = &self.id {
262            buffer.put_u8(PropertyType::SubscriptionIdentifier as u8);
263            write_remaining_length(buffer, *id)?;
264        }
265
266        for (key, value) in self.user_properties.iter() {
267            buffer.put_u8(PropertyType::UserProperty as u8);
268            write_mqtt_string(buffer, key);
269            write_mqtt_string(buffer, value);
270        }
271
272        Ok(())
273    }
274}
275
276#[cfg(test)]
277mod test {
278    use super::super::test::{USER_PROP_KEY, USER_PROP_VAL};
279    use super::*;
280    use bytes::BytesMut;
281    use pretty_assertions::assert_eq;
282
283    #[test]
284    fn length_calculation() {
285        let mut dummy_bytes = BytesMut::new();
286        // Use user_properties to pad the size to exceed ~128 bytes to make the
287        // remaining_length field in the packet be 2 bytes long.
288        let subscribe_props = SubscribeProperties {
289            id: None,
290            user_properties: vec![(USER_PROP_KEY.into(), USER_PROP_VAL.into())],
291        };
292
293        let subscribe_pkt = Subscribe::new(
294            Filter::new("hello/world", QoS::AtMostOnce),
295            Some(subscribe_props),
296        );
297
298        let size_from_size = subscribe_pkt.size();
299        let size_from_write = subscribe_pkt.write(&mut dummy_bytes).unwrap();
300        let size_from_bytes = dummy_bytes.len();
301
302        assert_eq!(size_from_write, size_from_bytes);
303        assert_eq!(size_from_size, size_from_bytes);
304    }
305}