diff --git a/src/node/encrypt_worker.rs b/src/node/encrypt_worker.rs index 8f5e462d..d8cd4da7 100644 --- a/src/node/encrypt_worker.rs +++ b/src/node/encrypt_worker.rs @@ -156,9 +156,7 @@ pub(crate) struct FspSealJob { struct QueuedFmpSendJob { job: FmpSendJob, #[cfg(target_os = "macos")] - macos_flow: Option>, - #[cfg(target_os = "macos")] - macos_seq: u64, + macos_ticket: Option, } 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) -> 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 MacWorkerSender { .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 MacWorkerSender { .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>, + seq: u64, +} + +#[cfg(target_os = "macos")] +impl MacSeqTicket { + fn reserve(flow: Arc) -> 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, 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) { + // 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, - flow: Arc, - 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 { + 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::(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); + } +}