mirror of
https://github.com/YGGverse/aquatic.git
synced 2026-03-31 17:55:36 +00:00
Improve CPU pinning
This commit is contained in:
parent
5057ba73bd
commit
fb607ac0c2
20 changed files with 219 additions and 143 deletions
2
Cargo.lock
generated
2
Cargo.lock
generated
|
|
@ -75,6 +75,7 @@ dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"aquatic_toml_config",
|
"aquatic_toml_config",
|
||||||
"arc-swap",
|
"arc-swap",
|
||||||
|
"glommio",
|
||||||
"hashbrown 0.12.0",
|
"hashbrown 0.12.0",
|
||||||
"hex",
|
"hex",
|
||||||
"hwloc",
|
"hwloc",
|
||||||
|
|
@ -101,6 +102,7 @@ dependencies = [
|
||||||
"futures-rustls",
|
"futures-rustls",
|
||||||
"glommio",
|
"glommio",
|
||||||
"itoa 1.0.1",
|
"itoa 1.0.1",
|
||||||
|
"libc",
|
||||||
"log",
|
"log",
|
||||||
"memchr",
|
"memchr",
|
||||||
"mimalloc",
|
"mimalloc",
|
||||||
|
|
|
||||||
|
|
@ -12,7 +12,8 @@ readme = "../README.md"
|
||||||
name = "aquatic_common"
|
name = "aquatic_common"
|
||||||
|
|
||||||
[features]
|
[features]
|
||||||
cpu-pinning = ["hwloc", "libc"]
|
with-glommio = ["glommio"]
|
||||||
|
with-hwloc = ["hwloc"]
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
aquatic_toml_config = "0.2.0"
|
aquatic_toml_config = "0.2.0"
|
||||||
|
|
@ -23,11 +24,12 @@ arc-swap = "1"
|
||||||
hashbrown = "0.12"
|
hashbrown = "0.12"
|
||||||
hex = "0.4"
|
hex = "0.4"
|
||||||
indexmap-amortized = "1"
|
indexmap-amortized = "1"
|
||||||
|
libc = "0.2"
|
||||||
log = "0.4"
|
log = "0.4"
|
||||||
privdrop = "0.5"
|
privdrop = "0.5"
|
||||||
rand = { version = "0.8", features = ["small_rng"] }
|
rand = { version = "0.8", features = ["small_rng"] }
|
||||||
serde = { version = "1", features = ["derive"] }
|
serde = { version = "1", features = ["derive"] }
|
||||||
|
|
||||||
# cpu-pinning
|
# Optional
|
||||||
|
glommio = { version = "0.7", optional = true }
|
||||||
hwloc = { version = "0.5", optional = true }
|
hwloc = { version = "0.5", optional = true }
|
||||||
libc = { version = "0.2", optional = true }
|
|
||||||
|
|
@ -1,5 +1,4 @@
|
||||||
use aquatic_toml_config::TomlConfig;
|
use aquatic_toml_config::TomlConfig;
|
||||||
use hwloc::{CpuSet, ObjectType, Topology, CPUBIND_THREAD};
|
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
|
|
||||||
#[derive(Clone, Debug, PartialEq, TomlConfig, Serialize, Deserialize)]
|
#[derive(Clone, Debug, PartialEq, TomlConfig, Serialize, Deserialize)]
|
||||||
|
|
@ -15,6 +14,7 @@ impl Default for CpuPinningMode {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Experimental CPU pinning
|
||||||
#[derive(Clone, Debug, PartialEq, TomlConfig, Deserialize)]
|
#[derive(Clone, Debug, PartialEq, TomlConfig, Deserialize)]
|
||||||
pub struct CpuPinningConfig {
|
pub struct CpuPinningConfig {
|
||||||
pub active: bool,
|
pub active: bool,
|
||||||
|
|
@ -45,37 +45,128 @@ impl CpuPinningConfig {
|
||||||
pub enum WorkerIndex {
|
pub enum WorkerIndex {
|
||||||
SocketWorker(usize),
|
SocketWorker(usize),
|
||||||
RequestWorker(usize),
|
RequestWorker(usize),
|
||||||
Other,
|
Util,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl WorkerIndex {
|
impl WorkerIndex {
|
||||||
fn get_core_index(
|
pub fn get_core_index(
|
||||||
self,
|
&self,
|
||||||
config: &CpuPinningConfig,
|
config: &CpuPinningConfig,
|
||||||
socket_workers: usize,
|
socket_workers: usize,
|
||||||
core_count: usize,
|
request_workers: usize,
|
||||||
|
num_cores: usize,
|
||||||
) -> usize {
|
) -> usize {
|
||||||
let ascending_index = match self {
|
let ascending_index = match self {
|
||||||
Self::Other => config.core_offset,
|
Self::SocketWorker(index) => config.core_offset + index,
|
||||||
Self::SocketWorker(index) => config.core_offset + 1 + index,
|
Self::RequestWorker(index) => config.core_offset + socket_workers + index,
|
||||||
Self::RequestWorker(index) => config.core_offset + 1 + socket_workers + index,
|
Self::Util => config.core_offset + socket_workers + request_workers,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
let max_core_index = num_cores - 1;
|
||||||
|
|
||||||
|
let ascending_index = ascending_index.min(max_core_index);
|
||||||
|
|
||||||
match config.mode {
|
match config.mode {
|
||||||
CpuPinningMode::Ascending => ascending_index,
|
CpuPinningMode::Ascending => ascending_index,
|
||||||
CpuPinningMode::Descending => core_count - 1 - ascending_index,
|
CpuPinningMode::Descending => max_core_index - ascending_index,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "with-glommio")]
|
||||||
|
pub mod glommio {
|
||||||
|
use ::glommio::{CpuSet, Placement};
|
||||||
|
|
||||||
|
use super::*;
|
||||||
|
|
||||||
|
fn get_cpu_set() -> anyhow::Result<CpuSet> {
|
||||||
|
CpuSet::online().map_err(|err| anyhow::anyhow!("Couldn't get CPU set: {:#}", err))
|
||||||
|
}
|
||||||
|
|
||||||
|
fn get_num_cpu_cores() -> anyhow::Result<usize> {
|
||||||
|
get_cpu_set()?
|
||||||
|
.iter()
|
||||||
|
.map(|l| l.core)
|
||||||
|
.max()
|
||||||
|
.ok_or(anyhow::anyhow!("CpuSet is empty"))
|
||||||
|
}
|
||||||
|
|
||||||
|
fn get_worker_cpu_set(
|
||||||
|
config: &CpuPinningConfig,
|
||||||
|
socket_workers: usize,
|
||||||
|
request_workers: usize,
|
||||||
|
worker_index: WorkerIndex,
|
||||||
|
) -> anyhow::Result<CpuSet> {
|
||||||
|
let num_cpu_cores = get_num_cpu_cores()?;
|
||||||
|
|
||||||
|
let core_index =
|
||||||
|
worker_index.get_core_index(&config, socket_workers, request_workers, num_cpu_cores);
|
||||||
|
|
||||||
|
Ok(get_cpu_set()?.filter(|l| l.core == core_index))
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn get_worker_placement(
|
||||||
|
config: &CpuPinningConfig,
|
||||||
|
socket_workers: usize,
|
||||||
|
request_workers: usize,
|
||||||
|
worker_index: WorkerIndex,
|
||||||
|
) -> anyhow::Result<Placement> {
|
||||||
|
if config.active {
|
||||||
|
let cpu_set =
|
||||||
|
get_worker_cpu_set(&config, socket_workers, request_workers, worker_index)?;
|
||||||
|
|
||||||
|
Ok(Placement::Fenced(cpu_set))
|
||||||
|
} else {
|
||||||
|
Ok(Placement::Unbound)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn set_affinity_for_util_worker(
|
||||||
|
config: &CpuPinningConfig,
|
||||||
|
socket_workers: usize,
|
||||||
|
request_workers: usize,
|
||||||
|
) -> anyhow::Result<()> {
|
||||||
|
let worker_cpu_set =
|
||||||
|
get_worker_cpu_set(&config, socket_workers, request_workers, WorkerIndex::Util)?;
|
||||||
|
|
||||||
|
let logical_cpus: Vec<usize> = worker_cpu_set.iter().map(|l| l.cpu).collect();
|
||||||
|
|
||||||
|
unsafe {
|
||||||
|
let mut set: libc::cpu_set_t = ::std::mem::zeroed();
|
||||||
|
|
||||||
|
for cpu in logical_cpus {
|
||||||
|
libc::CPU_SET(cpu, &mut set);
|
||||||
|
}
|
||||||
|
|
||||||
|
let status = libc::pthread_setaffinity_np(
|
||||||
|
libc::pthread_self(),
|
||||||
|
::std::mem::size_of::<libc::cpu_set_t>(),
|
||||||
|
&set,
|
||||||
|
);
|
||||||
|
|
||||||
|
if status != 0 {
|
||||||
|
return Err(anyhow::Error::new(::std::io::Error::from_raw_os_error(
|
||||||
|
status,
|
||||||
|
)));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Pin current thread to a suitable core
|
/// Pin current thread to a suitable core
|
||||||
///
|
///
|
||||||
/// Requires hwloc (`apt-get install libhwloc-dev`)
|
/// Requires hwloc (`apt-get install libhwloc-dev`)
|
||||||
|
#[cfg(feature = "with-hwloc")]
|
||||||
pub fn pin_current_if_configured_to(
|
pub fn pin_current_if_configured_to(
|
||||||
config: &CpuPinningConfig,
|
config: &CpuPinningConfig,
|
||||||
socket_workers: usize,
|
socket_workers: usize,
|
||||||
|
request_workers: usize,
|
||||||
worker_index: WorkerIndex,
|
worker_index: WorkerIndex,
|
||||||
) {
|
) {
|
||||||
|
use hwloc::{CpuSet, ObjectType, Topology, CPUBIND_THREAD};
|
||||||
|
|
||||||
if config.active {
|
if config.active {
|
||||||
let mut topology = Topology::new();
|
let mut topology = Topology::new();
|
||||||
|
|
||||||
|
|
@ -86,7 +177,10 @@ pub fn pin_current_if_configured_to(
|
||||||
.map(|core| core.allowed_cpuset().expect("hwloc: get core cpu set"))
|
.map(|core| core.allowed_cpuset().expect("hwloc: get core cpu set"))
|
||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
let core_index = worker_index.get_core_index(config, socket_workers, core_cpu_sets.len());
|
let num_cores = core_cpu_sets.len();
|
||||||
|
|
||||||
|
let core_index =
|
||||||
|
worker_index.get_core_index(config, socket_workers, request_workers, num_cores);
|
||||||
|
|
||||||
let cpu_set = core_cpu_sets
|
let cpu_set = core_cpu_sets
|
||||||
.get(core_index)
|
.get(core_index)
|
||||||
|
|
|
||||||
|
|
@ -5,7 +5,6 @@ use ahash::RandomState;
|
||||||
use rand::Rng;
|
use rand::Rng;
|
||||||
|
|
||||||
pub mod access_list;
|
pub mod access_list;
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
pub mod cpu_pinning;
|
pub mod cpu_pinning;
|
||||||
pub mod privileges;
|
pub mod privileges;
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -15,12 +15,9 @@ name = "aquatic_http"
|
||||||
[[bin]]
|
[[bin]]
|
||||||
name = "aquatic_http"
|
name = "aquatic_http"
|
||||||
|
|
||||||
[features]
|
|
||||||
cpu-pinning = ["aquatic_common/cpu-pinning"]
|
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
aquatic_cli_helpers = "0.2.0"
|
aquatic_cli_helpers = "0.2.0"
|
||||||
aquatic_common = "0.2.0"
|
aquatic_common = { version = "0.2.0", features = ["with-glommio"] }
|
||||||
aquatic_http_protocol = "0.2.0"
|
aquatic_http_protocol = "0.2.0"
|
||||||
aquatic_toml_config = "0.2.0"
|
aquatic_toml_config = "0.2.0"
|
||||||
|
|
||||||
|
|
@ -31,6 +28,7 @@ futures-lite = "1"
|
||||||
futures-rustls = "0.22"
|
futures-rustls = "0.22"
|
||||||
glommio = "0.7"
|
glommio = "0.7"
|
||||||
itoa = "1"
|
itoa = "1"
|
||||||
|
libc = "0.2"
|
||||||
log = "0.4"
|
log = "0.4"
|
||||||
mimalloc = { version = "0.1", default-features = false }
|
mimalloc = { version = "0.1", default-features = false }
|
||||||
memchr = "2"
|
memchr = "2"
|
||||||
|
|
|
||||||
|
|
@ -23,7 +23,6 @@ pub struct Config {
|
||||||
pub cleaning: CleaningConfig,
|
pub cleaning: CleaningConfig,
|
||||||
pub privileges: PrivilegeConfig,
|
pub privileges: PrivilegeConfig,
|
||||||
pub access_list: AccessListConfig,
|
pub access_list: AccessListConfig,
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
pub cpu_pinning: aquatic_common::cpu_pinning::CpuPinningConfig,
|
pub cpu_pinning: aquatic_common::cpu_pinning::CpuPinningConfig,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -38,7 +37,6 @@ impl Default for Config {
|
||||||
cleaning: CleaningConfig::default(),
|
cleaning: CleaningConfig::default(),
|
||||||
privileges: PrivilegeConfig::default(),
|
privileges: PrivilegeConfig::default(),
|
||||||
access_list: AccessListConfig::default(),
|
access_list: AccessListConfig::default(),
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
cpu_pinning: Default::default(),
|
cpu_pinning: Default::default(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,7 +1,10 @@
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
use aquatic_common::cpu_pinning::{pin_current_if_configured_to, WorkerIndex};
|
|
||||||
use aquatic_common::{
|
use aquatic_common::{
|
||||||
access_list::update_access_list, privileges::drop_privileges_after_socket_binding,
|
access_list::update_access_list,
|
||||||
|
cpu_pinning::{
|
||||||
|
glommio::{get_worker_placement, set_affinity_for_util_worker},
|
||||||
|
WorkerIndex,
|
||||||
|
},
|
||||||
|
privileges::drop_privileges_after_socket_binding,
|
||||||
};
|
};
|
||||||
use common::{State, TlsConfig};
|
use common::{State, TlsConfig};
|
||||||
use glommio::{channels::channel_mesh::MeshBuilder, prelude::*};
|
use glommio::{channels::channel_mesh::MeshBuilder, prelude::*};
|
||||||
|
|
@ -37,12 +40,13 @@ pub fn run(config: Config) -> ::anyhow::Result<()> {
|
||||||
::std::thread::spawn(move || run_inner(config, state));
|
::std::thread::spawn(move || run_inner(config, state));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(feature = "cpu-pinning")]
|
if config.cpu_pinning.active {
|
||||||
pin_current_if_configured_to(
|
set_affinity_for_util_worker(
|
||||||
&config.cpu_pinning,
|
&config.cpu_pinning,
|
||||||
config.socket_workers,
|
config.socket_workers,
|
||||||
WorkerIndex::Other,
|
config.request_workers,
|
||||||
);
|
)?;
|
||||||
|
}
|
||||||
|
|
||||||
for signal in &mut signals {
|
for signal in &mut signals {
|
||||||
match signal {
|
match signal {
|
||||||
|
|
@ -76,16 +80,15 @@ pub fn run_inner(config: Config, state: State) -> anyhow::Result<()> {
|
||||||
let response_mesh_builder = response_mesh_builder.clone();
|
let response_mesh_builder = response_mesh_builder.clone();
|
||||||
let num_bound_sockets = num_bound_sockets.clone();
|
let num_bound_sockets = num_bound_sockets.clone();
|
||||||
|
|
||||||
let builder = LocalExecutorBuilder::default().name("socket");
|
let placement = get_worker_placement(
|
||||||
|
&config.cpu_pinning,
|
||||||
|
config.socket_workers,
|
||||||
|
config.request_workers,
|
||||||
|
WorkerIndex::SocketWorker(i),
|
||||||
|
)?;
|
||||||
|
let builder = LocalExecutorBuilder::new(placement).name("socket");
|
||||||
|
|
||||||
let executor = builder.spawn(move || async move {
|
let executor = builder.spawn(move || async move {
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
pin_current_if_configured_to(
|
|
||||||
&config.cpu_pinning,
|
|
||||||
config.socket_workers,
|
|
||||||
WorkerIndex::SocketWorker(i),
|
|
||||||
);
|
|
||||||
|
|
||||||
workers::socket::run_socket_worker(
|
workers::socket::run_socket_worker(
|
||||||
config,
|
config,
|
||||||
state,
|
state,
|
||||||
|
|
@ -106,16 +109,15 @@ pub fn run_inner(config: Config, state: State) -> anyhow::Result<()> {
|
||||||
let request_mesh_builder = request_mesh_builder.clone();
|
let request_mesh_builder = request_mesh_builder.clone();
|
||||||
let response_mesh_builder = response_mesh_builder.clone();
|
let response_mesh_builder = response_mesh_builder.clone();
|
||||||
|
|
||||||
let builder = LocalExecutorBuilder::default().name("request");
|
let placement = get_worker_placement(
|
||||||
|
&config.cpu_pinning,
|
||||||
|
config.socket_workers,
|
||||||
|
config.request_workers,
|
||||||
|
WorkerIndex::RequestWorker(i),
|
||||||
|
)?;
|
||||||
|
let builder = LocalExecutorBuilder::new(placement).name("request");
|
||||||
|
|
||||||
let executor = builder.spawn(move || async move {
|
let executor = builder.spawn(move || async move {
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
pin_current_if_configured_to(
|
|
||||||
&config.cpu_pinning,
|
|
||||||
config.socket_workers,
|
|
||||||
WorkerIndex::RequestWorker(i),
|
|
||||||
);
|
|
||||||
|
|
||||||
workers::request::run_request_worker(
|
workers::request::run_request_worker(
|
||||||
config,
|
config,
|
||||||
state,
|
state,
|
||||||
|
|
@ -135,12 +137,13 @@ pub fn run_inner(config: Config, state: State) -> anyhow::Result<()> {
|
||||||
)
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
#[cfg(feature = "cpu-pinning")]
|
if config.cpu_pinning.active {
|
||||||
pin_current_if_configured_to(
|
set_affinity_for_util_worker(
|
||||||
&config.cpu_pinning,
|
&config.cpu_pinning,
|
||||||
config.socket_workers,
|
config.socket_workers,
|
||||||
WorkerIndex::Other,
|
config.request_workers,
|
||||||
);
|
)?;
|
||||||
|
}
|
||||||
|
|
||||||
for executor in executors {
|
for executor in executors {
|
||||||
executor
|
executor
|
||||||
|
|
|
||||||
|
|
@ -12,12 +12,9 @@ readme = "../README.md"
|
||||||
[[bin]]
|
[[bin]]
|
||||||
name = "aquatic_http_load_test"
|
name = "aquatic_http_load_test"
|
||||||
|
|
||||||
[features]
|
|
||||||
cpu-pinning = ["aquatic_common/cpu-pinning"]
|
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
aquatic_cli_helpers = "0.2.0"
|
aquatic_cli_helpers = "0.2.0"
|
||||||
aquatic_common = "0.2.0"
|
aquatic_common = { version = "0.2.0", features = ["with-glommio"] }
|
||||||
aquatic_http_protocol = "0.2.0"
|
aquatic_http_protocol = "0.2.0"
|
||||||
aquatic_toml_config = "0.2.0"
|
aquatic_toml_config = "0.2.0"
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -20,7 +20,6 @@ pub struct Config {
|
||||||
pub connection_creation_interval_ms: u64,
|
pub connection_creation_interval_ms: u64,
|
||||||
pub duration: usize,
|
pub duration: usize,
|
||||||
pub torrents: TorrentConfig,
|
pub torrents: TorrentConfig,
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
pub cpu_pinning: aquatic_common::cpu_pinning::CpuPinningConfig,
|
pub cpu_pinning: aquatic_common::cpu_pinning::CpuPinningConfig,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -58,7 +57,6 @@ impl Default for Config {
|
||||||
connection_creation_interval_ms: 10,
|
connection_creation_interval_ms: 10,
|
||||||
duration: 0,
|
duration: 0,
|
||||||
torrents: TorrentConfig::default(),
|
torrents: TorrentConfig::default(),
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
cpu_pinning: aquatic_common::cpu_pinning::CpuPinningConfig::default_for_load_test(),
|
cpu_pinning: aquatic_common::cpu_pinning::CpuPinningConfig::default_for_load_test(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -3,8 +3,8 @@ use std::thread;
|
||||||
use std::time::{Duration, Instant};
|
use std::time::{Duration, Instant};
|
||||||
|
|
||||||
use ::glommio::LocalExecutorBuilder;
|
use ::glommio::LocalExecutorBuilder;
|
||||||
#[cfg(feature = "cpu-pinning")]
|
use aquatic_common::cpu_pinning::glommio::{get_worker_placement, set_affinity_for_util_worker};
|
||||||
use aquatic_common::cpu_pinning::{pin_current_if_configured_to, WorkerIndex};
|
use aquatic_common::cpu_pinning::WorkerIndex;
|
||||||
use rand::prelude::*;
|
use rand::prelude::*;
|
||||||
use rand_distr::Pareto;
|
use rand_distr::Pareto;
|
||||||
|
|
||||||
|
|
@ -55,8 +55,6 @@ fn run(config: Config) -> ::anyhow::Result<()> {
|
||||||
pareto: Arc::new(pareto),
|
pareto: Arc::new(pareto),
|
||||||
};
|
};
|
||||||
|
|
||||||
// Start socket workers
|
|
||||||
|
|
||||||
let tls_config = create_tls_config().unwrap();
|
let tls_config = create_tls_config().unwrap();
|
||||||
|
|
||||||
for i in 0..config.num_workers {
|
for i in 0..config.num_workers {
|
||||||
|
|
@ -64,26 +62,23 @@ fn run(config: Config) -> ::anyhow::Result<()> {
|
||||||
let tls_config = tls_config.clone();
|
let tls_config = tls_config.clone();
|
||||||
let state = state.clone();
|
let state = state.clone();
|
||||||
|
|
||||||
LocalExecutorBuilder::default()
|
let placement = get_worker_placement(
|
||||||
.spawn(move || async move {
|
&config.cpu_pinning,
|
||||||
#[cfg(feature = "cpu-pinning")]
|
config.num_workers,
|
||||||
pin_current_if_configured_to(
|
0,
|
||||||
&config.cpu_pinning,
|
WorkerIndex::SocketWorker(i),
|
||||||
config.num_workers,
|
)?;
|
||||||
WorkerIndex::SocketWorker(i),
|
|
||||||
);
|
|
||||||
|
|
||||||
|
LocalExecutorBuilder::new(placement)
|
||||||
|
.spawn(move || async move {
|
||||||
run_socket_thread(config, tls_config, state).await.unwrap();
|
run_socket_thread(config, tls_config, state).await.unwrap();
|
||||||
})
|
})
|
||||||
.unwrap();
|
.unwrap();
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(feature = "cpu-pinning")]
|
if config.cpu_pinning.active {
|
||||||
pin_current_if_configured_to(
|
set_affinity_for_util_worker(&config.cpu_pinning, config.num_workers, 0)?;
|
||||||
&config.cpu_pinning,
|
}
|
||||||
config.num_workers as usize,
|
|
||||||
WorkerIndex::Other,
|
|
||||||
);
|
|
||||||
|
|
||||||
monitor_statistics(state, &config);
|
monitor_statistics(state, &config);
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -16,7 +16,7 @@ name = "aquatic_udp"
|
||||||
name = "aquatic_udp"
|
name = "aquatic_udp"
|
||||||
|
|
||||||
[features]
|
[features]
|
||||||
cpu-pinning = ["aquatic_common/cpu-pinning"]
|
cpu-pinning = ["aquatic_common/with-hwloc"]
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
aquatic_cli_helpers = "0.2.0"
|
aquatic_cli_helpers = "0.2.0"
|
||||||
|
|
|
||||||
|
|
@ -75,6 +75,7 @@ pub fn run(config: Config) -> ::anyhow::Result<()> {
|
||||||
pin_current_if_configured_to(
|
pin_current_if_configured_to(
|
||||||
&config.cpu_pinning,
|
&config.cpu_pinning,
|
||||||
config.socket_workers,
|
config.socket_workers,
|
||||||
|
config.request_workers,
|
||||||
WorkerIndex::RequestWorker(i),
|
WorkerIndex::RequestWorker(i),
|
||||||
);
|
);
|
||||||
|
|
||||||
|
|
@ -104,6 +105,7 @@ pub fn run(config: Config) -> ::anyhow::Result<()> {
|
||||||
pin_current_if_configured_to(
|
pin_current_if_configured_to(
|
||||||
&config.cpu_pinning,
|
&config.cpu_pinning,
|
||||||
config.socket_workers,
|
config.socket_workers,
|
||||||
|
config.request_workers,
|
||||||
WorkerIndex::SocketWorker(i),
|
WorkerIndex::SocketWorker(i),
|
||||||
);
|
);
|
||||||
|
|
||||||
|
|
@ -130,7 +132,8 @@ pub fn run(config: Config) -> ::anyhow::Result<()> {
|
||||||
pin_current_if_configured_to(
|
pin_current_if_configured_to(
|
||||||
&config.cpu_pinning,
|
&config.cpu_pinning,
|
||||||
config.socket_workers,
|
config.socket_workers,
|
||||||
WorkerIndex::Other,
|
config.request_workers,
|
||||||
|
WorkerIndex::Util,
|
||||||
);
|
);
|
||||||
|
|
||||||
workers::statistics::run_statistics_worker(config, state);
|
workers::statistics::run_statistics_worker(config, state);
|
||||||
|
|
@ -149,7 +152,8 @@ pub fn run(config: Config) -> ::anyhow::Result<()> {
|
||||||
pin_current_if_configured_to(
|
pin_current_if_configured_to(
|
||||||
&config.cpu_pinning,
|
&config.cpu_pinning,
|
||||||
config.socket_workers,
|
config.socket_workers,
|
||||||
WorkerIndex::Other,
|
config.request_workers,
|
||||||
|
WorkerIndex::Util,
|
||||||
);
|
);
|
||||||
|
|
||||||
for signal in &mut signals {
|
for signal in &mut signals {
|
||||||
|
|
|
||||||
|
|
@ -9,12 +9,12 @@ repository = "https://github.com/greatest-ape/aquatic"
|
||||||
keywords = ["udp", "benchmark", "peer-to-peer", "torrent", "bittorrent"]
|
keywords = ["udp", "benchmark", "peer-to-peer", "torrent", "bittorrent"]
|
||||||
readme = "../README.md"
|
readme = "../README.md"
|
||||||
|
|
||||||
|
[features]
|
||||||
|
cpu-pinning = ["aquatic_common/with-hwloc"]
|
||||||
|
|
||||||
[[bin]]
|
[[bin]]
|
||||||
name = "aquatic_udp_load_test"
|
name = "aquatic_udp_load_test"
|
||||||
|
|
||||||
[features]
|
|
||||||
cpu-pinning = ["aquatic_common/cpu-pinning"]
|
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
aquatic_cli_helpers = "0.2.0"
|
aquatic_cli_helpers = "0.2.0"
|
||||||
aquatic_common = "0.2.0"
|
aquatic_common = "0.2.0"
|
||||||
|
|
|
||||||
|
|
@ -84,6 +84,7 @@ fn run(config: Config) -> ::anyhow::Result<()> {
|
||||||
pin_current_if_configured_to(
|
pin_current_if_configured_to(
|
||||||
&config.cpu_pinning,
|
&config.cpu_pinning,
|
||||||
config.workers as usize,
|
config.workers as usize,
|
||||||
|
0,
|
||||||
WorkerIndex::SocketWorker(i as usize),
|
WorkerIndex::SocketWorker(i as usize),
|
||||||
);
|
);
|
||||||
|
|
||||||
|
|
@ -95,7 +96,8 @@ fn run(config: Config) -> ::anyhow::Result<()> {
|
||||||
pin_current_if_configured_to(
|
pin_current_if_configured_to(
|
||||||
&config.cpu_pinning,
|
&config.cpu_pinning,
|
||||||
config.workers as usize,
|
config.workers as usize,
|
||||||
WorkerIndex::Other,
|
0,
|
||||||
|
WorkerIndex::Util,
|
||||||
);
|
);
|
||||||
|
|
||||||
monitor_statistics(state, &config);
|
monitor_statistics(state, &config);
|
||||||
|
|
|
||||||
|
|
@ -9,19 +9,15 @@ repository = "https://github.com/greatest-ape/aquatic"
|
||||||
keywords = ["webtorrent", "websocket", "peer-to-peer", "torrent", "bittorrent"]
|
keywords = ["webtorrent", "websocket", "peer-to-peer", "torrent", "bittorrent"]
|
||||||
readme = "../README.md"
|
readme = "../README.md"
|
||||||
|
|
||||||
|
|
||||||
[lib]
|
[lib]
|
||||||
name = "aquatic_ws"
|
name = "aquatic_ws"
|
||||||
|
|
||||||
[[bin]]
|
[[bin]]
|
||||||
name = "aquatic_ws"
|
name = "aquatic_ws"
|
||||||
|
|
||||||
[features]
|
|
||||||
cpu-pinning = ["aquatic_common/cpu-pinning"]
|
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
aquatic_cli_helpers = "0.2.0"
|
aquatic_cli_helpers = "0.2.0"
|
||||||
aquatic_common = "0.2.0"
|
aquatic_common = { version = "0.2.0", features = ["with-glommio"] }
|
||||||
aquatic_toml_config = "0.2.0"
|
aquatic_toml_config = "0.2.0"
|
||||||
aquatic_ws_protocol = "0.2.0"
|
aquatic_ws_protocol = "0.2.0"
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,7 +1,6 @@
|
||||||
use std::net::SocketAddr;
|
use std::net::SocketAddr;
|
||||||
use std::path::PathBuf;
|
use std::path::PathBuf;
|
||||||
|
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
use aquatic_common::cpu_pinning::CpuPinningConfig;
|
use aquatic_common::cpu_pinning::CpuPinningConfig;
|
||||||
use aquatic_common::{access_list::AccessListConfig, privileges::PrivilegeConfig};
|
use aquatic_common::{access_list::AccessListConfig, privileges::PrivilegeConfig};
|
||||||
use serde::Deserialize;
|
use serde::Deserialize;
|
||||||
|
|
@ -26,7 +25,6 @@ pub struct Config {
|
||||||
pub cleaning: CleaningConfig,
|
pub cleaning: CleaningConfig,
|
||||||
pub privileges: PrivilegeConfig,
|
pub privileges: PrivilegeConfig,
|
||||||
pub access_list: AccessListConfig,
|
pub access_list: AccessListConfig,
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
pub cpu_pinning: CpuPinningConfig,
|
pub cpu_pinning: CpuPinningConfig,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -41,7 +39,6 @@ impl Default for Config {
|
||||||
cleaning: CleaningConfig::default(),
|
cleaning: CleaningConfig::default(),
|
||||||
privileges: PrivilegeConfig::default(),
|
privileges: PrivilegeConfig::default(),
|
||||||
access_list: AccessListConfig::default(),
|
access_list: AccessListConfig::default(),
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
cpu_pinning: Default::default(),
|
cpu_pinning: Default::default(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -6,12 +6,12 @@ use std::fs::File;
|
||||||
use std::io::BufReader;
|
use std::io::BufReader;
|
||||||
use std::sync::{atomic::AtomicUsize, Arc};
|
use std::sync::{atomic::AtomicUsize, Arc};
|
||||||
|
|
||||||
|
use aquatic_common::cpu_pinning::glommio::{get_worker_placement, set_affinity_for_util_worker};
|
||||||
|
use aquatic_common::cpu_pinning::WorkerIndex;
|
||||||
use glommio::{channels::channel_mesh::MeshBuilder, prelude::*};
|
use glommio::{channels::channel_mesh::MeshBuilder, prelude::*};
|
||||||
use signal_hook::{consts::SIGUSR1, iterator::Signals};
|
use signal_hook::{consts::SIGUSR1, iterator::Signals};
|
||||||
|
|
||||||
use aquatic_common::access_list::update_access_list;
|
use aquatic_common::access_list::update_access_list;
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
use aquatic_common::cpu_pinning::{pin_current_if_configured_to, WorkerIndex};
|
|
||||||
use aquatic_common::privileges::drop_privileges_after_socket_binding;
|
use aquatic_common::privileges::drop_privileges_after_socket_binding;
|
||||||
|
|
||||||
use common::*;
|
use common::*;
|
||||||
|
|
@ -36,12 +36,13 @@ pub fn run(config: Config) -> ::anyhow::Result<()> {
|
||||||
::std::thread::spawn(move || run_workers(config, state));
|
::std::thread::spawn(move || run_workers(config, state));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(feature = "cpu-pinning")]
|
if config.cpu_pinning.active {
|
||||||
pin_current_if_configured_to(
|
set_affinity_for_util_worker(
|
||||||
&config.cpu_pinning,
|
&config.cpu_pinning,
|
||||||
config.socket_workers,
|
config.socket_workers,
|
||||||
WorkerIndex::Other,
|
config.request_workers,
|
||||||
);
|
)?;
|
||||||
|
}
|
||||||
|
|
||||||
for signal in &mut signals {
|
for signal in &mut signals {
|
||||||
match signal {
|
match signal {
|
||||||
|
|
@ -75,16 +76,15 @@ fn run_workers(config: Config, state: State) -> anyhow::Result<()> {
|
||||||
let response_mesh_builder = response_mesh_builder.clone();
|
let response_mesh_builder = response_mesh_builder.clone();
|
||||||
let num_bound_sockets = num_bound_sockets.clone();
|
let num_bound_sockets = num_bound_sockets.clone();
|
||||||
|
|
||||||
let builder = LocalExecutorBuilder::default().name("socket");
|
let placement = get_worker_placement(
|
||||||
|
&config.cpu_pinning,
|
||||||
|
config.socket_workers,
|
||||||
|
config.request_workers,
|
||||||
|
WorkerIndex::SocketWorker(i),
|
||||||
|
)?;
|
||||||
|
let builder = LocalExecutorBuilder::new(placement).name("socket");
|
||||||
|
|
||||||
let executor = builder.spawn(move || async move {
|
let executor = builder.spawn(move || async move {
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
pin_current_if_configured_to(
|
|
||||||
&config.cpu_pinning,
|
|
||||||
config.socket_workers,
|
|
||||||
WorkerIndex::SocketWorker(i),
|
|
||||||
);
|
|
||||||
|
|
||||||
workers::socket::run_socket_worker(
|
workers::socket::run_socket_worker(
|
||||||
config,
|
config,
|
||||||
state,
|
state,
|
||||||
|
|
@ -105,16 +105,15 @@ fn run_workers(config: Config, state: State) -> anyhow::Result<()> {
|
||||||
let request_mesh_builder = request_mesh_builder.clone();
|
let request_mesh_builder = request_mesh_builder.clone();
|
||||||
let response_mesh_builder = response_mesh_builder.clone();
|
let response_mesh_builder = response_mesh_builder.clone();
|
||||||
|
|
||||||
let builder = LocalExecutorBuilder::default().name("request");
|
let placement = get_worker_placement(
|
||||||
|
&config.cpu_pinning,
|
||||||
|
config.socket_workers,
|
||||||
|
config.request_workers,
|
||||||
|
WorkerIndex::RequestWorker(i),
|
||||||
|
)?;
|
||||||
|
let builder = LocalExecutorBuilder::new(placement).name("request");
|
||||||
|
|
||||||
let executor = builder.spawn(move || async move {
|
let executor = builder.spawn(move || async move {
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
pin_current_if_configured_to(
|
|
||||||
&config.cpu_pinning,
|
|
||||||
config.socket_workers,
|
|
||||||
WorkerIndex::RequestWorker(i),
|
|
||||||
);
|
|
||||||
|
|
||||||
workers::request::run_request_worker(
|
workers::request::run_request_worker(
|
||||||
config,
|
config,
|
||||||
state,
|
state,
|
||||||
|
|
@ -134,12 +133,13 @@ fn run_workers(config: Config, state: State) -> anyhow::Result<()> {
|
||||||
)
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
#[cfg(feature = "cpu-pinning")]
|
if config.cpu_pinning.active {
|
||||||
pin_current_if_configured_to(
|
set_affinity_for_util_worker(
|
||||||
&config.cpu_pinning,
|
&config.cpu_pinning,
|
||||||
config.socket_workers,
|
config.socket_workers,
|
||||||
WorkerIndex::Other,
|
config.request_workers,
|
||||||
);
|
)?;
|
||||||
|
}
|
||||||
|
|
||||||
for executor in executors {
|
for executor in executors {
|
||||||
executor
|
executor
|
||||||
|
|
|
||||||
|
|
@ -12,13 +12,10 @@ readme = "../README.md"
|
||||||
[[bin]]
|
[[bin]]
|
||||||
name = "aquatic_ws_load_test"
|
name = "aquatic_ws_load_test"
|
||||||
|
|
||||||
[features]
|
|
||||||
cpu-pinning = ["aquatic_common/cpu-pinning"]
|
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
async-tungstenite = "0.17"
|
async-tungstenite = "0.17"
|
||||||
aquatic_cli_helpers = "0.2.0"
|
aquatic_cli_helpers = "0.2.0"
|
||||||
aquatic_common = "0.2.0"
|
aquatic_common = { version = "0.2.0", features = ["with-glommio"] }
|
||||||
aquatic_toml_config = "0.2.0"
|
aquatic_toml_config = "0.2.0"
|
||||||
aquatic_ws_protocol = "0.2.0"
|
aquatic_ws_protocol = "0.2.0"
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,7 +1,6 @@
|
||||||
use std::net::SocketAddr;
|
use std::net::SocketAddr;
|
||||||
|
|
||||||
use aquatic_cli_helpers::LogLevel;
|
use aquatic_cli_helpers::LogLevel;
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
use aquatic_common::cpu_pinning::CpuPinningConfig;
|
use aquatic_common::cpu_pinning::CpuPinningConfig;
|
||||||
use aquatic_toml_config::TomlConfig;
|
use aquatic_toml_config::TomlConfig;
|
||||||
use serde::Deserialize;
|
use serde::Deserialize;
|
||||||
|
|
@ -17,7 +16,6 @@ pub struct Config {
|
||||||
pub connection_creation_interval_ms: u64,
|
pub connection_creation_interval_ms: u64,
|
||||||
pub duration: usize,
|
pub duration: usize,
|
||||||
pub torrents: TorrentConfig,
|
pub torrents: TorrentConfig,
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
pub cpu_pinning: CpuPinningConfig,
|
pub cpu_pinning: CpuPinningConfig,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -37,7 +35,6 @@ impl Default for Config {
|
||||||
connection_creation_interval_ms: 10,
|
connection_creation_interval_ms: 10,
|
||||||
duration: 0,
|
duration: 0,
|
||||||
torrents: TorrentConfig::default(),
|
torrents: TorrentConfig::default(),
|
||||||
#[cfg(feature = "cpu-pinning")]
|
|
||||||
cpu_pinning: CpuPinningConfig::default_for_load_test(),
|
cpu_pinning: CpuPinningConfig::default_for_load_test(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -2,8 +2,8 @@ use std::sync::{atomic::Ordering, Arc};
|
||||||
use std::thread;
|
use std::thread;
|
||||||
use std::time::{Duration, Instant};
|
use std::time::{Duration, Instant};
|
||||||
|
|
||||||
#[cfg(feature = "cpu-pinning")]
|
use aquatic_common::cpu_pinning::glommio::{get_worker_placement, set_affinity_for_util_worker};
|
||||||
use aquatic_common::cpu_pinning::{pin_current_if_configured_to, WorkerIndex};
|
use aquatic_common::cpu_pinning::WorkerIndex;
|
||||||
use glommio::LocalExecutorBuilder;
|
use glommio::LocalExecutorBuilder;
|
||||||
use rand::prelude::*;
|
use rand::prelude::*;
|
||||||
use rand_distr::Pareto;
|
use rand_distr::Pareto;
|
||||||
|
|
@ -59,26 +59,23 @@ fn run(config: Config) -> ::anyhow::Result<()> {
|
||||||
let tls_config = tls_config.clone();
|
let tls_config = tls_config.clone();
|
||||||
let state = state.clone();
|
let state = state.clone();
|
||||||
|
|
||||||
LocalExecutorBuilder::default()
|
let placement = get_worker_placement(
|
||||||
.spawn(move || async move {
|
&config.cpu_pinning,
|
||||||
#[cfg(feature = "cpu-pinning")]
|
config.num_workers,
|
||||||
pin_current_if_configured_to(
|
0,
|
||||||
&config.cpu_pinning,
|
WorkerIndex::SocketWorker(i),
|
||||||
config.num_workers,
|
)?;
|
||||||
WorkerIndex::SocketWorker(i),
|
|
||||||
);
|
|
||||||
|
|
||||||
|
LocalExecutorBuilder::new(placement)
|
||||||
|
.spawn(move || async move {
|
||||||
run_socket_thread(config, tls_config, state).await.unwrap();
|
run_socket_thread(config, tls_config, state).await.unwrap();
|
||||||
})
|
})
|
||||||
.unwrap();
|
.unwrap();
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(feature = "cpu-pinning")]
|
if config.cpu_pinning.active {
|
||||||
pin_current_if_configured_to(
|
set_affinity_for_util_worker(&config.cpu_pinning, config.num_workers, 0)?;
|
||||||
&config.cpu_pinning,
|
}
|
||||||
config.num_workers as usize,
|
|
||||||
WorkerIndex::Other,
|
|
||||||
);
|
|
||||||
|
|
||||||
monitor_statistics(state, &config);
|
monitor_statistics(state, &config);
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue