1use crate::client::circuit::handshake::RelayCryptLayerProtocol;
5
6use crate::ccparams::CongestionControlParams;
7use crate::circuit::CircParameters;
8use crate::congestion::{CongestionControl, sendme};
9use crate::memquota::{SpecificAccount, StreamAccount};
10use crate::stream::CloseStreamBehavior;
11use crate::stream::SEND_WINDOW_INIT;
12use crate::stream::StreamMpscSender;
13use crate::stream::cmdcheck::{AnyCmdChecker, StreamStatus};
14use crate::stream::flow_ctrl::params::FlowCtrlParameters;
15use crate::stream::flow_ctrl::state::{
16 FlowCtrlHooks, StreamFlowCtrl, StreamRateLimit, WithSidechannelMitigations,
17};
18use crate::stream::flow_ctrl::xon_xoff::reader::DrainRateRequest;
19use crate::stream::queue::{StreamQueueReceiver, stream_queue};
20use crate::streammap::{
21 self, EndSentStreamEnt, OpenStreamEnt, ShouldSendEnd, StreamEntMut, StreamMap,
22};
23use crate::util::notify::{NotifyReceiver, NotifySender};
24use crate::{Error, HopNum, Result};
25
26use derive_deftly::Deftly;
27use postage::watch;
28use safelog::sensitive as sv;
29use tracing::{debug, trace};
30
31use tor_cell::chancell::{BoxedCellBody, CircId};
32use tor_cell::relaycell::extend::{CcRequest, CircRequestExt};
33use tor_cell::relaycell::flow_ctrl::{Xoff, Xon, XonKBpsEwma};
34use tor_cell::relaycell::msg::AnyRelayMsg;
35use tor_cell::relaycell::{
36 AnyRelayMsgOuter, RelayCellDecoder, RelayCellDecoderResult, RelayCellFormat, RelayCmd,
37 StreamId, UnparsedRelayMsg,
38};
39use tor_error::{Bug, ErrorKind, HasKind, internal, into_internal};
40use tor_memquota::derive_deftly_template_HasMemoryCost;
41use tor_memquota::mq_queue::{ChannelSpec as _, MpscSpec};
42use tor_protover::named;
43use tor_rtcompat::DynTimeProvider;
44
45use std::num::NonZeroU32;
46use std::pin::Pin;
47use std::result::Result as StdResult;
48use std::sync::{Arc, Mutex};
49use web_time_compat::Instant;
50
51#[cfg(test)]
52use tor_cell::relaycell::msg::SendmeTag;
53
54#[cfg(feature = "relay")]
55use {
56 crate::ccparams::{Algorithm, AlgorithmDiscriminants},
57 crate::circuit::HandshakeSubprotocols,
58 crate::relay::{CircNetParameters, CongestionControlNetParams},
59};
60
61use cfg_if::cfg_if;
62
63const CIRCUIT_BUFFER_SIZE: usize = 128;
66
67#[derive(Debug, Clone, Copy, Eq, PartialEq)]
75pub(crate) enum HopNegotiationType {
76 None,
78 HsV3,
87 Full,
89}
90
91#[derive(Clone, Debug)]
103pub(crate) struct HopSettings {
104 pub(crate) ccontrol: CongestionControlParams,
106
107 pub(crate) flow_ctrl_params: FlowCtrlParameters,
109
110 pub(crate) n_incoming_cells_permitted: Option<u32>,
112
113 pub(crate) n_outgoing_cells_permitted: Option<u32>,
115
116 relay_crypt_protocol: RelayCryptLayerProtocol,
118}
119
120impl HopSettings {
121 #[allow(clippy::unnecessary_wraps)] pub(crate) fn from_params_and_caps(
132 hoptype: HopNegotiationType,
133 params: &CircParameters,
134 caps: &tor_protover::Protocols,
135 ) -> Result<Self> {
136 let mut ccontrol = params.ccontrol.clone();
137 match ccontrol.alg() {
138 crate::ccparams::Algorithm::FixedWindow(_) => {}
139 crate::ccparams::Algorithm::Vegas(_) => {
140 if !caps.supports_named_subver(named::FLOWCTRL_CC) {
142 ccontrol.use_fallback_alg();
143 }
144 }
145 };
146 if hoptype == HopNegotiationType::None {
147 ccontrol.use_fallback_alg();
148 }
149 let ccontrol = ccontrol; let relay_crypt_protocol = match hoptype {
154 HopNegotiationType::None => RelayCryptLayerProtocol::Tor1(RelayCellFormat::V0),
155 HopNegotiationType::HsV3 => {
156 cfg_if! {
157 if #[cfg(feature = "hs-common")] {
158 if ccontrol.alg().compatible_with_cgo() && caps.supports_named_subver(named::RELAY_CRYPT_CGO) {
159 RelayCryptLayerProtocol::Cgo
160 } else {
161 RelayCryptLayerProtocol::HsV3(RelayCellFormat::V0)
162 }
163 } else {
164 return Err(
165 tor_error::internal!("Unexpectedly tried to negotiate HsV3 without support!").into(),
166 );
167 }
168 }
169 }
170 HopNegotiationType::Full => {
171 #[allow(clippy::overly_complex_bool_expr)]
172 if ccontrol.alg().compatible_with_cgo()
173 && caps.supports_named_subver(named::RELAY_NEGOTIATE_SUBPROTO)
174 && caps.supports_named_subver(named::RELAY_CRYPT_CGO)
175 {
176 RelayCryptLayerProtocol::Cgo
177 } else {
178 RelayCryptLayerProtocol::Tor1(RelayCellFormat::V0)
179 }
180 }
181 };
182
183 Ok(Self {
184 ccontrol,
185 flow_ctrl_params: params.flow_ctrl.clone(),
186 relay_crypt_protocol,
187 n_incoming_cells_permitted: params.n_incoming_cells_permitted,
188 n_outgoing_cells_permitted: params.n_outgoing_cells_permitted,
189 })
190 }
191
192 #[warn(unused)]
197 #[cfg(feature = "relay")]
198 pub(crate) fn from_handshake_params(
199 circ_net_params: CircNetParameters,
200 cc_algorithm: AlgorithmDiscriminants,
201 subprotos_requested: HandshakeSubprotocols,
202 ) -> StdResult<Self, HandshakeParamsError> {
203 let CircNetParameters {
206 cc:
207 CongestionControlNetParams {
208 fixed_window,
209 vegas_exit,
210 cwnd,
211 rtt,
212 flow_ctrl,
213 },
214 } = circ_net_params;
215
216 let HandshakeSubprotocols { relay_crypt_cgo } = subprotos_requested;
217
218 let (cc_algorithm, relay_crypt_protocol) = match (cc_algorithm, relay_crypt_cgo) {
223 (AlgorithmDiscriminants::FixedWindow, false) => (
224 Algorithm::FixedWindow(fixed_window),
225 RelayCryptLayerProtocol::Tor1(RelayCellFormat::V0),
226 ),
227 (AlgorithmDiscriminants::FixedWindow, true) => {
228 return Err(HandshakeParamsError::IncompatibleParams(
229 "requested CGO but not congestion control",
230 ));
231 }
232 (AlgorithmDiscriminants::Vegas, false) => (
233 Algorithm::Vegas(vegas_exit),
234 RelayCryptLayerProtocol::Tor1(RelayCellFormat::V0),
235 ),
236 (AlgorithmDiscriminants::Vegas, true) => {
237 (Algorithm::Vegas(vegas_exit), RelayCryptLayerProtocol::Cgo)
238 }
239 };
240
241 let ccontrol = CongestionControlParams::builder()
243 .alg(cc_algorithm)
244 .fixed_window_params(fixed_window)
245 .cwnd_params(cwnd)
246 .rtt_params(rtt)
247 .build()
248 .map_err(into_internal!("Could not build `CongestionControlParams`"))?;
249
250 Ok(Self {
251 ccontrol,
252 flow_ctrl_params: flow_ctrl,
253 relay_crypt_protocol,
254 n_incoming_cells_permitted: None,
255 n_outgoing_cells_permitted: None,
256 })
257 }
258
259 pub(crate) fn relay_crypt_protocol(&self) -> RelayCryptLayerProtocol {
261 self.relay_crypt_protocol
262 }
263
264 #[allow(clippy::unnecessary_wraps)]
267 pub(crate) fn circuit_request_extensions(&self) -> Result<Vec<CircRequestExt>> {
268 #[allow(unused_mut)]
270 let mut client_extensions = Vec::new();
271
272 #[allow(unused, unused_mut)]
273 let mut cc_extension_set = false;
274
275 if self.ccontrol.is_enabled() {
276 client_extensions.push(CircRequestExt::CcRequest(CcRequest::default()));
277 cc_extension_set = true;
278 }
279
280 #[allow(unused_mut)]
293 let mut required_protocol_capabilities: Vec<tor_protover::NamedSubver> = Vec::new();
294
295 if matches!(self.relay_crypt_protocol(), RelayCryptLayerProtocol::Cgo) {
296 if !cc_extension_set {
297 return Err(tor_error::internal!("Tried to negotiate CGO without CC.").into());
298 }
299 required_protocol_capabilities.push(tor_protover::named::RELAY_CRYPT_CGO);
300 }
301
302 if !required_protocol_capabilities.is_empty() {
303 client_extensions.push(CircRequestExt::SubprotocolRequest(
304 required_protocol_capabilities.into_iter().collect(),
305 ));
306 }
307
308 Ok(client_extensions)
309 }
310}
311
312#[cfg(test)]
313impl std::default::Default for CircParameters {
314 fn default() -> Self {
315 Self {
316 extend_by_ed25519_id: true,
317 ccontrol: crate::congestion::test_utils::params::build_cc_fixed_params(),
318 flow_ctrl: FlowCtrlParameters::defaults_for_tests(),
319 n_incoming_cells_permitted: None,
320 n_outgoing_cells_permitted: None,
321 }
322 }
323}
324
325#[derive(Clone, Debug, thiserror::Error)]
328pub(crate) enum HandshakeParamsError {
329 #[error("The provided handshake parameters are incompatible with each other: {0}")]
331 IncompatibleParams(&'static str),
332 #[error("Internal error")]
334 Internal(#[from] tor_error::Bug),
335}
336
337impl HasKind for HandshakeParamsError {
338 fn kind(&self) -> ErrorKind {
339 match self {
340 Self::IncompatibleParams(_) => ErrorKind::TorProtocolViolation,
341 Self::Internal(_) => ErrorKind::Internal,
342 }
343 }
344}
345
346impl CircParameters {
347 pub fn new(
349 extend_by_ed25519_id: bool,
350 ccontrol: CongestionControlParams,
351 flow_ctrl: FlowCtrlParameters,
352 ) -> Self {
353 Self {
354 extend_by_ed25519_id,
355 ccontrol,
356 flow_ctrl,
357 n_incoming_cells_permitted: None,
358 n_outgoing_cells_permitted: None,
359 }
360 }
361}
362
363#[derive(educe::Educe)]
368#[educe(Debug)]
369pub(crate) struct SendRelayCell {
370 pub(crate) hop: Option<HopNum>,
372 pub(crate) early: bool,
374 pub(crate) cell: AnyRelayMsgOuter,
376}
377
378pub(crate) struct CircHopInbound {
380 decoder: RelayCellDecoder,
382 n_incoming_cells_permitted: Option<NonZeroU32>,
390}
391
392pub(crate) struct CircHopOutbound {
394 ccontrol: Arc<Mutex<CongestionControl>>,
398 map: Arc<Mutex<StreamMap>>,
416 relay_format: RelayCellFormat,
420 flow_ctrl_params: Arc<FlowCtrlParameters>,
422 n_outgoing_cells_permitted: Option<NonZeroU32>,
426}
427
428impl CircHopInbound {
429 pub(crate) fn new(decoder: RelayCellDecoder, settings: &HopSettings) -> Self {
431 Self {
432 decoder,
433 n_incoming_cells_permitted: settings.n_incoming_cells_permitted.map(cvt),
434 }
435 }
436
437 pub(crate) fn decode(&mut self, cell: BoxedCellBody) -> Result<RelayCellDecoderResult> {
442 self.decoder
443 .decode(cell)
444 .map_err(|e| Error::from_bytes_err(e, "relay cell"))
445 }
446
447 pub(crate) fn decrement_cell_limit(&mut self) -> Result<()> {
450 try_decrement_cell_limit(&mut self.n_incoming_cells_permitted)
451 .map_err(|_| Error::ExcessInboundCells)
452 }
453}
454
455impl CircHopOutbound {
456 pub(crate) fn new(
458 ccontrol: Arc<Mutex<CongestionControl>>,
459 relay_format: RelayCellFormat,
460 flow_ctrl_params: Arc<FlowCtrlParameters>,
461 settings: &HopSettings,
462 ) -> Self {
463 Self {
464 ccontrol,
465 map: Arc::new(Mutex::new(StreamMap::new())),
466 relay_format,
467 flow_ctrl_params,
468 n_outgoing_cells_permitted: settings.n_outgoing_cells_permitted.map(cvt),
469 }
470 }
471
472 pub(crate) fn begin_stream(
475 &mut self,
476 hop: Option<HopNum>,
477 message: AnyRelayMsg,
478 time_prov: &DynTimeProvider,
479 cmd_checker: AnyCmdChecker,
480 memquota: &StreamAccount,
481 ) -> Result<(SendRelayCell, StreamId, ReactorStreamComponents)> {
482 let (rate_limit_tx, rate_limit_rx) = watch::channel_with(StreamRateLimit::MAX);
486
487 let mut drain_rate_request_tx = NotifySender::new_typed();
491 let drain_rate_request_rx = drain_rate_request_tx.subscribe();
492
493 let flow_ctrl = self.build_flow_ctrl(
494 WithSidechannelMitigations::Enabled,
497 rate_limit_tx,
498 drain_rate_request_tx,
499 )?;
500
501 let stream_queue_max_len = flow_ctrl.inbound_queue_max_len();
502
503 let (sender, receiver) = stream_queue(stream_queue_max_len, memquota, time_prov)?;
505
506 let (msg_tx, msg_rx) = MpscSpec::new(CIRCUIT_BUFFER_SIZE)
508 .new_mq(time_prov.clone(), memquota.as_raw_account())?;
509
510 let r = self.map.lock().expect("lock poisoned").add_ent(
511 sender,
512 msg_rx,
513 flow_ctrl,
514 cmd_checker,
515 )?;
516 let cell = AnyRelayMsgOuter::new(Some(r), message);
517
518 let stream_components = ReactorStreamComponents {
519 stream_inbound_rx: receiver,
520 stream_outbound_tx: msg_tx,
521 rate_limit_rx,
522 drain_rate_request_rx,
523 };
524
525 Ok((
526 SendRelayCell {
527 hop,
528 early: false,
529 cell,
530 },
531 r,
532 stream_components,
533 ))
534 }
535
536 #[allow(clippy::too_many_arguments)]
547 pub(crate) fn close_stream(
548 &mut self,
549 circ_uniq_id: impl std::fmt::Display,
550 circ_id: CircId,
551 id: StreamId,
552 hop: Option<HopNum>,
553 message: CloseStreamBehavior,
554 why: streammap::TerminateReason,
555 expiry: Instant,
556 ) -> Result<Option<SendRelayCell>> {
557 let should_send_end = self
558 .map
559 .lock()
560 .expect("lock poisoned")
561 .terminate(id, why, expiry)?;
562 trace!(
563 circ_uniq_id = %circ_uniq_id,
564 circ_id = %circ_id,
565 stream_id = %id,
566 should_send_end = ?should_send_end,
567 "Ending stream",
568 );
569 let end_message = match should_send_end {
572 ShouldSendEnd::Send => match message {
573 CloseStreamBehavior::SendEnd(end_message) => end_message.into(),
574 CloseStreamBehavior::SendResolved(resolved_message) => resolved_message.into(),
575 CloseStreamBehavior::SendNothing => return Ok(None),
576 },
577 ShouldSendEnd::DontSend => return Ok(None),
578 };
579
580 let end_cell = AnyRelayMsgOuter::new(Some(id), end_message);
581 let cell = SendRelayCell {
582 hop,
583 early: false,
584 cell: end_cell,
585 };
586
587 Ok(Some(cell))
588 }
589
590 pub(crate) fn maybe_send_xon(
594 &mut self,
595 rate: XonKBpsEwma,
596 id: StreamId,
597 ) -> Result<Option<Xon>> {
598 if !self
601 .ccontrol()
602 .lock()
603 .expect("poisoned lock")
604 .uses_xon_xoff()
605 {
606 return Ok(None);
607 }
608
609 let mut map = self.map.lock().expect("lock poisoned");
610 let Some(StreamEntMut::Open(ent)) = map.get_mut(id) else {
611 return Ok(None);
613 };
614
615 ent.maybe_send_xon(rate)
616 }
617
618 pub(crate) fn maybe_send_xoff(&mut self, id: StreamId) -> Result<Option<Xoff>> {
622 if !self
625 .ccontrol()
626 .lock()
627 .expect("poisoned lock")
628 .uses_xon_xoff()
629 {
630 return Ok(None);
631 }
632
633 let mut map = self.map.lock().expect("lock poisoned");
634 let Some(StreamEntMut::Open(ent)) = map.get_mut(id) else {
635 return Ok(None);
637 };
638
639 ent.maybe_send_xoff()
640 }
641
642 pub(crate) fn relay_cell_format(&self) -> RelayCellFormat {
647 self.relay_format
648 }
649
650 #[cfg(test)]
652 pub(crate) fn send_window_and_expected_tags(&self) -> (u32, Vec<SendmeTag>) {
653 self.ccontrol()
654 .lock()
655 .expect("poisoned lock")
656 .send_window_and_expected_tags()
657 }
658
659 pub(crate) fn n_open_streams(&self) -> usize {
664 self.map.lock().expect("lock poisoned").n_open_streams()
665 }
666
667 pub(crate) fn ccontrol(&self) -> &Arc<Mutex<CongestionControl>> {
669 &self.ccontrol
670 }
671
672 pub(crate) fn about_to_send(
683 &mut self,
684 circ_uniq_id: impl std::fmt::Display,
685 circ_id: CircId,
686 stream_id: StreamId,
687 msg: &AnyRelayMsg,
688 ) -> Result<()> {
689 let mut hop_map = self.map.lock().expect("lock poisoned");
690 let Some(StreamEntMut::Open(ent)) = hop_map.get_mut(stream_id) else {
691 debug!(
701 circ_uniq_id = %circ_uniq_id,
702 circ_id = %circ_id,
703 stream_id = %stream_id,
704 "sending a relay cell for non-existent or non-open stream!",
705 );
706 return Ok(());
707 };
708
709 ent.about_to_send(msg)
710 }
711
712 #[cfg(any(feature = "hs-service", feature = "relay"))]
714 pub(crate) fn add_ent_with_id(
715 &self,
716 time_prov: &DynTimeProvider,
717 stream_id: StreamId,
718 cmd_checker: AnyCmdChecker,
719 with_sidechannel_mitigations: WithSidechannelMitigations,
720 memquota: &StreamAccount,
721 ) -> Result<ReactorStreamComponents> {
722 let (rate_limit_tx, rate_limit_rx) = watch::channel_with(StreamRateLimit::MAX);
726
727 let mut drain_rate_request_tx = NotifySender::new_typed();
731 let drain_rate_request_rx = drain_rate_request_tx.subscribe();
732
733 let flow_ctrl = self.build_flow_ctrl(
734 with_sidechannel_mitigations,
735 rate_limit_tx,
736 drain_rate_request_tx,
737 )?;
738
739 let stream_queue_max_len = flow_ctrl.inbound_queue_max_len();
740
741 let (sender, receiver) = stream_queue(stream_queue_max_len, memquota, time_prov)?;
743
744 let (msg_tx, msg_rx) = MpscSpec::new(CIRCUIT_BUFFER_SIZE)
746 .new_mq(time_prov.clone(), memquota.as_raw_account())?;
747
748 let mut hop_map = self.map.lock().expect("lock poisoned");
749 hop_map.add_ent_with_id(sender, msg_rx, flow_ctrl, stream_id, cmd_checker)?;
750
751 Ok(ReactorStreamComponents {
752 stream_inbound_rx: receiver,
753 stream_outbound_tx: msg_tx,
754 rate_limit_rx,
755 drain_rate_request_rx,
756 })
757 }
758
759 #[expect(clippy::unnecessary_wraps)]
762 fn build_flow_ctrl(
763 &self,
764 with_sidechannel_mitigations: WithSidechannelMitigations,
765 rate_limit_updater: watch::Sender<StreamRateLimit>,
766 drain_rate_requester: NotifySender<DrainRateRequest>,
767 ) -> Result<StreamFlowCtrl> {
768 let params = Arc::clone(&self.flow_ctrl_params);
769
770 if self
771 .ccontrol()
772 .lock()
773 .expect("poisoned lock")
774 .uses_stream_sendme()
775 {
776 let window = sendme::StreamSendWindow::new(SEND_WINDOW_INIT);
777 Ok(StreamFlowCtrl::new_window(window))
778 } else {
779 Ok(StreamFlowCtrl::new_xon_xoff(
780 params,
781 with_sidechannel_mitigations,
782 rate_limit_updater,
783 drain_rate_requester,
784 ))
785 }
786 }
787
788 fn deliver_msg_to_stream(
790 streamid: StreamId,
791 ent: &mut OpenStreamEnt,
792 cell_counts_toward_windows: bool,
793 msg: UnparsedRelayMsg,
794 ) -> Result<bool> {
795 use tor_async_utils::SinkTrySend as _;
796 use tor_async_utils::SinkTrySendError as _;
797
798 match msg.cmd() {
803 RelayCmd::SENDME => {
804 ent.put_for_incoming_sendme(msg)?;
805 return Ok(false);
806 }
807 RelayCmd::XON => {
808 ent.handle_incoming_xon(msg)?;
809 return Ok(false);
810 }
811 RelayCmd::XOFF => {
812 ent.handle_incoming_xoff(msg)?;
813 return Ok(false);
814 }
815 _ => {}
816 }
817
818 let message_closes_stream = ent.cmd_checker.check_msg(&msg)? == StreamStatus::Closed;
819
820 if let Err(e) = Pin::new(&mut ent.sink).try_send(msg) {
821 if e.is_full() {
822 return Err(internal!(
823 "Stream (ID {}) uses an unbounded queue, but apparently it's full?",
824 sv(streamid),
825 )
826 .into());
827 }
828 if e.is_disconnected() && cell_counts_toward_windows {
829 ent.dropped += 1;
834 }
835 }
836
837 Ok(message_closes_stream)
838 }
839
840 #[cfg(feature = "hs-service")]
845 pub(crate) fn ending_msg_received(&self, stream_id: StreamId) -> Result<()> {
846 let mut hop_map = self.map.lock().expect("lock poisoned");
847
848 hop_map.ending_msg_received(stream_id)?;
849
850 Ok(())
851 }
852
853 pub(crate) fn handle_msg<F>(
861 &self,
862 possible_proto_violation_err: F,
863 cell_counts_toward_windows: bool,
864 streamid: StreamId,
865 msg: UnparsedRelayMsg,
866 now: Instant,
867 ) -> Result<Option<UnparsedRelayMsg>>
868 where
869 F: FnOnce(StreamId) -> Error,
870 {
871 let mut hop_map = self.map.lock().expect("lock poisoned");
872
873 match hop_map.get_mut(streamid) {
874 Some(StreamEntMut::Open(ent)) => {
875 let message_closes_stream =
877 Self::deliver_msg_to_stream(streamid, ent, cell_counts_toward_windows, msg)?;
878
879 if message_closes_stream {
880 hop_map.ending_msg_received(streamid)?;
881 }
882 }
883 Some(StreamEntMut::EndSent(EndSentStreamEnt { expiry, .. })) if now >= *expiry => {
884 return Err(possible_proto_violation_err(streamid));
885 }
886 Some(StreamEntMut::EndSent(_))
887 if matches!(
888 msg.cmd(),
889 RelayCmd::BEGIN | RelayCmd::BEGIN_DIR | RelayCmd::RESOLVE
890 ) =>
891 {
892 hop_map.ending_msg_received(streamid)?;
896 return Ok(Some(msg));
897 }
898 Some(StreamEntMut::EndSent(EndSentStreamEnt { half_stream, .. })) => {
899 match half_stream.handle_msg(msg)? {
902 StreamStatus::Open => {}
903 StreamStatus::Closed => {
904 hop_map.ending_msg_received(streamid)?;
905 }
906 }
907 }
908 None if matches!(
909 msg.cmd(),
910 RelayCmd::BEGIN | RelayCmd::BEGIN_DIR | RelayCmd::RESOLVE
911 ) =>
912 {
913 return Ok(Some(msg));
914 }
915 _ => {
916 return Err(possible_proto_violation_err(streamid));
918 }
919 }
920
921 Ok(None)
922 }
923
924 pub(crate) fn stream_map(&self) -> &Arc<Mutex<StreamMap>> {
926 &self.map
927 }
928
929 pub(crate) fn set_stream_map(&mut self, map: Arc<Mutex<StreamMap>>) -> StdResult<(), Bug> {
933 if self.n_open_streams() != 0 {
934 return Err(internal!("Tried to discard existing open streams?!"));
935 }
936
937 self.map = map;
938
939 Ok(())
940 }
941
942 pub(crate) fn decrement_cell_limit(&mut self) -> Result<()> {
945 try_decrement_cell_limit(&mut self.n_outgoing_cells_permitted)
946 .map_err(|_| Error::ExcessOutboundCells)
947 }
948}
949
950#[inline]
953fn try_decrement_cell_limit(val: &mut Option<NonZeroU32>) -> StdResult<(), ()> {
954 match val {
956 Some(x) => {
957 let z = u32::from(*x);
958 if z == 1 {
959 Err(())
960 } else {
961 *x = (z - 1).try_into().expect("NonZeroU32 was zero?!");
962 Ok(())
963 }
964 }
965 None => Ok(()),
966 }
967}
968
969fn cvt(limit: u32) -> NonZeroU32 {
972 limit
974 .saturating_add(1)
975 .try_into()
976 .expect("Adding one left it as zero?")
977}
978
979#[derive(Debug, Deftly)]
987#[derive_deftly(HasMemoryCost)]
988pub(crate) struct ReactorStreamComponents {
989 #[deftly(has_memory_cost(indirect_size = "0"))] pub(crate) stream_inbound_rx: StreamQueueReceiver,
992
993 #[deftly(has_memory_cost(indirect_size = "size_of::<AnyRelayMsg>()"))] pub(crate) stream_outbound_tx: StreamMpscSender<AnyRelayMsg>,
996
997 #[deftly(has_memory_cost(indirect_size = "0"))]
1000 pub(crate) rate_limit_rx: watch::Receiver<StreamRateLimit>,
1001
1002 #[deftly(has_memory_cost(indirect_size = "0"))]
1005 pub(crate) drain_rate_request_rx: NotifyReceiver<DrainRateRequest>,
1006}