Skip to content

WireGuard

Source files: 17 · checked against Etemenanki 596916d · katana v3.0.1
  • Etemenanki/protocols/src/wireguard/mod.rs
  • Etemenanki/protocols/src/wireguard/config.rs
  • Etemenanki/protocols/src/wireguard/device.rs
  • Etemenanki/protocols/src/wireguard/slot.rs
  • Etemenanki/protocols/src/wireguard/connector.rs
  • Etemenanki/protocols/src/helpers/address_family.rs
  • Etemenanki/concepts/src/link.rs
  • Etemenanki/app/src/outbound/mod.rs
  • Etemenanki/app/src/inbound/mod.rs
  • Etemenanki/protocols/tests/pipeline/wireguard.rs
  • Etemenanki/protocols/tests/support/wireguard.rs
  • Etemenanki/protocols/tests/unit/wireguard/device.rs
  • Etemenanki/protocols/tests/unit/wireguard/slot.rs
  • Etemenanki/protocols/tests/unit/helpers/address_family.rs
  • Etemenanki/app/tests/integration/e2e_wg.rs
  • Etemenanki/app/tests/unit/outbound.rs
  • katana/src/outbound/mod.rs

Unlike the proxy protocols in etemenanki-protocols, WireGuard is not built from a server core and a client codec. It is a Layer 3 VPN: the tunnel carries raw IP packets, not “a connection to target X”. The module therefore provides only the outbound direction. WgConnector implements the ordinary Connector trait, and behind it a single driver task runs a WireGuard state machine (boringtun) and a userspace TCP/IP stack (smoltcp), so every routed flow becomes a TCP connection or UDP association inside that stack.

This page is for contributors who change protocols/src/wireguard/. It follows one flow from the connector down to the UDP socket and back, names every channel, buffer and constant on the way, and lists the tests that pin the behaviour. The operator’s view of the same outbound is in the WireGuard user guide.

The module does four things:

  • Parse configuration. WgConfig holds the key material, the peer endpoint, the tunnel-local addresses and the transport knobs; parse_key decodes keys from base64 or hex.
  • Run the tunnel. WgDevice owns one driver task that connects smoltcp to boringtun and boringtun to a connected UDP socket.
  • Keep one tunnel alive per outbound. DeviceSlot builds the device lazily on the first flow, shares it with every later flow, and rebuilds it with backoff when its driver has stopped.
  • Present the tunnel as a connector. WgConnector turns a TCP flow into a WgStream and a UDP flow into a WgDatagramLink.

It deliberately leaves these to others:

Concern Where it lives
Reading the [outbound.settings] table etemenanki-app app/src/outbound/mod.rs → build_wg_config; katana src/outbound/mod.rs builds the same WgConfig from panel data
Resolving destination names the Resolver passed with WgConnector::with_resolver (see DNS)
Routing, sniffing, relaying the per-connection runtime that calls the connector (see Client runtime)
The WireGuard and Noise protocol the boringtun crate (Tunn)
TCP and UDP inside the tunnel the smoltcp crate (Interface, tcp::Socket, udp::Socket)

The driver binds its own tokio::net::UdpSocket to the peer. It does not go through the dialers in etemenanki-environment, so the socket policy described in Dialers does not apply to the tunnel’s outer socket.

A WireGuard inbound would terminate a peer at Layer 3 and then see a stream of raw IP packets. The other inbounds hand the runtime flows that each carry a decoded Destination; to produce those from IP packets, a netstack would have to intercept every new TCP or UDP flow and synthesise one flow per connection. That is a transparent-proxy netstack, which sits above a protocol rather than inside it, and building it here would mean bending the concepts contracts. The module docs in protocols/src/wireguard/mod.rs record the omission as intentional and point out that the driver in device.rs is the reusable building block.

The app enforces the omission: app/src/inbound/mod.rs rejects protocol = "wireguard" on an inbound with wireguard cannot be used as an inbound (no server implementation). Accepting traffic at the IP layer is the job of the TUN inbound.

  • Directoryprotocols/src/wireguard/
    • mod.rs module docs and re-exports
    • config.rs WgConfig, parse_key, KeyError, DEFAULT_MTU
    • device.rs WgDevice, the driver task, Uplink, InMemoryDevice
    • slot.rs DeviceSlot, acquire_device, connect_tcp_any
    • connector.rs WgConnector, WgStream, WgDatagramLink
  • Directoryprotocols/src/helpers/
    • address_family.rs AddressFamilyStrategy, FamilySupport, resolve_candidates

mod.rs re-exports AddressFamilyStrategy, DEFAULT_MTU, KeyError, WgConfig, parse_key, WgConnector, WgDatagramLink, WgStream, TunnelTcp, TunnelUdp and WgDevice. The WireGuard module is not behind a feature flag: boringtun and smoltcp are ordinary dependencies of the crate.

protocols/src/wireguard/config.rs describes one point-to-point association: one peer, one endpoint.

pub const DEFAULT_MTU: usize = 1420;
pub struct WgConfig {
pub private_key: [u8; 32],
pub peer_public_key: [u8; 32],
pub preshared_key: Option<[u8; 32]>,
pub endpoint: Destination,
pub local_addrs: Vec<IpAddr>,
pub mtu: usize,
pub persistent_keepalive: Option<u16>,
pub reserved: Option<[u8; 3]>,
}
impl WgConfig {
pub fn new(
private_key: [u8; 32],
peer_public_key: [u8; 32],
endpoint: Destination,
local_addrs: Vec<IpAddr>,
) -> Self;
}
Field Meaning Set by WgConfig::new
private_key Our Curve25519 static secret, passed to x25519::StaticSecret::from. argument
peer_public_key The peer’s Curve25519 public key. argument
preshared_key Optional key mixed into the Noise handshake. None
endpoint The peer’s UDP endpoint. The driver dials it directly. argument
local_addrs Tunnel-local addresses of the virtual interface. They also decide which address families the outbound can reach. argument
mtu IP-layer MTU of the virtual interface. DEFAULT_MTU (1420, the Xray-core default)
persistent_keepalive Keepalive interval in seconds, handed to Tunn::new. None
reserved The Xray-style 3-byte header field (see The reserved bytes). None

WgConfig derives Clone; the connector wraps it in an Arc so every rebuild starts from the same value.

