Skip to content

Commit a954228

Browse files
committed
fix flow management so that FlowMap is a ring buffer and reuses slots
1 parent 9443fca commit a954228

4 files changed

Lines changed: 66 additions & 28 deletions

File tree

lading/src/blackhole/tcp_crr.rs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,7 @@ use std::num::{NonZeroU16, NonZeroUsize};
2424
use serde::{Deserialize, Serialize};
2525

2626
use super::General;
27-
use crate::neper::rr::{self, ServerParams};
27+
use crate::neper::rr::{self, Mode, ServerParams};
2828

2929
fn default_nonzero_u16() -> NonZeroU16 {
3030
NonZeroU16::new(1).expect("1 is nonzero")
@@ -136,6 +136,7 @@ impl TcpCrr {
136136
response_size: self.config.response_size.get(),
137137
no_delay: self.config.no_delay,
138138
backlog: self.config.backlog,
139+
mode: Mode::Crr,
139140
};
140141
rr::run_server(params, self.metric_labels, self.shutdown, "tcp_crr").await?;
141142
Ok(())

lading/src/blackhole/tcp_rr.rs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@ use std::num::{NonZeroU16, NonZeroUsize};
2222
use serde::{Deserialize, Serialize};
2323

2424
use super::General;
25-
use crate::neper::rr::{self, ServerParams};
25+
use crate::neper::rr::{self, Mode, ServerParams};
2626

2727
fn default_nonzero_u16() -> NonZeroU16 {
2828
NonZeroU16::new(1).expect("1 is nonzero")
@@ -133,6 +133,7 @@ impl TcpRr {
133133
response_size: self.config.response_size.get(),
134134
no_delay: self.config.no_delay,
135135
backlog: self.config.backlog,
136+
mode: Mode::Rr,
136137
};
137138
rr::run_server(params, self.metric_labels, self.shutdown, "tcp_rr").await?;
138139
Ok(())

lading/src/neper/flow.rs

Lines changed: 26 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -37,28 +37,45 @@ pub(crate) struct FlowMap<S> {
3737
inner: Vec<Option<Flow<S>>>,
3838
}
3939

40+
/// Errors produced by `FlowMap`.
41+
#[derive(thiserror::Error, Debug)]
42+
pub(crate) enum FlowMapError {
43+
/// No capacity
44+
#[error("Server flow map is at capacity: {0}")]
45+
NoCapacity(usize),
46+
}
47+
4048
impl<S> FlowMap<S> {
41-
pub(crate) fn new() -> Self {
42-
Self { inner: Vec::new() }
49+
pub(crate) fn new(flows: usize) -> Self {
50+
Self {
51+
inner: Vec::with_capacity(flows * 2),
52+
}
4353
}
4454

45-
/// Insert a flow. Grows the backing vec if needed.
46-
pub(crate) fn insert(&mut self, flow: Flow<S>) {
47-
let idx = flow.token.0;
48-
if idx >= self.inner.len() {
55+
/// Insert a flow.
56+
pub(crate) fn insert(&mut self, flow: Flow<S>) -> Result<(), FlowMapError> {
57+
let idx = flow.token.0 % self.inner.capacity();
58+
if self.inner.len() < idx {
4959
self.inner.resize_with(idx + 1, || None);
5060
}
51-
self.inner[idx] = Some(flow);
61+
if self.inner[idx].is_none() {
62+
self.inner[idx] = Some(flow);
63+
return Ok(());
64+
}
65+
66+
Err(FlowMapError::NoCapacity(self.inner.capacity()))
5267
}
5368

5469
/// Get a mutable reference to the flow at the given token.
5570
pub(crate) fn get_mut(&mut self, token: Token) -> Option<&mut Flow<S>> {
56-
self.inner.get_mut(token.0).and_then(|slot| slot.as_mut())
71+
let idx = token.0 % self.inner.capacity();
72+
self.inner.get_mut(idx).and_then(|slot| slot.as_mut())
5773
}
5874

5975
/// Remove and return the flow at the given token.
6076
pub(crate) fn remove(&mut self, token: Token) -> Option<Flow<S>> {
61-
self.inner.get_mut(token.0).and_then(Option::take)
77+
let idx = token.0 % self.inner.capacity();
78+
self.inner.get_mut(idx).and_then(Option::take)
6279
}
6380
}
6481

lading/src/neper/rr.rs

Lines changed: 36 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -95,6 +95,8 @@ pub(crate) struct ServerParams {
9595
pub(crate) no_delay: bool,
9696
/// Listener backlog.
9797
pub(crate) backlog: i32,
98+
/// RR or CRR.
99+
pub(crate) mode: Mode,
98100
}
99101

100102
enum ClientState {
@@ -121,6 +123,7 @@ enum ClientAction {
121123
enum ServerState {
122124
RecvRequest,
123125
SendResponse,
126+
CloseStream,
124127
}
125128

126129
const LISTENER_TOKEN: Token = Token(0);
@@ -275,7 +278,7 @@ fn client_thread_main(
275278
let mut events = Events::with_capacity(num_flows as usize);
276279
let request_buf = vec![0u8; request_size];
277280
let mut response_buf = vec![0u8; response_size];
278-
let mut flows: FlowMap<ClientState> = FlowMap::new();
281+
let mut flows: FlowMap<ClientState> = FlowMap::new(num_flows as usize);
279282
let mut next_token: usize = 0;
280283

281284
for _ in 0..num_flows {
@@ -291,12 +294,14 @@ fn client_thread_main(
291294
poll.registry()
292295
.register(&mut stream, token, Interest::WRITABLE)
293296
.expect("failed to register flow");
294-
flows.insert(Flow {
295-
stream,
296-
token,
297-
state: ClientState::SendRequest,
298-
xfer: request_size,
299-
});
297+
flows
298+
.insert(Flow {
299+
stream,
300+
token,
301+
state: ClientState::SendRequest,
302+
xfer: request_size,
303+
})
304+
.expect("client should never be able to exceed FlowMap capacity");
300305
metrics.connections_initiated.add(1);
301306
}
302307
Err(e) => {
@@ -564,6 +569,7 @@ pub(crate) async fn run_server(
564569
&flag,
565570
&tm[i as usize],
566571
tx,
572+
params.mode,
567573
);
568574
});
569575
handles.push(handle);
@@ -713,6 +719,7 @@ fn server_thread_main(
713719
shutdown_flag: &AtomicBool,
714720
metrics: &ThreadMetrics,
715721
ready_tx: mpsc::UnboundedSender<()>,
722+
mode: Mode,
716723
) {
717724
// Thread 0 uses the pre-built listener (with BPF already attached); others
718725
// bind their own sockets that join the existing reuseport group.
@@ -736,7 +743,7 @@ fn server_thread_main(
736743

737744
let mut request_buf = vec![0u8; request_size];
738745
let response_buf = vec![0u8; response_size];
739-
let mut flows: FlowMap<ServerState> = FlowMap::new();
746+
let mut flows: FlowMap<ServerState> = FlowMap::new(num_flows as usize);
740747
let mut next_token: usize = 1;
741748

742749
loop {
@@ -754,16 +761,22 @@ fn server_thread_main(
754761
set_nodelay_mio(&stream, no_delay);
755762
let token = Token(next_token);
756763
next_token += 1;
757-
let mut mio_stream = stream;
758-
poll.registry()
759-
.register(&mut mio_stream, token, Interest::READABLE)
760-
.expect("failed to register flow");
761-
flows.insert(Flow {
762-
stream: mio_stream,
764+
// Insert first: `insert` takes the flow by value and
765+
// drops it (closing the fd) on failure, so there is
766+
// nothing registered to clean up on the error path.
767+
if let Err(err) = flows.insert(Flow {
768+
stream,
763769
token,
764770
state: ServerState::RecvRequest,
765771
xfer: request_size,
766-
});
772+
}) {
773+
warn!("failed to insert flow in server FlowMap: {err}");
774+
break;
775+
}
776+
let flow = flows.get_mut(token).expect("flow was just inserted");
777+
poll.registry()
778+
.register(&mut flow.stream, token, Interest::READABLE)
779+
.expect("failed to register flow");
767780
metrics.connections_accepted.add(1);
768781
}
769782
Err(ref e) if e.kind() == ErrorKind::WouldBlock => break,
@@ -782,7 +795,8 @@ fn server_thread_main(
782795
let Some(fl) = flows.get_mut(token) else {
783796
continue;
784797
};
785-
let action = handle_server_event(fl, &mut request_buf, &response_buf, metrics);
798+
let action =
799+
handle_server_event(fl, &mut request_buf, &response_buf, metrics, mode);
786800
flow::apply_action(action, token, &mut flows, poll.registry());
787801
}
788802
}
@@ -802,6 +816,7 @@ fn handle_server_event(
802816
request_buf: &mut [u8],
803817
response_buf: &[u8],
804818
metrics: &ThreadMetrics,
819+
mode: Mode,
805820
) -> Action {
806821
match flow.state {
807822
ServerState::RecvRequest => {
@@ -838,7 +853,10 @@ fn handle_server_event(
838853
flow.xfer -= n;
839854
if flow.xfer == 0 {
840855
flow.xfer = request_buf.len();
841-
flow.state = ServerState::RecvRequest;
856+
match mode {
857+
Mode::Rr => flow.state = ServerState::RecvRequest,
858+
Mode::Crr => flow.state = ServerState::CloseStream,
859+
}
842860
metrics.responses_sent.add(1);
843861
metrics.bytes_written.add(response_buf.len() as u64);
844862
Action::Reregister(Interest::READABLE)
@@ -854,5 +872,6 @@ fn handle_server_event(
854872
}
855873
}
856874
}
875+
ServerState::CloseStream => Action::Remove,
857876
}
858877
}

0 commit comments

Comments
 (0)