1use fidl_fuchsia_fdomain as proto;
6use std::collections::VecDeque;
7use std::num::NonZeroU32;
8use std::task::{Context, Poll, Waker};
9
10use proto::f_domain_ordinals as ordinals;
11
12#[pin_project::pin_project]
19pub struct FDomainCodec {
20 #[pin]
21 fdomain: crate::FDomain,
22 outgoing: VecDeque<Box<[u8]>>,
23 wakers: Vec<Waker>,
24}
25
26impl FDomainCodec {
27 pub fn new(fdomain: crate::FDomain) -> FDomainCodec {
29 FDomainCodec { fdomain, outgoing: VecDeque::new(), wakers: Vec::new() }
30 }
31
32 pub fn message(&mut self, data: &[u8]) -> fidl::Result<()> {
34 let (header, rest) = fidl_message::decode_transaction_header(data)?;
35 let Some(tx_id) = NonZeroU32::new(header.tx_id) else {
36 return Err(fidl::Error::UnknownOrdinal {
37 ordinal: header.ordinal,
38 protocol_name:
39 <proto::FDomainMarker as fidl::endpoints::ProtocolMarker>::DEBUG_NAME,
40 });
41 };
42
43 match header.ordinal {
44 ordinals::GET_NAMESPACE => {
45 let request = fidl_message::decode_message::<proto::FDomainGetNamespaceRequest>(
46 header, rest,
47 )?;
48 let result = self.fdomain.get_namespace(request);
49 self.send_response(tx_id, header.ordinal, result)?;
50 }
51 ordinals::CREATE_CHANNEL => {
52 let request = fidl_message::decode_message::<proto::ChannelCreateChannelRequest>(
53 header, rest,
54 )?;
55 let result = self.fdomain.create_channel(request);
56 self.send_response(tx_id, header.ordinal, result)?;
57 }
58 ordinals::CREATE_SOCKET => {
59 let request =
60 fidl_message::decode_message::<proto::SocketCreateSocketRequest>(header, rest)?;
61 let result = self.fdomain.create_socket(request);
62 self.send_response(tx_id, header.ordinal, result)?;
63 }
64 ordinals::CREATE_EVENT_PAIR => {
65 let request = fidl_message::decode_message::<proto::EventPairCreateEventPairRequest>(
66 header, rest,
67 )?;
68 let result = self.fdomain.create_event_pair(request);
69 self.send_response(tx_id, header.ordinal, result)?;
70 }
71 ordinals::CREATE_EVENT => {
72 let request =
73 fidl_message::decode_message::<proto::EventCreateEventRequest>(header, rest)?;
74 let result = self.fdomain.create_event(request);
75 self.send_response(tx_id, header.ordinal, result)?;
76 }
77 ordinals::SET_SOCKET_DISPOSITION => {
78 let request = fidl_message::decode_message::<
79 proto::SocketSetSocketDispositionRequest,
80 >(header, rest)?;
81 self.fdomain.set_socket_disposition(tx_id, request);
82 }
83 ordinals::READ_SOCKET => {
84 let request =
85 fidl_message::decode_message::<proto::SocketReadSocketRequest>(header, rest)?;
86 self.fdomain.read_socket(tx_id, request);
87 }
88 ordinals::READ_CHANNEL => {
89 let request =
90 fidl_message::decode_message::<proto::ChannelReadChannelRequest>(header, rest)?;
91 self.fdomain.read_channel(tx_id, request);
92 }
93 ordinals::WRITE_SOCKET => {
94 let request =
95 fidl_message::decode_message::<proto::SocketWriteSocketRequest>(header, rest)?;
96 self.fdomain.write_socket(tx_id, request);
97 }
98 ordinals::WRITE_CHANNEL => {
99 let request = fidl_message::decode_message::<proto::ChannelWriteChannelRequest>(
100 header, rest,
101 )?;
102 self.fdomain.write_channel(tx_id, request);
103 }
104 ordinals::WAIT_FOR_SIGNALS => {
105 let request = fidl_message::decode_message::<proto::FDomainWaitForSignalsRequest>(
106 header, rest,
107 )?;
108 self.fdomain.wait_for_signals(tx_id, request);
109 }
110 ordinals::CLOSE => {
111 let request =
112 fidl_message::decode_message::<proto::FDomainCloseRequest>(header, rest)?;
113 self.fdomain.close(tx_id, request);
114 }
115 ordinals::DUPLICATE => {
116 let request =
117 fidl_message::decode_message::<proto::FDomainDuplicateRequest>(header, rest)?;
118 let result = self.fdomain.duplicate(request);
119 self.send_response(tx_id, header.ordinal, result)?;
120 }
121 ordinals::REPLACE => {
122 let request =
123 fidl_message::decode_message::<proto::FDomainReplaceRequest>(header, rest)?;
124 let result = self.fdomain.replace(tx_id, request);
125 self.send_response(tx_id, header.ordinal, result)?;
126 }
127 ordinals::SIGNAL => {
128 let request =
129 fidl_message::decode_message::<proto::FDomainSignalRequest>(header, rest)?;
130 let result = self.fdomain.signal(request);
131 self.send_response(tx_id, header.ordinal, result)?;
132 }
133 ordinals::SIGNAL_PEER => {
134 let request =
135 fidl_message::decode_message::<proto::FDomainSignalPeerRequest>(header, rest)?;
136 let result = self.fdomain.signal_peer(request);
137 self.send_response(tx_id, header.ordinal, result)?;
138 }
139 ordinals::READ_CHANNEL_STREAMING_START => {
140 let request = fidl_message::decode_message::<
141 proto::ChannelReadChannelStreamingStartRequest,
142 >(header, rest)?;
143 self.fdomain.read_channel_streaming_start(tx_id, request);
144 }
145 ordinals::READ_CHANNEL_STREAMING_STOP => {
146 let request = fidl_message::decode_message::<
147 proto::ChannelReadChannelStreamingStopRequest,
148 >(header, rest)?;
149 self.fdomain.read_channel_streaming_stop(tx_id, request);
150 }
151 ordinals::READ_SOCKET_STREAMING_START => {
152 let request = fidl_message::decode_message::<
153 proto::SocketReadSocketStreamingStartRequest,
154 >(header, rest)?;
155 self.fdomain.read_socket_streaming_start(tx_id, request);
156 }
157 ordinals::READ_SOCKET_STREAMING_STOP => {
158 let request = fidl_message::decode_message::<
159 proto::SocketReadSocketStreamingStopRequest,
160 >(header, rest)?;
161 self.fdomain.read_socket_streaming_stop(tx_id, request);
162 }
163 ordinals::GET_KOID => {
164 let request =
165 fidl_message::decode_message::<proto::FDomainGetKoidRequest>(header, rest)?;
166 let result = self.fdomain.get_koid(request);
167 self.send_response(tx_id, header.ordinal, result)?;
168 }
169 ordinals::CREATE_VMO => {
170 let request =
171 fidl_message::decode_message::<proto::VmoCreateVmoRequest>(header, rest)?;
172 let result = self.fdomain.create_vmo(request);
173 self.send_response(tx_id, header.ordinal, result)?;
174 }
175 ordinals::READ_VMO => {
176 let request =
177 fidl_message::decode_message::<proto::VmoReadVmoRequest>(header, rest)?;
178 let result = self.fdomain.read_vmo(request);
179 self.send_response(tx_id, header.ordinal, result)?;
180 }
181 ordinals::WRITE_VMO => {
182 let request =
183 fidl_message::decode_message::<proto::VmoWriteVmoRequest>(header, rest)?;
184 let result = self.fdomain.write_vmo(request);
185 self.send_response(tx_id, header.ordinal, result)?;
186 }
187 ordinals::GET_VMO_SIZE => {
188 let request =
189 fidl_message::decode_message::<proto::VmoGetVmoSizeRequest>(header, rest)?;
190 let result = self.fdomain.get_vmo_size(request);
191 self.send_response(tx_id, header.ordinal, result)?;
192 }
193 ordinals::SET_VMO_SIZE => {
194 let request =
195 fidl_message::decode_message::<proto::VmoSetVmoSizeRequest>(header, rest)?;
196 let result = self.fdomain.set_vmo_size(request);
197 self.send_response(tx_id, header.ordinal, result)?;
198 }
199 ordinals::GET_VMO_STREAM_SIZE => {
200 let request = fidl_message::decode_message::<proto::VmoGetVmoStreamSizeRequest>(
201 header, rest,
202 )?;
203 let result = self.fdomain.get_vmo_stream_size(request);
204 self.send_response(tx_id, header.ordinal, result)?;
205 }
206 ordinals::SET_VMO_STREAM_SIZE => {
207 let request = fidl_message::decode_message::<proto::VmoSetVmoStreamSizeRequest>(
208 header, rest,
209 )?;
210 let result = self.fdomain.set_vmo_stream_size(request);
211 self.send_response(tx_id, header.ordinal, result)?;
212 }
213 unknown if header.dynamic_flags().contains(fidl_message::DynamicFlags::FLEXIBLE) => {
214 if header.tx_id != 0 {
215 let header = fidl_message::TransactionHeader::new(
216 header.tx_id,
217 unknown,
218 fidl_message::DynamicFlags::FLEXIBLE,
219 );
220 self.enqueue_outgoing::<Vec<u8>>(
221 fidl_message::encode_response_flexible_unknown(header)?.into(),
222 );
223 }
224 }
225 _ => {
226 return Err(fidl::Error::UnknownOrdinal {
227 ordinal: header.ordinal,
228 protocol_name:
229 <proto::FDomainMarker as fidl::endpoints::ProtocolMarker>::DEBUG_NAME,
230 });
231 }
232 }
233
234 Ok(())
235 }
236
237 fn enqueue_outgoing<T: Into<Box<[u8]>>>(&mut self, msg: T) {
240 self.outgoing.push_back(msg.into());
241 self.wakers.drain(..).for_each(Waker::wake);
242 }
243
244 fn send_response<T: fidl_message::Body, E: fidl_message::ErrorType>(
248 &mut self,
249 tx_id: NonZeroU32,
250 ordinal: u64,
251 body: Result<T, E>,
252 ) -> fidl::Result<()>
253where
254 for<'a> <<T as fidl_message::Body>::MarkerInResultUnion as fidl::encoding::ValueTypeMarker>::Borrowed<'a>:
255 fidl::encoding::Encode<T::MarkerInResultUnion, fidl::encoding::NoHandleResourceDialect>,
256 for<'a> <<E as fidl_message::ErrorType>::Marker as fidl::encoding::ValueTypeMarker>::Borrowed<'a>:
257 fidl::encoding::Encode<E::Marker, fidl::encoding::NoHandleResourceDialect>,
258{
259 let header = fidl_message::TransactionHeader::new(
260 tx_id.into(),
261 ordinal,
262 fidl_message::DynamicFlags::FLEXIBLE,
263 );
264 self.enqueue_outgoing(fidl_message::encode_response_result::<T, E>(header, body)?);
265
266 Ok(())
267 }
268
269 fn send_event<T: fidl_message::Body>(&mut self, ordinal: u64, body: T) -> fidl::Result<()>
273 where
274 for<'a> <<T as fidl_message::Body>::MarkerAtTopLevel as fidl::encoding::ValueTypeMarker>::Borrowed<
275 'a,
276 >: fidl::encoding::Encode<T::MarkerAtTopLevel, fidl::encoding::NoHandleResourceDialect>,
277 {
278 let header =
279 fidl_message::TransactionHeader::new(0, ordinal, fidl_message::DynamicFlags::empty());
280 self.enqueue_outgoing(fidl_message::encode_message(header, body)?);
281
282 Ok(())
283 }
284}
285
286impl futures::Stream for FDomainCodec {
287 type Item = fidl::Result<Box<[u8]>>;
288
289 fn poll_next(
290 mut self: std::pin::Pin<&mut Self>,
291 ctx: &mut Context<'_>,
292 ) -> Poll<Option<Self::Item>> {
293 while let Poll::Ready(Some(event)) = self.as_mut().project().fdomain.poll_next(ctx) {
294 let result = match event {
295 crate::FDomainEvent::ChannelStreamingReadStart(tx_id, msg) => {
296 self.send_response(tx_id, ordinals::READ_CHANNEL_STREAMING_START, msg)
297 }
298 crate::FDomainEvent::ChannelStreamingReadStop(tx_id, msg) => {
299 self.send_response(tx_id, ordinals::READ_CHANNEL_STREAMING_STOP, msg)
300 }
301 crate::FDomainEvent::SocketStreamingReadStart(tx_id, msg) => {
302 self.send_response(tx_id, ordinals::READ_SOCKET_STREAMING_START, msg)
303 }
304 crate::FDomainEvent::SocketStreamingReadStop(tx_id, msg) => {
305 self.send_response(tx_id, ordinals::READ_SOCKET_STREAMING_STOP, msg)
306 }
307 crate::FDomainEvent::WaitForSignals(tx_id, msg) => {
308 self.send_response(tx_id, ordinals::WAIT_FOR_SIGNALS, msg)
309 }
310 crate::FDomainEvent::SocketData(tx_id, msg) => {
311 self.send_response(tx_id, ordinals::READ_SOCKET, msg)
312 }
313 crate::FDomainEvent::SocketStreamingData(msg) => {
314 self.send_event(ordinals::ON_SOCKET_STREAMING_DATA, msg)
315 }
316 crate::FDomainEvent::SocketDispositionSet(tx_id, msg) => {
317 self.send_response(tx_id, ordinals::SET_SOCKET_DISPOSITION, msg)
318 }
319 crate::FDomainEvent::WroteSocket(tx_id, msg) => {
320 self.send_response(tx_id, ordinals::WRITE_SOCKET, msg)
321 }
322 crate::FDomainEvent::ChannelData(tx_id, msg) => {
323 self.send_response(tx_id, ordinals::READ_CHANNEL, msg)
324 }
325 crate::FDomainEvent::ChannelStreamingData(msg) => {
326 self.send_event(ordinals::ON_CHANNEL_STREAMING_DATA, msg)
327 }
328 crate::FDomainEvent::WroteChannel(tx_id, msg) => {
329 self.send_response(tx_id, ordinals::WRITE_CHANNEL, msg)
330 }
331 crate::FDomainEvent::ClosedHandle(tx_id, msg) => {
332 self.send_response(tx_id, ordinals::CLOSE, msg)
333 }
334 crate::FDomainEvent::ReplacedHandle(tx_id, msg) => {
335 self.send_response(tx_id, ordinals::REPLACE, msg)
336 }
337 };
338
339 if let Err(e) = result {
340 return Poll::Ready(Some(Err(e)));
341 }
342 }
343
344 if let Some(got) = self.outgoing.pop_front() {
345 Poll::Ready(Some(Ok(got)))
346 } else {
347 self.wakers.push(ctx.waker().clone());
348 Poll::Pending
349 }
350 }
351}