pub enum KeyError {
Invalid,
}
pub fn parse_key(text: &str) -> Result<[u8; 32], KeyError>;

parse_key trims surrounding whitespace, then tries two encodings in order:

  1. standard padded base64 (base64::engine::general_purpose::STANDARD), which is what wg genkey and wg pubkey print;
  2. hex, which Xray-core also accepts.

Each attempt succeeds only if it decodes to exactly 32 bytes. The two cannot collide: 32 bytes are 44 base64 characters including = padding, or 64 hex digits. Anything else returns KeyError::Invalid, whose message is invalid WireGuard key: expected base64 or hex encoding of 32 bytes. The callers wrap it: etemenanki-app reports outbound <tag>: invalid wireguard private_key (or peer_public_key, preshared_key), and katana reports wireguard outbound <tag> <field>: <KeyError>.

protocols/src/wireguard/device.rs exposes a handle to the running driver and two pairs of channels.

pub struct TunnelTcp {
pub tx: mpsc::Sender<Bytes>,
pub rx: mpsc::Receiver<Bytes>,
}
pub struct TunnelUdp {
pub tx: mpsc::Sender<(SocketAddr, Bytes)>,
pub rx: mpsc::Receiver<(SocketAddr, Bytes)>,
}
pub struct WgDevice {
cmd_tx: mpsc::Sender<Command>,
local_addrs: Vec<IpAddr>,
}
impl WgDevice {
pub async fn start(config: Arc<WgConfig>) -> io::Result<WgDevice>;
pub fn is_alive(&self) -> bool;
pub async fn connect_tcp(&self, ip: IpAddr, port: u16) -> io::Result<TunnelTcp>;
pub async fn connect_udp(&self) -> io::Result<TunnelUdp>;
}

WgDevice is deliberately not Clone; it is shared as Arc<WgDevice>. Its only link to the driver is cmd_tx, the sending half of the command channel. is_alive returns !self.cmd_tx.is_closed(): the driver owns the receiving half, so a closed channel means exactly that the driver task has returned.

The command channel carries two requests. Command is private to device.rs:

enum Command {
ConnectTcp {
dst: SocketAddr,
src: IpAddr,
out_rx: mpsc::Receiver<Bytes>,
in_tx: mpsc::Sender<Bytes>,
reply: oneshot::Sender<io::Result<()>>,
},
ConnectUdp {
out_rx: mpsc::Receiver<(SocketAddr, Bytes)>,
in_tx: mpsc::Sender<(SocketAddr, Bytes)>,
reply: oneshot::Sender<io::Result<()>>,
},
}

connect_tcp and connect_udp create two mpsc::channel(CHANNEL_CAP) pairs per connection: out for the uplink (application to tunnel) and in for the downlink. They move the driver’s ends (out_rx, in_tx) into the command, send it, and await the oneshot reply. The caller keeps the other ends, returned as TunnelTcp or TunnelUdp. Because the reply is awaited, a TCP connect failure surfaces as an error from connect_tcp rather than as an early EOF on a stream.

connect_tcp also picks the source address: local_addr_for returns the first entry of local_addrs with the same family as the destination, or AddrNotAvailable with wireguard: no tunnel-local address matching destination family. connect_udp takes no destination, because each datagram carries its own.

WgDevice::start spawns one task that owns everything below. The struct is private:

struct Driver {
tunn: Tunn,
device: InMemoryDevice,
iface: Interface,
sockets: SocketSet<'static>,
udp: UdpSocket,
cmd_rx: mpsc::Receiver<Command>,
reserved: Option<[u8; 3]>,
scratch: Vec<u8>,
recv_buf: Vec<u8>,
tcp_conns: HashMap<SocketHandle, TcpConn>,
udp_conns: HashMap<SocketHandle, UdpConn>,
next_port: u16,
}
struct TcpConn {
in_tx: mpsc::Sender<Bytes>,
uplink: Uplink<Bytes>,
connect_reply: Option<oneshot::Sender<io::Result<()>>>,
uplink_closed: bool,
}
struct UdpConn {
in_tx: mpsc::Sender<(SocketAddr, Bytes)>,
uplink: Uplink<(SocketAddr, Bytes)>,
}
Field Role
tunn The boringtun Tunn: Noise handshake, session keys, timers, encryption. A pure state machine with no I/O of its own.
device InMemoryDevice, the smoltcp phy::Device. Its rx queue holds decrypted IP packets for smoltcp to read; its tx queue collects IP packets smoltcp emits.
iface, sockets The smoltcp Interface and the SocketSet holding one smoltcp socket per tunnelled connection.
udp The UDP socket, connected to the peer endpoint.
cmd_rx Receiving half of the command channel.
scratch, recv_buf SCRATCH-sized buffers for boringtun output and for UDP receive.
tcp_conns, udp_conns Per-connection state, keyed by smoltcp SocketHandle. Each connection owns its uplink receiver inside its Uplink.
next_port The next local port for a tunnel-side socket.

Each connection’s uplink holds at most one item: one the driver has taken from the uplink channel but smoltcp has not yet accepted. connect_reply holds the oneshot sender until the TCP handshake settles.

Uplink is private to device.rs, together with the free function that waits on every connection’s uplink at once:

struct Uplink<T> {
rx: Option<mpsc::Receiver<T>>,
held: Option<T>,
}
impl<T> Uplink<T> {
fn new(rx: mpsc::Receiver<T>) -> Self;
fn take(&mut self) -> Option<T>;
fn hold(&mut self, item: T);
fn is_finished(&self) -> bool;
fn poll_refill(&mut self, cx: &mut Context<'_>) -> Poll<()>;
}
fn poll_uplinks(
tcp: &mut HashMap<SocketHandle, TcpConn>,
udp: &mut HashMap<SocketHandle, UdpConn>,
cx: &mut Context<'_>,
) -> Poll<()>;
Item Behaviour
rx The receiving half of the connection’s uplink channel. It becomes None once every sender is gone and the channel is drained.
held The one item the socket has not accepted yet.
new Wraps the uplink receiver, with nothing held.
take Returns the held item, else the next item already in the channel (try_recv). It never waits for an item. When the channel reports Disconnected, it records the end by setting rx to None.
hold Keeps an item the socket did not accept; take offers it first next time.
is_finished true when rx is gone and nothing is held: the application is done and everything it sent was handed to the socket.
poll_refill With nothing held, waits on the channel: its next item becomes the held one, or its end is recorded, and the call returns Ready. With an item held, it returns Pending (see Backpressure).
poll_uplinks Calls poll_refill on every TCP connection and every UDP association and returns Ready if any of them did.

Nothing outside the task touches this state, so the driver needs no locks. smoltcp and boringtun are both synchronous and poll-driven; the task is what drives them.

protocols/src/wireguard/slot.rs holds the tunnel that every flow of one outbound shares.

#[derive(Default)]
pub struct DeviceSlot {
pub device: Option<Arc<WgDevice>>,
pub failures: u32,
pub retry_at: Option<StdInstant>,
pub started_at: Option<StdInstant>,
}
impl DeviceSlot {
pub fn died_young(&self) -> bool;
pub fn note_failure(&mut self);
}
pub async fn acquire_device(
slot: &Mutex<DeviceSlot>,
config: &Arc<WgConfig>,
) -> io::Result<Arc<WgDevice>>;
pub fn tunnel_down() -> io::Error;
pub async fn connect_tcp_any(
device: &WgDevice,
candidates: Vec<IpAddr>,
port: u16,
) -> io::Result<TunnelTcp>;

The slot is an Option rather than a OnceCell so that a dead device can be dropped and replaced. Mutex here is tokio::sync::Mutex. The fields are public so the unit tests can inspect them; WgConnector::slot exposes the slot for the same reason.

protocols/src/wireguard/connector.rs adapts the tunnel to the Connector trait from concepts/src/link.rs (see Links and types).

pub struct WgConnector {
config: Arc<WgConfig>,
address_family: AddressFamilyStrategy,
resolver: Resolver,
device: Arc<Mutex<DeviceSlot>>,
}
impl WgConnector {
pub fn new(config: WgConfig) -> Self;
pub fn with_address_family(config: WgConfig, address_family: AddressFamilyStrategy) -> Self;
pub fn with_resolver(mut self, resolver: Resolver) -> Self;
pub fn slot(&self) -> &Arc<Mutex<DeviceSlot>>;
}
impl<T: Send + Sync + 'static> Connector<Flow<T>> for WgConnector {
type Stream = WgStream;
type Datagram = WgDatagramLink;
type Future = DialFuture;
fn connect(&mut self, flow: Flow<T>) -> DialFuture;
}
type DialFuture =
Pin<Box<dyn Future<Output = io::Result<Outbound<WgStream, WgDatagramLink>>> + Send>>;

WgConnector::new uses AddressFamilyStrategy::Auto and Resolver::default(), which is the system resolver. Both etemenanki-app and katana call with_address_family(...).with_resolver(resolver.clone()), so tunnelled names go through the configured resolver. The lookup itself runs on the host; only the connection to the resolved address travels through the tunnel.

Clone is implemented by hand and clones the Arcs, so every clone shares one DeviceSlot and therefore one tunnel. connect clones self into a boxed future that calls the private dial.

pub struct WgStream {
tx: PollSender<Bytes>,
rx: mpsc::Receiver<Bytes>,
pending: Bytes,
}
pub struct WgDatagramLink {
tx: PollSender<(SocketAddr, Bytes)>,
rx: mpsc::Receiver<(SocketAddr, Bytes)>,
resolver: Resolver,
strategy: AddressFamilyStrategy,
support: FamilySupport,
resolved: HashMap<CompactString, Option<IpAddr>>,
resolving: Option<(CompactString, ResolveFuture)>,
}

WgStream implements AsyncRead and AsyncWrite over the TunnelTcp channels; pending is the unread tail of the last chunk it received. WgDatagramLink implements DatagramLink<Addr = Destination> over the TunnelUdp channels and resolves domain destinations packet by packet.

An ordinary dialer lets the kernel pick a source address and fail fast when a family is unusable. The smoltcp netstack has no host routing table, so WireGuard answers the capability question itself, from its tunnel-local addresses (protocols/src/helpers/address_family.rs):

pub struct FamilySupport {
ipv4: bool,
ipv6: bool,
}
impl FamilySupport {
pub fn both() -> Self;
pub fn from_addrs(addrs: &[IpAddr]) -> Self;
pub fn supports(self, ip: IpAddr) -> bool;
pub fn describe(self) -> &'static str;
}
pub async fn resolve_candidates(
context: &str,
dest: &Destination,
strategy: AddressFamilyStrategy,
support: FamilySupport,
resolver: &Resolver,
) -> io::Result<Vec<IpAddr>>;

dial computes FamilySupport::from_addrs(&self.config.local_addrs) and passes it with the operator’s AddressFamilyStrategy to resolve_candidates. That keeps every resolved address, removes those the strategy forbids or the tunnel cannot source, and sorts stably for prefer_ipv4 and prefer_ipv6, so the other family stays available as a fallback. When nothing is left, the error names the capability: wireguard: no usable auto destination address for example.com:443 (local address supports IPv4 only).

Both front ends also check the combination when they build the outbound. etemenanki-app’s validate_wg_address_family rejects ipv4_only without an IPv4 address and ipv6_only without an IPv6 one (outbound <tag>: wireguard address_family ipv6_only needs an IPv6 address); katana has the same check against local_address.

Every WireGuard message starts with a 4-byte header. The WireGuard specification fixes bytes 1 to 3 at zero; Xray-core lets an operator put a 3-byte value there, which some peers use to identify a client. The driver supports it through WgConfig::reserved.

Offset Size Field On send On receive
0 1 message type (1 initiation, 2 response, 3 cookie reply, 4 transport data) written by boringtun read by boringtun
1 3 reserved apply_reserved overwrites with WgConfig::reserved when it is Some decapsulate_all zeroes it before parsing, when the datagram is longer than 3 bytes
4 rest message body written by boringtun read by boringtun

The zeroing on receive is required, not cosmetic: boringtun reads the first four bytes as one little-endian message type, so non-zero reserved bytes would make every datagram unparseable. apply_reserved runs on all three outgoing paths: encapsulated data (flush_tx), handshake and cookie replies produced while decapsulating (send_all), and timer packets.

The module docs in device.rs draw the whole path as a loop through one task. Redrawn:

