Merge branch 'master' into next

This commit is contained in:
Johnathan Corgan
2026-10-02 06:18:32 +00:00
17 changed files with 672 additions and 100 deletions
+12 -4
View File
@@ -437,10 +437,18 @@ jobs:
# Debug-only helpers (anything behind #[cfg(debug_assertions)]) vanish in
# a release build, so a test calling one without the same gate breaks a
# build no other job performs: every run above compiles the test target
# in debug. Compile it in release too, without running it — the point is
# that it builds at all. Mirrored in testing/ci-local.sh.
- name: Compile the library tests in release mode
run: cargo test --release --lib --no-run
# in debug. Compile it in release too. Of what it builds, run only the
# leg-slot residue test, which is release-only because only an optimised
# build leaves the slot in a state worth measuring; the grep fails the
# step if that test did not run, since a name filter that matches
# nothing passes. Mirrored in testing/ci-local.sh.
- name: Compile the library tests in release mode and run the leg-slot residue test
shell: bash
env:
RESIDUE_TEST: peer::machine::tests::take_leg_leaves_no_session_keys_in_the_slot_it_empties
run: |
cargo test --release --lib -- --exact "$RESIDUE_TEST" | tee "$RUNNER_TEMP/release-lib-tests.log"
grep -q '^test result: ok\. 1 passed;' "$RUNNER_TEMP/release-lib-tests.log"
# ─────────────────────────────────────────────────────────────────────────────
# Job 2b – Unit tests (macOS)
+17
View File
@@ -698,6 +698,9 @@ with v0.5.x or earlier peers.
the interface goes away between the presence check and the bind.
- TCP connections try the remaining addresses for a hostname after a
connection fails, within the existing overall connection timeout.
- An Ethernet peer's address is always shown as a colon-separated MAC address
in `fipsctl show links` and in log lines. An address whose six bytes happened
to be valid UTF-8 was printed as text instead.
#### Gateway
@@ -851,6 +854,10 @@ with v0.5.x or earlier peers.
- If an encrypt worker thread exits, the daemon no longer stops once that
worker's send queue fills. Packets for that worker are now dropped instead
of blocking forever. This applies to the default sender.
- With the opt-in ordered sender (`FIPS_MACOS_ORDERED_SENDER`), the daemon no
longer stops forwarding when an encrypt worker thread or a destination's send
thread exits. Packets that thread would have handled are dropped instead of
holding up every later packet to the same destination.
#### Native datagram API
@@ -909,6 +916,11 @@ with v0.5.x or earlier peers.
foreground run reads, which it omitted, and says to stop the service before
rerunning `install-service.ps1` to upgrade: with the service running, the
installer fails copying `fips.exe`.
- An ICMP port-unreachable answering a datagram the node sent no longer shows
up as a receive error on the UDP transport's socket. Windows reports one on
the socket's next receive by default, and that socket is shared by every UDP
peer, so one unreachable peer address put errors into the receive path that
every other UDP peer's traffic uses.
### Security
@@ -975,6 +987,11 @@ with v0.5.x or earlier peers.
handshake. The security reference now states what clearing key material in
memory does and does not cover in a release build, including copies left
behind by moves and the identity loaded from a secret string.
- A connection's handshake state is now cleared from its control machine when
the connection is promoted to an active peer, resolved as one side of a
cross-connection, or reaped as stale. After a promotion, the memory it left
held the session's two traffic keys for as long as the peer stayed
connected. The security reference is updated to match.
## [0.5.2] - 2026-09-28
Generated
+1
View File
@@ -1208,6 +1208,7 @@ dependencies = [
"tracing-subscriber",
"tun",
"windows-service",
"windows-sys 0.61.2",
"wintun",
"zeroize",
]
+3
View File
@@ -78,6 +78,9 @@ bluer = { version = "0.17", features = ["bluetoothd", "l2cap"] }
[target.'cfg(windows)'.dependencies]
wintun = "0.5"
windows-service = "0.8.1"
# For WSAIoctl(SIO_UDP_CONNRESET) on the UDP socket. The version socket2 and
# tokio already pull in, so the lock file gains no package.
windows-sys = { version = "0.61", features = ["Win32_Networking_WinSock", "Win32_System_IO"] }
[package.metadata.deb]
maintainer = "Johnathan Corgan <johnathan@corganlabs.com>"
+7 -7
View File
@@ -123,14 +123,14 @@ reach:
dropped. Each of those moves leaves behind, in a stack frame that is
no longer in use, a copy of the node's long-term private key and,
once the handshake has started, of its ephemeral key and chaining
key. Two moves out of the slots on a connection's control machine
key. Three moves out of the slots on a connection's control machine
are cleared: when a completed handshake leaves the slot that held it,
and when a session is taken out of its slot for a rekey, the slot is
overwritten as the value leaves. Other moves are not. When a
connection is promoted to an active peer, or reaped as stale, its
whole handshake state is moved off its control machine, which leaves
the session's two traffic keys, or an unfinished handshake's private
keys, in the heap memory the machine occupies.
when a session is taken out of its slot for a rekey, and when the
whole handshake state leaves the machine because the connection is
promoted to an active peer, resolved as one side of a
cross-connection, rejected by a leaf node that already has its peer,
or reaped as stale. In each case the slot is overwritten as the value
leaves. Other moves are not.
- **Loading the identity from a secret string.** Building the node's
identity from its key file or from `node.identity.nsec` leaves
copies of the private key, among them a whole intermediate identity,
+17 -3
View File
@@ -1092,7 +1092,12 @@ mod tests {
// already queued and no read here waits on anything.
(&daemon).write_all(REFUSAL).unwrap();
fdpass::send_once(daemon.as_raw_fd(), CONNECT_REPLY, Some(passed.as_fd())).unwrap();
drop(passed);
// `passed` stays alive until the client holds the descriptor: until
// then the message would be the socket's only reference, and Darwin's
// collector flushes such a socket. The provoked collection is there so
// a regression shows on macOS without waiting for a stray collection;
// whether it shows every run rests on `dgram_probe`'s fixed wait.
crate::native::dgram_probe::provoke_collection();
// **How many reads this takes is the platform's business, and asserting
// it was wrong.** Linux coalesces the plain write with the sendmsg that
@@ -1133,6 +1138,7 @@ mod tests {
let (second, second_fd) = wire.line().unwrap();
assert!(second.starts_with(br#"{"status":"ok""#));
let received = second_fd.expect("the reply must carry the descriptor");
drop(passed);
// A live socket rather than merely a number: the far half sees it.
let mut received = UnixStream::from(received);
@@ -1207,7 +1213,6 @@ mod tests {
let worker = thread::spawn(move || {
let command = read_line(&daemon);
fdpass::send_once(daemon.as_raw_fd(), LISTEN_REPLY, Some(passed.as_fd())).unwrap();
drop(passed);
// As for connect: setup is over, so the connection is over. A
// listener that held one would take its flows down with it.
// A deadline, or the defect this asserts against fails as a hang
@@ -1220,6 +1225,9 @@ mod tests {
let read = (&daemon)
.read(&mut rest)
.expect("the RPC connection outlived the setup call that opened it");
// Only now: the client closes once it holds the descriptor, and
// until then the message would be the socket's only reference.
drop(passed);
(command, read)
});
@@ -1308,9 +1316,15 @@ mod tests {
// client could read the arrival is already there when it holds it.
(&held).write_all(b"opening").unwrap();
fdpass::send_once(theirs.as_raw_fd(), ARRIVAL, Some(passed.as_fd())).unwrap();
drop(passed);
// `passed` stays alive until the client holds the descriptor: until
// then the message would be the socket's only reference, and Darwin's
// collector flushes such a socket. The provoked collection is there so
// a regression shows on macOS without waiting for a stray collection;
// whether it shows every run rests on `dgram_probe`'s fixed wait.
crate::native::dgram_probe::provoke_collection();
let (flow, peer) = listener.accept().unwrap();
drop(passed);
assert_eq!(peer.to_string(), format!("{PEER}:5001"));
assert_eq!(flow.peer_addr(), peer);
// From the arrival's own `node`, not from the listener: an accepted
+6 -4
View File
@@ -1595,10 +1595,12 @@ mod tests {
drop(client);
connection.settle_closed(flow).await;
// `settle_closed` observes the reader's flag, which it sets before it
// sends the release, so the registry may not have processed it yet.
// Retry rather than sleep: without the reclaim every attempt fails and
// the loop runs out, which is the failure this test exists to produce.
// `settle_closed` returns once the reader has flagged the flow closed or
// forgotten it. Neither means the registry has served the release: the
// flag is set before the release is queued, and the flow is forgotten
// once it is queued, not once it is processed. Retry rather than sleep:
// without the reclaim every attempt fails and the loop runs out, which
// is the failure this test exists to produce.
let mut last = serde_json::Value::Null;
for _ in 0..1000 {
// Not `ask`: a connect that succeeds carries a descriptor, and that
+360 -46
View File
@@ -156,9 +156,7 @@ pub(crate) struct FspSealJob {
struct QueuedFmpSendJob {
job: FmpSendJob,
#[cfg(target_os = "macos")]
macos_flow: Option<Arc<MacSequencedSendFlow>>,
#[cfg(target_os = "macos")]
macos_seq: u64,
macos_ticket: Option<MacSeqTicket>,
}
impl QueuedFmpSendJob {
@@ -167,19 +165,15 @@ impl QueuedFmpSendJob {
Self {
job,
#[cfg(target_os = "macos")]
macos_flow: None,
#[cfg(target_os = "macos")]
macos_seq: 0,
macos_ticket: None,
}
}
#[cfg(target_os = "macos")]
fn macos_sequenced(job: FmpSendJob, macos_flow: Arc<MacSequencedSendFlow>) -> Self {
let macos_seq = macos_flow.reserve_seq();
Self {
job,
macos_flow: Some(macos_flow),
macos_seq,
macos_ticket: Some(MacSeqTicket::reserve(macos_flow)),
}
}
}
@@ -283,6 +277,8 @@ impl<T> MacWorkerSender<T> {
.lock()
.expect("encrypt worker queue poisoned");
if state.closed {
// Outside the lock: dropping a sequenced job completes its slot.
drop(state);
drop(job);
return Err(MacWorkerTryPushError::Closed);
}
@@ -307,6 +303,7 @@ impl<T> MacWorkerSender<T> {
.expect("encrypt worker queue poisoned");
loop {
if state.closed {
drop(state);
drop(job);
return Err(MacWorkerPushError);
}
@@ -610,7 +607,11 @@ impl MacSequencedSendFlows {
let mut flows = self.flows.lock().expect("mac send flow map poisoned");
self.prune_idle_locked(&mut flows, now_ms);
if let Some(flow) = flows.get(&key) {
// A closed flow's sender thread has exited and sends nothing more, so
// it is replaced rather than handed further jobs.
if let Some(flow) = flows.get(&key)
&& !flow.is_closed()
{
flow.mark_used(now_ms);
return Arc::clone(flow);
}
@@ -652,7 +653,7 @@ impl MacSequencedSendFlows {
let idle_ms = mac_send_flow_idle_ms();
flows.retain(|_, flow| {
if flow.is_idle(now_ms, idle_ms) {
if flow.is_closed() || flow.is_idle(now_ms, idle_ms) {
flow.close();
false
} else {
@@ -764,6 +765,77 @@ struct MacCompletionGroup {
items: Vec<(u64, MacSendItem)>,
}
#[cfg(target_os = "macos")]
impl MacCompletionGroup {
/// Hand every item to the flow's sender.
fn deliver(mut self) {
let items = std::mem::take(&mut self.items);
self.flow.complete_many(items);
}
}
/// A group dropped before delivery, when its worker unwinds, still completes
/// its slots, each as a skip, so the flow moves past them.
#[cfg(target_os = "macos")]
impl Drop for MacCompletionGroup {
fn drop(&mut self) {
for (seq, _) in self.items.drain(..) {
self.flow.complete_skip(seq);
}
}
}
/// One reserved slot in a flow's send order, owed a completion.
///
/// Every slot the flow hands out must be completed, or its sender waits at
/// the gap for ever and every later packet for that destination piles up
/// behind it. A ticket dropped without being taken, because its job never
/// reached a worker or its worker died, completes its slot as a skip.
#[cfg(target_os = "macos")]
struct MacSeqTicket {
flow: Option<Arc<MacSequencedSendFlow>>,
seq: u64,
}
#[cfg(target_os = "macos")]
impl MacSeqTicket {
fn reserve(flow: Arc<MacSequencedSendFlow>) -> Self {
let seq = flow.reserve_seq();
Self {
flow: Some(flow),
seq,
}
}
/// Take the slot, leaving its completion to the caller.
fn take(mut self) -> (Arc<MacSequencedSendFlow>, u64) {
let flow = self.flow.take().expect("a ticket is taken once");
(flow, self.seq)
}
}
#[cfg(target_os = "macos")]
impl Drop for MacSeqTicket {
fn drop(&mut self) {
if let Some(flow) = self.flow.take() {
flow.complete_skip(self.seq);
}
}
}
/// Closes its flow when the flow's sender thread leaves `run`, including by
/// a panic, so completions waiting for room are released instead of waiting
/// on a sender that will never drain.
#[cfg(target_os = "macos")]
struct MacSenderExit<'a>(&'a MacSequencedSendFlow);
#[cfg(target_os = "macos")]
impl Drop for MacSenderExit<'_> {
fn drop(&mut self) {
self.0.close();
}
}
#[cfg(target_os = "macos")]
enum MacSendItem {
Packet {
@@ -817,37 +889,63 @@ impl MacSequencedSendFlow {
return false;
}
let state = self.state.lock().expect("mac send flow state poisoned");
let state = self.state.lock().unwrap_or_else(PoisonError::into_inner);
state.pending.is_empty()
&& state.next_send_seq == self.next_seq.load(std::sync::atomic::Ordering::Relaxed)
}
fn close(&self) {
let mut state = self.state.lock().expect("mac send flow state poisoned");
let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner);
state.closed = true;
drop(state);
self.ready_cv.notify_one();
self.space_cv.notify_all();
}
fn is_closed(&self) -> bool {
self.state
.lock()
.unwrap_or_else(PoisonError::into_inner)
.closed
}
/// Complete one slot as a skip without waiting for room, so a caller on
/// the rx_loop never blocks here.
fn complete_skip(&self, seq: u64) {
let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner);
if state.closed {
return;
}
let wakes_sender = seq == state.next_send_seq;
state.pending.insert(seq, MacSendItem::Skip);
drop(state);
if wakes_sender {
self.ready_cv.notify_one();
}
}
fn complete_many(&self, items: Vec<(u64, MacSendItem)>) {
const PENDING_CAP: usize = 4096;
if items.is_empty() {
return;
}
let mut state = self.state.lock().expect("mac send flow state poisoned");
if state.closed {
return;
}
let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner);
let mut wakes_sender = false;
for (seq, item) in items {
while state.pending.len() >= PENDING_CAP && seq != state.next_send_seq && !wakes_sender
while !state.closed
&& state.pending.len() >= PENDING_CAP
&& seq != state.next_send_seq
&& !wakes_sender
{
state = self
.space_cv
.wait(state)
.expect("mac send flow state poisoned");
.unwrap_or_else(PoisonError::into_inner);
}
// The sender has exited; what is left of `items` is dropped.
if state.closed {
return;
}
if seq == state.next_send_seq {
wakes_sender = true;
@@ -861,6 +959,8 @@ impl MacSequencedSendFlow {
}
fn run(self: Arc<Self>) {
// Closes the flow however this returns, a panic included.
let _exit = MacSenderExit(&self);
trace!(
socket_fd = self.key.socket_fd,
connected_fd = ?self.key.connected_fd,
@@ -927,10 +1027,10 @@ impl MacSequencedSendFlow {
#[cfg(target_os = "macos")]
fn push_mac_completion(
groups: &mut Vec<MacCompletionGroup>,
flow: Arc<MacSequencedSendFlow>,
seq: u64,
ticket: MacSeqTicket,
item: MacSendItem,
) {
let (flow, seq) = ticket.take();
if let Some(group) = groups
.iter_mut()
.find(|group| Arc::ptr_eq(&group.flow, &flow))
@@ -1058,11 +1158,7 @@ fn flush_batch_sync(
for queued in batch.drain(..) {
#[cfg(target_os = "macos")]
let QueuedFmpSendJob {
job,
macos_flow,
macos_seq,
} = queued;
let QueuedFmpSendJob { job, macos_ticket } = queued;
#[cfg(not(target_os = "macos"))]
let QueuedFmpSendJob { job } = queued;
@@ -1087,13 +1183,8 @@ fn flush_batch_sync(
|| fsp.plaintext_offset > wire_buf.len()
{
#[cfg(target_os = "macos")]
if let Some(flow) = macos_flow.as_ref() {
push_mac_completion(
&mut macos_completions,
Arc::clone(flow),
macos_seq,
MacSendItem::Skip,
);
if let Some(ticket) = macos_ticket {
push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip);
}
continue;
}
@@ -1111,13 +1202,8 @@ fn flush_batch_sync(
Ok(tag) => tag,
Err(_) => {
#[cfg(target_os = "macos")]
if let Some(flow) = macos_flow.as_ref() {
push_mac_completion(
&mut macos_completions,
Arc::clone(flow),
macos_seq,
MacSendItem::Skip,
);
if let Some(ticket) = macos_ticket {
push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip);
}
continue;
}
@@ -1142,8 +1228,8 @@ fn flush_batch_sync(
Ok(tag) => tag,
Err(_) => {
#[cfg(target_os = "macos")]
if let Some(flow) = macos_flow {
push_mac_completion(&mut macos_completions, flow, macos_seq, MacSendItem::Skip);
if let Some(ticket) = macos_ticket {
push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip);
}
continue;
}
@@ -1152,11 +1238,10 @@ fn flush_batch_sync(
wire_buf.extend_from_slice(tag.as_ref());
#[cfg(target_os = "macos")]
if let Some(flow) = macos_flow {
if let Some(ticket) = macos_ticket {
push_mac_completion(
&mut macos_completions,
flow,
macos_seq,
ticket,
MacSendItem::Packet {
packet: wire_buf,
drop_on_backpressure,
@@ -1215,7 +1300,7 @@ fn flush_batch_sync(
#[cfg(target_os = "macos")]
for group in macos_completions {
group.flow.complete_many(group.items);
group.deliver();
}
drop(_t); // close the encrypt timer before we open the send timer
@@ -2578,3 +2663,232 @@ mod mac_queue_tests {
assert!(dropped.is_ok(), "receiver drop panicked on a poisoned lock");
}
}
/// The opt-in ordered sender's completion contract: every reserved slot is
/// completed, a dead sender releases its waiters, and nothing on the rx_loop
/// waits on either. Every wait is bounded so a regression fails, not hangs.
#[cfg(all(test, target_os = "macos"))]
mod mac_ordered_tests {
use super::*;
use crate::transport::udp::io::UdpRawSocket;
use ring::aead::{LessSafeKey, UnboundKey};
use std::net::UdpSocket;
use std::sync::mpsc;
use std::thread;
use std::time::{Duration, Instant};
const WAIT: Duration = Duration::from_secs(5);
const PENDING_CAP: u64 = 4096;
fn job(socket: &AsyncUdpSocket, dest: SocketAddr, counter: u64) -> FmpSendJob {
let key = UnboundKey::new(&ring::aead::CHACHA20_POLY1305, &[7u8; 32]).expect("key");
let mut wire_buf = Vec::with_capacity(ESTABLISHED_HEADER_SIZE + 4 + crate::noise::TAG_SIZE);
wire_buf.extend_from_slice(&[0xA5; ESTABLISHED_HEADER_SIZE]);
wire_buf.extend_from_slice(&counter.to_le_bytes()[..4]);
FmpSendJob {
cipher: LessSafeKey::new(key),
counter,
wire_buf,
fsp_seal: None,
socket: socket.clone(),
dest_addr: dest,
connected_socket: None,
drop_on_backpressure: false,
queued_at: None,
}
}
/// Reserve `n` slots on `flow` and complete them as skips, which fills
/// it to `n` pending items when an earlier slot is still open.
fn skip_reserved(flow: &MacSequencedSendFlow, n: u64) {
let skips = (0..n)
.map(|_| (flow.reserve_seq(), MacSendItem::Skip))
.collect();
flow.complete_many(skips);
}
struct Rig {
_rt: tokio::runtime::Runtime,
recv: UdpSocket,
dest: SocketAddr,
socket: AsyncUdpSocket,
flows: MacSequencedSendFlows,
}
impl Rig {
fn new() -> Self {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_io()
.build()
.expect("tokio rt");
let enter = rt.enter();
let recv = UdpSocket::bind("127.0.0.1:0").expect("bind recv");
recv.set_read_timeout(Some(WAIT)).expect("read timeout");
let dest = recv.local_addr().expect("recv addr");
let socket = UdpRawSocket::open("127.0.0.1:0".parse().unwrap(), 1 << 20, 1 << 20)
.expect("open send socket")
.into_async()
.expect("into_async");
drop(enter);
Self {
_rt: rt,
recv,
dest,
socket,
flows: MacSequencedSendFlows::default(),
}
}
fn job(&self, counter: u64) -> FmpSendJob {
job(&self.socket, self.dest, counter)
}
fn flow(&self) -> Arc<MacSequencedSendFlow> {
self.flows.flow_for(&self.job(0))
}
fn sequenced(&self, counter: u64) -> QueuedFmpSendJob {
let job = self.job(counter);
let flow = self.flows.flow_for(&job);
QueuedFmpSendJob::macos_sequenced(job, flow)
}
fn received(&self) -> bool {
let mut buf = [0u8; 256];
self.recv.recv_from(&mut buf).is_ok()
}
}
#[test]
fn a_sequenced_job_dropped_unsent_lets_its_flow_send_what_follows() {
let rig = Rig::new();
drop(rig.sequenced(1));
let mut batch = vec![rig.sequenced(2)];
flush_batch_sync(&mut batch).expect("flush");
assert!(rig.received(), "the flow stalled at the dropped job's slot");
}
#[test]
fn a_dead_workers_queue_releases_the_slots_it_held() {
let rig = Rig::new();
let (tx, rx) = mac_worker_channel::<QueuedFmpSendJob>(4);
assert!(tx.try_push(rig.sequenced(1)).is_ok());
assert!(tx.try_push(rig.sequenced(2)).is_ok());
drop(rx);
assert!(matches!(
tx.try_push(rig.sequenced(3)),
Err(MacWorkerTryPushError::Closed)
));
let mut batch = vec![rig.sequenced(4)];
flush_batch_sync(&mut batch).expect("flush");
assert!(
rig.received(),
"the flow stalled at a slot the dead queue held"
);
}
#[test]
fn a_completion_group_dropped_undelivered_lets_its_flow_send_what_follows() {
let rig = Rig::new();
let flow = rig.flow();
let group = MacCompletionGroup {
flow: Arc::clone(&flow),
items: vec![(
flow.reserve_seq(),
MacSendItem::Packet {
packet: b"undelivered".to_vec(),
drop_on_backpressure: false,
},
)],
};
drop(group);
let mut batch = vec![rig.sequenced(2)];
flush_batch_sync(&mut batch).expect("flush");
assert!(
rig.received(),
"the flow stalled at a slot an undelivered group held"
);
}
#[test]
fn a_completion_waiting_for_room_returns_when_its_flow_closes() {
let rig = Rig::new();
let flow = rig.flow();
let gap = QueuedFmpSendJob::macos_sequenced(rig.job(0), Arc::clone(&flow));
skip_reserved(&flow, PENDING_CAP);
let late = flow.reserve_seq();
let (done_tx, done_rx) = mpsc::channel();
let waiter = Arc::clone(&flow);
thread::spawn(move || {
waiter.complete_many(vec![(late, MacSendItem::Skip)]);
let _ = done_tx.send(());
});
thread::sleep(Duration::from_millis(200));
assert!(
done_rx.try_recv().is_err(),
"a full flow took a completion without room"
);
flow.close();
done_rx
.recv_timeout(WAIT)
.expect("the completion still waited after its flow closed");
drop(gap);
}
#[test]
fn a_sender_thread_that_panics_closes_its_flow() {
let rig = Rig::new();
let flow = rig.flow();
// Poison the flow's state lock, then wake the sender: its own
// `expect` on the lock panics inside `run`.
let poisoner = Arc::clone(&flow);
let poisoned = thread::spawn(move || {
let _state = poisoner.state.lock().unwrap();
panic!("poison the flow state lock");
})
.join();
assert!(poisoned.is_err());
flow.complete_skip(0);
let deadline = Instant::now() + WAIT;
while !flow.is_closed() && Instant::now() < deadline {
thread::sleep(Duration::from_millis(10));
}
assert!(
flow.is_closed(),
"a sender thread that panicked left its flow open"
);
}
#[test]
fn a_flow_whose_sender_has_exited_is_replaced() {
let rig = Rig::new();
let first = rig.flow();
first.close();
let second = rig.flow();
assert!(
!Arc::ptr_eq(&first, &second),
"a closed flow was handed out again"
);
let mut batch = vec![QueuedFmpSendJob::macos_sequenced(rig.job(3), second)];
flush_batch_sync(&mut batch).expect("flush");
assert!(rig.received(), "the replacement flow sent nothing");
}
#[test]
fn dropping_a_sequenced_job_never_waits_for_room() {
let rig = Rig::new();
let flow = rig.flow();
let gap = QueuedFmpSendJob::macos_sequenced(rig.job(0), Arc::clone(&flow));
skip_reserved(&flow, PENDING_CAP);
let dropped = QueuedFmpSendJob::macos_sequenced(rig.job(1), Arc::clone(&flow));
let (done_tx, done_rx) = mpsc::channel();
thread::spawn(move || {
drop(dropped);
let _ = done_tx.send(());
});
done_rx
.recv_timeout(WAIT)
.expect("dropping a job waited for room in a full flow");
drop(gap);
}
}
+1 -1
View File
@@ -1453,7 +1453,7 @@ impl Node {
NodeError::NoTransportForType(format!("invalid MAC in '{}': {}", addr_str, e))
})?;
Ok((transport_id, TransportAddr::from_bytes(&mac)))
Ok((transport_id, TransportAddr::from_mac(mac)))
}
#[cfg(not(any(target_os = "linux", target_os = "macos")))]
{
+47 -1
View File
@@ -766,7 +766,12 @@ impl PeerMachine {
/// Take the handshake crypto carrier off the machine (promotion and
/// teardown consume it by value).
pub(crate) fn take_leg(&mut self) -> Option<HandshakeCrypto> {
self.leg.take()
// On promotion the machine lives on as the peer's control machine, so
// the slot is cleared as the leg leaves, or it would keep the session's
// traffic keys. Callers unwrap the result at once, which `take_cleared`
// warns may leave a second stack copy; clearing here still covers every
// caller's heap slot, and a stack copy is a move like any other.
noise::take_cleared(&mut self.leg)
}
/// Attach a handshake crypto carrier to the machine.
@@ -3989,6 +3994,47 @@ mod tests {
assert_eq!(decrypted, plaintext);
}
/// Non-zero bytes left in a leg slot, read in place.
#[cfg(not(debug_assertions))]
fn leg_slot_nonzero_bytes(slot: &Option<HandshakeCrypto>) -> usize {
let base = (slot as *const Option<HandshakeCrypto>).cast::<u8>();
(0..std::mem::size_of::<Option<HandshakeCrypto>>())
// SAFETY: `base` comes from a live reference, so every byte of the
// slot is in bounds; volatile so the read is of memory as it is.
.filter(|&i| unsafe { std::ptr::read_volatile(base.add(i)) } != 0)
.count()
}
/// Release only: in a debug build `Option::take` leaves stack garbage in
/// the payload of the `None` it writes, so the count means nothing there.
#[cfg(not(debug_assertions))]
#[test]
fn take_leg_leaves_no_session_keys_in_the_slot_it_empties() {
let initiator = Identity::generate();
let responder = Identity::generate();
let responder_id = PeerIdentity::from_pubkey_full(responder.pubkey_full());
let mut ini = outbound_leg(LinkId::new(1), responder_id, 1000);
let mut res = Box::new(inbound_leg(LinkId::new(2), 1000));
let msg1 = ini
.start_handshake(initiator.keypair(), make_epoch(), 1100)
.unwrap();
let msg2 = res
.receive_handshake_init(responder.keypair(), make_epoch(), &msg1, None, 1200)
.unwrap();
let (msg3, _) = ini.complete_handshake(&msg2, None, 1300).unwrap();
res.complete_handshake_msg3(&msg3, 1400).unwrap();
assert!(res.has_session(), "the responder's leg holds the session");
assert!(leg_slot_nonzero_bytes(&res.leg) > 64);
let leg = res.take_leg().expect("the leg was attached");
assert!(leg.noise_session.is_some(), "the session left with the leg");
let left = leg_slot_nonzero_bytes(&res.leg);
assert!(
left <= 8,
"{left} non-zero bytes of the session stayed in the emptied leg slot"
);
}
#[test]
fn test_connection_failure() {
// `mark_failed` releases the leg's Noise handshake handle. The failure
+1 -1
View File
@@ -1274,7 +1274,7 @@ async fn ethernet_receive_loop(
}
};
let bytes = data.len();
let addr = TransportAddr::from_bytes(&src_mac);
let addr = TransportAddr::from_mac(src_mac);
let packet = ReceivedPacket::new(transport_id, addr, data);
trace!(
+12 -1
View File
@@ -155,7 +155,7 @@ impl NeighborBuffer {
/// This line's beacon carries no key, so there is no pubkey hint to
/// attach: the identity is learned from the handshake instead.
fn peer(&self, src_mac: [u8; 6]) -> DiscoveredPeer {
let addr = TransportAddr::from_bytes(&src_mac);
let addr = TransportAddr::from_mac(src_mac);
DiscoveredPeer::new(self.transport_id, addr)
}
}
@@ -243,6 +243,17 @@ mod tests {
assert!(peers.is_empty());
}
#[test]
fn a_discovered_peer_whose_mac_bytes_are_valid_utf8_still_displays_as_a_mac() {
// "2|\u{46c}Zd": six bytes that decode as UTF-8, seen on a veth MAC.
let mac = [0x32, 0x7c, 0xd1, 0xac, 0x5a, 0x64];
assert!(core::str::from_utf8(&mac).is_ok());
let buffer = NeighborBuffer::new(TransportId::new(1));
assert!(buffer.add_peer(mac));
let peers = buffer.take();
assert_eq!(peers[0].addr.to_string(), "32:7c:d1:ac:5a:64");
}
#[test]
fn test_neighbor_buffer_dedup() {
let buffer = NeighborBuffer::new(TransportId::new(1));
+23 -2
View File
@@ -1372,12 +1372,33 @@ mod tests {
#[test]
fn test_transport_addr_mac_display() {
// Raw 6-byte MACs (as Ethernet stores via from_bytes) display in
// standard colon-separated notation, not bare hex.
// Raw 6-byte non-UTF-8 values from from_bytes display in standard
// colon-separated notation, not bare hex.
let mac = TransportAddr::from_bytes(&[0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff]);
assert_eq!(format!("{}", mac), "aa:bb:cc:dd:ee:ff");
}
#[test]
fn a_mac_address_whose_bytes_are_valid_utf8_displays_as_a_mac() {
let bytes = [0x32, 0x7c, 0xd1, 0xac, 0x5a, 0x64];
assert!(core::str::from_utf8(&bytes).is_ok());
let mac = TransportAddr::from_mac(bytes);
assert_eq!(mac.to_string(), "32:7c:d1:ac:5a:64");
assert_eq!(format!("{:?}", mac), "TransportAddr(32:7c:d1:ac:5a:64)");
assert_eq!(mac.as_bytes(), &bytes);
}
#[test]
fn a_mac_address_equals_and_hashes_as_its_bytes() {
use std::hash::BuildHasher;
let bytes = [0x32, 0x7c, 0xd1, 0xac, 0x5a, 0x64];
let mac = TransportAddr::from_mac(bytes);
let raw = TransportAddr::from_bytes(&bytes);
assert_eq!(mac, raw);
let hasher = std::collections::hash_map::RandomState::new();
assert_eq!(hasher.hash_one(&mac), hasher.hash_one(&raw));
}
#[test]
fn test_transport_addr_non_mac_binary_is_bare_hex() {
// Non-6-byte non-UTF-8 payloads stay bare hex (no separators).
+71 -24
View File
@@ -91,11 +91,31 @@ impl fmt::Display for LinkDirection {
///
/// Each transport type interprets this differently:
/// - UDP/TCP: "host:port" (IP address or DNS hostname)
/// - Ethernet: MAC address (6 bytes)
/// - Ethernet: MAC address (6 bytes), built with [`TransportAddr::from_mac`]
/// so it always displays as a MAC, whatever its bytes decode as
///
/// The immutable bytes are shared across clones.
#[derive(Clone, PartialEq, Eq, Hash)]
pub struct TransportAddr(Arc<[u8]>);
#[derive(Clone)]
pub struct TransportAddr {
bytes: Arc<[u8]>,
/// Set only by [`TransportAddr::from_mac`]: the bytes are a MAC address
/// and always display as one. Not part of equality or hashing.
mac: bool,
}
impl PartialEq for TransportAddr {
fn eq(&self, other: &Self) -> bool {
self.bytes == other.bytes
}
}
impl Eq for TransportAddr {}
impl core::hash::Hash for TransportAddr {
fn hash<H: core::hash::Hasher>(&self, state: &mut H) {
self.bytes.hash(state);
}
}
impl TransportAddr {
/// Create a transport address from raw bytes.
@@ -103,12 +123,26 @@ impl TransportAddr {
/// Copies the bytes into shared storage; the vector's allocation is not
/// reused. Prefer [`Self::from_bytes`] when a byte slice is already available.
pub fn new(bytes: Vec<u8>) -> Self {
Self(bytes.into())
Self {
bytes: bytes.into(),
mac: false,
}
}
/// Create a transport address from a byte slice.
pub fn from_bytes(bytes: &[u8]) -> Self {
Self(Arc::from(bytes))
Self {
bytes: Arc::from(bytes),
mac: false,
}
}
/// Create an Ethernet transport address from a MAC address.
pub fn from_mac(mac: [u8; 6]) -> Self {
Self {
bytes: Arc::from(&mac[..]),
mac: true,
}
}
/// Create a transport address from a string.
@@ -118,53 +152,66 @@ impl TransportAddr {
/// Get the raw bytes.
pub fn as_bytes(&self) -> &[u8] {
&self.0
&self.bytes
}
/// Try to interpret as a UTF-8 string.
pub fn as_str(&self) -> Option<&str> {
core::str::from_utf8(&self.0).ok()
core::str::from_utf8(&self.bytes).ok()
}
/// Get the length in bytes.
pub fn len(&self) -> usize {
self.0.len()
self.bytes.len()
}
/// Check if empty.
pub fn is_empty(&self) -> bool {
self.0.is_empty()
self.bytes.is_empty()
}
}
impl TransportAddr {
/// Write the bytes as a colon-separated MAC address.
fn write_mac(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
for (i, byte) in self.bytes.iter().enumerate() {
if i > 0 {
write!(f, ":")?;
}
write!(f, "{:02x}", byte)?;
}
Ok(())
}
}
impl fmt::Debug for TransportAddr {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
if self.mac {
write!(f, "TransportAddr(")?;
self.write_mac(f)?;
return write!(f, ")");
}
match self.as_str() {
Some(s) => write!(f, "TransportAddr(\"{}\")", s),
None => write!(f, "TransportAddr({:?})", self.0),
None => write!(f, "TransportAddr({:?})", self.bytes),
}
}
}
impl fmt::Display for TransportAddr {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
// Best-effort display as string if valid UTF-8. Otherwise render a
// 6-byte payload as a colon-separated MAC (standard Unix notation,
// matching BLE addrs, `ip link`/`ip neigh`, and packet logs), and
// any other non-UTF-8 byte string as bare hex.
// An Ethernet address is a MAC whatever its bytes decode as. Any
// other address displays as a string if it is valid UTF-8; otherwise
// a 6-byte payload is rendered as a colon-separated MAC and any other
// byte string as bare hex.
if self.mac {
return self.write_mac(f);
}
match self.as_str() {
Some(s) => write!(f, "{}", s),
None if self.0.len() == 6 => {
for (i, byte) in self.0.iter().enumerate() {
if i > 0 {
write!(f, ":")?;
}
write!(f, "{:02x}", byte)?;
}
Ok(())
}
None if self.bytes.len() == 6 => self.write_mac(f),
None => {
for byte in self.0.iter() {
for byte in self.bytes.iter() {
write!(f, "{:02x}", byte)?;
}
Ok(())
+37
View File
@@ -143,6 +143,43 @@ mod tests {
assert!(reuse_port, "the adopted socket must carry SO_REUSEPORT");
}
/// A datagram sent to a port nobody holds draws an ICMP port-unreachable.
/// Windows reports one, by default, as `WSAECONNRESET` on the socket's next
/// receive; this socket is shared by every peer, so the next datagram from
/// any of them must come through instead. Checked for both ways a socket
/// is set up.
#[cfg(windows)]
#[tokio::test]
async fn a_port_unreachable_answer_is_not_a_receive_error_on_windows() {
let opened = UdpRawSocket::open("127.0.0.1:0".parse().unwrap(), 65536, 65536)
.expect("failed to open");
let plain = std::net::UdpSocket::bind("127.0.0.1:0").expect("failed to bind");
let adopted = UdpRawSocket::adopt(plain, 65536, 65536).expect("failed to adopt");
for (how, raw) in [("open", opened), ("adopt", adopted)] {
let addr = raw.local_addr();
let sock = raw.into_async().expect("into_async");
let gone = std::net::UdpSocket::bind("127.0.0.1:0").expect("bind the closed port");
let gone_addr = gone.local_addr().expect("closed port address");
drop(gone);
sock.send_to(b"nobody", &gone_addr).await.expect("send_to");
// Let the port-unreachable arrive before the datagram that follows.
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
let peer = std::net::UdpSocket::bind("127.0.0.1:0").expect("bind the peer");
peer.send_to(b"peer", addr).expect("peer send");
let mut buf = [0u8; 16];
let got =
tokio::time::timeout(std::time::Duration::from_secs(5), sock.recv_from(&mut buf))
.await
.unwrap_or_else(|_| panic!("{how}: nothing received"));
let (n, from, _) = got.unwrap_or_else(|e| panic!("{how}: receive failed: {e}"));
assert_eq!(&buf[..n], b"peer", "{how}");
assert_eq!(from, peer.local_addr().unwrap(), "{how}");
}
}
#[tokio::test]
async fn test_async_udp_socket_send_recv() {
let sock1 = UdpRawSocket::open("127.0.0.1:0".parse().unwrap(), 65536, 65536)
+43
View File
@@ -8,7 +8,10 @@
use crate::transport::TransportError;
use socket2::{Domain, Protocol, Socket, Type};
use std::net::SocketAddr;
use std::os::windows::io::AsRawSocket;
use std::sync::Arc;
use tracing::warn;
use windows_sys::Win32::Networking::WinSock::{SIO_UDP_CONNRESET, SOCKET, SOCKET_ERROR, WSAIoctl};
/// UDP socket wrapper (Windows).
///
@@ -44,6 +47,7 @@ impl UdpRawSocket {
sock.bind(&bind_addr.into())
.map_err(|e| TransportError::StartFailed(format!("bind failed: {}", e)))?;
ignore_port_unreachable(&sock);
// Set socket buffer sizes
sock.set_recv_buffer_size(recv_buf_size)
@@ -75,6 +79,7 @@ impl UdpRawSocket {
sock.set_nonblocking(true)
.map_err(|e| TransportError::StartFailed(format!("set nonblocking failed: {}", e)))?;
ignore_port_unreachable(&sock);
sock.set_recv_buffer_size(recv_buf_size)
.map_err(|e| TransportError::StartFailed(format!("set recv buffer: {}", e)))?;
@@ -126,6 +131,44 @@ impl UdpRawSocket {
}
}
/// Stop an ICMP port-unreachable from failing the socket's next receive.
///
/// By default Windows reports a port-unreachable answering any datagram the
/// socket sent as `WSAECONNRESET` on its next `recv_from`. The socket is shared
/// by every peer, so one unreachable address would put errors into the receive
/// loop that all the other peers' traffic arrives through. `SIO_UDP_CONNRESET`
/// set to false turns the report off. A socket that refuses it still works and
/// still sees the errors, so a failure is logged rather than returned.
fn ignore_port_unreachable(sock: &Socket) {
let report: u32 = 0;
let mut returned: u32 = 0;
// SAFETY: the socket is open for the duration of the call, the input
// buffer is a live u32 of the stated size, no output buffer is passed, and
// the call is synchronous (no OVERLAPPED, no completion routine).
let rc = unsafe {
WSAIoctl(
sock.as_raw_socket() as SOCKET,
SIO_UDP_CONNRESET,
(&report as *const u32).cast(),
std::mem::size_of::<u32>() as u32,
std::ptr::null_mut(),
0,
&mut returned,
std::ptr::null_mut(),
None,
)
};
if rc == SOCKET_ERROR {
// Read before `warn!`, whose level and dispatcher checks run first and
// could replace the error code WSAIoctl left.
let err = std::io::Error::last_os_error();
warn!(
error = %err,
"Could not turn off ICMP port-unreachable reports on the UDP socket"
);
}
}
/// Async UDP socket wrapper (Windows).
///
/// Uses `tokio::net::UdpSocket` directly. Kernel drop counting
+14 -6
View File
@@ -699,16 +699,24 @@ run_tests() {
# Debug-only helpers (anything behind #[cfg(debug_assertions)]) vanish in a
# release build, so a test calling one without the same gate breaks a build
# nothing here ever performs: every run above compiles the test target in
# debug. Compile it in release too, without running it — the point is that
# it builds at all. Mirrored in .github/workflows/ci.yml; check-ci-parity.sh
# compares integration suites only and would not catch a stage added to one
# runner and not the other.
info "cargo test --release --lib --no-run"
if cargo test --release --lib --no-run 2>&1; then
# debug. Compile it in release too. Of what it builds, run only the leg-slot
# residue test, which is release-only because only an optimised build
# leaves the slot in a state worth measuring; the grep fails the stage if
# that test did not run, since a name filter that matches nothing passes.
# Mirrored in .github/workflows/ci.yml; check-ci-parity.sh compares
# integration suites only and would not catch a stage added to one runner
# and not the other.
local residue_test="peer::machine::tests::take_leg_leaves_no_session_keys_in_the_slot_it_empties"
local release_log
release_log=$(mktemp "/tmp/ci-release-lib-tests.XXXXXX")
info "cargo test --release --lib -- --exact $residue_test"
if cargo test --release --lib -- --exact "$residue_test" 2>&1 | tee "$release_log" \
&& grep -q '^test result: ok\. 1 passed;' "$release_log"; then
record "release-test-compile" 0
else
record "release-test-compile" 1
fi
rm -f "$release_log"
}
# ── Stage 3: Integration Tests ─────────────────────────────────────────────