generated from PlexSheep/rs-base
concurrency
cargo devel CI / cargo CI (push) Successful in 2m29s
Details
cargo devel CI / cargo CI (push) Successful in 2m29s
Details
This commit is contained in:
parent
b6c9b660f6
commit
b0fb7de6e9
|
@ -1,3 +1,4 @@
|
|||
workspace = { members = ["spammer"] }
|
||||
[package]
|
||||
name = "netpong"
|
||||
version = "0.1.1"
|
||||
|
|
|
@ -1,11 +1,14 @@
|
|||
import socket
|
||||
|
||||
HOST = "127.0.0.1"
|
||||
PORT = 9999
|
||||
|
||||
payload = b"ping\0"
|
||||
def ping():
|
||||
|
||||
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
|
||||
HOST = "127.0.0.1"
|
||||
PORT = 9999
|
||||
|
||||
payload = b"ping\0"
|
||||
|
||||
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
|
||||
s.connect((HOST, PORT))
|
||||
while True:
|
||||
try:
|
||||
|
@ -19,3 +22,5 @@ with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
|
|||
print(f"< {reply}")
|
||||
print("connection shut down")
|
||||
|
||||
if __name__ == "__main__":
|
||||
ping()
|
||||
|
|
|
@ -0,0 +1,12 @@
|
|||
const MAX: usize = 50;
|
||||
use std::process::Command;
|
||||
|
||||
fn main() {
|
||||
let mut pool = ThreadPool::new(MAX);
|
||||
|
||||
loop {
|
||||
pool.execute(||{
|
||||
Command::new("python3").args(["scripts/client.py"]).output().unwrap();
|
||||
});
|
||||
}
|
||||
}
|
|
@ -0,0 +1,9 @@
|
|||
[package]
|
||||
name = "spammer"
|
||||
version = "0.1.0"
|
||||
edition = "2021"
|
||||
|
||||
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
|
||||
|
||||
[dependencies]
|
||||
threadpool = "1.8.1"
|
|
@ -0,0 +1,22 @@
|
|||
use threadpool::ThreadPool;
|
||||
const MAX: usize = 500;
|
||||
use std::process::{exit, Command};
|
||||
|
||||
fn main() {
|
||||
let pool = ThreadPool::new(MAX);
|
||||
|
||||
loop {
|
||||
if pool.queued_count() < MAX {
|
||||
pool.execute(|| {
|
||||
let mut cmd = Command::new("/usr/bin/python3");
|
||||
cmd.args(["../scripts/client.py"]);
|
||||
let o = cmd.output().unwrap();
|
||||
let s = cmd.status().unwrap();
|
||||
});
|
||||
}
|
||||
else {
|
||||
std::thread::sleep(std::time::Duration::from_millis(400));
|
||||
println!("pool: {pool:?}")
|
||||
}
|
||||
}
|
||||
}
|
|
@ -25,13 +25,10 @@ async fn main() -> Result<()> {
|
|||
#[cfg(feature = "server")]
|
||||
if cli.server {
|
||||
info!("starting server");
|
||||
return Server::build(cfg).await?.start().await;
|
||||
loop {
|
||||
// now we wait
|
||||
}
|
||||
return Server::build(cfg).await?.run().await;
|
||||
}
|
||||
// implicit else, so we can work without the server feature
|
||||
info!("starting client");
|
||||
todo!();
|
||||
Ok(())
|
||||
loop {}
|
||||
}
|
||||
|
|
|
@ -1,5 +1,5 @@
|
|||
#![cfg(feature = "server")]
|
||||
use std::time::Duration;
|
||||
use std::{time::Duration, sync::Arc};
|
||||
|
||||
use libpt::log::{debug, info, trace, warn};
|
||||
use tokio::{
|
||||
|
@ -18,46 +18,32 @@ pub struct Server {
|
|||
cfg: Config,
|
||||
pub timeout: Option<Duration>,
|
||||
server: TcpListener,
|
||||
accept_runtime: Runtime,
|
||||
request_runtime: Runtime,
|
||||
}
|
||||
|
||||
impl Server {
|
||||
pub async fn build(cfg: Config) -> anyhow::Result<Self> {
|
||||
let server = TcpListener::bind(cfg.addr).await?;
|
||||
let timeout = Some(Duration::from_secs(5));
|
||||
let accept_runtime = Builder::new_multi_thread()
|
||||
.worker_threads(1)
|
||||
.thread_stack_size(3 * 1024 * 1024)
|
||||
.enable_io()
|
||||
.enable_time()
|
||||
.build()?;
|
||||
let request_runtime = Builder::new_multi_thread()
|
||||
.worker_threads(1)
|
||||
.thread_stack_size(3 * 1024 * 1024)
|
||||
.enable_io()
|
||||
.enable_time()
|
||||
.build()?;
|
||||
Ok(Server {
|
||||
cfg,
|
||||
timeout,
|
||||
server,
|
||||
accept_runtime,
|
||||
request_runtime,
|
||||
})
|
||||
}
|
||||
pub async fn start(self) -> anyhow::Result<()> {
|
||||
self.accept_runtime.block_on(async {
|
||||
pub async fn run(self) -> anyhow::Result<()> {
|
||||
let rc_self = Arc::new(self);
|
||||
loop {
|
||||
let (stream, addr) = match self.server.accept().await {
|
||||
let rc_self = rc_self.clone();
|
||||
let (stream, addr) = match rc_self.server.accept().await {
|
||||
Ok(s) => s,
|
||||
Err(err) => {
|
||||
warn!("could not accept stream: {err:?}");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let _guard = self.request_runtime.enter();
|
||||
match self.handle_stream(stream).await {
|
||||
tokio::spawn(async move {
|
||||
|
||||
match rc_self.handle_stream(stream).await {
|
||||
Ok(_) => (),
|
||||
Err(err) => {
|
||||
match err {
|
||||
|
@ -68,12 +54,10 @@ impl Server {
|
|||
warn!("error while handling stream: {:?}", err)
|
||||
}
|
||||
}
|
||||
continue;
|
||||
}
|
||||
};
|
||||
}
|
||||
});
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_stream(&self, stream: TcpStream) -> Result<()> {
|
||||
|
|
Loading…
Reference in New Issue