flowchart LR
  app["WgStream / WgDatagramLink"] -->|"uplink channel"| held["Uplink held slot, one item"]
  held -->|"send_slice while the socket has room"| sock["smoltcp socket"]
  sock -->|"Interface::poll"| devtx["InMemoryDevice tx"]
  devtx -->|"Tunn::encapsulate, apply_reserved"| udp["connected UdpSocket"]
  udp -->|"UDP"| peer["WireGuard peer"]
  peer -->|"UDP"| udp
  udp -->|"zero reserved, Tunn::decapsulate"| devrx["InMemoryDevice rx"]
  devrx -->|"Interface::poll"| sock
  sock -->|"try_reserve on downlink channel"| app

WgDevice::start(config) performs these steps and returns as soon as the task is spawned:

  1. Reject an empty local_addrs with InvalidInput: wireguard: no tunnel-local addresses configured.
  2. Resolve the endpoint (resolve_endpoint). An IP literal is used as is; a domain goes through tokio::net::lookup_host, not the configured Resolver, and the first result is taken. No result gives NotFound: wireguard: endpoint domain did not resolve. The address is fixed for the life of the device; a rebuilt device resolves the name again.
  3. Bind a UDP socket to 0.0.0.0:0 or [::]:0, matching the endpoint’s family, and connect it to the endpoint. A connected socket receives only the peer’s datagrams and reports ICMP errors for them.
  4. Build the Tunn with the static secret, the peer’s public key, the pre-shared key, the keepalive, a random u32 session index and None for the rate limiter, so boringtun creates its own handshake rate limiter for this tunnel.
  5. Build the smoltcp Interface (build_interface): HardwareAddress::Ip on Medium::Ip, a random seed, each local address as a host route (/32 for IPv4, /128 for IPv6), and, for each family that has a local address, a default route via the first such address. On Medium::Ip the gateway is never used for next-hop resolution; the default route only makes off-link destinations routable. The device reports config.mtu as max_transmission_unit.
  6. Create the command channel (CHANNEL_CAP) and spawn the driver.

start performs no handshake. boringtun starts one when the first IP packet is encapsulated, which happens when the first flow sends its SYN or first datagram, or, when persistent_keepalive is set, when the first keepalive falls due in update_timers. So start succeeds even against a peer that is unreachable; that is why DeviceSlot judges a tunnel by how long it lived (MIN_HEALTHY_LIFETIME), not by whether start returned Ok.

etemenanki-app builds the connector at config time but never calls start then. A settings table with address = [] passes --test, and the InvalidInput error above appears on the first flow. katana rejects an empty local_address list when it builds the outbound (wireguard outbound <tag> needs at least one local_address).

Driver::run repeats the same three phases, then waits on five event sources.

flowchart TB
  poll["Interface::poll: rx packets into sockets, sockets emit into tx"] --> svc["service_sockets: connect replies, uplinks into sockets, sockets into downlink channels, closes"]
  svc --> flush["flush_tx: encapsulate every tx packet, send over UDP"]
  flush --> wait{"select!"}
  wait -->|"udp.recv"| dec["decapsulate_all, send handshake replies"]
  wait -->|"cmd_rx.recv"| cmd["open_tcp / open_udp, or return on None"]
  wait -->|"poll_fn(poll_uplinks)"| push["fill the held slot of each ready connection"]
  wait -->|"timer.tick"| tim["Tunn::update_timers, send its packet"]
  wait -->|"sleep"| idle["nothing"]
  dec --> poll
  cmd --> poll
  push --> poll
  tim --> poll
  idle --> poll

The phases:

  • Poll. self.iface.poll(now, &mut self.device, &mut self.sockets) feeds every packet in InMemoryDevice::rx to smoltcp and lets smoltcp emit segments, ACKs and datagrams into InMemoryDevice::tx.
  • Service sockets. service_sockets walks every connection. It resolves pending TCP connects, hands uplink items to the smoltcp socket while it has room, at most CHANNEL_CAP items per connection per pass so that no producer can keep the driver to itself, moves received data into the downlink channel while the channel has room, and collects sockets to remove. It returns needs_repoll, which is true when a downlink channel was full.
  • Flush. flush_tx pops each IP packet from InMemoryDevice::tx, encapsulates it into scratch, copies the result out, applies the reserved bytes and sends it with send_datagram. Without a session, boringtun queues the packet and returns a handshake initiation as WriteToNetwork, or TunnResult::Done when a handshake is already in flight; the queued packets go out later from the decapsulate loop. Done and TunnResult::Err are skipped.
  • Wait. compute_sleep turns Interface::poll_delay into a tokio::time::Sleep: 2 ms when needs_repoll is set, the smoltcp delay when there is one, and 3600 s otherwise. The select! then waits for the first of:
Branch Action
self.udp.recv(&mut self.recv_buf) decapsulate_all zeroes the reserved bytes and calls Tunn::decapsulate. A WriteToNetwork result (handshake response, cookie reply, or packets boringtun had queued) is collected, and decapsulate is called again with an empty datagram until it stops producing them, as boringtun’s contract requires. A WriteToTunnelV4 or WriteToTunnelV6 result is pushed onto InMemoryDevice::rx. send_all then sends the collected packets.
self.cmd_rx.recv() Some(cmd) goes to handle_command. None means every WgDevice handle is gone, and the driver returns.
std::future::poll_fn(poll_uplinks) Resolves once any connection that holds nothing has an item in its uplink channel or has seen the channel end. poll_uplinks polls every such connection, not only the first, so the next pass serves them all. The branch only fills the held slot; service_sockets hands the item to the socket on the next pass.
timer.tick() Every TIMER_TICK (250 ms, MissedTickBehavior::Delay), call Tunn::update_timers. A WriteToNetwork result (handshake retry, keepalive) gets the reserved bytes and is sent directly with udp.send, not through send_datagram; any error from that send ends the driver. Other results are ignored.
sleep Wake to re-poll smoltcp for its own timers (retransmission, delayed ACK) or a full downlink channel.
sequenceDiagram
  participant R as Runtime
  participant C as WgConnector
  participant S as DeviceSlot
  participant D as WgDevice
  participant T as Driver task
  participant P as Peer
  R->>C: connect(flow)
  C->>S: acquire_device (lock)
  S-->>C: Arc of WgDevice
  C->>C: resolve_candidates with FamilySupport
  C->>D: connect_tcp_any, one candidate at a time
  D->>T: Command::ConnectTcp over cmd channel
  T->>T: open_tcp: tcp::Socket::connect, keep reply
  T->>P: SYN, encapsulated (handshake first if needed)
  P-->>T: SYN-ACK, decapsulated into smoltcp
  T-->>D: reply Ok once may_send()
  D-->>C: TunnelTcp
  C-->>R: Outbound::Stream(WgStream)

