1mod declare;
15mod interest;
16mod oam;
17mod push;
18mod request;
19mod response;
20mod timestamp_stack;
21
22use zenoh_buffers::{
23 reader::{BacktrackableReader, DidntRead, Reader},
24 writer::{DidntWrite, Writer},
25};
26use zenoh_protocol::{
27 common::{imsg, ZExtZ64, ZExtZBufHeader},
28 core::{EntityId, Reliability, ZenohIdProto},
29 network::{
30 ext::{self, EntityGlobalIdType},
31 id, NetworkBody, NetworkBodyRef, NetworkMessage, NetworkMessageExt, NetworkMessageRef,
32 },
33};
34
35use crate::{
36 LCodec, RCodec, WCodec, Zenoh080, Zenoh080Header, Zenoh080Length, Zenoh080Reliability,
37};
38
39impl<W> WCodec<NetworkMessageRef<'_>, &mut W> for Zenoh080
41where
42 W: Writer,
43{
44 type Output = Result<(), DidntWrite>;
45
46 #[inline(always)]
47 fn write(self, writer: &mut W, x: NetworkMessageRef) -> Self::Output {
48 let NetworkMessageRef { body, .. } = x;
49 if let NetworkBodyRef::Push(b) = body {
50 return self.write(&mut *writer, b);
51 }
52 #[cold]
53 fn write_not_push<W: Writer>(
54 codec: Zenoh080,
55 writer: &mut W,
56 body: NetworkBodyRef,
57 ) -> Result<(), DidntWrite> {
58 match body {
59 NetworkBodyRef::Push(_) => unreachable!(),
60 NetworkBodyRef::Request(b) => codec.write(&mut *writer, b),
61 NetworkBodyRef::Response(b) => codec.write(&mut *writer, b),
62 NetworkBodyRef::ResponseFinal(b) => codec.write(&mut *writer, b),
63 NetworkBodyRef::Interest(b) => codec.write(&mut *writer, b),
64 NetworkBodyRef::Declare(b) => codec.write(&mut *writer, b),
65 NetworkBodyRef::OAM(b) => codec.write(&mut *writer, b),
66 }
67 }
68 write_not_push(self, writer, body)
69 }
70}
71
72impl<W> WCodec<&NetworkMessage, &mut W> for Zenoh080
73where
74 W: Writer,
75{
76 type Output = Result<(), DidntWrite>;
77
78 fn write(self, writer: &mut W, x: &NetworkMessage) -> Self::Output {
79 self.write(writer, x.as_ref())
80 }
81}
82
83impl<R> RCodec<NetworkMessage, &mut R> for Zenoh080
84where
85 R: Reader,
86{
87 type Error = DidntRead;
88
89 fn read(self, reader: &mut R) -> Result<NetworkMessage, Self::Error> {
90 let codec = Zenoh080Reliability::new(Reliability::DEFAULT);
91 codec.read(reader)
92 }
93}
94
95impl<R> RCodec<NetworkMessage, &mut R> for Zenoh080Reliability
96where
97 R: Reader,
98{
99 type Error = DidntRead;
100
101 #[inline(always)]
102 fn read(self, reader: &mut R) -> Result<NetworkMessage, Self::Error> {
103 let header: u8 = self.codec.read(&mut *reader)?;
104
105 let codec = Zenoh080Header::new(header);
106 let mut msg: NetworkMessage = codec.read(&mut *reader)?;
107 msg.reliability = self.reliability;
108 Ok(msg)
109 }
110}
111
112impl<R> RCodec<NetworkMessage, &mut R> for Zenoh080Header
113where
114 R: Reader,
115{
116 type Error = DidntRead;
117
118 #[inline(always)]
119 fn read(self, reader: &mut R) -> Result<NetworkMessage, Self::Error> {
120 if imsg::mid(self.header) == id::PUSH {
121 return Ok(NetworkBody::Push(self.read(&mut *reader)?).into());
122 }
123 #[cold]
124 fn read_not_push<R: Reader>(
125 header: Zenoh080Header,
126 reader: &mut R,
127 ) -> Result<NetworkMessage, DidntRead> {
128 let body = match imsg::mid(header.header) {
129 id::REQUEST => NetworkBody::Request(header.read(&mut *reader)?),
130 id::RESPONSE => NetworkBody::Response(header.read(&mut *reader)?),
131 id::RESPONSE_FINAL => NetworkBody::ResponseFinal(header.read(&mut *reader)?),
132 id::INTEREST => NetworkBody::Interest(header.read(&mut *reader)?),
133 id::DECLARE => NetworkBody::Declare(header.read(&mut *reader)?),
134 id::OAM => NetworkBody::OAM(header.read(&mut *reader)?),
135 _ => return Err(DidntRead),
136 };
137
138 Ok(body.into())
139 }
140 read_not_push(self, reader)
141 }
142}
143
144#[derive(Debug)]
145pub struct NetworkMessageIter<R> {
146 codec: Zenoh080Reliability,
147 reader: R,
148}
149
150impl<R> NetworkMessageIter<R> {
151 pub fn new(reliability: Reliability, reader: R) -> Self {
152 let codec = Zenoh080Reliability::new(reliability);
153 Self { codec, reader }
154 }
155}
156
157impl<R: BacktrackableReader> Iterator for NetworkMessageIter<R> {
158 type Item = NetworkMessage;
159
160 fn next(&mut self) -> Option<Self::Item> {
161 let mark = self.reader.mark();
162 let msg = self.codec.read(&mut self.reader).ok();
163 if msg.is_none() {
164 self.reader.rewind(mark);
165 }
166 msg
167 }
168}
169
170impl<W, const ID: u8> WCodec<(ext::QoSType<{ ID }>, bool), &mut W> for Zenoh080
172where
173 W: Writer,
174{
175 type Output = Result<(), DidntWrite>;
176
177 fn write(self, writer: &mut W, x: (ext::QoSType<{ ID }>, bool)) -> Self::Output {
178 let (x, more) = x;
179 let ext: ZExtZ64<{ ID }> = x.into();
180 self.write(&mut *writer, (&ext, more))
181 }
182}
183
184impl<R, const ID: u8> RCodec<(ext::QoSType<{ ID }>, bool), &mut R> for Zenoh080
185where
186 R: Reader,
187{
188 type Error = DidntRead;
189
190 fn read(self, reader: &mut R) -> Result<(ext::QoSType<{ ID }>, bool), Self::Error> {
191 let header: u8 = self.read(&mut *reader)?;
192 let codec = Zenoh080Header::new(header);
193 codec.read(reader)
194 }
195}
196
197impl<R, const ID: u8> RCodec<(ext::QoSType<{ ID }>, bool), &mut R> for Zenoh080Header
198where
199 R: Reader,
200{
201 type Error = DidntRead;
202
203 fn read(self, reader: &mut R) -> Result<(ext::QoSType<{ ID }>, bool), Self::Error> {
204 let (ext, more): (ZExtZ64<{ ID }>, bool) = self.read(&mut *reader)?;
205 Ok((ext.into(), more))
206 }
207}
208
209impl<W, const ID: u8> WCodec<(&ext::TimestampType<{ ID }>, bool), &mut W> for Zenoh080
211where
212 W: Writer,
213{
214 type Output = Result<(), DidntWrite>;
215
216 fn write(self, writer: &mut W, x: (&ext::TimestampType<{ ID }>, bool)) -> Self::Output {
217 let (tstamp, more) = x;
218 let header: ZExtZBufHeader<{ ID }> = ZExtZBufHeader::new(self.w_len(&tstamp.timestamp));
219 self.write(&mut *writer, (&header, more))?;
220 self.write(&mut *writer, &tstamp.timestamp)
221 }
222}
223
224impl<R, const ID: u8> RCodec<(ext::TimestampType<{ ID }>, bool), &mut R> for Zenoh080
225where
226 R: Reader,
227{
228 type Error = DidntRead;
229
230 fn read(self, reader: &mut R) -> Result<(ext::TimestampType<{ ID }>, bool), Self::Error> {
231 let header: u8 = self.read(&mut *reader)?;
232 let codec = Zenoh080Header::new(header);
233 codec.read(reader)
234 }
235}
236
237impl<R, const ID: u8> RCodec<(ext::TimestampType<{ ID }>, bool), &mut R> for Zenoh080Header
238where
239 R: Reader,
240{
241 type Error = DidntRead;
242
243 fn read(self, reader: &mut R) -> Result<(ext::TimestampType<{ ID }>, bool), Self::Error> {
244 let (_, more): (ZExtZBufHeader<{ ID }>, bool) = self.read(&mut *reader)?;
245 let timestamp: uhlc::Timestamp = self.codec.read(&mut *reader)?;
246 Ok((ext::TimestampType { timestamp }, more))
247 }
248}
249
250impl<W, const ID: u8> WCodec<(ext::NodeIdType<{ ID }>, bool), &mut W> for Zenoh080
252where
253 W: Writer,
254{
255 type Output = Result<(), DidntWrite>;
256
257 fn write(self, writer: &mut W, x: (ext::NodeIdType<{ ID }>, bool)) -> Self::Output {
258 let (x, more) = x;
259 let ext: ZExtZ64<{ ID }> = x.into();
260 self.write(&mut *writer, (&ext, more))
261 }
262}
263
264impl<R, const ID: u8> RCodec<(ext::NodeIdType<{ ID }>, bool), &mut R> for Zenoh080
265where
266 R: Reader,
267{
268 type Error = DidntRead;
269
270 fn read(self, reader: &mut R) -> Result<(ext::NodeIdType<{ ID }>, bool), Self::Error> {
271 let header: u8 = self.read(&mut *reader)?;
272 let codec = Zenoh080Header::new(header);
273 codec.read(reader)
274 }
275}
276
277impl<R, const ID: u8> RCodec<(ext::NodeIdType<{ ID }>, bool), &mut R> for Zenoh080Header
278where
279 R: Reader,
280{
281 type Error = DidntRead;
282
283 fn read(self, reader: &mut R) -> Result<(ext::NodeIdType<{ ID }>, bool), Self::Error> {
284 let (ext, more): (ZExtZ64<{ ID }>, bool) = self.read(&mut *reader)?;
285 Ok((ext.into(), more))
286 }
287}
288
289impl<const ID: u8> LCodec<&ext::EntityGlobalIdType<{ ID }>> for Zenoh080 {
291 fn w_len(self, x: &ext::EntityGlobalIdType<{ ID }>) -> usize {
292 let EntityGlobalIdType { zid, eid } = x;
293
294 1 + self.w_len(zid) + self.w_len(*eid)
295 }
296}
297
298impl<W, const ID: u8> WCodec<(&ext::EntityGlobalIdType<{ ID }>, bool), &mut W> for Zenoh080
299where
300 W: Writer,
301{
302 type Output = Result<(), DidntWrite>;
303
304 fn write(self, writer: &mut W, x: (&ext::EntityGlobalIdType<{ ID }>, bool)) -> Self::Output {
305 let (x, more) = x;
306 let header: ZExtZBufHeader<{ ID }> = ZExtZBufHeader::new(self.w_len(x));
307 self.write(&mut *writer, (&header, more))?;
308
309 let flags: u8 = (x.zid.size() as u8 - 1) << 4;
310 self.write(&mut *writer, flags)?;
311
312 let lodec = Zenoh080Length::new(x.zid.size());
313 lodec.write(&mut *writer, &x.zid)?;
314
315 self.write(&mut *writer, x.eid)?;
316 Ok(())
317 }
318}
319
320impl<R, const ID: u8> RCodec<(ext::EntityGlobalIdType<{ ID }>, bool), &mut R> for Zenoh080Header
321where
322 R: Reader,
323{
324 type Error = DidntRead;
325
326 fn read(self, reader: &mut R) -> Result<(ext::EntityGlobalIdType<{ ID }>, bool), Self::Error> {
327 let (_, more): (ZExtZBufHeader<{ ID }>, bool) = self.read(&mut *reader)?;
328
329 let flags: u8 = self.codec.read(&mut *reader)?;
330 let length = 1 + ((flags >> 4) as usize);
331
332 let lodec = Zenoh080Length::new(length);
333 let zid: ZenohIdProto = lodec.read(&mut *reader)?;
334
335 let eid: EntityId = self.codec.read(&mut *reader)?;
336
337 Ok((ext::EntityGlobalIdType { zid, eid }, more))
338 }
339}