Skip to main content

sunset_async/
async_channel.rs

1//! Presents SSH channels as async
2use core::future::poll_fn;
3
4#[allow(unused_imports)]
5use log::{debug, error, info, log, trace, warn};
6
7use embedded_io_async::{ErrorType, Read, Write};
8
9use crate::*;
10use sunset::{ChanData, ChanNum, Result};
11
12/// Common implementation
13pub(crate) struct ChanIO<'g> {
14    num: ChanNum,
15    dt: ChanData,
16    sunset: &'g dyn async_sunset::ChanCore,
17}
18
19impl<'g> ChanIO<'g> {
20    /// Create a new Normal ChanIO.
21    ///
22    /// Only to be called by add_channel(), which has already set
23    /// the initial refcount = 1.
24    pub(crate) fn new_normal(
25        num: ChanNum,
26        sunset: &'g dyn async_sunset::ChanCore,
27    ) -> Self {
28        Self { num, dt: ChanData::Normal, sunset }
29    }
30
31    pub(crate) fn clone_stderr(&self) -> Self {
32        let mut c = self.clone();
33        c.dt = ChanData::Stderr;
34        c
35    }
36}
37
38impl core::fmt::Debug for ChanIO<'_> {
39    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
40        f.debug_struct("ChanIO")
41            .field("num", &self.num)
42            .field("dt", &self.dt)
43            .finish_non_exhaustive()
44    }
45}
46
47impl ChanIO<'_> {
48    pub async fn until_closed(&self) -> Result<()> {
49        poll_fn(|cx| self.sunset.poll_until_channel_closed(cx, self.num)).await
50    }
51
52    pub async fn term_window_change(
53        &self,
54        winch: sunset::packets::WinChange,
55    ) -> Result<()> {
56        poll_fn(|cx| self.sunset.poll_term_window_change(cx, self.num, &winch)).await
57    }
58}
59
60impl Drop for ChanIO<'_> {
61    fn drop(&mut self) {
62        self.sunset.dec_chan(self.num)
63    }
64}
65
66// ChanIO implements Clone to share between ChanIn/ChanOut/ChanInOut.
67// There's only one waker for each of in/out/ext, so allowing clone
68// on the ChanInOut etc isn't desirable - having two instances polling
69// the same direction/dt will just result in churn between wakers if they're
70// in different tasks.
71impl Clone for ChanIO<'_> {
72    fn clone(&self) -> Self {
73        self.sunset.inc_chan(self.num);
74        Self { num: self.num, dt: self.dt, sunset: self.sunset }
75    }
76}
77
78impl ErrorType for ChanIO<'_> {
79    type Error = sunset::Error;
80}
81
82impl Read for ChanIO<'_> {
83    async fn read(&mut self, buf: &mut [u8]) -> Result<usize, sunset::Error> {
84        poll_fn(|cx| self.sunset.poll_read_channel(cx, self.num, self.dt, buf)).await
85    }
86}
87
88impl Write for ChanIO<'_> {
89    async fn write(&mut self, buf: &[u8]) -> Result<usize, sunset::Error> {
90        poll_fn(|cx| self.sunset.poll_write_channel(cx, self.num, self.dt, buf))
91            .await
92    }
93
94    async fn flush(&mut self) -> Result<()> {
95        // TODO: could this wait for the packet to get sent out?
96        Ok(())
97    }
98}
99
100// Public wrappers for In only
101
102/// An input-only SSH channel.
103///
104/// This is used as stderr for a client.
105///
106/// <div class="warning">
107///
108/// This must be read, otherwise the SSH session will block.
109///
110/// </div>
111///
112/// `Clone` is implemented for convenience, but only one instance each
113/// should be read from.
114/// Otherwise ordering will be arbitrary, and if competing readers or writers
115/// are in different tasks, there will be churn as they continually wake
116/// each other up. Simultaneous single-reader and single-writer is fine.
117#[derive(Debug)]
118pub struct ChanIn<'g>(ChanIO<'g>);
119
120impl<'g> ChanIn<'g> {
121    pub(crate) fn new(io: ChanIO<'g>) -> Self {
122        io.sunset.inc_read_chan(io.num, io.dt);
123        Self(io)
124    }
125
126    /// Return the channel number.
127    pub fn num(&self) -> ChanNum {
128        self.0.num
129    }
130
131    /// Wait until the channel closes.
132    pub async fn until_closed(&self) -> Result<()> {
133        self.0.until_closed().await
134    }
135}
136
137impl Drop for ChanIn<'_> {
138    fn drop(&mut self) {
139        self.0.sunset.dec_read_chan(self.0.num, self.0.dt)
140    }
141}
142
143impl Clone for ChanIn<'_> {
144    fn clone(&self) -> Self {
145        Self::new(self.0.clone())
146    }
147}
148
149/// An output-only SSH channel.
150///
151/// This is used as stderr for a server, or can also be obtained using
152/// [`ChanInOut::split()`] for cases where a channel's input should
153/// be discarded.
154///
155/// `Clone` is implemented for convenience, but only one instance each
156/// should be read from or written to (this applies to `split()` instances too).
157/// Otherwise ordering will be arbitrary, and if competing readers or writers
158/// are in different tasks, there will be churn as they continually wake
159/// each other up. Simultaneous single-reader and single-writer is fine.
160#[derive(Debug, Clone)]
161pub struct ChanOut<'g>(ChanIO<'g>);
162
163impl<'g> ChanOut<'g> {
164    pub(crate) fn new(io: ChanIO<'g>) -> Self {
165        Self(io)
166    }
167
168    /// Return the channel number.
169    pub fn num(&self) -> ChanNum {
170        self.0.num
171    }
172
173    /// Wait until the channel closes.
174    pub async fn until_closed(&self) -> Result<()> {
175        self.0.until_closed().await
176    }
177
178    /// Send a terminal size change notification
179    ///
180    /// Only applicable to client shell channels with a PTY
181    pub async fn term_window_change(
182        &self,
183        winch: sunset::packets::WinChange,
184    ) -> Result<()> {
185        self.0.term_window_change(winch).await
186    }
187}
188
189/// A bidirectional SSH channel.
190///
191/// Used as stdin/stdout for a shell/exec/subsystem.
192/// Represents other forwarded transports.
193///
194/// <div class="warning">
195///
196/// This must be read, otherwise the SSH session will block.
197/// If input isn't required, use [`split()`](Self::split) and
198/// discard the input half.
199///
200/// </div>
201///
202/// `Clone` is implemented for convenience, but only one instance each
203/// should be read from or written to (this applies to `split()` instances too).
204/// Otherwise ordering will be arbitrary, and if competing readers or writers
205/// are in different tasks, there will be churn as they continually wake
206/// each other up. Simultaneous single-reader and single-writer is fine.
207#[derive(Debug)]
208pub struct ChanInOut<'g>(ChanIO<'g>);
209
210impl<'g> ChanInOut<'g> {
211    pub(crate) fn new(io: ChanIO<'g>) -> Self {
212        io.sunset.inc_read_chan(io.num, io.dt);
213        Self(io)
214    }
215
216    /// Return the channel number.
217    pub fn num(&self) -> ChanNum {
218        self.0.num
219    }
220
221    /// Convert this into separate input and output.
222    ///
223    /// Note the warning above against simultaneous use and `Clone`.
224    pub fn split(&self) -> (ChanIn<'g>, ChanOut<'g>) {
225        (ChanIn::new(self.0.clone()), ChanOut::new(self.0.clone()))
226    }
227
228    /// Wait until the channel closes.
229    pub async fn until_closed(&self) -> Result<()> {
230        self.0.until_closed().await
231    }
232
233    /// Send a terminal size change notification
234    ///
235    /// Only applicable to client shell channels with a PTY
236    pub async fn term_window_change(
237        &self,
238        winch: sunset::packets::WinChange,
239    ) -> Result<()> {
240        self.0.term_window_change(winch).await
241    }
242}
243
244impl Drop for ChanInOut<'_> {
245    fn drop(&mut self) {
246        self.0.sunset.dec_read_chan(self.0.num, self.0.dt)
247    }
248}
249
250impl Clone for ChanInOut<'_> {
251    fn clone(&self) -> Self {
252        Self::new(self.0.clone())
253    }
254}
255
256impl ErrorType for ChanInOut<'_> {
257    type Error = sunset::Error;
258}
259
260impl ErrorType for ChanIn<'_> {
261    type Error = sunset::Error;
262}
263
264impl ErrorType for ChanOut<'_> {
265    type Error = sunset::Error;
266}
267
268impl Read for ChanInOut<'_> {
269    async fn read(&mut self, buf: &mut [u8]) -> Result<usize, sunset::Error> {
270        self.0.read(buf).await
271    }
272}
273
274impl Write for ChanInOut<'_> {
275    async fn write(&mut self, buf: &[u8]) -> Result<usize, sunset::Error> {
276        self.0.write(buf).await
277    }
278
279    async fn flush(&mut self) -> Result<()> {
280        // TODO: could this wait for the packet to get sent out?
281        Ok(())
282    }
283}
284
285impl Read for ChanIn<'_> {
286    async fn read(&mut self, buf: &mut [u8]) -> Result<usize, sunset::Error> {
287        self.0.read(buf).await
288    }
289}
290
291impl Write for ChanOut<'_> {
292    async fn write(&mut self, buf: &[u8]) -> Result<usize, sunset::Error> {
293        self.0.write(buf).await
294    }
295
296    async fn flush(&mut self) -> Result<()> {
297        // TODO: could this wait for the packet to get sent out?
298        Ok(())
299    }
300}