WgConnector::dial branches on dest.network. For TCP it resolves the destination with resolve_candidates("wireguard", ...) and calls connect_tcp_any, which tries the candidates in order. Each attempt is tokio::time::timeout(TCP_CONNECT_ATTEMPT_TIMEOUT, device.connect_tcp(ip, port)) with a 10 s limit; the first success wins. If all fail, the error is TimedOut with every attempt’s reason: wireguard: tunnel TCP connect failed for all resolved addresses (192.0.2.1: timed out; ...).

Inside the driver, open_tcp creates a tcp::Socket with a TCP_BUFFER receive and send buffer, takes a local port from next_ephemeral_port, and calls socket.connect from src to dst. If smoltcp refuses the connect call, the reply is InvalidInput (wireguard: connect failed: <smoltcp error>) and nothing is stored. Otherwise the socket joins the SocketSet, and a TcpConn holding the reply and Uplink::new(out_rx) joins tcp_conns. On a later pass, service_sockets answers the reply:

  • socket.may_send() is true (the handshake reached Established): Ok(()).
  • socket.state() == tcp::State::Closed (the peer reset it, or smoltcp gave up): ConnectionRefused with wireguard: tunnel TCP connect failed.

Local ports come from a counter that starts at EPHEMERAL_BASE (49152) and wraps back to it after 65535. TCP sockets and UDP associations share the counter.

Uplink. WgStream::poll_write waits in PollSender::poll_reserve for a slot in the uplink channel, copies the whole buffer into one Bytes chunk and reports it as written. poll_flush does nothing. On each pass, service_sockets calls uplink.take() for the next chunk, up to CHANNEL_CAP times. take runs before the can_send() check, so an uplink that has ended is noticed in any socket state. Then:

  • the socket cannot send: the chunk is held;
  • send_slice accepts the whole chunk: the loop takes the next one;
  • send_slice accepts only part of it: the rest, chunk.split_off(sent), is held and the loop stops;
  • send_slice returns an error: the chunk is held and the loop stops.

A held chunk is offered first on the next pass. Nothing goes back into a queue, because the driver keeps at most that one chunk per connection.

Downlink. For each readable socket, service_sockets first takes a permit with in_tx.try_reserve(), then copies the contiguous readable part of the smoltcp receive buffer into one Bytes and sends it, repeating while can_recv() holds. If the channel is full, the data stays in the smoltcp receive buffer and needs_repoll is set, so the driver comes back after 2 ms. While the data stays there, smoltcp’s advertised receive window shrinks, which slows the remote sender. WgStream::poll_read hands out a received chunk across as many reads as the caller’s buffer needs, through pending.

UDP associations and per-packet resolution

Section titled “UDP associations and per-packet resolution”

A UDP flow does not resolve anything at dial time. dial calls device.connect_udp(), and the driver’s open_udp creates a smoltcp udp::Socket with 64 metadata slots and TCP_BUFFER bytes of payload storage per direction, binds it to the next ephemeral port with no fixed address, stores a UdpConn holding Uplink::new(out_rx), and replies Ok at once.

WgDatagramLink::poll_send_to(cx, buf, to) first maps the Destination to a SocketAddr in poll_target:

  • An IP literal is used directly; the strategy and FamilySupport filters apply only to names.
  • A domain already in resolved uses the cached answer. The cache stores failures as None, and neither answers nor failures expire for the life of the link.
  • Otherwise a lookup starts in resolving: a boxed future running resolve_candidates with the link’s strategy and FamilySupport, keeping only the first candidate. Only one lookup is held at a time. A packet for a different name returns Pending without polling the lookup in flight; that lookup advances only when poll_send_to is called again for its own name.

When the name does not resolve, the packet is dropped (logged at trace) and poll_send_to still returns Ok(buf.len()): dropping is what UDP does. Otherwise the link reserves a slot in the uplink channel and sends (target, payload). On each pass, service_sockets takes up to CHANNEL_CAP datagrams with uplink.take() and sends each with udp::Socket::send_slice(payload, target):

  • Ok: the datagram is in the socket’s send ring; the loop takes the next one.
  • SendError::BufferFull and the payload fits the ring (payload.len() <= socket.payload_send_capacity()): the datagram is held until the next poll empties the ring, and the loop stops.
  • Any other error, meaning a datagram larger than the whole send ring or SendError::Unaddressable: the datagram is dropped with a trace line, wireguard: dropping a <n>-byte datagram to <target>: <error>, and the loop moves on. A datagram that can never be sent therefore does not hold up the ones behind it.

Received datagrams travel back as (source, payload). poll_recv_from copies as much of the payload as fits in the caller’s buffer and returns the source as Destination { network: DialNetwork::Udp, remote: Remote::IpAddr(..), port }.

The # Backpressure section of the module docs in device.rs states the bound, and Uplink implements it. The driver reads a connection’s uplink channel only while that connection holds no item. When a socket stops accepting, because the remote stops reading or the tunnel stops draining, the connection keeps its held item and the uplink channel (CHANNEL_CAP, 256 items) fills behind it. WgStream::poll_write then waits in PollSender::poll_reserve, so the pressure reaches the producer and the driver takes nothing more for that connection. Every other connection has its own channel and held slot, and poll_uplinks keeps refilling those that hold nothing, so they keep moving. A UDP association works the same way, with WgDatagramLink::poll_send_to waiting for the channel slot.

With an item held, poll_refill returns Pending on purpose, without registering a waker. The connection is then waiting on its socket, not on its channel, and the socket’s progress (an ACK that opens the window, a poll that empties the UDP send ring) always comes through the driver loop, whose next pass offers the held item again. Leaving the channel unread is what pushes back on the application.

For one TCP connection, the uplink holds at most:

Stage Bound
smoltcp send ring TCP_BUFFER, 64 KiB
uplink channel CHANNEL_CAP, 256 items
held slot in Uplink 1 item

