tor_proto/channel.rs
1//! Code for talking directly (over a TLS connection) to a Tor client or relay.
2//!
3//! Channels form the basis of the rest of the Tor protocol: they are
4//! the only way for two Tor instances to talk.
5//!
6//! Channels are not useful directly for application requests: after
7//! making a channel, it needs to get used to build circuits, and the
8//! circuits are used to anonymize streams. The streams are the
9//! objects corresponding to directory requests.
10//!
11//! In general, you shouldn't try to manage channels on your own;
12//! use the `tor-chanmgr` crate instead.
13//!
14//! To launch a channel:
15//!
16//! * Create a TLS connection as an object that implements AsyncRead +
17//! AsyncWrite + StreamOps, and pass it to a channel builder. This will
18//! yield an [crate::client::channel::handshake::ClientInitiatorHandshake] that represents
19//! the state of the handshake.
20//! * Call [crate::client::channel::handshake::ClientInitiatorHandshake::connect] on the result
21//! to negotiate the rest of the handshake. This will verify
22//! syntactic correctness of the handshake, but not its cryptographic
23//! integrity.
24//! * Call handshake::UnverifiedChannel::check on the result. This
25//! finishes the cryptographic checks.
26//! * Call handshake::VerifiedChannel::finish on the result. This
27//! completes the handshake and produces an open channel and Reactor.
28//! * Launch an asynchronous task to call the reactor's run() method.
29//!
30//! One you have a running channel, you can create circuits on it with
31//! its [Channel::new_tunnel] method. See
32//! [crate::client::circuit::PendingClientTunnel] for information on how to
33//! proceed from there.
34//!
35//! # Design
36//!
37//! For now, this code splits the channel into two pieces: a "Channel"
38//! object that can be used by circuits to write cells onto the
39//! channel, and a "Reactor" object that runs as a task in the
40//! background, to read channel cells and pass them to circuits as
41//! appropriate.
42//!
43//! I'm not at all sure that's the best way to do that, but it's what
44//! I could think of.
45//!
46//! # Limitations
47//!
48//! TODO: There is no rate limiting or fairness.
49
50/// The size of the channel buffer for communication between `Channel` and its reactor.
51pub const CHANNEL_BUFFER_SIZE: usize = 128;
52
53pub(crate) mod circmap;
54pub(crate) mod handler;
55pub(crate) mod handshake;
56pub mod kist;
57mod msg;
58pub mod padding;
59pub mod params;
60mod reactor;
61mod unique_id;
62
63#[cfg(test)]
64pub(crate) mod test_utils;
65
66pub use crate::channel::params::*;
67pub(crate) use crate::channel::reactor::Reactor;
68use crate::channel::reactor::{BoxedChannelSink, BoxedChannelStream};
69pub use crate::channel::unique_id::UniqId;
70use crate::client::circuit::PendingClientTunnel;
71use crate::client::circuit::padding::{PaddingController, QueuedCellPaddingInfo};
72use crate::memquota::{ChannelAccount, CircuitAccount, SpecificAccount as _};
73use crate::peer::PeerInfo;
74use crate::util::err::ChannelClosed;
75use crate::util::oneshot_broadcast;
76use crate::util::timeout::TimeoutEstimator;
77use crate::util::ts::AtomicOptTimestamp;
78use crate::{ClockSkew, client};
79use crate::{Error, Result};
80use cfg_if::cfg_if;
81use reactor::BoxedChannelStreamOps;
82use safelog::{MaybeSensitive, sensitive as sv};
83use std::future::{Future, IntoFuture};
84use std::net::IpAddr;
85use std::pin::Pin;
86use std::sync::{Mutex, MutexGuard};
87use std::time::Duration;
88use tor_cell::chancell::ChanMsg;
89use tor_cell::chancell::{AnyChanCell, CircId, msg::Netinfo, msg::PaddingNegotiate};
90use tor_error::internal;
91use tor_linkspec::{HasRelayIds, OwnedChanTarget};
92use tor_memquota::mq_queue::{self, ChannelSpec as _, MpscSpec};
93use tor_rtcompat::{CoarseTimeProvider, DynTimeProvider, Runtime, SleepProvider};
94
95#[cfg(feature = "circ-padding")]
96use tor_async_utils::counting_streams::{self, CountingSink, CountingStream};
97
98#[cfg(feature = "relay")]
99use {
100 crate::channel::reactor::CreateRequestHandlerAndData, crate::circuit::CircuitRxReceiver,
101 crate::relay::channel::create_handler::CreateRequestHandler,
102 tor_llcrypto::pk::ed25519::Ed25519Identity, tor_llcrypto::pk::rsa::RsaIdentity,
103};
104
105/// Imports that are re-exported pub if feature `testing` is enabled
106///
107/// Putting them together in a little module like this allows us to select the
108/// visibility for all of these things together.
109mod testing_exports {
110 #![allow(unreachable_pub)]
111 pub use super::reactor::CtrlMsg;
112 pub use crate::circuit::celltypes::CreateResponse;
113}
114#[cfg(feature = "testing")]
115pub use testing_exports::*;
116#[cfg(not(feature = "testing"))]
117use testing_exports::*;
118
119use asynchronous_codec;
120use futures::channel::mpsc;
121use futures::io::{AsyncRead, AsyncWrite};
122use oneshot_fused_workaround as oneshot;
123
124use educe::Educe;
125use futures::{FutureExt as _, Sink};
126use std::result::Result as StdResult;
127use std::sync::Arc;
128use std::task::{Context, Poll};
129
130use tracing::{instrument, trace};
131
132// reexport
133pub use super::client::channel::handshake::ClientInitiatorHandshake;
134#[cfg(feature = "relay")]
135pub use super::relay::channel::handshake::RelayInitiatorHandshake;
136pub(crate) use crate::channel::handler::{ClogDigest, SlogDigest};
137use crate::channel::unique_id::CircUniqIdContext;
138
139use kist::KistParams;
140
141/// This indicate what type of channel it is. It allows us to decide for the correct channel cell
142/// state machines and authentication process (if any).
143///
144/// It is created when a channel is requested for creation which means the subsystem wanting to
145/// open a channel needs to know what type it wants.
146#[derive(Clone, Copy, Debug, derive_more::Display)]
147#[non_exhaustive]
148pub enum ChannelType {
149 /// Client: Initiated from a client to a relay. Client is unauthenticated and relay is
150 /// authenticated.
151 ClientInitiator,
152 /// Relay: Initiating as a relay to a relay. Both sides are authenticated.
153 RelayInitiator,
154 /// Relay: Responding as a relay to a relay or client. Authenticated or Unauthenticated.
155 RelayResponder {
156 /// Indicate if the channel is authenticated. Responding as a relay can be either from a
157 /// Relay (authenticated) or a Client/Bridge (Unauthenticated). We only know this
158 /// information once the handshake is completed.
159 ///
160 /// This side is always authenticated, the other side can be if a relay or not if
161 /// bridge/client. This is set to false unless we end up authenticating the other side
162 /// meaning a relay.
163 authenticated: bool,
164 },
165}
166
167impl ChannelType {
168 /// Set that this channel type is now authenticated. This only applies to RelayResponder.
169 pub(crate) fn set_authenticated(&mut self) {
170 if let Self::RelayResponder { authenticated } = self {
171 *authenticated = true;
172 }
173 }
174}
175
176/// A channel cell frame used for sending and receiving cells on a channel. The handler takes care
177/// of the cell codec transition depending in which state the channel is.
178///
179/// ChannelFrame is used to basically handle all in and outbound cells on a channel for its entire
180/// lifetime.
181pub(crate) type ChannelFrame<T> = asynchronous_codec::Framed<T, handler::ChannelCellHandler>;
182
183/// An entry in a channel's queue of cells to be flushed.
184pub(crate) type ChanCellQueueEntry = (AnyChanCell, Option<QueuedCellPaddingInfo>);
185
186/// Helper: Return a new channel frame [ChannelFrame] from an object implementing AsyncRead + AsyncWrite. In the
187/// tor context, it is always a TLS stream.
188///
189/// The ty (type) argument needs to be able to transform into a [handler::ChannelCellHandler] which would
190/// generally be a [ChannelType].
191pub(crate) fn new_frame<T, I>(tls: T, ty: I) -> ChannelFrame<T>
192where
193 T: AsyncRead + AsyncWrite,
194 I: Into<handler::ChannelCellHandler>,
195{
196 let mut framed = asynchronous_codec::Framed::new(tls, ty.into());
197 framed.set_send_high_water_mark(32 * 1024);
198 framed
199}
200
201/// Canonical state between this channel and its peer. This is inferred from the [`Netinfo`]
202/// received during the channel handshake.
203///
204/// A connection is "canonical" if the TCP connection's peer IP address matches an address
205/// that the relay itself claims in its [`Netinfo`] cell.
206#[derive(Debug)]
207pub(crate) struct Canonicity {
208 /// The peer has proven this connection is canonical for its address: at least one NETINFO "my
209 /// address" matches the observed TCP peer address.
210 pub(crate) peer_is_canonical: bool,
211 /// We appear canonical from the peer's perspective: its NETINFO "other address" matches our
212 /// advertised relay address.
213 pub(crate) canonical_to_peer: bool,
214}
215
216impl Canonicity {
217 /// Using a [`Netinfo`], build the canonicity object with the given addresses.
218 ///
219 /// The `my_addrs` are the advertised address of this relay or empty if a client/bridge as they
220 /// do not advertise or expose a reachable address.
221 ///
222 /// The `peer_addr` is the IP address we believe the peer has. In other words, it is either the
223 /// IP we used to connect to or the address we see in the accept() phase of the connection.
224 ///
225 /// It can be None if we used a non-IP address to connect to the peer (PT).
226 pub(crate) fn from_netinfo(
227 netinfo: &Netinfo,
228 my_addrs: &[IpAddr],
229 peer_addr: Option<IpAddr>,
230 ) -> Self {
231 Self {
232 // The "other addr" (our address as seen by the peer) matches the one we advertised.
233 canonical_to_peer: netinfo
234 .their_addr()
235 .is_some_and(|a: &IpAddr| my_addrs.contains(a)),
236 // The "my addresses" (the peer addresses that it claims to have) matches the one we
237 // see on the connection or that we attempted to connect to.
238 peer_is_canonical: peer_addr
239 .map(|a| netinfo.my_addrs().contains(&a))
240 .unwrap_or_default(),
241 }
242 }
243
244 /// Construct a fully canonical object.
245 #[cfg(any(test, feature = "testing"))]
246 pub(crate) fn new_canonical() -> Self {
247 Self {
248 peer_is_canonical: true,
249 canonical_to_peer: true,
250 }
251 }
252}
253
254/// An open client channel, ready to send and receive Tor cells.
255///
256/// A channel is a direct connection to a Tor relay, implemented using TLS.
257///
258/// This struct is a frontend that can be used to send cells
259/// and otherwise control the channel. The main state is
260/// in the Reactor object.
261///
262/// (Users need a mutable reference because of the types in `Sink`, and
263/// ultimately because `cell_tx: mpsc::Sender` doesn't work without mut.
264///
265/// # Channel life cycle
266///
267/// Channels can be created directly here through a channel builder (client or relay) API.
268/// For a higher-level API (with better support for TLS, pluggable transports,
269/// and channel reuse) see the `tor-chanmgr` crate.
270///
271/// After a channel is created, it will persist until it is closed in one of
272/// four ways:
273/// 1. A remote error occurs.
274/// 2. The other side of the channel closes the channel.
275/// 3. Someone calls [`Channel::terminate`] on the channel.
276/// 4. The last reference to the `Channel` is dropped. (Note that every circuit
277/// on a `Channel` keeps a reference to it, which will in turn keep the
278/// channel from closing until all those circuits have gone away.)
279///
280/// Note that in cases 1-3, the [`Channel`] object itself will still exist: it
281/// will just be unusable for most purposes. Most operations on it will fail
282/// with an error.
283pub struct Channel {
284 /// A channel used to send control messages to the Reactor.
285 control: mpsc::UnboundedSender<CtrlMsg>,
286 /// A channel used to send cells to the Reactor.
287 cell_tx: CellTx,
288
289 /// A receiver that indicates whether the channel is closed.
290 ///
291 /// Awaiting will return a `CancelledError` event when the reactor is dropped.
292 /// Read to decide if operations may succeed, and is returned by `wait_for_close`.
293 reactor_closed_rx: oneshot_broadcast::Receiver<Result<CloseInfo>>,
294
295 /// Padding controller, used to report when data is queued for this channel.
296 padding_ctrl: PaddingController,
297
298 /// A unique identifier for this channel.
299 unique_id: UniqId,
300 /// Target identity and address information for this peer.
301 peer_id: OwnedChanTarget,
302 /// Validated information for this peer.
303 peer: MaybeSensitive<Arc<PeerInfo>>,
304 /// The declared clock skew on this channel, at the time when this channel was
305 /// created.
306 clock_skew: ClockSkew,
307 /// The time when this channel was successfully completed
308 opened_at: coarsetime::Instant,
309 /// Mutable state used by the `Channel.
310 mutable: Mutex<MutableDetails>,
311 /// Information shared with the reactor
312 details: Arc<ChannelDetails>,
313 /// Canonicity of this channel.
314 canonicity: Canonicity,
315}
316
317/// This is information shared between the reactor and the frontend (`Channel` object).
318///
319/// `control` can't be here because we rely on it getting dropped when the last user goes away.
320#[derive(Debug)]
321pub(crate) struct ChannelDetails {
322 /// Since when the channel became unused.
323 ///
324 /// If calling `time_since_update` returns None,
325 /// this channel is still in use by at least one circuit.
326 ///
327 /// Set by reactor when a circuit is added or removed.
328 /// Read from `Channel::duration_unused`.
329 unused_since: AtomicOptTimestamp,
330 /// Memory quota account
331 ///
332 /// This is here partly because we need to ensure it lives as long as the channel,
333 /// as otherwise the memquota system will tear the account down.
334 #[allow(dead_code)]
335 memquota: ChannelAccount,
336}
337
338/// Mutable details (state) used by the `Channel` (frontend)
339#[derive(Debug, Default)]
340struct MutableDetails {
341 /// State used to control padding
342 padding: PaddingControlState,
343}
344
345/// State used to control padding
346///
347/// We store this here because:
348///
349/// 1. It must be per-channel, because it depends on channel usage. So it can't be in
350/// (for example) `ChannelPaddingInstructionsUpdate`.
351///
352/// 2. It could be in the channel manager's per-channel state but (for code flow reasons
353/// there, really) at the point at which the channel manager concludes for a pending
354/// channel that it ought to update the usage, it has relinquished the lock on its own data
355/// structure.
356/// And there is actually no need for this to be global: a per-channel lock is better than
357/// reacquiring the global one.
358///
359/// 3. It doesn't want to be in the channel reactor since that's super hot.
360///
361/// See also the overview at [`tor_proto::channel::padding`](padding)
362#[derive(Debug, Educe)]
363#[educe(Default)]
364enum PaddingControlState {
365 /// No usage of this channel, so far, implies sending or negotiating channel padding.
366 ///
367 /// This means we do not send (have not sent) any `ChannelPaddingInstructionsUpdates` to the reactor,
368 /// with the following consequences:
369 ///
370 /// * We don't enable our own padding.
371 /// * We don't do any work to change the timeout distribution in the padding timer,
372 /// (which is fine since this timer is not enabled).
373 /// * We don't send any PADDING_NEGOTIATE cells. The peer is supposed to come to the
374 /// same conclusions as us, based on channel usage: it should also not send padding.
375 #[educe(Default)]
376 UsageDoesNotImplyPadding {
377 /// The last padding parameters (from reparameterize)
378 ///
379 /// We keep this so that we can send it if and when
380 /// this channel starts to be used in a way that implies (possibly) sending padding.
381 padding_params: ChannelPaddingInstructionsUpdates,
382 },
383
384 /// Some usage of this channel implies possibly sending channel padding
385 ///
386 /// The required padding timer, negotiation cell, etc.,
387 /// have been communicated to the reactor via a `CtrlMsg::ConfigUpdate`.
388 ///
389 /// Once we have set this variant, it remains this way forever for this channel,
390 /// (the spec speaks of channels "only used for" certain purposes not getting padding).
391 PaddingConfigured,
392}
393
394use PaddingControlState as PCS;
395
396cfg_if! {
397 if #[cfg(feature="circ-padding")] {
398 /// Implementation type for a ChannelSender.
399 type CellTx = CountingSink<mq_queue::Sender<ChanCellQueueEntry, mq_queue::MpscSpec>>;
400
401 /// Implementation type for a cell queue held by a reactor.
402 type CellRx = CountingStream<mq_queue::Receiver<ChanCellQueueEntry, mq_queue::MpscSpec>>;
403 } else {
404 /// Implementation type for a ChannelSender.
405 type CellTx = mq_queue::Sender<ChanCellQueueEntry, mq_queue::MpscSpec>;
406
407 /// Implementation type for a cell queue held by a reactor.
408 type CellRx = mq_queue::Receiver<ChanCellQueueEntry, mq_queue::MpscSpec>;
409 }
410}
411
412/// A handle to a [`Channel`]` that can be used, by circuits, to send channel cells.
413#[derive(Debug)]
414pub(crate) struct ChannelSender {
415 /// MPSC sender to send cells.
416 cell_tx: CellTx,
417 /// A receiver used to check if the channel is closed.
418 reactor_closed_rx: oneshot_broadcast::Receiver<Result<CloseInfo>>,
419 /// Unique ID for this channel. For logging.
420 unique_id: UniqId,
421 /// Padding controller for this channel:
422 /// used to report when we queue data that will eventually wind up on the channel.
423 padding_ctrl: PaddingController,
424}
425
426impl ChannelSender {
427 /// Check whether a cell type is permissible to be _sent_ on an
428 /// open client channel.
429 fn check_cell(&self, cell: &AnyChanCell) -> Result<()> {
430 use tor_cell::chancell::msg::AnyChanMsg::*;
431 let msg = cell.msg();
432 match msg {
433 Created(_) | Created2(_) | CreatedFast(_) => Err(Error::from(internal!(
434 "Can't send {} cell on client channel",
435 msg.cmd()
436 ))),
437 Certs(_) | Versions(_) | Authenticate(_) | AuthChallenge(_) | Netinfo(_) => {
438 Err(Error::from(internal!(
439 "Can't send {} cell after handshake is done",
440 msg.cmd()
441 )))
442 }
443 _ => Ok(()),
444 }
445 }
446
447 /// Obtain a reference to the `ChannelSender`'s [`DynTimeProvider`]
448 ///
449 /// (This can sometimes be used to avoid having to keep
450 /// a separate clone of the time provider.)
451 pub(crate) fn time_provider(&self) -> &DynTimeProvider {
452 cfg_if! {
453 if #[cfg(feature="circ-padding")] {
454 self.cell_tx.inner().time_provider()
455 } else {
456 self.cell_tx.time_provider()
457 }
458 }
459 }
460
461 /// Return an approximate count of the number of outbound cells queued for this channel.
462 ///
463 /// This count is necessarily approximate,
464 /// because the underlying count can be modified by other senders and receivers
465 /// between when this method is called and when its return value is used.
466 ///
467 /// Does not include cells that have already been passed to the TLS connection.
468 ///
469 /// Circuit padding uses this count to determine
470 /// when messages are already outbound for the first hop of a circuit.
471 #[cfg(feature = "circ-padding")]
472 pub(crate) fn approx_count(&self) -> usize {
473 self.cell_tx.approx_count()
474 }
475
476 /// Note that a cell has been queued that will eventually be placed onto this sender.
477 ///
478 /// We use this as an input for padding machines.
479 pub(crate) fn note_cell_queued(&self) {
480 self.padding_ctrl.queued_data(crate::HopNum::from(0));
481 }
482}
483
484impl Sink<ChanCellQueueEntry> for ChannelSender {
485 type Error = Error;
486
487 fn poll_ready(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<()>> {
488 let this = self.get_mut();
489 Pin::new(&mut this.cell_tx)
490 .poll_ready(cx)
491 .map_err(|_| ChannelClosed.into())
492 }
493
494 fn start_send(self: Pin<&mut Self>, cell: ChanCellQueueEntry) -> Result<()> {
495 let this = self.get_mut();
496 if this.reactor_closed_rx.is_ready() {
497 return Err(ChannelClosed.into());
498 }
499 this.check_cell(&cell.0)?;
500 {
501 use tor_cell::chancell::msg::AnyChanMsg::*;
502 match cell.0.msg() {
503 Relay(_) | Padding(_) | Vpadding(_) => {} // too frequent to log.
504 _ => trace!(
505 channel_id = %this.unique_id,
506 "Sending {} for {}",
507 cell.0.msg().cmd(),
508 CircId::get_or_zero(cell.0.circid())
509 ),
510 }
511 }
512
513 Pin::new(&mut this.cell_tx)
514 .start_send(cell)
515 .map_err(|_| ChannelClosed.into())
516 }
517
518 fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<()>> {
519 let this = self.get_mut();
520 Pin::new(&mut this.cell_tx)
521 .poll_flush(cx)
522 .map_err(|_| ChannelClosed.into())
523 }
524
525 fn poll_close(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<()>> {
526 let this = self.get_mut();
527 Pin::new(&mut this.cell_tx)
528 .poll_close(cx)
529 .map_err(|_| ChannelClosed.into())
530 }
531}
532
533impl Channel {
534 /// Construct a channel and reactor.
535 ///
536 /// Internal method, called to finalize the channel when we've
537 /// sent our netinfo cell, received the peer's netinfo cell, and
538 /// we're finally ready to create circuits.
539 ///
540 /// Quick note on the allow clippy. This is has one call site so for now, it is fine that we
541 /// bust the mighty 7 arguments.
542 #[allow(clippy::too_many_arguments)] // TODO consider if we want a builder
543 fn new<R>(
544 channel_mode: ChannelMode,
545 link_protocol: u16,
546 sink: BoxedChannelSink,
547 stream: BoxedChannelStream,
548 streamops: BoxedChannelStreamOps,
549 unique_id: UniqId,
550 peer_id: OwnedChanTarget,
551 peer: MaybeSensitive<PeerInfo>,
552 clock_skew: ClockSkew,
553 runtime: R,
554 memquota: ChannelAccount,
555 canonicity: Canonicity,
556 ) -> Result<(Arc<Self>, reactor::Reactor<R>)>
557 where
558 R: Runtime,
559 {
560 use circmap::{CircIdRange, CircMap};
561 let circid_range = match channel_mode {
562 // client channels always originate here
563 ChannelMode::Client => CircIdRange::High,
564 #[cfg(feature = "relay")]
565 ChannelMode::Relay { circ_id_range, .. } => circ_id_range,
566 };
567 let circmap = CircMap::new(circid_range);
568 let dyn_time = DynTimeProvider::new(runtime.clone());
569
570 let (control_tx, control_rx) = mpsc::unbounded();
571 let (cell_tx, cell_rx) = mq_queue::MpscSpec::new(CHANNEL_BUFFER_SIZE)
572 .new_mq(dyn_time.clone(), memquota.as_raw_account())?;
573 #[cfg(feature = "circ-padding")]
574 let (cell_tx, cell_rx) = counting_streams::channel(cell_tx, cell_rx);
575 let unused_since = AtomicOptTimestamp::new();
576 unused_since.update();
577
578 let mutable = MutableDetails::default();
579 let (reactor_closed_tx, reactor_closed_rx) = oneshot_broadcast::channel();
580
581 let details = ChannelDetails {
582 unused_since,
583 memquota,
584 };
585 let details = Arc::new(details);
586
587 // We might be using experimental maybenot padding; this creates the padding framework for that.
588 //
589 // TODO: This backend is currently optimized for circuit padding,
590 // so it might allocate a bit more than necessary to account for multiple hops.
591 // We should tune it when we deploy padding in production.
592 let (padding_ctrl, padding_event_stream) =
593 client::circuit::padding::new_padding(DynTimeProvider::new(runtime.clone()));
594
595 let channel = Arc::new(Channel {
596 control: control_tx,
597 cell_tx,
598 reactor_closed_rx,
599 padding_ctrl: padding_ctrl.clone(),
600 unique_id,
601 peer_id,
602 peer: peer.map(Arc::new),
603 clock_skew,
604 opened_at: coarsetime::Instant::now(),
605 mutable: Mutex::new(mutable),
606 details: Arc::clone(&details),
607 canonicity,
608 });
609
610 // We start disabled; the channel manager will `reconfigure` us soon after creation.
611 let padding_timer = Box::pin(padding::Timer::new_disabled(runtime.clone(), None)?);
612
613 cfg_if! {
614 if #[cfg(feature = "circ-padding")] {
615 use crate::util::sink_blocker::{SinkBlocker,CountingPolicy};
616 let sink = SinkBlocker::new(sink, CountingPolicy::new_unlimited());
617 }
618 }
619
620 #[cfg(feature = "relay")]
621 let create_request_handler: Option<_> = match channel_mode {
622 ChannelMode::Relay {
623 create_request_handler,
624 our_ed25519_id,
625 our_rsa_id,
626 ..
627 } => Some(CreateRequestHandlerAndData {
628 handler: create_request_handler,
629 channel: Arc::downgrade(&channel),
630 our_ed25519_id,
631 our_rsa_id,
632 }),
633 ChannelMode::Client => None,
634 };
635 // clippy wants us to consume `channel_mode` (`needless_pass_by_value`)
636 #[cfg(not(feature = "relay"))]
637 #[expect(clippy::drop_non_drop)]
638 drop(channel_mode);
639
640 let reactor = Reactor {
641 runtime,
642 control: control_rx,
643 cells: cell_rx,
644 reactor_closed_tx,
645 input: futures::StreamExt::fuse(stream),
646 output: sink,
647 streamops,
648 circs: circmap,
649 circ_unique_id_ctx: CircUniqIdContext::new(),
650 link_protocol,
651 unique_id,
652 details,
653 #[cfg(feature = "relay")]
654 create_request_handler,
655 padding_timer,
656 padding_ctrl,
657 padding_event_stream,
658 padding_blocker: None,
659 special_outgoing: Default::default(),
660 };
661
662 Ok((channel, reactor))
663 }
664
665 /// Return a process-unique identifier for this channel.
666 pub fn unique_id(&self) -> UniqId {
667 self.unique_id
668 }
669
670 /// Return a reference to the memory tracking account for this Channel
671 pub fn mq_account(&self) -> &ChannelAccount {
672 &self.details.memquota
673 }
674
675 /// Obtain a reference to the `Channel`'s [`DynTimeProvider`]
676 ///
677 /// (This can sometimes be used to avoid having to keep
678 /// a separate clone of the time provider.)
679 pub fn time_provider(&self) -> &DynTimeProvider {
680 cfg_if! {
681 if #[cfg(feature="circ-padding")] {
682 self.cell_tx.inner().time_provider()
683 } else {
684 self.cell_tx.time_provider()
685 }
686 }
687 }
688
689 /// Return an OwnedChanTarget representing the actual handshake used to
690 /// create this channel.
691 pub fn target(&self) -> &OwnedChanTarget {
692 &self.peer_id
693 }
694
695 /// Return the amount of time that has passed since this channel became open.
696 pub fn age(&self) -> Duration {
697 self.opened_at.elapsed().into()
698 }
699
700 /// Return a ClockSkew declaring how much clock skew the other side of this channel
701 /// claimed that we had when we negotiated the connection.
702 pub fn clock_skew(&self) -> ClockSkew {
703 self.clock_skew
704 }
705
706 /// Send a control message
707 #[instrument(level = "trace", skip_all)]
708 fn send_control(&self, msg: CtrlMsg) -> StdResult<(), ChannelClosed> {
709 self.control
710 .unbounded_send(msg)
711 .map_err(|_| ChannelClosed)?;
712 Ok(())
713 }
714
715 /// Acquire the lock on `mutable` (and handle any poison error)
716 fn mutable(&self) -> MutexGuard<MutableDetails> {
717 self.mutable.lock().expect("channel details poisoned")
718 }
719
720 /// Specify that this channel should do activities related to channel padding
721 ///
722 /// Initially, the channel does nothing related to channel padding:
723 /// it neither sends any padding, nor sends any PADDING_NEGOTIATE cells.
724 ///
725 /// After this function has been called, it will do both,
726 /// according to the parameters specified through `reparameterize`.
727 /// Note that this might include *disabling* padding
728 /// (for example, by sending a `PADDING_NEGOTIATE`).
729 ///
730 /// Idempotent.
731 ///
732 /// There is no way to undo the effect of this call.
733 #[instrument(level = "trace", skip_all)]
734 pub fn engage_padding_activities(&self) {
735 let mut mutable = self.mutable();
736
737 match &mutable.padding {
738 PCS::UsageDoesNotImplyPadding {
739 padding_params: params,
740 } => {
741 // Well, apparently the channel usage *does* imply padding now,
742 // so we need to (belatedly) enable the timer,
743 // send the padding negotiation cell, etc.
744 let mut params = params.clone();
745
746 // Except, maybe the padding we would be requesting is precisely default,
747 // so we wouldn't actually want to send that cell.
748 if params.padding_negotiate == Some(PaddingNegotiate::start_default()) {
749 params.padding_negotiate = None;
750 }
751
752 match self.send_control(CtrlMsg::ConfigUpdate(Arc::new(params))) {
753 Ok(()) => {}
754 Err(ChannelClosed) => return,
755 }
756
757 mutable.padding = PCS::PaddingConfigured;
758 }
759
760 PCS::PaddingConfigured => {
761 // OK, nothing to do
762 }
763 }
764
765 drop(mutable); // release the lock now: lock span covers the send, ensuring ordering
766 }
767
768 /// Reparameterise (update parameters; reconfigure)
769 ///
770 /// Returns `Err` if the channel was closed earlier
771 #[instrument(level = "trace", skip_all)]
772 pub fn reparameterize(&self, params: Arc<ChannelPaddingInstructionsUpdates>) -> Result<()> {
773 let mut mutable = self
774 .mutable
775 .lock()
776 .map_err(|_| internal!("channel details poisoned"))?;
777
778 match &mut mutable.padding {
779 PCS::PaddingConfigured => {
780 self.send_control(CtrlMsg::ConfigUpdate(params))?;
781 }
782 PCS::UsageDoesNotImplyPadding { padding_params } => {
783 padding_params.combine(¶ms);
784 }
785 }
786
787 drop(mutable); // release the lock now: lock span covers the send, ensuring ordering
788 Ok(())
789 }
790
791 /// Update the KIST parameters.
792 ///
793 /// Returns `Err` if the channel is closed.
794 #[instrument(level = "trace", skip_all)]
795 pub fn reparameterize_kist(&self, kist_params: KistParams) -> Result<()> {
796 Ok(self.send_control(CtrlMsg::KistConfigUpdate(kist_params))?)
797 }
798
799 /// Return an error if this channel is somehow mismatched with the
800 /// given target.
801 pub fn check_match<T: HasRelayIds + ?Sized>(&self, target: &T) -> Result<()> {
802 check_id_match_helper(&self.peer_id, target)
803 }
804
805 /// Return true if this channel is closed and therefore unusable.
806 pub fn is_closing(&self) -> bool {
807 self.reactor_closed_rx.is_ready()
808 }
809
810 /// Return true iff this channel is considered canonical by us.
811 pub fn is_canonical(&self) -> bool {
812 self.canonicity.peer_is_canonical
813 }
814
815 /// Return true if we think the peer considers this channel as canonical.
816 pub fn is_canonical_to_peer(&self) -> bool {
817 self.canonicity.canonical_to_peer
818 }
819
820 /// If the channel is not in use, return the amount of time
821 /// it has had with no circuits.
822 ///
823 /// Return `None` if the channel is currently in use.
824 pub fn duration_unused(&self) -> Option<std::time::Duration> {
825 self.details
826 .unused_since
827 .time_since_update()
828 .map(Into::into)
829 }
830
831 /// Return a new [`ChannelSender`] to transmit cells on this channel.
832 pub(crate) fn sender(&self) -> ChannelSender {
833 ChannelSender {
834 cell_tx: self.cell_tx.clone(),
835 reactor_closed_rx: self.reactor_closed_rx.clone(),
836 unique_id: self.unique_id,
837 padding_ctrl: self.padding_ctrl.clone(),
838 }
839 }
840
841 /// Return the [`PeerInfo`] of this channel.
842 #[cfg(feature = "relay")]
843 pub(crate) fn peer_info(&self) -> &Arc<PeerInfo> {
844 &self.peer
845 }
846
847 /// Return a newly allocated PendingClientTunnel object with
848 /// a corresponding tunnel reactor. A circuit ID is allocated, but no
849 /// messages are sent, and no cryptography is done.
850 ///
851 /// To use the results of this method, call Reactor::run() in a
852 /// new task, then use the methods of
853 /// [crate::client::circuit::PendingClientTunnel] to build the circuit.
854 #[instrument(level = "trace", skip_all)]
855 pub async fn new_tunnel(
856 self: &Arc<Self>,
857 timeouts: Arc<dyn TimeoutEstimator>,
858 ) -> Result<(PendingClientTunnel, client::reactor::Reactor)> {
859 if self.is_closing() {
860 return Err(ChannelClosed.into());
861 }
862
863 let time_prov = self.time_provider().clone();
864 let memquota = CircuitAccount::new(&self.details.memquota)?;
865
866 // TODO: blocking is risky, but so is unbounded.
867 let (sender, receiver) =
868 MpscSpec::new(128).new_mq(time_prov.clone(), memquota.as_raw_account())?;
869 let (sender, receiver) = crate::circuit::circ_sender::channel(sender, receiver);
870 let (createdsender, createdreceiver) = oneshot::channel::<CreateResponse>();
871
872 let (tx, rx) = oneshot::channel();
873
874 self.send_control(CtrlMsg::AllocateCircuit {
875 created_sender: createdsender,
876 sender,
877 tx,
878 })?;
879 let (circ_id, circ_unique_id, padding_ctrl, padding_stream) =
880 rx.await.map_err(|_| ChannelClosed)??;
881
882 trace!("{}: Allocated CircId {}", circ_unique_id, circ_id);
883
884 Ok(PendingClientTunnel::new(
885 circ_id,
886 self.clone(),
887 createdreceiver,
888 receiver,
889 circ_unique_id,
890 time_prov,
891 memquota,
892 padding_ctrl,
893 padding_stream,
894 timeouts,
895 ))
896 }
897
898 /// Return a newly allocated outbound relay circuit with.
899 ///
900 /// A circuit ID is allocated, but no messages are sent, and no cryptography is done.
901 ///
902 // TODO(relay): this duplicates much of new_tunnel above, but I expect
903 // the implementations to diverge once we introduce a new CtrlMsg for
904 // allocating relay circuits.
905 #[cfg(feature = "relay")]
906 pub(crate) async fn new_outbound_circ(
907 self: &Arc<Self>,
908 memquota: CircuitAccount,
909 ) -> Result<(CircId, CircuitRxReceiver, oneshot::Receiver<CreateResponse>)> {
910 if self.is_closing() {
911 return Err(ChannelClosed.into());
912 }
913
914 let time_prov = self.time_provider().clone();
915
916 // TODO: blocking is risky, but so is unbounded.
917 let (sender, receiver) =
918 MpscSpec::new(128).new_mq(time_prov.clone(), memquota.as_raw_account())?;
919 let (sender, receiver) = crate::circuit::circ_sender::channel(sender, receiver);
920 let (createdsender, createdreceiver) = oneshot::channel::<CreateResponse>();
921
922 let (tx, rx) = oneshot::channel();
923
924 self.send_control(CtrlMsg::AllocateCircuit {
925 created_sender: createdsender,
926 sender,
927 tx,
928 })?;
929
930 // TODO(relay): I don't think we need circuit-level padding on this side of the circuit.
931 // This just drops the padding controller and corresponding event stream,
932 // but maybe it would be better to just not set it up in the first place?
933 // This suggests we might need a different control command for allocating
934 // the outbound relay circuits...
935 let (id, circ_unique_id, _padding_ctrl, _padding_stream) =
936 rx.await.map_err(|_| ChannelClosed)??;
937
938 let channel_account = self.details.memquota.as_raw_account();
939 // Link the memquota circuit account with the outbound channel account:
940 memquota.as_raw_account().add_parent(channel_account)?;
941
942 trace!("{}: Allocated CircId {}", circ_unique_id, id);
943
944 Ok((id, receiver, createdreceiver))
945 }
946
947 /// Shut down this channel immediately, along with all circuits that
948 /// are using it.
949 ///
950 /// Note that other references to this channel may exist. If they
951 /// do, they will stop working after you call this function.
952 ///
953 /// It's not necessary to call this method if you're just done
954 /// with a channel: the channel should close on its own once nothing
955 /// is using it any more.
956 #[instrument(level = "trace", skip_all)]
957 pub fn terminate(&self) {
958 let _ = self.send_control(CtrlMsg::Shutdown);
959 }
960
961 /// Tell the reactor that the circuit with the given ID has gone away.
962 #[instrument(level = "trace", skip_all)]
963 pub fn close_circuit(&self, circid: CircId) -> Result<()> {
964 self.send_control(CtrlMsg::CloseCircuit(circid))?;
965 Ok(())
966 }
967
968 /// Return a future that will resolve once this channel has closed.
969 ///
970 /// Note that this method does not _cause_ the channel to shut down on its own.
971 pub fn wait_for_close(
972 &self,
973 ) -> impl Future<Output = StdResult<CloseInfo, ClosedUnexpectedly>> + Send + Sync + 'static + use<>
974 {
975 self.reactor_closed_rx
976 .clone()
977 .into_future()
978 .map(|recv| match recv {
979 Ok(Ok(info)) => Ok(info),
980 Ok(Err(e)) => Err(ClosedUnexpectedly::ReactorError(e)),
981 Err(oneshot_broadcast::SenderDropped) => Err(ClosedUnexpectedly::ReactorDropped),
982 })
983 }
984
985 /// Install a [`CircuitPadder`](client::CircuitPadder) for this channel.
986 ///
987 /// Replaces any previous padder installed.
988 #[cfg(feature = "circ-padding-manual")]
989 pub async fn start_padding(self: &Arc<Self>, padder: client::CircuitPadder) -> Result<()> {
990 self.set_padder_impl(Some(padder)).await
991 }
992
993 /// Remove any [`CircuitPadder`](client::CircuitPadder) installed for this channel.
994 ///
995 /// Does nothing if there was not a padder installed there.
996 #[cfg(feature = "circ-padding-manual")]
997 pub async fn stop_padding(self: &Arc<Self>) -> Result<()> {
998 self.set_padder_impl(None).await
999 }
1000
1001 /// Replace the [`CircuitPadder`](client::CircuitPadder) installed for this channel with `padder`.
1002 #[cfg(feature = "circ-padding-manual")]
1003 async fn set_padder_impl(
1004 self: &Arc<Self>,
1005 padder: Option<client::CircuitPadder>,
1006 ) -> Result<()> {
1007 let (tx, rx) = oneshot::channel();
1008 let msg = CtrlMsg::SetChannelPadder { padder, sender: tx };
1009 self.control
1010 .unbounded_send(msg)
1011 .map_err(|_| Error::ChannelClosed(ChannelClosed))?;
1012 rx.await.map_err(|_| Error::ChannelClosed(ChannelClosed))?
1013 }
1014
1015 /// Make a new fake reactor-less channel. For testing only, obviously.
1016 ///
1017 /// Returns the receiver end of the control message mpsc.
1018 ///
1019 /// Suitable for external callers who want to test behaviour
1020 /// of layers including the logic in the channel frontend
1021 /// (`Channel` object methods).
1022 //
1023 // This differs from test::fake_channel as follows:
1024 // * It returns the mpsc Receiver
1025 // * It does not require explicit specification of details
1026 #[cfg(feature = "testing")]
1027 pub fn new_fake(
1028 rt: impl SleepProvider + CoarseTimeProvider,
1029 _channel_type: ChannelType,
1030 ) -> (Channel, mpsc::UnboundedReceiver<CtrlMsg>) {
1031 let (control, control_recv) = mpsc::unbounded();
1032 let details = fake_channel_details();
1033
1034 let unique_id = UniqId::new();
1035 let peer_id = OwnedChanTarget::builder()
1036 .ed_identity([6_u8; 32].into())
1037 .rsa_identity([10_u8; 20].into())
1038 .build()
1039 .expect("Couldn't construct peer id");
1040
1041 // This will make rx trigger immediately.
1042 let (_tx, rx) = oneshot_broadcast::channel();
1043 let (padding_ctrl, _) = client::circuit::padding::new_padding(DynTimeProvider::new(rt));
1044
1045 let channel = Channel {
1046 control,
1047 cell_tx: fake_mpsc().0,
1048 reactor_closed_rx: rx,
1049 padding_ctrl,
1050 unique_id,
1051 peer_id,
1052 peer: MaybeSensitive::not_sensitive(Arc::new(PeerInfo::EMPTY)),
1053 clock_skew: ClockSkew::None,
1054 opened_at: coarsetime::Instant::now(),
1055 mutable: Default::default(),
1056 details,
1057 canonicity: Canonicity::new_canonical(),
1058 };
1059 (channel, control_recv)
1060 }
1061}
1062
1063/// If there is any identity in `wanted_ident` that is not present in
1064/// `my_ident`, return a ChanMismatch error.
1065///
1066/// This is a helper for [`Channel::check_match`] and
1067/// UnverifiedChannel::check_internal.
1068fn check_id_match_helper<T, U>(my_ident: &T, wanted_ident: &U) -> Result<()>
1069where
1070 T: HasRelayIds + ?Sized,
1071 U: HasRelayIds + ?Sized,
1072{
1073 for desired in wanted_ident.identities() {
1074 let id_type = desired.id_type();
1075 match my_ident.identity(id_type) {
1076 Some(actual) if actual == desired => {}
1077 Some(actual) => {
1078 return Err(Error::ChanMismatch(format!(
1079 "Identity {} does not match target {}",
1080 sv(actual),
1081 sv(desired)
1082 )));
1083 }
1084 None => {
1085 return Err(Error::ChanMismatch(format!(
1086 "Peer does not have {} identity",
1087 id_type
1088 )));
1089 }
1090 }
1091 }
1092 Ok(())
1093}
1094
1095impl HasRelayIds for Channel {
1096 fn identity(
1097 &self,
1098 key_type: tor_linkspec::RelayIdType,
1099 ) -> Option<tor_linkspec::RelayIdRef<'_>> {
1100 self.peer_id.identity(key_type)
1101 }
1102}
1103
1104/// The status of a channel which was closed successfully.
1105///
1106/// **Note:** This doesn't have any associated data,
1107/// but may be expanded in the future.
1108// I can't think of any info we'd want to return to waiters,
1109// but this type leaves the possibility open without requiring any backwards-incompatible changes.
1110#[derive(Clone, Debug)]
1111#[non_exhaustive]
1112pub struct CloseInfo;
1113
1114/// The status of a channel which closed unexpectedly.
1115#[derive(Clone, Debug, thiserror::Error)]
1116#[non_exhaustive]
1117pub enum ClosedUnexpectedly {
1118 /// The channel reactor was dropped or panicked before completing.
1119 #[error("channel reactor was dropped or panicked before completing")]
1120 ReactorDropped,
1121 /// The channel reactor had an internal error.
1122 #[error("channel reactor had an internal error")]
1123 ReactorError(Error),
1124}
1125
1126/// Whether the channel is operating in "client" or "relay" mode,
1127/// and some mode-specific parameters.
1128pub(crate) enum ChannelMode {
1129 /// An incoming channel,
1130 /// or an outgoing channel made by a non-bridge relay.
1131 #[cfg(feature = "relay")]
1132 Relay {
1133 /// A handler for CREATE2/CREATE_FAST messages.
1134 create_request_handler: Arc<CreateRequestHandler>,
1135 /// Our Ed25519 identity.
1136 our_ed25519_id: Ed25519Identity,
1137 /// Our RSA identity.
1138 our_rsa_id: RsaIdentity,
1139 /// The range of circuit IDs that we allocate for new circuits.
1140 circ_id_range: circmap::CircIdRange,
1141 },
1142 /// An outgoing channel made by a client or bridge relay.
1143 Client,
1144}
1145
1146impl ChannelMode {
1147 /// Returns an error if the mode doesn't agree with the channel type.
1148 pub(crate) fn check_agrees_with_type(
1149 &self,
1150 channel_type: ChannelType,
1151 ) -> StdResult<(), tor_error::Bug> {
1152 use ChannelType::*;
1153 use circmap::CircIdRange::*;
1154
1155 match (channel_type, self) {
1156 (ClientInitiator, Self::Client) => {}
1157 #[cfg(feature = "relay")]
1158 #[rustfmt::skip]
1159 (RelayInitiator, Self::Relay { circ_id_range: High, .. }) => {}
1160 #[cfg(feature = "relay")]
1161 #[rustfmt::skip]
1162 (RelayResponder { .. }, Self::Relay { circ_id_range: Low, .. }) => {}
1163 _ => return Err(internal!("`ChannelMode` doesn't agree with `ChannelType`")),
1164 }
1165
1166 Ok(())
1167 }
1168}
1169
1170/// Make some fake channel details (for testing only!)
1171#[cfg(any(test, feature = "testing"))]
1172fn fake_channel_details() -> Arc<ChannelDetails> {
1173 let unused_since = AtomicOptTimestamp::new();
1174
1175 Arc::new(ChannelDetails {
1176 unused_since,
1177 memquota: crate::util::fake_mq(),
1178 })
1179}
1180
1181/// Make an MPSC queue, of the type we use in Channels, but a fake one for testing
1182#[cfg(any(test, feature = "testing"))] // Used by Channel::new_fake which is also feature=testing
1183pub(crate) fn fake_mpsc() -> (CellTx, CellRx) {
1184 let (tx, rx) = crate::fake_mpsc(CHANNEL_BUFFER_SIZE);
1185 #[cfg(feature = "circ-padding")]
1186 let (tx, rx) = counting_streams::channel(tx, rx);
1187 (tx, rx)
1188}
1189
1190#[cfg(test)]
1191pub(crate) mod test {
1192 // Most of this module is tested via tests that also check on the
1193 // reactor code; there are just a few more cases to examine here.
1194 #![allow(clippy::unwrap_used)]
1195 use super::*;
1196 pub(crate) use crate::channel::reactor::test::{CodecResult, new_reactor};
1197 use tor_cell::chancell::msg::HandshakeType;
1198 use tor_cell::chancell::{AnyChanCell, msg};
1199 use tor_rtcompat::test_with_one_runtime;
1200
1201 /// Make a new fake reactor-less channel. For testing only, obviously.
1202 pub(crate) fn fake_channel(
1203 rt: impl SleepProvider + CoarseTimeProvider,
1204 _channel_type: ChannelType,
1205 ) -> Channel {
1206 let unique_id = UniqId::new();
1207 let peer_id = OwnedChanTarget::builder()
1208 .ed_identity([6_u8; 32].into())
1209 .rsa_identity([10_u8; 20].into())
1210 .build()
1211 .expect("Couldn't construct peer id");
1212 // This will make rx trigger immediately.
1213 let (_tx, rx) = oneshot_broadcast::channel();
1214 let (padding_ctrl, _) = client::circuit::padding::new_padding(DynTimeProvider::new(rt));
1215 Channel {
1216 control: mpsc::unbounded().0,
1217 cell_tx: fake_mpsc().0,
1218 reactor_closed_rx: rx,
1219 padding_ctrl,
1220 unique_id,
1221 peer_id,
1222 peer: MaybeSensitive::not_sensitive(Arc::new(PeerInfo::EMPTY)),
1223 clock_skew: ClockSkew::None,
1224 opened_at: coarsetime::Instant::now(),
1225 mutable: Default::default(),
1226 details: fake_channel_details(),
1227 canonicity: Canonicity::new_canonical(),
1228 }
1229 }
1230
1231 #[test]
1232 fn send_bad() {
1233 tor_rtcompat::test_with_all_runtimes!(|rt| async move {
1234 use std::error::Error;
1235 let chan = fake_channel(rt, ChannelType::ClientInitiator);
1236
1237 let cell = AnyChanCell::new(CircId::new(7), msg::Created2::new(&b"hihi"[..]).into());
1238 let e = chan.sender().check_cell(&cell);
1239 assert!(e.is_err());
1240 assert!(
1241 format!("{}", e.unwrap_err().source().unwrap())
1242 .contains("Can't send CREATED2 cell on client channel")
1243 );
1244 let cell = AnyChanCell::new(None, msg::Certs::new_empty().into());
1245 let e = chan.sender().check_cell(&cell);
1246 assert!(e.is_err());
1247 assert!(
1248 format!("{}", e.unwrap_err().source().unwrap())
1249 .contains("Can't send CERTS cell after handshake is done")
1250 );
1251
1252 let cell = AnyChanCell::new(
1253 CircId::new(5),
1254 msg::Create2::new(HandshakeType::NTOR, &b"abc"[..]).into(),
1255 );
1256 let e = chan.sender().check_cell(&cell);
1257 assert!(e.is_ok());
1258 // FIXME(eta): more difficult to test that sending works now that it has to go via reactor
1259 // let got = output.next().await.unwrap();
1260 // assert!(matches!(got.msg(), ChanMsg::Create2(_)));
1261 });
1262 }
1263
1264 #[test]
1265 fn check_match() {
1266 test_with_one_runtime!(|rt| async move {
1267 let chan = fake_channel(rt, ChannelType::ClientInitiator);
1268
1269 let t1 = OwnedChanTarget::builder()
1270 .ed_identity([6; 32].into())
1271 .rsa_identity([10; 20].into())
1272 .build()
1273 .unwrap();
1274 let t2 = OwnedChanTarget::builder()
1275 .ed_identity([1; 32].into())
1276 .rsa_identity([3; 20].into())
1277 .build()
1278 .unwrap();
1279 let t3 = OwnedChanTarget::builder()
1280 .ed_identity([3; 32].into())
1281 .rsa_identity([2; 20].into())
1282 .build()
1283 .unwrap();
1284
1285 assert!(chan.check_match(&t1).is_ok());
1286 assert!(chan.check_match(&t2).is_err());
1287 assert!(chan.check_match(&t3).is_err());
1288 });
1289 }
1290
1291 #[test]
1292 fn unique_id() {
1293 test_with_one_runtime!(|rt| async move {
1294 let ch1 = fake_channel(rt.clone(), ChannelType::ClientInitiator);
1295 let ch2 = fake_channel(rt, ChannelType::ClientInitiator);
1296 assert_ne!(ch1.unique_id(), ch2.unique_id());
1297 });
1298 }
1299
1300 #[test]
1301 fn duration_unused_at() {
1302 test_with_one_runtime!(|rt| async move {
1303 let details = fake_channel_details();
1304 let mut ch = fake_channel(rt, ChannelType::ClientInitiator);
1305 ch.details = details.clone();
1306 details.unused_since.update();
1307 assert!(ch.duration_unused().is_some());
1308 });
1309 }
1310}