Skip to main content

microsandbox_protocol_client/
stream.rs

1//! Owned subscriptions, split senders, and one consuming receiver per ID.
2
3use 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
10//--------------------------------------------------------------------------------------------------
11// Types
12//--------------------------------------------------------------------------------------------------
13
14/// Opaque frame stream whose lifetime retains its connection and ID lease.
15pub struct RawStream<P: Protocol> {
16    sender: RawStreamSender<P>,
17    pub(crate) receiver: RawStreamReceiver<P>,
18}
19
20/// Cloneable send permission bound to one specific lease, not merely its number.
21pub struct RawStreamSender<P: Protocol> {
22    client: Client<P>,
23    lease: Arc<Lease>,
24}
25
26/// Sole consuming raw receiver. Dropping it disables sends and starts draining.
27pub struct RawStreamReceiver<P: Protocol> {
28    client: Client<P>,
29    pub(crate) lease: Arc<Lease>,
30    receiver: mpsc::Receiver<QueuedFrame>,
31    done: bool,
32}
33
34/// Decoded message stream over the same opaque router and lease.
35pub struct Stream<P: Protocol> {
36    raw: RawStream<P>,
37}
38
39/// Cloneable native-message sender bound to an owned stream lease.
40pub struct StreamSender<P: Protocol> {
41    raw: RawStreamSender<P>,
42}
43
44/// Sole decoded receiver; no competing reader is created by conversion/split.
45pub struct StreamReceiver<P: Protocol> {
46    raw: RawStreamReceiver<P>,
47}
48
49//--------------------------------------------------------------------------------------------------
50// Methods
51//--------------------------------------------------------------------------------------------------
52
53impl<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    /// Current stream's correlation ID.
74    pub fn id(&self) -> u32 {
75        self.sender.id()
76    }
77
78    /// Send a follow-up opaque frame; local close never sends a protocol signal.
79    pub async fn send(&self, flags: u8, body: &[u8]) -> ClientResult<()> {
80        self.sender.send(flags, body).await
81    }
82
83    /// Receive each frame, terminal included, then clean exhaustion or an error.
84    pub async fn recv(&mut self) -> ClientResult<Option<RawFrame>> {
85        self.receiver.recv().await
86    }
87
88    /// Abandon local consumption and retain bounded drain state until terminal.
89    pub fn close(&mut self) {
90        self.receiver.close();
91    }
92
93    /// Move ownership into a cloneable sender and one consuming receiver.
94    pub fn into_parts(self) -> (RawStreamSender<P>, RawStreamReceiver<P>) {
95        (self.sender, self.receiver)
96    }
97}
98
99impl<P: Protocol> RawStreamSender<P> {
100    /// Correlation ID associated with this specific lease.
101    pub fn id(&self) -> u32 {
102        self.lease.id
103    }
104
105    /// Send while this lease remains live; a reused numeric ID gives no authority.
106    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    /// Correlation ID whose frames this receiver alone consumes.
113    pub fn id(&self) -> u32 {
114        self.lease.id
115    }
116
117    /// Receive the next frame. EOF before terminal is an error exactly once.
118    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    /// Stop local consumption and disable sends without canceling remote work.
146    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    /// Correlation ID of the owned stream.
160    pub fn id(&self) -> u32 {
161        self.raw.id()
162    }
163
164    /// Send a native or already-encoded payload using the selected codec.
165    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    /// Receive an inspectable message, including unknown names and peer errors.
174    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    /// Abandon consumption; no domain-specific cancellation or EOF is sent.
191    pub fn close(&mut self) {
192        self.raw.close();
193    }
194
195    /// Split into owned handles while preserving one receiver and the ID lease.
196    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    /// Correlation ID of this lease.
207    pub fn id(&self) -> u32 {
208        self.raw.id()
209    }
210
211    /// Send on this lease, rejecting stale senders even after numeric ID reuse.
212    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    /// Correlation ID consumed by this receiver.
219    pub fn id(&self) -> u32 {
220        self.raw.id()
221    }
222
223    /// Receive messages on the sole raw receiver, decoding only on demand.
224    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    /// Disable further sends and drain until terminal completion.
240    pub fn close(&mut self) {
241        self.raw.close();
242    }
243}
244
245//--------------------------------------------------------------------------------------------------
246// Trait Implementations
247//--------------------------------------------------------------------------------------------------
248
249impl<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}