The channel and the held slot count items, not bytes. A poll_write chunk is the caller’s whole buffer, so the number of bytes behind a full channel follows the size of the writer’s buffers.

The downlink is bounded the other way round: a socket is read only while its downlink channel has room (see Moving TCP bytes).

TCP close is driven by the channels:

stateDiagram-v2
  [*] --> Connecting: open_tcp stores TcpConn
  Connecting --> Open: may_send, reply Ok
  Connecting --> Closed: state Closed, reply ConnectionRefused
  Connecting --> Closed: uplink ended in SYN-SENT (dial abandoned)
  Open --> LocalClosed: uplink finished, FIN sent
  Open --> LocalClosed: downlink receiver dropped, seen on next data
  Open --> RemoteClosed: FIN from the remote, no EOF yet
  LocalClosed --> Closed: remote FIN, then 10 s TIME-WAIT
  RemoteClosed --> Closed: local close, LAST-ACK acknowledged
  Open --> Closed: reset
  Closed --> [*]: socket and TcpConn removed
  • WgStream::poll_shutdown closes the PollSender. Once uplink.is_finished() is true (the sender is gone, the channel is drained and the last item has been handed to the socket), service_sockets calls socket.close() (a FIN) once and sets uplink_closed. Dropping a WgStream has the same effect.
  • If the application dropped the downlink receiver, the driver notices the next time the socket has data to deliver: try_reserve returns Closed and the driver calls socket.close().
  • When the socket reaches tcp::State::Closed, the driver removes it from the SocketSet and tcp_conns. Removing the TcpConn drops in_tx, so WgStream::poll_read returns end of file, and a later write fails with BrokenPipe (wireguard: the tunnelled connection is closed).
  • in_tx is dropped only at tcp::State::Closed, not when the remote’s FIN arrives. After a FIN from the remote end the socket sits in CLOSE-WAIT and the reader sees no end of file until the local side also closes (LAST-ACK, then Closed). When the local side closes first, the socket passes through TIME-WAIT, which smoltcp holds for 10 s, before it reaches Closed.
  • If a dial is abandoned while its socket is still in SYN-SENT, socket.close() moves it straight to Closed, and the driver removes it on the same pass.

A UDP association is removed in two cases:

  • Its uplink has finished (every sender dropped, the last datagram handed to the socket). It retires on the pass after the one in which that happened: service_sockets checks uplink.is_finished() at the top of the association’s turn, so the Interface::poll that starts the pass has already sent its last datagrams.
  • Its downlink receiver has been dropped and a datagram arrives for it.

Dropping a WgDatagramLink drops both ends, so the first case applies. poll_recv_from on a link whose driver side is gone returns BrokenPipe with the same message.

acquire_device runs under the slot’s mutex and handles three cases.

stateDiagram-v2
  [*] --> Empty
  Empty --> Live: start Ok, started_at = now
  Empty --> Waiting: start Err, note_failure
  Live --> Live: is_alive, hand out Arc
  Live --> Empty: driver stopped after MIN_HEALTHY_LIFETIME, reset failures
  Live --> Waiting: driver stopped young, note_failure
  Waiting --> Waiting: before retry_at, tunnel_down error
  Waiting --> Empty: retry_at passed
  1. Live device. device.is_alive() is true: return a clone of the Arc.
  2. Dead device. Log wireguard: tunnel driver stopped, rebuilding at warn and drop it. If it died_young (never recorded a start, or ran less than MIN_HEALTHY_LIFETIME, 10 s), count a failure with note_failure and return tunnel_down(): note_failure has just set retry_at at least 2 s in the future, so this request never rebuilds. If it ran longer, the death is a fresh problem: failures and retry_at are reset and the rebuild happens immediately.
  3. No device. If retry_at is in the future, return tunnel_down(), which is BrokenPipe with wireguard: tunnel is down, waiting before the next attempt. Otherwise call WgDevice::start. Success stores the device and started_at; failure calls note_failure, logs wireguard: tunnel start failed: <error> and returns the error.

note_failure increments failures (saturating) and sets retry_at to now plus REBUILD_BACKOFF_BASE * 2^min(failures, 5), capped at REBUILD_BACKOFF_MAX:

failures after the call Wait before the next attempt
1 2 s
2 4 s
3 8 s
4 16 s
5 or more 30 s (REBUILD_BACKOFF_MAX)

A successful start does not reset failures; only the death of a tunnel that lived at least MIN_HEALTHY_LIFETIME does.

The mutex is held across WgDevice::start, including the endpoint lookup. Flows that arrive while a tunnel is being built wait for that build and then receive the same device, instead of each starting a tunnel of their own.

