microsandbox_protocol_client/
stream.rs1use std::sync::{Arc, atomic::Ordering};
4
5use tokio::sync::mpsc;
6
7use crate::router::{Lease, QueuedFrame};
8use crate::{Client, ClientResult, Delivery, IntoOutboundMessage, Message, Protocol, RawFrame};
9
10pub struct RawStream<P: Protocol> {
16 sender: RawStreamSender<P>,
17 pub(crate) receiver: RawStreamReceiver<P>,
18}
19
20pub struct RawStreamSender<P: Protocol> {
22 client: Client<P>,
23 lease: Arc<Lease>,
24}
25
26pub struct RawStreamReceiver<P: Protocol> {
28 client: Client<P>,
29 pub(crate) lease: Arc<Lease>,
30 receiver: mpsc::Receiver<QueuedFrame>,
31 done: bool,
32}
33
34pub struct Stream<P: Protocol> {
36 raw: RawStream<P>,
37}
38
39pub struct StreamSender<P: Protocol> {
41 raw: RawStreamSender<P>,
42}
43
44pub struct StreamReceiver<P: Protocol> {
46 raw: RawStreamReceiver<P>,
47}
48
49impl<P: Protocol> RawStream<P> {
54 pub(crate) fn new(
55 client: Client<P>,
56 lease: Arc<Lease>,
57 receiver: mpsc::Receiver<QueuedFrame>,
58 ) -> Self {
59 Self {
60 sender: RawStreamSender {
61 client: client.clone(),
62 lease: Arc::clone(&lease),
63 },
64 receiver: RawStreamReceiver {
65 client,
66 lease,
67 receiver,
68 done: false,
69 },
70 }
71 }
72
73 pub fn id(&self) -> u32 {
75 self.sender.id()
76 }
77
78 pub async fn send(&self, flags: u8, body: &[u8]) -> ClientResult<()> {
80 self.sender.send(flags, body).await
81 }
82
83 pub async fn recv(&mut self) -> ClientResult<Option<RawFrame>> {
85 self.receiver.recv().await
86 }
87
88 pub fn close(&mut self) {
90 self.receiver.close();
91 }
92
93 pub fn into_parts(self) -> (RawStreamSender<P>, RawStreamReceiver<P>) {
95 (self.sender, self.receiver)
96 }
97}
98
99impl<P: Protocol> RawStreamSender<P> {
100 pub fn id(&self) -> u32 {
102 self.lease.id
103 }
104
105 pub async fn send(&self, flags: u8, body: &[u8]) -> ClientResult<()> {
107 self.client.send_raw_owned(&self.lease, flags, body).await
108 }
109}
110
111impl<P: Protocol> RawStreamReceiver<P> {
112 pub fn id(&self) -> u32 {
114 self.lease.id
115 }
116
117 pub async fn recv(&mut self) -> ClientResult<Option<RawFrame>> {
119 if self.done {
120 return Ok(None);
121 }
122 match self.receiver.recv().await {
123 Some(queued) => {
124 if queued.frame.flags & microsandbox_protocol::message::FLAG_TERMINAL != 0 {
125 self.done = true;
126 }
127 Ok(Some(queued.frame))
128 }
129 None => {
130 self.done = true;
131 if self.lease.terminal.load(Ordering::Acquire) {
132 Ok(None)
133 } else {
134 Err(self
135 .client
136 .inner
137 .state
138 .error()
139 .with_delivery(self.lease.delivery()))
140 }
141 }
142 }
143 }
144
145 pub fn close(&mut self) {
147 self.client.inner.state.abandon(&self.lease);
148 self.receiver.close();
149 while self.receiver.try_recv().is_ok() {}
150 self.done = true;
151 }
152}
153
154impl<P: Protocol> Stream<P> {
155 pub(crate) fn from_raw(raw: RawStream<P>) -> Self {
156 Self { raw }
157 }
158
159 pub fn id(&self) -> u32 {
161 self.raw.id()
162 }
163
164 pub async fn send<M: IntoOutboundMessage<P>>(&self, message: M) -> ClientResult<()> {
166 self.raw
167 .sender
168 .client
169 .send_owned(&self.raw.sender.lease, message)
170 .await
171 }
172
173 pub async fn recv(&mut self) -> ClientResult<Option<Message>> {
175 self.raw
176 .recv()
177 .await?
178 .map(|frame| {
179 self.raw
180 .sender
181 .client
182 .inner
183 .codec
184 .decode(frame)
185 .map_err(|error| error.with_delivery(Delivery::Unknown))
186 })
187 .transpose()
188 }
189
190 pub fn close(&mut self) {
192 self.raw.close();
193 }
194
195 pub fn into_parts(self) -> (StreamSender<P>, StreamReceiver<P>) {
197 let (sender, receiver) = self.raw.into_parts();
198 (
199 StreamSender { raw: sender },
200 StreamReceiver { raw: receiver },
201 )
202 }
203}
204
205impl<P: Protocol> StreamSender<P> {
206 pub fn id(&self) -> u32 {
208 self.raw.id()
209 }
210
211 pub async fn send<M: IntoOutboundMessage<P>>(&self, message: M) -> ClientResult<()> {
213 self.raw.client.send_owned(&self.raw.lease, message).await
214 }
215}
216
217impl<P: Protocol> StreamReceiver<P> {
218 pub fn id(&self) -> u32 {
220 self.raw.id()
221 }
222
223 pub async fn recv(&mut self) -> ClientResult<Option<Message>> {
225 self.raw
226 .recv()
227 .await?
228 .map(|frame| {
229 self.raw
230 .client
231 .inner
232 .codec
233 .decode(frame)
234 .map_err(|error| error.with_delivery(Delivery::Unknown))
235 })
236 .transpose()
237 }
238
239 pub fn close(&mut self) {
241 self.raw.close();
242 }
243}
244
245impl<P: Protocol> Clone for RawStreamSender<P> {
250 fn clone(&self) -> Self {
251 Self {
252 client: self.client.clone(),
253 lease: Arc::clone(&self.lease),
254 }
255 }
256}
257
258impl<P: Protocol> Clone for StreamSender<P> {
259 fn clone(&self) -> Self {
260 Self {
261 raw: self.raw.clone(),
262 }
263 }
264}
265
266impl<P: Protocol> Drop for RawStreamReceiver<P> {
267 fn drop(&mut self) {
268 self.client.inner.state.abandon(&self.lease);
269 }
270}