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}