1use super::*;
2use bytes::{Buf, Bytes};
3
4#[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 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 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 buffer.put_u8(0x82);
77
78 let remaining_len = self.len();
80 let remaining_len_bytes = write_remaining_length(buffer, remaining_len)?;
81
82 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 for f in self.filters.iter() {
93 f.write(buffer);
94 }
95
96 Ok(1 + remaining_len_bytes + remaining_len)
97 }
98}
99
100#[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 2 + self.path.len() + 1
122 }
123
124 pub fn read(bytes: &mut Bytes) -> Result<Vec<Filter>, Error> {
125 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 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 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 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}