Skip to main content

fdomain_container/
wire.rs

1// Copyright 2024 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use 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/// Wraps an [`FDomain`] and provides an interface that operates on
13/// binary-encoded FIDL messages, as opposed to the FIDL request/response
14/// structs [`FDomain`] itself deals with.
15///
16/// Request messages are passed in with the [`message`] method, and polling the
17/// `FDomainCodec` as a stream will yield the responses.
18#[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    /// Construct a new [`FDomainCodec`] around the given [`FDomain`]
28    pub fn new(fdomain: crate::FDomain) -> FDomainCodec {
29        FDomainCodec { fdomain, outgoing: VecDeque::new(), wakers: Vec::new() }
30    }
31
32    /// Process an incoming message.
33    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    /// Add an outgoing message to our outgoing queue and wake any wakers that
238    /// were waiting for that.
239    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    /// Encode and enqueue a fallible response message to the client. The
245    /// `ordinal` field should correspond correctly to the type of message given
246    /// by the type argument.
247    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    /// Encode and enqueue an event message to the client. The `ordinal` field
270    /// should correspond correctly to the type of message given by the type
271    /// argument.
272    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}