Invariant Mechanism Pinned by
All flows routed to one outbound share one tunnel and one UDP socket. WgConnector clones share Arc<Mutex<DeviceSlot>>; acquire_device returns the stored Arc<WgDevice> while it is alive. a_dead_device_is_replaced_on_the_next_request (Arc::ptr_eq on two acquire_device calls) in protocols/tests/unit/wireguard/slot.rs; a_tcp_flow_rides_the_tunnel in protocols/tests/pipeline/wireguard.rs opens a second flow on the same connector
A stopped driver is replaced on the next request, not reused. WgDevice::is_alive checks cmd_tx.is_closed(); acquire_device drops a dead device and starts a new one. a_dead_device_is_replaced_on_the_next_request clears the slot itself and checks that the next request builds a live device; the is_alive() == false branch has no dedicated test
While the backoff gate is closed, no socket or driver task is created. acquire_device checks retry_at before calling WgDevice::start. the_backoff_gate_refuses_without_building_another_tunnel, a_failed_start_backs_off_before_retrying
The backoff grows with repeated failures and never exceeds 30 s. DeviceSlot::note_failure with REBUILD_BACKOFF_BASE, the min(5) exponent and REBUILD_BACKOFF_MAX. repeated_failures_escalate_the_backoff
A tunnel that dies within 10 s of starting counts as a failure. DeviceSlot::died_young against MIN_HEALTHY_LIFETIME. a_short_lived_tunnel_counts_as_a_failure
An ICMP error from an unreachable peer does not end the driver when it surfaces on the receive, flush_tx or send_all path. is_transient_send_error treats ConnectionRefused, ConnectionReset, HostUnreachable, NetworkUnreachable and WouldBlock as survivable in the udp.recv branch and in send_datagram. The timer branch sends without this filter. an_unreachable_peer_does_not_kill_the_driver
A tunnelled TCP connect reports failure as an error, not as a stream that closes at once. The oneshot reply in TcpConn::connect_reply, answered on may_send() or tcp::State::Closed; the per-attempt timeout in connect_tcp_any. a_tcp_flow_rides_the_tunnel covers the success path
Destinations are only tried in families the tunnel has a local address for. FamilySupport::from_addrs and select_candidate_ips; local_addr_for as a backstop; validate_wg_address_family at build time. auto_skips_families_without_a_local_address, the_capability_clause_is_omitted_for_a_kernel_routed_dialer in protocols/tests/unit/helpers/address_family.rs; wireguard_ipv6_only_requires_ipv6_address, wireguard_ipv4_only_builds_with_ipv4_address in app/tests/unit/outbound.rs
A connection’s uplink is bounded, and a stalled flow blocks its own writer, not other flows. Uplink holds at most one item; the channel is read only while nothing is held; service_sockets hands a socket at most CHANNEL_CAP items per pass. a_held_item_keeps_the_channel_unread, the_uplink_finishes_only_after_its_last_item in protocols/tests/unit/wireguard/device.rs; a_stalled_tcp_flow_blocks_its_writer in protocols/tests/pipeline/wireguard.rs
An unsendable datagram is dropped, not retried. A UDP send_slice error drops the datagram, unless it is BufferFull for a payload that fits the send ring. an_oversized_datagram_does_not_wedge_the_association
The driver never blocks on a slow reader. Downlink uses try_reserve; a full channel leaves data in smoltcp and schedules a 2 ms re-poll. Not pinned by a dedicated test
Reserved bytes are written on every outgoing datagram and cleared on every incoming one before boringtun parses it. apply_reserved on the flush_tx, send_all and timer paths; the zeroing at the top of decapsulate_all. Only wireguard_outbound_live_env_tcp (ignored, see below) can set a non-zero reserved, through ETEMENANKI_WG_RESERVED; the in-process peers use none
A TCP flow becomes a stream and a UDP flow becomes a datagram link. dial branches on dest.network == DialNetwork::Udp. The app turns the other combination into wireguard dialed a TCP flow as UDP (or the reverse). a_tcp_flow_rides_the_tunnel, a_udp_flow_rides_the_tunnel
Where io::ErrorKind Message
WgDevice::start InvalidInput wireguard: no tunnel-local addresses configured
WgDevice::start NotFound wireguard: endpoint domain did not resolve
WgDevice::start as reported by the OS the lookup_host, bind or connect error itself
acquire_device BrokenPipe wireguard: tunnel is down, waiting before the next attempt
resolve_candidates NotFound wireguard: destination did not resolve (an empty answer; a resolver error is passed through as is)
resolve_candidates AddrNotAvailable wireguard: no usable <strategy> destination address for <host>:<port> (local address supports <families>)
WgDevice::connect_tcp AddrNotAvailable wireguard: no tunnel-local address matching destination family
WgDevice::connect_tcp, connect_udp BrokenPipe wireguard driver stopped (the command channel or the reply was dropped)
driver, open_tcp InvalidInput wireguard: connect failed: <smoltcp error>
driver, open_udp InvalidInput wireguard: udp bind failed: <smoltcp error>
driver, service_sockets ConnectionRefused wireguard: tunnel TCP connect failed
connect_tcp_any TimedOut, whatever the individual reasons wireguard: tunnel TCP connect failed for all resolved addresses (<ip>: <reason>; ...)
WgStream, WgDatagramLink BrokenPipe wireguard: the tunnelled connection is closed

The driver task returns, ending the tunnel, when one of these happens:

Cause Log
Every WgDevice handle is dropped, so cmd_rx.recv() yields None. This is the normal shutdown. none
flush_tx or send_all hits a send error that is_transient_send_error does not cover. warn: wireguard: tunnel driver stopping, send failed: <error>
udp.recv returns an error that is_transient_send_error does not cover. warn: wireguard: tunnel driver stopping, receive failed: <error>
Sending a timer packet (handshake retry or keepalive) fails with any error. This path does not use is_transient_send_error. none

A transient error on the receive path is logged at debug (wireguard: transient receive error: <error>); on the flush_tx and send_all paths the datagram is dropped with a debug line, wireguard: dropping datagram, transient send error: <error>. WireGuard is connectionless and retransmits its handshake from the timers, so a peer that comes back is reached again without a rebuild.

When the driver returns, it drops its Driver value: every smoltcp socket, every in_tx and every uplink receiver. Readers of existing WgStreams see end of file, writers get BrokenPipe, and WgDevice::is_alive turns false, so the next acquire_device rebuilds.

  • Dropping a dial future while it waits for the reply. The oneshot receiver goes away, and the driver’s later reply.send result is ignored. The uplink sender is dropped with the future, so the uplink channel disconnects, the connection’s Uplink finishes, and service_sockets closes and removes the socket as described in Closing.
  • Dropping a dial future inside connect_tcp_any. The same as above, for the attempt in flight. The per-attempt timeout drops the attempt in exactly this way before moving to the next candidate.
  • Dropping a dial future during acquire_device. The mutex guard is released with the future. If WgDevice::start had not yet spawned the task, its socket is simply dropped.
  • Dropping the outbound. The last WgConnector clone owns the DeviceSlot, and the slot owns the Arc<WgDevice>. Once no in-flight dial holds another clone of that Arc, the command channel closes and the driver shuts down, ending every connection in the tunnel. In etemenanki-app this happens when a generation is retired on reload (see Generations and reload).
