diff --git a/src/kube_cache.rs b/src/kube_cache.rs index bba747b..d0f5eba 100644 --- a/src/kube_cache.rs +++ b/src/kube_cache.rs @@ -33,7 +33,7 @@ pub struct KubeCache { dep_cache: kube::runtime::reflector::Store, dep_api: Api, srv_api: Api, - in_cluster: bool, + running_in_cluster: bool, } impl KubeCache { /// This initializes the creation of a "kubernetes client" @@ -86,7 +86,7 @@ impl KubeCache { dep_cache: dep_reader, dep_api: dep_api, srv_api, - in_cluster, + running_in_cluster: in_cluster, }, handle, )); @@ -117,7 +117,7 @@ impl KubeCache { Some(t) => t == "ClusterIP", None => false, }; - let incorrect_type = in_cluster ^ self.in_cluster; + let incorrect_type = in_cluster ^ self.running_in_cluster; !incorrect_type && filter_label_value(x, addr, port) }) } @@ -164,7 +164,7 @@ impl MinecraftAPI for McApi { } }; - let internal_id = match deployment.metadata.name.clone() { + let resource_name = match deployment.metadata.name.clone() { Some(x) => x, None => { return Err(OpaqueError::create_with_kind( @@ -183,14 +183,8 @@ impl MinecraftAPI for McApi { .ok_or(OpaqueError::create( "Could not find \"mc-router\" nodePort for server", ))?; - let inter_addr = match self.cache.in_cluster { - false => { - let node_port = port_srv - .node_port - .map(|x| x.to_string()) - .ok_or(OpaqueError::create("Could not map nodePort to port string"))?; - format!("localhost:{}", node_port) - } + + let inter_addr = match self.cache.running_in_cluster { true => { let target_port = port_srv.port; format!( @@ -201,12 +195,19 @@ impl MinecraftAPI for McApi { target_port ) } + false => { + let node_port = port_srv + .node_port + .map(|x| x.to_string()) + .ok_or(OpaqueError::create("Could not map nodePort to port string"))?; + format!("localhost:{}", node_port) + } }; return Ok(Server { server_addr: addr.to_string(), server_port: port.to_string(), internal_address: inter_addr, - internal_id, + resource_name, cache: self.cache.clone(), }); } @@ -236,13 +237,13 @@ pub struct Server { server_addr: String, server_port: String, internal_address: String, - internal_id: String, + resource_name: String, cache: KubeCache, } impl fmt::Debug for Server { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.debug_struct("KubeServer") - .field("internal_id", &self.internal_id) + .field("internal_id", &self.resource_name) .field( "join_addr", &format!("{}:{}", self.server_addr, self.server_port), @@ -268,7 +269,7 @@ impl MinecraftServerHandle for Server { async fn query_status(&self) -> Result { let dep = self .cache - .get_dep(&self.internal_id) + .get_dep(&self.resource_name) .ok_or_else(|| OpaqueError::create("Failed to get deployment from cache"))?; let status = match &dep.status { @@ -319,7 +320,7 @@ impl MinecraftServerHandle for Server { self.internal_address.as_str() } fn get_internal_id(&self) -> &str { - self.internal_id.as_str() + self.resource_name.as_str() } fn get_join_addr(&self) -> &str { @@ -331,7 +332,7 @@ impl MinecraftServerHandle for Server { } fn get_motd(&self) -> Option { - let dep = self.cache.get_dep(&self.internal_id)?; + let dep = self.cache.get_dep(&self.resource_name)?; let all_container_motds = dep .spec @@ -357,7 +358,7 @@ impl MinecraftServerHandle for Server { impl Server { async fn set_scale(&self, num: i32) -> Result<(), kube::Error> { - let _res = self.cache.set_dep_scale(&self.internal_id, num).await?; + let _res = self.cache.set_dep_scale(&self.resource_name, num).await?; tracing::info!( "scaled replicas of {}:{} to {num}", self.server_addr, diff --git a/src/mc_server/helpers.rs b/src/mc_server/helpers.rs new file mode 100644 index 0000000..277c5bf --- /dev/null +++ b/src/mc_server/helpers.rs @@ -0,0 +1,57 @@ +use tokio::{io::AsyncWriteExt, net::TcpStream}; + +use crate::{ + packets::{clientbound::status::StatusStructNew, Packet, SendPacket}, + OpaqueError, +}; + +#[tracing::instrument(skip(client_stream))] +pub async fn handle_ping(client_stream: &mut TcpStream) -> Result<(), OpaqueError> { + // --- Respond to ping packet --- + let ping_packet = Packet::parse(client_stream).await?; + match ping_packet.id.get_int() { + 1 => Ok(ping_packet + .send_packet(client_stream) + .await + .map_err(|_| "Failed to send ping")?), + _ => Err(OpaqueError::create(&format!( + "Expected ping packet, got: {}", + ping_packet.id.get_int() + ))), + } +} +/// small helper function for sending `status request` to client +pub async fn complete_status_request( + client_stream: &mut TcpStream, + status_struct: StatusStructNew, +) -> Result<(), OpaqueError> { + let status_res = + crate::packets::clientbound::status::StatusResponse::set_json(Box::new(status_struct)) + .await; + status_res + .send_packet(client_stream) + .await + .map_err(|_| "Failed to send status packet")?; + Ok(()) +} + +/// Disconnects the client. +/// +/// It works if the client is in the login state, and it +/// has *already* and *only* sent the **handshake** and **login_start** packet. +#[tracing::instrument(skip(client_stream))] +pub async fn send_disconnect( + client_stream: &mut TcpStream, + reason: &str, +) -> Result<(), OpaqueError> { + let disconnect_packet = + crate::packets::clientbound::login::Disconnect::set_reason(reason.to_owned()) + .await + .ok_or_else(|| "failed to *create* disconnect packet")?; + disconnect_packet + .send_packet(client_stream) + .await + .map_err(|_| "failed to *send* disconnect packet")?; + client_stream.flush().await.map_err(|e| e.to_string())?; + Ok(()) +} diff --git a/src/mc_server.rs b/src/mc_server/mod.rs similarity index 81% rename from src/mc_server.rs rename to src/mc_server/mod.rs index e403b71..ff6b544 100644 --- a/src/mc_server.rs +++ b/src/mc_server/mod.rs @@ -1,68 +1,17 @@ use std::{collections::HashMap, sync::Arc}; use async_trait::async_trait; -use tokio::{io::AsyncWriteExt, net::TcpStream}; +use tokio::net::TcpStream; use tracing::Instrument; use crate::{ packets::{ - clientbound::status::{StatusStructNew, StatusTrait}, - serverbound::handshake::Handshake, - Packet, SendPacket, + clientbound::status::StatusTrait, serverbound::handshake::Handshake, Packet, SendPacket, }, OpaqueError, }; -#[tracing::instrument(skip(client_stream))] -pub async fn handle_ping(client_stream: &mut TcpStream) -> Result<(), OpaqueError> { - // --- Respond to ping packet --- - let ping_packet = Packet::parse(client_stream).await?; - match ping_packet.id.get_int() { - 1 => Ok(ping_packet - .send_packet(client_stream) - .await - .map_err(|_| "Failed to send ping")?), - _ => Err(OpaqueError::create(&format!( - "Expected ping packet, got: {}", - ping_packet.id.get_int() - ))), - } -} -/// small helper function for sending `status request` to client -pub async fn complete_status_request( - client_stream: &mut TcpStream, - status_struct: StatusStructNew, -) -> Result<(), OpaqueError> { - let status_res = - crate::packets::clientbound::status::StatusResponse::set_json(Box::new(status_struct)) - .await; - status_res - .send_packet(client_stream) - .await - .map_err(|_| "Failed to send status packet")?; - Ok(()) -} - -/// Disconnects the client. -/// -/// It works if the client is in the login state, and it -/// has *already* and *only* sent the **handshake** and **login_start** packet. -#[tracing::instrument(skip(client_stream))] -pub async fn send_disconnect( - client_stream: &mut TcpStream, - reason: &str, -) -> Result<(), OpaqueError> { - let disconnect_packet = - crate::packets::clientbound::login::Disconnect::set_reason(reason.to_owned()) - .await - .ok_or_else(|| "failed to *create* disconnect packet")?; - disconnect_packet - .send_packet(client_stream) - .await - .map_err(|_| "failed to *send* disconnect packet")?; - client_stream.flush().await.map_err(|e| e.to_string())?; - Ok(()) -} +pub mod helpers; #[async_trait::async_trait] pub trait MinecraftServerHandle: Send + Sync + 'static + Clone { diff --git a/src/opaque_error.rs b/src/opaque_error.rs index 07c65d7..d5071b1 100644 --- a/src/opaque_error.rs +++ b/src/opaque_error.rs @@ -23,8 +23,8 @@ impl fmt::Display for OpaqueError { vec.push(metadata.name()); true }); - if vec.len() > 1 { - write!(f, "trace = {}", vec.pop().unwrap())?; + if let Some(trace) = vec.pop() { + write!(f, "trace = {}", trace)?; } vec.reverse(); for s in vec { diff --git a/src/proxy.rs b/src/proxy.rs index 24a6aca..5c1066c 100644 --- a/src/proxy.rs +++ b/src/proxy.rs @@ -208,7 +208,7 @@ async fn handle_status( "Could not find §kserver§r: §f§o{join_addr}§r\nMinecraft Ingress - {BYE_MESSAGE}" ); - mc_server::complete_status_request(client_stream, status_struct) + mc_server::helpers::complete_status_request(client_stream, status_struct) .instrument(span.clone()) .await?; @@ -261,8 +261,8 @@ async fn handle_status( ServerDeploymentStatus::Unavailable(_) => unreachable!(), }; - mc_server::complete_status_request(client_stream, status_struct).await?; - return mc_server::handle_ping(client_stream).await; + mc_server::helpers::complete_status_request(client_stream, status_struct).await?; + return mc_server::helpers::handle_ping(client_stream).await; } #[tracing::instrument(level = "info", fields(join_addr = handshake.get_server_address(),join_port = handshake.server_port.get_value(),username = login_start.name.get_value()),skip(client_stream, handshake, api, login_start))] @@ -337,13 +337,13 @@ where } ServerDeploymentStatus::PodOk | ServerDeploymentStatus::Starting => { tracing::info!(?status, "server is starting... disconnecting client"); - mc_server::send_disconnect(client_stream, format!("[\"\",{{\"text\":\"The server is still starting up...\n wait a bit more please ^^\n\n\"}},{{\"text\":\"{BYE_MESSAGE}\"}}]").as_str()).await?; + mc_server::helpers::send_disconnect(client_stream, format!("[\"\",{{\"text\":\"The server is still starting up...\n wait a bit more please ^^\n\n\"}},{{\"text\":\"{BYE_MESSAGE}\"}}]").as_str()).await?; } ServerDeploymentStatus::Offline => { let server = server?; server.start().await?; api.start_watch(server.clone(), OFFLINE_TIMER).await?; - mc_server::send_disconnect(client_stream, format!("[\"\",{{\"text\":\"Okayy, §2starting§r the server!\n\n\"}},{{\"text\":\"{BYE_MESSAGE}\"}}]").as_str()).await?; + mc_server::helpers::send_disconnect(client_stream, format!("[\"\",{{\"text\":\"Okayy, §2starting§r the server!\n\n\"}},{{\"text\":\"{BYE_MESSAGE}\"}}]").as_str()).await?; } ServerDeploymentStatus::Unavailable(_) => { tracing::info!(?status, "tried connecting, droppping connection");