mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-05 11:08:25 +00:00
Keep the macOS ordered sender moving when a worker or sender thread dies
With FIPS_MACOS_ORDERED_SENDER set, each packet reserves a slot in its destination's send order and that destination's sender thread sends the slots strictly in turn. A slot that was never completed stopped the sender for good: completions for later slots then filled the flow, the encrypt workers blocked waiting for room, their queues filled, and the rx loop blocked behind them. Two ways led there: an encrypt worker thread that exited dropped every job it held without completing its slot, and a sender thread that exited left its flow open, so completions kept waiting for room it would never make, and later packets for that destination were still handed to it. Each reserved slot is now owned by a value that completes it as a skip when dropped unsent, whether it is dropped in a worker's batch, in its queue, in an undelivered completion group, or on the rx loop when a worker's queue has closed. That skip never waits for room, so the rx loop cannot block on it. A sender thread that leaves its loop, a panic included, now closes its flow; completions waiting for room return when their flow closes; and a closed flow is replaced the next time its destination is used and pruned from the flow map. The flow's lock is taken poison-tolerantly wherever another thread reads it, and a closed worker queue drops a refused job only after releasing its own lock. After an encrypt worker has died, the packets round-robined to it are still lost for the rest of the process, but the loss no longer holds up every later packet to the same destination. The default macOS sender is unchanged.
This commit is contained in:
+360
-46
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user