diff --git a/CHANGELOG.md b/CHANGELOG.md index 2e30a11d..ef445ad3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -496,6 +496,23 @@ with v0.5.x or earlier peers. identity msg2 carried, dropping the leg on a mismatch. Two situations that used to end in a connection no longer do, both of them intended. +- Tracing targets of the TUN adapter, DNS responder and ICMPv6 code move from + `fips::upper::*` to `fips::ipv6tun::*`, because the module they log from is + now `ipv6tun` (for example `fips::upper::tun` becomes `fips::ipv6tun::tun`). + The hosts-file loader and reloader log as `fips::hosts` rather than + `fips::upper::hosts`, since the hosts file is now a top-level module. The + ICMPv6 Packet Too Big debug lines ("Sending ICMP Packet Too Big", "Rate + limiting ICMP Packet Too Big") log as `fips::ipv6tun::icmp` rather than + `fips::node::handlers::session`, so a `fips::node=debug` filter no longer + shows them. All logging from TUN and DNS start and stop (for example "TUN + device active", "Shutting down TUN interface" and "DNS responder started", + with their failure warnings) logs as `fips::ipv6tun::lifecycle` rather than + `fips::node::lifecycle`, so a filter on `fips::node::lifecycle` or + `fips::node` no longer selects it. An existing `RUST_LOG` filter naming an + old target still parses and simply stops matching, so the symptom is missing + log lines rather than an error. Update `RUST_LOG` filters, journal-watch + recipes and any log-scraping alert accordingly. + #### Packaging - The glibc floor is 2.34, one step below the 2.35 Ubuntu 22.04 sets. Both diff --git a/docs/how-to/diagnose-mtu-issues.md b/docs/how-to/diagnose-mtu-issues.md index c81f727f..fff84b44 100644 --- a/docs/how-to/diagnose-mtu-issues.md +++ b/docs/how-to/diagnose-mtu-issues.md @@ -87,7 +87,7 @@ fully-qualified module paths under the `fips` crate. sudo systemctl edit fips # Add: # [Service] -# Environment=RUST_LOG=info,fips::upper::tun=trace,fips::node::handlers::mmp=debug +# Environment=RUST_LOG=info,fips::ipv6tun::tun=trace,fips::node::handlers::mmp=debug sudo systemctl restart fips sudo journalctl -u fips -f ``` diff --git a/src/upper/hosts.rs b/src/hosts.rs similarity index 100% rename from src/upper/hosts.rs rename to src/hosts.rs diff --git a/src/upper/config.rs b/src/ipv6tun/config.rs similarity index 100% rename from src/upper/config.rs rename to src/ipv6tun/config.rs diff --git a/src/upper/dns.rs b/src/ipv6tun/dns.rs similarity index 93% rename from src/upper/dns.rs rename to src/ipv6tun/dns.rs index 69500c21..497caf61 100644 --- a/src/upper/dns.rs +++ b/src/ipv6tun/dns.rs @@ -164,7 +164,7 @@ fn is_mesh_interface_query(arrival_ifindex: Option, mesh_ifindex: Option Result { + use socket2::{Domain, Protocol, Socket, Type}; + let domain = if addr.is_ipv4() { + Domain::IPV4 + } else { + Domain::IPV6 + }; + let sock = Socket::new(domain, Type::DGRAM, Some(Protocol::UDP))?; + if addr.is_ipv6() { + sock.set_only_v6(false)?; + #[cfg(unix)] + set_recv_pktinfo_v6(&sock)?; + } + sock.set_nonblocking(true)?; + sock.bind(&addr.into())?; + tokio::net::UdpSocket::from_std(sock.into()) +} + +/// Enable `IPV6_RECVPKTINFO` on an IPv6 UDP socket. +/// +/// After this setsockopt, each `recvmsg()` call on the socket receives +/// an `IPV6_PKTINFO` control message containing the arrival interface +/// index, which the DNS responder uses for its mesh-interface filter. +#[cfg(unix)] +fn set_recv_pktinfo_v6(sock: &socket2::Socket) -> Result<(), std::io::Error> { + use std::os::fd::AsRawFd; + let enable: libc::c_int = 1; + let ret = unsafe { + libc::setsockopt( + sock.as_raw_fd(), + libc::IPPROTO_IPV6, + libc::IPV6_RECVPKTINFO, + &enable as *const _ as *const libc::c_void, + std::mem::size_of::() as libc::socklen_t, + ) + }; + if ret < 0 { + return Err(std::io::Error::last_os_error()); + } + Ok(()) +} + +/// Resolve an interface index by name. +/// +/// Returns `None` if the interface does not exist. +pub(crate) fn lookup_mesh_ifindex(name: &str) -> Option { + #[cfg(unix)] + { + let c_name = std::ffi::CString::new(name).ok()?; + let idx = unsafe { libc::if_nametoindex(c_name.as_ptr()) }; + if idx == 0 { None } else { Some(idx) } + } + #[cfg(not(unix))] + { + let _ = name; + None + } +} + /// Receive a UDP datagram with arrival-interface info via `IPV6_PKTINFO`. /// /// Returns `(len, src, arrival_ifindex)`. The ifindex is `Some` when the @@ -860,7 +933,7 @@ mod tests { } /// Build a socket bound to `[::1]:0` with `IPV6_RECVPKTINFO` enabled, - /// mirroring the setup done in `Node::bind_dns_socket`. + /// mirroring the setup done in `bind_dns_socket`. #[cfg(unix)] fn bind_loopback_v6_with_pktinfo() -> tokio::net::UdpSocket { use socket2::{Domain, Protocol, Socket, Type}; diff --git a/src/upper/icmp.rs b/src/ipv6tun/icmp.rs similarity index 87% rename from src/upper/icmp.rs rename to src/ipv6tun/icmp.rs index 080e2534..088d714d 100644 --- a/src/upper/icmp.rs +++ b/src/ipv6tun/icmp.rs @@ -4,7 +4,10 @@ //! Currently supports Destination Unreachable (Type 1) for //! packets that cannot be routed. +use super::icmp_rate_limit::IcmpRateLimiter; +use super::tun::TunTx; use std::net::Ipv6Addr; +use tracing::debug; /// ICMPv6 message types. #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -370,6 +373,102 @@ pub fn build_packet_too_big( Some(response) } +/// Sends ICMPv6 errors back to the host through the TUN writer. +/// +/// Holds what both replies need: the channel to the TUN writer, this +/// node's address, and the per-source Packet Too Big limiter. It borrows +/// all three for the duration of one use, so it never keeps the TUN +/// channel open past the adapter's teardown. With no TUN channel the +/// replies are built and dropped. +pub(crate) struct IcmpContext<'a> { + /// Channel to the TUN writer, `None` when no TUN is up. + tun_tx: Option<&'a TunTx>, + /// This node's address, the source of Destination Unreachable replies. + our_addr: Ipv6Addr, + /// Per-source limiter, applied to Packet Too Big only. + limiter: &'a mut IcmpRateLimiter, +} + +impl<'a> IcmpContext<'a> { + /// Borrow the TUN channel and limiter, with our address as `our_addr`. + pub(crate) fn new( + tun_tx: Option<&'a TunTx>, + our_addr: Ipv6Addr, + limiter: &'a mut IcmpRateLimiter, + ) -> Self { + Self { + tun_tx, + our_addr, + limiter, + } + } + + /// Send ICMPv6 Destination Unreachable back through TUN. + /// + /// Not rate limited: only Packet Too Big goes through the limiter. + pub(crate) fn dest_unreachable(&self, original_packet: &[u8]) { + if !should_send_icmp_error(original_packet) { + return; + } + + if let Some(response) = + build_dest_unreachable(original_packet, DestUnreachableCode::NoRoute, self.our_addr) + && let Some(tun_tx) = self.tun_tx + { + let _ = tun_tx.send(response); + } + } + + /// Answer packets that were queued for a destination found to have no + /// route, with one Destination Unreachable each. + pub(crate) fn no_route>(&self, packets: impl IntoIterator) { + for packet in packets { + self.dest_unreachable(packet.as_ref()); + } + } + + /// Send ICMPv6 Packet Too Big back through TUN. + /// + /// Rate-limited per source address to prevent ICMP floods from + /// misconfigured applications sending repeated oversized packets. + pub(crate) fn too_big(&mut self, original_packet: &[u8], mtu: u32) { + // Extract source address for rate limiting + if original_packet.len() < 40 { + return; + } + // SAFETY: slice is exactly 16 bytes; length validated above (>= 40) + let src_addr = Ipv6Addr::from(<[u8; 16]>::try_from(&original_packet[8..24]).unwrap()); + + // Rate limit ICMP PTB messages per source + if !self.limiter.should_send(src_addr) { + debug!( + src = %src_addr, + "Rate limiting ICMP Packet Too Big" + ); + return; + } + + // Use the original packet's *destination* as the ICMP source so the + // kernel sees the PTB coming from a remote router, not from itself. + // Linux ignores PTBs whose source matches a local address, which + // causes a PMTUD blackhole when both src and ICMP-src are local. + // SAFETY: slice is exactly 16 bytes; length validated above (>= 40) + let dest_addr = Ipv6Addr::from(<[u8; 16]>::try_from(&original_packet[24..40]).unwrap()); + if let Some(response) = build_packet_too_big(original_packet, mtu, dest_addr) + && let Some(tun_tx) = self.tun_tx + { + debug!( + original_src = %src_addr, + original_dst = %dest_addr, + packet_size = original_packet.len(), + reported_mtu = mtu, + "Sending ICMP Packet Too Big" + ); + let _ = tun_tx.send(response); + } + } +} + /// Calculate ICMPv6 checksum per RFC 4443. /// /// The checksum is calculated over a pseudo-header plus the ICMPv6 message. @@ -709,7 +808,7 @@ mod tests { let remote_addr: Ipv6Addr = "fddf::2".parse().unwrap(); let original = make_ipv6_packet(local_addr, remote_addr, 6, &[0u8; 1200]); // TCP - // Pass remote_addr as our_addr — this is what send_icmpv6_packet_too_big + // Pass remote_addr as our_addr — this is what IcmpContext::too_big // does after the fix (original packet's dst = remote peer). let response = build_packet_too_big(&original, 1203, remote_addr); assert!(response.is_some()); diff --git a/src/upper/icmp_rate_limit.rs b/src/ipv6tun/icmp_rate_limit.rs similarity index 100% rename from src/upper/icmp_rate_limit.rs rename to src/ipv6tun/icmp_rate_limit.rs diff --git a/src/upper/ipv6_shim.rs b/src/ipv6tun/ipv6_shim.rs similarity index 100% rename from src/upper/ipv6_shim.rs rename to src/ipv6tun/ipv6_shim.rs diff --git a/src/ipv6tun/lifecycle.rs b/src/ipv6tun/lifecycle.rs new file mode 100644 index 00000000..16d89abe --- /dev/null +++ b/src/ipv6tun/lifecycle.rs @@ -0,0 +1,424 @@ +//! Start and stop of the TUN and DNS children. +//! +//! The node's supervisor decides when each child starts and stops, and in +//! what order. The bodies that bring the TUN device and the `.fips` DNS +//! responder up and take them down live here, together with the handles +//! they leave behind for the node to drive and, later, to tear down. + +use std::net::{IpAddr, SocketAddr}; +use std::path::PathBuf; +use std::thread::{self, JoinHandle}; + +use tokio::sync::mpsc::Sender; +use tracing::{debug, info, warn}; + +use super::config::{DnsConfig, TunConfig}; +use super::dns::{DnsIdentityRx, bind_dns_socket, lookup_mesh_ifindex, run_responder}; +use super::tun::{ + MssCeiling, TunDevice, TunError, TunOutboundRx, TunTx, run_tun_reader, shutdown_tun_interface, +}; +use crate::FipsAddress; +use crate::config::PeerConfig; +use crate::hosts::{DEFAULT_HOSTS_PATH, HostMap, HostMapReloader}; +use crate::node::lifecycle::supervisor::Child; +use crate::node::lifecycle::{report_exit, report_thread}; +use crate::node::path_mtu::PathMtuLookup; + +/// Runtime handles of the TUN and DNS children, held by the node's +/// supervisor. +/// +/// Every field is empty until its child starts. The TUN channels are the +/// exception: an embedder that owns the TUN installs them before start +/// through `Node::enable_app_owned_tun`, without a device name. +#[derive(Default)] +pub(crate) struct Handles { + /// Kernel name of the TUN device this node created, used to delete or + /// down it at teardown. Set only while a system TUN is up; an app-owned + /// TUN leaves it unset. + pub(crate) tun_name: Option, + /// TUN packet sender channel. + pub(crate) tun_tx: Option, + /// Receiver for outbound packets from the TUN reader. + pub(crate) tun_outbound_rx: Option, + /// TUN reader thread handle. + pub(crate) tun_reader_handle: Option>, + /// TUN writer thread handle. + pub(crate) tun_writer_handle: Option>, + /// Shutdown pipe: writing to this fd unblocks the TUN reader thread on + /// macOS and FreeBSD. On Linux, deleting the interface via netlink + /// serves the same purpose. + #[cfg(any(target_os = "macos", target_os = "freebsd"))] + pub(crate) tun_shutdown_fd: Option, + + /// Receiver for resolved identities from the DNS responder. + pub(crate) dns_identity_rx: Option, + /// DNS responder task handle. + pub(crate) dns_task: Option>, + /// Address the DNS responder actually bound, read back from the socket + /// after `bind` so a port-0 config resolves to the assigned port. `Some` + /// only while the responder is up; published to embedders through + /// [`Node::dns_local_addr`](crate::Node::dns_local_addr). + pub(crate) dns_local_addr: Option, + /// Sends a new peer-alias base to the running DNS responder; `Some` only + /// while it runs. Written through [`Handles::publish_aliases`]. + dns_aliases: Option>, +} + +/// What the TUN threads take from the node when they start. +pub(crate) struct TunThreads { + /// Shared TCP MSS ceiling, already refreshed by the node. Both threads + /// read it per packet from then on. + pub(crate) ceiling: MssCeiling, + /// The node's effective IPv6 MTU at start, logged with the ceiling. + pub(crate) effective: u16, + /// Per-destination path MTU, read by the clamp in both threads. + pub(crate) path_mtu: PathMtuLookup, + /// Capacity of the host-to-mesh channel (`node.buffers.tun_channel`). + pub(crate) channel: usize, + /// Sender each TUN thread reports `Child::Tun` on when it exits. + pub(crate) exit_tx: Option>, +} + +/// Create the TUN device, logging it, or log the failure and return `None`. +/// +/// A failure here is not fatal to the node: it continues without a TUN. +pub(crate) async fn open_tun(config: &TunConfig, address: FipsAddress) -> Option { + match TunDevice::create(config, address).await { + Ok(device) => { + info!("TUN device active:"); + info!(" name: {}", device.name()); + info!(" address: {}", device.address()); + info!(" mtu: {}", device.mtu()); + Some(device) + } + Err(e) => { + warn!(error = %e, "Failed to initialize TUN, continuing without it"); + None + } + } +} + +impl Handles { + /// Whether the TUN child is up. + /// + /// Follows the device name, not the TUN sender: an app-owned TUN + /// installs the sender but creates no device, so there is no TUN child + /// to tear down. + pub(crate) fn tun_up(&self) -> bool { + self.tun_name.is_some() + } + + /// Whether the DNS child is up: its responder task handle exists. + pub(crate) fn dns_up(&self) -> bool { + self.dns_task.is_some() + } + + /// Resolve the index of the mesh TUN device this node actually created. + /// + /// Reads the device name recorded when the TUN was brought up, which is + /// the kernel's name rather than the configured one. Returns `None` when + /// no TUN is up, which disables the DNS responder's mesh-interface + /// filter: with no mesh interface there is no mesh exposure to defend. + /// An app-owned TUN also leaves the name unset, so the filter stays off + /// there even though a mesh interface exists. + pub(crate) fn mesh_ifindex(&self) -> Option { + self.tun_name.as_deref().and_then(lookup_mesh_ifindex) + } + + /// Start the TUN reader and writer threads on an opened device and keep + /// their handles. + /// + /// An error here (the macOS/FreeBSD shutdown pipe, or duplicating the + /// device fd for the writer) is fatal to the node's start, unlike a + /// failure to create the device. + pub(crate) fn spawn_tun( + &mut self, + device: TunDevice, + threads: TunThreads, + ) -> Result<(), TunError> { + let TunThreads { + ceiling: max_mss, + effective: effective_mtu, + path_mtu: path_mtu_lookup, + channel: tun_channel_size, + exit_tx, + } = threads; + let mtu = device.mtu(); + let name = device.name().to_string(); + let our_addr = *device.address(); + + info!("effective MTU: {} bytes", effective_mtu); + debug!( + " max TCP MSS: {} bytes", + max_mss.load(std::sync::atomic::Ordering::Relaxed) + ); + + // On macOS and FreeBSD, create a shutdown pipe. Writing to it + // unblocks the reader thread's select() loop without closing + // the TUN fd (which would cause a double-close when TunDevice + // drops). Linux instead unblocks the reader by deleting the + // interface; on macOS/FreeBSD downing the interface does not + // wake a blocked read. + #[cfg(any(target_os = "macos", target_os = "freebsd"))] + let (shutdown_read_fd, shutdown_write_fd) = { + let mut fds = [0i32; 2]; + if unsafe { libc::pipe(fds.as_mut_ptr()) } < 0 { + return Err(TunError::Configure("failed to create shutdown pipe".into())); + } + (fds[0], fds[1]) + }; + + // Create writer (dups the fd for independent write access). + // Pass path_mtu_lookup so inbound SYN-ACK clamp can read + // per-destination path MTU learned via discovery. + let (writer, tun_tx) = device.create_writer(max_mss.clone(), path_mtu_lookup.clone())?; + + // Spawn writer thread. On exit, including a panic, + // it self-reports `Child::Tun` (sync context → + // `blocking_send`); TUN is one compound child, so + // both threads reporting is fine (the FSM de-dups + // via `up.remove`). + let writer_child_tx = exit_tx.clone(); + let writer_handle = thread::spawn(move || { + report_thread(Child::Tun, move || writer.run(), writer_child_tx.as_ref()); + }); + + // Clone tun_tx for the reader + let reader_tun_tx = tun_tx.clone(); + + // Create outbound channel for TUN reader → Node + let (outbound_tx, outbound_rx) = tokio::sync::mpsc::channel(tun_channel_size); + + // Spawn reader thread. Like the writer, it + // self-reports `Child::Tun` on exit or panic (sync + // context → `blocking_send`). Exactly one cfg + // variant compiles, so the exit sender is moved + // into that closure. + let reader_child_tx = exit_tx; + #[cfg(any(target_os = "macos", target_os = "freebsd"))] + let reader_handle = thread::spawn(move || { + report_thread( + Child::Tun, + move || { + run_tun_reader( + device, + mtu, + our_addr, + reader_tun_tx, + outbound_tx, + max_mss, + path_mtu_lookup, + shutdown_read_fd, + ) + }, + reader_child_tx.as_ref(), + ); + }); + #[cfg(not(any(target_os = "macos", target_os = "freebsd")))] + let reader_handle = thread::spawn(move || { + report_thread( + Child::Tun, + move || { + run_tun_reader( + device, + mtu, + our_addr, + reader_tun_tx, + outbound_tx, + max_mss, + path_mtu_lookup, + ) + }, + reader_child_tx.as_ref(), + ); + }); + + self.tun_name = Some(name); + self.tun_tx = Some(tun_tx); + self.tun_outbound_rx = Some(outbound_rx); + self.tun_reader_handle = Some(reader_handle); + self.tun_writer_handle = Some(writer_handle); + #[cfg(any(target_os = "macos", target_os = "freebsd"))] + { + self.tun_shutdown_fd = Some(shutdown_write_fd); + } + Ok(()) + } + + /// Stop the TUN child: close the writer's channel, delete or down the + /// interface, wake the reader, and join both threads. + /// + /// Returns whether there was a TUN child to stop. With no device name + /// (never started, or app-owned) it does nothing, and an app-owned + /// sender is left in place. + pub(crate) async fn stop_tun(&mut self) -> bool { + let Some(name) = self.tun_name.take() else { + return false; + }; + info!(name = %name, "Shutting down TUN interface"); + + // Drop the tun_tx to signal the writer to stop + self.tun_tx.take(); + + // Delete the interface (on Linux, causes reader to get + // EFAULT; on macOS/FreeBSD this downs it — the kernel + // destroys the device once the reader closes the fd). + if let Err(e) = shutdown_tun_interface(&name).await { + warn!(name = %name, error = %e, "Failed to shutdown TUN interface"); + } + + // On macOS and FreeBSD, signal the reader thread to exit by + // writing to the shutdown pipe. The reader's select() will + // wake up and break. + #[cfg(any(target_os = "macos", target_os = "freebsd"))] + if let Some(fd) = self.tun_shutdown_fd.take() { + unsafe { + libc::write(fd, b"x".as_ptr() as *const libc::c_void, 1); + libc::close(fd); + } + } + + // Wait for threads to finish + if let Some(handle) = self.tun_reader_handle.take() { + let _ = handle.join(); + } + if let Some(handle) = self.tun_writer_handle.take() { + let _ = handle.join(); + } + true + } + + /// Start the `.fips` DNS responder and keep its handles, returning + /// whether it came up. A failure is logged and is not fatal to the node. + /// + /// `peers` seeds the responder's own hosts map, which it reloads from + /// the hosts file on its own; `channel` is the capacity of the identity + /// channel back to the node (`node.buffers.dns_channel`). Reads the TUN + /// device name for the mesh-interface filter, so the TUN child, when + /// enabled, starts first. + pub(crate) fn start_dns( + &mut self, + config: &DnsConfig, + peers: &[PeerConfig], + channel: usize, + exit_tx: Option>, + ) -> bool { + // Initialize DNS responder (independent of TUN). + // + // Default bind_addr is "::1" (IPv6 loopback). The shipped + // fips-dns-setup configures systemd-resolved via a global + // /etc/systemd/resolved.conf.d/fips.conf drop-in pointing at + // [::1]:5354, which sidesteps a Linux IPV6_PKTINFO behaviour + // where self-destined traffic to fips0's address is attributed + // to fips0 in PKTINFO and gets silently dropped by the + // mesh-interface filter in src/ipv6tun/dns.rs. + // + // For mesh-reachable resolution (rare), set bind_addr: "::" + // in fips.yaml. The mesh-interface filter remains active to + // prevent hosts-file alias enumeration in that mode. + // `IPV6_V6ONLY=0` is set explicitly so IPv4 clients on + // 127.0.0.1 still reach us regardless of kernel sysctl + // defaults — but only when bind is on a wildcard / IPv6 path. + let addr_str = config.bind_addr(); + match addr_str.parse::() { + Ok(ip) => { + let bind = SocketAddr::new(ip, config.port()); + match bind_dns_socket(bind) { + Ok(socket) => { + // Read the bound address back off the socket + // rather than reusing `bind`: a port-0 config + // resolves to the kernel-assigned port here, + // and this is the address an embedder that + // proxies queries to us has to dial. + let local_addr = socket.local_addr().unwrap_or(bind); + let (identity_tx, identity_rx) = tokio::sync::mpsc::channel(channel); + let dns_ttl = config.ttl(); + let base_hosts = HostMap::from_peer_configs(peers); + let hosts_path = PathBuf::from(DEFAULT_HOSTS_PATH); + let (aliases_tx, aliases_rx) = + tokio::sync::watch::channel(base_hosts.clone()); + let reloader = HostMapReloader::new(base_hosts, hosts_path); + // Resolve the TUN ifindex so the responder can + // drop queries arriving on the mesh interface + // Without this, the `::` bind exposes the + // hosts file's alias space to any mesh peer. + // The name comes from the device the TUN + // path actually created, not the configured + // one: macOS and FreeBSD assign utunN/tunN + // of their own choosing and the configured + // name resolves to nothing there, which left + // the filter permanently off. + let mesh_ifindex = self.mesh_ifindex(); + if self.tun_name.is_some() && mesh_ifindex.is_none() { + warn!( + device = ?self.tun_name, + "Mesh interface index unresolved; DNS mesh filter disabled" + ); + } + info!( + bind = %local_addr, + hosts = reloader.hosts().len(), + mesh_ifindex = ?mesh_ifindex, + "DNS responder started for .fips domain (auto-reload enabled)" + ); + // Self-report on exit so the supervisor FSM + // routes health when the DNS task dies at + // runtime. The responder never returns, so + // in practice that is a panic. On a + // deliberate stop the task is `.abort()`ed, + // which drops the report with it; even if + // one fired, the FSM ignores it outside + // `Running`. + let handle = tokio::spawn(report_exit( + Child::Dns, + run_responder( + socket, + identity_tx, + dns_ttl, + reloader, + Some(aliases_rx), + mesh_ifindex, + ), + exit_tx, + )); + self.dns_identity_rx = Some(identity_rx); + self.dns_task = Some(handle); + self.dns_aliases = Some(aliases_tx); + self.dns_local_addr = Some(local_addr); + true + } + Err(e) => { + warn!(bind = %bind, error = %e, "Failed to start DNS responder"); + false + } + } + } + Err(e) => { + warn!(addr = %addr_str, error = %e, "Invalid dns.bind_addr; DNS responder not started"); + false + } + } + } + + /// Give the running DNS responder a new peer-alias base, which it applies + /// before answering its next query. Does nothing while it is not running. + pub(crate) fn publish_aliases(&self, base: HostMap) { + if let Some(tx) = &self.dns_aliases { + tx.send_replace(base); + } + } + + /// Stop the DNS responder and retract its published address. + pub(crate) fn stop_dns(&mut self) { + // Stop DNS responder + if let Some(handle) = self.dns_task.take() { + handle.abort(); + debug!("DNS responder stopped"); + } + self.dns_aliases.take(); + // Retract the published address in the same step that kills + // the listener, so an embedder polling `dns_local_addr()` + // never dials a socket that is already gone. + self.dns_local_addr.take(); + } +} diff --git a/src/upper/mod.rs b/src/ipv6tun/mod.rs similarity index 67% rename from src/upper/mod.rs rename to src/ipv6tun/mod.rs index 99488a66..9060595c 100644 --- a/src/upper/mod.rs +++ b/src/ipv6tun/mod.rs @@ -7,9 +7,14 @@ pub mod config; pub mod dns; -pub mod hosts; pub mod icmp; pub mod icmp_rate_limit; pub mod ipv6_shim; +pub(crate) mod lifecycle; +pub(crate) mod outbound; pub mod tcp_mss; pub mod tun; + +// The hosts file moved to `crate::hosts`; re-exported so `crate::upper::hosts` +// and `fips::upper::hosts` still resolve. +pub use crate::hosts; diff --git a/src/ipv6tun/outbound.rs b/src/ipv6tun/outbound.rs new file mode 100644 index 00000000..cd7cc3fe --- /dev/null +++ b/src/ipv6tun/outbound.rs @@ -0,0 +1,302 @@ +//! Host-side half of forwarding an IPv6 packet read from the TUN. +//! +//! This side checks the packet, answers the host with ICMPv6 when the mesh +//! cannot take it, and turns the destination address into the 15-byte +//! prefix the mesh resolves. What happens to the packet inside the mesh +//! (session lookup, queueing while a session is set up, discovery) is the +//! mesh's business, reached through the [`Mesh`] trait. +//! +//! The decisions are the synchronous functions [`admit`] and +//! [`path_limit`]; [`forward`] drives them against a [`Mesh`]. + +use super::icmp::{IcmpContext, effective_ipv6_mtu}; +use std::future::Future; + +/// What the outbound path needs from the mesh side. +pub(crate) trait Mesh { + /// The mesh's handle on a resolved destination, passed back to `send`. + type Dest; + + /// Largest IPv6 packet, header included, the mesh carries on its + /// narrowest transport. + fn ipv6_mtu(&self) -> u16; + + /// Resolve a destination from bytes 1-15 of its IPv6 address, with the + /// session's current path MTU if a session to it is established. + fn resolve(&mut self, prefix: &[u8; 15]) -> Option>; + + /// Send the packet on the destination's session, or hold it while one + /// is set up. + fn send(&mut self, dest: Self::Dest, packet: Vec) -> impl Future + Send; + + /// Borrow the sender of ICMPv6 replies to the host. + fn icmp(&mut self) -> IcmpContext<'_>; +} + +/// A destination the mesh resolved from an address prefix. +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct Route { + /// The mesh's handle on the destination. + pub dest: D, + /// Current path MTU of an established session to it, `None` when no + /// session is established or it has no path MTU state. + pub path_mtu: Option, +} + +/// What the mesh did with a packet handed to [`Mesh::send`]. +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) enum Outcome { + /// Handed to the established session. A send error is logged by the + /// mesh and not answered to the host. + Sent, + /// Held until the destination's session is established. + Queued, + /// Refused because the session table is full; the packet comes back so + /// the host can be told the destination is unreachable. + TableFull(Vec), +} + +/// The verdict on a packet read from the TUN, before the mesh is asked. +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) enum Admit { + /// Not an IPv6 packet, or shorter than its fixed header. + Drop, + /// Larger than the mesh carries; answer Packet Too Big with this MTU. + TooBig(u32), + /// Forward to the destination with this address prefix. + Forward([u8; 15]), +} + +/// Check a packet from the TUN against the node-wide IPv6 MTU and extract +/// its destination prefix. +pub(crate) fn admit(packet: &[u8], ipv6_mtu: u16) -> Admit { + if packet.len() < 40 || packet[0] >> 4 != 6 { + return Admit::Drop; + } + + // Check if packet will fit after FIPS encapsulation + let effective_mtu = ipv6_mtu as usize; + if packet.len() > effective_mtu { + return Admit::TooBig(effective_mtu as u32); + } + + // Extract destination FipsAddress prefix (IPv6 dest bytes 1-15) + // IPv6 header: bytes 24-39 are dest addr, so prefix = bytes 25-39 + let mut prefix = [0u8; 15]; + prefix.copy_from_slice(&packet[25..40]); + Admit::Forward(prefix) +} + +/// The Packet Too Big MTU for a packet of `len` bytes on a session whose +/// path MTU is `path_mtu`, or `None` if it fits. +/// +/// Applies only when the path is narrower than the node-wide `ipv6_mtu`, +/// which [`admit`] has already enforced. +pub(crate) fn path_limit(len: usize, path_mtu: u16, ipv6_mtu: u16) -> Option { + let path_ipv6_mtu = effective_ipv6_mtu(path_mtu) as usize; + if path_ipv6_mtu < ipv6_mtu as usize && len > path_ipv6_mtu { + Some(path_ipv6_mtu as u32) + } else { + None + } +} + +/// Forward one IPv6 packet read from the TUN into the mesh. +/// +/// Packets that are not IPv6, too large for the node, too large for an +/// established session's path, to an unknown destination, or refused for a +/// full session table are answered or dropped here. The rest go to +/// [`Mesh::send`]. +pub(crate) async fn forward(mesh: &mut M, packet: Vec) { + let ipv6_mtu = mesh.ipv6_mtu(); + let prefix = match admit(&packet, ipv6_mtu) { + Admit::Drop => return, + Admit::TooBig(mtu) => { + mesh.icmp().too_big(&packet, mtu); + return; + } + Admit::Forward(prefix) => prefix, + }; + + let Some(route) = mesh.resolve(&prefix) else { + mesh.icmp().dest_unreachable(&packet); + return; + }; + + // Check per-destination path MTU learned from MtuExceeded signals. + // The first oversized packet is forwarded normally and triggers + // the MtuExceeded signal; subsequent packets are caught here and + // generate ICMPv6 Packet Too Big back to the application. + if let Some(mtu) = route + .path_mtu + .and_then(|path_mtu| path_limit(packet.len(), path_mtu, ipv6_mtu)) + { + mesh.icmp().too_big(&packet, mtu); + return; + } + + if let Outcome::TableFull(packet) = mesh.send(route.dest, packet).await { + mesh.icmp().dest_unreachable(&packet); + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::ipv6tun::icmp_rate_limit::IcmpRateLimiter; + use crate::ipv6tun::tun::TunTx; + use std::net::Ipv6Addr; + use std::sync::mpsc; + + /// A mesh with one known destination prefix and a scripted send outcome. + struct FakeMesh { + ipv6_mtu: u16, + known: [u8; 15], + path_mtu: Option, + refuse: bool, + sent: Vec>, + tun_tx: TunTx, + limiter: IcmpRateLimiter, + } + + impl Mesh for FakeMesh { + type Dest = (); + + fn ipv6_mtu(&self) -> u16 { + self.ipv6_mtu + } + + fn resolve(&mut self, prefix: &[u8; 15]) -> Option> { + (*prefix == self.known).then_some(Route { + dest: (), + path_mtu: self.path_mtu, + }) + } + + fn send(&mut self, _dest: (), packet: Vec) -> impl Future + Send { + let outcome = if self.refuse { + Outcome::TableFull(packet) + } else { + self.sent.push(packet); + Outcome::Sent + }; + async move { outcome } + } + + fn icmp(&mut self) -> IcmpContext<'_> { + IcmpContext::new( + Some(&self.tun_tx), + "fd00::1".parse().unwrap(), + &mut self.limiter, + ) + } + } + + /// Build a fake mesh and the receiving end of its TUN channel. + fn fake(ipv6_mtu: u16) -> (FakeMesh, mpsc::Receiver>) { + let (tun_tx, rx) = mpsc::channel(); + let mesh = FakeMesh { + ipv6_mtu, + known: dest().octets()[1..16].try_into().unwrap(), + path_mtu: None, + refuse: false, + sent: Vec::new(), + tun_tx, + limiter: IcmpRateLimiter::new(), + }; + (mesh, rx) + } + + /// The destination the fake mesh knows. + fn dest() -> Ipv6Addr { + "fd12:3456:789a::2".parse().unwrap() + } + + /// An IPv6 UDP packet of `len` bytes from a host address to `dst`. + fn packet(dst: Ipv6Addr, len: usize) -> Vec { + let mut p = vec![0u8; len]; + p[0] = 0x60; + p[4..6].copy_from_slice(&((len - 40) as u16).to_be_bytes()); + p[6] = 17; + p[7] = 64; + let src: Ipv6Addr = "fd00::5".parse().unwrap(); + p[8..24].copy_from_slice(&src.octets()); + p[24..40].copy_from_slice(&dst.octets()); + p + } + + /// The ICMPv6 type of a reply written to the TUN. + fn icmp_type(reply: &[u8]) -> u8 { + reply[40] + } + + #[test] + fn admit_drops_short_and_non_ipv6_packets() { + assert_eq!(admit(&[0x60; 39], 1280), Admit::Drop); + let mut p = packet(dest(), 60); + p[0] = 0x45; + assert_eq!(admit(&p, 1280), Admit::Drop); + } + + #[test] + fn admit_reports_the_node_mtu_for_an_oversized_packet() { + assert_eq!(admit(&packet(dest(), 1281), 1280), Admit::TooBig(1280)); + assert!(matches!( + admit(&packet(dest(), 1280), 1280), + Admit::Forward(_) + )); + } + + #[test] + fn admit_extracts_destination_bytes_one_to_fifteen() { + let expected: [u8; 15] = dest().octets()[1..16].try_into().unwrap(); + assert_eq!(admit(&packet(dest(), 60), 1280), Admit::Forward(expected)); + } + + #[test] + fn path_limit_applies_only_below_the_node_mtu() { + let node = 1280; + let narrow = effective_ipv6_mtu(1000); + assert_eq!( + path_limit(narrow as usize + 1, 1000, node), + Some(narrow as u32) + ); + assert_eq!(path_limit(narrow as usize, 1000, node), None); + // A path at least as wide as the node MTU never answers. + assert_eq!(path_limit(1280, 1280 + 77, node), None); + } + + #[tokio::test] + async fn forward_answers_unknown_destination_with_unreachable() { + let (mut mesh, rx) = fake(1280); + forward(&mut mesh, packet("fd99::1".parse().unwrap(), 60)).await; + assert!(mesh.sent.is_empty()); + assert_eq!(icmp_type(&rx.try_recv().unwrap()), 1); + } + + #[tokio::test] + async fn forward_answers_oversized_for_session_path_with_packet_too_big() { + let (mut mesh, rx) = fake(1280); + mesh.path_mtu = Some(1000); + forward(&mut mesh, packet(dest(), 1200)).await; + assert!(mesh.sent.is_empty()); + assert_eq!(icmp_type(&rx.try_recv().unwrap()), 2); + } + + #[tokio::test] + async fn forward_sends_a_fitting_packet_without_a_reply() { + let (mut mesh, rx) = fake(1280); + mesh.path_mtu = Some(1400); + forward(&mut mesh, packet(dest(), 1200)).await; + assert_eq!(mesh.sent.len(), 1); + assert!(rx.try_recv().is_err()); + } + + #[tokio::test] + async fn forward_answers_a_full_session_table_with_unreachable() { + let (mut mesh, rx) = fake(1280); + mesh.refuse = true; + forward(&mut mesh, packet(dest(), 60)).await; + assert_eq!(icmp_type(&rx.try_recv().unwrap()), 1); + } +} diff --git a/src/upper/tcp_mss.rs b/src/ipv6tun/tcp_mss.rs similarity index 100% rename from src/upper/tcp_mss.rs rename to src/ipv6tun/tcp_mss.rs diff --git a/src/upper/tun.rs b/src/ipv6tun/tun.rs similarity index 100% rename from src/upper/tun.rs rename to src/ipv6tun/tun.rs diff --git a/src/lib.rs b/src/lib.rs index 6f32de3b..256f2cab 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -14,12 +14,14 @@ pub mod config; pub mod control; #[cfg(target_os = "linux")] pub mod gateway; +pub mod hosts; pub mod identity; // Declared before `node` (and named to sort there) because it carries // `#[macro_use]`: the tick instrumentation macro must be in scope for the // modules that follow. #[macro_use] pub(crate) mod instr; +pub mod ipv6tun; pub mod mdns; pub mod native; pub mod node; @@ -34,10 +36,13 @@ pub(crate) mod proto; pub(crate) mod testutil; mod time; pub mod transport; -pub mod upper; pub mod utils; pub mod version; +// `upper` is the former name of `ipv6tun`; the alias keeps `crate::upper::` +// and `fips::upper::` paths resolving. +pub use ipv6tun as upper; + // Re-export identity types pub use identity::{ AuthChallenge, AuthResponse, FipsAddress, Identity, IdentityError, NodeAddr, PeerIdentity, diff --git a/src/node/dataplane/rx_loop.rs b/src/node/dataplane/rx_loop.rs index de3f6a81..1e3f6c38 100644 --- a/src/node/dataplane/rx_loop.rs +++ b/src/node/dataplane/rx_loop.rs @@ -79,7 +79,8 @@ impl Node { // Take the TUN outbound receiver, or create a dummy channel that never // produces messages (when TUN is disabled). Holding the sender prevents // the channel from closing. - let (mut tun_outbound_rx, _tun_guard) = match self.supervisor.tun_outbound_rx.take() { + let (mut tun_outbound_rx, _tun_guard) = match self.supervisor.ipv6tun.tun_outbound_rx.take() + { Some(rx) => (rx, None), None => { let (tx, rx) = tokio::sync::mpsc::channel(1); @@ -89,7 +90,8 @@ impl Node { // Take the DNS identity receiver, or create a dummy channel (when DNS // is disabled). Same pattern as TUN outbound. - let (mut dns_identity_rx, _dns_guard) = match self.supervisor.dns_identity_rx.take() { + let (mut dns_identity_rx, _dns_guard) = match self.supervisor.ipv6tun.dns_identity_rx.take() + { Some(rx) => (rx, None), None => { let (tx, rx) = tokio::sync::mpsc::channel(1); diff --git a/src/node/handlers/lookup.rs b/src/node/handlers/lookup.rs index 37beeaae..e86d3846 100644 --- a/src/node/handlers/lookup.rs +++ b/src/node/handlers/lookup.rs @@ -820,9 +820,7 @@ impl Node { "Discovery lookup timed out, destination unreachable" ); if let Some(packets) = queued { - for pkt in &packets { - self.send_icmpv6_dest_unreachable(pkt); - } + self.host_icmp().no_route(&packets); } } } diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index cff3e0e4..703edffc 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -6,6 +6,8 @@ //! encrypted data, and error signals (CoordsRequired, PathBroken). use crate::NodeAddr; +use crate::ipv6tun::icmp::IcmpContext; +use crate::ipv6tun::outbound::{Mesh, Outcome, Route}; use crate::node::handlers::mmp::format_throughput; use crate::node::rate_limit::Msg1Class; use crate::node::reject::{RejectReason, SessionReject}; @@ -453,7 +455,7 @@ impl Node { mark_ipv6_ecn_ce(&mut packet); self.metrics().congestion.ce_received.inc(); } - if let Some(tun_tx) = &self.supervisor.tun_tx { + if let Some(tun_tx) = &self.supervisor.ipv6tun.tun_tx { if let Err(e) = tun_tx.send(packet) { debug!(error = %e, "Failed to deliver decompressed IPv6 packet to TUN"); } @@ -3449,63 +3451,44 @@ impl Node { /// Handle an outbound IPv6 packet from the TUN reader. /// - /// Extracts the destination FipsAddress, looks up the NodeAddr and PublicKey - /// from the identity cache, and either sends through an established session - /// or initiates a new one (queuing the packet until established). + /// The host-side checks and ICMPv6 replies are in + /// [`crate::ipv6tun::outbound::forward`], which reaches back into the + /// node through its [`Mesh`] impl: the destination prefix is resolved + /// in the identity cache, and the packet is sent on an established + /// session or queued while one is set up. /// /// Also performs MTU checking: if the packet (plus FIPS overhead) exceeds /// the transport MTU, an ICMP Packet Too Big message is sent back to the /// source and the packet is dropped. pub(in crate::node) async fn handle_tun_outbound(&mut self, ipv6_packet: Vec) { - // Validate IPv6 header - if ipv6_packet.len() < 40 || ipv6_packet[0] >> 4 != 6 { - return; - } - - // Check if packet will fit after FIPS encapsulation - let effective_mtu = self.effective_ipv6_mtu() as usize; - if ipv6_packet.len() > effective_mtu { - self.send_icmpv6_packet_too_big(&ipv6_packet, effective_mtu as u32); - return; - } - - // Extract destination FipsAddress prefix (IPv6 dest bytes 1-15) - // IPv6 header: bytes 24-39 are dest addr, so prefix = bytes 25-39 - let mut prefix = [0u8; 15]; - prefix.copy_from_slice(&ipv6_packet[25..40]); - - // Look up in identity cache - let (dest_addr, dest_pubkey) = match self.lookup_by_fips_prefix(&prefix) { - Some((addr, pk)) => (addr, pk), - None => { - self.send_icmpv6_dest_unreachable(&ipv6_packet); - return; - } - }; + crate::ipv6tun::outbound::forward(self, ipv6_packet).await; + } + /// Send a TUN packet to a resolved destination, or queue it while the + /// destination's session is set up. + /// + /// Sends through an established session, queues behind one still being + /// set up, and otherwise initiates a session and queues. With no route + /// for the initiation it starts discovery and still queues. Refuses the + /// packet, handing it back, when a new session would exceed the session + /// table. + async fn send_outbound( + &mut self, + dest_addr: NodeAddr, + dest_pubkey: PublicKey, + ipv6_packet: Vec, + ) -> Outcome { // Check for established session if let Some(entry) = self.sessions.get(&dest_addr) { if entry.is_established() { - // Check per-destination path MTU learned from MtuExceeded signals. - // The first oversized packet is forwarded normally and triggers - // the MtuExceeded signal; subsequent packets are caught here and - // generate ICMPv6 Packet Too Big back to the application. - if let Some(mmp) = entry.mmp() { - let path_mtu = mmp.path_mtu.current_mtu(); - let path_ipv6_mtu = crate::upper::icmp::effective_ipv6_mtu(path_mtu) as usize; - if path_ipv6_mtu < effective_mtu && ipv6_packet.len() > path_ipv6_mtu { - self.send_icmpv6_packet_too_big(&ipv6_packet, path_ipv6_mtu as u32); - return; - } - } if let Err(e) = self.send_ipv6_packet(&dest_addr, &ipv6_packet).await { debug!(dest = %self.peer_display_name(&dest_addr), error = %e, "Failed to send TUN packet via session"); } - return; + return Outcome::Sent; } // Session exists but not yet established — queue the packet self.queue_pending_packet(dest_addr, ipv6_packet); - return; + return Outcome::Queued; } // No session, so this one would grow the table. Answer the local @@ -3515,8 +3498,7 @@ impl Node { // queued packet, which is outbound traffic on a node already at its // limit. if !self.admit_new_session(&dest_addr) { - self.send_icmpv6_dest_unreachable(&ipv6_packet); - return; + return Outcome::TableFull(ipv6_packet); } // No session: initiate one and queue the packet. @@ -3526,73 +3508,21 @@ impl Node { debug!(dest = %self.peer_display_name(&dest_addr), error = %e, "Failed to initiate session, trying discovery"); self.maybe_initiate_lookup(&dest_addr).await; self.queue_pending_packet(dest_addr, ipv6_packet); - return; + return Outcome::Queued; } self.queue_pending_packet(dest_addr, ipv6_packet); + Outcome::Queued } - /// Send ICMPv6 Destination Unreachable back through TUN. - pub(in crate::node) fn send_icmpv6_dest_unreachable(&self, original_packet: &[u8]) { - use crate::FipsAddress; - use crate::upper::icmp::{ - DestUnreachableCode, build_dest_unreachable, should_send_icmp_error, - }; - - if !should_send_icmp_error(original_packet) { - return; - } - - let our_ipv6 = FipsAddress::from_node_addr(self.node_addr()).to_ipv6(); - if let Some(response) = - build_dest_unreachable(original_packet, DestUnreachableCode::NoRoute, our_ipv6) - && let Some(tun_tx) = &self.supervisor.tun_tx - { - let _ = tun_tx.send(response); - } - } - - /// Send ICMPv6 Packet Too Big back through TUN. - /// - /// Rate-limited per source address to prevent ICMP floods from - /// misconfigured applications sending repeated oversized packets. - pub(in crate::node) fn send_icmpv6_packet_too_big(&mut self, original_packet: &[u8], mtu: u32) { - use crate::upper::icmp::build_packet_too_big; - use std::net::Ipv6Addr; - - // Extract source address for rate limiting - if original_packet.len() < 40 { - return; - } - // SAFETY: slice is exactly 16 bytes; length validated above (>= 40) - let src_addr = Ipv6Addr::from(<[u8; 16]>::try_from(&original_packet[8..24]).unwrap()); - - // Rate limit ICMP PTB messages per source - if !self.icmp_rate_limiter.should_send(src_addr) { - debug!( - src = %src_addr, - "Rate limiting ICMP Packet Too Big" - ); - return; - } - - // Use the original packet's *destination* as the ICMP source so the - // kernel sees the PTB coming from a remote router, not from itself. - // Linux ignores PTBs whose source matches a local address, which - // causes a PMTUD blackhole when both src and ICMP-src are local. - // SAFETY: slice is exactly 16 bytes; length validated above (>= 40) - let dest_addr = Ipv6Addr::from(<[u8; 16]>::try_from(&original_packet[24..40]).unwrap()); - if let Some(response) = build_packet_too_big(original_packet, mtu, dest_addr) - && let Some(tun_tx) = &self.supervisor.tun_tx - { - debug!( - original_src = %src_addr, - original_dst = %dest_addr, - packet_size = original_packet.len(), - reported_mtu = mtu, - "Sending ICMP Packet Too Big" - ); - let _ = tun_tx.send(response); - } + /// Borrow the host-facing ICMPv6 sender: the TUN channel, our address + /// and the Packet Too Big rate limiter. + pub(in crate::node) fn host_icmp(&mut self) -> IcmpContext<'_> { + let our_ipv6 = crate::FipsAddress::from_node_addr(self.node_addr()).to_ipv6(); + IcmpContext::new( + self.supervisor.ipv6tun.tun_tx.as_ref(), + our_ipv6, + &mut self.icmp_rate_limiter, + ) } /// Queue a packet while waiting for session establishment. @@ -3662,3 +3592,36 @@ impl Node { } } } + +/// The mesh side of TUN outbound forwarding. +impl Mesh for Node { + type Dest = (NodeAddr, PublicKey); + + fn ipv6_mtu(&self) -> u16 { + self.effective_ipv6_mtu() + } + + fn resolve(&mut self, prefix: &[u8; 15]) -> Option> { + // Look up in identity cache + let (dest_addr, dest_pubkey) = self.lookup_by_fips_prefix(prefix)?; + let path_mtu = self + .sessions + .get(&dest_addr) + .filter(|entry| entry.is_established()) + .and_then(|entry| entry.mmp()) + .map(|mmp| mmp.path_mtu.current_mtu()); + Some(Route { + dest: (dest_addr, dest_pubkey), + path_mtu, + }) + } + + fn send(&mut self, dest: Self::Dest, packet: Vec) -> impl Future + Send { + let (dest_addr, dest_pubkey) = dest; + self.send_outbound(dest_addr, dest_pubkey, packet) + } + + fn icmp(&mut self) -> IcmpContext<'_> { + self.host_icmp() + } +} diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index 9f9eb687..7de6680c 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -12,6 +12,7 @@ use super::peering::reconcile::{ use super::peering::retry::MAX_RETRY_CONNECTIONS_PER_TICK; use crate::config::{ConnectPolicy, PeerAddress, PeerConfig}; +use crate::ipv6tun::lifecycle::{TunThreads, open_tun}; use crate::node::acl::PeerAclContext; use crate::node::dataplane::PeerActionCtx; use crate::nostr::{BootstrapEvent, NostrRendezvous}; @@ -20,11 +21,10 @@ use crate::peer::machine::{HandshakeCrypto, PeerEvent, PeerMachine}; use crate::proto::fmp::wire::build_msg1; use crate::proto::fmp::{Disconnect, DisconnectReason}; use crate::transport::{Link, LinkDirection, LinkId, TransportAddr, TransportId, packet_channel}; -use crate::upper::tun::{TunDevice, TunState, run_tun_reader, shutdown_tun_interface}; +use crate::upper::tun::TunState; use crate::{NodeAddr, PeerIdentity}; use std::collections::{HashMap, HashSet}; use std::net::SocketAddr; -use std::thread; use std::time::Duration; use tracing::{debug, error, info, warn}; @@ -45,7 +45,7 @@ fn socket_addr_families_compatible(local: SocketAddr, remote: SocketAddr) -> boo /// A panic would otherwise unwind past the report and leave the node healthy /// with the child gone. Aborting the task still reports nothing: the abort /// drops this whole future, so a deliberate stop does not read as a death. -pub(in crate::node) async fn report_exit( +pub(crate) async fn report_exit( child: Child, body: impl std::future::Future, tx: Option>, @@ -67,7 +67,7 @@ pub(in crate::node) async fn report_exit( /// Run a supervised child's thread body, then report the child's exit on `tx`, /// whether the body returned or panicked. -pub(in crate::node) fn report_thread( +pub(crate) fn report_thread( child: Child, body: impl FnOnce(), tx: Option<&tokio::sync::mpsc::Sender>, @@ -1732,7 +1732,7 @@ impl Node { // No Tun child when the TUN is app-owned (the embedder pre-set // `tun_tx` via `enable_app_owned_tun`) — FIPS does no system-TUN ops; // the channels installed before `start` carry both directions. - let tun = self.config().tun.enabled && self.supervisor.tun_tx.is_none(); + let tun = self.config().tun.enabled && self.supervisor.ipv6tun.tun_tx.is_none(); let dns = self.config().dns.enabled; // Worker-pool booleans + counts. Unix only — the workers issue @@ -1954,248 +1954,43 @@ impl Node { // Initialize TUN interface after transports and peers are // ready. let address = *self.identity().address(); - match TunDevice::create(&self.config().tun, address).await { - Ok(device) => { - let mtu = device.mtu(); - let name = device.name().to_string(); - let our_addr = *device.address(); - - info!("TUN device active:"); - info!(" name: {}", name); - info!(" address: {}", device.address()); - info!(" mtu: {}", mtu); - + match open_tun(&self.config().tun, address).await { + Some(device) => { // Seed the shared MSS ceiling from whatever is bound // right now. Both TUN threads read it live from here // on, so a transport binding or unbinding later moves // the clamp instead of leaving it at this instant's // value — see `crate::upper::tun::MssCeiling`. self.refresh_tun_mss_ceiling(); - let max_mss = self.tun_mss_ceiling.clone(); - let effective_mtu = self.effective_ipv6_mtu(); - - info!("effective MTU: {} bytes", effective_mtu); - debug!( - " max TCP MSS: {} bytes", - max_mss.load(std::sync::atomic::Ordering::Relaxed) - ); - - // On macOS and FreeBSD, create a shutdown pipe. Writing to it - // unblocks the reader thread's select() loop without closing - // the TUN fd (which would cause a double-close when TunDevice - // drops). Linux instead unblocks the reader by deleting the - // interface; on macOS/FreeBSD downing the interface does not - // wake a blocked read. - #[cfg(any(target_os = "macos", target_os = "freebsd"))] - let (shutdown_read_fd, shutdown_write_fd) = { - let mut fds = [0i32; 2]; - if unsafe { libc::pipe(fds.as_mut_ptr()) } < 0 { - return Err(NodeError::Tun( - crate::upper::tun::TunError::Configure( - "failed to create shutdown pipe".into(), - ), - )); - } - (fds[0], fds[1]) + let threads = TunThreads { + ceiling: self.tun_mss_ceiling.clone(), + effective: self.effective_ipv6_mtu(), + path_mtu: self.path_mtu_lookup.clone(), + channel: self.config().node.buffers.tun_channel, + exit_tx: self.child_exit_tx.clone(), }; - - // Create writer (dups the fd for independent write access). - // Pass path_mtu_lookup so inbound SYN-ACK clamp can read - // per-destination path MTU learned via discovery. - let (writer, tun_tx) = device - .create_writer(max_mss.clone(), self.path_mtu_lookup.clone())?; - - // Spawn writer thread. On exit, including a panic, - // it self-reports `Child::Tun` (sync context → - // `blocking_send`); TUN is one compound child, so - // both threads reporting is fine (the FSM de-dups - // via `up.remove`). - let writer_child_tx = self.child_exit_tx.clone(); - let writer_handle = thread::spawn(move || { - report_thread( - Child::Tun, - move || writer.run(), - writer_child_tx.as_ref(), - ); - }); - - // Clone tun_tx for the reader - let reader_tun_tx = tun_tx.clone(); - - // Create outbound channel for TUN reader → Node - let tun_channel_size = self.config().node.buffers.tun_channel; - let (outbound_tx, outbound_rx) = - tokio::sync::mpsc::channel(tun_channel_size); - - // Spawn reader thread. Like the writer, it - // self-reports `Child::Tun` on exit or panic (sync - // context → `blocking_send`). Exactly one cfg - // variant compiles, so the single clone is moved - // into that closure. - let path_mtu_lookup = self.path_mtu_lookup.clone(); - let reader_child_tx = self.child_exit_tx.clone(); - #[cfg(any(target_os = "macos", target_os = "freebsd"))] - let reader_handle = thread::spawn(move || { - report_thread( - Child::Tun, - move || { - run_tun_reader( - device, - mtu, - our_addr, - reader_tun_tx, - outbound_tx, - max_mss, - path_mtu_lookup, - shutdown_read_fd, - ) - }, - reader_child_tx.as_ref(), - ); - }); - #[cfg(not(any(target_os = "macos", target_os = "freebsd")))] - let reader_handle = thread::spawn(move || { - report_thread( - Child::Tun, - move || { - run_tun_reader( - device, - mtu, - our_addr, - reader_tun_tx, - outbound_tx, - max_mss, - path_mtu_lookup, - ) - }, - reader_child_tx.as_ref(), - ); - }); - + self.supervisor.ipv6tun.spawn_tun(device, threads)?; self.tun_state = TunState::Active; - self.tun_name = Some(name); - self.supervisor.tun_tx = Some(tun_tx); - self.supervisor.tun_outbound_rx = Some(outbound_rx); - self.supervisor.tun_reader_handle = Some(reader_handle); - self.supervisor.tun_writer_handle = Some(writer_handle); - #[cfg(any(target_os = "macos", target_os = "freebsd"))] - { - self.supervisor.tun_shutdown_fd = Some(shutdown_write_fd); - } Event::SubstrateUp { child } } - Err(e) => { + None => { self.tun_state = TunState::Failed; - warn!(error = %e, "Failed to initialize TUN, continuing without it"); Event::SubstrateFailed { child } } } } Child::Dns => { - // Initialize DNS responder (independent of TUN). - // - // Default bind_addr is "::1" (IPv6 loopback). The shipped - // fips-dns-setup configures systemd-resolved via a global - // /etc/systemd/resolved.conf.d/fips.conf drop-in pointing at - // [::1]:5354, which sidesteps a Linux IPV6_PKTINFO behaviour - // where self-destined traffic to fips0's address is attributed - // to fips0 in PKTINFO and gets silently dropped by the - // mesh-interface filter in src/upper/dns.rs. - // - // For mesh-reachable resolution (rare), set bind_addr: "::" - // in fips.yaml. The mesh-interface filter remains active to - // prevent hosts-file alias enumeration in that mode. - // `IPV6_V6ONLY=0` is set explicitly so IPv4 clients on - // 127.0.0.1 still reach us regardless of kernel sysctl - // defaults — but only when bind is on a wildcard / IPv6 path. - let addr_str = self.config().dns.bind_addr(); - match addr_str.parse::() { - Ok(ip) => { - let bind = std::net::SocketAddr::new(ip, self.config().dns.port()); - match Self::bind_dns_socket(bind) { - Ok(socket) => { - // Read the bound address back off the socket - // rather than reusing `bind`: a port-0 config - // resolves to the kernel-assigned port here, - // and this is the address an embedder that - // proxies queries to us has to dial. - let local_addr = socket.local_addr().unwrap_or(bind); - let dns_channel_size = self.config().node.buffers.dns_channel; - let (identity_tx, identity_rx) = - tokio::sync::mpsc::channel(dns_channel_size); - let dns_ttl = self.config().dns.ttl(); - let base_hosts = - crate::upper::hosts::HostMap::from_peer_configs( - self.config().peers(), - ); - let hosts_path = std::path::PathBuf::from( - crate::upper::hosts::DEFAULT_HOSTS_PATH, - ); - let (aliases_tx, aliases_rx) = - tokio::sync::watch::channel(base_hosts.clone()); - let reloader = crate::upper::hosts::HostMapReloader::new( - base_hosts, hosts_path, - ); - // Resolve the TUN ifindex so the responder can - // drop queries arriving on the mesh interface - // Without this, the `::` bind exposes the - // hosts file's alias space to any mesh peer. - // The name comes from the device the TUN - // path actually created, not the configured - // one: macOS and FreeBSD assign utunN/tunN - // of their own choosing and the configured - // name resolves to nothing there, which left - // the filter permanently off. - let mesh_ifindex = self.mesh_ifindex(); - if self.tun_name.is_some() && mesh_ifindex.is_none() { - warn!( - device = ?self.tun_name, - "Mesh interface index unresolved; DNS mesh filter disabled" - ); - } - info!( - bind = %local_addr, - hosts = reloader.hosts().len(), - mesh_ifindex = ?mesh_ifindex, - "DNS responder started for .fips domain (auto-reload enabled)" - ); - // Self-report on exit so the supervisor FSM - // routes health when the DNS task dies at - // runtime. The responder never returns, so - // in practice that is a panic. On a - // deliberate stop the task is `.abort()`ed, - // which drops the report with it; even if - // one fired, the FSM ignores it outside - // `Running`. - let dns_child_tx = self.child_exit_tx.clone(); - let handle = tokio::spawn(report_exit( - Child::Dns, - crate::upper::dns::run_responder( - socket, - identity_tx, - dns_ttl, - reloader, - Some(aliases_rx), - mesh_ifindex, - ), - dns_child_tx, - )); - self.supervisor.dns_identity_rx = Some(identity_rx); - self.supervisor.dns_task = Some(handle); - self.supervisor.dns_aliases = Some(aliases_tx); - self.supervisor.dns_local_addr = Some(local_addr); - Event::SubstrateUp { child } - } - Err(e) => { - warn!(bind = %bind, error = %e, "Failed to start DNS responder"); - Event::SubstrateFailed { child } - } - } - } - Err(e) => { - warn!(addr = %addr_str, error = %e, "Invalid dns.bind_addr; DNS responder not started"); - Event::SubstrateFailed { child } - } + let config = &self.context.config; + let up = self.supervisor.ipv6tun.start_dns( + &config.dns, + config.peers(), + config.node.buffers.dns_channel, + self.child_exit_tx.clone(), + ); + if up { + Event::SubstrateUp { child } + } else { + Event::SubstrateFailed { child } } } }; @@ -2324,91 +2119,6 @@ impl Node { Ok(()) } - /// Bind a UDP socket for the DNS responder. - /// - /// For IPv6 binds (including `::`), sets `IPV6_V6ONLY=0` so the socket - /// also accepts IPv4-mapped addresses. This guarantees dual-stack - /// delivery regardless of `net.ipv6.bindv6only` sysctl on the host — - /// v4 clients on 127.0.0.1 and v6 clients on the fips0 address both - /// land on the same socket. - /// - /// Also enables `IPV6_RECVPKTINFO` on IPv6 sockets so the responder - /// can learn the arrival interface per packet. The responder uses that - /// to drop queries arriving on the mesh TUN, closing the hosts-file - /// probing side-channel created by the `::` bind. - fn bind_dns_socket( - addr: std::net::SocketAddr, - ) -> Result { - use socket2::{Domain, Protocol, Socket, Type}; - let domain = if addr.is_ipv4() { - Domain::IPV4 - } else { - Domain::IPV6 - }; - let sock = Socket::new(domain, Type::DGRAM, Some(Protocol::UDP))?; - if addr.is_ipv6() { - sock.set_only_v6(false)?; - #[cfg(unix)] - Self::set_recv_pktinfo_v6(&sock)?; - } - sock.set_nonblocking(true)?; - sock.bind(&addr.into())?; - tokio::net::UdpSocket::from_std(sock.into()) - } - - /// Enable `IPV6_RECVPKTINFO` on an IPv6 UDP socket. - /// - /// After this setsockopt, each `recvmsg()` call on the socket receives - /// an `IPV6_PKTINFO` control message containing the arrival interface - /// index, which the DNS responder uses for its mesh-interface filter. - #[cfg(unix)] - fn set_recv_pktinfo_v6(sock: &socket2::Socket) -> Result<(), std::io::Error> { - use std::os::fd::AsRawFd; - let enable: libc::c_int = 1; - let ret = unsafe { - libc::setsockopt( - sock.as_raw_fd(), - libc::IPPROTO_IPV6, - libc::IPV6_RECVPKTINFO, - &enable as *const _ as *const libc::c_void, - std::mem::size_of::() as libc::socklen_t, - ) - }; - if ret < 0 { - return Err(std::io::Error::last_os_error()); - } - Ok(()) - } - - /// Resolve the index of the mesh TUN device this node actually created. - /// - /// Reads the device name recorded when the TUN was brought up, which is - /// the kernel's name rather than the configured one. Returns `None` when - /// no TUN is up, which disables the DNS responder's mesh-interface - /// filter: with no mesh interface there is no mesh exposure to defend. - /// An app-owned TUN also leaves the name unset, so the filter stays off - /// there even though a mesh interface exists. - pub(crate) fn mesh_ifindex(&self) -> Option { - self.tun_name.as_deref().and_then(Self::lookup_mesh_ifindex) - } - - /// Resolve an interface index by name. - /// - /// Returns `None` if the interface does not exist. - fn lookup_mesh_ifindex(name: &str) -> Option { - #[cfg(unix)] - { - let c_name = std::ffi::CString::new(name).ok()?; - let idx = unsafe { libc::if_nametoindex(c_name.as_ptr()) }; - if idx == 0 { None } else { Some(idx) } - } - #[cfg(not(unix))] - { - let _ = name; - None - } - } - /// Stop the node. /// /// Shuts down TUN interface, stops I/O threads, and transitions to @@ -2484,16 +2194,7 @@ impl Node { match child { Child::Dns => { - // Stop DNS responder - if let Some(handle) = self.supervisor.dns_task.take() { - handle.abort(); - debug!("DNS responder stopped"); - } - self.supervisor.dns_aliases.take(); - // Retract the published address in the same step that kills - // the listener, so an embedder polling `dns_local_addr()` - // never dials a socket that is already gone. - self.supervisor.dns_local_addr.take(); + self.supervisor.ipv6tun.stop_dns(); } Child::Nostr => { // Stop Nostr overlay discovery background work and withdraw @@ -2533,38 +2234,7 @@ impl Node { } Child::Tun => { // Shutdown TUN interface - if let Some(name) = self.tun_name.take() { - info!(name = %name, "Shutting down TUN interface"); - - // Drop the tun_tx to signal the writer to stop - self.supervisor.tun_tx.take(); - - // Delete the interface (on Linux, causes reader to get - // EFAULT; on macOS/FreeBSD this downs it — the kernel - // destroys the device once the reader closes the fd). - if let Err(e) = shutdown_tun_interface(&name).await { - warn!(name = %name, error = %e, "Failed to shutdown TUN interface"); - } - - // On macOS and FreeBSD, signal the reader thread to exit by - // writing to the shutdown pipe. The reader's select() will - // wake up and break. - #[cfg(any(target_os = "macos", target_os = "freebsd"))] - if let Some(fd) = self.supervisor.tun_shutdown_fd.take() { - unsafe { - libc::write(fd, b"x".as_ptr() as *const libc::c_void, 1); - libc::close(fd); - } - } - - // Wait for threads to finish - if let Some(handle) = self.supervisor.tun_reader_handle.take() { - let _ = handle.join(); - } - if let Some(handle) = self.supervisor.tun_writer_handle.take() { - let _ = handle.join(); - } - + if self.supervisor.ipv6tun.stop_tun().await { self.tun_state = TunState::Disabled; } } @@ -2622,7 +2292,7 @@ impl Node { /// is what reaches this. pub(in crate::node) fn retract_child_publications(&mut self, child: Child) { if matches!(child, Child::Dns) { - self.supervisor.dns_local_addr.take(); + self.supervisor.ipv6tun.dns_local_addr.take(); } } @@ -2688,7 +2358,7 @@ impl Node { /// stops them. Shared by [`Self::stop`] and [`Self::enter_drain`]. fn reconstruct_supervised_up(&self) -> Vec { let mut up: Vec = Vec::new(); - if self.supervisor.dns_task.is_some() { + if self.supervisor.ipv6tun.dns_up() { up.push(Child::Dns); } if self.supervisor.nostr_rendezvous.engine().is_some() { @@ -2700,7 +2370,7 @@ impl Node { for id in self.transports.keys() { up.push(Child::Transport(*id)); } - if self.tun_name.is_some() { + if self.supervisor.ipv6tun.tun_up() { up.push(Child::Tun); } up diff --git a/src/node/lifecycle/supervisor.rs b/src/node/lifecycle/supervisor.rs index edb4e658..e5d6bee7 100644 --- a/src/node/lifecycle/supervisor.rs +++ b/src/node/lifecycle/supervisor.rs @@ -104,11 +104,9 @@ use std::collections::HashSet; use std::sync::Arc; -use std::thread::JoinHandle; use crate::node::NodeState; use crate::transport::{PacketTx, TransportId}; -use crate::upper::tun::{TunOutboundRx, TunTx}; /// A supervised substrate child. /// @@ -397,7 +395,7 @@ impl SupervisorFsm { /// A supervisor seeded directly into `Running` with a known up-set. /// /// The teardown driver (`stop()`) reconstructs the up-set from observed - /// runtime presence (`dns_task.is_some()`, transports keys, etc.) rather + /// runtime presence (`ipv6tun.dns_up()`, transports keys, etc.) rather /// than relying on a live machine persisted across start/stop, so that /// teardown ordering is authored here regardless of how the node reached /// `Running`. Feeding `Event::Stop` then yields the ordered `StopChild` @@ -813,33 +811,11 @@ pub(crate) struct Supervisor { /// Packet sender for transports. pub(in crate::node) packet_tx: Option, - /// TUN packet sender channel. - pub(in crate::node) tun_tx: Option, - /// Receiver for outbound packets from the TUN reader. - pub(in crate::node) tun_outbound_rx: Option, - /// TUN reader thread handle. - pub(in crate::node) tun_reader_handle: Option>, - /// TUN writer thread handle. - pub(in crate::node) tun_writer_handle: Option>, - /// Shutdown pipe: writing to this fd unblocks the TUN reader thread on - /// macOS and FreeBSD. On Linux, deleting the interface via netlink - /// serves the same purpose. - #[cfg(any(target_os = "macos", target_os = "freebsd"))] - pub(in crate::node) tun_shutdown_fd: Option, - - /// Receiver for resolved identities from the DNS responder. - pub(in crate::node) dns_identity_rx: Option, - /// DNS responder task handle. - pub(in crate::node) dns_task: Option>, - /// Address the DNS responder actually bound, read back from the socket - /// after `bind` so a port-0 config resolves to the assigned port. `Some` - /// only while the responder is up; published to embedders through - /// [`Node::dns_local_addr`](crate::Node::dns_local_addr). - pub(in crate::node) dns_local_addr: Option, - /// Sends a new peer-alias base to the running DNS responder; `Some` only - /// while it runs. - pub(in crate::node) dns_aliases: - Option>, + /// TUN and DNS child handles: the TUN device name, channels, reader and + /// writer threads and (macOS/FreeBSD) shutdown pipe, and the DNS + /// responder task, identity receiver, bound address and peer-alias + /// sender. + pub(in crate::node) ipv6tun: crate::ipv6tun::lifecycle::Handles, /// Sender for each UDP listen socket the transport spawn binds — its raw /// fd and the instance name it was configured under — armed by @@ -914,16 +890,7 @@ impl Supervisor { Self { state: NodeState::Created, packet_tx: None, - tun_tx: None, - tun_outbound_rx: None, - tun_reader_handle: None, - tun_writer_handle: None, - #[cfg(any(target_os = "macos", target_os = "freebsd"))] - tun_shutdown_fd: None, - dns_identity_rx: None, - dns_task: None, - dns_local_addr: None, - dns_aliases: None, + ipv6tun: Default::default(), #[cfg(unix)] udp_fd_tx: None, nostr_rendezvous: crate::nostr::RendezvousDriver::default(), diff --git a/src/node/mod.rs b/src/node/mod.rs index 5f63cdfc..11efe5fc 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -15,7 +15,7 @@ pub(crate) mod decrypt_worker; #[cfg(unix)] pub(crate) mod encrypt_worker; mod handlers; -mod lifecycle; +pub(crate) mod lifecycle; pub(crate) mod metrics; pub(crate) mod netmon; pub use netmon::NetmonTrigger; @@ -634,8 +634,6 @@ pub struct Node { // === TUN Interface === /// TUN device state. tun_state: TunState, - /// TUN interface name (for cleanup). - tun_name: Option, /// Slot the embedder installs its BLE radio into, armed by /// [`Self::enable_app_owned_ble_radio`]. `None` unless armed. @@ -950,7 +948,6 @@ impl Node { crate::control::snapshot::NativeSnapshot::empty(), )), tun_state, - tun_name: None, #[cfg(all(ble_available, any(target_os = "android", test)))] ble_radio: None, index_allocator: IndexAllocator::new(), @@ -1128,7 +1125,6 @@ impl Node { crate::control::snapshot::NativeSnapshot::empty(), )), tun_state, - tun_name: None, #[cfg(all(ble_available, any(target_os = "android", test)))] ble_radio: None, index_allocator: IndexAllocator::new(), @@ -1546,9 +1542,7 @@ impl Node { tracing::debug!(entries = base.len(), "Rebuilding peer alias maps"); self.host_map.set_base(base.clone()); self.peer_acl.rebase(base.clone()).await; - if let Some(tx) = &self.supervisor.dns_aliases { - tx.send_replace(base); - } + self.supervisor.ipv6tun.publish_aliases(base); } /// Return a human-readable display name for a NodeAddr. @@ -2018,7 +2012,7 @@ impl Node { estimated_mesh_size: self.estimated_mesh_size, state: self.supervisor.state, tun_state: self.tun_state, - tun_name: self.tun_name.clone(), + tun_name: self.supervisor.ipv6tun.tun_name.clone(), effective_ipv6_mtu: self.effective_ipv6_mtu(), connection_count: self.connection_count(), peer_count: self.peers.len(), @@ -2723,7 +2717,7 @@ impl Node { /// Get the TUN interface name, if active. pub fn tun_name(&self) -> Option<&str> { - self.tun_name.as_deref() + self.supervisor.ipv6tun.tun_name.as_deref() } // === Resource Limits === @@ -3693,7 +3687,14 @@ impl Node { /// /// Returns None if TUN is not active or the node hasn't been started. pub fn tun_tx(&self) -> Option<&TunTx> { - self.supervisor.tun_tx.as_ref() + self.supervisor.ipv6tun.tun_tx.as_ref() + } + + /// Install a TUN packet sender, standing in for the TUN writer, so a + /// test can read what the node delivers toward the host. + #[cfg(test)] + pub(crate) fn install_tun(&mut self, tun_tx: TunTx) { + self.supervisor.ipv6tun.tun_tx = Some(tun_tx); } /// Set up an **app-owned TUN**: rather than FIPS creating a system TUN @@ -3720,8 +3721,8 @@ impl Node { let (outbound_tx, outbound_rx) = tokio::sync::mpsc::channel(tun_channel_size); // mesh → app: the node writes inbound packets to `tun_tx`; the app pulls. let (tun_tx, tun_rx) = std::sync::mpsc::channel(); - self.supervisor.tun_tx = Some(tun_tx); - self.supervisor.tun_outbound_rx = Some(outbound_rx); + self.supervisor.ipv6tun.tun_tx = Some(tun_tx); + self.supervisor.ipv6tun.tun_outbound_rx = Some(outbound_rx); self.tun_state = TunState::Active; (outbound_tx, tun_rx) } @@ -3894,7 +3895,7 @@ impl Node { /// Reading live node state from a backgrounded loop is a general gap, not /// one this accessor tries to close. pub fn dns_local_addr(&self) -> Option { - self.supervisor.dns_local_addr + self.supervisor.ipv6tun.dns_local_addr } /// A handle that wakes the medium-change detector (`node.netmon.*`) now diff --git a/src/node/tests/discovery.rs b/src/node/tests/discovery.rs index f5a39ae1..6b4bef13 100644 --- a/src/node/tests/discovery.rs +++ b/src/node/tests/discovery.rs @@ -2029,9 +2029,9 @@ async fn test_check_pending_lookups_default_sequence_unreachable() { "test pins the [1,2,4,8] default; update the test if the default changes" ); - // Inject a TUN sender so `send_icmpv6_dest_unreachable` is observable. + // Inject a TUN sender so the no-route Destination Unreachable is observable. let (tun_tx, tun_rx) = mpsc::channel::>(); - node.supervisor.tun_tx = Some(tun_tx); + node.install_tun(tun_tx); // Build a target identity (the unreachable destination). let target_identity = Identity::generate(); diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index 1c636c8d..39123a80 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -616,7 +616,7 @@ async fn mixed_profile_nodes_converge_and_forward() { let mut tun_rx = Vec::with_capacity(nodes.len()); for tn in nodes.iter_mut() { let (tx, rx) = std::sync::mpsc::channel(); - tn.node.supervisor.tun_tx = Some(tx); + tn.node.install_tun(tx); tun_rx.push(rx); } @@ -726,7 +726,7 @@ async fn leaf_smallest_addr_does_not_partition_multihop() { // TUN receiver on B so delivered plaintext can be observed. let (tx, b_rx) = std::sync::mpsc::channel(); - nodes[1].node.supervisor.tun_tx = Some(tx); + nodes[1].node.install_tun(tx); let d_addr = *nodes[2].node.node_addr(); let (b_addr, b_pubkey) = ( @@ -849,7 +849,7 @@ async fn test_session_100_nodes() { let mut tun_receivers: Vec>> = Vec::with_capacity(NUM_NODES); for tn in nodes.iter_mut() { let (tx, rx) = mpsc::channel(); - tn.node.supervisor.tun_tx = Some(tx); + tn.node.install_tun(tx); tun_receivers.push(rx); } @@ -1263,7 +1263,7 @@ async fn test_tun_outbound_established_session() { // Install TUN receiver on Node 1 let (tun_tx, tun_rx) = std::sync::mpsc::channel(); - nodes[1].node.supervisor.tun_tx = Some(tun_tx); + nodes[1].node.install_tun(tun_tx); // Build and inject an IPv6 packet let test_payload = b"data-plane-test-12345"; @@ -1352,7 +1352,7 @@ async fn rekey_cutover_preserves_data_plane() { // node 1's TUN receiver observes decoded plaintext. let (tun_tx, tun_rx) = std::sync::mpsc::channel(); - nodes[1].node.supervisor.tun_tx = Some(tun_tx); + nodes[1].node.install_tun(tun_tx); let src_fips = crate::FipsAddress::from_node_addr(&node0_addr); let dst_fips = crate::FipsAddress::from_node_addr(&node1_addr); @@ -1501,9 +1501,9 @@ async fn aged_link_pair( // Each node's TUN receiver observes the plaintext the other one sent. let (tun0_tx, tun0_rx) = std::sync::mpsc::channel(); - nodes[0].node.supervisor.tun_tx = Some(tun0_tx); + nodes[0].node.install_tun(tun0_tx); let (tun1_tx, tun1_rx) = std::sync::mpsc::channel(); - nodes[1].node.supervisor.tun_tx = Some(tun1_tx); + nodes[1].node.install_tun(tun1_tx); let fips0 = crate::FipsAddress::from_node_addr(&node0_addr); let fips1 = crate::FipsAddress::from_node_addr(&node1_addr); @@ -2041,7 +2041,7 @@ async fn test_tun_outbound_triggers_session_initiation() { // Install TUN receiver on Node 1 let (tun_tx, tun_rx) = std::sync::mpsc::channel(); - nodes[1].node.supervisor.tun_tx = Some(tun_tx); + nodes[1].node.install_tun(tun_tx); // Build and inject an IPv6 packet (identity cache populated at peer promotion) let test_payload = b"trigger-session-test"; @@ -2094,7 +2094,7 @@ async fn test_tun_outbound_unknown_destination() { // Install TUN receiver on Node 0 (for ICMPv6 response) let (tun_tx, tun_rx) = std::sync::mpsc::channel(); - nodes[0].node.supervisor.tun_tx = Some(tun_tx); + nodes[0].node.install_tun(tun_tx); let src_fips = crate::FipsAddress::from_node_addr(nodes[0].node.node_addr()); @@ -2146,7 +2146,7 @@ async fn test_tun_outbound_3node_forwarded() { // Install TUN receiver on Node 2 let (tun_tx, tun_rx) = std::sync::mpsc::channel(); - nodes[2].node.supervisor.tun_tx = Some(tun_tx); + nodes[2].node.install_tun(tun_tx); // Build and inject an IPv6 packet (triggers session initiation to Node 2) let test_payload = b"forwarded-data-plane"; @@ -2191,7 +2191,7 @@ async fn test_tun_outbound_pending_queue_flush() { // Install TUN receiver on Node 1 let (tun_tx, tun_rx) = std::sync::mpsc::channel(); - nodes[1].node.supervisor.tun_tx = Some(tun_tx); + nodes[1].node.install_tun(tun_tx); // Send 5 packets before any session exists let mut packets = Vec::new(); @@ -2842,7 +2842,7 @@ async fn test_tun_outbound_path_mtu_generates_ptb() { // Install TUN receiver on source node to capture ICMPv6 PTB let (tun_tx, tun_rx) = std::sync::mpsc::channel(); - nodes[0].node.supervisor.tun_tx = Some(tun_tx); + nodes[0].node.install_tun(tun_tx); // Build an IPv6 packet that fits local MTU but exceeds path MTU let reduced_ipv6_mtu = crate::upper::icmp::effective_ipv6_mtu(reduced_mtu) as usize; @@ -2899,7 +2899,7 @@ async fn test_tun_outbound_path_mtu_generates_ptb() { // Verify a packet that fits within path MTU passes through (no PTB) let (tun_tx2, tun_rx2) = std::sync::mpsc::channel(); - nodes[0].node.supervisor.tun_tx = Some(tun_tx2); + nodes[0].node.install_tun(tun_tx2); let fitting_payload = vec![0u8; reduced_ipv6_mtu - 41]; // fits within path MTU let fitting_packet = build_ipv6_packet(&src_fips, &dst_fips, &fitting_payload); assert!(fitting_packet.len() <= reduced_ipv6_mtu); @@ -3038,7 +3038,7 @@ async fn test_multihop_pmtud_heterogeneous_mtu() { // should check PathMtuState and generate ICMPv6 PTB on TUN instead // of forwarding. let (tun_tx2, tun_rx2) = std::sync::mpsc::channel(); - nodes[0].node.supervisor.tun_tx = Some(tun_tx2); + nodes[0].node.install_tun(tun_tx2); nodes[0].node.handle_tun_outbound(ipv6_packet.clone()).await; @@ -3082,7 +3082,7 @@ async fn test_multihop_pmtud_heterogeneous_mtu() { // Verify a fitting packet still passes through without PTB let (tun_tx3, tun_rx3) = std::sync::mpsc::channel(); - nodes[0].node.supervisor.tun_tx = Some(tun_tx3); + nodes[0].node.install_tun(tun_tx3); let fitting_payload = vec![0xCDu8; 600 - 40]; // 600-byte IPv6 packet, well within 694 let fitting_packet = build_ipv6_packet(&src_fips, &dst_fips, &fitting_payload); @@ -3410,7 +3410,7 @@ async fn test_forged_mtu_exceeded_of_zero_does_not_blackhole_the_session() { nodes[0].node.handle_mtu_exceeded(&reporter, &inner).await; let (tun_tx, tun_rx) = std::sync::mpsc::channel(); - nodes[0].node.supervisor.tun_tx = Some(tun_tx); + nodes[0].node.install_tun(tun_tx); let payload = vec![0u8; 560]; let ipv6_packet = build_ipv6_packet(&src_fips, &dst_fips, &payload); @@ -5972,9 +5972,9 @@ async fn a_lost_initial_msg3_is_resent_and_the_responder_completes_the_session() let node1_addr = *nodes[1].node.node_addr(); let (tun0_tx, tun0_rx) = std::sync::mpsc::channel(); - nodes[0].node.supervisor.tun_tx = Some(tun0_tx); + nodes[0].node.install_tun(tun0_tx); let (tun1_tx, tun1_rx) = std::sync::mpsc::channel(); - nodes[1].node.supervisor.tun_tx = Some(tun1_tx); + nodes[1].node.install_tun(tun1_tx); let fips0 = crate::FipsAddress::from_node_addr(&node0_addr); let fips1 = crate::FipsAddress::from_node_addr(&node1_addr); @@ -6027,7 +6027,7 @@ async fn an_initiator_stops_resending_msg3_once_a_responder_frame_authenticates( let node0_addr = *nodes[0].node.node_addr(); let node1_addr = *nodes[1].node.node_addr(); let (tun0_tx, tun0_rx) = std::sync::mpsc::channel(); - nodes[0].node.supervisor.tun_tx = Some(tun0_tx); + nodes[0].node.install_tun(tun0_tx); let fips0 = crate::FipsAddress::from_node_addr(&node0_addr); let fips1 = crate::FipsAddress::from_node_addr(&node1_addr); let interval_ms = nodes[0] @@ -6082,9 +6082,9 @@ async fn a_resent_msg3_reaching_an_established_responder_is_refused_and_the_sess let node0_addr = *nodes[0].node.node_addr(); let node1_addr = *nodes[1].node.node_addr(); let (tun0_tx, tun0_rx) = std::sync::mpsc::channel(); - nodes[0].node.supervisor.tun_tx = Some(tun0_tx); + nodes[0].node.install_tun(tun0_tx); let (tun1_tx, tun1_rx) = std::sync::mpsc::channel(); - nodes[1].node.supervisor.tun_tx = Some(tun1_tx); + nodes[1].node.install_tun(tun1_tx); let fips0 = crate::FipsAddress::from_node_addr(&node0_addr); let fips1 = crate::FipsAddress::from_node_addr(&node1_addr); let interval_ms = nodes[0] diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index f07d8efb..6cff5903 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -3675,6 +3675,7 @@ async fn dns_responder_serves_a_proxying_embedder() { let identity = tokio::time::timeout( std::time::Duration::from_secs(2), node.supervisor + .ipv6tun .dns_identity_rx .as_mut() .expect("responder installed the identity receiver") @@ -4934,10 +4935,10 @@ fn mesh_filter_resolves_the_live_tun_device_rather_than_the_configured_name() { config.tun.name = Some("fips-absent-dev".to_string()); let mut node = Node::new(config).unwrap(); - assert_eq!(node.mesh_ifindex(), None); + assert_eq!(node.supervisor.ipv6tun.mesh_ifindex(), None); - node.tun_name = Some(loopback.to_string()); - assert_eq!(node.mesh_ifindex(), Some(expected)); + node.supervisor.ipv6tun.tun_name = Some(loopback.to_string()); + assert_eq!(node.supervisor.ipv6tun.mesh_ifindex(), Some(expected)); } /// The msg1 handler keeps its pending slot for as long as it is running. diff --git a/testing/acl-allowlist/README.md b/testing/acl-allowlist/README.md index 95d0f2ee..866861d8 100644 --- a/testing/acl-allowlist/README.md +++ b/testing/acl-allowlist/README.md @@ -157,7 +157,7 @@ Rejected peer by ACL ... context=inbound_handshake decision=denylist match ``` Those messages are now emitted at debug level. This harness enables -`RUST_LOG=info,fips::node=debug` so the ACL rejection details stay visible in +`fips::node=debug` in `RUST_LOG` so the ACL rejection details stay visible in test logs, and operators can temporarily raise log level the same way when diagnosing ACL issues locally. diff --git a/testing/acl-allowlist/docker-compose.yml b/testing/acl-allowlist/docker-compose.yml index a20428a0..63f58050 100644 --- a/testing/acl-allowlist/docker-compose.yml +++ b/testing/acl-allowlist/docker-compose.yml @@ -35,7 +35,7 @@ x-fips-common: &fips-common restart: "no" environment: - FIPS_TEST_MODE=default - - RUST_LOG=info,fips::node=debug + - RUST_LOG=info,fips::node=debug,fips::ipv6tun::icmp=debug,fips::ipv6tun::lifecycle=debug volumes: - ../docker/resolv.conf:/etc/resolv.conf:ro diff --git a/testing/firewall/docker-compose.yml b/testing/firewall/docker-compose.yml index fe657732..07deb44f 100644 --- a/testing/firewall/docker-compose.yml +++ b/testing/firewall/docker-compose.yml @@ -36,7 +36,7 @@ x-fips-common: &fips-common restart: "no" environment: - FIPS_TEST_MODE=default - - RUST_LOG=info,fips::node=debug + - RUST_LOG=info,fips::node=debug,fips::ipv6tun::icmp=debug,fips::ipv6tun::lifecycle=debug services: service-a: diff --git a/testing/iface-binding/docker-compose.yml b/testing/iface-binding/docker-compose.yml index fd87ba6b..440585c8 100644 --- a/testing/iface-binding/docker-compose.yml +++ b/testing/iface-binding/docker-compose.yml @@ -33,7 +33,7 @@ x-fips-common: &fips-common # the daemon, which is precisely the workaround this mechanism retires. The # daemon must do its own waiting here or the suite proves nothing. - FIPS_TEST_MODE=default - - RUST_LOG=info,fips::transport::ethernet=debug,fips::node=debug + - RUST_LOG=info,fips::transport::ethernet=debug,fips::node=debug,fips::ipv6tun::icmp=debug,fips::ipv6tun::lifecycle=debug networks: - ifb-net diff --git a/testing/mesh-lab/README.md b/testing/mesh-lab/README.md index 4887907a..15d3b020 100644 --- a/testing/mesh-lab/README.md +++ b/testing/mesh-lab/README.md @@ -102,7 +102,8 @@ rep does, set them in the invoking shell: `handshake`, `forwarding`, `session`, `encrypted`, `mmp` (via `compose-trace.yml`). - nat-lan — `fips::nostr`, `transport::udp`, - `node::lifecycle`, `handlers::handshake`, `dataplane::forwarding` + `node::lifecycle`, `ipv6tun::lifecycle`, `handlers::handshake`, + `dataplane::forwarding` (via `compose-trace-nat.yml`, picked up by `testing/nat/scripts/nat-test.sh` through the `FIPS_NAT_EXTRA_COMPOSE` env-var hook). diff --git a/testing/mesh-lab/compose-trace-nat.yml b/testing/mesh-lab/compose-trace-nat.yml index 70c6c700..13b5ac6b 100644 --- a/testing/mesh-lab/compose-trace-nat.yml +++ b/testing/mesh-lab/compose-trace-nat.yml @@ -8,6 +8,9 @@ # (where the punch packets flow) # - fips::node::lifecycle — daemon bootstrap, peer state # machine, adoption transitions +# - fips::ipv6tun::lifecycle — TUN and DNS child start/stop +# (formerly logged under +# node::lifecycle) # - fips::node::handlers::handshake — Noise handshake msg1/2/3, # cross-init tie-breaker # ("Ignoring established NAT @@ -33,7 +36,7 @@ # is set and the suite is nat-lan. x-trace-rust-log: &trace-rust-log - RUST_LOG: "info,fips::nostr=trace,fips::transport::udp=trace,fips::node::lifecycle=trace,fips::node::handlers::handshake=trace,fips::node::dataplane::forwarding=trace" + RUST_LOG: "info,fips::nostr=trace,fips::transport::udp=trace,fips::node::lifecycle=trace,fips::ipv6tun::lifecycle=trace,fips::node::handlers::handshake=trace,fips::node::dataplane::forwarding=trace" services: lan-a: diff --git a/testing/mesh-lab/compose-trace.yml b/testing/mesh-lab/compose-trace.yml index feb0f7ae..dbf7a69b 100644 --- a/testing/mesh-lab/compose-trace.yml +++ b/testing/mesh-lab/compose-trace.yml @@ -9,6 +9,9 @@ # - fips::node::handlers::session — FSP K-bit cutover, drain # - fips::node::dataplane::encrypted — FMP K-bit flip detection # - fips::node::handlers::mmp — link liveness, SRTT, ETX +# - fips::ipv6tun::icmp — ICMPv6 Packet Too Big toward +# the host (formerly logged +# under handlers::session) # # Other modules stay at info to keep log volume manageable. The base # docker-compose.yml's per-service RUST_LOG values (currently @@ -33,7 +36,7 @@ # environment variable FIPS_MESH_LAB_TRACE=1 is set. x-trace-rust-log: &trace-rust-log - RUST_LOG: "info,fips::node::handlers::rekey=trace,fips::node::handlers::handshake=trace,fips::node::dataplane::forwarding=trace,fips::node::handlers::session=trace,fips::node::dataplane::encrypted=trace,fips::node::handlers::mmp=trace" + RUST_LOG: "info,fips::node::handlers::rekey=trace,fips::node::handlers::handshake=trace,fips::node::dataplane::forwarding=trace,fips::node::handlers::session=trace,fips::node::dataplane::encrypted=trace,fips::node::handlers::mmp=trace,fips::ipv6tun::icmp=trace" services: # rekey profile diff --git a/testing/mesh-lab/run-loop.sh b/testing/mesh-lab/run-loop.sh index a4e0c987..cc6f5ca3 100755 --- a/testing/mesh-lab/run-loop.sh +++ b/testing/mesh-lab/run-loop.sh @@ -433,9 +433,10 @@ run_nat_lan() { # nat-test.sh at the nat-specific trace overlay via the # FIPS_NAT_EXTRA_COMPOSE env-var hook in nat-test.sh. The overlay # bumps RUST_LOG to trace on discovery::nostr, transport::udp, - # node::lifecycle, handlers::handshake, dataplane::forwarding — - # the modules covering the cross-init / adoption / handshake - # path that the NAT-traversal flake exhibits. Path is repo-relative. + # node::lifecycle, ipv6tun::lifecycle, handlers::handshake, + # dataplane::forwarding — the modules covering the cross-init / + # adoption / handshake path that the NAT-traversal flake exhibits. + # Path is repo-relative. local -a env_args=(FIPS_NAT_SKIP_FINAL_CLEANUP=1) if [ -n "${FIPS_MESH_LAB_TRACE:-}" ]; then env_args+=(FIPS_NAT_EXTRA_COMPOSE=testing/mesh-lab/compose-trace-nat.yml) diff --git a/testing/nat/docker-compose.yml b/testing/nat/docker-compose.yml index 93c2fb39..3189c471 100644 --- a/testing/nat/docker-compose.yml +++ b/testing/nat/docker-compose.yml @@ -27,7 +27,7 @@ x-fips-common: &fips-common - net.ipv6.conf.all.disable_ipv6=0 restart: "no" environment: - - RUST_LOG=info,fips::nostr=debug,fips::node::lifecycle=debug + - RUST_LOG=info,fips::nostr=debug,fips::node::lifecycle=debug,fips::ipv6tun::lifecycle=debug services: relay: @@ -134,7 +134,7 @@ services: entrypoint: - /usr/local/bin/nat-node-entrypoint.sh environment: - - RUST_LOG=info,fips::nostr=debug,fips::node::lifecycle=debug + - RUST_LOG=info,fips::nostr=debug,fips::node::lifecycle=debug,fips::ipv6tun::lifecycle=debug - DATA_IF=eth0 - ROUTE_SUBNET=${NAT_WAN_PREFIX:-172.31.254}.0/24 - ROUTE_VIA=172.31.1.254 @@ -160,7 +160,7 @@ services: entrypoint: - /usr/local/bin/nat-node-entrypoint.sh environment: - - RUST_LOG=info,fips::nostr=debug,fips::node::lifecycle=debug + - RUST_LOG=info,fips::nostr=debug,fips::node::lifecycle=debug,fips::ipv6tun::lifecycle=debug - DATA_IF=eth0 - ROUTE_SUBNET=${NAT_WAN_PREFIX:-172.31.254}.0/24 - ROUTE_VIA=172.31.2.254 @@ -186,7 +186,7 @@ services: entrypoint: - /usr/local/bin/nat-node-entrypoint.sh environment: - - RUST_LOG=info,fips::nostr=debug,fips::node::lifecycle=debug + - RUST_LOG=info,fips::nostr=debug,fips::node::lifecycle=debug,fips::ipv6tun::lifecycle=debug - DATA_IF=eth0 - ROUTE_SUBNET=${NAT_WAN_PREFIX:-172.31.254}.0/24 - ROUTE_VIA=172.31.1.254 @@ -212,7 +212,7 @@ services: entrypoint: - /usr/local/bin/nat-node-entrypoint.sh environment: - - RUST_LOG=info,fips::nostr=debug,fips::node::lifecycle=debug + - RUST_LOG=info,fips::nostr=debug,fips::node::lifecycle=debug,fips::ipv6tun::lifecycle=debug - DATA_IF=eth0 - ROUTE_SUBNET=${NAT_WAN_PREFIX:-172.31.254}.0/24 - ROUTE_VIA=172.31.2.254