Skip to main content

zenoh_codec/network/
response.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//
14use zenoh_buffers::{
15    reader::{DidntRead, Reader},
16    writer::{DidntWrite, Writer},
17};
18use zenoh_protocol::{
19    common::{iext, imsg},
20    core::WireExpr,
21    network::{
22        id,
23        response::{ext, flag},
24        Mapping, RequestId, Response, ResponseFinal,
25    },
26    zenoh::ResponseBody,
27};
28
29use crate::{
30    common::extension, RCodec, WCodec, Zenoh080, Zenoh080Bounded, Zenoh080Condition, Zenoh080Header,
31};
32
33// Response
34impl<W> WCodec<&Response, &mut W> for Zenoh080
35where
36    W: Writer,
37{
38    type Output = Result<(), DidntWrite>;
39
40    fn write(self, writer: &mut W, x: &Response) -> Self::Output {
41        let Response {
42            rid,
43            wire_expr,
44            payload,
45            ext_qos,
46            ext_tstamp,
47            ext_respid,
48            ext_ts_stack,
49        } = x;
50
51        // Header
52        let mut header = id::RESPONSE;
53        let mut n_exts = ((ext_qos != &ext::QoSType::DEFAULT) as u8)
54            + (ext_tstamp.is_some() as u8)
55            + (ext_respid.is_some() as u8)
56            + (ext_ts_stack.is_some() as u8);
57        if n_exts != 0 {
58            header |= flag::Z;
59        }
60        if wire_expr.mapping != Mapping::DEFAULT {
61            header |= flag::M;
62        }
63        if wire_expr.has_suffix() {
64            header |= flag::N;
65        }
66        self.write(&mut *writer, header)?;
67
68        // Body
69        self.write(&mut *writer, rid)?;
70        self.write(&mut *writer, wire_expr)?;
71
72        // Extensions
73        if ext_qos != &ext::QoSType::DEFAULT {
74            n_exts -= 1;
75            self.write(&mut *writer, (*ext_qos, n_exts != 0))?;
76        }
77        if let Some(ts) = ext_tstamp.as_ref() {
78            n_exts -= 1;
79            self.write(&mut *writer, (ts, n_exts != 0))?;
80        }
81        if let Some(ri) = ext_respid.as_ref() {
82            n_exts -= 1;
83            self.write(&mut *writer, (ri, n_exts != 0))?;
84        }
85        if let Some(ts_stack) = ext_ts_stack.as_ref() {
86            n_exts -= 1;
87            self.write(&mut *writer, (ts_stack, n_exts != 0))?;
88        }
89
90        // Payload
91        self.write(&mut *writer, payload)?;
92
93        Ok(())
94    }
95}
96
97impl<R> RCodec<Response, &mut R> for Zenoh080
98where
99    R: Reader,
100{
101    type Error = DidntRead;
102
103    fn read(self, reader: &mut R) -> Result<Response, Self::Error> {
104        let header: u8 = self.read(&mut *reader)?;
105        let codec = Zenoh080Header::new(header);
106        codec.read(reader)
107    }
108}
109
110impl<R> RCodec<Response, &mut R> for Zenoh080Header
111where
112    R: Reader,
113{
114    type Error = DidntRead;
115
116    fn read(self, reader: &mut R) -> Result<Response, Self::Error> {
117        if imsg::mid(self.header) != id::RESPONSE {
118            return Err(DidntRead);
119        }
120
121        // Body
122        let bodec = Zenoh080Bounded::<RequestId>::new();
123        let rid: RequestId = bodec.read(&mut *reader)?;
124        let ccond = Zenoh080Condition::new(imsg::has_flag(self.header, flag::N));
125        let mut wire_expr: WireExpr<'static> = ccond.read(&mut *reader)?;
126        wire_expr.mapping = if imsg::has_flag(self.header, flag::M) {
127            Mapping::Sender
128        } else {
129            Mapping::Receiver
130        };
131
132        // Extensions
133        let mut ext_qos = ext::QoSType::DEFAULT;
134        let mut ext_tstamp = None;
135        let mut ext_respid = None;
136        let mut ext_ts_stack = None;
137
138        let mut has_ext = imsg::has_flag(self.header, flag::Z);
139        while has_ext {
140            let ext: u8 = self.codec.read(&mut *reader)?;
141            let eodec = Zenoh080Header::new(ext);
142            match iext::eid(ext) {
143                ext::QoS::ID => {
144                    let (q, ext): (ext::QoSType, bool) = eodec.read(&mut *reader)?;
145                    ext_qos = q;
146                    has_ext = ext;
147                }
148                ext::Timestamp::ID => {
149                    let (t, ext): (ext::TimestampType, bool) = eodec.read(&mut *reader)?;
150                    ext_tstamp = Some(t);
151                    has_ext = ext;
152                }
153                ext::ResponderId::ID => {
154                    let (t, ext): (ext::ResponderIdType, bool) = eodec.read(&mut *reader)?;
155                    ext_respid = Some(t);
156                    has_ext = ext;
157                }
158                ext::TsStack::ID => {
159                    let (ts, ext): (ext::TsStackType, bool) = eodec.read(&mut *reader)?;
160                    ext_ts_stack = Some(ts);
161                    has_ext = ext;
162                }
163                _ => {
164                    has_ext = extension::skip(reader, "Response", ext)?;
165                }
166            }
167        }
168
169        // Payload
170        let payload: ResponseBody = self.codec.read(&mut *reader)?;
171
172        Ok(Response {
173            rid,
174            wire_expr,
175            payload,
176            ext_qos,
177            ext_tstamp,
178            ext_respid,
179            ext_ts_stack,
180        })
181    }
182}
183
184// ResponseFinal
185impl<W> WCodec<&ResponseFinal, &mut W> for Zenoh080
186where
187    W: Writer,
188{
189    type Output = Result<(), DidntWrite>;
190
191    fn write(self, writer: &mut W, x: &ResponseFinal) -> Self::Output {
192        let ResponseFinal {
193            rid,
194            ext_qos,
195            ext_tstamp,
196        } = x;
197
198        // Header
199        let mut header = id::RESPONSE_FINAL;
200        let mut n_exts = ((ext_qos != &ext::QoSType::DEFAULT) as u8) + (ext_tstamp.is_some() as u8);
201        if n_exts != 0 {
202            header |= flag::Z;
203        }
204        self.write(&mut *writer, header)?;
205
206        // Body
207        self.write(&mut *writer, rid)?;
208
209        // Extensions
210        if ext_qos != &ext::QoSType::DEFAULT {
211            n_exts -= 1;
212            self.write(&mut *writer, (*ext_qos, n_exts != 0))?;
213        }
214        if let Some(ts) = ext_tstamp.as_ref() {
215            n_exts -= 1;
216            self.write(&mut *writer, (ts, n_exts != 0))?;
217        }
218
219        Ok(())
220    }
221}
222
223impl<R> RCodec<ResponseFinal, &mut R> for Zenoh080
224where
225    R: Reader,
226{
227    type Error = DidntRead;
228
229    fn read(self, reader: &mut R) -> Result<ResponseFinal, Self::Error> {
230        let header: u8 = self.read(&mut *reader)?;
231        let codec = Zenoh080Header::new(header);
232        codec.read(reader)
233    }
234}
235
236impl<R> RCodec<ResponseFinal, &mut R> for Zenoh080Header
237where
238    R: Reader,
239{
240    type Error = DidntRead;
241
242    fn read(self, reader: &mut R) -> Result<ResponseFinal, Self::Error> {
243        if imsg::mid(self.header) != id::RESPONSE_FINAL {
244            return Err(DidntRead);
245        }
246
247        // Body
248        let bodec = Zenoh080Bounded::<RequestId>::new();
249        let rid: RequestId = bodec.read(&mut *reader)?;
250
251        // Extensions
252        let mut ext_qos = ext::QoSType::DEFAULT;
253        let mut ext_tstamp = None;
254
255        let mut has_ext = imsg::has_flag(self.header, flag::Z);
256        while has_ext {
257            let ext: u8 = self.codec.read(&mut *reader)?;
258            let eodec = Zenoh080Header::new(ext);
259            match iext::eid(ext) {
260                ext::QoS::ID => {
261                    let (q, ext): (ext::QoSType, bool) = eodec.read(&mut *reader)?;
262                    ext_qos = q;
263                    has_ext = ext;
264                }
265                ext::Timestamp::ID => {
266                    let (t, ext): (ext::TimestampType, bool) = eodec.read(&mut *reader)?;
267                    ext_tstamp = Some(t);
268                    has_ext = ext;
269                }
270                _ => {
271                    has_ext = extension::skip(reader, "ResponseFinal", ext)?;
272                }
273            }
274        }
275
276        Ok(ResponseFinal {
277            rid,
278            ext_qos,
279            ext_tstamp,
280        })
281    }
282}