Skip to main content

tor_proto/
relay.rs

1//! This module contains a WIP relay tunnel reactor.
2//!
3//! The initial version will duplicate some of the logic from
4//! the client tunnel reactor.
5//!
6//! TODO(relay): refactor the relay tunnel
7//! to share the same base tunnel implementation
8//! as the client tunnel (to reduce code duplication).
9//!
10//! See the design notes at doc/dev/notes/relay-reactor.md
11
12pub(crate) mod channel;
13#[allow(unreachable_pub)] // TODO(relay): use in tor-chanmgr(?)
14pub mod channel_provider;
15pub(crate) mod reactor;
16
17pub use channel::MaybeVerifiableRelayResponderChannel;
18pub use channel::create_handler::{
19    CircNetParameters, CircuitIncomingStreamReceiver, CongestionControlNetParams,
20    CreateRequestHandler, IncomingStreamRequestFilterFactory,
21};
22
23use derive_deftly::Deftly;
24use oneshot_fused_workaround as oneshot;
25
26use tor_cell::chancell::msg::{self as chanmsg};
27use tor_cell::relaycell::StreamId;
28use tor_cell::relaycell::flow_ctrl::XonKBpsEwma;
29use tor_memquota::derive_deftly_template_HasMemoryCost;
30
31use crate::Error;
32use crate::circuit::celltypes::derive_deftly_template_RestrictedChanMsgSet;
33use crate::circuit::reactor::{CircReactorHandle, CtrlCmd, backward, forward};
34use crate::relay::reactor::backward::Backward;
35use crate::relay::reactor::forward::Forward;
36use crate::stream::incoming::IncomingStreamRequestFilter;
37
38#[cfg(doc)]
39use crate::stream::StreamTarget;
40
41/// A subclass of ChanMsg that can correctly arrive on a live relay
42/// circuit (one where a CREATE* has been received).
43#[derive(Debug, Deftly)]
44#[derive_deftly(HasMemoryCost)]
45#[derive_deftly(RestrictedChanMsgSet)]
46#[deftly(usage = "on an open relay circuit")]
47#[cfg(feature = "relay")]
48#[cfg_attr(not(test), allow(unused))] // TODO(relay)
49pub(crate) enum RelayCircChanMsg {
50    /// A relay cell telling us some kind of remote command from some
51    /// party on the circuit.
52    Relay(chanmsg::Relay),
53    /// A relay early cell that is allowed to contain a CREATE message.
54    RelayEarly(chanmsg::RelayEarly),
55    /// A cell telling us to destroy the circuit.
56    Destroy(chanmsg::Destroy),
57    /// A cell telling us to enable/disable channel padding.
58    PaddingNegotiate(chanmsg::PaddingNegotiate),
59}
60
61/// A handle for interacting with a relay circuit.
62#[allow(unused)] // TODO(relay)
63#[derive(Debug)]
64pub struct RelayCirc(pub(crate) CircReactorHandle<Forward, Backward>);
65
66impl RelayCirc {
67    /// Shut down this circuit, along with all streams that are using it.
68    /// Happens asynchronously (i.e. the tunnel won't necessarily be done shutting down
69    /// immediately after this function returns!).
70    ///
71    /// Note that other references to this tunnel may exist.
72    /// If they do, they will stop working after you call this function.
73    ///
74    /// It's not necessary to call this method if you're just done with a circuit:
75    /// the circuit should close on its own once nothing is using it any more.
76    pub fn terminate(&self) {
77        let _ = self.0.command.unbounded_send(CtrlCmd::Shutdown);
78    }
79
80    /// Return true if this circuit is closed and therefore unusable.
81    pub fn is_closing(&self) -> bool {
82        self.0.control.is_closed()
83    }
84
85    /// Inform the circuit reactor that there has been a change in the drain rate for this stream.
86    ///
87    /// See [`StreamTarget::drain_rate_update`].
88    pub(crate) fn drain_rate_update(
89        &self,
90        stream_id: StreamId,
91        rate: XonKBpsEwma,
92    ) -> crate::Result<()> {
93        let msg = backward::CtrlMsg::FlowCtrlUpdate {
94            hop: None,
95            msg: backward::FlowCtrlMsg::Xon(rate),
96            stream_id,
97        };
98        self.0
99            .control
100            .unbounded_send(msg.into())
101            .map_err(|_| Error::CircuitClosed)?;
102
103        Ok(())
104    }
105
106    /// Request to send a SENDME cell for this stream.
107    ///
108    /// See [`StreamTarget::send_sendme`].
109    pub(crate) fn send_sendme(&self, stream_id: StreamId) -> crate::Result<()> {
110        let msg = backward::CtrlMsg::FlowCtrlUpdate {
111            hop: None,
112            msg: backward::FlowCtrlMsg::Sendme,
113            stream_id,
114        };
115        self.0
116            .control
117            .unbounded_send(msg.into())
118            .map_err(|_| Error::CircuitClosed)?;
119
120        Ok(())
121    }
122
123    /// Close the pending stream that owns this StreamTarget, delivering the specified
124    /// END message (if any)
125    ///
126    /// The stream is closed by sending a control message (`ClosePendingStream`)
127    /// to the reactor.
128    ///
129    /// Returns a [`oneshot::Receiver`] that can be used to await the reactor's response.
130    ///
131    /// The StreamTarget will set the correct stream ID and pick the
132    /// right hop, but will not validate that the message is well-formed
133    /// or meaningful in context.
134    ///
135    /// Note that in many cases, the actual contents of an END message can leak unwanted
136    /// information. Please consider carefully before sending anything but an
137    /// [`End::new_misc()`](tor_cell::relaycell::msg::End::new_misc) message over a `ClientTunnel`.
138    /// (For onion services, we send [`DONE`](tor_cell::relaycell::msg::EndReason::DONE) )
139    ///
140    /// In addition to sending the END message, this function also ensures
141    /// the state of the stream map entry of this stream is updated
142    /// accordingly.
143    ///
144    /// Normally, you shouldn't need to call this function, as streams are implicitly closed by the
145    /// reactor when their corresponding `StreamTarget` is dropped. The only valid use of this
146    /// function is for closing pending incoming streams (a stream is said to be pending if we have
147    /// received the message initiating the stream but have not responded to it yet).
148    ///
149    /// **NOTE**: This function should be called at most once per request.
150    /// Calling it twice is an error.
151    //
152    // TODO(relay): this duplicates the ClientTunnel API and docs. Do we care?
153    pub(crate) fn close_pending(
154        &self,
155        stream_id: StreamId,
156        message: crate::stream::CloseStreamBehavior,
157    ) -> crate::Result<oneshot::Receiver<crate::Result<()>>> {
158        let (tx, rx) = oneshot::channel();
159
160        let cmd = forward::CtrlCmd::ClosePendingStream {
161            stream_id,
162            hop: None,
163            message,
164            done: tx,
165        };
166
167        self.0
168            .command
169            .unbounded_send(CtrlCmd::Forward(cmd))
170            .map_err(|_| Error::CircuitClosed)?;
171
172        Ok(rx)
173    }
174}
175
176#[cfg(test)]
177mod test {
178    // @@ begin test lint list maintained by maint/add_warning @@
179    #![allow(clippy::bool_assert_comparison)]
180    #![allow(clippy::clone_on_copy)]
181    #![allow(clippy::dbg_macro)]
182    #![allow(clippy::mixed_attributes_style)]
183    #![allow(clippy::print_stderr)]
184    #![allow(clippy::print_stdout)]
185    #![allow(clippy::single_char_pattern)]
186    #![allow(clippy::unwrap_used)]
187    #![allow(clippy::unchecked_time_subtraction)]
188    #![allow(clippy::useless_vec)]
189    #![allow(clippy::needless_pass_by_value)]
190    #![allow(clippy::string_slice)] // See arti#2571
191    //! <!-- @@ end test lint list maintained by maint/add_warning @@ -->
192
193    #[test]
194    fn relay_circ_chan_msg() {
195        use tor_cell::chancell::msg::{self, AnyChanMsg};
196        fn good(m: AnyChanMsg) {
197            use crate::relay::RelayCircChanMsg;
198            assert!(RelayCircChanMsg::try_from(m).is_ok());
199        }
200        fn bad(m: AnyChanMsg) {
201            use crate::relay::RelayCircChanMsg;
202            assert!(RelayCircChanMsg::try_from(m).is_err());
203        }
204
205        good(msg::Destroy::new(2.into()).into());
206        bad(msg::CreatedFast::new(&b"The great globular mass"[..]).into());
207        bad(msg::Created2::new(&b"of protoplasmic slush"[..]).into());
208        good(msg::Relay::new(&b"undulated slightly,"[..]).into());
209        good(msg::AnyChanMsg::RelayEarly(
210            msg::Relay::new(&b"as if aware of him"[..]).into(),
211        ));
212        bad(msg::Versions::new([1, 2, 3]).unwrap().into());
213        good(msg::PaddingNegotiate::start_default().into());
214        good(msg::RelayEarly::from(msg::Relay::new(b"snail-like unipedular organism")).into());
215    }
216}