mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-05 11:08:25 +00:00
Encrypt packets on the main loop when their encrypt worker has exited
A packet hashed to an encrypt worker that had exited was dropped after its counters were reserved and, on the session data path, after the link and session statistics had already counted it as sent. Every peer that hashed to the dead worker was cut off until the process restarted. Dispatch now hands a job its worker refused back to the caller instead of dropping it. Both send sites seal the job on the main loop with the counters it already reserved, so no counter is skipped and the receiver sees no gap, and send it through the transport the way the inline path does. On the link-message path the statistics count the bytes actually sent; on the session data path they were recorded before dispatch and now describe a packet that was sent. A job is handed back only when no worker holds it, so each reserved counter is still used at most once. In the macOS ordered sender the job's place in its flow is released before it is handed back, so later packets for that flow are not held behind it. The worker and the main loop share one seal function. It refuses a job whose offsets do not fit its buffer instead of indexing out of bounds, since a panic there would now end the node rather than one worker.
This commit is contained in:
+390
-104
@@ -57,7 +57,7 @@ use crate::proto::fmp::wire::ESTABLISHED_HEADER_SIZE;
|
|||||||
use crate::proto::fsp::wire::FSP_HEADER_SIZE;
|
use crate::proto::fsp::wire::FSP_HEADER_SIZE;
|
||||||
use crate::transport::udp::io::AsyncUdpSocket;
|
use crate::transport::udp::io::AsyncUdpSocket;
|
||||||
#[cfg(not(target_os = "macos"))]
|
#[cfg(not(target_os = "macos"))]
|
||||||
use crossbeam_channel::{Receiver, Sender, TrySendError, bounded};
|
use crossbeam_channel::{Receiver, SendError, Sender, TrySendError, bounded};
|
||||||
use ring::aead::{Aad, LessSafeKey, Nonce};
|
use ring::aead::{Aad, LessSafeKey, Nonce};
|
||||||
#[cfg(any(target_os = "macos", test))]
|
#[cfg(any(target_os = "macos", test))]
|
||||||
use std::collections::VecDeque;
|
use std::collections::VecDeque;
|
||||||
@@ -181,6 +181,24 @@ impl QueuedFmpSendJob {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A queued job its worker refused because the worker has exited. Boxed so
|
||||||
|
/// the dispatch path's `Result` stays small; the allocation happens only on a
|
||||||
|
/// refusal.
|
||||||
|
struct Refused(Box<QueuedFmpSendJob>);
|
||||||
|
|
||||||
|
impl Refused {
|
||||||
|
/// The job, with any macOS ordered-flow slot it held released as a skip.
|
||||||
|
fn into_job(self) -> Box<FmpSendJob> {
|
||||||
|
#[cfg(target_os = "macos")]
|
||||||
|
let QueuedFmpSendJob { job, macos_ticket } = *self.0;
|
||||||
|
#[cfg(not(target_os = "macos"))]
|
||||||
|
let QueuedFmpSendJob { job } = *self.0;
|
||||||
|
#[cfg(target_os = "macos")]
|
||||||
|
drop(macos_ticket);
|
||||||
|
Box::new(job)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Handle to the encrypt worker pool. Dispatches jobs **hash-by-
|
/// Handle to the encrypt worker pool. Dispatches jobs **hash-by-
|
||||||
/// destination** across N worker tasks via per-worker bounded
|
/// destination** across N worker tasks via per-worker bounded
|
||||||
/// crossbeam channels. The bounded queue intentionally backpressures
|
/// crossbeam channels. The bounded queue intentionally backpressures
|
||||||
@@ -242,14 +260,17 @@ struct MacWorkerQueueState<T> {
|
|||||||
closed: bool,
|
closed: bool,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Why `try_push` did not queue a job. Both variants hand the job back.
|
||||||
#[cfg(any(target_os = "macos", test))]
|
#[cfg(any(target_os = "macos", test))]
|
||||||
enum MacWorkerTryPushError<T> {
|
enum MacWorkerTryPushError<T> {
|
||||||
Full(Box<T>),
|
Full(Box<T>),
|
||||||
Closed,
|
/// The receiver is gone: the worker has exited.
|
||||||
|
Closed(Box<T>),
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// `push_blocking` found the receiver gone; the job is handed back.
|
||||||
#[cfg(any(target_os = "macos", test))]
|
#[cfg(any(target_os = "macos", test))]
|
||||||
struct MacWorkerPushError;
|
struct MacWorkerPushError<T>(Box<T>);
|
||||||
|
|
||||||
#[cfg(any(target_os = "macos", test))]
|
#[cfg(any(target_os = "macos", test))]
|
||||||
fn mac_worker_channel<T>(cap: usize) -> (MacWorkerSender<T>, MacWorkerReceiver<T>) {
|
fn mac_worker_channel<T>(cap: usize) -> (MacWorkerSender<T>, MacWorkerReceiver<T>) {
|
||||||
@@ -280,10 +301,10 @@ impl<T> MacWorkerSender<T> {
|
|||||||
.lock()
|
.lock()
|
||||||
.expect("encrypt worker queue poisoned");
|
.expect("encrypt worker queue poisoned");
|
||||||
if state.closed {
|
if state.closed {
|
||||||
// Outside the lock: dropping a sequenced job completes its slot.
|
// The caller drops or reuses the job, outside the lock: dropping a
|
||||||
|
// sequenced job completes its slot.
|
||||||
drop(state);
|
drop(state);
|
||||||
drop(job);
|
return Err(MacWorkerTryPushError::Closed(Box::new(job)));
|
||||||
return Err(MacWorkerTryPushError::Closed);
|
|
||||||
}
|
}
|
||||||
if state.queue.len() >= self.inner.cap {
|
if state.queue.len() >= self.inner.cap {
|
||||||
return Err(MacWorkerTryPushError::Full(Box::new(job)));
|
return Err(MacWorkerTryPushError::Full(Box::new(job)));
|
||||||
@@ -298,7 +319,7 @@ impl<T> MacWorkerSender<T> {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn push_blocking(&self, job: T) -> Result<(), MacWorkerPushError> {
|
fn push_blocking(&self, job: T) -> Result<(), MacWorkerPushError<T>> {
|
||||||
let mut state = self
|
let mut state = self
|
||||||
.inner
|
.inner
|
||||||
.state
|
.state
|
||||||
@@ -307,8 +328,7 @@ impl<T> MacWorkerSender<T> {
|
|||||||
loop {
|
loop {
|
||||||
if state.closed {
|
if state.closed {
|
||||||
drop(state);
|
drop(state);
|
||||||
drop(job);
|
return Err(MacWorkerPushError(Box::new(job)));
|
||||||
return Err(MacWorkerPushError);
|
|
||||||
}
|
}
|
||||||
if state.queue.len() < self.inner.cap {
|
if state.queue.len() < self.inner.cap {
|
||||||
let was_empty = state.queue.is_empty();
|
let was_empty = state.queue.is_empty();
|
||||||
@@ -503,8 +523,15 @@ impl EncryptWorkerPool {
|
|||||||
/// kernel `SocketAddr` lands on the same worker and stays in
|
/// kernel `SocketAddr` lands on the same worker and stays in
|
||||||
/// order — required for TCP's fast-retransmit logic above to
|
/// order — required for TCP's fast-retransmit logic above to
|
||||||
/// behave on a single-flow run. The worker handles send errors
|
/// behave on a single-flow run. The worker handles send errors
|
||||||
/// itself via stats counters. A job whose worker has exited is
|
/// itself via stats counters.
|
||||||
/// dropped, counted, and logged at WARN.
|
///
|
||||||
|
/// A job whose worker has exited is never queued: it is counted,
|
||||||
|
/// logged at WARN, and handed back as `Err` so the caller can seal
|
||||||
|
/// and send it on its own path with the counters it already
|
||||||
|
/// reserved. A job is handed back only when no worker has it, so
|
||||||
|
/// each reserved counter is still used at most once. In the macOS
|
||||||
|
/// ordered mode the job's place in its flow is released as a skip
|
||||||
|
/// before it is handed back.
|
||||||
///
|
///
|
||||||
/// Uses `try_send` for the common uncontended case, then blocks
|
/// Uses `try_send` for the common uncontended case, then blocks
|
||||||
/// only when the bounded worker channel is full. These jobs carry
|
/// only when the bounded worker channel is full. These jobs carry
|
||||||
@@ -512,19 +539,22 @@ impl EncryptWorkerPool {
|
|||||||
/// this internal queue makes TCP-over-TUN collapse with avoidable
|
/// this internal queue makes TCP-over-TUN collapse with avoidable
|
||||||
/// retransmits. Blocking here pushes back toward the TUN reader
|
/// retransmits. Blocking here pushes back toward the TUN reader
|
||||||
/// and lets the kernel/app TCP stack pace the flow instead.
|
/// and lets the kernel/app TCP stack pace the flow instead.
|
||||||
pub fn dispatch(&self, job: FmpSendJob) {
|
#[must_use = "a job handed back was not sent; the caller must send it another way"]
|
||||||
|
pub fn dispatch(&self, job: FmpSendJob) -> Result<(), Box<FmpSendJob>> {
|
||||||
let (idx, job) = self.prepare_dispatch(job);
|
let (idx, job) = self.prepare_dispatch(job);
|
||||||
if !self.dispatch_to_worker(idx, job) {
|
let Err(refused) = self.dispatch_to_worker(idx, job) else {
|
||||||
let n = self.workers.note_refused();
|
return Ok(());
|
||||||
if worth_logging(n) {
|
};
|
||||||
warn!(
|
let n = self.workers.note_refused();
|
||||||
pool = "encrypt",
|
if worth_logging(n) {
|
||||||
worker = idx,
|
warn!(
|
||||||
refused = n + 1,
|
pool = "encrypt",
|
||||||
"Encrypt worker has exited; dropping packet"
|
worker = idx,
|
||||||
);
|
refused = n + 1,
|
||||||
}
|
"Encrypt worker has exited; encrypting the packet on the main loop"
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
Err(refused.into_job())
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(target_os = "macos")]
|
#[cfg(target_os = "macos")]
|
||||||
@@ -573,13 +603,13 @@ impl EncryptWorkerPool {
|
|||||||
(h.finish() as usize) % self.workers.len()
|
(h.finish() as usize) % self.workers.len()
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Queue `job` on worker `idx`, returning `false` when that worker has
|
/// Queue `job` on worker `idx`, or hand it back when that worker has
|
||||||
/// exited and the job was not queued.
|
/// exited.
|
||||||
#[cfg(target_os = "macos")]
|
#[cfg(target_os = "macos")]
|
||||||
fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> bool {
|
fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> Result<(), Refused> {
|
||||||
let sender = self.workers.sender(idx);
|
let sender = self.workers.sender(idx);
|
||||||
match sender.try_push(job) {
|
match sender.try_push(job) {
|
||||||
Ok(()) => true,
|
Ok(()) => Ok(()),
|
||||||
Err(MacWorkerTryPushError::Full(job)) => {
|
Err(MacWorkerTryPushError::Full(job)) => {
|
||||||
static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0);
|
static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0);
|
||||||
let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
||||||
@@ -590,19 +620,21 @@ impl EncryptWorkerPool {
|
|||||||
"EncryptWorker channel full; applying outbound backpressure"
|
"EncryptWorker channel full; applying outbound backpressure"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
sender.push_blocking(*job).is_ok()
|
sender
|
||||||
|
.push_blocking(*job)
|
||||||
|
.map_err(|MacWorkerPushError(job)| Refused(job))
|
||||||
}
|
}
|
||||||
Err(MacWorkerTryPushError::Closed) => false,
|
Err(MacWorkerTryPushError::Closed(job)) => Err(Refused(job)),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Queue `job` on worker `idx`, returning `false` when that worker has
|
/// Queue `job` on worker `idx`, or hand it back when that worker has
|
||||||
/// exited and the job was not queued.
|
/// exited.
|
||||||
#[cfg(not(target_os = "macos"))]
|
#[cfg(not(target_os = "macos"))]
|
||||||
fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> bool {
|
fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> Result<(), Refused> {
|
||||||
let sender = self.workers.sender(idx);
|
let sender = self.workers.sender(idx);
|
||||||
match sender.try_send(job) {
|
match sender.try_send(job) {
|
||||||
Ok(()) => true,
|
Ok(()) => Ok(()),
|
||||||
Err(TrySendError::Full(job)) => {
|
Err(TrySendError::Full(job)) => {
|
||||||
static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0);
|
static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0);
|
||||||
let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
||||||
@@ -613,9 +645,11 @@ impl EncryptWorkerPool {
|
|||||||
"EncryptWorker channel full; applying outbound backpressure"
|
"EncryptWorker channel full; applying outbound backpressure"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
sender.send(job).is_ok()
|
sender
|
||||||
|
.send(job)
|
||||||
|
.map_err(|SendError(job)| Refused(Box::new(job)))
|
||||||
}
|
}
|
||||||
Err(TrySendError::Disconnected(_)) => false,
|
Err(TrySendError::Disconnected(job)) => Err(Refused(Box::new(job))),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -627,6 +661,21 @@ impl EncryptWorkerPool {
|
|||||||
pub(crate) fn for_test(plan: Vec<TestWorker>) -> Self {
|
pub(crate) fn for_test(plan: Vec<TestWorker>) -> Self {
|
||||||
Self::start_with(plan.len(), test_spawner(plan, spawn_worker))
|
Self::start_with(plan.len(), test_spawner(plan, spawn_worker))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The worker a job to `dest` is dispatched to, where that depends on
|
||||||
|
/// `dest` alone. On macOS it also depends on the sending sockets, or on a
|
||||||
|
/// round-robin in the ordered mode, so there it is `None`.
|
||||||
|
pub(crate) fn worker_index_for_dest(&self, dest: SocketAddr) -> Option<usize> {
|
||||||
|
#[cfg(target_os = "macos")]
|
||||||
|
{
|
||||||
|
let _ = dest;
|
||||||
|
None
|
||||||
|
}
|
||||||
|
#[cfg(not(target_os = "macos"))]
|
||||||
|
{
|
||||||
|
Some(self.worker_index_for(dest))
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(target_os = "macos")]
|
#[cfg(target_os = "macos")]
|
||||||
@@ -1145,6 +1194,89 @@ fn run_worker_macos(idx: usize, rx: MacWorkerReceiver<QueuedFmpSendJob>) {
|
|||||||
trace!(worker = idx, "FMP encrypt worker thread exiting");
|
trace!(worker = idx, "FMP encrypt worker thread exiting");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Why a job could not be sealed.
|
||||||
|
#[derive(Debug, thiserror::Error)]
|
||||||
|
pub(crate) enum SealError {
|
||||||
|
/// The offsets the job carries do not fit its buffer.
|
||||||
|
#[error("job layout does not fit its buffer")]
|
||||||
|
Layout,
|
||||||
|
/// The AEAD refused to seal.
|
||||||
|
#[error("AEAD seal failed")]
|
||||||
|
Aead,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl FmpSendJob {
|
||||||
|
/// Seal this job on the calling thread, as its worker would have, and
|
||||||
|
/// return the wire packet. For a job its worker refused: the counters it
|
||||||
|
/// carries were reserved for it and are used here, once.
|
||||||
|
pub(crate) fn seal_inline(self) -> Result<Vec<u8>, SealError> {
|
||||||
|
let FmpSendJob {
|
||||||
|
cipher,
|
||||||
|
counter,
|
||||||
|
mut wire_buf,
|
||||||
|
fsp_seal,
|
||||||
|
..
|
||||||
|
} = self;
|
||||||
|
seal_wire(&cipher, counter, &mut wire_buf, fsp_seal)?;
|
||||||
|
Ok(wire_buf)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Seal one job's wire buffer in place: the inner FSP seal first when the job
|
||||||
|
/// carries one, then the outer FMP seal over `[16..]` with the header as AAD.
|
||||||
|
/// Each tag is appended into capacity the builder reserved, so the buffer
|
||||||
|
/// becomes the wire packet without reallocating.
|
||||||
|
///
|
||||||
|
/// A layout that does not fit the buffer is refused rather than indexed out of
|
||||||
|
/// bounds: this runs on a worker thread and, for a job that worker refused, on
|
||||||
|
/// the rx loop, where a panic would end the node rather than one worker.
|
||||||
|
fn seal_wire(
|
||||||
|
cipher: &LessSafeKey,
|
||||||
|
counter: u64,
|
||||||
|
wire_buf: &mut Vec<u8>,
|
||||||
|
fsp_seal: Option<FspSealJob>,
|
||||||
|
) -> Result<(), SealError> {
|
||||||
|
if let Some(fsp) = fsp_seal {
|
||||||
|
let aad_end = fsp
|
||||||
|
.aad_offset
|
||||||
|
.checked_add(FSP_HEADER_SIZE)
|
||||||
|
.ok_or(SealError::Layout)?;
|
||||||
|
if aad_end > fsp.plaintext_offset || fsp.plaintext_offset > wire_buf.len() {
|
||||||
|
return Err(SealError::Layout);
|
||||||
|
}
|
||||||
|
|
||||||
|
let mut nonce_bytes = [0u8; 12];
|
||||||
|
nonce_bytes[4..12].copy_from_slice(&fsp.counter.to_le_bytes());
|
||||||
|
let nonce = Nonce::assume_unique_for_key(nonce_bytes);
|
||||||
|
let (prefix, plaintext_slice) = wire_buf.split_at_mut(fsp.plaintext_offset);
|
||||||
|
let aad = &prefix[fsp.aad_offset..aad_end];
|
||||||
|
let tag = fsp
|
||||||
|
.cipher
|
||||||
|
.seal_in_place_separate_tag(nonce, Aad::from(aad), plaintext_slice)
|
||||||
|
.map_err(|_| SealError::Aead)?;
|
||||||
|
wire_buf.extend_from_slice(tag.as_ref());
|
||||||
|
}
|
||||||
|
|
||||||
|
if wire_buf.len() < ESTABLISHED_HEADER_SIZE {
|
||||||
|
return Err(SealError::Layout);
|
||||||
|
}
|
||||||
|
let mut nonce_bytes = [0u8; 12];
|
||||||
|
nonce_bytes[4..12].copy_from_slice(&counter.to_le_bytes());
|
||||||
|
let nonce = Nonce::assume_unique_for_key(nonce_bytes);
|
||||||
|
// Split-borrow: AAD reads from header bytes [0..16], seal writes
|
||||||
|
// into the plaintext slice [16..]. ring::aead's `seal_in_place_
|
||||||
|
// separate_tag` takes `&mut [u8]` so we can hand it the
|
||||||
|
// post-header slice while AAD references the header slice.
|
||||||
|
// `split_at_mut` is the standard way to do this safely.
|
||||||
|
let (header_slice, plaintext_slice) = wire_buf.split_at_mut(ESTABLISHED_HEADER_SIZE);
|
||||||
|
let tag = cipher
|
||||||
|
.seal_in_place_separate_tag(nonce, Aad::from(&*header_slice), plaintext_slice)
|
||||||
|
.map_err(|_| SealError::Aead)?;
|
||||||
|
// wire_buf already has `+16` capacity reserved → no realloc.
|
||||||
|
wire_buf.extend_from_slice(tag.as_ref());
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
/// Encrypt every job in `batch` in place, then issue one or more
|
/// Encrypt every job in `batch` in place, then issue one or more
|
||||||
/// bulk-send syscalls grouped **by exact send target**. Clears
|
/// bulk-send syscalls grouped **by exact send target**. Clears
|
||||||
/// `batch` on return. Sync version — operates directly on the raw
|
/// `batch` on return. Sync version — operates directly on the raw
|
||||||
@@ -1227,65 +1359,14 @@ fn flush_batch_sync(
|
|||||||
crate::perf_profile::Stage::FmpWorkerQueueWait,
|
crate::perf_profile::Stage::FmpWorkerQueueWait,
|
||||||
queued_at,
|
queued_at,
|
||||||
);
|
);
|
||||||
if let Some(fsp) = fsp_seal {
|
if seal_wire(&cipher, counter, &mut wire_buf, fsp_seal).is_err() {
|
||||||
if fsp.aad_offset + FSP_HEADER_SIZE > fsp.plaintext_offset
|
#[cfg(target_os = "macos")]
|
||||||
|| fsp.plaintext_offset > wire_buf.len()
|
if let Some(ticket) = macos_ticket {
|
||||||
{
|
push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip);
|
||||||
#[cfg(target_os = "macos")]
|
|
||||||
if let Some(ticket) = macos_ticket {
|
|
||||||
push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip);
|
|
||||||
}
|
|
||||||
continue;
|
|
||||||
}
|
}
|
||||||
|
continue;
|
||||||
let mut nonce_bytes = [0u8; 12];
|
|
||||||
nonce_bytes[4..12].copy_from_slice(&fsp.counter.to_le_bytes());
|
|
||||||
let nonce = Nonce::assume_unique_for_key(nonce_bytes);
|
|
||||||
let (prefix, plaintext_slice) = wire_buf.split_at_mut(fsp.plaintext_offset);
|
|
||||||
let aad = &prefix[fsp.aad_offset..fsp.aad_offset + FSP_HEADER_SIZE];
|
|
||||||
let tag =
|
|
||||||
match fsp
|
|
||||||
.cipher
|
|
||||||
.seal_in_place_separate_tag(nonce, Aad::from(aad), plaintext_slice)
|
|
||||||
{
|
|
||||||
Ok(tag) => tag,
|
|
||||||
Err(_) => {
|
|
||||||
#[cfg(target_os = "macos")]
|
|
||||||
if let Some(ticket) = macos_ticket {
|
|
||||||
push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip);
|
|
||||||
}
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
};
|
|
||||||
wire_buf.extend_from_slice(tag.as_ref());
|
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut nonce_bytes = [0u8; 12];
|
|
||||||
nonce_bytes[4..12].copy_from_slice(&counter.to_le_bytes());
|
|
||||||
let nonce = Nonce::assume_unique_for_key(nonce_bytes);
|
|
||||||
// Split-borrow: AAD reads from header bytes [0..16], seal writes
|
|
||||||
// into the plaintext slice [16..]. ring::aead's `seal_in_place_
|
|
||||||
// separate_tag` takes `&mut [u8]` so we can hand it the
|
|
||||||
// post-header slice while AAD references the header slice.
|
|
||||||
// `split_at_mut` is the standard way to do this safely.
|
|
||||||
let (header_slice, plaintext_slice) = wire_buf.split_at_mut(ESTABLISHED_HEADER_SIZE);
|
|
||||||
let tag = match cipher.seal_in_place_separate_tag(
|
|
||||||
nonce,
|
|
||||||
Aad::from(&*header_slice),
|
|
||||||
plaintext_slice,
|
|
||||||
) {
|
|
||||||
Ok(tag) => tag,
|
|
||||||
Err(_) => {
|
|
||||||
#[cfg(target_os = "macos")]
|
|
||||||
if let Some(ticket) = macos_ticket {
|
|
||||||
push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip);
|
|
||||||
}
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
};
|
|
||||||
// wire_buf already has `+16` capacity reserved → no realloc.
|
|
||||||
wire_buf.extend_from_slice(tag.as_ref());
|
|
||||||
|
|
||||||
#[cfg(target_os = "macos")]
|
#[cfg(target_os = "macos")]
|
||||||
if let Some(ticket) = macos_ticket {
|
if let Some(ticket) = macos_ticket {
|
||||||
push_mac_completion(
|
push_mac_completion(
|
||||||
@@ -2202,6 +2283,150 @@ mod unix_tests {
|
|||||||
assert_eq!(recovered_fsp_plaintext, fsp_plaintext);
|
assert_eq!(recovered_fsp_plaintext, fsp_plaintext);
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The seal a refused job gets on the main loop produces a packet the
|
||||||
|
/// canonical receive-side decoders accept, inner FSP layer included.
|
||||||
|
#[test]
|
||||||
|
fn an_inline_seal_matches_the_worker_wire_layout() {
|
||||||
|
use crate::NodeAddr;
|
||||||
|
use crate::noise::TAG_SIZE;
|
||||||
|
use crate::proto::fmp::wire::{EncryptedHeader, build_established_header};
|
||||||
|
use crate::proto::fsp::wire::build_fsp_header;
|
||||||
|
use crate::proto::link::{
|
||||||
|
LinkMessageType, SESSION_DATAGRAM_HEADER_SIZE, SessionDatagramRef,
|
||||||
|
};
|
||||||
|
use crate::utils::index::SessionIndex;
|
||||||
|
|
||||||
|
let rt = tokio::runtime::Builder::new_current_thread()
|
||||||
|
.enable_io()
|
||||||
|
.build()
|
||||||
|
.expect("tokio rt");
|
||||||
|
let _enter = rt.enter();
|
||||||
|
let socket = UdpRawSocket::open("127.0.0.1:0".parse().unwrap(), 1 << 20, 1 << 20)
|
||||||
|
.expect("open send socket")
|
||||||
|
.into_async()
|
||||||
|
.expect("into_async");
|
||||||
|
|
||||||
|
let fmp_cipher = test_cipher(0x31);
|
||||||
|
let fsp_cipher = test_cipher(0x32);
|
||||||
|
let (fmp_counter, fsp_counter) = (900u64, 77u64);
|
||||||
|
let fsp_plaintext = b"sealed on the main loop".to_vec();
|
||||||
|
let link_plaintext_len =
|
||||||
|
SESSION_DATAGRAM_HEADER_SIZE + FSP_HEADER_SIZE + fsp_plaintext.len();
|
||||||
|
let fmp_inner_len = 4 + link_plaintext_len + TAG_SIZE;
|
||||||
|
let fsp_header = build_fsp_header(fsp_counter, 0, fsp_plaintext.len() as u16);
|
||||||
|
let fmp_header =
|
||||||
|
build_established_header(SessionIndex::new(5), fmp_counter, 0, fmp_inner_len as u16);
|
||||||
|
|
||||||
|
let mut wire_buf = Vec::with_capacity(ESTABLISHED_HEADER_SIZE + fmp_inner_len + TAG_SIZE);
|
||||||
|
wire_buf.extend_from_slice(&fmp_header);
|
||||||
|
wire_buf.extend_from_slice(&7u32.to_le_bytes());
|
||||||
|
wire_buf.push(LinkMessageType::SessionDatagram.to_byte());
|
||||||
|
wire_buf.push(16);
|
||||||
|
wire_buf.extend_from_slice(&1280u16.to_le_bytes());
|
||||||
|
wire_buf.extend_from_slice(NodeAddr::from_bytes([0xAA; 16]).as_bytes());
|
||||||
|
wire_buf.extend_from_slice(NodeAddr::from_bytes([0xBB; 16]).as_bytes());
|
||||||
|
let aad_offset = wire_buf.len();
|
||||||
|
wire_buf.extend_from_slice(&fsp_header);
|
||||||
|
let plaintext_offset = wire_buf.len();
|
||||||
|
wire_buf.extend_from_slice(&fsp_plaintext);
|
||||||
|
|
||||||
|
let job = FmpSendJob {
|
||||||
|
cipher: fmp_cipher.clone(),
|
||||||
|
counter: fmp_counter,
|
||||||
|
wire_buf,
|
||||||
|
fsp_seal: Some(FspSealJob {
|
||||||
|
cipher: fsp_cipher.clone(),
|
||||||
|
counter: fsp_counter,
|
||||||
|
aad_offset,
|
||||||
|
plaintext_offset,
|
||||||
|
}),
|
||||||
|
socket: socket.clone(),
|
||||||
|
dest_addr: "127.0.0.1:9".parse().unwrap(),
|
||||||
|
#[cfg(any(target_os = "linux", target_os = "macos"))]
|
||||||
|
connected_socket: None,
|
||||||
|
drop_on_backpressure: true,
|
||||||
|
queued_at: None,
|
||||||
|
};
|
||||||
|
let wire = job.seal_inline().expect("inline seal");
|
||||||
|
|
||||||
|
let parsed = EncryptedHeader::parse(&wire).expect("FMP header parses");
|
||||||
|
assert_eq!(parsed.counter, fmp_counter);
|
||||||
|
let fmp_plaintext = crate::noise::open(
|
||||||
|
Some(&fmp_cipher),
|
||||||
|
fmp_counter,
|
||||||
|
&parsed.header_bytes,
|
||||||
|
&wire[ESTABLISHED_HEADER_SIZE..],
|
||||||
|
)
|
||||||
|
.expect("FMP open");
|
||||||
|
let datagram = SessionDatagramRef::decode(&fmp_plaintext[5..]).expect("datagram decodes");
|
||||||
|
let inner = crate::noise::open(
|
||||||
|
Some(&fsp_cipher),
|
||||||
|
fsp_counter,
|
||||||
|
&datagram.payload[..FSP_HEADER_SIZE],
|
||||||
|
&datagram.payload[FSP_HEADER_SIZE..],
|
||||||
|
)
|
||||||
|
.expect("FSP open");
|
||||||
|
assert_eq!(inner, fsp_plaintext);
|
||||||
|
|
||||||
|
// FMP-only twin: a link message carries no inner seal.
|
||||||
|
let link_plaintext = b"\x10link message".to_vec();
|
||||||
|
let header = build_established_header(
|
||||||
|
SessionIndex::new(5),
|
||||||
|
fmp_counter + 1,
|
||||||
|
0,
|
||||||
|
(4 + link_plaintext.len()) as u16,
|
||||||
|
);
|
||||||
|
let mut wire_buf = Vec::with_capacity(ESTABLISHED_HEADER_SIZE + 4 + 64);
|
||||||
|
wire_buf.extend_from_slice(&header);
|
||||||
|
wire_buf.extend_from_slice(&9u32.to_le_bytes());
|
||||||
|
wire_buf.extend_from_slice(&link_plaintext);
|
||||||
|
let job = FmpSendJob {
|
||||||
|
cipher: fmp_cipher.clone(),
|
||||||
|
counter: fmp_counter + 1,
|
||||||
|
wire_buf,
|
||||||
|
fsp_seal: None,
|
||||||
|
socket,
|
||||||
|
dest_addr: "127.0.0.1:9".parse().unwrap(),
|
||||||
|
#[cfg(any(target_os = "linux", target_os = "macos"))]
|
||||||
|
connected_socket: None,
|
||||||
|
drop_on_backpressure: false,
|
||||||
|
queued_at: None,
|
||||||
|
};
|
||||||
|
let wire = job.seal_inline().expect("inline seal");
|
||||||
|
let opened = crate::noise::open(
|
||||||
|
Some(&fmp_cipher),
|
||||||
|
fmp_counter + 1,
|
||||||
|
&wire[..ESTABLISHED_HEADER_SIZE],
|
||||||
|
&wire[ESTABLISHED_HEADER_SIZE..],
|
||||||
|
)
|
||||||
|
.expect("FMP open");
|
||||||
|
assert_eq!(&opened[4..], &link_plaintext[..]);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A job whose offsets do not fit its buffer is refused, not indexed out
|
||||||
|
/// of bounds. On the main loop a panic here would end the node.
|
||||||
|
#[test]
|
||||||
|
fn a_job_whose_layout_does_not_fit_is_refused_not_a_panic() {
|
||||||
|
let cipher = test_cipher(9);
|
||||||
|
let mut short = vec![0u8; ESTABLISHED_HEADER_SIZE - 1];
|
||||||
|
assert!(matches!(
|
||||||
|
seal_wire(&cipher, 1, &mut short, None),
|
||||||
|
Err(SealError::Layout)
|
||||||
|
));
|
||||||
|
|
||||||
|
let mut buf = vec![0u8; 64];
|
||||||
|
let overflowing = FspSealJob {
|
||||||
|
cipher: test_cipher(8),
|
||||||
|
counter: 1,
|
||||||
|
aad_offset: usize::MAX - 2,
|
||||||
|
plaintext_offset: 40,
|
||||||
|
};
|
||||||
|
assert!(matches!(
|
||||||
|
seal_wire(&cipher, 1, &mut buf, Some(overflowing)),
|
||||||
|
Err(SealError::Layout)
|
||||||
|
));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Standalone tests for the GSO-eligibility predicate. The full
|
/// Standalone tests for the GSO-eligibility predicate. The full
|
||||||
@@ -2581,7 +2806,7 @@ mod mac_queue_tests {
|
|||||||
fn spawn_pusher<T: Send + 'static>(
|
fn spawn_pusher<T: Send + 'static>(
|
||||||
tx: MacWorkerSender<T>,
|
tx: MacWorkerSender<T>,
|
||||||
item: T,
|
item: T,
|
||||||
) -> mpsc::Receiver<Result<(), MacWorkerPushError>> {
|
) -> mpsc::Receiver<Result<(), MacWorkerPushError<T>>> {
|
||||||
let (done_tx, done_rx) = mpsc::channel();
|
let (done_tx, done_rx) = mpsc::channel();
|
||||||
thread::spawn(move || {
|
thread::spawn(move || {
|
||||||
let result = tx.push_blocking(item);
|
let result = tx.push_blocking(item);
|
||||||
@@ -2621,7 +2846,10 @@ mod mac_queue_tests {
|
|||||||
let result = done
|
let result = done
|
||||||
.recv_timeout(WAIT)
|
.recv_timeout(WAIT)
|
||||||
.expect("push_blocking still blocked after the worker thread died");
|
.expect("push_blocking still blocked after the worker thread died");
|
||||||
assert!(matches!(result, Err(MacWorkerPushError)));
|
match result {
|
||||||
|
Err(MacWorkerPushError(job)) => assert_eq!(*job, 3, "the refused job comes back"),
|
||||||
|
Ok(()) => panic!("push_blocking queued onto a dead worker"),
|
||||||
|
}
|
||||||
assert!(worker.join().is_err(), "worker thread should have panicked");
|
assert!(worker.join().is_err(), "worker thread should have panicked");
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2629,12 +2857,18 @@ mod mac_queue_tests {
|
|||||||
fn try_push_returns_closed_after_receiver_dropped() {
|
fn try_push_returns_closed_after_receiver_dropped() {
|
||||||
let (tx, rx) = mac_worker_channel::<u32>(2);
|
let (tx, rx) = mac_worker_channel::<u32>(2);
|
||||||
drop(rx);
|
drop(rx);
|
||||||
assert!(matches!(tx.try_push(1), Err(MacWorkerTryPushError::Closed)));
|
match tx.try_push(1) {
|
||||||
|
Err(MacWorkerTryPushError::Closed(job)) => assert_eq!(*job, 1),
|
||||||
|
_ => panic!("try_push on a closed queue should hand the job back"),
|
||||||
|
}
|
||||||
let done = spawn_pusher(tx, 2);
|
let done = spawn_pusher(tx, 2);
|
||||||
let result = done
|
let result = done
|
||||||
.recv_timeout(WAIT)
|
.recv_timeout(WAIT)
|
||||||
.expect("push_blocking blocked on a queue whose receiver is gone");
|
.expect("push_blocking blocked on a queue whose receiver is gone");
|
||||||
assert!(matches!(result, Err(MacWorkerPushError)));
|
match result {
|
||||||
|
Err(MacWorkerPushError(job)) => assert_eq!(*job, 2),
|
||||||
|
Ok(()) => panic!("push_blocking queued onto a closed queue"),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -2824,9 +3058,11 @@ mod mac_ordered_tests {
|
|||||||
assert!(tx.try_push(rig.sequenced(1)).is_ok());
|
assert!(tx.try_push(rig.sequenced(1)).is_ok());
|
||||||
assert!(tx.try_push(rig.sequenced(2)).is_ok());
|
assert!(tx.try_push(rig.sequenced(2)).is_ok());
|
||||||
drop(rx);
|
drop(rx);
|
||||||
|
// The refused job comes back and is dropped here, which releases its
|
||||||
|
// slot as the dispatcher's caller would.
|
||||||
assert!(matches!(
|
assert!(matches!(
|
||||||
tx.try_push(rig.sequenced(3)),
|
tx.try_push(rig.sequenced(3)),
|
||||||
Err(MacWorkerTryPushError::Closed)
|
Err(MacWorkerTryPushError::Closed(_))
|
||||||
));
|
));
|
||||||
let mut batch = vec![rig.sequenced(4)];
|
let mut batch = vec![rig.sequenced(4)];
|
||||||
flush_batch_sync(&mut batch).expect("flush");
|
flush_batch_sync(&mut batch).expect("flush");
|
||||||
@@ -3049,7 +3285,10 @@ mod pool_tests {
|
|||||||
assert_eq!(pool.liveness().live_workers(), 1);
|
assert_eq!(pool.liveness().live_workers(), 1);
|
||||||
|
|
||||||
let recv = receiver_on_worker(&pool, 0);
|
let recv = receiver_on_worker(&pool, 0);
|
||||||
pool.dispatch(rig.job(recv.local_addr().unwrap(), 1));
|
assert!(
|
||||||
|
pool.dispatch(rig.job(recv.local_addr().unwrap(), 1))
|
||||||
|
.is_ok()
|
||||||
|
);
|
||||||
let mut buf = [0u8; 128];
|
let mut buf = [0u8; 128];
|
||||||
recv.recv_from(&mut buf)
|
recv.recv_from(&mut buf)
|
||||||
.expect("the live worker did not send the job dispatched to it");
|
.expect("the live worker did not send the job dispatched to it");
|
||||||
@@ -3062,9 +3301,9 @@ mod pool_tests {
|
|||||||
let rig = Rig::new();
|
let rig = Rig::new();
|
||||||
let pool = EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]);
|
let pool = EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]);
|
||||||
let recv = receiver_on_worker(&pool, 1);
|
let recv = receiver_on_worker(&pool, 1);
|
||||||
let ((), logs) = crate::testutil::capture_logs(|| {
|
let (dispatched, logs) =
|
||||||
pool.dispatch(rig.job(recv.local_addr().unwrap(), 1));
|
crate::testutil::capture_logs(|| pool.dispatch(rig.job(recv.local_addr().unwrap(), 1)));
|
||||||
});
|
assert!(dispatched.is_err(), "a dead worker's job must come back");
|
||||||
let warnings = logs.warnings();
|
let warnings = logs.warnings();
|
||||||
assert_eq!(warnings.len(), 1, "{warnings:?}");
|
assert_eq!(warnings.len(), 1, "{warnings:?}");
|
||||||
assert!(warnings[0].contains(" worker=1"), "{warnings:?}");
|
assert!(warnings[0].contains(" worker=1"), "{warnings:?}");
|
||||||
@@ -3091,15 +3330,14 @@ mod pool_tests {
|
|||||||
let recv = UdpSocket::bind("127.0.0.1:0").expect("bind receiver");
|
let recv = UdpSocket::bind("127.0.0.1:0").expect("bind receiver");
|
||||||
let dest = recv.local_addr().unwrap();
|
let dest = recv.local_addr().unwrap();
|
||||||
for counter in 0..WORKER_CHANNEL_CAP as u64 {
|
for counter in 0..WORKER_CHANNEL_CAP as u64 {
|
||||||
pool.dispatch(rig.job(dest, counter));
|
assert!(pool.dispatch(rig.job(dest, counter)).is_ok());
|
||||||
}
|
}
|
||||||
|
|
||||||
let (done_tx, done_rx) = mpsc::channel::<()>();
|
let (done_tx, done_rx) = mpsc::channel::<bool>();
|
||||||
let blocked_pool = pool.clone();
|
let blocked_pool = pool.clone();
|
||||||
let last = rig.job(dest, WORKER_CHANNEL_CAP as u64);
|
let last = rig.job(dest, WORKER_CHANNEL_CAP as u64);
|
||||||
std::thread::spawn(move || {
|
std::thread::spawn(move || {
|
||||||
blocked_pool.dispatch(last);
|
let _ = done_tx.send(blocked_pool.dispatch(last).is_ok());
|
||||||
let _ = done_tx.send(());
|
|
||||||
});
|
});
|
||||||
std::thread::sleep(Duration::from_millis(200));
|
std::thread::sleep(Duration::from_millis(200));
|
||||||
assert!(
|
assert!(
|
||||||
@@ -3108,9 +3346,57 @@ mod pool_tests {
|
|||||||
);
|
);
|
||||||
|
|
||||||
release_tx.send(()).expect("worker gone before release");
|
release_tx.send(()).expect("worker gone before release");
|
||||||
done_rx
|
let queued = done_rx
|
||||||
.recv_timeout(Duration::from_secs(5))
|
.recv_timeout(Duration::from_secs(5))
|
||||||
.expect("dispatch still blocked after the worker drained");
|
.expect("dispatch still blocked after the worker drained");
|
||||||
|
assert!(queued, "the blocked job was refused, not queued");
|
||||||
assert_eq!(pool.liveness().refused_dispatches(), 0);
|
assert_eq!(pool.liveness().refused_dispatches(), 0);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A job refused by a dead worker comes back whole, with the counter and
|
||||||
|
/// buffer the caller reserved, whether the worker was already gone or
|
||||||
|
/// died while the dispatch waited on its full queue.
|
||||||
|
#[cfg(not(target_os = "macos"))]
|
||||||
|
#[test]
|
||||||
|
fn dispatch_to_an_exited_worker_hands_the_job_back() {
|
||||||
|
let rig = Rig::new();
|
||||||
|
let dest: SocketAddr = "127.0.0.1:9".parse().unwrap();
|
||||||
|
|
||||||
|
let pool = EncryptWorkerPool::for_test(vec![TestWorker::FailSpawn]);
|
||||||
|
let job = rig.job(dest, 41);
|
||||||
|
let wire_buf = job.wire_buf.clone();
|
||||||
|
let back = match pool.dispatch(job) {
|
||||||
|
Err(back) => back,
|
||||||
|
Ok(()) => panic!("a dead worker took the job"),
|
||||||
|
};
|
||||||
|
assert_eq!(back.counter, 41);
|
||||||
|
assert_eq!(back.wire_buf, wire_buf);
|
||||||
|
|
||||||
|
// Worker alive but not draining; it exits while a dispatch waits.
|
||||||
|
let (exit_tx, exit_rx) = mpsc::channel::<()>();
|
||||||
|
let pool = EncryptWorkerPool::for_test(vec![TestWorker::ExitOn(exit_rx)]);
|
||||||
|
for counter in 0..WORKER_CHANNEL_CAP as u64 {
|
||||||
|
assert!(pool.dispatch(rig.job(dest, counter)).is_ok());
|
||||||
|
}
|
||||||
|
let (done_tx, done_rx) = mpsc::channel::<Option<u64>>();
|
||||||
|
let blocked_pool = pool.clone();
|
||||||
|
let last = rig.job(dest, 7_000);
|
||||||
|
std::thread::spawn(move || {
|
||||||
|
let _ = done_tx.send(blocked_pool.dispatch(last).err().map(|job| job.counter));
|
||||||
|
});
|
||||||
|
std::thread::sleep(Duration::from_millis(100));
|
||||||
|
assert!(
|
||||||
|
matches!(done_rx.try_recv(), Err(mpsc::TryRecvError::Empty)),
|
||||||
|
"dispatch returned while the queue was full and the worker alive"
|
||||||
|
);
|
||||||
|
exit_tx.send(()).expect("worker gone before its signal");
|
||||||
|
let back = done_rx
|
||||||
|
.recv_timeout(Duration::from_secs(5))
|
||||||
|
.expect("dispatch still blocked after the worker exited");
|
||||||
|
assert_eq!(
|
||||||
|
back,
|
||||||
|
Some(7_000),
|
||||||
|
"the job blocked on a dying worker must come back"
|
||||||
|
);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2875,7 +2875,7 @@ impl Node {
|
|||||||
entry.touch(send.now_ms);
|
entry.touch(send.now_ms);
|
||||||
}
|
}
|
||||||
|
|
||||||
workers.dispatch(crate::node::encrypt_worker::FmpSendJob {
|
let dispatched = workers.dispatch(crate::node::encrypt_worker::FmpSendJob {
|
||||||
cipher: fmp_cipher,
|
cipher: fmp_cipher,
|
||||||
counter: fmp_counter,
|
counter: fmp_counter,
|
||||||
wire_buf,
|
wire_buf,
|
||||||
@@ -2895,9 +2895,44 @@ impl Node {
|
|||||||
drop_on_backpressure: true,
|
drop_on_backpressure: true,
|
||||||
queued_at: None,
|
queued_at: None,
|
||||||
});
|
});
|
||||||
|
if let Err(job) = dispatched {
|
||||||
|
self.send_refused_job_inline(*job, transport_id, &remote_addr, next_hop_addr)
|
||||||
|
.await;
|
||||||
|
}
|
||||||
Ok(true)
|
Ok(true)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Seal and send a job the encrypt worker for its next hop refused
|
||||||
|
/// because that worker has exited, using the FSP and FMP counters the
|
||||||
|
/// job already reserved so neither counter is skipped. Stats were
|
||||||
|
/// recorded before dispatch and now describe this packet.
|
||||||
|
///
|
||||||
|
/// A failure is logged and swallowed, as the worker does with its own:
|
||||||
|
/// the caller sees the same result whichever of the two sent the packet.
|
||||||
|
#[cfg(unix)]
|
||||||
|
async fn send_refused_job_inline(
|
||||||
|
&self,
|
||||||
|
job: crate::node::encrypt_worker::FmpSendJob,
|
||||||
|
transport_id: crate::transport::TransportId,
|
||||||
|
remote_addr: &crate::transport::TransportAddr,
|
||||||
|
next_hop_addr: NodeAddr,
|
||||||
|
) {
|
||||||
|
let wire = match job.seal_inline() {
|
||||||
|
Ok(wire) => wire,
|
||||||
|
Err(error) => {
|
||||||
|
debug!(next_hop = %next_hop_addr, %error, "Inline seal of session data failed");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let Some(transport) = self.transports.get(&transport_id) else {
|
||||||
|
debug!(next_hop = %next_hop_addr, "Transport gone before inline send of session data");
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
if let Err(error) = transport.send(remote_addr, &wire).await {
|
||||||
|
debug!(next_hop = %next_hop_addr, %error, "Inline send of session data failed");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Send an IPv6 packet through the IPv6 shim (port 256) with header compression.
|
/// Send an IPv6 packet through the IPv6 shim (port 256) with header compression.
|
||||||
///
|
///
|
||||||
/// Compresses the IPv6 header (format 0x00), then sends via `send_session_data`
|
/// Compresses the IPv6 header (format 0x00), then sends via `send_session_data`
|
||||||
|
|||||||
@@ -3,7 +3,8 @@
|
|||||||
//! exited.
|
//! exited.
|
||||||
//!
|
//!
|
||||||
//! The pools are a performance offload, and Windows never starts them at
|
//! The pools are a performance offload, and Windows never starts them at
|
||||||
//! all, so losing workers is `Degraded` at most and never fatal. Inbound
|
//! all, so losing workers is `Degraded` at most and never fatal. An outbound
|
||||||
|
//! packet for a missing encrypt worker is sealed on the main loop. Inbound
|
||||||
//! packets for a session already held by a missing decrypt worker are
|
//! packets for a session already held by a missing decrypt worker are
|
||||||
//! dropped until the session rekeys or the link is re-established; a session
|
//! dropped until the session rekeys or the link is re-established; a session
|
||||||
//! that would be registered on the missing worker after the loss is decrypted
|
//! that would be registered on the missing worker after the loss is decrypted
|
||||||
|
|||||||
+44
-23
@@ -3914,7 +3914,7 @@ impl Node {
|
|||||||
// Drop bulk endpoint data on UDP backpressure to
|
// Drop bulk endpoint data on UDP backpressure to
|
||||||
// keep the queue moving; control frames retry.
|
// keep the queue moving; control frames retry.
|
||||||
let drop_on_backpressure = plaintext.first().is_some_and(|t| *t == 0x00);
|
let drop_on_backpressure = plaintext.first().is_some_and(|t| *t == 0x00);
|
||||||
workers.dispatch(crate::node::encrypt_worker::FmpSendJob {
|
let dispatched = workers.dispatch(crate::node::encrypt_worker::FmpSendJob {
|
||||||
cipher: fmp_cipher,
|
cipher: fmp_cipher,
|
||||||
counter,
|
counter,
|
||||||
wire_buf,
|
wire_buf,
|
||||||
@@ -3926,12 +3926,27 @@ impl Node {
|
|||||||
drop_on_backpressure,
|
drop_on_backpressure,
|
||||||
queued_at: None,
|
queued_at: None,
|
||||||
});
|
});
|
||||||
|
let sent_bytes = match dispatched {
|
||||||
|
Ok(()) => predicted_bytes,
|
||||||
|
// The worker for this destination has exited. Seal
|
||||||
|
// here with the counter already reserved, so no
|
||||||
|
// counter is skipped, and send as the inline path does.
|
||||||
|
Err(job) => {
|
||||||
|
let wire = job.seal_inline().map_err(|e| NodeError::SendFailed {
|
||||||
|
node_addr: *node_addr,
|
||||||
|
reason: format!("encryption failed: {}", e),
|
||||||
|
})?;
|
||||||
|
transport
|
||||||
|
.send(&remote_addr, &wire)
|
||||||
|
.await
|
||||||
|
.map_err(|e| link_send_error(*node_addr, e))?
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
if let Some(peer) = self.peers.get_mut(node_addr) {
|
if let Some(peer) = self.peers.get_mut(node_addr) {
|
||||||
peer.link_stats_mut().record_sent(predicted_bytes);
|
peer.link_stats_mut().record_sent(sent_bytes);
|
||||||
if let Some(mmp) = peer.mmp_mut() {
|
if let Some(mmp) = peer.mmp_mut() {
|
||||||
mmp.sender
|
mmp.sender.record_sent(counter, timestamp_ms, sent_bytes);
|
||||||
.record_sent(counter, timestamp_ms, predicted_bytes);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return Ok(());
|
return Ok(());
|
||||||
@@ -3989,25 +4004,7 @@ impl Node {
|
|||||||
let bytes_sent = transport
|
let bytes_sent = transport
|
||||||
.send(&remote_addr, &wire_packet)
|
.send(&remote_addr, &wire_packet)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| match e {
|
.map_err(|e| link_send_error(*node_addr, e))?;
|
||||||
TransportError::MtuExceeded { packet_size, mtu } => NodeError::MtuExceeded {
|
|
||||||
node_addr: *node_addr,
|
|
||||||
packet_size,
|
|
||||||
mtu,
|
|
||||||
},
|
|
||||||
// Preserve the transport's own classification instead of
|
|
||||||
// flattening every non-MTU failure into one string. A caller
|
|
||||||
// that wants to keep its half-built state across an interface
|
|
||||||
// flap can only do that if the distinction survives to it.
|
|
||||||
other if other.is_transient() => NodeError::SendUnavailable {
|
|
||||||
node_addr: *node_addr,
|
|
||||||
reason: format!("transport send: {}", other),
|
|
||||||
},
|
|
||||||
other => NodeError::SendFailed {
|
|
||||||
node_addr: *node_addr,
|
|
||||||
reason: format!("transport send: {}", other),
|
|
||||||
},
|
|
||||||
})?;
|
|
||||||
|
|
||||||
// Update send statistics
|
// Update send statistics
|
||||||
if let Some(peer) = self.peers.get_mut(node_addr) {
|
if let Some(peer) = self.peers.get_mut(node_addr) {
|
||||||
@@ -4022,6 +4019,30 @@ impl Node {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Map a transport's refusal of an encrypted link frame to `node_addr` onto
|
||||||
|
/// the error the link-send path reports.
|
||||||
|
fn link_send_error(node_addr: NodeAddr, e: TransportError) -> NodeError {
|
||||||
|
match e {
|
||||||
|
TransportError::MtuExceeded { packet_size, mtu } => NodeError::MtuExceeded {
|
||||||
|
node_addr,
|
||||||
|
packet_size,
|
||||||
|
mtu,
|
||||||
|
},
|
||||||
|
// Preserve the transport's own classification instead of
|
||||||
|
// flattening every non-MTU failure into one string. A caller
|
||||||
|
// that wants to keep its half-built state across an interface
|
||||||
|
// flap can only do that if the distinction survives to it.
|
||||||
|
other if other.is_transient() => NodeError::SendUnavailable {
|
||||||
|
node_addr,
|
||||||
|
reason: format!("transport send: {}", other),
|
||||||
|
},
|
||||||
|
other => NodeError::SendFailed {
|
||||||
|
node_addr,
|
||||||
|
reason: format!("transport send: {}", other),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Shell-side [`routing::RoutingView`] seam over live `Node` state — the sole
|
/// Shell-side [`routing::RoutingView`] seam over live `Node` state — the sole
|
||||||
/// routing read adapter the shell retains. It hands the sans-IO routing core
|
/// routing read adapter the shell retains. It hands the sans-IO routing core
|
||||||
/// borrowed peers plus raw `may_reach` / `link_cost` / `coords`
|
/// borrowed peers plus raw `may_reach` / `link_cost` / `coords`
|
||||||
|
|||||||
@@ -2377,3 +2377,209 @@ async fn a_transient_msg2_failure_on_the_restart_path_leaves_the_fresh_leg_pendi
|
|||||||
"a local interface flap must not be recorded as the peer's misbehaviour"
|
"a local interface flap must not be recorded as the peer's misbehaviour"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Link messages to a peer whose encrypt worker has exited are still sent,
|
||||||
|
/// sealed on the main loop with the counter the worker path reserved.
|
||||||
|
#[cfg(unix)]
|
||||||
|
mod dead_encrypt_worker {
|
||||||
|
use super::*;
|
||||||
|
use crate::config::UdpConfig;
|
||||||
|
use crate::node::encrypt_worker::EncryptWorkerPool;
|
||||||
|
use crate::node::tests::pool_with_dead_worker_for;
|
||||||
|
use crate::node::worker_set::TestWorker;
|
||||||
|
use crate::proto::fmp::wire::{EncryptedHeader, build_msg1};
|
||||||
|
use crate::transport::ReceivedPacket;
|
||||||
|
use crate::transport::udp::UdpTransport;
|
||||||
|
use tokio::time::{Duration, timeout};
|
||||||
|
|
||||||
|
/// Two nodes on real UDP sockets with an established link from A to B.
|
||||||
|
struct Pair {
|
||||||
|
node_a: Node,
|
||||||
|
node_b: Node,
|
||||||
|
packet_rx_b: crate::transport::PacketRx,
|
||||||
|
peer_a: NodeAddr,
|
||||||
|
peer_b: NodeAddr,
|
||||||
|
addr_b: std::net::SocketAddr,
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn linked_pair() -> Pair {
|
||||||
|
let mut node_a = make_node();
|
||||||
|
let mut node_b = make_node();
|
||||||
|
let tid = TransportId::new(1);
|
||||||
|
let udp_config = UdpConfig {
|
||||||
|
bind_addr: Some("127.0.0.1:0".to_string()),
|
||||||
|
mtu: Some(1280),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let (packet_tx_a, mut packet_rx_a) = packet_channel(64);
|
||||||
|
let (packet_tx_b, mut packet_rx_b) = packet_channel(64);
|
||||||
|
let mut transport_a = UdpTransport::new(tid, None, udp_config.clone(), packet_tx_a);
|
||||||
|
let mut transport_b = UdpTransport::new(tid, None, udp_config, packet_tx_b);
|
||||||
|
transport_a.start_async().await.unwrap();
|
||||||
|
transport_b.start_async().await.unwrap();
|
||||||
|
let addr_b = transport_b.local_addr().unwrap();
|
||||||
|
let remote_addr_b = TransportAddr::from_string(&addr_b.to_string());
|
||||||
|
node_a
|
||||||
|
.transports
|
||||||
|
.insert(tid, TransportHandle::Udp(transport_a));
|
||||||
|
node_b
|
||||||
|
.transports
|
||||||
|
.insert(tid, TransportHandle::Udp(transport_b));
|
||||||
|
|
||||||
|
let peer_b_identity = PeerIdentity::from_pubkey_full(node_b.identity().pubkey_full());
|
||||||
|
let peer_b = *peer_b_identity.node_addr();
|
||||||
|
let peer_a = *PeerIdentity::from_pubkey_full(node_a.identity().pubkey_full()).node_addr();
|
||||||
|
let link_id = node_a.allocate_link_id();
|
||||||
|
let our_index = node_a.index_allocator.allocate().unwrap();
|
||||||
|
node_a
|
||||||
|
.seed_handshake_machine(
|
||||||
|
HandshakeSeed::outbound(link_id, peer_b_identity, 1000)
|
||||||
|
.with_our_index(our_index)
|
||||||
|
.with_transport_id(tid)
|
||||||
|
.with_source_addr(remote_addr_b.clone()),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
let keypair = node_a.identity().keypair();
|
||||||
|
let epoch = node_a.startup_epoch();
|
||||||
|
let msg1 = node_a
|
||||||
|
.peer_machines
|
||||||
|
.get_mut(&link_id)
|
||||||
|
.unwrap()
|
||||||
|
.start_handshake(keypair, epoch, 1000)
|
||||||
|
.unwrap();
|
||||||
|
node_a.links.insert(
|
||||||
|
link_id,
|
||||||
|
Link::connectionless(
|
||||||
|
link_id,
|
||||||
|
tid,
|
||||||
|
remote_addr_b.clone(),
|
||||||
|
LinkDirection::Outbound,
|
||||||
|
Duration::from_millis(100),
|
||||||
|
),
|
||||||
|
);
|
||||||
|
node_a
|
||||||
|
.pending_outbound
|
||||||
|
.insert((tid, our_index.as_u32()), link_id);
|
||||||
|
node_a
|
||||||
|
.transports
|
||||||
|
.get(&tid)
|
||||||
|
.unwrap()
|
||||||
|
.send(&remote_addr_b, &build_msg1(our_index, &msg1))
|
||||||
|
.await
|
||||||
|
.expect("send msg1");
|
||||||
|
|
||||||
|
let msg1 = next_packet(&mut packet_rx_b).await;
|
||||||
|
node_b.handle_msg1(msg1).await;
|
||||||
|
let msg2 = next_packet(&mut packet_rx_a).await;
|
||||||
|
node_a.handle_msg2(msg2).await;
|
||||||
|
assert!(node_a.get_peer(&peer_b).is_some(), "A promoted B");
|
||||||
|
assert!(node_b.get_peer(&peer_a).is_some(), "B promoted A");
|
||||||
|
|
||||||
|
Pair {
|
||||||
|
node_a,
|
||||||
|
node_b,
|
||||||
|
packet_rx_b,
|
||||||
|
peer_a,
|
||||||
|
peer_b,
|
||||||
|
addr_b,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn next_packet(rx: &mut crate::transport::PacketRx) -> ReceivedPacket {
|
||||||
|
timeout(Duration::from_secs(5), rx.recv())
|
||||||
|
.await
|
||||||
|
.expect("no packet within the bound")
|
||||||
|
.expect("packet channel closed")
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Send one link message from A to B through `pool` and have B process
|
||||||
|
/// it. Returns (the frame's FMP counter, A's send counter before the
|
||||||
|
/// send, A's packets-sent delta, B's packets-received delta, encrypt
|
||||||
|
/// WARN lines).
|
||||||
|
async fn send_through(
|
||||||
|
pair: &mut Pair,
|
||||||
|
pool: &EncryptWorkerPool,
|
||||||
|
) -> (u64, u64, u64, u64, Vec<String>) {
|
||||||
|
// Let B take whatever A sent on promotion before measuring.
|
||||||
|
while let Ok(Some(packet)) =
|
||||||
|
timeout(Duration::from_millis(100), pair.packet_rx_b.recv()).await
|
||||||
|
{
|
||||||
|
pair.node_b.handle_encrypted_frame(packet).await;
|
||||||
|
}
|
||||||
|
pair.node_a.supervisor.encrypt_workers = Some(pool.clone());
|
||||||
|
let peer = pair.node_a.get_peer(&pair.peer_b).unwrap();
|
||||||
|
let counter_before = peer.noise_session().unwrap().current_send_counter();
|
||||||
|
let sent_before = peer.link_stats().packets_sent;
|
||||||
|
let recv_before = pair
|
||||||
|
.node_b
|
||||||
|
.get_peer(&pair.peer_a)
|
||||||
|
.unwrap()
|
||||||
|
.link_stats()
|
||||||
|
.packets_recv;
|
||||||
|
|
||||||
|
let (logs, guard) = crate::testutil::capture_logs_scoped();
|
||||||
|
pair.node_a
|
||||||
|
.send_encrypted_link_message(&pair.peer_b, b"\x10dead worker test")
|
||||||
|
.await
|
||||||
|
.expect("link message send");
|
||||||
|
drop(guard);
|
||||||
|
|
||||||
|
let packet = next_packet(&mut pair.packet_rx_b).await;
|
||||||
|
let frame_counter = EncryptedHeader::parse(&packet.data)
|
||||||
|
.expect("an established frame")
|
||||||
|
.counter;
|
||||||
|
pair.node_b.handle_encrypted_frame(packet).await;
|
||||||
|
|
||||||
|
let sent = pair
|
||||||
|
.node_a
|
||||||
|
.get_peer(&pair.peer_b)
|
||||||
|
.unwrap()
|
||||||
|
.link_stats()
|
||||||
|
.packets_sent
|
||||||
|
- sent_before;
|
||||||
|
let recv = pair
|
||||||
|
.node_b
|
||||||
|
.get_peer(&pair.peer_a)
|
||||||
|
.unwrap()
|
||||||
|
.link_stats()
|
||||||
|
.packets_recv
|
||||||
|
- recv_before;
|
||||||
|
let warnings = logs
|
||||||
|
.warnings()
|
||||||
|
.into_iter()
|
||||||
|
.filter(|line| line.contains("pool=\"encrypt\""))
|
||||||
|
.collect();
|
||||||
|
(frame_counter, counter_before, sent, recv, warnings)
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn a_link_message_for_a_dead_workers_peer_is_still_sent() {
|
||||||
|
let mut pair = linked_pair().await;
|
||||||
|
let pool = pool_with_dead_worker_for(pair.addr_b);
|
||||||
|
let (frame_counter, counter_before, sent, recv, warnings) =
|
||||||
|
send_through(&mut pair, &pool).await;
|
||||||
|
|
||||||
|
assert_eq!(recv, 1, "B did not authenticate the frame");
|
||||||
|
assert_eq!(
|
||||||
|
frame_counter, counter_before,
|
||||||
|
"the frame must carry the counter reserved for it, not a fresh one"
|
||||||
|
);
|
||||||
|
assert_eq!(sent, 1, "A counted the packet other than once");
|
||||||
|
assert_eq!(pool.liveness().refused_dispatches(), 1);
|
||||||
|
assert_eq!(warnings.len(), 1, "{warnings:?}");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn a_link_message_through_live_workers_is_sent_by_the_worker() {
|
||||||
|
let mut pair = linked_pair().await;
|
||||||
|
let pool = EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]);
|
||||||
|
let (frame_counter, counter_before, sent, recv, warnings) =
|
||||||
|
send_through(&mut pair, &pool).await;
|
||||||
|
|
||||||
|
assert_eq!(recv, 1, "B did not authenticate the frame");
|
||||||
|
assert_eq!(frame_counter, counter_before);
|
||||||
|
assert_eq!(sent, 1);
|
||||||
|
assert_eq!(pool.liveness().refused_dispatches(), 0);
|
||||||
|
assert!(warnings.is_empty(), "{warnings:?}");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -94,6 +94,31 @@ pub(super) fn install_connected_udp(
|
|||||||
.set_connected_udp(socket, drain);
|
.set_connected_udp(socket, drain);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// An encrypt pool of two whose worker for `dest` has exited. Where the worker
|
||||||
|
/// for a destination is not a function of the address alone (macOS), both
|
||||||
|
/// have.
|
||||||
|
#[cfg(unix)]
|
||||||
|
pub(super) fn pool_with_dead_worker_for(
|
||||||
|
dest: std::net::SocketAddr,
|
||||||
|
) -> crate::node::encrypt_worker::EncryptWorkerPool {
|
||||||
|
use crate::node::encrypt_worker::EncryptWorkerPool;
|
||||||
|
use crate::node::worker_set::TestWorker;
|
||||||
|
|
||||||
|
let dead = EncryptWorkerPool::for_test(vec![TestWorker::FailSpawn, TestWorker::FailSpawn])
|
||||||
|
.worker_index_for_dest(dest);
|
||||||
|
EncryptWorkerPool::for_test(
|
||||||
|
(0..2)
|
||||||
|
.map(|idx| {
|
||||||
|
if dead.is_none_or(|d| d == idx) {
|
||||||
|
TestWorker::FailSpawn
|
||||||
|
} else {
|
||||||
|
TestWorker::Run
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.collect(),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
/// Build a test node with an explicit `max_peers` limit (replaces the removed
|
/// Build a test node with an explicit `max_peers` limit (replaces the removed
|
||||||
/// `set_max_peers` setter; resource limits are immutable post-construction).
|
/// `set_max_peers` setter; resource limits are immutable post-construction).
|
||||||
pub(super) fn make_node_with_max_peers(max_peers: usize) -> Node {
|
pub(super) fn make_node_with_max_peers(max_peers: usize) -> Node {
|
||||||
|
|||||||
@@ -9016,3 +9016,110 @@ async fn a_path_broken_flood_releases_the_stored_path_mtu_only_once_per_interval
|
|||||||
"a second release for the same destination inside the interval is refused"
|
"a second release for the same destination inside the interval is refused"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Session data whose next hop's encrypt worker has exited is still
|
||||||
|
/// delivered, sealed on the main loop with the FSP and FMP counters the
|
||||||
|
/// worker path reserved.
|
||||||
|
#[cfg(unix)]
|
||||||
|
mod dead_encrypt_worker {
|
||||||
|
use super::*;
|
||||||
|
use crate::config::UdpConfig;
|
||||||
|
use crate::node::tests::pool_with_dead_worker_for;
|
||||||
|
use crate::transport::udp::UdpTransport;
|
||||||
|
use crate::transport::{TransportAddr, TransportHandle, TransportId, packet_channel};
|
||||||
|
|
||||||
|
/// A test node on a real UDP socket: the worker path needs one.
|
||||||
|
async fn make_test_node_udp() -> TestNode {
|
||||||
|
let mut node = make_node();
|
||||||
|
let transport_id = TransportId::new(1);
|
||||||
|
let config = UdpConfig {
|
||||||
|
bind_addr: Some("127.0.0.1:0".to_string()),
|
||||||
|
mtu: Some(1280),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let (packet_tx, packet_rx) = packet_channel(256);
|
||||||
|
let mut transport = UdpTransport::new(transport_id, None, config, packet_tx);
|
||||||
|
transport.start_async().await.unwrap();
|
||||||
|
let addr = TransportAddr::from_string(&transport.local_addr().unwrap().to_string());
|
||||||
|
node.transports
|
||||||
|
.insert(transport_id, TransportHandle::Udp(transport));
|
||||||
|
TestNode {
|
||||||
|
node,
|
||||||
|
transport_id,
|
||||||
|
packet_rx: crate::node::tests::spanning_tree::bridge_to_unbounded(packet_rx),
|
||||||
|
addr,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn fsp_send_counter(node: &Node, dest: &NodeAddr) -> u64 {
|
||||||
|
match node.get_session(dest).expect("session").state() {
|
||||||
|
EndToEndState::Established(session) => session.current_send_counter(),
|
||||||
|
_ => panic!("session not established"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn fmp_send_counter(node: &Node, peer: &NodeAddr) -> u64 {
|
||||||
|
node.get_peer(peer)
|
||||||
|
.expect("peer")
|
||||||
|
.noise_session()
|
||||||
|
.expect("link session")
|
||||||
|
.current_send_counter()
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn session_data_for_a_dead_workers_next_hop_is_still_delivered() {
|
||||||
|
let mut nodes = vec![make_test_node_udp().await, make_test_node_udp().await];
|
||||||
|
initiate_handshake(&mut nodes, 0, 1).await;
|
||||||
|
drain_all_packets(&mut nodes, false).await;
|
||||||
|
verify_tree_convergence(&nodes);
|
||||||
|
populate_all_coord_caches(&mut nodes);
|
||||||
|
establish_pair_session(&mut nodes).await;
|
||||||
|
drain_all_packets(&mut nodes, false).await;
|
||||||
|
|
||||||
|
let node0 = *nodes[0].node.node_addr();
|
||||||
|
let node1 = *nodes[1].node.node_addr();
|
||||||
|
let dest: std::net::SocketAddr = nodes[1].addr.to_string().parse().unwrap();
|
||||||
|
let pool = pool_with_dead_worker_for(dest);
|
||||||
|
nodes[0].node.supervisor.encrypt_workers = Some(pool.clone());
|
||||||
|
|
||||||
|
let fsp_before = fsp_send_counter(&nodes[0].node, &node1);
|
||||||
|
let fmp_before = fmp_send_counter(&nodes[0].node, &node1);
|
||||||
|
let recv_before = nodes[1]
|
||||||
|
.node
|
||||||
|
.get_session(&node0)
|
||||||
|
.unwrap()
|
||||||
|
.traffic_counters()
|
||||||
|
.1;
|
||||||
|
|
||||||
|
nodes[0]
|
||||||
|
.node
|
||||||
|
.send_session_data(&node1, 0, 0, b"for a dead worker")
|
||||||
|
.await
|
||||||
|
.expect("send_session_data");
|
||||||
|
// Read before anything else runs on A: one packet, one counter each.
|
||||||
|
assert_eq!(fsp_send_counter(&nodes[0].node, &node1), fsp_before + 1);
|
||||||
|
assert_eq!(fmp_send_counter(&nodes[0].node, &node1), fmp_before + 1);
|
||||||
|
assert_eq!(pool.liveness().refused_dispatches(), 1);
|
||||||
|
|
||||||
|
let delivered = |nodes: &[TestNode]| {
|
||||||
|
nodes[1]
|
||||||
|
.node
|
||||||
|
.get_session(&node0)
|
||||||
|
.unwrap()
|
||||||
|
.traffic_counters()
|
||||||
|
.1
|
||||||
|
};
|
||||||
|
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
|
||||||
|
while delivered(&nodes) == recv_before && tokio::time::Instant::now() < deadline {
|
||||||
|
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||||
|
process_available_packets(&mut nodes).await;
|
||||||
|
}
|
||||||
|
assert_eq!(
|
||||||
|
delivered(&nodes),
|
||||||
|
recv_before + 1,
|
||||||
|
"the payload never reached the destination"
|
||||||
|
);
|
||||||
|
|
||||||
|
cleanup_nodes(&mut nodes).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user