去你妈的线程池,新增use_thread_pool参数

This commit is contained in:
Zihaoxu2008
2025-01-25 20:28:46 +08:00
parent b170844fc4
commit 043ea80cbd
4 changed files with 44 additions and 21 deletions

View File

@@ -4,6 +4,7 @@
"proxy_url": "", "proxy_url": "",
"username_prefix": "test", "username_prefix": "test",
"worker_threads": 100, "worker_threads": 100,
"use_thread_pool": false,
"count": 300, "count": 300,
"reco_words": ["验证程序已启用","AntiAttack"], "reco_words": ["验证程序已启用","AntiAttack"],
"login_commands": ["reg helloa helloa"], "login_commands": ["reg helloa helloa"],

View File

@@ -269,8 +269,9 @@ impl<T: StreamType> Bot<T> {
let mut login_command_flag = false; let mut login_command_flag = false;
let mut status = 0; let mut status = 0;
let mut iter = 0; let mut iter = 0;
let mut flag = true;
println!("[{}/WORKER] Worker started!", self.username); println!("[{}/WORKER] Worker started!", self.username);
loop { while flag {
if !status_rx.is_empty() { if !status_rx.is_empty() {
status = status_rx.recv().await.unwrap(); status = status_rx.recv().await.unwrap();
} }
@@ -320,7 +321,9 @@ impl<T: StreamType> Bot<T> {
iter += 1; iter += 1;
} }
sleep(Duration::from_millis(1)); sleep(Duration::from_millis(1));
flag = self.conn.lock().await.is_alive();
} }
println!("[{}/WORKER] Worker stopped.", self.username);
} }
// 运行bot // 运行bot

View File

@@ -14,7 +14,8 @@ pub struct Configuration {
spam_file: String, spam_file: String,
spam_delay_ms: i32, spam_delay_ms: i32,
random_prefix: bool, random_prefix: bool,
pub worker_threads: i32 pub worker_threads: i32,
pub use_thread_pool: bool
} }
impl Configuration { impl Configuration {
@@ -57,7 +58,8 @@ impl Configuration {
spam_file: "".to_string(), spam_file: "".to_string(),
spam_delay_ms: 0, spam_delay_ms: 0,
random_prefix: false, random_prefix: false,
worker_threads: 100 worker_threads: 100,
use_thread_pool: true
} }
} }

View File

@@ -9,7 +9,7 @@ use reqwest;
use serde_json::Value; use serde_json::Value;
// use tokio::task; // use tokio::task;
use std::sync::Arc; use std::sync::Arc;
// use std::thread; use std::thread;
use tokio::sync::Mutex; use tokio::sync::Mutex;
use crate::utils::generate_username; use crate::utils::generate_username;
use crate::configuration::Configuration; use crate::configuration::Configuration;
@@ -31,7 +31,7 @@ async fn main() {
// 一个互斥锁,防止有多个线程同时请求uuproxy // 一个互斥锁,防止有多个线程同时请求uuproxy
let thread_lock = Arc::new(Mutex::new(())); let thread_lock = Arc::new(Mutex::new(()));
let mut handles = Vec::new(); // let mut handles = Vec::new();
let configuration = Configuration::new("config.json"); let configuration = Configuration::new("config.json");
let configuration_rc = Arc::new(configuration); let configuration_rc = Arc::new(configuration);
let mut spam_file = None; let mut spam_file = None;
@@ -60,19 +60,23 @@ async fn main() {
// }); // });
// 生成一个线程 // 生成一个线程
// let handle = tokio::task::spawn(attack_thread(lock_clone,configuration_rc_clone,variables_clone)); // let handle = tokio::task::spawn(attack_thread(lock_clone,configuration_rc_clone,variables_clone));
let handle = rt.spawn(attack_thread(lock_clone,configuration_rc_clone,variables_clone)); // let handle;
// let handle = thread::spawn(move || { if configuration_rc.use_thread_pool {
// tokio::runtime::Builder::new_current_thread() rt.spawn(attack_thread(lock_clone,configuration_rc_clone,variables_clone));
// .enable_all() }else{
// .thread_stack_size(8 * 1024 * 1024) thread::spawn(move || {
// .build() tokio::runtime::Builder::new_current_thread()
// .unwrap() .enable_all()
// .block_on(async move { .thread_stack_size(8 * 1024 * 1024)
// attack_thread(lock_clone,configuration_rc_clone,variables_clone).await; .build()
// // println!("[MAIN] Thread {} created", i); .unwrap()
// }); .block_on(async move {
//}); attack_thread(lock_clone,configuration_rc_clone,variables_clone).await;
handles.push(handle); // println!("[MAIN] Thread {} created", i);
});
});
}
// handles.push(handle);
// println!("[MAIN] Thread {} created", i) // 这个太吵了( // println!("[MAIN] Thread {} created", i) // 这个太吵了(
} }
@@ -141,9 +145,22 @@ async fn attack_thread(lock: Arc<Mutex<()>>, config: Arc<Configuration>, variabl
// 插入生成update_worker线程的代码 // 插入生成update_worker线程的代码
let mut bot_clone = bot.clone(); let mut bot_clone = bot.clone();
let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<u8>(); let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<u8>();
if config.use_thread_pool {
tokio::spawn(async move { tokio::spawn(async move {
bot_clone.update_worker(rx).await; bot_clone.update_worker(rx).await;
}); });
}else {
thread::spawn(|| {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.thread_stack_size(8 * 1024 * 1024)
.build()
.unwrap()
.block_on(async move {
bot_clone.update_worker(rx).await
});
});
}
match bot.run(tx).await { match bot.run(tx).await {
Ok(_) => { Ok(_) => {
if bot.should_restart() { if bot.should_restart() {