Skip to main content

zenoh_codec/network/
mod.rs

1//
2// Copyright (c) 2022 ZettaScale Technology
3//
4// This program and the accompanying materials are made available under the
5// terms of the Eclipse Public License 2.0 which is available at
6// http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0
7// which is available at https://www.apache.org/licenses/LICENSE-2.0.
8//
9// SPDX-License-Identifier: EPL-2.0 OR Apache-2.0
10//
11// Contributors:
12//   ZettaScale Zenoh Team, <zenoh@zettascale.tech>
13//
14mod 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
39// NetworkMessage
40impl<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
170// Extensions: QoS
171impl<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
209// Extensions: Timestamp
210impl<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
250// Extensions: NodeId
251impl<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
289// Extension: EntityId
290impl<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}