Moved RaaS handling to its own service
This commit is contained in:
+1
-5
@@ -1,7 +1,6 @@
|
||||
use crate::album_manager::AlbumManager;
|
||||
use crate::error::Error;
|
||||
use crate::migrations::{CURRENT_DB_VERSION, do_migration};
|
||||
use crate::rass::RaaSHandler;
|
||||
use config::{Config, File};
|
||||
use j_db::database::Database;
|
||||
use reqwest::Url;
|
||||
@@ -44,6 +43,7 @@ pub struct BotConfig {
|
||||
pub announcement_channel: ChannelId,
|
||||
pub toys: Vec<String>,
|
||||
pub effect_role_duration: i64,
|
||||
pub raas_server: String,
|
||||
pub picox: PicOxConfig,
|
||||
}
|
||||
|
||||
@@ -60,18 +60,14 @@ impl BotConfig {
|
||||
#[derive(Debug)]
|
||||
pub struct BotState {
|
||||
pub accepted_nsfw: Option<UserId>,
|
||||
pub bad_apple_running: bool,
|
||||
pub speak_lock: Mutex<()>,
|
||||
pub raas_handler: Option<RaaSHandler>,
|
||||
}
|
||||
|
||||
impl BotState {
|
||||
pub async fn new() -> Result<Self, Error> {
|
||||
Ok(Self {
|
||||
accepted_nsfw: None,
|
||||
bad_apple_running: false,
|
||||
speak_lock: Mutex::new(()),
|
||||
raas_handler: None,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
+44
-24
@@ -2,7 +2,6 @@ use crate::album_manager::{Album, AlbumQuery, ImageQuery, ImageSort};
|
||||
use crate::error::Error;
|
||||
use crate::models::insult_compliment::{RandomResponseTemplate, ResponseType};
|
||||
use crate::models::random::{Random, RandomConfig};
|
||||
use crate::rass::RaaSCmd;
|
||||
use crate::{command, group, GlobalData, BAD_APPLE};
|
||||
use reqwest::Client;
|
||||
use serde::{Deserialize, Serialize};
|
||||
@@ -15,6 +14,10 @@ use serenity::model::channel::Message;
|
||||
use serenity::utils::MessageBuilder;
|
||||
use std::collections::HashMap;
|
||||
use std::time::Duration;
|
||||
use chrono::Utc;
|
||||
use raas_types::raas::bot::roll::{Roll, roll_response, RollCmd};
|
||||
use raas_types::raas::resp::response::Resp;
|
||||
use raas_types::raas::service::raas_client::RaasClient;
|
||||
|
||||
#[group]
|
||||
#[commands(dad_joke, roll, real_roll, bad_apple, insult, add_random, list_random)]
|
||||
@@ -244,32 +247,49 @@ async fn roll(ctx: &Context, msg: &Message, args: Args) -> CommandResult {
|
||||
#[aliases("real_roll")]
|
||||
#[description("Roll a real die!")]
|
||||
async fn real_roll(ctx: &Context, msg: &Message, _args: Args) -> CommandResult {
|
||||
let mut data = ctx.data.write().await;
|
||||
let global = data.get_mut::<GlobalData>().unwrap();
|
||||
let data = ctx.data.read().await;
|
||||
let global = data.get::<GlobalData>().unwrap();
|
||||
|
||||
let addr = global.cfg.raas_server.clone();
|
||||
|
||||
if let Some(raas_handler) = &mut global.bot_state.raas_handler {
|
||||
raas_handler.send_msg_queue.send(RaaSCmd::Roll(3)).await?;
|
||||
msg.reply(&ctx.http, "Sent request to Roll Bot...").await?;
|
||||
if let Some(img) = raas_handler.recv_msg_queue.recv().await {
|
||||
match img {
|
||||
RaaSCmd::Roll(_) => {}
|
||||
RaaSCmd::Img(img_data) => {
|
||||
msg.channel_id
|
||||
.send_message(
|
||||
&ctx.http,
|
||||
CreateMessage::new()
|
||||
.content("Your roll my friend, hope its good I can't read!")
|
||||
.add_file(CreateAttachment::bytes(img_data, "roll.jpg")),
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
let mut client = RaasClient::connect(addr).await?;
|
||||
|
||||
let rolls = 3;
|
||||
|
||||
let roll_request = raas_types::raas::cmd::Request {
|
||||
timestamp: Utc::now().timestamp() as u64,
|
||||
cmd: Some(raas_types::raas::cmd::request::Cmd::RollCmd(RollCmd {
|
||||
cmd: Some(raas_types::raas::bot::roll::roll_cmd::Cmd::Roll(Roll {
|
||||
rotations: rolls,
|
||||
})),
|
||||
})),
|
||||
};
|
||||
|
||||
let grpc_request = tonic::Request::new(roll_request);
|
||||
|
||||
|
||||
let response = client.send_request(grpc_request).await?;
|
||||
|
||||
let raas_response = response.into_inner();
|
||||
|
||||
println!("Got resp: {} @ {}",raas_response.id, raas_response.timestamp);
|
||||
|
||||
match raas_response.resp.unwrap() {
|
||||
Resp::RollResp(roll_resp) => {
|
||||
if let roll_response::Response::RollImage(img) = roll_resp.response.unwrap() {
|
||||
msg.channel_id
|
||||
.send_message(
|
||||
&ctx.http,
|
||||
CreateMessage::new()
|
||||
.content("Your roll my friend, hope its good I can't read!")
|
||||
.add_file(CreateAttachment::bytes(img.img, "roll.jpg")),
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
} else {
|
||||
msg.reply(&ctx.http, "Roll Bot is gone... oh god!").await?;
|
||||
}
|
||||
} else {
|
||||
msg.reply(&ctx.http, "Looks like I can't reach my real flesh")
|
||||
.await?;
|
||||
Resp::Error(err) => {
|
||||
msg.reply(&ctx.http, format!("My real flesh encountered an error. Get Dad to fix it. `{}`", err.msg)).await?;
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
|
||||
@@ -17,7 +17,6 @@ use crate::discord::fren_coin::give_coin;
|
||||
use crate::discord::joke::random;
|
||||
use crate::models::lil_fren::lil_fren_task;
|
||||
use crate::models::task::Task;
|
||||
use crate::rass::{RaaS, RaaSCmd, RaaSHandler};
|
||||
use crate::{help, hook, GlobalData};
|
||||
use rand::prelude::IteratorRandom;
|
||||
use rand::thread_rng;
|
||||
@@ -49,29 +48,6 @@ impl EventHandler for Handler {
|
||||
lil_fren_task(&ctx1).await;
|
||||
}
|
||||
});
|
||||
let ctx2 = ctx.clone();
|
||||
tokio::spawn(async move {
|
||||
let (tx_client, rx_server) = tokio::sync::mpsc::channel::<RaaSCmd>(10);
|
||||
let (tx_server, rx_client) = tokio::sync::mpsc::channel::<RaaSCmd>(10);
|
||||
|
||||
let handler = RaaSHandler {
|
||||
recv_msg_queue: rx_client,
|
||||
send_msg_queue: tx_client,
|
||||
};
|
||||
|
||||
let mut raas = RaaS {
|
||||
recv_msg_queue: rx_server,
|
||||
send_msg_queue: tx_server,
|
||||
};
|
||||
|
||||
{
|
||||
let mut data = ctx2.data.write().await;
|
||||
let global_data = data.get_mut::<GlobalData>().unwrap();
|
||||
global_data.bot_state.raas_handler = Some(handler);
|
||||
}
|
||||
|
||||
raas.worker("0.0.0.0:50000").await.unwrap();
|
||||
});
|
||||
|
||||
tokio::spawn(async move {
|
||||
Task::create_reoccurring_tasks(&ctx).await.unwrap();
|
||||
|
||||
@@ -6,7 +6,6 @@ mod error;
|
||||
mod inventory;
|
||||
mod migrations;
|
||||
mod models;
|
||||
mod rass;
|
||||
mod user;
|
||||
|
||||
use crate::config::{Args, BotConfig, Channel, GlobalData};
|
||||
|
||||
-128
@@ -1,128 +0,0 @@
|
||||
use prost::Message;
|
||||
use raas_types::raas;
|
||||
use raas_types::raas::bot::roll::Roll;
|
||||
use raas_types::raas::cmd::Command;
|
||||
use raas_types::raas::register::{Register, RegisterResponse};
|
||||
use raas_types::raas::resp::response::Resp;
|
||||
use raas_types::raas::resp::Response;
|
||||
use raas_types::raas::RaasMessage;
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
use tokio::net::{TcpListener, TcpStream};
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum RaaSCmd {
|
||||
Roll(u32),
|
||||
Img(Vec<u8>),
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct RaaSHandler {
|
||||
pub recv_msg_queue: tokio::sync::mpsc::Receiver<RaaSCmd>,
|
||||
pub send_msg_queue: tokio::sync::mpsc::Sender<RaaSCmd>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct RaaS {
|
||||
pub recv_msg_queue: tokio::sync::mpsc::Receiver<RaaSCmd>,
|
||||
pub send_msg_queue: tokio::sync::mpsc::Sender<RaaSCmd>,
|
||||
}
|
||||
|
||||
impl RaaS {
|
||||
async fn receive_packet(socket: &mut TcpStream) -> Result<RaasMessage, std::io::Error> {
|
||||
let mut message_len_bytes = [0u8; 4];
|
||||
socket.read_exact(&mut message_len_bytes).await?;
|
||||
|
||||
let len = u32::from_be_bytes(message_len_bytes);
|
||||
|
||||
let mut message = vec![0u8; len as usize];
|
||||
|
||||
socket.read_exact(&mut message).await?;
|
||||
|
||||
Ok(RaasMessage { len, msg: message })
|
||||
}
|
||||
|
||||
async fn send_packet(socket: &mut TcpStream, data: Vec<u8>) -> Result<(), std::io::Error> {
|
||||
let msg = RaasMessage::new(data);
|
||||
|
||||
socket.write_all(&msg.into_bytes()).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn handle_register(
|
||||
&self,
|
||||
tcp_listener: TcpListener,
|
||||
) -> Result<TcpStream, std::io::Error> {
|
||||
let (mut socket, _) = tcp_listener.accept().await?;
|
||||
|
||||
let register_msg = Self::receive_packet(&mut socket).await?;
|
||||
|
||||
let register = Register::decode(&*register_msg.msg).unwrap();
|
||||
|
||||
let register_resp = RegisterResponse {
|
||||
name: register.name.to_string(),
|
||||
r#type: register.bot_type,
|
||||
id: 0,
|
||||
};
|
||||
|
||||
let mut msg = Vec::new();
|
||||
|
||||
register_resp.encode(&mut msg).unwrap();
|
||||
|
||||
Self::send_packet(&mut socket, msg).await?;
|
||||
|
||||
Ok(socket)
|
||||
}
|
||||
|
||||
pub async fn worker(&mut self, addr: &str) -> Result<(), std::io::Error> {
|
||||
let tcp_listener = TcpListener::bind(addr).await.unwrap();
|
||||
println!("Waiting for Rollbot to connect to {}...", addr);
|
||||
let mut stream = self.handle_register(tcp_listener).await?;
|
||||
println!("Rollbot has entered the chat!");
|
||||
|
||||
loop {
|
||||
if self.recv_msg_queue.recv().await.is_none() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let timestamp = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap()
|
||||
.as_secs();
|
||||
|
||||
let roll = Roll { rotations: 3 };
|
||||
|
||||
let cmd = Command {
|
||||
id: 0,
|
||||
timestamp,
|
||||
cmd: Some(raas::cmd::command::Cmd::RollCmd(raas::bot::roll::RollCmd {
|
||||
cmd: Some(raas::bot::roll::roll_cmd::Cmd::Roll(roll)),
|
||||
})),
|
||||
};
|
||||
|
||||
let mut msg = Vec::new();
|
||||
|
||||
cmd.encode(&mut msg).unwrap();
|
||||
Self::send_packet(&mut stream, msg).await?;
|
||||
|
||||
let recv = Self::receive_packet(&mut stream).await?;
|
||||
|
||||
let resp = Response::decode(&*recv.msg).unwrap();
|
||||
|
||||
println!("Got {} {}", resp.id, resp.timestamp);
|
||||
|
||||
match resp.resp.unwrap() {
|
||||
Resp::RollResp(roll) => match roll.response.unwrap() {
|
||||
raas::bot::roll::roll_response::Response::Pong(_) => {}
|
||||
raas::bot::roll::roll_response::Response::RollImage(img) => {
|
||||
println!("Got img!");
|
||||
self.send_msg_queue
|
||||
.send(RaaSCmd::Img(img.img))
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
},
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user