chore: some code cleanup
- remove one unwrap - move `mc_server` module to folder and create `::helpers` submodule
This commit is contained in:
parent
eac1c5d8fa
commit
241504ac3b
5 changed files with 87 additions and 80 deletions
|
|
@ -33,7 +33,7 @@ pub struct KubeCache {
|
||||||
dep_cache: kube::runtime::reflector::Store<Deployment>,
|
dep_cache: kube::runtime::reflector::Store<Deployment>,
|
||||||
dep_api: Api<Deployment>,
|
dep_api: Api<Deployment>,
|
||||||
srv_api: Api<Service>,
|
srv_api: Api<Service>,
|
||||||
in_cluster: bool,
|
running_in_cluster: bool,
|
||||||
}
|
}
|
||||||
impl KubeCache {
|
impl KubeCache {
|
||||||
/// This initializes the creation of a "kubernetes client"
|
/// This initializes the creation of a "kubernetes client"
|
||||||
|
|
@ -86,7 +86,7 @@ impl KubeCache {
|
||||||
dep_cache: dep_reader,
|
dep_cache: dep_reader,
|
||||||
dep_api: dep_api,
|
dep_api: dep_api,
|
||||||
srv_api,
|
srv_api,
|
||||||
in_cluster,
|
running_in_cluster: in_cluster,
|
||||||
},
|
},
|
||||||
handle,
|
handle,
|
||||||
));
|
));
|
||||||
|
|
@ -117,7 +117,7 @@ impl KubeCache {
|
||||||
Some(t) => t == "ClusterIP",
|
Some(t) => t == "ClusterIP",
|
||||||
None => false,
|
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)
|
!incorrect_type && filter_label_value(x, addr, port)
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
@ -164,7 +164,7 @@ impl MinecraftAPI<Server> for McApi {
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
let internal_id = match deployment.metadata.name.clone() {
|
let resource_name = match deployment.metadata.name.clone() {
|
||||||
Some(x) => x,
|
Some(x) => x,
|
||||||
None => {
|
None => {
|
||||||
return Err(OpaqueError::create_with_kind(
|
return Err(OpaqueError::create_with_kind(
|
||||||
|
|
@ -183,14 +183,8 @@ impl MinecraftAPI<Server> for McApi {
|
||||||
.ok_or(OpaqueError::create(
|
.ok_or(OpaqueError::create(
|
||||||
"Could not find \"mc-router\" nodePort for server",
|
"Could not find \"mc-router\" nodePort for server",
|
||||||
))?;
|
))?;
|
||||||
let inter_addr = match self.cache.in_cluster {
|
|
||||||
false => {
|
let inter_addr = match self.cache.running_in_cluster {
|
||||||
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)
|
|
||||||
}
|
|
||||||
true => {
|
true => {
|
||||||
let target_port = port_srv.port;
|
let target_port = port_srv.port;
|
||||||
format!(
|
format!(
|
||||||
|
|
@ -201,12 +195,19 @@ impl MinecraftAPI<Server> for McApi {
|
||||||
target_port
|
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 {
|
return Ok(Server {
|
||||||
server_addr: addr.to_string(),
|
server_addr: addr.to_string(),
|
||||||
server_port: port.to_string(),
|
server_port: port.to_string(),
|
||||||
internal_address: inter_addr,
|
internal_address: inter_addr,
|
||||||
internal_id,
|
resource_name,
|
||||||
cache: self.cache.clone(),
|
cache: self.cache.clone(),
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
@ -236,13 +237,13 @@ pub struct Server {
|
||||||
server_addr: String,
|
server_addr: String,
|
||||||
server_port: String,
|
server_port: String,
|
||||||
internal_address: String,
|
internal_address: String,
|
||||||
internal_id: String,
|
resource_name: String,
|
||||||
cache: KubeCache,
|
cache: KubeCache,
|
||||||
}
|
}
|
||||||
impl fmt::Debug for Server {
|
impl fmt::Debug for Server {
|
||||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||||
f.debug_struct("KubeServer")
|
f.debug_struct("KubeServer")
|
||||||
.field("internal_id", &self.internal_id)
|
.field("internal_id", &self.resource_name)
|
||||||
.field(
|
.field(
|
||||||
"join_addr",
|
"join_addr",
|
||||||
&format!("{}:{}", self.server_addr, self.server_port),
|
&format!("{}:{}", self.server_addr, self.server_port),
|
||||||
|
|
@ -268,7 +269,7 @@ impl MinecraftServerHandle for Server {
|
||||||
async fn query_status(&self) -> Result<crate::mc_server::ServerDeploymentStatus, OpaqueError> {
|
async fn query_status(&self) -> Result<crate::mc_server::ServerDeploymentStatus, OpaqueError> {
|
||||||
let dep = self
|
let dep = self
|
||||||
.cache
|
.cache
|
||||||
.get_dep(&self.internal_id)
|
.get_dep(&self.resource_name)
|
||||||
.ok_or_else(|| OpaqueError::create("Failed to get deployment from cache"))?;
|
.ok_or_else(|| OpaqueError::create("Failed to get deployment from cache"))?;
|
||||||
|
|
||||||
let status = match &dep.status {
|
let status = match &dep.status {
|
||||||
|
|
@ -319,7 +320,7 @@ impl MinecraftServerHandle for Server {
|
||||||
self.internal_address.as_str()
|
self.internal_address.as_str()
|
||||||
}
|
}
|
||||||
fn get_internal_id(&self) -> &str {
|
fn get_internal_id(&self) -> &str {
|
||||||
self.internal_id.as_str()
|
self.resource_name.as_str()
|
||||||
}
|
}
|
||||||
|
|
||||||
fn get_join_addr(&self) -> &str {
|
fn get_join_addr(&self) -> &str {
|
||||||
|
|
@ -331,7 +332,7 @@ impl MinecraftServerHandle for Server {
|
||||||
}
|
}
|
||||||
|
|
||||||
fn get_motd(&self) -> Option<String> {
|
fn get_motd(&self) -> Option<String> {
|
||||||
let dep = self.cache.get_dep(&self.internal_id)?;
|
let dep = self.cache.get_dep(&self.resource_name)?;
|
||||||
|
|
||||||
let all_container_motds = dep
|
let all_container_motds = dep
|
||||||
.spec
|
.spec
|
||||||
|
|
@ -357,7 +358,7 @@ impl MinecraftServerHandle for Server {
|
||||||
|
|
||||||
impl Server {
|
impl Server {
|
||||||
async fn set_scale(&self, num: i32) -> Result<(), kube::Error> {
|
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!(
|
tracing::info!(
|
||||||
"scaled replicas of {}:{} to {num}",
|
"scaled replicas of {}:{} to {num}",
|
||||||
self.server_addr,
|
self.server_addr,
|
||||||
|
|
|
||||||
57
src/mc_server/helpers.rs
Normal file
57
src/mc_server/helpers.rs
Normal file
|
|
@ -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(())
|
||||||
|
}
|
||||||
|
|
@ -1,68 +1,17 @@
|
||||||
use std::{collections::HashMap, sync::Arc};
|
use std::{collections::HashMap, sync::Arc};
|
||||||
|
|
||||||
use async_trait::async_trait;
|
use async_trait::async_trait;
|
||||||
use tokio::{io::AsyncWriteExt, net::TcpStream};
|
use tokio::net::TcpStream;
|
||||||
use tracing::Instrument;
|
use tracing::Instrument;
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
packets::{
|
packets::{
|
||||||
clientbound::status::{StatusStructNew, StatusTrait},
|
clientbound::status::StatusTrait, serverbound::handshake::Handshake, Packet, SendPacket,
|
||||||
serverbound::handshake::Handshake,
|
|
||||||
Packet, SendPacket,
|
|
||||||
},
|
},
|
||||||
OpaqueError,
|
OpaqueError,
|
||||||
};
|
};
|
||||||
|
|
||||||
#[tracing::instrument(skip(client_stream))]
|
pub mod helpers;
|
||||||
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(())
|
|
||||||
}
|
|
||||||
|
|
||||||
#[async_trait::async_trait]
|
#[async_trait::async_trait]
|
||||||
pub trait MinecraftServerHandle: Send + Sync + 'static + Clone {
|
pub trait MinecraftServerHandle: Send + Sync + 'static + Clone {
|
||||||
|
|
@ -23,8 +23,8 @@ impl fmt::Display for OpaqueError {
|
||||||
vec.push(metadata.name());
|
vec.push(metadata.name());
|
||||||
true
|
true
|
||||||
});
|
});
|
||||||
if vec.len() > 1 {
|
if let Some(trace) = vec.pop() {
|
||||||
write!(f, "trace = {}", vec.pop().unwrap())?;
|
write!(f, "trace = {}", trace)?;
|
||||||
}
|
}
|
||||||
vec.reverse();
|
vec.reverse();
|
||||||
for s in vec {
|
for s in vec {
|
||||||
|
|
|
||||||
10
src/proxy.rs
10
src/proxy.rs
|
|
@ -208,7 +208,7 @@ async fn handle_status<T: MinecraftServerHandle>(
|
||||||
"Could not find §kserver§r: §f§o{join_addr}§r\nMinecraft Ingress - {BYE_MESSAGE}"
|
"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())
|
.instrument(span.clone())
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
|
|
@ -261,8 +261,8 @@ async fn handle_status<T: MinecraftServerHandle>(
|
||||||
ServerDeploymentStatus::Unavailable(_) => unreachable!(),
|
ServerDeploymentStatus::Unavailable(_) => unreachable!(),
|
||||||
};
|
};
|
||||||
|
|
||||||
mc_server::complete_status_request(client_stream, status_struct).await?;
|
mc_server::helpers::complete_status_request(client_stream, status_struct).await?;
|
||||||
return mc_server::handle_ping(client_stream).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))]
|
#[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 => {
|
ServerDeploymentStatus::PodOk | ServerDeploymentStatus::Starting => {
|
||||||
tracing::info!(?status, "server is starting... disconnecting client");
|
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 => {
|
ServerDeploymentStatus::Offline => {
|
||||||
let server = server?;
|
let server = server?;
|
||||||
server.start().await?;
|
server.start().await?;
|
||||||
api.start_watch(server.clone(), OFFLINE_TIMER).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(_) => {
|
ServerDeploymentStatus::Unavailable(_) => {
|
||||||
tracing::info!(?status, "tried connecting, droppping connection");
|
tracing::info!(?status, "tried connecting, droppping connection");
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue