tor_proto/circuit/reactor/stream.rs
1//! The stream reactor.
2
3use crate::circuit::circhop::CircHopOutbound;
4use crate::circuit::reactor::macros::derive_deftly_template_CircuitReactor;
5use crate::circuit::{CircHopSyncView, UniqId};
6use crate::congestion::{CongestionControl, sendme};
7use crate::memquota::{CircuitAccount, SpecificAccount as _, StreamAccount};
8use crate::stream::CloseStreamBehavior;
9use crate::stream::cmdcheck::StreamStatus;
10use crate::stream::flow_ctrl::state::WithSidechannelMitigations;
11use crate::streammap;
12use crate::util::err::ReactorError;
13use crate::{Error, HopNum};
14
15#[cfg(any(feature = "hs-service", feature = "relay"))]
16use crate::stream::incoming::{
17 InboundDataCmdChecker, IncomingStreamRequest, IncomingStreamRequestContext,
18 IncomingStreamRequestDisposition, IncomingStreamRequestHandler, StreamReqInfo,
19};
20
21use tor_async_utils::{SinkTrySend as _, SinkTrySendError as _};
22use tor_cell::chancell::CircId;
23use tor_cell::relaycell::msg::{AnyRelayMsg, Begin, BeginDir, End, EndReason, Resolve};
24use tor_cell::relaycell::{
25 AnyRelayMsgOuter, RelayCellFormat, RelayCmd, StreamId, UnparsedRelayMsg,
26};
27use tor_error::{internal, into_internal};
28use tor_log_ratelim::log_ratelim;
29use tor_rtcompat::{DynTimeProvider, Runtime, SleepProvider as _};
30
31use derive_deftly::Deftly;
32use futures::SinkExt;
33use futures::channel::mpsc;
34use futures::{FutureExt as _, StreamExt as _, future, select_biased};
35use tracing::debug;
36
37use std::pin::Pin;
38use std::result::Result as StdResult;
39use std::sync::{Arc, Mutex};
40use std::task::Poll;
41use std::time::Duration;
42
43/// Trait for customizing the behavior of the stream reactor.
44///
45/// Used for plugging in the implementation-dependent (client vs relay)
46/// parts of the implementation into the generic one.
47pub(crate) trait StreamHandler: Send + Sync + 'static {
48 /// Return the amount of time a newly closed stream
49 /// should be kept in the stream map for.
50 ///
51 /// This is the amount of time we are willing to wait for
52 /// an END ack before removing the half-stream from the map.
53 fn halfstream_expiry(&self, hop: &CircHopOutbound) -> Duration;
54
55 /// Whether sidechannel mitigations should be enabled for incoming streams.
56 fn flowctrl_sidechannel_mitigations(&self) -> WithSidechannelMitigations;
57}
58
59/// The stream reactor for a given hop.
60///
61/// Drives the application streams.
62///
63/// This reactor accepts [`CtrlMsg`]s from the forward reactor over its [`Self::cell_rx`]
64/// MPSC channel, and delivers them to the corresponding stream entries in the stream map.
65///
66/// The local streams are polled from the main loop, and any ready messages are sent
67/// to the backward reactor over the `bwd_tx` MPSC channel for packaging and delivery.
68///
69/// Shuts downs down if an error occurs, or if the sending end
70/// of the `cell_rx` MPSC channel, i.e. the forward reactor, closes.
71#[derive(Deftly)]
72#[derive_deftly(CircuitReactor)]
73#[deftly(reactor_name = "stream reactor")]
74#[deftly(run_inner_fn = "Self::run_once")]
75#[must_use = "If you don't call run() on a reactor, the circuit won't work."]
76pub(crate) struct StreamReactor {
77 /// The hop this stream reactor is for.
78 ///
79 /// This is `None` for relays.
80 hopnum: Option<HopNum>,
81 /// The state of this circuit hop.
82 hop: CircHopOutbound,
83 /// The time provider.
84 time_provider: DynTimeProvider,
85 /// An identifier for logging about this reactor's circuit.
86 unique_id: UniqId,
87 /// The circuit identifier on the inbound Tor channel.
88 circ_id: CircId,
89 /// Receiver for Tor stream data that need to be delivered to a Tor stream.
90 ///
91 /// The sender is in the [`HopMgr`](super::hop_mgr::HopMgr) of the
92 /// [`ForwardReactor`](super::ForwardReactor), which will forward all cells
93 /// carrying Tor stream data to us.
94 ///
95 /// This serves a dual purpose:
96 ///
97 /// * it enables the `ForwardReactor` to deliver Tor stream data received from the client
98 /// * it lets the `StreamReactor` know if the `ForwardReactor` has shut down:
99 /// we select! on this MPSC channel in the main loop, so if the `ForwardReactor`
100 /// shuts down, we will get EOS upon calling `.next()`)
101 cell_rx: mpsc::Receiver<CtrlMsg>,
102 /// Sender for sending Tor stream data to [`BackwardReactor`](super::BackwardReactor).
103 bwd_tx: mpsc::Sender<ReadyStreamMsg>,
104 /// A handler for incoming streams.
105 ///
106 /// Set to `None` if incoming streams are not allowed on this circuit.
107 ///
108 /// This handler is shared with the [`HopMgr`](super::hop_mgr::HopMgr) of this reactor,
109 /// which can install a new handler at runtime (for example, in response to a CtrlMsg).
110 /// The ability to update the handler after the reactor is launched is needed
111 /// for onion services, where the incoming stream request handler only gets installed
112 /// after the virtual hop is created.
113 #[cfg(any(feature = "hs-service", feature = "relay"))]
114 incoming: Arc<Mutex<Option<IncomingStreamRequestHandler>>>,
115 /// A handler for customizing the stream reactor behavior.
116 inner: Arc<dyn StreamHandler>,
117 /// Memory quota account
118 memquota: CircuitAccount,
119}
120
121#[allow(unused)] // TODO(relay)
122impl StreamReactor {
123 /// Create a new [`StreamReactor`].
124 #[allow(clippy::too_many_arguments)] // TODO
125 pub(crate) fn new<R: Runtime>(
126 runtime: R,
127 hopnum: Option<HopNum>,
128 hop: CircHopOutbound,
129 unique_id: UniqId,
130 circ_id: CircId,
131 cell_rx: mpsc::Receiver<CtrlMsg>,
132 bwd_tx: mpsc::Sender<ReadyStreamMsg>,
133 inner: Arc<dyn StreamHandler>,
134 #[cfg(any(feature = "hs-service", feature = "relay"))] //
135 incoming: Arc<Mutex<Option<IncomingStreamRequestHandler>>>,
136 memquota: CircuitAccount,
137 ) -> Self {
138 Self {
139 hopnum,
140 hop,
141 time_provider: DynTimeProvider::new(runtime),
142 unique_id,
143 circ_id,
144 #[cfg(any(feature = "hs-service", feature = "relay"))]
145 incoming,
146 cell_rx,
147 bwd_tx,
148 inner,
149 memquota,
150 }
151 }
152
153 /// Helper for [`run`](Self::run).
154 ///
155 /// Polls the stream map for messages
156 /// that need to be delivered to the other endpoint,
157 /// and the `cells_rx` MPSC stream for stream messages received
158 /// from the `ForwardReactor` that need to be delivered to the application streams.
159 async fn run_once(&mut self) -> StdResult<(), ReactorError> {
160 use postage::prelude::{Sink as _, Stream as _};
161
162 // Garbage-collect all halfstreams that have expired.
163 //
164 // Note: this will iterate over the closed streams of this hop.
165 // If we think this will cause perf issues, one idea would be to make
166 // StreamMap::closed_streams into a min-heap, and add a branch to the
167 // select_biased! below to sleep until the first expiry is due
168 // (but my gut feeling is that iterating is cheaper)
169 self.hop
170 .stream_map()
171 .lock()
172 .expect("poisoned lock")
173 .remove_expired_halfstreams(self.time_provider.now());
174
175 let mut streams = Arc::clone(self.hop.stream_map());
176 let can_send = self
177 .hop
178 .ccontrol()
179 .lock()
180 .expect("poisoned lock")
181 .can_send();
182 let mut ready_streams_fut = future::poll_fn(move |cx| {
183 if !can_send {
184 // We can't send anything on this hop that counts towards SENDME windows.
185 //
186 // Note: this does not block outgoing flow-control messages:
187 //
188 // * circuit SENDMEs are initiated by the forward reactor,
189 // by sending a BackwardReactorCmd::SendRelayMsg to BWD,
190 // * stream SENDMEs will be initiated by StreamTarget::send_sendme(),
191 // by sending a control message to the reactor
192 // (TODO(relay): not yet implemented)
193 // * XOFFs are sent in response to messages on streams
194 // (i.e. RELAY messages with non-zero stream IDs).
195 // These messages are delivered to us by the forward reactor
196 // inside BackwardReactorCmd::HandleMsg
197 // * XON will be initiated by StreamTarget::drain_rate_update(),
198 // by sending a control message to the reactor
199 // (TODO(relay): not yet implemented)\
200 return Poll::Pending;
201 }
202
203 let mut streams = streams.lock().expect("lock poisoned");
204 let Some((sid, msg)) = streams.poll_ready_streams_iter(cx).next() else {
205 // No ready streams
206 //
207 // TODO(flushing): if there are no ready Tor streams, we might want to defer
208 // flushing until stream data becomes available (or until a timeout elapses).
209 // The deferred flushing approach should enable us to send
210 // more than one message at a time to the channel reactor.
211 return Poll::Pending;
212 };
213
214 if msg.is_none() {
215 // This means the local sender has been dropped,
216 // which presumably can only happen if an error occurs,
217 // or if the Tor stream ends. In both cases, we're going to
218 // want to send an END to the client to let them know,
219 // and to remove the stream from the stream map.
220 //
221 // TODO(relay): the local sender part is not implemented yet
222 return Poll::Ready(StreamEvent::ApplicationStreamClosed(sid));
223 };
224
225 let msg = streams.take_ready_msg(sid).expect("msg disappeared");
226
227 Poll::Ready(StreamEvent::ReadyMsg { sid, msg })
228 });
229
230 select_biased! {
231 res = self.cell_rx.next().fuse() => {
232 let Some(cmd) = res else {
233 // The forward reactor has shut down
234 return Err(ReactorError::Shutdown);
235 };
236
237 self.handle_reactor_cmd(cmd).await?;
238 }
239 event = ready_streams_fut.fuse() => {
240 self.handle_stream_event(event).await?;
241 }
242 }
243
244 Ok(())
245 }
246
247 /// Handle a stream message sent to us by the forward reactor.
248 ///
249 /// Delivers the message to its corresponding application stream.
250 async fn handle_reactor_cmd(&mut self, msg: CtrlMsg) -> StdResult<(), ReactorError> {
251 match msg {
252 CtrlMsg::DeliverStreamMsg {
253 sid,
254 msg,
255 cell_counts_toward_windows,
256 } => {
257 self.deliver_message_to_stream(sid, msg, cell_counts_toward_windows)
258 .await
259 }
260 #[cfg(any(feature = "hs-service", feature = "relay"))]
261 CtrlMsg::ClosePendingStream { stream_id, behav } => {
262 self.close_stream(stream_id, behav, streammap::TerminateReason::ExplicitEnd)
263 .await
264 }
265 }
266 }
267
268 /// Deliver `msg` to the specified stream
269 async fn deliver_message_to_stream(
270 &mut self,
271 sid: StreamId,
272 msg: UnparsedRelayMsg,
273 cell_counts_toward_windows: bool,
274 ) -> StdResult<(), ReactorError> {
275 // We need to apply stream-level flow control *before* encoding the message.
276 // May optionally return a message that needs to be sent back to the client.
277 let bwd_msg = self.handle_msg(sid, msg, cell_counts_toward_windows)?;
278
279 if let Some(bwd_msg) = bwd_msg {
280 self.send_msg_to_bwd(bwd_msg).await?;
281 }
282
283 Ok(())
284 }
285
286 /// Handle a RELAY message that has a non-zero stream ID.
287 ///
288 /// A returned message is one that we need to send back to the client.
289 //
290 // TODO(relay): this is very similar to the client impl from
291 // Circuit::handle_in_order_relay_msg()
292 fn handle_msg(
293 &mut self,
294 streamid: StreamId,
295 msg: UnparsedRelayMsg,
296 cell_counts_toward_windows: bool,
297 ) -> StdResult<Option<AnyRelayMsgOuter>, ReactorError> {
298 let cmd = msg.cmd();
299 let possible_proto_violation_err = move |streamid: StreamId| {
300 Error::StreamProto(format!(
301 "Unexpected {cmd:?} message on unknown stream {streamid}"
302 ))
303 };
304 let now = self.time_provider.now();
305
306 // Check if any of our already-open streams want this message
307 let res = self.hop.handle_msg(
308 possible_proto_violation_err,
309 cell_counts_toward_windows,
310 streamid,
311 msg,
312 now,
313 )?;
314
315 // If it was an incoming stream request, we don't need to worry about
316 // sending an XOFF as there's no stream data within this message.
317 if let Some(msg) = res {
318 cfg_if::cfg_if! {
319 if #[cfg(any(feature = "hs-service", feature = "relay"))] {
320 return self.handle_incoming_stream_request(streamid, msg);
321 } else {
322 return Err(
323 Error::CircProto(format!("Cannot handle {} cells on this circuit", msg.cmd())).into(),
324 );
325 }
326 }
327 }
328
329 // We may want to send an XOFF if the incoming buffer is too large.
330 if let Some(cell) = self.hop.maybe_send_xoff(streamid)? {
331 let cell = AnyRelayMsgOuter::new(Some(streamid), cell.into());
332 return Ok(Some(cell));
333 }
334
335 Ok(None)
336 }
337
338 /// A helper for handling incoming stream requests.
339 ///
340 /// Accepts the specified incoming stream request,
341 /// by adding a new entry to our stream map.
342 ///
343 /// Returns the cell we need to send back to the client,
344 /// if an error occurred and the stream cannot be opened.
345 ///
346 /// Returns None if everything went well
347 /// (the CONNECTED response only comes if the external
348 /// consumer of our [Stream](futures::Stream) of incoming Tor streams
349 /// is able to actually establish the connection to the address
350 /// specified in the BEGIN).
351 ///
352 /// Any error returned from this function will shut down the reactor.
353 #[cfg(any(feature = "hs-service", feature = "relay"))]
354 fn handle_incoming_stream_request(
355 &mut self,
356 sid: StreamId,
357 msg: UnparsedRelayMsg,
358 ) -> StdResult<Option<AnyRelayMsgOuter>, ReactorError> {
359 let mut lock = self.incoming.lock().expect("poisoned lock");
360 let Some(handler) = lock.as_mut() else {
361 return Err(Error::CircProto(format!(
362 "Cannot handle {} cells on this circuit",
363 msg.cmd()
364 ))
365 .into());
366 };
367
368 if self.hopnum != handler.hop_num {
369 let expected_hopnum = match handler.hop_num {
370 Some(hopnum) => hopnum.display().to_string(),
371 None => "client".to_string(),
372 };
373
374 let actual_hopnum = match self.hopnum {
375 Some(hopnum) => hopnum.display().to_string(),
376 None => "None".to_string(),
377 };
378
379 return Err(Error::CircProto(format!(
380 "Expecting incoming streams from {}, but received {} cell from unexpected hop {}",
381 expected_hopnum,
382 msg.cmd(),
383 actual_hopnum,
384 ))
385 .into());
386 }
387
388 let message_closes_stream = handler.cmd_checker.check_msg(&msg)? == StreamStatus::Closed;
389
390 if message_closes_stream {
391 self.hop
392 .stream_map()
393 .lock()
394 .expect("poisoned lock")
395 .ending_msg_received(sid)?;
396
397 return Ok(None);
398 }
399
400 let req = parse_incoming_stream_req(msg)?;
401 let view = CircHopSyncView::new(&self.hop);
402
403 if let Some(reject) = Self::should_reject_incoming(handler, sid, &req, &view)? {
404 // We can't honor this request, so we bail by sending an END.
405 return Ok(Some(reject));
406 };
407
408 let memquota =
409 StreamAccount::new(&self.memquota).map_err(|e| ReactorError::Err(e.into()))?;
410
411 let cmd_checker = InboundDataCmdChecker::new_connected();
412 let stream_components = self.hop.add_ent_with_id(
413 &self.time_provider,
414 sid,
415 cmd_checker,
416 self.inner.flowctrl_sidechannel_mitigations(),
417 &memquota,
418 )?;
419
420 let outcome = Pin::new(&mut handler.incoming_sender).try_send(StreamReqInfo {
421 req,
422 stream_id: sid,
423 hop: None,
424 stream_components,
425 memquota,
426 relay_cell_format: self.hop.relay_cell_format(),
427 });
428
429 log_ratelim!("Delivering message to incoming stream handler"; outcome);
430
431 if let Err(e) = outcome {
432 if e.is_full() {
433 // The IncomingStreamRequestHandler's stream is full; it isn't
434 // handling requests fast enough. So instead, we reply with an
435 // END cell.
436 let end_msg = AnyRelayMsgOuter::new(
437 Some(sid),
438 End::new_with_reason(EndReason::RESOURCELIMIT).into(),
439 );
440
441 return Ok(Some(end_msg));
442 } else if e.is_disconnected() {
443 // The IncomingStreamRequestHandler's stream has been dropped.
444 // In the Tor protocol as it stands, this always means that the
445 // circuit itself is out-of-use and should be closed.
446 //
447 // Note that we will _not_ reach this point immediately after
448 // the IncomingStreamRequestHandler is dropped; we won't hit it
449 // until we next get an incoming request. Thus, if we later
450 // want to add early detection for a dropped
451 // IncomingStreamRequestHandler, we need to do it elsewhere, in
452 // a different way.
453 debug!(
454 circ_uniq_id = %self.unique_id,
455 backward_circ_id = %self.circ_id,
456 "Incoming stream request receiver dropped",
457 );
458 // This will _cause_ the circuit to get closed.
459 return Err(ReactorError::Err(Error::CircuitClosed));
460 } else {
461 // There are no errors like this with the current design of
462 // futures::mpsc, but we shouldn't just ignore the possibility
463 // that they'll be added later.
464 return Err(
465 Error::from((into_internal!("try_send failed unexpectedly"))(e)).into(),
466 );
467 }
468 }
469
470 Ok(None)
471 }
472
473 /// Check if we should reject this incoming stream request or not.
474 ///
475 /// Returns a cell we need to send back to the client if we must reject the request,
476 /// or `None` if we are allowed to accept it.
477 ///`
478 /// Any error returned from this function will shut down the reactor.
479 #[cfg(any(feature = "hs-service", feature = "relay"))]
480 fn should_reject_incoming<'a>(
481 handler: &mut IncomingStreamRequestHandler,
482 sid: StreamId,
483 request: &IncomingStreamRequest,
484 view: &CircHopSyncView<'a>,
485 ) -> StdResult<Option<AnyRelayMsgOuter>, ReactorError> {
486 use IncomingStreamRequestDisposition::*;
487
488 let ctx = IncomingStreamRequestContext { request };
489
490 // Run the externally provided filter to check if we should
491 // open the stream or not.
492 match handler.filter.as_mut().disposition(&ctx, view)? {
493 Accept => {
494 // All is well, we can accept the stream request
495 Ok(None)
496 }
497 CloseCircuit => Err(ReactorError::Shutdown),
498 RejectRequest(end) => {
499 let end_msg = AnyRelayMsgOuter::new(Some(sid), end.into());
500
501 Ok(Some(end_msg))
502 }
503 }
504 }
505
506 /// Handle a [`StreamEvent`].
507 async fn handle_stream_event(&mut self, event: StreamEvent) -> StdResult<(), ReactorError> {
508 match event {
509 StreamEvent::ApplicationStreamClosed(sid) => {
510 self.close_stream(
511 sid,
512 CloseStreamBehavior::default(),
513 streammap::TerminateReason::StreamTargetClosed,
514 )
515 .await
516 }
517 StreamEvent::ReadyMsg { sid, msg } => {
518 self.send_msg_to_bwd(AnyRelayMsgOuter::new(Some(sid), msg))
519 .await
520 }
521 }
522 }
523
524 /// Close the stream that has the specified `sid`.
525 ///
526 /// The `behav` controls whether an `END` will be sent or not.
527 ///
528 /// This calls [`CircHopOutbound::close_stream`] under the hood,
529 /// which removes the stream from the stream map,
530 /// and returns an optional `END` cell to send back to the other party.
531 async fn close_stream(
532 &mut self,
533 sid: StreamId,
534 behav: CloseStreamBehavior,
535 reason: streammap::TerminateReason,
536 ) -> StdResult<(), ReactorError> {
537 let timeout = self.inner.halfstream_expiry(&self.hop);
538 let expire_at = self.time_provider.now() + timeout;
539 let res = self.hop.close_stream(
540 self.unique_id,
541 self.circ_id,
542 sid,
543 None,
544 behav,
545 reason,
546 expire_at,
547 )?;
548 let Some(msg) = res else {
549 // We may not need to send anything at all...
550 return Ok(());
551 };
552
553 self.send_msg_to_bwd(msg.cell).await
554 }
555
556 /// Wrap `msg` in [`ReadyStreamMsg`], and send it to the backward reactor.
557 async fn send_msg_to_bwd(&mut self, msg: AnyRelayMsgOuter) -> StdResult<(), ReactorError> {
558 // TODO(DEDUP): this contains parts of Circuit::send_relay_cell_inner()
559
560 // We might be out of capacity entirely; see if we are about to hit a limit.
561 //
562 // TODO: If we ever add a notion of _recoverable_ errors below, we'll
563 // need a way to restore this limit, and similarly for about_to_send().
564 self.hop.decrement_cell_limit()?;
565
566 // We need to apply stream-level flow control *before* encoding the message
567 // (the BWD handles the encoding)
568 if sendme::cmd_counts_towards_windows(msg.cmd()) {
569 if let Some(stream_id) = msg.stream_id() {
570 self.hop
571 .about_to_send(self.unique_id, self.circ_id, stream_id, msg.msg())?;
572 }
573 }
574
575 // NOTE: on the client side, we call note_data_sent()
576 // just before writing the cell to the channel.
577 // We can't do that here, because we're not the ones
578 // encoding the cell, so we don't have the SENDME tag
579 // which is needed for note_data_sent().
580 //
581 // Instead, we notify the CC algorithm in the BWD,
582 // right after we've finished sending the cell.
583
584 let msg = ReadyStreamMsg {
585 hop: self.hopnum,
586 relay_cell_format: self.hop.relay_cell_format(),
587 ccontrol: Arc::clone(self.hop.ccontrol()),
588 msg,
589 };
590
591 self.bwd_tx
592 .send(msg)
593 .await
594 .map_err(|_| ReactorError::Shutdown)?;
595
596 Ok(())
597 }
598}
599
600/// A Tor stream-related event.
601enum StreamEvent {
602 /// An application stream was closed.
603 ///
604 /// The corresponding entry needs to be removed from the reactor's stream map.
605 ApplicationStreamClosed(StreamId),
606 /// A stream has a ready message.
607 ReadyMsg {
608 /// The ID of the stream to close.
609 sid: StreamId,
610 /// The message.
611 msg: AnyRelayMsg,
612 },
613}
614
615/// Convert an incoming stream request message (BEGIN, BEGIN_DIR, RESOLVE, etc.)
616/// to an [`IncomingStreamRequest`]
617///
618// TODO(dedup): when we rewrite the client reactor in the multi-reactor register,
619// we should rethink this part a bit: ideally, onion services shouldn't even
620// try to parse BEGIN_DIR, RESOLVE.
621//
622// We will likely need an implementation-specific hook for this,
623// similar to the `{Forward,Backward}Handler` implementation-specific handlers
624// we have for the FWD and BWD reactors.
625//
626// See https://gitlab.torproject.org/tpo/core/arti/-/merge_requests/4188#note_3432579
627#[cfg(any(feature = "hs-service", feature = "relay"))]
628fn parse_incoming_stream_req(msg: UnparsedRelayMsg) -> crate::Result<IncomingStreamRequest> {
629 /// Helper for parsing an incoming stream request
630 /// (BEGIN, BEGIN_DIR, or RESOLVE)
631 macro_rules! parse_stream_req {
632 ($msg:expr, $type:tt) => {{
633 let req = $msg
634 .decode::<$type>()
635 .map_err(|e| {
636 Error::from_bytes_err(e, concat!("Invalid ", stringify!($type), " message"))
637 })?
638 .into_msg();
639
640 IncomingStreamRequest::$type(req)
641 }};
642 }
643
644 let req = match msg.cmd() {
645 RelayCmd::BEGIN => parse_stream_req!(msg, Begin),
646 RelayCmd::BEGIN_DIR => parse_stream_req!(msg, BeginDir),
647 RelayCmd::RESOLVE => parse_stream_req!(msg, Resolve),
648 cmd => {
649 // It's a bug if we reach this point, because CircHopOutbound::handle_msg()
650 // should have consumed the message (by forwarding it to the appropriate stream
651 // in its stream map)
652 return Err(internal!("{cmd} is not an incoming stream request").into());
653 }
654 };
655
656 Ok(req)
657}
658
659/// A stream message to be sent to the backward reactor for delivery.
660pub(crate) struct ReadyStreamMsg {
661 /// The hop number, or `None` if we are a relay.
662 pub(crate) hop: Option<HopNum>,
663 /// The message to send.
664 pub(crate) msg: AnyRelayMsgOuter,
665 /// The cell format used with the hop the message should be sent to.
666 pub(crate) relay_cell_format: RelayCellFormat,
667 /// The CC object to use.
668 pub(crate) ccontrol: Arc<Mutex<CongestionControl>>,
669}
670
671/// A control message
672/// that needs to be handled by [`StreamReactor`].
673pub(crate) enum CtrlMsg {
674 /// Stream data received from the other endpoint
675 /// that needs to be delivered to a Tor stream
676 DeliverStreamMsg {
677 /// The ID of the stream this message is for.
678 sid: StreamId,
679 /// The message.
680 msg: UnparsedRelayMsg,
681 /// Whether the cell this message came from counts towards flow-control windows.
682 cell_counts_toward_windows: bool,
683 },
684
685 /// Close the specified pending incoming stream, sending the provided END message.
686 #[cfg(any(feature = "hs-service", feature = "relay"))]
687 ClosePendingStream {
688 /// The stream ID to send the END for.
689 stream_id: StreamId,
690 /// The END message to send, if any.
691 behav: CloseStreamBehavior,
692 },
693}