- add: handle channel list request
- add: handle server select request - add: handle character list request stub - add: start health check function for consul
This commit is contained in:
@@ -1,26 +1,30 @@
|
||||
use std::collections::HashMap;
|
||||
use std::env;
|
||||
use tonic::{Code, Status};
|
||||
use std::error::Error;
|
||||
use std::sync::Arc;
|
||||
use tokio::net::TcpStream;
|
||||
use tokio::sync::Mutex;
|
||||
use tracing::{debug, error, info, warn};
|
||||
use utils::service_discovery;
|
||||
use crate::auth_client::AuthClient;
|
||||
use crate::handlers::null_string::NullTerminatedString;
|
||||
use crate::packet::{send_packet, Packet, PacketPayload};
|
||||
use crate::packet_type::PacketType;
|
||||
use crate::packets::cli_accept_req::CliAcceptReq;
|
||||
use crate::packets::cli_channel_list_req::CliChannelListReq;
|
||||
use crate::packets::cli_join_server_req::CliJoinServerReq;
|
||||
use crate::packets::cli_login_req::CliLoginReq;
|
||||
use crate::packets::{srv_accept_reply, srv_login_reply};
|
||||
use crate::packets::cli_srv_select_req::CliSrvSelectReq;
|
||||
use crate::packets::srv_accept_reply::SrvAcceptReply;
|
||||
use crate::packets::srv_channel_list_reply::{ChannelInfo, SrvChannelListReply};
|
||||
use crate::packets::srv_login_reply::{ServerInfo, SrvLoginReply};
|
||||
use crate::packets::srv_srv_select_reply::SrvSrvSelectReply;
|
||||
use crate::packets::*;
|
||||
use std::collections::HashMap;
|
||||
use std::env;
|
||||
use std::error::Error;
|
||||
use std::sync::Arc;
|
||||
use tokio::net::TcpStream;
|
||||
use tokio::sync::Mutex;
|
||||
use tonic::{Code, Status};
|
||||
use tracing::{debug, error, info, warn};
|
||||
use utils::service_discovery;
|
||||
|
||||
pub(crate) async fn handle_accept_req(stream: &mut TcpStream, packet: Packet) -> Result<(), Box<dyn Error + Send + Sync>> {
|
||||
let data = CliAcceptReq::decode(packet.payload.as_slice());
|
||||
debug!("{:?}", data);
|
||||
let request = CliAcceptReq::decode(packet.payload.as_slice());
|
||||
debug!("{:?}", request);
|
||||
|
||||
// We need to do reply to this packet
|
||||
let data = SrvAcceptReply { result: srv_accept_reply::Result::Accepted, rand_value: 0 };
|
||||
@@ -32,8 +36,8 @@ pub(crate) async fn handle_accept_req(stream: &mut TcpStream, packet: Packet) ->
|
||||
}
|
||||
|
||||
pub(crate) async fn handle_join_server_req(stream: &mut TcpStream, packet: Packet) -> Result<(), Box<dyn Error + Send + Sync>> {
|
||||
let data = CliJoinServerReq::decode(packet.payload.as_slice());
|
||||
debug!("{:?}", data);
|
||||
let request = CliJoinServerReq::decode(packet.payload.as_slice());
|
||||
debug!("{:?}", request);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -53,29 +57,35 @@ pub(crate) async fn handle_login_req(stream: &mut TcpStream, packet: Packet, aut
|
||||
send_packet(stream, &response_packet).await?;
|
||||
} else {
|
||||
debug!("Successfully logged in");
|
||||
|
||||
|
||||
let consul_url = env::var("CONSUL_URL").unwrap_or_else(|_| "http://127.0.0.1:8500".to_string());
|
||||
let servers = service_discovery::get_service_address(&consul_url, "character-service").await.unwrap_or_else(|err| {
|
||||
warn!(err);
|
||||
Vec::new()
|
||||
});
|
||||
|
||||
|
||||
if servers.len() == 0 {
|
||||
let data = SrvLoginReply { result: srv_login_reply::Result::Failed, right: 0, type_: 0, servers_info: Vec::new() };
|
||||
let response_packet = Packet::new(PacketType::PaklcLoginReply, &data)?;
|
||||
send_packet(stream, &response_packet).await?;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
|
||||
let mut server_info: Vec<ServerInfo> = Vec::new();
|
||||
let mut id = 0;
|
||||
for server in servers {
|
||||
let name = server.ServiceMeta.get("name").unwrap_or(&"".to_string()).clone();
|
||||
let mut name = server.ServiceMeta.get("name").unwrap_or(&"".to_string()).clone();
|
||||
let is_test = server.ServiceTags.contains(&"test".to_string()) || server.ServiceTags.contains(&"staging".to_string());
|
||||
if is_test {
|
||||
name = format!("@{}", name);
|
||||
} else {
|
||||
name = format!(" {}", name);
|
||||
}
|
||||
server_info.push(ServerInfo { test: u8::from(is_test), name: NullTerminatedString::new(&name), id});
|
||||
id = id + 1;
|
||||
}
|
||||
|
||||
debug!("Server info: {:?}", server_info);
|
||||
|
||||
let data = SrvLoginReply { result: srv_login_reply::Result::Ok, right: 0, type_: 0, servers_info: server_info };
|
||||
let response_packet = Packet::new(PacketType::PaklcLoginReply, &data)?;
|
||||
send_packet(stream, &response_packet).await?;
|
||||
@@ -112,13 +122,51 @@ pub(crate) async fn handle_login_req(stream: &mut TcpStream, packet: Packet, aut
|
||||
}
|
||||
|
||||
pub(crate) async fn handle_server_select_req(stream: &mut TcpStream, packet: Packet) -> Result<(), Box<dyn Error + Send + Sync>> {
|
||||
let data = CliJoinServerReq::decode(packet.payload.as_slice());
|
||||
debug!("{:?}", data);
|
||||
let request = CliSrvSelectReq::decode(packet.payload.as_slice());
|
||||
debug!("{:?}", request);
|
||||
|
||||
let data = SrvSrvSelectReply {
|
||||
result: srv_srv_select_reply::Result::Failed,
|
||||
session_id: 0,
|
||||
crypt_val: 0,
|
||||
ip: NullTerminatedString::new(""),
|
||||
port: 0,
|
||||
};
|
||||
|
||||
let response_packet = Packet::new(PacketType::PaklcSrvSelectReply, &data)?;
|
||||
send_packet(stream, &response_packet).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn handle_channel_list_req(stream: &mut TcpStream, packet: Packet) -> Result<(), Box<dyn Error + Send + Sync>> {
|
||||
let data = CliJoinServerReq::decode(packet.payload.as_slice());
|
||||
debug!("{:?}", data);
|
||||
let request = CliChannelListReq::decode(packet.payload.as_slice());
|
||||
debug!("{:?}", request);
|
||||
|
||||
let consul_url = env::var("CONSUL_URL").unwrap_or_else(|_| "http://127.0.0.1:8500".to_string());
|
||||
let channels = service_discovery::get_service_address(&consul_url, "world-service").await.unwrap_or_else(|err| {
|
||||
warn!(err);
|
||||
Vec::new()
|
||||
});
|
||||
|
||||
if channels.len() == 0 {
|
||||
let data = SrvChannelListReply { id: request?.server_id, channels: Vec::new() };
|
||||
let response_packet = Packet::new(PacketType::PaklcChannelListReply, &data)?;
|
||||
send_packet(stream, &response_packet).await?;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
debug!("Server info: {:?}", channels);
|
||||
let mut channel_info: Vec<ChannelInfo> = Vec::new();
|
||||
let mut id = 0;
|
||||
for channel in channels {
|
||||
let name = format!("{}", channel.ServiceMeta.get("name").unwrap_or(&"".to_string()).clone());
|
||||
channel_info.push(ChannelInfo { id: id, low_age: 0, high_age: 0, capacity: 0, name: NullTerminatedString::new(&name) });
|
||||
id = id + 1;
|
||||
}
|
||||
debug!("Channel info: {:?}", channel_info);
|
||||
|
||||
let data = SrvChannelListReply { id: request?.server_id, channels: channel_info };
|
||||
let response_packet = Packet::new(PacketType::PaklcChannelListReply, &data)?;
|
||||
send_packet(stream, &response_packet).await?;
|
||||
Ok(())
|
||||
}
|
||||
Reference in New Issue
Block a user