Skip to content

Commit 4c55a2d

Browse files
committed
ts_derp: factor out peer id lookup
The baseline derp client doesn't need to know about translating nodekeys to transport `PeerId`s, that's a higher level concern now provided by `ts_transport::UnderlayTransportExt::with_lookup`, which the runtime provides. Signed-off-by: Nathan Perry <nathan@tailscale.com> Change-Id: Ib965f787e92880ac3d74c364760acc546a6a6964
1 parent fc08131 commit 4c55a2d

9 files changed

Lines changed: 72 additions & 185 deletions

File tree

Cargo.lock

Lines changed: 0 additions & 2 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

ts_derp/examples/listen.rs

Lines changed: 7 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,6 @@
33
//! Intended to test ping/pong/keepalive.
44
55
use ts_keys::NodeKeyPair;
6-
use ts_transport::{BatchRecvIter, UnderlayTransport};
76

87
mod common;
98

@@ -16,20 +15,16 @@ async fn main() -> ts_cli_util::Result<()> {
1615

1716
let keypair = NodeKeyPair::new();
1817

19-
let client =
20-
ts_derp::Client::connect(region, &keypair, ts_derp::DummyStaticLookup::default()).await?;
18+
let client = ts_derp::Client::connect(region, &keypair).await?;
2119
tracing::info!("derp handshake done");
2220

2321
loop {
24-
for result in client.recv().await.batch_iter() {
25-
match result {
26-
Ok((peer, pkts)) => {
27-
let pkts = pkts.into_iter().collect::<Vec<_>>();
28-
tracing::info!(?peer, ?pkts);
29-
}
30-
Err(e) => {
31-
tracing::error!(err = %e, "recv");
32-
}
22+
match client.recv_one().await {
23+
Ok((peer, pkt)) => {
24+
tracing::info!(?peer, ?pkt);
25+
}
26+
Err(e) => {
27+
tracing::error!(err = %e, "recv");
3328
}
3429
}
3530
}

ts_derp/examples/ping.rs

Lines changed: 4 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -3,9 +3,7 @@
33
use std::{sync::Arc, time::Duration};
44

55
use tokio::task::JoinSet;
6-
use ts_derp::PeerLookup;
76
use ts_keys::NodeKeyPair;
8-
use ts_transport::UnderlayTransport;
97

108
mod common;
119

@@ -18,10 +16,7 @@ async fn main() -> ts_cli_util::Result<()> {
1816

1917
let keypair = NodeKeyPair::new();
2018

21-
let peer_map = &*Box::leak(Box::new(ts_derp::DummyStaticLookup::default()));
22-
let self_id = peer_map.key_to_id(&keypair.public).unwrap();
23-
24-
let client = ts_derp::Client::connect(region, &keypair, peer_map).await?;
19+
let client = ts_derp::Client::connect(region, &keypair).await?;
2520
tracing::info!("derp handshake done");
2621

2722
let client = Arc::new(client);
@@ -33,7 +28,7 @@ async fn main() -> ts_cli_util::Result<()> {
3328
let mut ticker = tokio::time::interval(Duration::from_secs(1));
3429

3530
loop {
36-
if let Err(e) = pinger.send([(self_id, vec![vec![1].into()])]).await {
31+
if let Err(e) = pinger.send_one(keypair.public, &[1]).await {
3732
tracing::error!(err = %e, "ping");
3833
} else {
3934
tracing::info!("ping");
@@ -47,10 +42,8 @@ async fn main() -> ts_cli_util::Result<()> {
4742
js.spawn(async move {
4843
loop {
4944
match recv.recv_one().await {
50-
Ok((peer_id, pkt)) => {
51-
let peer_key = peer_map.id_to_key(peer_id);
52-
53-
tracing::info!(?pkt, %peer_id, ?peer_key, "pong");
45+
Ok((peer_key, pkt)) => {
46+
tracing::info!(?pkt, %peer_key, "pong");
5447
}
5548
Err(e) => {
5649
tracing::error!(err = %e, "recv");
Lines changed: 29 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -8,27 +8,25 @@ use tokio::{
88
};
99
use tokio_util::codec::{FramedRead, FramedWrite};
1010
use ts_http_util::Client as _;
11-
use ts_keys::NodeKeyPair;
11+
use ts_keys::{NodeKeyPair, NodePublicKey};
1212
use ts_packet::PacketMut;
13-
use ts_transport::{BatchRecvIter, BatchSendIter, PeerId, UnderlayTransport};
13+
use ts_transport::{BatchRecvIter, BatchSendIter, UnderlayTransport};
1414
use url::Url;
1515

1616
use crate::{
1717
Error, ServerConnInfo, frame,
1818
frame::{ClientInfo, FrameType, PeerGone, Ping, RawFrame, ServerInfo, ServerKey},
19-
peer_lookup::PeerLookup,
2019
};
2120

2221
type DefaultIo = ts_http_util::Upgraded;
2322

2423
/// Type alias for the default derp client over upgraded HTTP on a tokio executor.
25-
pub type DefaultClient<Lookup> = Client<DefaultIo, Lookup>;
24+
pub type DefaultClient = Client<DefaultIo>;
2625

27-
/// Asynchronous DERP transport for a single DERP region.
28-
pub struct Client<Io, Lookup> {
26+
/// Single-region DERP client.
27+
pub struct Client<Io> {
2928
read_conn: Mutex<FramedRead<ReadHalf<Io>, frame::Codec>>,
3029
write_conn: Mutex<FramedWrite<WriteHalf<Io>, frame::Codec>>,
31-
peer_lookup: Lookup,
3230
}
3331

3432
/// Establish and upgrade a http connection to the derp region.
@@ -56,18 +54,13 @@ pub async fn connect<'c>(
5654
Ok(Some(upgraded))
5755
}
5856

59-
impl<Io, Lookup> Client<Io, Lookup>
57+
impl<Io> Client<Io>
6058
where
6159
Io: AsyncRead + AsyncWrite,
62-
Lookup: PeerLookup,
6360
{
6461
/// Perform a derp handshake over the given transport and return a [`Client`].
6562
#[tracing::instrument(skip_all)]
66-
pub async fn handshake(
67-
conn: Io,
68-
node_keypair: &NodeKeyPair,
69-
peer_lookup: Lookup,
70-
) -> Result<Self, Error> {
63+
pub async fn handshake(conn: Io, node_keypair: &NodeKeyPair) -> Result<Self, Error> {
7164
let (read_conn, write_conn) = tokio::io::split(conn);
7265

7366
let mut fw = FramedWrite::new(write_conn, frame::Codec);
@@ -121,10 +114,15 @@ where
121114
Ok(Self {
122115
read_conn: Mutex::new(fr),
123116
write_conn: Mutex::new(fw),
124-
peer_lookup,
125117
})
126118
}
127119

120+
/// Send a message to a nodekey on the derp server.
121+
pub async fn send_one(&self, node_key: NodePublicKey, msg: &[u8]) -> Result<(), Error> {
122+
self.send_frame_with_extra(&frame::SendPacket { dest: node_key }, msg)
123+
.await
124+
}
125+
128126
/// Send a frame to the derp server.
129127
pub async fn send_frame(
130128
&self,
@@ -151,7 +149,7 @@ where
151149

152150
/// Waits for a single data packet from a peer to arrive via this DERP server and returns it.
153151
/// DERP control messages (KeepAlive, Ping, etc) are handled inline and are not returned.
154-
pub async fn recv_one(&self) -> Result<(PeerId, PacketMut), Error> {
152+
pub async fn recv_one(&self) -> Result<(NodePublicKey, PacketMut), Error> {
155153
// DERP exchanges control messages (KeepAlives, Pings, etc) in-band with data messages
156154
// (SendPacket, RecvPacket, etc). The caller only cares about the payloads of data
157155
// messages, so we recv_one_raw() in a loop to handle any control messages while waiting
@@ -199,12 +197,8 @@ where
199197
}
200198
FrameType::RecvPacket => {
201199
let (recv, payload) = frame.as_type::<frame::RecvPacket>().unwrap();
202-
let Some(id) = self.peer_lookup.key_to_id(&recv.src) else {
203-
tracing::trace!(src = %recv.src, "no known peer for node key");
204-
continue;
205-
};
206200

207-
return Ok((id, payload.into()));
201+
return Ok((recv.src, payload.into()));
208202
}
209203
t => {
210204
return Err(Error::UnexpectedRecvFrameType(t));
@@ -214,29 +208,25 @@ where
214208
}
215209
}
216210

217-
impl<Lookup> Client<DefaultIo, Lookup>
218-
where
219-
Lookup: PeerLookup,
220-
{
211+
impl Client<DefaultIo> {
221212
/// Connect to and handshake with the derp server with the given URL over HTTP.
222213
pub async fn connect<'c>(
223214
region: impl IntoIterator<Item = &'c ServerConnInfo>,
224215
node_keypair: &NodeKeyPair,
225-
lookup: Lookup,
226216
) -> Result<Self, Error> {
227217
let conn = connect(region).await?.unwrap();
228218

229-
Client::handshake(conn, node_keypair, lookup).await
219+
Client::handshake(conn, node_keypair).await
230220
}
231221
}
232222

233-
impl<Io, Lookup> fmt::Debug for Client<Io, Lookup> {
223+
impl<Io> fmt::Debug for Client<Io> {
234224
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
235225
fmt::Display::fmt(self, f)
236226
}
237227
}
238228

239-
impl<Io, Lookup> fmt::Display for Client<Io, Lookup> {
229+
impl<Io> fmt::Display for Client<Io> {
240230
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
241231
f.debug_tuple("Client").finish()
242232
}
@@ -291,36 +281,27 @@ fn decrypt_server_info(
291281
Ok(sip)
292282
}
293283

294-
impl<Io, Lookup> UnderlayTransport for Client<Io, Lookup>
284+
impl<Io> UnderlayTransport for Client<Io>
295285
where
296286
Io: AsyncRead + AsyncWrite + Send,
297-
Lookup: PeerLookup,
298287
{
299-
type PeerKey = PeerId;
288+
type PeerKey = NodePublicKey;
300289
type Error = Error;
301290

302-
#[tracing::instrument(fields(%self))]
303-
async fn recv(&self) -> impl BatchRecvIter<Self::PeerKey, Error = Self::Error> {
304-
[self.recv_one().await.map(|(k, pkt)| (k, [pkt]))]
305-
}
306-
307-
/// Send a batch of packets to a peer via this DERP server.
308291
async fn send(
309292
&self,
310-
peer_packets: impl BatchSendIter<Self::PeerKey>,
293+
packet_batch: impl BatchSendIter<Self::PeerKey>,
311294
) -> Result<(), Self::Error> {
312-
for (peer, packets) in peer_packets.batch_iter() {
313-
let Some(node_key) = self.peer_lookup.id_to_key(peer) else {
314-
tracing::warn!(peer_id = %peer, "no node key known for peer");
315-
continue;
316-
};
317-
318-
for packet in packets {
319-
self.send_frame_with_extra(&frame::SendPacket { dest: node_key }, packet.as_ref())
320-
.await?;
295+
for (key, pkt) in packet_batch.batch_iter() {
296+
for pkt in pkt {
297+
self.send_one(key, pkt.as_ref()).await?;
321298
}
322299
}
323300

324301
Ok(())
325302
}
303+
304+
async fn recv(&self) -> impl BatchRecvIter<Self::PeerKey, Error = Self::Error> {
305+
[self.recv_one().await.map(|(k, pkt)| (k, [pkt]))]
306+
}
326307
}

ts_derp/src/lib.rs

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -8,15 +8,13 @@ use core::{
88

99
use zerocopy::{FromBytes, Immutable, IntoBytes, KnownLayout};
1010

11-
mod async_tokio;
11+
mod client;
1212
pub mod dial;
1313
mod error;
1414
pub mod frame;
15-
mod peer_lookup;
1615

17-
pub use async_tokio::{Client, DefaultClient};
16+
pub use client::{Client, DefaultClient};
1817
pub use error::Error;
19-
pub use peer_lookup::{DummyStaticLookup, PeerLookup};
2018

2119
/// A 24-byte nonce for symmetric encryption with ChaCha20Poly1305.
2220
#[repr(C)]

ts_derp/src/peer_lookup.rs

Lines changed: 0 additions & 71 deletions
This file was deleted.

ts_devtools/Cargo.toml

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,10 +20,8 @@ tracing.workspace = true
2020

2121
ts_cli_util.workspace = true
2222
ts_keys.workspace = true
23-
ts_packet.workspace = true
2423
ts_packetfilter.workspace = true
2524
ts_packetfilter_state.workspace = true
26-
ts_transport.workspace = true
2725
ts_derp.workspace = true
2826

2927
[lints]

0 commit comments

Comments
 (0)