Constant Value Where What it limits
DEFAULT_MTU 1420 config.rs Interface MTU when none is configured
CHANNEL_CAP 256 device.rs Items in the command channel, and in each per-connection uplink and downlink channel; also the most items one connection hands its socket per pass
held uplink item 1 per connection device.rs → Uplink Items the driver has taken from an uplink channel that the socket has not accepted
TCP_BUFFER 64 KiB device.rs Each direction of a smoltcp TCP socket; also the payload storage of each direction of a UDP socket
UDP packet metadata 64 device.rs → open_udp Datagrams per direction a smoltcp UDP socket can hold
SCRATCH 64 KiB device.rs scratch and recv_buf, well above an MTU-sized packet plus 32 bytes of WireGuard overhead
TIMER_TICK 250 ms device.rs Cadence of Tunn::update_timers
EPHEMERAL_BASE 49152 device.rs First local port for tunnel-side sockets; the counter wraps back here after 65535
re-poll delay 2 ms device.rs → compute_sleep Wait before retrying a full downlink channel
idle delay 3600 s device.rs → compute_sleep Sleep when smoltcp has no timer pending
TCP_CONNECT_ATTEMPT_TIMEOUT 10 s slot.rs Each candidate address in connect_tcp_any
REBUILD_BACKOFF_BASE 1 s slot.rs Base of the rebuild backoff
REBUILD_BACKOFF_MAX 30 s slot.rs Cap of the rebuild backoff
MIN_HEALTHY_LIFETIME 10 s slot.rs Lifetime below which a stopped tunnel counts as a failure

The tests that carry traffic through the tunnel dial an in-process peer rather than a real server; the slot tests point the device at 127.0.0.1:1, where nothing listens, or fail before any socket exists. protocols/tests/support/wireguard.rs → spawn_echo_peer(server_priv, client_pub, sink_reads: Arc<AtomicBool>) builds the mirror image of the driver: a second boringtun Tunn bridged to its own smoltcp interface at SERVER_TUN_IP (10.0.0.1/24), with LISTENERS (4) TCP listeners that re-arm after each connection and a UDP socket, all echoing on ECHO_PORT (5555). It also runs a TCP sink at SINK_PORT (5556) with 64 KiB buffers. The sink reads nothing until sink_reads is set, so its 64 KiB receive window closes and the tunnel stops draining a writer to it; it re-arms its listener once closed. The peer binds 127.0.0.1:0, takes the client’s address from each datagram it receives (it sends nothing before the first), zeroes bytes 1 to 3 of every received datagram longer than 3 bytes, and drives its timers every 200 ms. spawn_peer generates two fixed key pairs and returns a Peer with the endpoint, the client’s private key, the server’s public key and sink_reads, the flag a test sets to let the sink read. The peer notices a change within one of its timer ticks. The client uses CLIENT_TUN_IP (10.0.0.2).

Test File What it pins
a_tcp_flow_rides_the_tunnel protocols/tests/pipeline/wireguard.rs Handshake, TCP connect and echo through WgConnector; a second flow on the same connector also echoes; shutdown succeeds
a_udp_flow_rides_the_tunnel protocols/tests/pipeline/wireguard.rs A UDP flow becomes a WgDatagramLink; the echo returns, and its source equals the destination sent to
a_stalled_tcp_flow_blocks_its_writer protocols/tests/pipeline/wireguard.rs Writing 1 KiB chunks to the sink that does not read, the writer blocks after at most 512 KiB (the expected bound is 385 KiB: the sink’s 64 KiB window, the socket’s 64 KiB, 256 channel chunks and 1 held chunk); a second flow on the same tunnel still echoes; once sink_reads is set, the blocked write completes and 1024 more chunks go through
an_oversized_datagram_does_not_wedge_the_association protocols/tests/pipeline/wireguard.rs A 70,000-byte datagram, larger than the 64 KiB send ring, is dropped, and the next datagram on the same association is echoed
a_held_item_keeps_the_channel_unread protocols/tests/unit/wireguard/device.rs With an item held, poll_refill stays Pending and the channel fills until try_send reports Full; take then returns the held item first, then the channel in order
the_uplink_finishes_only_after_its_last_item protocols/tests/unit/wireguard/device.rs After the sender is dropped, is_finished stays false while an item is held or queued and turns true once the last one is taken and the channel’s end is seen; poll_refill on a disconnected channel returns Ready and records the end
an_unreachable_peer_does_not_kill_the_driver protocols/tests/unit/wireguard/slot.rs A 300 ms connect attempt against 127.0.0.1:1, which answers with ICMP port unreachable, leaves the driver alive
a_dead_device_is_replaced_on_the_next_request protocols/tests/unit/wireguard/slot.rs A live device is reused (Arc::ptr_eq); after the test empties the slot and drops its handles, the next request builds a live device
a_failed_start_backs_off_before_retrying protocols/tests/unit/wireguard/slot.rs A failed start (InvalidInput) is followed by the BrokenPipe gate error
a_short_lived_tunnel_counts_as_a_failure protocols/tests/unit/wireguard/slot.rs died_young for no start, a fresh start and a long-lived tunnel
repeated_failures_escalate_the_backoff protocols/tests/unit/wireguard/slot.rs Strictly growing waits for four failures, capped at REBUILD_BACKOFF_MAX after many
the_backoff_gate_refuses_without_building_another_tunnel protocols/tests/unit/wireguard/slot.rs No device is stored while the gate is closed
auto_skips_families_without_a_local_address and the other address_family tests protocols/tests/unit/helpers/address_family.rs Candidate filtering by FamilySupport and strategy, and the capability clause in the error
wireguard_ipv4_only_builds_with_ipv4_address, wireguard_ipv6_only_requires_ipv6_address app/tests/unit/outbound.rs Build-time check of address_family against address
app_socks_to_wireguard_outbound_tcp app/tests/integration/e2e_wg.rs The real etemenanki-app binary: SOCKS inbound, WireGuard outbound with hex keys, TCP echo through the test file’s own copy of the peer (one listener, TCP only)
wireguard_outbound_live_env_tcp app/tests/integration/e2e_wg.rs #[ignore]d. Runs the same app path against a real peer configured through ETEMENANKI_WG_* environment variables (keys, endpoint, addresses, optional MTU, keepalive, reserved, probe target and host) and expects an HTTP response

Run them from the Etemenanki workspace:

Terminal window
cargo test -p etemenanki-protocols --test pipeline wireguard
cargo test -p etemenanki-protocols --lib wireguard
cargo test -p etemenanki-protocols --lib address_family
cargo test -p etemenanki-app --bin etemenanki-app wireguard
cargo test -p etemenanki-app --test integration e2e_wg

The --lib wireguard run includes the device unit tests, which device.rs mounts from protocols/tests/unit/wireguard/device.rs with #[path], as slot.rs does for the slot tests. The app unit tests live in the binary target, because etemenanki-app has no library target.