diff --git a/migrations/20260327_init.sql b/migrations/20260327_init.sql index c34f5a9..0b5ad6a 100644 --- a/migrations/20260327_init.sql +++ b/migrations/20260327_init.sql @@ -18,14 +18,14 @@ create table if not exists servers ( port integer, first_seen timestamp without time zone not null default now (), last_seen timestamp without time zone not null, - last_time_player_online timestamp without time zone not null, - last_time_no_players_online timestamp without time zone not null, + last_time_player_online timestamp without time zone, + last_time_no_players_online timestamp without time zone, version_protocol integer not null, version_name text not null, enforces_secure_chat boolean not null, previews_chat boolean not null, is_online_mode boolean not null, - is_whitelisted boolean not null, + is_whitelisted boolean, favicon_hash bytea references favicons (hash), max_players integer, online_players integer, @@ -33,7 +33,7 @@ create table if not exists servers ( description_raw jsonb, -- neoforge - is_modded boolean, + neoforge_is_modded boolean, -- forge fml_network_version integer, @@ -53,8 +53,8 @@ create table if not exists players ( address inet, port integer, uuid uuid, - username text not null, is_online_mode boolean not null, + username text not null, first_seen timestamp without time zone not null default now (), last_seen timestamp without time zone not null, primary key (address, port, uuid), @@ -65,11 +65,9 @@ create table if not exists mods ( address inet, port integer, id text, - mod_marker text, + version text, primary key (address, port, id), foreign key (address, port) references servers (address, port) on delete cascade ); --- Indexes - -create index if not exists countries_index on countries using gist (network inet_ops); +create index if not exists countries_index on countries using gist (network inet_ops); \ No newline at end of file diff --git a/src/database/country_tracking.rs b/src/database/country_tracking.rs index a8ef6fb..4fb24a4 100644 --- a/src/database/country_tracking.rs +++ b/src/database/country_tracking.rs @@ -10,7 +10,7 @@ use std::fs::File; use std::io::{Read, Write}; use std::str::FromStr; use std::time::Duration; -use tracing::{debug, info}; +use tracing::info; const DOWNLOAD_URL: &str = "https://ipinfo.io/data/ipinfo_lite.json.gz?token="; diff --git a/src/database/mod.rs b/src/database/mod.rs index 6499558..298d93e 100644 --- a/src/database/mod.rs +++ b/src/database/mod.rs @@ -2,19 +2,16 @@ pub mod country_tracking; use std::str::FromStr; -use crate::protocol::response::MinecraftPlayers; - use super::protocol::response::MinecraftServer; -use sqlx::postgres::PgQueryResult; -use sqlx::types::ipnet::IpNet; +use chrono::NaiveDateTime; +use sqlx::postgres::{PgArguments, PgQueryResult}; use sqlx::types::Uuid; -use sqlx::{Pool, Postgres, QueryBuilder, Row}; - -const INSERT_SERVERS_QUERY: &str = "INSERT INTO servers VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23) ON CONFLICT (address, port) DO UPDATE SET last_seen = EXCLUDED.last_seen, version_protocol = EXCLUDED.version_protocol, version_name = EXCLUDED.version_name, enforces_secure_chat = EXCLUDED.enforces_secure_chat, previews_chat = EXCLUDED.previews_chat, is_online_mode = EXCLUDED.is_online_mode, max_players = EXCLUDED.max_players, online_players = EXCLUDED.online_players, description_formatted = EXCLUDED.description_formatted, description_raw = EXCLUDED.description_raw, is_modded = EXCLUDED.is_modded, fml_network_version = EXCLUDED.fml_network_version, prevents_chat_reports = EXCLUDED.prevents_chat_reports, bcc_modpack_projectid = EXCLUDED.bcc_modpack_projectid, bcc_modpack_version = EXCLUDED.bcc_modpack_version, bcc_modpack_name = EXCLUDED.bcc_modpack_name"; +use sqlx::types::ipnet::IpNet; +use sqlx::{Encode, Execute, Pool, Postgres, QueryBuilder, Row}; #[derive(Debug, Clone)] pub struct Database { - pub connection: Pool + pub connection: Pool, } impl Database { @@ -40,60 +37,148 @@ pub struct ServerUpdateOperation { pub server: MinecraftServer, pub address: IpNet, pub port: i32, - pub timestamp: chrono::NaiveDateTime, + pub timestamp: NaiveDateTime, pub database: Database, } impl ServerUpdateOperation { /// Inserts a single server into the database pub async fn update_or_insert_server(&self) -> anyhow::Result<()> { + let mut query = String::from("INSERT INTO servers (address, port, first_seen, last_seen, "); + + // Used to determine if we should update last time player seen online or last time no player seen online field + let has_players_online = self + .server + .players + .sample + .as_ref() + .is_some_and(|s| !s.is_empty()); + + query.push_str(match has_players_online { + true => "last_time_player_online, ", + false => "last_time_no_players_online, ", + }); + + query.push_str("version_protocol, version_name, enforces_secure_chat, previews_chat, is_online_mode, favicon_hash, max_players, online_players, description_formatted, description_raw, neoforge_is_modded, fml_network_version, prevents_chat_reports, bcc_modpack_projectid, bcc_modpack_version, bcc_modpack_name) VALUES ("); + + let mut query_builder = sqlx::QueryBuilder::new(query); + + query_builder + .push_bind(self.address) + .push(", ") + .push_bind(self.port) + .push(", ") + .push_bind(self.timestamp) + .push(", ") + .push_bind(self.timestamp) + .push(", ") + .push_bind(self.timestamp) + .push(", ") + .push_bind(self.server.version.protocol) + .push(", ") + .push_bind(&self.server.version.name) + .push(", ") + .push_bind(self.server.enforces_secure_chat.unwrap_or(false)) + .push(", ") + .push_bind(self.server.previews_chat.unwrap_or(false)) + .push(", ") + .push_bind(self.server.is_server_online_mode()) + .push(", "); + + // Hash the icon using blake3 let favicon_hash = self .server .favicon .as_ref() - .map(|x| MinecraftServer::get_favicon_hash(x)); + .and_then(|a| a.split("data:image/png;base64,").nth(1)) + .map(|str| blake3::hash(str.as_bytes())) + .map(|a| *a.as_bytes()); + // The favicon must be inserted before the server due to foreign key constraints + if let Some(hash) = favicon_hash { + self.update_or_insert_favicon(hash).await.unwrap(); + } + + query_builder + .push_bind(favicon_hash) + .push(", ") + .push_bind(self.server.players.max) + .push(", ") + .push_bind(self.server.players.online) + .push(", "); + + // Format description if it exists let formatted_description = self .server .description_raw .as_ref() .map(|value| self.server.format_description(value)); + query_builder + .push_bind(formatted_description) + .push(", ") + .push_bind(&self.server.description_raw) + .push(", ") + // Modded fields + .push_bind(self.server.is_modded) + .push(", "); + + // Forge mod loader network version + let fml_network_version = self + .server + .forge_data + .as_ref() + .map(|f| f.fml_network_version); + + query_builder + .push_bind(fml_network_version) + .push(", ") + .push_bind(self.server.prevents_chat_reports.unwrap_or(false)) + .push(", "); + + // Better compatibility checker let bcc = self.server.bcc.as_ref(); - sqlx::query(INSERT_SERVERS_QUERY) - .bind(self.address) - .bind(self.port) - // first_seen - .bind(self.timestamp) - // last_seen - .bind(self.timestamp) - // last_time_player_seen_online - .bind(self.timestamp) - // last_time_no_players_seen_online - .bind(self.timestamp) - .bind(&self.server.version.protocol) - .bind(&self.server.version.name) - .bind(self.server.enforces_secure_chat.unwrap_or(false)) - .bind(self.server.previews_chat.unwrap_or(false)) - .bind(self.server.is_server_online_mode()) - // Whitelist checking, not implemented - .bind(false) - // TODO! Verify this works - .bind(self.server.players.max) - .bind(self.server.players.online) - .bind(formatted_description) - .bind(&self.server.description_raw) - // Modded fields + query_builder + .push_bind(bcc.map(|v| v.project_id)) + .push(", ") + .push_bind(bcc.map(|v| &v.version)) + .push(", ") + .push_bind(bcc.map(|v| &v.name)); - .bind(self.server.is_modded) - .bind(self.server.forge_data.as_ref().map(|f| f.fml_network_version)) - .bind(self.server.prevents_chat_reports.unwrap_or(false)) + query_builder.push( + ") ON CONFLICT (address, port) DO UPDATE SET + last_seen = EXCLUDED.last_seen, ", + ); - // Better Compatibility checker - .bind(bcc.map(|v| v.project_id)) - .bind(bcc.map(|v| &v.version)) - .bind(bcc.map(|v| &v.name)) + // Update timestamps for existing servers + query_builder.push(match has_players_online { + true => "last_time_player_online = EXCLUDED.last_time_player_online, ", + false => "last_time_no_players_online = EXCLUDED.last_time_no_players_online, ", + }); + + // Push remaining conflict updates + query_builder.push( + "version_protocol = EXCLUDED.version_protocol, + version_name = EXCLUDED.version_name, + enforces_secure_chat = EXCLUDED.enforces_secure_chat, + previews_chat = EXCLUDED.previews_chat, + is_online_mode = EXCLUDED.is_online_mode, + favicon_hash = EXCLUDED.favicon_hash, + max_players = EXCLUDED.max_players, + online_players = EXCLUDED.online_players, + description_formatted = EXCLUDED.description_formatted, + description_raw = EXCLUDED.description_raw, + neoforge_is_modded = EXCLUDED.neoforge_is_modded, + fml_network_version = EXCLUDED.fml_network_version, + prevents_chat_reports = EXCLUDED.prevents_chat_reports, + bcc_modpack_projectid = EXCLUDED.bcc_modpack_projectid, + bcc_modpack_version = EXCLUDED.bcc_modpack_version, + bcc_modpack_name = EXCLUDED.bcc_modpack_name", + ); + + query_builder + .build() .execute(&self.database.connection) .await .unwrap(); @@ -104,25 +189,22 @@ impl ServerUpdateOperation { /// Bulk insert players into the players table pub async fn update_or_insert_players(&self) -> anyhow::Result<()> { if let Some(player_sample) = &self.server.players.sample { - let mut query_builder = QueryBuilder::new( - "INSERT INTO players (address, port, uuid, username, is_online_mode, first_seen, last_seen) ", - ); + let mut query_builder = QueryBuilder::new("INSERT INTO players "); query_builder.push_values(player_sample, |mut b, player| { - // Ignore player if uuid fails to parse if let Ok(parsed_uuid) = Uuid::from_str(&player.id) { b.push_bind(self.address) .push_bind(self.port) .push_bind(parsed_uuid) - .push_bind(&player.name) .push_bind(parsed_uuid.get_version_num() == 4) + .push_bind(&player.name) .push_bind(self.timestamp) .push_bind(self.timestamp); } }); - query_builder.push("ON CONFLICT (address, port, uuid) DO UPDATE SET last_seen = EXCLUDED.last_seen, username = EXCLUDED.username"); + query_builder.push(" ON CONFLICT (address, port, uuid) DO UPDATE SET last_seen = EXCLUDED.last_seen, username = EXCLUDED.username"); query_builder .build() .execute(&self.database.connection) @@ -135,18 +217,16 @@ impl ServerUpdateOperation { /// Bulk insert forge mods into the mods table pub async fn update_or_insert_mods(&self) -> anyhow::Result<()> { if let Some(forge_data) = &self.server.forge_data { - let mut query_builder = - QueryBuilder::new("INSERT INTO mods (address, port, id, mod_marker) "); + let mut query_builder = QueryBuilder::new("INSERT INTO mods "); query_builder.push_values(&forge_data.mods, |mut b, forge_mod| { b.push_bind(self.address) .push_bind(self.port) .push_bind(&forge_mod.id) - .push_bind(&forge_mod.version) - .push_bind(self.timestamp); + .push_bind(&forge_mod.version); }); - query_builder.push("ON CONFLICT (address, port, id) DO NOTHING"); + query_builder.push(" ON CONFLICT (address, port, id) DO NOTHING"); query_builder .build() .execute(&self.database.connection) @@ -156,17 +236,14 @@ impl ServerUpdateOperation { Ok(()) } - pub async fn update_or_insert_favicon(&self) -> Result<(), sqlx::Error> { - if let Some(base64) = &self.server.favicon { - let favicon_hash = MinecraftServer::get_favicon_hash(base64); - - sqlx::query("INSERT INTO favicons (hash, data, first_seen) VALUES ($1, $2, $3) ON CONFLICT (hash) DO NOTHING") - .bind(favicon_hash.as_slice()) - .bind(base64) - .bind(self.timestamp) - .execute(&self.database.connection) - .await?; - } + /// Insert a servers icon into the favicons table + pub async fn update_or_insert_favicon(&self, hash: [u8; 32]) -> anyhow::Result<()> { + sqlx::query("INSERT INTO favicons VALUES ($1, $2, $3) ON CONFLICT (hash) DO UPDATE set last_seen = EXCLUDED.last_seen") + .bind(hash.as_slice()) + .bind(&self.server.favicon) + .bind(self.timestamp) + .execute(&self.database.connection) + .await?; Ok(()) } diff --git a/src/main.rs b/src/main.rs index 4c54a7e..0476b8c 100644 --- a/src/main.rs +++ b/src/main.rs @@ -9,16 +9,11 @@ use clap::Parser; use config::load_config; use sqlx::ConnectOptions; use sqlx::postgres::{PgConnectOptions, PgPoolOptions}; -use sqlx::types::ipnet::IpNet; -use std::net::{Ipv4Addr, SocketAddrV4}; -use std::str::FromStr; use std::time::Duration; use tracing::error; use tracing::log::LevelFilter; -use crate::database::{Database, ServerUpdateOperation, country_tracking}; -use crate::protocol::minecraft::simple_ping; -use crate::protocol::response::MinecraftServer; +use crate::database::Database; use crate::scanning::rescanner::{Rescanner, ServerRescanPriority}; #[derive(Parser, Debug)] @@ -79,40 +74,16 @@ async fn main() { .await .expect("Failed to run migrations on database!"); - let mut stream = tokio::net::TcpStream::connect(SocketAddrV4::new(Ipv4Addr::from_str("104.218.52.246").unwrap(), 25567)) - .await - .unwrap(); - - if let Ok(ping_response) = simple_ping(&mut stream).await { - - let server = serde_json::from_str::(&ping_response).unwrap(); - - let update_operation = ServerUpdateOperation { - server, - address: IpNet::from_str("104.218.52.246/32").unwrap(), - port: 25565, - timestamp: chrono::Utc::now().naive_utc(), + match arguments.mode { + Mode::Scanning => todo!(), + Mode::Rescanner => { + let rescanner = Rescanner { + is_active: true, database: Database { connection: pool }, + rescan_priority: ServerRescanPriority::OldestFirst, }; - update_operation.update_or_insert_favicon().await.unwrap(); - update_operation.update_or_insert_server().await.unwrap(); - update_operation.update_or_insert_players().await.unwrap(); - update_operation.update_or_insert_mods().await.unwrap(); - - println!("asda"); + rescanner.rescan_database().await; + } } - - // match arguments.mode { - // Mode::Scanning => todo!(), - // Mode::Rescanner => { - // let rescanner = Rescanner { - // is_active: true, - // database: Database { connection: pool }, - // rescan_priority: ServerRescanPriority::OldestFirst, - // }; - - // rescanner.rescan_database().await; - // } - // } } diff --git a/src/protocol/mod.rs b/src/protocol/mod.rs index b2b1975..321946d 100644 --- a/src/protocol/mod.rs +++ b/src/protocol/mod.rs @@ -76,4 +76,4 @@ impl MinecraftColorCodes { UnknownValue => 'r', } } -} \ No newline at end of file +} diff --git a/src/protocol/response.rs b/src/protocol/response.rs index 00a9294..8caeca3 100644 --- a/src/protocol/response.rs +++ b/src/protocol/response.rs @@ -55,10 +55,6 @@ pub struct MinecraftPlayers { #[derive(Deserialize, Serialize, PartialEq, Clone, Debug)] pub struct MinecraftPlayer { - // Some servers could send a player sample with empty name or id fields. - // This is useful for fingerprinting but currently ignored - // To add it in the future all we would need to do is change these fields to be an Option - // and check if either is None pub id: String, pub name: String, } @@ -66,8 +62,8 @@ pub struct MinecraftPlayer { #[derive(Deserialize, Serialize, PartialEq, Clone, Debug)] pub struct ForgeData { #[serde(rename = "fmlNetworkVersion")] - pub fml_network_version: i32, - pub truncated: bool, + pub fml_network_version: Option, + pub truncated: Option, #[serde(rename = "d")] pub forge_encoded_data: Option, @@ -93,11 +89,6 @@ pub struct BetterCompatibilityChecker { } impl MinecraftServer { - /// Hash the favicons data - pub fn get_favicon_hash(favicon: &str) -> [u8; 32] { - *blake3::hash(favicon.as_bytes()).as_bytes() - } - // Checks if this server is running in offline mode by checking if any of the players uuid version is anything other than 4 pub fn is_server_online_mode(&self) -> bool { if let Some(sample) = &self.players.sample { @@ -126,12 +117,8 @@ impl MinecraftServer { } // Check for duplicate uuids - self.players - .sample - .clone() - .unwrap() - .into_iter() - .for_each(|a| { + self.players.sample.as_ref().map(|a| { + a.iter().for_each(|a| { let uuid = Uuid::parse_str(&a.id).unwrap(); if seen_uuids.contains(&uuid) { @@ -139,7 +126,8 @@ impl MinecraftServer { } seen_uuids.insert(uuid); - }); + }) + }); is_fake_sample } diff --git a/src/scanning/discovery.rs b/src/scanning/discovery.rs index d706dda..9605f75 100644 --- a/src/scanning/discovery.rs +++ b/src/scanning/discovery.rs @@ -1,3 +1,4 @@ +use chrono::Timelike; use sqlx::types::ipnet::{IpNet, Ipv4Net}; use std::{ net::{Ipv4Addr, SocketAddrV4}, @@ -66,15 +67,21 @@ impl DiscoveryScanner { if let Ok(response) = simple_ping(&mut stream).await { if let Ok(server) = serde_json::from_str::(&response) { + let address = IpNet::from(Ipv4Net::from(address)); + + if server.has_opted_out() { + println!("Deleting server!"); + database_clone.delete_server(address).await.unwrap(); + } + let update_operation = ServerUpdateOperation { server, - address: IpNet::from(Ipv4Net::from(address)), + address, port: port as i32, - timestamp: chrono::Utc::now().naive_utc(), + timestamp: chrono::Utc::now().naive_utc().with_nanosecond(0).unwrap(), database: database_clone, }; - update_operation.update_or_insert_favicon().await.unwrap(); update_operation.update_or_insert_server().await.unwrap(); update_operation.update_or_insert_players().await.unwrap(); update_operation.update_or_insert_mods().await.unwrap(); diff --git a/src/scanning/rescanner.rs b/src/scanning/rescanner.rs index 6c73089..945b46d 100644 --- a/src/scanning/rescanner.rs +++ b/src/scanning/rescanner.rs @@ -1,5 +1,6 @@ use std::net::SocketAddrV4; +use chrono::Timelike; use futures_util::StreamExt; use indicatif::{ProgressBar, ProgressStyle}; use sqlx::{Row, types::ipnet::IpNet}; @@ -46,7 +47,7 @@ impl Rescanner { // Need to be able to pause this process at anytime when the countries table gets updated while let Some(Ok(row)) = servers_stream.next().await && self.is_active { let address = row.get::("address"); - + let connect_address = match address { IpNet::V4(i) => i.addr(), _ => continue, @@ -58,15 +59,19 @@ impl Rescanner { if let Ok(ping_response) = simple_ping(&mut stream).await { if let Ok(server) = serde_json::from_str::(&ping_response) { + + if server.has_opted_out() { + database_clone.delete_server(address).await.unwrap(); + } + let update_operation = ServerUpdateOperation { server, address, port: 25565, - timestamp: chrono::Utc::now().naive_utc(), + timestamp: chrono::Utc::now().naive_utc().with_nanosecond(0).unwrap(), database: database_clone, }; - - update_operation.update_or_insert_favicon().await.unwrap(); + update_operation.update_or_insert_server().await.unwrap(); update_operation.update_or_insert_players().await.unwrap(); update_operation.update_or_insert_mods().await.unwrap();