v1.1更新,详见changelogs
This commit is contained in:
@@ -1,31 +1,132 @@
|
||||
// bot的实现
|
||||
|
||||
use std::{io::Error, thread::sleep, time::Duration};
|
||||
use std::{io::{Error, Read, Write}, net::TcpStream, sync::Arc, thread::{self, sleep}, time::{self, Duration}};
|
||||
|
||||
use rust_mc_proto::{DataBufferReader, DataBufferWriter, MinecraftConnection, Packet, ProtocolError};
|
||||
use socks::Socks5Stream;
|
||||
use crate::utils;
|
||||
use tokio::{io::{AsyncBufReadExt, AsyncSeekExt}, sync::Mutex};
|
||||
use crate::{configuration::Configuration, utils};
|
||||
|
||||
pub struct Socks5Bot {
|
||||
username: String,
|
||||
proxy_addr: String,
|
||||
server_addr: String,
|
||||
conn: MinecraftConnection<Socks5Stream>,
|
||||
should_restart: bool,
|
||||
// 让所有bot能够共享变量
|
||||
pub struct BotVariable {
|
||||
pub spam_file: Arc<Mutex<Option<tokio::fs::File>>>,
|
||||
pub spam_cursor: Arc<Mutex<u64>>
|
||||
}
|
||||
|
||||
impl Socks5Bot {
|
||||
pub fn new(username: String, proxy_addr: String, server_addr: String) -> Result<Socks5Bot, Error> {
|
||||
// impl Clone for BotVariable {
|
||||
// fn clone(&self) -> Self {
|
||||
// Self { spam_file: self.spam_file.clone() }
|
||||
// }
|
||||
// }
|
||||
|
||||
// 定义一个StreamType的特性,要求实现Read和Write,且含有connect函数
|
||||
pub trait StreamType: Read + Write {
|
||||
fn connect(server_addr: &str,proxy_addr: &str) -> Result<Self, Error> where Self: Sized;
|
||||
}
|
||||
|
||||
// *由于下面的Stream已经实现过connect了,所以这里没必要重复一遍
|
||||
// // 然后让Socks5Stream实现StreamType
|
||||
// impl StreamType for Socks5Stream{
|
||||
// fn connect(server_addr: &str,proxy_addr: &str) -> Result<Self, Error> where Self: Sized {
|
||||
// return Socks5Stream::connect(proxy_addr, server_addr);
|
||||
// }
|
||||
// }
|
||||
// // 再让TcpStream实现StreamType
|
||||
// impl StreamType for TcpStream {
|
||||
// fn connect(server_addr: &str,_: &str) -> Result<Self, Error> where Self: Sized {
|
||||
// return TcpStream::connect(server_addr);
|
||||
// }
|
||||
// }
|
||||
|
||||
// Bot泛型版
|
||||
pub struct Bot<T: StreamType> {
|
||||
pub username: String,
|
||||
pub proxy_addr: String,
|
||||
pub server_addr: String,
|
||||
conn: Arc<Mutex<MinecraftConnection<T>>>,
|
||||
should_restart: bool,
|
||||
config: Arc<Configuration>,
|
||||
status: i32,
|
||||
var: Arc<BotVariable>,
|
||||
}
|
||||
|
||||
// 然后用一个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<usize> {
|
||||
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<usize> {
|
||||
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<Self, Error> 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<Stream> {
|
||||
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()}
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: StreamType> Bot<T> {
|
||||
pub fn new(username: String, proxy_addr: String, server_addr: String, config: Arc<Configuration>, var: Arc<BotVariable>) -> Result<Bot<T>, Error> {
|
||||
println!("[{}] Creating bot on {} with proxy {}", username, server_addr, proxy_addr);
|
||||
let stream = Socks5Stream::connect(proxy_addr.as_str(), server_addr.as_str());
|
||||
let stream = T::connect(server_addr.as_str(), proxy_addr.as_str());
|
||||
let should_restart = false;
|
||||
match stream {
|
||||
Ok(stream) => Ok(Socks5Bot { username, proxy_addr, server_addr, conn: MinecraftConnection::new(stream), should_restart }),
|
||||
Ok(stream) => Ok(Bot {
|
||||
username,
|
||||
proxy_addr,
|
||||
server_addr,
|
||||
conn: Arc::new(Mutex::new(MinecraftConnection::new(stream))),
|
||||
should_restart,
|
||||
config,
|
||||
status: 0,
|
||||
var: var}),
|
||||
Err(e) => Err(e),
|
||||
}
|
||||
}
|
||||
// bot登录函数
|
||||
pub fn login(&mut self) -> Result<(), ProtocolError> {
|
||||
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);
|
||||
@@ -33,7 +134,7 @@ impl Socks5Bot {
|
||||
server_host = server_host.strip_suffix(':').unwrap().to_string();
|
||||
let server_port_num: u16 = server_port.parse::<u16>().unwrap();
|
||||
// 发一个handshake包,设置next_state为2(Login)
|
||||
self.conn.write_packet(&Packet::build(0x00, |packet| {
|
||||
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
|
||||
@@ -41,7 +142,7 @@ impl Socks5Bot {
|
||||
})?)?; // handshake packet
|
||||
|
||||
// 再发一个,把用户名发过去(不知道为什么,反正文档上写要发两个)
|
||||
self.conn.write_packet(&Packet::build(0x00, |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())
|
||||
@@ -50,67 +151,98 @@ impl Socks5Bot {
|
||||
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 fn handle_keepalive(&mut self,packet: &mut Packet) -> Result<(), ProtocolError> {
|
||||
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.write_packet(&Packet::build(0x12, |packet| { // respond with a same keep alive packet
|
||||
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.write_packet(&Packet::build(0x07, |packet| {
|
||||
self.conn.lock().await.write_packet(&Packet::build(0x07, |packet| {
|
||||
packet.write_u8_varint(0)
|
||||
})?)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// 处理断开连接事件(这个指的是建立连接后再断开连接)
|
||||
pub fn handle_disconnect(&mut self,packet: &mut Packet) -> Result<(), ProtocolError> {
|
||||
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());
|
||||
println!("[{}] 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;
|
||||
// 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()) {
|
||||
println!("[{}] Restart flag setted, will restart after 1.5 min", self.username);
|
||||
sleep(Duration::from_secs(90));
|
||||
self.should_restart = true;
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// 处理网络数据包
|
||||
pub fn handle_packets(&mut self) -> Result<(), ProtocolError> {
|
||||
let mut status = 0;
|
||||
|
||||
while self.conn.is_alive() {
|
||||
let mut packet = self.conn.read_packet()?;
|
||||
pub async fn handle_packets(&mut self, status_tx: tokio::sync::mpsc::UnboundedSender<u8>) -> Result<(), ProtocolError> {
|
||||
while self.conn.lock().await.is_alive() {
|
||||
let mut packet = self.conn.lock().await.read_packet()?;
|
||||
match packet.id() {
|
||||
0x02 => { // 成功登录
|
||||
// login success
|
||||
println!("[{}] Successfully logged in!",self.username);
|
||||
status = 1;
|
||||
self.status = 1;
|
||||
status_tx.send(1).unwrap();
|
||||
}
|
||||
0x03 => {// 设置压缩CompressionThreshold
|
||||
// set compression
|
||||
let threshold = packet.read_i32_varint()?;
|
||||
if threshold >= 0 {
|
||||
self.conn.set_compression(Some(threshold as usize));
|
||||
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)?;
|
||||
self.handle_keepalive(&mut packet).await?;
|
||||
// println!("{}",self.status);
|
||||
}
|
||||
0x1A => {// 断连包
|
||||
// disconnect
|
||||
self.handle_disconnect(&mut packet)?;
|
||||
self.handle_disconnect(&mut packet).await?;
|
||||
break;
|
||||
}
|
||||
0x00 => {// 版本不匹配啥的就会走这
|
||||
if status == 0{
|
||||
if self.status == 0{
|
||||
let text = packet.read_string()?;
|
||||
println!("[{}] Failed to login: {}", self.username, utils::parse_json_component(text.as_str()));
|
||||
break;
|
||||
@@ -119,7 +251,7 @@ impl Socks5Bot {
|
||||
0x38 => {// 去世包(自动复活)
|
||||
// player died
|
||||
println!("[{}] Player died,respawning...",self.username);
|
||||
self.conn.write_packet(&Packet::build(0x07, |packet| {
|
||||
self.conn.lock().await.write_packet(&Packet::build(0x07, |packet| {
|
||||
packet.write_u8_varint(0)
|
||||
})?)?
|
||||
}
|
||||
@@ -131,15 +263,96 @@ impl Socks5Bot {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// !这里的self不与主线程的self相同
|
||||
pub async fn update_worker(&mut self, mut status_rx: tokio::sync::mpsc::UnboundedReceiver<u8>){
|
||||
let mut login_timestamp = 0;
|
||||
let mut login_command_flag = false;
|
||||
let mut status = 0;
|
||||
let mut iter = 0;
|
||||
println!("[{}/WORKER] Worker started!", self.username);
|
||||
loop {
|
||||
if !status_rx.is_empty() {
|
||||
status = status_rx.recv().await.unwrap();
|
||||
}
|
||||
if status == 1 && login_timestamp == 0 { // 刚刚登录完毕
|
||||
login_timestamp = time::SystemTime::now().duration_since(time::SystemTime::UNIX_EPOCH).unwrap().as_millis() as i64;
|
||||
}else if status == 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();
|
||||
thread::sleep(Duration::from_secs(1));
|
||||
}
|
||||
println!("[{}/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();
|
||||
}
|
||||
{ // 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();
|
||||
}
|
||||
iter += 1;
|
||||
}
|
||||
sleep(Duration::from_millis(1));
|
||||
}
|
||||
}
|
||||
|
||||
// 运行bot
|
||||
pub fn run(&mut self) -> Result<(), ProtocolError> {
|
||||
println!("[{}] Connecting to {} with proxy {}", self.username, self.server_addr,self.proxy_addr);
|
||||
self.login().unwrap();
|
||||
self.handle_packets().unwrap_or(());
|
||||
pub async fn run(&mut self, status_tx: tokio::sync::mpsc::UnboundedSender<u8>) -> Result<(), ProtocolError> {
|
||||
println!("[{}] 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<Configuration>, var: Arc<BotVariable>) -> Result<Bot<Stream>, Error>{
|
||||
if !proxy_addr.is_empty() {
|
||||
Bot::<Stream>::new(
|
||||
username,
|
||||
proxy_addr,
|
||||
server_addr,
|
||||
config,
|
||||
var
|
||||
)
|
||||
}else {
|
||||
Bot::<Stream>::new(
|
||||
username,
|
||||
"".to_string(),
|
||||
server_addr,
|
||||
config,
|
||||
var
|
||||
)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user