1use 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
33impl<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 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 self.write(&mut *writer, rid)?;
70 self.write(&mut *writer, wire_expr)?;
71
72 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 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 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 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 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
184impl<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 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 self.write(&mut *writer, rid)?;
208
209 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 let bodec = Zenoh080Bounded::<RequestId>::new();
249 let rid: RequestId = bodec.read(&mut *reader)?;
250
251 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}