1pub(crate) mod circhop;
4pub(super) mod extender;
5
6use crate::channel::Channel;
7use crate::circuit::cell_sender::CircuitCellSender;
8use crate::circuit::celltypes::CreateResponse;
9use crate::circuit::circhop::{HopSettings, ReactorStreamComponents};
10use crate::circuit::create::{Create2Wrap, CreateFastWrap, CreateHandshakeWrap};
11use crate::circuit::padding::CircPaddingDisposition;
12use crate::circuit::{CircuitRxReceiver, UniqId};
13use crate::client::circuit::handshake::{BoxedClientLayer, HandshakeRole};
14use crate::client::circuit::padding::{
15 self, PaddingController, PaddingEventStream, QueuedCellPaddingInfo,
16};
17use crate::client::circuit::{ClientCircChanMsg, MutableState, path};
18use crate::client::reactor::MetaCellDisposition;
19use crate::congestion::CongestionSignals;
20use crate::congestion::sendme;
21use crate::crypto::binding::CircuitBinding;
22use crate::crypto::cell::{
23 HopNum, InboundClientCrypt, InboundClientLayer, OutboundClientCrypt, OutboundClientLayer,
24 RelayCellBody,
25};
26use crate::crypto::handshake::fast::CreateFastClient;
27use crate::crypto::handshake::ntor::{NtorClient, NtorPublicKey};
28use crate::crypto::handshake::ntor_v3::{NtorV3Client, NtorV3PublicKey};
29use crate::crypto::handshake::{ClientHandshake, KeyGenerator};
30use crate::memquota::{CircuitAccount, SpecificAccount as _, StreamAccount};
31use crate::stream::cmdcheck::{AnyCmdChecker, StreamStatus};
32use crate::stream::flow_ctrl::state::WithSidechannelMitigations;
33use crate::stream::msg_streamid;
34use crate::streammap;
35use crate::tunnel::TunnelScopedCircId;
36use crate::util::err::ReactorError;
37use crate::util::timeout::TimeoutEstimator;
38use crate::{ClockSkew, Error, Result};
39
40use tor_async_utils::{SinkTrySend as _, SinkTrySendError as _};
41use tor_basic_utils::onionperf_types::{OnionperfCircuitStatus, OnionperfEvent};
42use tor_cell::chancell::msg::{AnyChanMsg, HandshakeType, Relay};
43use tor_cell::chancell::{AnyChanCell, ChanCmd, CircId};
44use tor_cell::chancell::{BoxedCellBody, ChanMsg};
45use tor_cell::relaycell::msg::{AnyRelayMsg, End, Sendme, SendmeTag, Truncated};
46use tor_cell::relaycell::{
47 AnyRelayMsgOuter, RelayCellDecoderResult, RelayCellFormat, RelayCmd, StreamId, UnparsedRelayMsg,
48};
49use tor_error::{Bug, internal};
50use tor_linkspec::RelayIds;
51use tor_llcrypto::pk;
52use web_time_compat::{Duration, Instant, SystemTime};
53
54use futures::SinkExt as _;
55use oneshot_fused_workaround as oneshot;
56use tor_rtcompat::{DynTimeProvider, SleepProvider as _};
57use tracing::{debug, instrument, trace, warn};
58
59use super::{
60 CellHandlers, CircuitHandshake, CloseStreamBehavior, ReactorResultChannel, SendRelayCell,
61};
62
63use crate::conflux::msghandler::ConfluxStatus;
64
65use std::borrow::Borrow;
66use std::pin::Pin;
67use std::result::Result as StdResult;
68use std::sync::Arc;
69
70use extender::HandshakeAuxDataHandler;
71
72#[cfg(feature = "hs-service")]
73use {
74 crate::circuit::CircHopSyncView,
75 crate::stream::{InboundDataCmdChecker, IncomingStreamRequest},
76 tor_cell::relaycell::msg::Begin,
77};
78
79#[cfg(feature = "conflux")]
80use {
81 crate::conflux::msghandler::{ConfluxAction, ConfluxCmd, ConfluxMsgHandler, OooRelayMsg},
82 crate::tunnel::TunnelId,
83};
84
85pub(super) use circhop::{CircHop, CircHopList};
86
87pub(crate) struct Circuit {
92 runtime: DynTimeProvider,
94 channel: Arc<Channel>,
96 pub(super) chan_sender: CircuitCellSender,
101 pub(super) input: CircuitRxReceiver,
106 crypto_in: InboundClientCrypt,
110 crypto_out: OutboundClientCrypt,
112 pub(super) hops: CircHopList,
114 mutable: Arc<MutableState>,
117 circ_id: CircId,
119 unique_id: TunnelScopedCircId,
121 #[cfg(feature = "conflux")]
126 conflux_handler: Option<ConfluxMsgHandler>,
127 padding_ctrl: PaddingController,
129 pub(super) padding_event_stream: PaddingEventStream,
137 #[cfg(feature = "circ-padding")]
139 padding_block: Option<padding::StartBlocking>,
140 timeouts: Arc<dyn TimeoutEstimator>,
144 #[allow(dead_code)] memquota: CircuitAccount,
147}
148
149#[derive(Debug, derive_more::From)]
157pub(super) enum CircuitCmd {
158 Send(SendRelayCell),
160 HandleSendMe {
162 hop: HopNum,
164 sendme: Sendme,
166 },
167 CloseStream {
169 hop: HopNum,
171 sid: StreamId,
173 behav: CloseStreamBehavior,
175 reason: streammap::TerminateReason,
177 },
178 #[cfg(feature = "conflux")]
180 Conflux(ConfluxCmd),
181 CleanShutdown,
183 #[cfg(feature = "conflux")]
185 Enqueue(OooRelayMsg),
186}
187
188macro_rules! unsupported_client_cell {
195 ($msg:expr) => {{
196 unsupported_client_cell!(@ $msg, "")
197 }};
198
199 ($msg:expr, $hopnum:expr) => {{
200 let hop: HopNum = $hopnum;
201 let hop_display = format!(" from hop {}", hop.display());
202 unsupported_client_cell!(@ $msg, hop_display)
203 }};
204
205 (@ $msg:expr, $hopnum_display:expr) => {
206 Err(crate::Error::CircProto(format!(
207 "Unexpected {} cell{} on client circuit",
208 $msg.cmd(),
209 $hopnum_display,
210 )))
211 };
212}
213
214pub(super) use unsupported_client_cell;
215
216impl Circuit {
217 #[allow(clippy::too_many_arguments)]
219 pub(super) fn new(
220 runtime: DynTimeProvider,
221 channel: Arc<Channel>,
222 circ_id: CircId,
223 unique_id: TunnelScopedCircId,
224 input: CircuitRxReceiver,
225 memquota: CircuitAccount,
226 mutable: Arc<MutableState>,
227 padding_ctrl: PaddingController,
228 padding_event_stream: PaddingEventStream,
229 timeouts: Arc<dyn TimeoutEstimator>,
230 ) -> Self {
231 let chan_sender = CircuitCellSender::from_channel_sender(channel.sender());
232
233 let crypto_out = OutboundClientCrypt::new();
234 Circuit {
235 runtime,
236 channel,
237 chan_sender,
238 input,
239 crypto_in: InboundClientCrypt::new(),
240 hops: CircHopList::default(),
241 unique_id,
242 circ_id,
243 crypto_out,
244 mutable,
245 #[cfg(feature = "conflux")]
246 conflux_handler: None,
247 padding_ctrl,
248 padding_event_stream,
249 #[cfg(feature = "circ-padding")]
250 padding_block: None,
251 timeouts,
252 memquota,
253 }
254 }
255
256 pub(super) fn unique_id(&self) -> UniqId {
258 self.unique_id.unique_id()
259 }
260
261 pub(super) fn circ_id(&self) -> CircId {
263 self.circ_id
264 }
265
266 pub(super) fn mutable(&self) -> &Arc<MutableState> {
268 &self.mutable
269 }
270
271 #[cfg(feature = "conflux")]
276 pub(super) fn add_to_conflux_tunnel(
277 &mut self,
278 tunnel_id: TunnelId,
279 conflux_handler: ConfluxMsgHandler,
280 ) {
281 self.unique_id = TunnelScopedCircId::new(tunnel_id, self.unique_id.unique_id());
282 self.conflux_handler = Some(conflux_handler);
283 }
284
285 #[cfg(feature = "conflux")]
290 pub(super) async fn begin_conflux_link(
291 &mut self,
292 hop: HopNum,
293 cell: AnyRelayMsgOuter,
294 runtime: &tor_rtcompat::DynTimeProvider,
295 ) -> Result<()> {
296 use tor_rtcompat::SleepProvider as _;
297
298 if self.conflux_handler.is_none() {
299 return Err(internal!(
300 "tried to send LINK cell before installing a ConfluxMsgHandler?!"
301 )
302 .into());
303 }
304
305 let cell = SendRelayCell {
306 hop: Some(hop),
307 early: false,
308 cell,
309 };
310 self.send_relay_cell(cell).await?;
311
312 let Some(conflux_handler) = self.conflux_handler.as_mut() else {
313 return Err(internal!("ConfluxMsgHandler disappeared?!").into());
314 };
315
316 Ok(conflux_handler.note_link_sent(runtime.wallclock())?)
317 }
318
319 pub(super) fn conflux_hs_timeout(&self) -> Option<SystemTime> {
323 cfg_if::cfg_if! {
324 if #[cfg(feature = "conflux")] {
325 self.conflux_handler.as_ref().map(|handler| handler.handshake_timeout())?
326 } else {
327 None
328 }
329 }
330 }
331
332 #[cfg(test)]
334 pub(super) fn handle_add_fake_hop(
335 &mut self,
336 format: RelayCellFormat,
337 fwd_lasthop: bool,
338 rev_lasthop: bool,
339 dummy_peer_id: path::HopDetail,
340 params: &crate::client::circuit::CircParameters,
344 done: ReactorResultChannel<()>,
345 ) {
346 use tor_protover::{Protocols, named};
347
348 use crate::client::circuit::test::DummyCrypto;
349
350 assert!(matches!(format, RelayCellFormat::V0));
351 let _ = format; let fwd = Box::new(DummyCrypto::new(fwd_lasthop));
354 let rev = Box::new(DummyCrypto::new(rev_lasthop));
355 let binding = None;
356
357 let settings = HopSettings::from_params_and_caps(
358 crate::circuit::circhop::HopNegotiationType::Full,
360 params,
361 &[named::FLOWCTRL_CC].into_iter().collect::<Protocols>(),
362 )
363 .expect("Can't construct HopSettings");
364 self.add_hop(dummy_peer_id, fwd, rev, binding, &settings)
365 .expect("could not add hop to circuit");
366 let _ = done.send(Ok(()));
367 }
368
369 fn encode_relay_cell(
373 crypto_out: &mut OutboundClientCrypt,
374 relay_format: RelayCellFormat,
375 hop: HopNum,
376 early: bool,
377 msg: AnyRelayMsgOuter,
378 ) -> Result<(AnyChanMsg, SendmeTag)> {
379 let mut body: RelayCellBody = msg
380 .encode(relay_format, &mut rand::rng())
381 .map_err(|e| Error::from_cell_enc(e, "relay cell body"))?
382 .into();
383 let cmd = if early {
384 ChanCmd::RELAY_EARLY
385 } else {
386 ChanCmd::RELAY
387 };
388 let tag = crypto_out.encrypt(cmd, &mut body, hop)?;
389 let msg = Relay::from(BoxedCellBody::from(body));
390 let msg = if early {
391 AnyChanMsg::RelayEarly(msg.into())
392 } else {
393 AnyChanMsg::Relay(msg)
394 };
395
396 Ok((msg, tag))
397 }
398
399 #[instrument(level = "trace", skip_all)]
410 pub(super) async fn send_relay_cell(&mut self, msg: SendRelayCell) -> Result<()> {
411 self.send_relay_cell_inner(msg, None).await
412 }
413
414 #[instrument(level = "trace", skip_all)]
420 async fn send_relay_cell_inner(
421 &mut self,
422 msg: SendRelayCell,
423 padding_info: Option<QueuedCellPaddingInfo>,
424 ) -> Result<()> {
425 let SendRelayCell {
426 hop,
427 early,
428 cell: msg,
429 } = msg;
430
431 let is_conflux_link = msg.cmd() == RelayCmd::CONFLUX_LINK;
432 if !is_conflux_link && self.is_conflux_pending() {
433 return Err(internal!("tried to send cell on unlinked circuit").into());
436 }
437
438 trace!(
439 circ_uniq_id = %self.unique_id,
440 forward_circ_id = %self.circ_id,
441 cell = ?msg,
442 "sending relay cell"
443 );
444
445 let runtime = self.runtime.clone();
447 let c_t_w = sendme::cmd_counts_towards_windows(msg.cmd());
448 let stream_id = msg.stream_id();
449 let hop = hop.expect("missing hop in client SendRelayCell?!");
450 let circhop = self.hops.get_mut(hop).ok_or(Error::NoSuchHop)?;
451
452 circhop.decrement_outbound_cell_limit()?;
457
458 if c_t_w {
460 if let Some(stream_id) = stream_id {
461 circhop.about_to_send(stream_id, msg.msg())?;
462 }
463 }
464
465 let relay_cmd = msg.cmd();
469
470 let (msg, tag) = Self::encode_relay_cell(
473 &mut self.crypto_out,
474 circhop.relay_cell_format(),
475 hop,
476 early,
477 msg,
478 )?;
479 if c_t_w {
482 circhop.ccontrol().note_data_sent(&runtime, &tag)?;
483 }
484
485 let padding_info = padding_info.or_else(|| self.padding_ctrl.queued_data(hop));
487
488 self.send_msg(msg, padding_info).await?;
489
490 #[cfg(feature = "conflux")]
491 if let Some(conflux) = self.conflux_handler.as_mut() {
492 conflux.note_cell_sent(relay_cmd);
493 }
494
495 Ok(())
496 }
497
498 pub(super) fn handle_cell(
511 &mut self,
512 handlers: &mut CellHandlers,
513 leg: UniqId,
514 cell: ClientCircChanMsg,
515 ) -> Result<Vec<CircuitCmd>> {
516 trace!(
517 circ_uniq_id = %self.unique_id,
518 forward_circ_id = %self.circ_id,
519 cell = ?cell,
520 "handling cell"
521 );
522 use ClientCircChanMsg::*;
523 match cell {
524 Relay(r) => self.handle_relay_cell(handlers, leg, r),
525 Destroy(d) => {
526 let reason = d.reason();
527 debug!(
528 circ_uniq_id = %self.unique_id,
529 forward_circ_id = %self.circ_id,
530 "Received DESTROY cell. Reason: {} [{}]",
531 reason.human_str(),
532 reason
533 );
534
535 self.handle_destroy_cell().map(|c| vec![c])
536 }
537 }
538 }
539
540 fn decode_relay_cell(
543 &mut self,
544 cell: Relay,
545 ) -> Result<(HopNum, SendmeTag, RelayCellDecoderResult)> {
546 let cmd = cell.cmd();
548 let mut body = cell.into_relay_body().into();
549
550 let (hopnum, tag) = self.crypto_in.decrypt(cmd, &mut body)?;
553
554 let decode_res = self
556 .hop_mut(hopnum)
557 .ok_or_else(|| {
558 Error::from(internal!(
559 "Trying to decode cell from nonexistent hop {:?}",
560 hopnum
561 ))
562 })?
563 .decode(body.into())?;
564
565 Ok((hopnum, tag, decode_res))
566 }
567
568 fn handle_relay_cell(
570 &mut self,
571 handlers: &mut CellHandlers,
572 leg: UniqId,
573 cell: Relay,
574 ) -> Result<Vec<CircuitCmd>> {
575 let (hopnum, tag, decode_res) = self.decode_relay_cell(cell)?;
576
577 if decode_res.is_padding() {
578 self.padding_ctrl.decrypted_padding(hopnum)?;
579 } else {
580 self.padding_ctrl.decrypted_data(hopnum);
581 }
582
583 self.hop_mut(hopnum)
585 .ok_or_else(|| internal!("nonexistent hop {:?}", hopnum))?
586 .decrement_inbound_cell_limit()?;
587
588 let c_t_w = decode_res.cmds().any(sendme::cmd_counts_towards_windows);
589
590 let send_circ_sendme = if c_t_w {
593 self.hop_mut(hopnum)
594 .ok_or_else(|| Error::CircProto("Sendme from nonexistent hop".into()))?
595 .ccontrol()
596 .note_data_received()?
597 } else {
598 false
599 };
600
601 let mut circ_cmds = vec![];
602 if send_circ_sendme {
604 let sendme = Sendme::from(tag);
609 let cell = AnyRelayMsgOuter::new(None, sendme.into());
610 circ_cmds.push(CircuitCmd::Send(SendRelayCell {
611 hop: Some(hopnum),
612 early: false,
613 cell,
614 }));
615
616 self.hop_mut(hopnum)
618 .ok_or_else(|| {
619 Error::from(internal!(
620 "Trying to send SENDME to nonexistent hop {:?}",
621 hopnum
622 ))
623 })?
624 .ccontrol()
625 .note_sendme_sent()?;
626 }
627
628 let (mut msgs, incomplete) = decode_res.into_parts();
629 while let Some(msg) = msgs.next() {
630 let msg_status = self.handle_relay_msg(handlers, hopnum, leg, c_t_w, msg)?;
631
632 match msg_status {
633 None => continue,
634 Some(msg @ CircuitCmd::CleanShutdown) => {
635 for m in msgs {
636 debug!(
637 "{id}: Ignoring relay msg received after triggering shutdown: {m:?}",
638 id = self.unique_id
639 );
640 }
641 if let Some(incomplete) = incomplete {
642 debug!(
643 "{id}: Ignoring partial relay msg received after triggering shutdown: {:?}",
644 incomplete,
645 id = self.unique_id,
646 );
647 }
648 circ_cmds.push(msg);
649 return Ok(circ_cmds);
650 }
651 Some(msg) => {
652 circ_cmds.push(msg);
653 }
654 }
655 }
656
657 Ok(circ_cmds)
658 }
659
660 fn handle_relay_msg(
662 &mut self,
663 handlers: &mut CellHandlers,
664 hopnum: HopNum,
665 leg: UniqId,
666 cell_counts_toward_windows: bool,
667 msg: UnparsedRelayMsg,
668 ) -> Result<Option<CircuitCmd>> {
669 let streamid = msg_streamid(&msg)?;
672
673 let Some(streamid) = streamid else {
676 return self.handle_meta_cell(handlers, hopnum, msg);
677 };
678
679 #[cfg(feature = "conflux")]
680 let msg = if let Some(conflux) = self.conflux_handler.as_mut() {
681 match conflux.action_for_msg(hopnum, cell_counts_toward_windows, streamid, msg)? {
682 ConfluxAction::Deliver(msg) => {
683 msg
690 }
691 ConfluxAction::Enqueue(msg) => {
692 return Ok(Some(CircuitCmd::Enqueue(msg)));
694 }
695 }
696 } else {
697 msg
700 };
701
702 self.handle_in_order_relay_msg(
703 handlers,
704 hopnum,
705 leg,
706 cell_counts_toward_windows,
707 streamid,
708 msg,
709 )
710 }
711
712 pub(super) fn handle_in_order_relay_msg(
714 &mut self,
715 handlers: &mut CellHandlers,
716 hopnum: HopNum,
717 leg: UniqId,
718 cell_counts_toward_windows: bool,
719 streamid: StreamId,
720 msg: UnparsedRelayMsg,
721 ) -> Result<Option<CircuitCmd>> {
722 let now = self.runtime.now();
723
724 #[cfg(feature = "conflux")]
725 if let Some(conflux) = self.conflux_handler.as_mut() {
726 conflux.inc_last_seq_delivered(&msg);
727 }
728
729 let path = self.mutable.path();
730
731 let nonexistent_hop_err = || Error::CircProto("Cell from nonexistent hop!".into());
732 let hop = self.hop_mut(hopnum).ok_or_else(nonexistent_hop_err)?;
733
734 let hop_detail = path
735 .iter()
736 .nth(usize::from(hopnum))
737 .ok_or_else(nonexistent_hop_err)?;
738
739 let res = hop.handle_msg(hop_detail, cell_counts_toward_windows, streamid, msg, now)?;
742
743 if let Some(msg) = res {
746 cfg_if::cfg_if! {
747 if #[cfg(feature = "hs-service")] {
748 return self.handle_incoming_stream_request(
749 handlers,
750 msg,
751 streamid,
752 hopnum,
753 leg,
754 WithSidechannelMitigations::Enabled,
757 );
758 } else {
759 return Err(
760 Error::CircProto(format!("Cannot handle {} cells on this circuit", msg.cmd())),
761 );
762 }
763 }
764 }
765
766 if let Some(cell) = hop.maybe_send_xoff(streamid)? {
768 let cell = AnyRelayMsgOuter::new(Some(streamid), cell.into());
769 let cell = SendRelayCell {
770 hop: Some(hopnum),
771 early: false,
772 cell,
773 };
774 return Ok(Some(CircuitCmd::Send(cell)));
775 }
776
777 Ok(None)
778 }
779
780 #[cfg(feature = "conflux")]
788 fn handle_conflux_msg(
789 &mut self,
790 hop: HopNum,
791 msg: UnparsedRelayMsg,
792 ) -> Result<Option<ConfluxCmd>> {
793 let Some(conflux_handler) = self.conflux_handler.as_mut() else {
794 return Err(Error::CircProto(format!(
797 "Received {} cell from hop {} on non-conflux client circuit?!",
798 msg.cmd(),
799 hop.display(),
800 )));
801 };
802
803 Ok(conflux_handler.handle_conflux_msg(msg, hop))
804 }
805
806 #[cfg(feature = "conflux")]
810 pub(super) fn last_seq_sent(&self) -> Result<u64> {
811 let handler = self
812 .conflux_handler
813 .as_ref()
814 .ok_or_else(|| internal!("tried to get last_seq_sent of non-conflux circ"))?;
815
816 Ok(handler.last_seq_sent())
817 }
818
819 #[cfg(feature = "conflux")]
823 pub(super) fn set_last_seq_sent(&mut self, n: u64) -> Result<()> {
824 let handler = self
825 .conflux_handler
826 .as_mut()
827 .ok_or_else(|| internal!("tried to get last_seq_sent of non-conflux circ"))?;
828
829 handler.set_last_seq_sent(n);
830 Ok(())
831 }
832
833 #[cfg(feature = "conflux")]
837 pub(super) fn last_seq_recv(&self) -> Result<u64> {
838 let handler = self
839 .conflux_handler
840 .as_ref()
841 .ok_or_else(|| internal!("tried to get last_seq_recv of non-conflux circ"))?;
842
843 Ok(handler.last_seq_recv())
844 }
845
846 #[cfg(feature = "hs-service")]
850 fn handle_incoming_stream_request(
851 &mut self,
852 handlers: &mut CellHandlers,
853 msg: UnparsedRelayMsg,
854 stream_id: StreamId,
855 hop_num: HopNum,
856 leg: UniqId,
857 with_sidechannel_mitigations: WithSidechannelMitigations,
858 ) -> Result<Option<CircuitCmd>> {
859 use tor_cell::relaycell::msg::EndReason;
860 use tor_error::into_internal;
861 use tor_log_ratelim::log_ratelim;
862
863 use crate::stream::incoming::StreamReqInfo;
864
865 let Some(handler) = handlers.incoming_stream_req_handler.as_mut() else {
868 return Err(Error::CircProto(
869 "Cannot handle BEGIN cells on this circuit".into(),
870 ));
871 };
872
873 let expected_hop_num = handler
875 .hop_num
876 .ok_or_else(|| internal!("Handler HopNum is None in client impl?!"))?;
877
878 if hop_num != expected_hop_num {
879 return Err(Error::CircProto(format!(
880 "Expecting incoming streams from {}, but received {} cell from unexpected hop {}",
881 expected_hop_num.display(),
882 msg.cmd(),
883 hop_num.display()
884 )));
885 }
886
887 let message_closes_stream = handler.cmd_checker.check_msg(&msg)? == StreamStatus::Closed;
888
889 let hop = self.hops.get_mut(hop_num).ok_or(Error::CircuitClosed)?;
897
898 if message_closes_stream {
899 hop.ending_msg_received(stream_id)?;
900
901 return Ok(None);
902 }
903
904 let begin = msg
905 .decode::<Begin>()
906 .map_err(|e| Error::from_bytes_err(e, "Invalid Begin message"))?
907 .into_msg();
908
909 let req = IncomingStreamRequest::Begin(begin);
910
911 {
912 use crate::stream::IncomingStreamRequestDisposition::*;
913
914 let ctx = crate::stream::IncomingStreamRequestContext { request: &req };
915 let view = CircHopSyncView::new(hop.outbound());
921
922 match handler.filter.as_mut().disposition(&ctx, &view)? {
923 Accept => {}
924 CloseCircuit => return Ok(Some(CircuitCmd::CleanShutdown)),
925 RejectRequest(end) => {
926 let end_msg = AnyRelayMsgOuter::new(Some(stream_id), end.into());
927 let cell = SendRelayCell {
928 hop: Some(hop_num),
929 early: false,
930 cell: end_msg,
931 };
932 return Ok(Some(CircuitCmd::Send(cell)));
933 }
934 }
935 }
936
937 let hop = self.hops.get_mut(hop_num).ok_or(Error::CircuitClosed)?;
940 let relay_cell_format = hop.relay_cell_format();
941
942 let memquota = StreamAccount::new(&self.memquota)?;
943
944 let cmd_checker = InboundDataCmdChecker::new_connected();
945 let stream_components = hop.add_ent_with_id(
946 self.chan_sender.time_provider(),
947 stream_id,
948 cmd_checker,
949 with_sidechannel_mitigations,
950 &memquota,
951 )?;
952
953 let outcome = Pin::new(&mut handler.incoming_sender).try_send(StreamReqInfo {
954 req,
955 stream_id,
956 hop: Some((leg, hop_num).into()),
957 stream_components,
958 memquota,
959 relay_cell_format,
960 });
961
962 log_ratelim!("Delivering message to incoming stream handler"; outcome);
963
964 if let Err(e) = outcome {
965 if e.is_full() {
966 let end_msg = AnyRelayMsgOuter::new(
970 Some(stream_id),
971 End::new_with_reason(EndReason::RESOURCELIMIT).into(),
972 );
973
974 let cell = SendRelayCell {
975 hop: Some(hop_num),
976 early: false,
977 cell: end_msg,
978 };
979 return Ok(Some(CircuitCmd::Send(cell)));
980 } else if e.is_disconnected() {
981 debug!(
993 circ_uniq_id = %self.unique_id,
994 forward_circ_id = %self.circ_id,
995 "Incoming stream request receiver dropped",
996 );
997 return Err(Error::CircuitClosed);
999 } else {
1000 return Err(Error::from((into_internal!(
1004 "try_send failed unexpectedly"
1005 ))(e)));
1006 }
1007 }
1008
1009 Ok(None)
1010 }
1011
1012 #[allow(clippy::unnecessary_wraps)]
1014 fn handle_destroy_cell(&mut self) -> Result<CircuitCmd> {
1015 Ok(CircuitCmd::CleanShutdown)
1017 }
1018
1019 pub(super) async fn handle_create(
1021 &mut self,
1022 recv_created: oneshot::Receiver<CreateResponse>,
1023 handshake: CircuitHandshake,
1024 settings: HopSettings,
1025 done: ReactorResultChannel<()>,
1026 ) -> StdResult<(), ReactorError> {
1027 let ret = match handshake {
1028 CircuitHandshake::CreateFast => self.create_firsthop_fast(recv_created, settings).await,
1029 CircuitHandshake::Ntor {
1030 public_key,
1031 ed_identity,
1032 } => {
1033 self.create_firsthop_ntor(recv_created, ed_identity, public_key, settings)
1034 .await
1035 }
1036 CircuitHandshake::NtorV3 { public_key } => {
1037 self.create_firsthop_ntor_v3(recv_created, public_key, settings)
1038 .await
1039 }
1040 };
1041 let _ = done.send(ret); self.chan_sender.flush().await?;
1046
1047 Ok(())
1048 }
1049
1050 async fn create_impl<H, W, M>(
1056 &mut self,
1057 recvcreated: oneshot::Receiver<CreateResponse>,
1058 wrap: &W,
1059 key: &H::KeyType,
1060 mut settings: HopSettings,
1061 msg: &M,
1062 ) -> Result<()>
1063 where
1064 H: ClientHandshake + HandshakeAuxDataHandler,
1065 W: CreateHandshakeWrap,
1066 H::KeyGen: KeyGenerator,
1067 M: Borrow<H::ClientAuxData>,
1068 {
1069 let (state, msg) = H::client1(&mut rand::rng(), key, msg)?;
1074 let create_cell = wrap.to_chanmsg(msg);
1075 trace!(
1076 circ_uniq_id = %self.unique_id,
1077 forward_circ_id = %self.circ_id,
1078 create = %create_cell.cmd(),
1079 "Extending to hop 1",
1080 );
1081 self.send_msg(create_cell, None).await?;
1082
1083 let reply = recvcreated
1084 .await
1085 .map_err(|_| Error::CircProto("Circuit closed while waiting".into()))?;
1086
1087 let relay_handshake = wrap.decode_chanmsg(reply)?;
1088 let (server_msg, keygen) = H::client2(state, relay_handshake)?;
1089
1090 H::handle_server_aux_data(&mut settings, &server_msg)?;
1091
1092 let BoxedClientLayer { fwd, back, binding } = settings
1093 .relay_crypt_protocol()
1094 .construct_client_layers(HandshakeRole::Initiator, keygen)?;
1095
1096 trace!(
1097 circ_uniq_id = %self.unique_id,
1098 forward_circ_id = %self.circ_id,
1099 "Handshake complete; circuit created."
1100 );
1101
1102 trace!(
1103 onionperf = true,
1104 circ_uniq_id = %self.unique_id,
1105 forward_circ_id = %self.circ_id,
1106 event = ?OnionperfEvent::Circuit(OnionperfCircuitStatus::Extended),
1107 );
1108
1109 let peer_id = self.channel.target().clone();
1110
1111 self.add_hop(
1112 path::HopDetail::Relay(peer_id),
1113 fwd,
1114 back,
1115 binding,
1116 &settings,
1117 )?;
1118 Ok(())
1119 }
1120
1121 async fn create_firsthop_fast(
1128 &mut self,
1129 recvcreated: oneshot::Receiver<CreateResponse>,
1130 settings: HopSettings,
1131 ) -> Result<()> {
1132 let wrap = CreateFastWrap;
1134 self.create_impl::<CreateFastClient, _, _>(recvcreated, &wrap, &(), settings, &())
1135 .await
1136 }
1137
1138 async fn create_firsthop_ntor(
1143 &mut self,
1144 recvcreated: oneshot::Receiver<CreateResponse>,
1145 ed_identity: pk::ed25519::Ed25519Identity,
1146 pubkey: NtorPublicKey,
1147 settings: HopSettings,
1148 ) -> Result<()> {
1149 let target = RelayIds::builder()
1151 .ed_identity(ed_identity)
1152 .rsa_identity(pubkey.id)
1153 .build()
1154 .expect("Unable to build RelayIds");
1155 self.channel.check_match(&target)?;
1156
1157 let wrap = Create2Wrap {
1158 handshake_type: HandshakeType::NTOR,
1159 };
1160 self.create_impl::<NtorClient, _, _>(recvcreated, &wrap, &pubkey, settings, &())
1161 .await
1162 }
1163
1164 async fn create_firsthop_ntor_v3(
1169 &mut self,
1170 recvcreated: oneshot::Receiver<CreateResponse>,
1171 pubkey: NtorV3PublicKey,
1172 settings: HopSettings,
1173 ) -> Result<()> {
1174 let target = RelayIds::builder()
1176 .ed_identity(pubkey.id)
1177 .build()
1178 .expect("Unable to build RelayIds");
1179 self.channel.check_match(&target)?;
1180
1181 let client_extensions = settings.circuit_request_extensions()?;
1183 let wrap = Create2Wrap {
1184 handshake_type: HandshakeType::NTOR_V3,
1185 };
1186
1187 self.create_impl::<NtorV3Client, _, _>(
1188 recvcreated,
1189 &wrap,
1190 &pubkey,
1191 settings,
1192 &client_extensions,
1193 )
1194 .await
1195 }
1196
1197 pub(super) fn add_hop(
1201 &mut self,
1202 peer_id: path::HopDetail,
1203 fwd: Box<dyn OutboundClientLayer + 'static + Send>,
1204 rev: Box<dyn InboundClientLayer + 'static + Send>,
1205 binding: Option<CircuitBinding>,
1206 settings: &HopSettings,
1207 ) -> StdResult<(), Bug> {
1208 let hop_num = self.hops.len();
1209 debug_assert_eq!(hop_num, usize::from(self.num_hops()));
1210
1211 if hop_num == usize::from(u8::MAX) {
1215 return Err(internal!(
1216 "cannot add more hops to a circuit with `u8::MAX` hops"
1217 ));
1218 }
1219
1220 let hop_num = (hop_num as u8).into();
1221
1222 let hop = CircHop::new(self.unique_id, self.circ_id, hop_num, settings);
1223 self.hops.push(hop);
1224 self.crypto_in.add_layer(rev);
1225 self.crypto_out.add_layer(fwd);
1226 self.mutable.add_hop(peer_id, binding);
1227
1228 Ok(())
1229 }
1230
1231 fn handle_meta_cell(
1247 &mut self,
1248 handlers: &mut CellHandlers,
1249 hopnum: HopNum,
1250 msg: UnparsedRelayMsg,
1251 ) -> Result<Option<CircuitCmd>> {
1252 if msg.cmd() == RelayCmd::SENDME {
1262 let sendme = msg
1263 .decode::<Sendme>()
1264 .map_err(|e| Error::from_bytes_err(e, "sendme message"))?
1265 .into_msg();
1266
1267 return Ok(Some(CircuitCmd::HandleSendMe {
1268 hop: hopnum,
1269 sendme,
1270 }));
1271 }
1272 if msg.cmd() == RelayCmd::TRUNCATED {
1273 let truncated = msg
1274 .decode::<Truncated>()
1275 .map_err(|e| Error::from_bytes_err(e, "truncated message"))?
1276 .into_msg();
1277 let reason = truncated.reason();
1278 debug!(
1279 circ_uniq_id = %self.unique_id,
1280 forward_circ_id = %self.circ_id,
1281 "Truncated from hop {}. Reason: {} [{}]",
1282 hopnum.display(),
1283 reason.human_str(),
1284 reason
1285 );
1286
1287 return Ok(Some(CircuitCmd::CleanShutdown));
1288 }
1289
1290 if msg.cmd() == RelayCmd::DROP {
1291 cfg_if::cfg_if! {
1292 if #[cfg(feature = "circ-padding")] {
1293 return Ok(None);
1294 } else {
1295 use crate::util::err::ExcessPadding;
1296 return Err(Error::ExcessPadding(ExcessPadding::NoPaddingNegotiated, hopnum));
1297 }
1298 }
1299 }
1300
1301 trace!(
1302 circ_uniq_id = %self.unique_id,
1303 forward_circ_id = %self.circ_id,
1304 cell = ?msg,
1305 "Received meta-cell"
1306 );
1307
1308 #[cfg(feature = "conflux")]
1309 if matches!(
1310 msg.cmd(),
1311 RelayCmd::CONFLUX_LINK
1312 | RelayCmd::CONFLUX_LINKED
1313 | RelayCmd::CONFLUX_LINKED_ACK
1314 | RelayCmd::CONFLUX_SWITCH
1315 ) {
1316 let cmd = self.handle_conflux_msg(hopnum, msg)?;
1317 return Ok(cmd.map(CircuitCmd::from));
1318 }
1319
1320 if self.is_conflux_pending() {
1321 warn!(
1322 circ_uniq_id = %self.unique_id,
1323 forward_circ_id = %self.circ_id,
1324 "received unexpected cell {msg:?} on unlinked conflux circuit",
1325 );
1326 return Err(Error::CircProto(
1327 "Received unexpected cell on unlinked circuit".into(),
1328 ));
1329 }
1330
1331 if let Some(mut handler) = handlers.meta_handler.take() {
1339 if handler.expected_hop() == (self.unique_id(), hopnum).into() {
1341 let ret = handler.handle_msg(msg, self);
1343 trace!(
1344 circ_uniq_id = %self.unique_id,
1345 forward_circ_id = %self.circ_id,
1346 result = ?ret,
1347 "meta handler completed",
1348 );
1349 match ret {
1350 #[cfg(feature = "send-control-msg")]
1351 Ok(MetaCellDisposition::Consumed) => {
1352 handlers.meta_handler = Some(handler);
1353 Ok(None)
1354 }
1355 Ok(MetaCellDisposition::ConversationFinished) => Ok(None),
1356 #[cfg(feature = "send-control-msg")]
1357 Ok(MetaCellDisposition::CloseCirc) => Ok(Some(CircuitCmd::CleanShutdown)),
1358 Err(e) => Err(e),
1359 }
1360 } else {
1361 handlers.meta_handler = Some(handler);
1364
1365 unsupported_client_cell!(msg, hopnum)
1366 }
1367 } else {
1368 unsupported_client_cell!(msg)
1371 }
1372 }
1373
1374 #[instrument(level = "trace", skip_all)]
1376 pub(super) fn handle_sendme(
1377 &mut self,
1378 hopnum: HopNum,
1379 msg: Sendme,
1380 signals: CongestionSignals,
1381 ) -> Result<Option<CircuitCmd>> {
1382 let runtime = self.runtime.clone();
1384
1385 let hop = self
1388 .hop_mut(hopnum)
1389 .ok_or_else(|| Error::CircProto(format!("Couldn't find hop {}", hopnum.display())))?;
1390
1391 let tag = msg.into_sendme_tag().ok_or_else(||
1392 Error::CircProto("missing tag on circuit sendme".into()))?;
1395 hop.ccontrol()
1397 .note_sendme_received(&runtime, tag, signals)?;
1398 Ok(None)
1399 }
1400
1401 #[instrument(level = "trace", skip_all)]
1416 async fn send_msg(
1417 &mut self,
1418 msg: AnyChanMsg,
1419 info: Option<QueuedCellPaddingInfo>,
1420 ) -> Result<()> {
1421 let cell = AnyChanCell::new(Some(self.circ_id), msg);
1422 Pin::new(&mut self.chan_sender)
1424 .send_unbounded((cell, info))
1425 .await?;
1426 Ok(())
1427 }
1428
1429 pub(super) fn remove_expired_halfstreams(&mut self, now: Instant) {
1431 self.hops.remove_expired_halfstreams(now);
1432 }
1433
1434 pub(super) fn hop(&self, hopnum: HopNum) -> Option<&CircHop> {
1436 self.hops.hop(hopnum)
1437 }
1438
1439 pub(super) fn hop_mut(&mut self, hopnum: HopNum) -> Option<&mut CircHop> {
1441 self.hops.get_mut(hopnum)
1442 }
1443
1444 #[allow(clippy::too_many_arguments)]
1447 pub(super) fn begin_stream(
1448 &mut self,
1449 hop_num: HopNum,
1450 message: AnyRelayMsg,
1451 time_prov: &DynTimeProvider,
1452 cmd_checker: AnyCmdChecker,
1453 memquota: &StreamAccount,
1454 ) -> Result<(SendRelayCell, StreamId, ReactorStreamComponents)> {
1455 let Some(hop) = self.hop_mut(hop_num) else {
1456 return Err(internal!(
1457 "{}: Attempting to send a BEGIN cell to an unknown hop {hop_num:?}",
1458 self.unique_id,
1459 )
1460 .into());
1461 };
1462
1463 hop.begin_stream(message, time_prov, cmd_checker, memquota)
1464 }
1465
1466 #[instrument(level = "trace", skip_all)]
1468 pub(super) async fn close_stream(
1469 &mut self,
1470 hop_num: HopNum,
1471 sid: StreamId,
1472 behav: CloseStreamBehavior,
1473 reason: streammap::TerminateReason,
1474 expiry: Instant,
1475 ) -> Result<()> {
1476 if let Some(hop) = self.hop_mut(hop_num) {
1477 let res = hop.close_stream(sid, behav, reason, expiry)?;
1478 if let Some(cell) = res {
1479 self.send_relay_cell(cell).await?;
1480 }
1481 }
1482 Ok(())
1483 }
1484
1485 pub(super) fn has_streams(&self) -> bool {
1491 self.hops.has_streams()
1492 }
1493
1494 pub(super) fn num_hops(&self) -> u8 {
1496 self.hops
1501 .len()
1502 .try_into()
1503 .expect("`hops.len()` has more than `u8::MAX` hops")
1504 }
1505
1506 pub(super) fn has_hops(&self) -> bool {
1508 !self.hops.is_empty()
1509 }
1510
1511 pub(super) fn last_hop_num(&self) -> Option<HopNum> {
1515 let num_hops = self.num_hops();
1516 if num_hops == 0 {
1517 return None;
1519 }
1520 Some(HopNum::from(num_hops - 1))
1521 }
1522
1523 pub(super) fn path(&self) -> Arc<path::Path> {
1527 self.mutable.path()
1528 }
1529
1530 pub(super) fn clock_skew(&self) -> ClockSkew {
1533 self.channel.clock_skew()
1534 }
1535
1536 pub(super) fn uses_stream_sendme(&self, hop: HopNum) -> Option<bool> {
1540 let hop = self.hop(hop)?;
1541 Some(hop.ccontrol().uses_stream_sendme())
1542 }
1543
1544 pub(super) fn is_conflux_pending(&self) -> bool {
1546 let Some(status) = self.conflux_status() else {
1547 return false;
1548 };
1549
1550 status != ConfluxStatus::Linked
1551 }
1552
1553 pub(super) fn conflux_status(&self) -> Option<ConfluxStatus> {
1557 cfg_if::cfg_if! {
1558 if #[cfg(feature = "conflux")] {
1559 self.conflux_handler
1560 .as_ref()
1561 .map(|handler| handler.status())
1562 } else {
1563 None
1564 }
1565 }
1566 }
1567
1568 #[cfg(feature = "conflux")]
1570 pub(super) fn init_rtt(&self) -> Option<Duration> {
1571 self.conflux_handler
1572 .as_ref()
1573 .map(|handler| handler.init_rtt())?
1574 }
1575
1576 #[cfg(feature = "circ-padding-manual")]
1582 pub(super) fn set_padding_at_hop(
1583 &self,
1584 hop: HopNum,
1585 padder: Option<padding::CircuitPadder>,
1586 ) -> Result<()> {
1587 if self.hop(hop).is_none() {
1588 return Err(Error::NoSuchHop);
1589 }
1590 self.padding_ctrl.install_padder_padding_at_hop(hop, padder);
1591 Ok(())
1592 }
1593
1594 #[cfg(feature = "circ-padding")]
1603 fn padding_disposition(&self, send_padding: &padding::SendPadding) -> CircPaddingDisposition {
1604 crate::circuit::padding::padding_disposition(
1605 send_padding,
1606 &self.chan_sender,
1607 self.padding_block.as_ref(),
1608 )
1609 }
1610
1611 #[cfg(feature = "circ-padding")]
1613 pub(super) async fn send_padding(&mut self, send_padding: padding::SendPadding) -> Result<()> {
1614 use CircPaddingDisposition::*;
1615
1616 let target_hop = send_padding.hop;
1617
1618 match self.padding_disposition(&send_padding) {
1619 QueuePaddingNormally => {
1620 let queue_info = self.padding_ctrl.queued_padding(target_hop, send_padding);
1621 self.queue_padding_cell_for_hop(target_hop, queue_info)
1622 .await?;
1623 }
1624 QueuePaddingAndBypass => {
1625 let queue_info = self.padding_ctrl.queued_padding(target_hop, send_padding);
1626 self.queue_padding_cell_for_hop(target_hop, queue_info)
1627 .await?;
1628 }
1629 TreatQueuedCellAsPadding => {
1630 self.padding_ctrl
1631 .replaceable_padding_already_queued(target_hop, send_padding);
1632 }
1633 }
1634 Ok(())
1635 }
1636
1637 #[cfg(feature = "circ-padding")]
1641 async fn queue_padding_cell_for_hop(
1642 &mut self,
1643 target_hop: HopNum,
1644 queue_info: Option<QueuedCellPaddingInfo>,
1645 ) -> Result<()> {
1646 use tor_cell::relaycell::msg::Drop as DropMsg;
1647 let msg = SendRelayCell {
1648 hop: Some(target_hop),
1649 early: false,
1651 cell: AnyRelayMsgOuter::new(None, DropMsg::default().into()),
1652 };
1653 self.send_relay_cell_inner(msg, queue_info).await
1654 }
1655
1656 #[cfg(feature = "circ-padding")]
1659 pub(super) fn start_blocking_for_padding(&mut self, block: padding::StartBlocking) {
1660 self.chan_sender.start_blocking();
1661 self.padding_block = Some(block);
1662 }
1663
1664 #[cfg(feature = "circ-padding")]
1666 pub(super) fn stop_blocking_for_padding(&mut self) {
1667 self.chan_sender.stop_blocking();
1668 self.padding_block = None;
1669 }
1670
1671 pub(super) fn estimate_cbt(&self, length: usize) -> Duration {
1673 self.timeouts.circuit_build_timeout(length)
1674 }
1675}
1676
1677impl Drop for Circuit {
1678 fn drop(&mut self) {
1679 let _ = self.channel.close_circuit(self.circ_id);
1680 }
1681}