// bot的实现 use std::{io::{Error, Read, Write}, net::TcpStream, sync::Arc, thread::{self, sleep}, time::{self, Duration}}; use log::{info, warn}; use rust_mc_proto::{DataBufferReader, DataBufferWriter, MinecraftConnection, Packet, ProtocolError}; use socks::Socks5Stream; use tokio::{io::{AsyncBufReadExt, AsyncSeekExt}, sync::Mutex}; use crate::{configuration::Configuration, utils}; // 让所有bot能够共享变量 pub struct BotVariable { pub spam_file: Arc>>, pub spam_cursor: Arc> } // Bot主线程和worker线程间的通信 pub struct BotMessage { pub state: i32, pub last_keepalive_sec: u64, } // 定义一个StreamType的特性,要求实现Read和Write,且含有connect函数 pub trait StreamType: Read + Write { fn connect(server_addr: &str,proxy_addr: &str) -> Result where Self: Sized; } // Bot泛型版 pub struct Bot { pub username: String, pub proxy_addr: String, pub server_addr: String, conn: Arc>>, should_restart: bool, config: Arc, status: i32, var: Arc, alive: Arc> } // 然后用一个Stream枚举来包装TcpStream和Socks5Stream(trait不能拿来创建bot) pub enum Stream { Socks5(Socks5Stream), Tcp(TcpStream), } //给Stream实现Read和Write impl Read for Stream { fn read(&mut self, buf: &mut [u8]) -> std::io::Result { match self { Stream::Socks5(t) => t.read(buf), Stream::Tcp(t) => t.read(buf), } } } impl Write for Stream { fn write(&mut self, buf: &[u8]) -> std::io::Result { match self { Stream::Socks5(t) => t.write(buf), Stream::Tcp(t) => t.write(buf), } } fn flush(&mut self) -> std::io::Result<()> { match self { Stream::Socks5(t) => t.flush(), Stream::Tcp(t) => t.flush(), } } } // 再让Stream实现StreamType,这样Stream就能拿来搞bot的泛型了 impl StreamType for Stream { fn connect(server_addr: &str,proxy_addr: &str) -> Result where Self: Sized { if !proxy_addr.is_empty() { Ok(Stream::Socks5(Socks5Stream::connect(proxy_addr, server_addr)?)) }else { Ok(Stream::Tcp(TcpStream::connect(server_addr)?)) } } } impl Clone for Bot { fn clone(&self) -> Self { Self { username: self.username.clone(), proxy_addr: self.proxy_addr.clone(), server_addr: self.server_addr.clone(), conn: self.conn.clone(), should_restart: self.should_restart.clone(), config: self.config.clone(), status: self.status.clone(), var: self.var.clone(), alive: self.alive.clone()} } } impl Bot { pub fn new(username: String, proxy_addr: String, server_addr: String, config: Arc, var: Arc) -> Result, Error> { info!("[{}] Creating bot on {} with proxy {}", username, server_addr, proxy_addr); let stream = T::connect(server_addr.as_str(), proxy_addr.as_str()); let should_restart = false; match stream { Ok(stream) => Ok(Bot { username, proxy_addr, server_addr, conn: Arc::new(Mutex::new(MinecraftConnection::new(stream))), should_restart, config, status: 0, var: var, alive: Arc::new(Mutex::new(true))}), Err(e) => Err(e), } } // bot登录函数 pub async fn login(&mut self) -> Result<(), ProtocolError> { // 这里其实传的server host是什么都不要紧,但保险起见还是弄一下 let mut server_host: String = self.server_addr.clone(); let server_port: String = server_host.split_off(server_host.find(":").unwrap()+1); // println!("{}",server_host); server_host = server_host.strip_suffix(':').unwrap().to_string(); let server_port_num: u16 = server_port.parse::().unwrap(); // 发一个handshake包,设置next_state为2(Login) self.conn.lock().await.write_packet(&Packet::build(0x00, |packet| { packet.write_u16_varint(763)?; // protocol_version packet.write_string(&server_host)?; // server_address packet.write_unsigned_short(server_port_num)?; // server_port packet.write_u8_varint(2) // next_state })?)?; // handshake packet // 再发一个,把用户名发过去(不知道为什么,反正文档上写要发两个) self.conn.lock().await.write_packet(&Packet::build(0x00, |packet| { packet.write_string(&self.username)?; packet.write_boolean(false)// ?; // packet.write_uuid(&Uuid::parse_str("550e8400-e29b-41d4-a716-446655440000").unwrap()) })?)?; // login start packet Ok(()) } pub async fn send_message(&mut self, message: String) -> Result<(), ProtocolError> { self.conn.lock().await.write_packet(&Packet::build(0x05, |packet| { packet.write_string(message.as_str())?; // message packet.write_long(time::SystemTime::now().duration_since(time::SystemTime::UNIX_EPOCH).unwrap().as_millis() as i64)?; // timestamp packet.write_long(0)?; // salt packet.write_boolean(false)?; // has signature packet.write_u8_varint(0)?; // message count packet.write_bytes(&[0,0,0]) })?)?; Ok(()) } pub async fn send_command(&mut self, command: String) -> Result<(), ProtocolError> { self.conn.lock().await.write_packet(&Packet::build(0x04, |packet| { packet.write_string(&command.as_str())?; packet.write_long(time::SystemTime::now().duration_since(time::SystemTime::UNIX_EPOCH).unwrap().as_millis() as i64)?; packet.write_long(0)?; packet.write_u8_varint(0)?; packet.write_u8_varint(0)?; // message count packet.write_bytes(&[0,0,0]) })?)?; Ok(()) } // 处理服务器发的keepalive pub async fn handle_keepalive(&mut self,packet: &mut Packet) -> Result<(), ProtocolError> { // keep alive let id = packet.read_long()?; // keepalive原封不动丢回去就行 // TODO 在20秒未进行Keepalive时断开连接 self.conn.lock().await.write_packet(&Packet::build(0x12, |packet| { // respond with a same keep alive packet packet.write_long(id) })?)?; // println!("Server keep alived"); // 这个包是重生用的,因为我发现如果玩家上线的时候就是死的,那它就没法自动重生 self.conn.lock().await.write_packet(&Packet::build(0x07, |packet| { packet.write_u8_varint(0) })?)?; Ok(()) } // 处理断开连接事件(这个指的是建立连接后再断开连接) pub async fn handle_disconnect(&mut self,packet: &mut Packet) -> Result<(), ProtocolError> { // disconnect let message = packet.read_string()?; let text = utils::parse_json_component(message.as_str()); info!("[{}] Server disconnected: {}", self.username, text); // if text.contains("验证程序已启用") { // TODO 把这个搞进配置文件里面 // println!("[{}] Restart flag setted, will restart after 1.5 min", self.username); // sleep(Duration::from_secs(90)); // self.should_restart = true; // } for i in self.config.reco_words(){ if text.contains(i.as_str()) { info!("[{}] Restart flag setted, will restart after 1.5 min", self.username); sleep(Duration::from_secs(90)); self.should_restart = true; } } Ok(()) } // 处理网络数据包 pub async fn handle_packets(&mut self, status_tx: tokio::sync::mpsc::UnboundedSender) -> Result<(), ProtocolError> { while self.conn.lock().await.is_alive() && *self.alive.lock().await { let mut packet = self.conn.lock().await.read_packet()?; let mut message = BotMessage { state: 0, last_keepalive_sec: 0 }; match packet.id() { 0x02 => { // 成功登录 // login success info!("[{}] Successfully logged in!",self.username); self.status = 1; message.state = 1; message.last_keepalive_sec = time::SystemTime::now().duration_since(time::SystemTime::UNIX_EPOCH).unwrap().as_secs(); status_tx.send(message).unwrap(); } 0x03 => {// 设置压缩CompressionThreshold // set compression let threshold = packet.read_i32_varint()?; if threshold >= 0 { self.conn.lock().await.set_compression(Some(threshold as usize)); // println!("[{}] Compression threshold set to {}", self.username, threshold) } } 0x23 => { // KeepAlive包 self.handle_keepalive(&mut packet).await?; message.last_keepalive_sec = time::SystemTime::now().duration_since(time::SystemTime::UNIX_EPOCH).unwrap().as_secs(); status_tx.send(message).unwrap(); // println!("{}",self.status); } 0x1A => {// 断连包 // disconnect self.handle_disconnect(&mut packet).await?; message.state = 2; status_tx.send(message).unwrap(); break; } 0x00 => {// 版本不匹配啥的就会走这 if self.status == 0{ let text = packet.read_string()?; warn!("[{}] Failed to login: {}", self.username, utils::parse_json_component(text.as_str())); message.state = 2; status_tx.send(message).unwrap(); break; } } 0x38 => {// 去世包(自动复活) // player died info!("[{}] Player died,respawning...",self.username); self.conn.lock().await.write_packet(&Packet::build(0x07, |packet| { packet.write_u8_varint(0) })?)? } _ => { } } } Ok(()) } // !这里的self不与主线程的self相同 pub async fn update_worker(&mut self, mut status_rx: tokio::sync::mpsc::UnboundedReceiver){ let mut login_timestamp = 0; let mut login_command_flag = false; let mut message = BotMessage { state: 0, last_keepalive_sec: 0 }; let mut iter = 0; let mut flag = true; info!("[{}/WORKER] Worker started!", self.username); while flag { if !status_rx.is_empty() { message = status_rx.recv().await.unwrap(); } if message.state == 1 && login_timestamp == 0 { // 刚刚登录完毕 login_timestamp = time::SystemTime::now().duration_since(time::SystemTime::UNIX_EPOCH).unwrap().as_millis() as i64; }else if message.state == 1 { // 已登录 if login_timestamp + 3000 <= time::SystemTime::now().duration_since(time::SystemTime::UNIX_EPOCH).unwrap().as_millis() as i64 && !login_command_flag { // 登录3s后且未执行过登录命令 let commands = self.config.login_commands().clone(); for i in commands{ self.send_command(i.to_string()).await.unwrap_or(()); thread::sleep(Duration::from_secs(1)); } info!("[{}/WORKER] Login commands executed!",self.username); login_command_flag = true; } } if login_command_flag { // 执行完登录命令后 if !self.config.spam_file().is_empty() && iter % self.config.spam_delay_ms() == 0{ let mut msg = String::new(); if self.config.random_prefix(){ msg += &utils::generate_string(2,10); } { // TODO 优化这里 let var = self.var.clone(); let mut file_lock = var.spam_file.lock().await; let file = file_lock.as_mut().unwrap(); let mut reader = tokio::io::BufReader::new(file); let mut line = String::new(); let mut flag = true; reader.seek(std::io::SeekFrom::Start(*var.spam_cursor.lock().await)).await.unwrap(); while flag { flag = false; if let Ok(size) = reader.read_line(&mut line).await { if size == 0 { reader.seek(std::io::SeekFrom::Start(0)).await.unwrap(); flag = true; }else { msg += line.trim(); } } } (*var.spam_cursor.lock().await) = reader.stream_position().await.unwrap(); } self.send_message(msg).await.unwrap_or(()); } let timestamp = time::SystemTime::now().duration_since(time::SystemTime::UNIX_EPOCH).unwrap().as_secs(); if message.last_keepalive_sec + 60 < timestamp { // 连接超时 (*self.alive.lock().await) = false; warn!("[{}/WORKER] Bot timed out", self.username); break; } iter += 1; } sleep(Duration::from_millis(1)); flag = self.conn.lock().await.is_alive() && message.state!=2; } warn!("[{}/WORKER] Worker stopped.", self.username); } // 运行bot pub async fn run(&mut self, status_tx: tokio::sync::mpsc::UnboundedSender) -> Result<(), ProtocolError> { info!("[{}] Connecting to {} with proxy {}", self.username, self.server_addr, self.proxy_addr); self.login().await.unwrap(); self.handle_packets(status_tx).await.unwrap_or(()); Ok(()) } pub fn should_restart(&self) -> bool { self.should_restart } } // 把创建bot的方法提取出来力 pub fn create_new_bot(username: String, proxy_addr: String, server_addr: String, config: Arc, var: Arc) -> Result, Error>{ if !proxy_addr.is_empty() { Bot::::new( username, proxy_addr, server_addr, config, var ) }else { Bot::::new( username, "".to_string(), server_addr, config, var ) } }