From 8acb4a64e698e2bb76b71009a2bfe08792341c01 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 30 Sep 2026 17:09:31 +0000 Subject: [PATCH] Add central live map backend: chunk/player ingest, Vantage sidecar supervisor and proxy Co-Authored-By: Claude Sonnet 5.5 Claude-Session: https://claude.ai/code/session_01DjMbLQujBHunCCu5GpsHaT --- Cargo.lock | 68 ++- panel/server/Cargo.toml | 2 +- panel/server/src/config.rs | 9 + panel/server/src/db.rs | 3 + panel/server/src/lib.rs | 2 + panel/server/src/livemap.rs | 686 +++++++++++++++++++++++++++++ panel/server/src/main.rs | 1 + panel/server/src/routes/livemap.rs | 268 +++++++++++ panel/server/src/routes/mod.rs | 16 +- panel/server/src/routes/servers.rs | 13 +- panel/server/src/state.rs | 1 + panel/server/tests/common/mod.rs | 11 +- panel/server/tests/livemap.rs | 182 ++++++++ 13 files changed, 1255 insertions(+), 7 deletions(-) create mode 100644 panel/server/src/livemap.rs create mode 100644 panel/server/src/routes/livemap.rs create mode 100644 panel/server/tests/livemap.rs diff --git a/Cargo.lock b/Cargo.lock index 08207af..d47ca4e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -279,6 +279,7 @@ checksum = "31b698c5f9a010f6573133b09e0de5408834d0c82f8d7475a89fc1867a71cd90" dependencies = [ "axum-core", "axum-macros", + "base64 0.22.1", "bytes", "form_urlencoded", "futures-util", @@ -298,8 +299,10 @@ dependencies = [ "serde_json", "serde_path_to_error", "serde_urlencoded", + "sha1", "sync_wrapper", "tokio", + "tokio-tungstenite", "tower", "tower-layer", "tower-service", @@ -916,6 +919,12 @@ dependencies = [ "syn 3.0.6", ] +[[package]] +name = "data-encoding" +version = "2.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4583a4551df46e2792f82ceeac45e850d2e2d5debba0b91f102385cda5b11f06" + [[package]] name = "dbus" version = "0.9.12" @@ -3558,10 +3567,20 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e058c7de0b26af77780c769414d6257830bb240f3c38477dbc2c16e5f54d6d4c" dependencies = [ "libc", - "rand_chacha", + "rand_chacha 0.3.1", "rand_core 0.6.4", ] +[[package]] +name = "rand" +version = "0.9.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b9ef1d0d795eb7d84685bca4f72f3649f064e6641543d3a8c415898726a57b41" +dependencies = [ + "rand_chacha 0.9.0", + "rand_core 0.9.5", +] + [[package]] name = "rand" version = "0.10.3" @@ -3583,6 +3602,16 @@ dependencies = [ "rand_core 0.6.4", ] +[[package]] +name = "rand_chacha" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" +dependencies = [ + "ppv-lite86", + "rand_core 0.9.5", +] + [[package]] name = "rand_core" version = "0.6.4" @@ -3592,6 +3621,15 @@ dependencies = [ "getrandom 0.2.17", ] +[[package]] +name = "rand_core" +version = "0.9.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "76afc826de14238e6e8c374ddcc1fa19e374fd8dd986b0d2af0d02377261d83c" +dependencies = [ + "getrandom 0.3.4", +] + [[package]] name = "rand_core" version = "0.10.1" @@ -5238,6 +5276,18 @@ dependencies = [ "tokio", ] +[[package]] +name = "tokio-tungstenite" +version = "0.29.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f72a05e828585856dacd553fba484c242c46e391fb0e58917c942ee9202915c" +dependencies = [ + "futures-util", + "log", + "tokio", + "tungstenite", +] + [[package]] name = "tokio-util" version = "0.7.19" @@ -5495,6 +5545,22 @@ version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" +[[package]] +name = "tungstenite" +version = "0.29.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6c01152af293afb9c7c2a57e4b559c5620b421f6d133261c60dd2d0cdb38e6b8" +dependencies = [ + "bytes", + "data-encoding", + "http", + "httparse", + "log", + "rand 0.9.5", + "sha1", + "thiserror 2.0.21", +] + [[package]] name = "typeid" version = "1.0.3" diff --git a/panel/server/Cargo.toml b/panel/server/Cargo.toml index 6cd3f51..9ff5df1 100644 --- a/panel/server/Cargo.toml +++ b/panel/server/Cargo.toml @@ -21,7 +21,7 @@ zip.workspace = true tracing.workspace = true chrono.workspace = true rand.workspace = true -axum = { version = "0.8", features = ["multipart", "macros"] } +axum = { version = "0.8", features = ["multipart", "macros", "ws"] } tower = { version = "0.5", features = ["util"] } tower-http = { version = "0.6", features = ["fs", "trace", "compression-gzip", "set-header"] } tracing-subscriber = { version = "0.3", features = ["env-filter", "fmt"] } diff --git a/panel/server/src/config.rs b/panel/server/src/config.rs index 44df8f4..a7b3c4b 100644 --- a/panel/server/src/config.rs +++ b/panel/server/src/config.rs @@ -14,6 +14,12 @@ pub struct Config { pub max_upload_mb: usize, pub public_url: Option, pub trusted_proxies: Vec, + /// Vantage generator binary (defaults to `vantage` on PATH or /opt/vantage). + pub vantage_bin: Option, + /// Minecraft client assets directory handed to Vantage. + pub vantage_assets: Option, + /// Extra arguments appended to `vantage server`. + pub vantage_args: Vec, } fn var(name: &str) -> Option { @@ -35,6 +41,9 @@ impl Config { trusted_proxies: var("SCOPENET_TRUSTED_PROXIES") .map(|s| s.split(',').map(|v| v.trim().parse().expect("SCOPENET_TRUSTED_PROXIES must contain IP addresses")).collect()) .unwrap_or_default(), + vantage_bin: var("SCOPENET_VANTAGE_BIN").map(Into::into), + vantage_assets: var("SCOPENET_VANTAGE_ASSETS").map(Into::into), + vantage_args: var("SCOPENET_VANTAGE_ARGS").map(|s| s.split_whitespace().map(str::to_owned).collect()).unwrap_or_default(), } } diff --git a/panel/server/src/db.rs b/panel/server/src/db.rs index fd51089..93fa0ad 100644 --- a/panel/server/src/db.rs +++ b/panel/server/src/db.rs @@ -555,6 +555,9 @@ const MIGRATIONS: &[&str] = &[ ); CREATE INDEX guild_wallet_transactions_recent ON guild_wallet_transactions(guild_id,server_id,id DESC); "#, + r#" + ALTER TABLE game_servers ADD COLUMN live_map_enabled INTEGER NOT NULL DEFAULT 0; + "#, ]; pub async fn connect(data_dir: &Path) -> Result { diff --git a/panel/server/src/lib.rs b/panel/server/src/lib.rs index 07ab1c1..056ccf9 100644 --- a/panel/server/src/lib.rs +++ b/panel/server/src/lib.rs @@ -7,6 +7,7 @@ pub mod auth; pub mod config; pub mod db; pub mod error; +pub mod livemap; pub mod net; pub mod packs; pub mod routes; @@ -59,6 +60,7 @@ pub async fn build_state_with_keys(cfg: config::Config, db: sqlx::SqlitePool, yg keys: Arc::new(auth::Keys::new(&secret)), http: scopenet_core::http::client(), login_guard: Arc::new(auth::LoginGuard::default()), + livemap: Arc::new(livemap::LiveMap::new(&cfg.data_dir, cfg.vantage_bin.clone(), cfg.vantage_assets.clone(), cfg.vantage_args.clone())), cfg: Arc::new(cfg), }) } diff --git a/panel/server/src/livemap.rs b/panel/server/src/livemap.rs new file mode 100644 index 0000000..52b17f1 --- /dev/null +++ b/panel/server/src/livemap.rs @@ -0,0 +1,686 @@ +//! Central live map. +//! +//! Game servers (the SCOPENET plugin/mod) push the chunks that changed and +//! their players' positions to the panel. The panel keeps a mirror of each +//! server's Anvil region files under `/livemap//world`, runs a +//! Vantage `server` sidecar over that mirror on demand, and proxies the +//! sidecar to authenticated launcher/admin sessions. Nobody has to run +//! Vantage next to their Minecraft server. + +use serde::{Deserialize, Serialize}; +use std::collections::{HashMap, VecDeque}; +use std::io::{Read, Seek, SeekFrom, Write}; +use std::path::{Path, PathBuf}; +use std::process::Stdio; +use std::sync::Mutex; +use std::time::{Duration, Instant}; +use tokio::io::{AsyncBufReadExt, BufReader}; + +pub const SECTOR: usize = 4096; +const HEADER: usize = 2 * SECTOR; +/// A chunk that needs more sectors than a location entry can address lives in +/// an external `.mcc` file; the mirror skips those (they are vanishingly rare). +const MAX_SECTORS: usize = 255; +/// Sidecars idle for this long are stopped and restarted on the next view. +pub const IDLE_STOP: Duration = Duration::from_secs(15 * 60); +/// Player positions older than this are treated as "server stopped sending". +const PLAYERS_FRESH: Duration = Duration::from_secs(45); + +// --------------------------------------------------------------------------- +// Dimensions +// --------------------------------------------------------------------------- + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub enum Dim { + Overworld, + Nether, + End, +} + +impl Dim { + pub const ALL: [Dim; 3] = [Dim::Overworld, Dim::Nether, Dim::End]; + + /// Accepts the short slug (`the_nether`) or the namespaced id. + pub fn parse(s: &str) -> Option { + match s.trim().trim_start_matches("minecraft:") { + "overworld" => Some(Dim::Overworld), + "the_nether" | "nether" => Some(Dim::Nether), + "the_end" | "end" => Some(Dim::End), + _ => None, + } + } + pub fn slug(self) -> &'static str { + match self { + Dim::Overworld => "overworld", + Dim::Nether => "the_nether", + Dim::End => "the_end", + } + } + pub fn id(self) -> String { + format!("minecraft:{}", self.slug()) + } + /// Folder inside a vanilla save. + fn folder(self) -> &'static str { + match self { + Dim::Overworld => "", + Dim::Nether => "DIM-1", + Dim::End => "DIM1", + } + } +} + +// --------------------------------------------------------------------------- +// Anvil region writer +// --------------------------------------------------------------------------- + +#[derive(Debug, Clone)] +pub struct Chunk { + pub x: i32, + pub z: i32, + /// Seconds since epoch, as stored in the region header. + pub timestamp: u32, + pub compression: u8, + pub data: Vec, +} + +pub fn region_path(dir: &Path, rx: i32, rz: i32) -> PathBuf { + dir.join(format!("r.{rx}.{rz}.mca")) +} + +fn slot(x: i32, z: i32) -> usize { + (x.rem_euclid(32) + z.rem_euclid(32) * 32) as usize +} + +fn sectors_for(len: usize) -> usize { + len.div_ceil(SECTOR) +} + +fn be32(b: &[u8]) -> u32 { + u32::from_be_bytes([b[0], b[1], b[2], b[3]]) +} + +/// Write `chunks` (all in one region) into `path`, creating it if needed. +/// Chunks older than what the file already holds are ignored. Returns the +/// number of chunks stored. +pub fn write_region(path: &Path, chunks: &[Chunk]) -> std::io::Result { + if let Some(parent) = path.parent() { + std::fs::create_dir_all(parent)?; + } + let mut file = std::fs::OpenOptions::new().read(true).write(true).create(true).truncate(false).open(path)?; + let mut len = file.metadata()?.len() as usize; + let mut header = vec![0u8; HEADER]; + if len >= HEADER { + file.read_exact(&mut header)?; + } else { + file.set_len(0)?; + file.write_all(&header)?; + len = HEADER; + } + let mut end = len.div_ceil(SECTOR); + let mut stored = 0; + for c in chunks { + let i = slot(c.x, c.z); + let old_ts = be32(&header[SECTOR + i * 4..]); + if old_ts != 0 && old_ts > c.timestamp { + continue; + } + let body = c.data.len() + 1; + let count = sectors_for(body + 4); + if count > MAX_SECTORS || c.compression & 0x80 != 0 || c.data.is_empty() { + continue; + } + let mut record = Vec::with_capacity(count * SECTOR); + record.extend_from_slice(&(body as u32).to_be_bytes()); + record.push(c.compression); + record.extend_from_slice(&c.data); + record.resize(count * SECTOR, 0); + file.seek(SeekFrom::Start((end * SECTOR) as u64))?; + file.write_all(&record)?; + let loc = ((end as u32) << 8) | count as u32; + header[i * 4..i * 4 + 4].copy_from_slice(&loc.to_be_bytes()); + header[SECTOR + i * 4..SECTOR + i * 4 + 4].copy_from_slice(&c.timestamp.to_be_bytes()); + end += count; + stored += 1; + } + file.seek(SeekFrom::Start(0))?; + file.write_all(&header)?; + file.flush()?; + drop(file); + + // Rewriting a chunk appends it, leaving its old sectors behind. Compact + // once more than half the file is dead space. + let live: usize = (0..1024).map(|i| (be32(&header[i * 4..]) & 0xff) as usize).sum(); + if stored > 0 && end > 2 * (live + 2) + 64 { + compact(path)?; + } + Ok(stored) +} + +fn compact(path: &Path) -> std::io::Result<()> { + let data = std::fs::read(path)?; + if data.len() < HEADER { + return Ok(()); + } + let mut header = vec![0u8; HEADER]; + header[SECTOR..].copy_from_slice(&data[SECTOR..HEADER]); + let mut body: Vec = Vec::new(); + let mut next = 2usize; + for i in 0..1024 { + let loc = be32(&data[i * 4..]); + let (off, count) = ((loc >> 8) as usize, (loc & 0xff) as usize); + if count == 0 || (off + count) * SECTOR > data.len() { + continue; + } + body.extend_from_slice(&data[off * SECTOR..(off + count) * SECTOR]); + header[i * 4..i * 4 + 4].copy_from_slice(&(((next as u32) << 8) | count as u32).to_be_bytes()); + next += count; + } + let tmp = path.with_extension("mca.tmp"); + let mut out = std::fs::File::create(&tmp)?; + out.write_all(&header)?; + out.write_all(&body)?; + out.flush()?; + drop(out); + std::fs::rename(tmp, path) +} + +/// The 4 KiB timestamp table of a region file, or `None` if it doesn't exist. +pub fn read_timestamps(path: &Path) -> Option> { + let mut f = std::fs::File::open(path).ok()?; + f.seek(SeekFrom::Start(SECTOR as u64)).ok()?; + let mut buf = vec![0u8; SECTOR]; + f.read_exact(&mut buf).ok()?; + Some(buf) +} + +/// Parse a batch upload: repeated `i32 x, i32 z, u32 timestamp, u8 +/// compression, u32 length, bytes` records, big-endian. +pub fn parse_chunks(mut b: &[u8]) -> Result, &'static str> { + let mut out = Vec::new(); + while !b.is_empty() { + if b.len() < 17 { + return Err("truncated chunk record"); + } + let x = be32(b) as i32; + let z = be32(&b[4..]) as i32; + let timestamp = be32(&b[8..]); + let compression = b[12]; + let len = be32(&b[13..]) as usize; + b = &b[17..]; + if len > MAX_SECTORS * SECTOR || len > b.len() { + return Err("chunk record exceeds its payload"); + } + if x.unsigned_abs() > 3_750_000 || z.unsigned_abs() > 3_750_000 { + return Err("chunk coordinates are outside the world"); + } + out.push(Chunk { x, z, timestamp, compression, data: b[..len].to_vec() }); + b = &b[len..]; + } + Ok(out) +} + +// --------------------------------------------------------------------------- +// Live players +// --------------------------------------------------------------------------- + +#[derive(Debug, Clone, Deserialize, Serialize)] +pub struct LivePlayer { + pub uuid: String, + pub name: String, + #[serde(default)] + pub dimension: String, + pub x: f64, + pub y: f64, + pub z: f64, + #[serde(default)] + pub yaw: f64, + #[serde(default)] + pub pitch: f64, +} + +struct Snapshot { + at: Instant, + players: Vec, +} + +// --------------------------------------------------------------------------- +// Sidecar supervisor +// --------------------------------------------------------------------------- + +struct Sidecar { + child: tokio::process::Child, + port: u16, + token: String, + last_used: Instant, + log: std::sync::Arc>>, +} + +#[derive(Default, Clone, Serialize)] +pub struct ServerStats { + pub chunks_received: u64, + pub last_chunk_at: Option, + pub last_players_at: Option, +} + +pub struct LiveMap { + root: PathBuf, + bin: Option, + assets: Option, + extra_args: Vec, + players: Mutex>, + stats: Mutex>, + sidecars: tokio::sync::Mutex>, + last_error: Mutex>, + write_lock: tokio::sync::Mutex<()>, + players_written: Mutex>, +} + +fn find_bin(configured: Option<&Path>) -> Option { + let candidates: Vec = match configured { + Some(p) => vec![p.to_path_buf()], + None => { + let mut v: Vec = std::env::var_os("PATH").map(|p| std::env::split_paths(&p).map(|d| d.join("vantage")).collect()).unwrap_or_default(); + v.push("/opt/vantage/vantage".into()); + v + } + }; + candidates.into_iter().find(|p| p.is_file()) +} + +impl LiveMap { + pub fn new(data_dir: &Path, bin: Option, assets: Option, extra_args: Vec) -> Self { + let assets = assets.or_else(|| { + let p = data_dir.join("livemap").join("assets"); + p.is_dir().then_some(p) + }); + Self { + root: data_dir.join("livemap"), + bin: find_bin(bin.as_deref()), + assets, + extra_args, + players: Mutex::default(), + stats: Mutex::default(), + sidecars: Default::default(), + last_error: Mutex::default(), + write_lock: Default::default(), + players_written: Mutex::default(), + } + } + + pub fn server_dir(&self, id: i64) -> PathBuf { + self.root.join(id.to_string()) + } + pub fn world_dir(&self, id: i64) -> PathBuf { + self.server_dir(id).join("world") + } + pub fn region_dir(&self, id: i64, dim: Dim) -> PathBuf { + let w = self.world_dir(id); + if dim.folder().is_empty() { w.join("region") } else { w.join(dim.folder()).join("region") } + } + fn cache_dir(&self, id: i64, dim: Dim) -> PathBuf { + self.server_dir(id).join("cache").join(dim.slug()) + } + fn players_file(&self, id: i64) -> PathBuf { + self.server_dir(id).join("players.json") + } + + pub fn generator_ready(&self) -> (bool, bool) { + (self.bin.is_some(), self.assets.as_ref().is_some_and(|a| a.is_dir())) + } + pub fn generator_path(&self) -> Option { + self.bin.as_ref().map(|p| p.display().to_string()) + } + pub fn assets_path(&self) -> Option { + self.assets.as_ref().map(|p| p.display().to_string()) + } + + // ---- ingest ---- + + /// Write a batch of chunks into the server's mirror. Returns how many + /// chunks were stored. + pub async fn ingest_chunks(&self, id: i64, dim: Dim, chunks: Vec) -> std::io::Result { + let dir = self.region_dir(id, dim); + let _guard = self.write_lock.lock().await; + let stored = tokio::task::spawn_blocking(move || { + let mut by_region: HashMap<(i32, i32), Vec> = HashMap::new(); + for c in chunks { + by_region.entry((c.x.div_euclid(32), c.z.div_euclid(32))).or_default().push(c); + } + let mut stored = 0; + for ((rx, rz), list) in by_region { + stored += write_region(®ion_path(&dir, rx, rz), &list)?; + } + Ok::<_, std::io::Error>(stored) + }) + .await + .map_err(std::io::Error::other)??; + let mut stats = self.stats.lock().unwrap(); + let s = stats.entry(id).or_default(); + s.chunks_received += stored as u64; + if stored > 0 { + s.last_chunk_at = Some(crate::db::now()); + } + Ok(stored) + } + + pub async fn ingest_level(&self, id: i64, bytes: Vec) -> std::io::Result<()> { + let dir = self.world_dir(id); + tokio::fs::create_dir_all(&dir).await?; + let tmp = dir.join("level.dat.tmp"); + tokio::fs::write(&tmp, bytes).await?; + tokio::fs::rename(tmp, dir.join("level.dat")).await + } + + /// Timestamp tables for every mirrored region of one dimension, so the + /// game server only uploads chunks that are newer. + pub fn manifest(&self, id: i64, dim: Dim) -> Vec<(i32, i32, Vec)> { + let mut out = Vec::new(); + let Ok(rd) = std::fs::read_dir(self.region_dir(id, dim)) else { return out }; + for e in rd.flatten() { + let name = e.file_name().to_string_lossy().to_string(); + let Some(rest) = name.strip_prefix("r.").and_then(|r| r.strip_suffix(".mca")) else { continue }; + let Some((a, b)) = rest.split_once('.') else { continue }; + let (Ok(rx), Ok(rz)) = (a.parse::(), b.parse::()) else { continue }; + if let Some(ts) = read_timestamps(&e.path()) { + out.push((rx, rz, ts)); + } + } + out + } + + pub fn set_players(&self, id: i64, players: Vec) { + self.players.lock().unwrap().insert(id, Snapshot { at: Instant::now(), players: players.clone() }); + self.stats.lock().unwrap().entry(id).or_default().last_players_at = Some(crate::db::now()); + // Vantage reads this file to decide what to pre-render; throttle it. + let mut written = self.players_written.lock().unwrap(); + if written.get(&id).is_none_or(|t| t.elapsed() > Duration::from_secs(5)) { + written.insert(id, Instant::now()); + let doc = players_doc(&players, None); + let path = self.players_file(id); + if std::fs::create_dir_all(self.server_dir(id)).is_ok() { + let _ = std::fs::write(path, doc.to_string()); + } + } + } + + pub fn live_players(&self, id: i64) -> Vec { + self.players.lock().unwrap().get(&id).filter(|s| s.at.elapsed() < PLAYERS_FRESH).map(|s| s.players.clone()).unwrap_or_default() + } + + /// `players.json` for one dimension (BlueMap-compatible, which Vantage reads). + pub fn players_json(&self, id: i64, dim: Dim) -> serde_json::Value { + players_doc(&self.live_players(id), Some(dim)) + } + + pub fn stats(&self, id: i64) -> ServerStats { + self.stats.lock().unwrap().get(&id).cloned().unwrap_or_default() + } + + /// Disk usage and regions per dimension. + pub fn world_summary(&self, id: i64) -> serde_json::Value { + let mut dims = Vec::new(); + for dim in Dim::ALL { + let (mut regions, mut bytes) = (0u64, 0u64); + if let Ok(rd) = std::fs::read_dir(self.region_dir(id, dim)) { + for e in rd.flatten() { + if let Ok(m) = e.metadata() { + regions += 1; + bytes += m.len(); + } + } + } + dims.push(serde_json::json!({ "id": dim.id(), "slug": dim.slug(), "regions": regions, "bytes": bytes })); + } + serde_json::json!({ "level_dat": self.world_dir(id).join("level.dat").is_file(), "dimensions": dims }) + } + + pub async fn sidecars_running(&self, id: i64) -> Vec { + self.sidecars.lock().await.keys().filter(|(s, _)| *s == id).map(|(_, d)| d.slug().to_string()).collect() + } + + pub fn last_error(&self, id: i64) -> Option { + let errs = self.last_error.lock().unwrap(); + Dim::ALL.iter().find_map(|d| errs.get(&(id, *d)).cloned()) + } + + // ---- sidecars ---- + + /// Make sure a Vantage sidecar is serving this server/dimension and + /// return its loopback port and bearer token. + pub async fn ensure_sidecar(&self, id: i64, dim: Dim) -> Result<(u16, String), String> { + let bin = self.bin.clone().ok_or("The Vantage generator isn't installed on this panel. See docs/live-map.md.")?; + let assets = self.assets.clone().filter(|a| a.is_dir()).ok_or("Minecraft client assets for Vantage aren't configured. See docs/live-map.md.")?; + let world = self.world_dir(id); + if !world.join("level.dat").is_file() || !self.region_dir(id, dim).is_dir() { + return Err("Waiting for the game server to upload this world.".into()); + } + + let mut sidecars = self.sidecars.lock().await; + if let Some(s) = sidecars.get_mut(&(id, dim)) { + if matches!(s.child.try_wait(), Ok(None)) { + s.last_used = Instant::now(); + return Ok((s.port, s.token.clone())); + } + sidecars.remove(&(id, dim)); + } + + let port = { + let l = std::net::TcpListener::bind("127.0.0.1:0").map_err(|e| e.to_string())?; + l.local_addr().map_err(|e| e.to_string())?.port() + }; + let token = { + use rand::RngCore; + let mut b = [0u8; 32]; + rand::thread_rng().fill_bytes(&mut b); + hex::encode(b) + }; + let cache = self.cache_dir(id, dim); + tokio::fs::create_dir_all(&cache).await.map_err(|e| e.to_string())?; + let players = self.players_file(id); + if !players.exists() { + let _ = std::fs::write(&players, players_doc(&[], None).to_string()); + } + let mut cmd = tokio::process::Command::new(&bin); + cmd.arg("server") + .arg(&world) + .arg("--assets") + .arg(&assets) + .arg("--out") + .arg(&cache) + .args(["--host", "127.0.0.1", "--port"]) + .arg(port.to_string()) + .arg("--players-file") + .arg(&players) + .args(["--dimension", dim.slug()]) + .args(&self.extra_args) + .env("VANTAGE_SERVER_TOKEN", &token) + .stdin(Stdio::null()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .kill_on_drop(true); + let mut child = cmd.spawn().map_err(|e| format!("couldn't start the Vantage generator: {e}"))?; + let log: std::sync::Arc>> = Default::default(); + if let Some(out) = child.stdout.take() { + tail(out, log.clone()); + } + if let Some(err) = child.stderr.take() { + tail(err, log.clone()); + } + let mut sidecar = Sidecar { child, port, token: token.clone(), last_used: Instant::now(), log }; + + // Wait for it to answer health checks. + let client = reqwest::Client::new(); + let url = format!("http://127.0.0.1:{port}/v1/health"); + let deadline = Instant::now() + Duration::from_secs(30); + loop { + if let Ok(Some(status)) = sidecar.child.try_wait() { + let tail = sidecar.log.lock().unwrap().iter().rev().take(4).cloned().collect::>().into_iter().rev().collect::>().join(" | "); + let msg = format!("Vantage exited ({status}) {tail}"); + self.last_error.lock().unwrap().insert((id, dim), msg.clone()); + return Err(msg); + } + if client.get(&url).timeout(Duration::from_secs(2)).send().await.is_ok_and(|r| r.status().is_success()) { + break; + } + if Instant::now() > deadline { + let _ = sidecar.child.start_kill(); + let msg = "Vantage didn't become ready within 30 seconds".to_string(); + self.last_error.lock().unwrap().insert((id, dim), msg.clone()); + return Err(msg); + } + tokio::time::sleep(Duration::from_millis(250)).await; + } + self.last_error.lock().unwrap().remove(&(id, dim)); + sidecars.insert((id, dim), sidecar); + Ok((port, token)) + } + + /// Stop sidecars nobody has looked at recently. + pub async fn reap_idle(&self) { + let mut sidecars = self.sidecars.lock().await; + let idle: Vec<_> = sidecars.iter().filter(|(_, s)| s.last_used.elapsed() > IDLE_STOP).map(|(k, _)| *k).collect(); + for key in idle { + if let Some(mut s) = sidecars.remove(&key) { + let _ = s.child.start_kill(); + tracing::info!("stopped idle live-map generator for server {} ({})", key.0, key.1.slug()); + } + } + } + + pub fn spawn_reaper(self: &std::sync::Arc) { + let me = std::sync::Arc::downgrade(self); + tokio::spawn(async move { + loop { + tokio::time::sleep(Duration::from_secs(60)).await; + let Some(lm) = me.upgrade() else { break }; + lm.reap_idle().await; + } + }); + } + + /// Stop every sidecar of a server and delete its mirror and caches. + pub async fn reset(&self, id: i64) -> std::io::Result<()> { + { + let mut sidecars = self.sidecars.lock().await; + let keys: Vec<_> = sidecars.keys().filter(|(s, _)| *s == id).copied().collect(); + for k in keys { + if let Some(mut s) = sidecars.remove(&k) { + let _ = s.child.kill().await; + } + } + } + let _guard = self.write_lock.lock().await; + self.players.lock().unwrap().remove(&id); + self.stats.lock().unwrap().remove(&id); + match tokio::fs::remove_dir_all(self.server_dir(id)).await { + Err(e) if e.kind() != std::io::ErrorKind::NotFound => Err(e), + _ => Ok(()), + } + } +} + +fn tail(r: R, log: std::sync::Arc>>) { + tokio::spawn(async move { + let mut lines = BufReader::new(r).lines(); + while let Ok(Some(line)) = lines.next_line().await { + let mut l = log.lock().unwrap(); + if l.len() >= 30 { + l.pop_front(); + } + l.push_back(line.chars().take(300).collect()); + } + }); +} + +fn players_doc(players: &[LivePlayer], dim: Option) -> serde_json::Value { + let list: Vec<_> = players + .iter() + .map(|p| { + let here = dim.is_none_or(|d| Dim::parse(&p.dimension).unwrap_or(Dim::Overworld) == d); + serde_json::json!({ + "uuid": p.uuid, + "name": p.name, + "foreign": !here, + "position": { "x": p.x, "y": p.y, "z": p.z }, + "rotation": { "pitch": p.pitch, "yaw": p.yaw, "roll": 0 }, + }) + }) + .collect(); + serde_json::json!({ "players": list }) +} + +/// Sidecar port and token for the proxy, without starting anything. +impl LiveMap { + pub async fn sidecar_log(&self, id: i64, dim: Dim) -> Vec { + self.sidecars.lock().await.get(&(id, dim)).map(|s| s.log.lock().unwrap().iter().cloned().collect()).unwrap_or_default() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn chunk(x: i32, z: i32, ts: u32, byte: u8, len: usize) -> Chunk { + Chunk { x, z, timestamp: ts, compression: 2, data: vec![byte; len] } + } + + fn read_chunk(path: &Path, x: i32, z: i32) -> Option<(u32, Vec)> { + let data = std::fs::read(path).ok()?; + let i = slot(x, z); + let loc = be32(&data[i * 4..]); + if loc == 0 { + return None; + } + let off = (loc >> 8) as usize * SECTOR; + let len = be32(&data[off..]) as usize; + Some((be32(&data[SECTOR + i * 4..]), data[off + 5..off + 4 + len].to_vec())) + } + + #[test] + fn writes_reads_and_rewrites_chunks() { + let dir = tempfile::tempdir().unwrap(); + let path = region_path(dir.path(), 0, 0); + assert_eq!(write_region(&path, &[chunk(1, 2, 10, 7, 5000), chunk(-1 + 32, 0, 10, 9, 10)]).unwrap(), 2); + assert_eq!(read_chunk(&path, 1, 2), Some((10, vec![7; 5000]))); + // Newer replaces, older is ignored. + assert_eq!(write_region(&path, &[chunk(1, 2, 20, 8, 100)]).unwrap(), 1); + assert_eq!(write_region(&path, &[chunk(1, 2, 5, 1, 100)]).unwrap(), 0); + assert_eq!(read_chunk(&path, 1, 2), Some((20, vec![8; 100]))); + assert_eq!(read_chunk(&path, 31, 0), Some((10, vec![9; 10]))); + assert_eq!(be32(&read_timestamps(&path).unwrap()[slot(1, 2) * 4..]), 20); + } + + #[test] + fn compaction_keeps_live_chunks() { + let dir = tempfile::tempdir().unwrap(); + let path = region_path(dir.path(), -1, 3); + for ts in 1..40u32 { + write_region(&path, &[chunk(-32, 96, ts, ts as u8, 20_000)]).unwrap(); + } + write_region(&path, &[chunk(-31, 96, 1, 5, 50)]).unwrap(); + assert!(std::fs::metadata(&path).unwrap().len() < 40 * 20_000 / 2); + assert_eq!(read_chunk(&path, -32, 96), Some((39, vec![39; 20_000]))); + assert_eq!(read_chunk(&path, -31, 96), Some((1, vec![5; 50]))); + } + + #[test] + fn parses_upload_batches() { + let mut b = Vec::new(); + b.extend_from_slice(&(-3i32).to_be_bytes()); + b.extend_from_slice(&4i32.to_be_bytes()); + b.extend_from_slice(&99u32.to_be_bytes()); + b.push(2); + b.extend_from_slice(&3u32.to_be_bytes()); + b.extend_from_slice(&[1, 2, 3]); + let chunks = parse_chunks(&b).unwrap(); + assert_eq!((chunks[0].x, chunks[0].z, chunks[0].timestamp, chunks[0].data.clone()), (-3, 4, 99, vec![1, 2, 3])); + assert!(parse_chunks(&b[..b.len() - 1]).is_err()); + } + + #[test] + fn dimensions_parse() { + assert_eq!(Dim::parse("minecraft:the_nether"), Some(Dim::Nether)); + assert_eq!(Dim::parse("overworld"), Some(Dim::Overworld)); + assert_eq!(Dim::parse("../etc"), None); + } +} diff --git a/panel/server/src/main.rs b/panel/server/src/main.rs index 6d66b60..496e4ea 100644 --- a/panel/server/src/main.rs +++ b/panel/server/src/main.rs @@ -22,6 +22,7 @@ async fn main() -> anyhow::Result<()> { let bind = cfg.bind.clone(); let state = build_state(cfg, pool).await?; bootstrap_admin(&state).await?; + state.livemap.spawn_reaper(); let listener = tokio::net::TcpListener::bind(&bind).await?; tracing::info!("SCOPENET panel v{} listening on http://{bind}", env!("CARGO_PKG_VERSION")); diff --git a/panel/server/src/routes/livemap.rs b/panel/server/src/routes/livemap.rs new file mode 100644 index 0000000..10ebe80 --- /dev/null +++ b/panel/server/src/routes/livemap.rs @@ -0,0 +1,268 @@ +//! Live map endpoints: ingest from game servers, the authenticated Vantage +//! proxy for launchers/admin, and admin status. + +use super::servers::{get_server, GameServer}; +use crate::auth::{AdminUser, AuthUser}; +use crate::error::{AppError, AppResult}; +use crate::livemap::{parse_chunks, Dim, LivePlayer}; +use crate::state::AppState; +use axum::body::{Body, Bytes}; +use axum::extract::ws::{Message, WebSocket, WebSocketUpgrade}; +use axum::extract::{Path, Query, State}; +use axum::http::{header, HeaderMap, HeaderValue, StatusCode}; +use axum::response::{IntoResponse, Response}; +use axum::Json; +use base64::Engine; +use serde::Deserialize; +use serde_json::{json, Value}; + +pub const CHUNK_INTERVAL_SECS: u64 = 20; +pub const PLAYER_INTERVAL_MS: u64 = 1000; +const MAX_PLAYERS: usize = 500; + +/// Settings handed to the game server in `/hello` and `/livemap/config`. +pub fn game_config(enabled: bool) -> Value { + json!({ + "enabled": enabled, + "chunk_interval_secs": CHUNK_INTERVAL_SECS, + "player_interval_ms": PLAYER_INTERVAL_MS, + "max_batch_bytes": 8 * 1024 * 1024, + }) +} + +fn require_enabled(enabled: bool) -> AppResult<()> { + if enabled { + Ok(()) + } else { + Err(AppError::forbidden("Live Map is turned off for this server. Enable it under Servers in the admin panel.")) + } +} + +fn dim_of(s: &str) -> AppResult { + Dim::parse(s).ok_or_else(|| AppError::bad_request("unknown dimension (use overworld, the_nether or the_end)")) +} + +#[derive(Deserialize)] +pub struct DimQuery { + #[serde(default = "overworld")] + dim: String, +} +fn overworld() -> String { + "overworld".into() +} + +// --------------------------------------------------------------------------- +// Game server → panel +// --------------------------------------------------------------------------- + +pub async fn game_config_route(GameServer(server): GameServer) -> Json { + Json(game_config(server.live_map_enabled)) +} + +/// Timestamp tables of the regions the panel already has, base64 per region. +pub async fn game_manifest(GameServer(server): GameServer, State(state): State, Query(q): Query) -> AppResult> { + require_enabled(server.live_map_enabled)?; + let dim = dim_of(&q.dim)?; + let b64 = base64::engine::general_purpose::STANDARD; + let regions: Vec = state.livemap.manifest(server.id, dim).into_iter().map(|(x, z, ts)| json!({ "x": x, "z": z, "timestamps": b64.encode(ts) })).collect(); + Ok(Json(json!({ "dimension": dim.slug(), "regions": regions }))) +} + +pub async fn game_chunks(GameServer(server): GameServer, State(state): State, Query(q): Query, body: Bytes) -> AppResult> { + require_enabled(server.live_map_enabled)?; + let dim = dim_of(&q.dim)?; + let chunks = parse_chunks(&body).map_err(AppError::bad_request)?; + let stored = state.livemap.ingest_chunks(server.id, dim, chunks).await?; + Ok(Json(json!({ "stored": stored }))) +} + +pub async fn game_level(GameServer(server): GameServer, State(state): State, body: Bytes) -> AppResult> { + require_enabled(server.live_map_enabled)?; + // level.dat is a gzip-compressed NBT file. + if body.len() < 18 || body.len() > 2 * 1024 * 1024 || body[0] != 0x1f || body[1] != 0x8b { + return Err(AppError::bad_request("level.dat must be the gzip file from the world folder (up to 2 MiB)")); + } + state.livemap.ingest_level(server.id, body.to_vec()).await?; + Ok(Json(json!({ "ok": true }))) +} + +#[derive(Deserialize)] +pub struct PlayersBody { + players: Vec, +} + +fn sanitize(mut players: Vec) -> Vec { + players.truncate(MAX_PLAYERS); + players.retain(|p| p.x.is_finite() && p.y.is_finite() && p.z.is_finite() && p.x.abs() < 3e7 && p.z.abs() < 3e7); + for p in &mut players { + p.name = p.name.chars().filter(|c| !c.is_control()).take(32).collect(); + p.uuid = p.uuid.chars().filter(|c| c.is_ascii_hexdigit() || *c == '-').take(36).collect(); + p.dimension = p.dimension.chars().take(40).collect(); + } + players +} + +pub async fn game_players(GameServer(server): GameServer, State(state): State, Json(b): Json) -> AppResult> { + require_enabled(server.live_map_enabled)?; + state.livemap.set_players(server.id, sanitize(b.players)); + Ok(Json(json!({ "ok": true }))) +} + +/// Streaming channel for the integration: JSON text frames carry player +/// positions (`{"type":"players","players":[…]}`), binary frames carry chunk +/// batches (`[dimension byte][chunk records…]`, see `parse_chunks`). +pub async fn game_ws(GameServer(server): GameServer, State(state): State, ws: WebSocketUpgrade) -> AppResult { + require_enabled(server.live_map_enabled)?; + Ok(ws.max_message_size(9 * 1024 * 1024).on_upgrade(move |socket| ws_loop(socket, state, server.id))) +} + +async fn ws_loop(mut socket: WebSocket, state: AppState, id: i64) { + let mut last_check = std::time::Instant::now(); + while let Some(Ok(msg)) = socket.recv().await { + if last_check.elapsed() > std::time::Duration::from_secs(30) { + last_check = std::time::Instant::now(); + let enabled = get_server(&state, id).await.map(|s| s.live_map_enabled).unwrap_or(false); + if !enabled { + let _ = socket.send(Message::Text(json!({ "type": "disabled" }).to_string().into())).await; + break; + } + } + match msg { + Message::Text(t) => { + #[derive(Deserialize)] + struct Frame { + #[serde(default)] + players: Vec, + } + if let Ok(f) = serde_json::from_str::(&t) { + state.livemap.set_players(id, sanitize(f.players)); + } + } + Message::Binary(b) => { + let reply = match b.split_first() { + Some((d, rest)) => match Dim::ALL.get(*d as usize).copied().ok_or("unknown dimension").and_then(|dim| parse_chunks(rest).map(|c| (dim, c))) { + Ok((dim, chunks)) => match state.livemap.ingest_chunks(id, dim, chunks).await { + Ok(n) => json!({ "type": "stored", "count": n }), + Err(e) => json!({ "type": "error", "error": e.to_string() }), + }, + Err(e) => json!({ "type": "error", "error": e }), + }, + None => json!({ "type": "error", "error": "empty frame" }), + }; + if socket.send(Message::Text(reply.to_string().into())).await.is_err() { + break; + } + } + Message::Ping(p) => { + let _ = socket.send(Message::Pong(p)).await; + } + Message::Close(_) => break, + _ => {} + } + } +} + +// --------------------------------------------------------------------------- +// Launcher / admin viewers +// --------------------------------------------------------------------------- + +fn label(d: Dim) -> &'static str { + match d { + Dim::Overworld => "Overworld", + Dim::Nether => "The Nether", + Dim::End => "The End", + } +} + +async fn info(state: &AppState, id: i64, enabled: bool) -> Value { + let summary = state.livemap.world_summary(id); + let (generator, assets) = state.livemap.generator_ready(); + let dims: Vec = summary["dimensions"] + .as_array() + .into_iter() + .flatten() + .filter_map(|d| { + let dim = Dim::parse(d["slug"].as_str()?)?; + Some(json!({ "id": dim.id(), "slug": dim.slug(), "label": label(dim), "available": d["regions"].as_u64().unwrap_or(0) > 0 })) + }) + .collect(); + let has_world = summary["level_dat"].as_bool().unwrap_or(false) && dims.iter().any(|d| d["available"] == true); + let message = if !enabled { + "Live Map is off for this server." + } else if !generator || !assets { + "The map generator isn't set up on the panel yet." + } else if !has_world { + "Waiting for the game server to upload its world." + } else { + "" + }; + json!({ + "enabled": enabled, + "ready": enabled && generator && assets && has_world, + "message": message, + "base": format!("/api/livemap/{id}"), + "dimensions": dims, + "players": state.livemap.live_players(id).len(), + }) +} + +/// What the launcher needs to decide whether to offer the Live Map button. +pub async fn viewer_info(_: AuthUser, State(state): State, Path(id): Path) -> AppResult> { + let server = get_server(&state, id).await?; + Ok(Json(info(&state, id, server.live_map_enabled).await)) +} + +pub async fn admin_status(_: AdminUser, State(state): State, Path(id): Path) -> AppResult> { + let server = get_server(&state, id).await?; + let mut v = info(&state, id, server.live_map_enabled).await; + let (generator, assets) = state.livemap.generator_ready(); + let stats = state.livemap.stats(id); + v["generator"] = json!({ "installed": generator, "path": state.livemap.generator_path() }); + v["assets"] = json!({ "configured": assets, "path": state.livemap.assets_path() }); + v["world"] = state.livemap.world_summary(id); + v["stats"] = serde_json::to_value(stats)?; + v["running"] = json!(state.livemap.sidecars_running(id).await); + v["last_error"] = json!(state.livemap.last_error(id)); + v["roster"] = serde_json::to_value(state.livemap.live_players(id))?; + Ok(Json(v)) +} + +pub async fn admin_reset(_: AdminUser, State(state): State, Path(id): Path) -> AppResult> { + get_server(&state, id).await?; + state.livemap.reset(id).await?; + Ok(Json(json!({ "ok": true }))) +} + +/// Authenticated reverse proxy to the per-server Vantage sidecar. Clients +/// use `…/api/livemap/{server}/{dimension}/v1/worlds/default/manifest.json` +/// and everything the manifest references. +pub async fn proxy(_: AuthUser, State(state): State, Path((id, dim, path)): Path<(i64, String, String)>, headers: HeaderMap) -> AppResult { + let server = get_server(&state, id).await?; + require_enabled(server.live_map_enabled)?; + let dim = dim_of(&dim)?; + if path.split('/').any(|s| s == ".." || s.is_empty() || s.contains('\\')) || !path.starts_with("v1/worlds/default/") { + return Err(AppError::not_found("not found")); + } + if path == "v1/worlds/default/players.json" { + let mut resp = Json(state.livemap.players_json(id, dim)).into_response(); + resp.headers_mut().insert(header::CACHE_CONTROL, HeaderValue::from_static("no-store")); + return Ok(resp); + } + let (port, token) = state.livemap.ensure_sidecar(id, dim).await.map_err(|m| AppError::new(StatusCode::SERVICE_UNAVAILABLE, m))?; + let mut req = state.http.get(format!("http://127.0.0.1:{port}/{path}")).bearer_auth(token).timeout(std::time::Duration::from_secs(120)); + if let Some(etag) = headers.get(header::IF_NONE_MATCH) { + req = req.header(header::IF_NONE_MATCH, etag); + } + let upstream = req.send().await.map_err(|e| AppError::new(StatusCode::BAD_GATEWAY, format!("map generator unreachable: {e}")))?; + let mut out = Response::builder().status(upstream.status()); + for name in [header::CONTENT_TYPE, header::ETAG, header::CACHE_CONTROL, header::LAST_MODIFIED] { + if let Some(v) = upstream.headers().get(&name) { + out = out.header(name, v); + } + } + // Tiles are immutable per etag; make sure shared caches never keep them. + Ok(out + .header(header::VARY, "Authorization") + .body(Body::from_stream(upstream.bytes_stream())) + .map_err(|e| AppError::new(StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?) +} diff --git a/panel/server/src/routes/mod.rs b/panel/server/src/routes/mod.rs index 3fafda1..b874bd9 100644 --- a/panel/server/src/routes/mod.rs +++ b/panel/server/src/routes/mod.rs @@ -6,6 +6,7 @@ pub mod connections; pub mod economy; pub mod guilds; pub mod landing; +pub mod livemap; pub mod leveling; pub mod meta; pub mod public; @@ -45,6 +46,7 @@ pub fn api(state: &AppState) -> Router { .route("/avatar/{id}", get(account::avatar)) .route("/servers/public", get(servers::public_servers)) .route("/servers/{id}/leaderboard", get(servers::public_server_leaderboard)) + .route("/servers/{id}/livemap", get(livemap::viewer_info)) .route("/leaderboard", get(servers::global_leaderboard)) .route("/landing", get(landing::public_landing)) .route("/launcher/authlib-injector.json", get(account::authlib_index)) @@ -140,6 +142,7 @@ pub fn api(state: &AppState) -> Router { .route("/servers", get(servers::list).post(servers::create)) .route("/servers/{id}", get(servers::detail).put(servers::update).delete(servers::remove)) .route("/servers/{id}/token", post(servers::regenerate_token)) + .route("/servers/{id}/livemap", get(livemap::admin_status).delete(livemap::admin_reset)) .route("/users/{id}/activity", get(servers::user_activity)) .route("/meta/minecraft", get(meta::minecraft)) .route("/meta/loaders/{loader}", get(meta::loaders)) @@ -183,9 +186,20 @@ pub fn api(state: &AppState) -> Router { .route("/economy/market/buy", post(economy::server_market_buy)) .layer(DefaultBodyLimit::max(2 * 1024 * 1024)); + // Live map ingest: chunk batches are larger than the other game calls. + let livemap_game = Router::new() + .route("/livemap/config", get(livemap::game_config_route)) + .route("/livemap/manifest", get(livemap::game_manifest)) + .route("/livemap/chunks", post(livemap::game_chunks)) + .route("/livemap/level", post(livemap::game_level)) + .route("/livemap/players", post(livemap::game_players)) + .route("/livemap/ws", get(livemap::game_ws)) + .layer(DefaultBodyLimit::max(10 * 1024 * 1024)); + Router::new() + .route("/api/livemap/{id}/{dim}/{*path}", get(livemap::proxy)) .nest("/api/v1", launcher) - .nest("/api/server/v1", game) + .nest("/api/server/v1", game.merge(livemap_game)) .nest("/api/admin", admin) .merge(crate::yggdrasil::routes().layer(DefaultBodyLimit::max(4 * 1024 * 1024))) .layer(axum::middleware::from_fn_with_state(state.clone(), activity::audit)) diff --git a/panel/server/src/routes/servers.rs b/panel/server/src/routes/servers.rs index 37a0708..3a466e6 100644 --- a/panel/server/src/routes/servers.rs +++ b/panel/server/src/routes/servers.rs @@ -62,6 +62,7 @@ fn clip(s: &str, max: usize) -> String { pub struct ServerRow { pub instance_id: String, pub map_url: String, + pub live_map_enabled: bool, pub id: i64, pub name: String, #[serde(skip)] @@ -168,6 +169,7 @@ pub async fn hello( "sync_interval_secs": SYNC_INTERVAL_SECS, "access": server.access, "require_launcher": server.require_launcher, + "live_map": super::livemap::game_config(server.live_map_enabled), }))) } @@ -835,6 +837,8 @@ pub struct ServerInput { require_launcher: bool, #[serde(default)] map_url: String, + #[serde(default)] + live_map_enabled: bool, } fn default_access() -> String { @@ -868,7 +872,7 @@ impl ServerInput { } } -async fn get_server(state: &AppState, id: i64) -> AppResult { +pub async fn get_server(state: &AppState, id: i64) -> AppResult { sqlx::query_as("SELECT * FROM game_servers WHERE id = ?") .bind(id) .fetch_optional(&state.db) @@ -880,7 +884,7 @@ pub async fn create(_: AdminUser, State(state): State, Json(input): Js input.validate()?; let token = new_token(); let id: i64 = sqlx::query_scalar( - "INSERT INTO game_servers (name, token_hash, token_hint, access, allowed_groups, require_launcher, created_at, instance_id, map_url) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) RETURNING id", + "INSERT INTO game_servers (name, token_hash, token_hint, access, allowed_groups, require_launcher, created_at, instance_id, map_url, live_map_enabled) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) RETURNING id", ) .bind(input.name.trim()) .bind(hash_token(&token)) @@ -891,6 +895,7 @@ pub async fn create(_: AdminUser, State(state): State, Json(input): Js .bind(now()) .bind(&input.instance_id) .bind(input.map_url.trim()) + .bind(input.live_map_enabled) .fetch_one(&state.db) .await?; let server = get_server(&state, id).await?; @@ -905,13 +910,14 @@ pub async fn update( ) -> AppResult> { input.validate()?; get_server(&state, id).await?; - sqlx::query("UPDATE game_servers SET name = ?, access = ?, allowed_groups = ?, require_launcher = ?, instance_id = ?, map_url = ? WHERE id = ?") + sqlx::query("UPDATE game_servers SET name = ?, access = ?, allowed_groups = ?, require_launcher = ?, instance_id = ?, map_url = ?, live_map_enabled = ? WHERE id = ?") .bind(input.name.trim()) .bind(&input.access) .bind(serde_json::to_string(&input.allowed_groups).unwrap_or_else(|_| "[]".into())) .bind(input.require_launcher) .bind(&input.instance_id) .bind(input.map_url.trim()) + .bind(input.live_map_enabled) .bind(id) .execute(&state.db) .await?; @@ -1096,6 +1102,7 @@ pub async fn public_servers(State(state): State) -> AppResult, pub http: reqwest::Client, pub login_guard: Arc, + pub livemap: Arc, } diff --git a/panel/server/tests/common/mod.rs b/panel/server/tests/common/mod.rs index a4b0f72..7949600 100644 --- a/panel/server/tests/common/mod.rs +++ b/panel/server/tests/common/mod.rs @@ -22,8 +22,13 @@ pub struct TestApp { } pub async fn setup() -> TestApp { + setup_with(|_, _| {}).await +} + +/// Like [`setup`], letting the test adjust the config (given the data dir). +pub async fn setup_with(tweak: impl FnOnce(&mut Config, &std::path::Path)) -> TestApp { let dir = tempfile::tempdir().unwrap(); - let cfg = Config { + let mut cfg = Config { bind: "127.0.0.1:0".into(), data_dir: dir.path().to_path_buf(), web_dir: dir.path().join("web"), @@ -34,7 +39,11 @@ pub async fn setup() -> TestApp { max_upload_mb: 64, public_url: Some("https://panel.test".into()), trusted_proxies: vec!["127.0.0.1".parse().unwrap()], + vantage_bin: None, + vantage_assets: None, + vantage_args: vec![], }; + tweak(&mut cfg, dir.path()); let pool = db::connect_memory().await.unwrap(); let state = build_state_with_keys(cfg, pool, test_keys()).await.unwrap(); bootstrap_admin(&state).await.unwrap(); diff --git a/panel/server/tests/livemap.rs b/panel/server/tests/livemap.rs new file mode 100644 index 0000000..96ea54a --- /dev/null +++ b/panel/server/tests/livemap.rs @@ -0,0 +1,182 @@ +//! Central live map: ingest from game servers, status for admins, viewer proxy. + +mod common; +use common::*; + +fn record(x: i32, z: i32, ts: u32, data: &[u8]) -> Vec { + let mut b = Vec::new(); + b.extend_from_slice(&x.to_be_bytes()); + b.extend_from_slice(&z.to_be_bytes()); + b.extend_from_slice(&ts.to_be_bytes()); + b.push(2); + b.extend_from_slice(&(data.len() as u32).to_be_bytes()); + b.extend_from_slice(data); + b +} + +async fn post_bytes(t: &TestApp, uri: &str, token: &str, body: Vec) -> (StatusCode, Value) { + let req = Request::builder().method("POST").uri(uri).header("authorization", format!("Bearer {token}")).body(Body::from(body)).unwrap(); + t.send(req).await +} + +#[tokio::test] +async fn live_map_ingest_status_and_proxy() { + let t = setup().await; + let admin = t.login("admin", "supersecret").await; + let (_, v) = t.call("POST", "/api/admin/servers", Some(&admin), Some(json!({"name": "Survival"}))).await; + let id = v["server"]["id"].as_i64().unwrap(); + let token = v["token"].as_str().unwrap().to_string(); + + // Off by default: nothing is accepted. + let (s, _) = post_bytes(&t, "/api/server/v1/livemap/chunks?dim=overworld", &token, record(0, 0, 5, b"abc")).await; + assert_eq!(s, StatusCode::FORBIDDEN); + let (_, hello) = t.call("POST", "/api/server/v1/hello", Some(&token), Some(json!({}))).await; + assert_eq!(hello["live_map"]["enabled"], false); + + let (s, v) = t.call("PUT", &format!("/api/admin/servers/{id}"), Some(&admin), Some(json!({"name": "Survival", "live_map_enabled": true}))).await; + assert_eq!(s, StatusCode::OK, "{v}"); + assert_eq!(v["live_map_enabled"], true); + let (_, hello) = t.call("POST", "/api/server/v1/hello", Some(&token), Some(json!({}))).await; + assert_eq!(hello["live_map"]["enabled"], true); + + // Chunks land in the region mirror and show up in the manifest. + let mut batch = record(1, 2, 100, b"hello-chunk"); + batch.extend(record(-1, 40, 101, b"other-region")); + let (s, v) = post_bytes(&t, "/api/server/v1/livemap/chunks?dim=overworld", &token, batch).await; + assert_eq!(s, StatusCode::OK, "{v}"); + assert_eq!(v["stored"], 2); + let (s, _) = post_bytes(&t, "/api/server/v1/livemap/chunks?dim=../x", &token, record(0, 0, 1, b"x")).await; + assert_eq!(s, StatusCode::BAD_REQUEST); + let (s, _) = post_bytes(&t, "/api/server/v1/livemap/chunks", &token, vec![1, 2, 3]).await; + assert_eq!(s, StatusCode::BAD_REQUEST); + let (_, m) = t.call("GET", "/api/server/v1/livemap/manifest?dim=overworld", Some(&token), None).await; + assert_eq!(m["regions"].as_array().unwrap().len(), 2); + + // level.dat must be a gzip file. + let (s, _) = post_bytes(&t, "/api/server/v1/livemap/level", &token, b"not gzip at all, sorry!".to_vec()).await; + assert_eq!(s, StatusCode::BAD_REQUEST); + let mut gz = vec![0x1f, 0x8b]; + gz.extend([0u8; 30]); + let (s, _) = post_bytes(&t, "/api/server/v1/livemap/level", &token, gz).await; + assert_eq!(s, StatusCode::OK); + + // Players, sanitised. + let (s, _) = t + .call( + "POST", + "/api/server/v1/livemap/players", + Some(&token), + Some(json!({"players": [ + {"uuid": "b50ad385-829d-3141-a216-7e7d7539ba7f", "name": "Notch", "dimension": "minecraft:the_nether", "x": 10.5, "y": 70.0, "z": -4.0, "yaw": 90.0}, + {"uuid": "bad", "name": "Nowhere", "x": 1e99, "y": 0, "z": 0} + ]})), + ) + .await; + assert_eq!(s, StatusCode::OK); + + // Viewers need a session. + let (s, _) = t.call("GET", &format!("/api/v1/servers/{id}/livemap"), None, None).await; + assert_eq!(s, StatusCode::UNAUTHORIZED); + let (s, info) = t.call("GET", &format!("/api/v1/servers/{id}/livemap"), Some(&admin), None).await; + assert_eq!(s, StatusCode::OK); + assert_eq!(info["enabled"], true); + assert_eq!(info["players"], 1); + assert_eq!(info["base"], format!("/api/livemap/{id}")); + assert!(info["dimensions"][0]["available"].as_bool().unwrap()); + // No generator is installed in the test environment. + assert_eq!(info["ready"], false); + + let (_, pj) = t.call("GET", &format!("/api/livemap/{id}/the_nether/v1/worlds/default/players.json"), Some(&admin), None).await; + assert_eq!(pj["players"][0]["name"], "Notch"); + assert_eq!(pj["players"][0]["foreign"], false); + let (_, pj) = t.call("GET", &format!("/api/livemap/{id}/overworld/v1/worlds/default/players.json"), Some(&admin), None).await; + assert_eq!(pj["players"][0]["foreign"], true); + + let (s, e) = t.call("GET", &format!("/api/livemap/{id}/overworld/v1/worlds/default/manifest.json"), Some(&admin), None).await; + assert_eq!(s, StatusCode::SERVICE_UNAVAILABLE); + assert!(e["error"].as_str().unwrap().contains("generator")); + let (s, _) = t.call("GET", &format!("/api/livemap/{id}/overworld/v1/openapi.json"), Some(&admin), None).await; + assert_eq!(s, StatusCode::NOT_FOUND); + let (s, _) = t.call("GET", &format!("/api/livemap/{id}/overworld/v1/worlds/default/manifest.json"), None, None).await; + assert_eq!(s, StatusCode::UNAUTHORIZED); + + let (s, st) = t.call("GET", &format!("/api/admin/servers/{id}/livemap"), Some(&admin), None).await; + assert_eq!(s, StatusCode::OK); + assert_eq!(st["stats"]["chunks_received"], 2); + assert_eq!(st["generator"]["installed"], false); + + // Public list advertises the map; reset wipes it. + let (_, list) = t.call("GET", "/api/v1/servers/public", None, None).await; + assert_eq!(list["servers"][0]["live_map"], true); + let (s, _) = t.call("DELETE", &format!("/api/admin/servers/{id}/livemap"), Some(&admin), None).await; + assert_eq!(s, StatusCode::OK); + let (_, m) = t.call("GET", "/api/server/v1/livemap/manifest?dim=overworld", Some(&token), None).await; + assert!(m["regions"].as_array().unwrap().is_empty()); +} + +/// A stand-in for `vantage server` that answers health checks and serves a +/// manifest, to exercise the supervisor and proxy end to end. +#[cfg(unix)] +#[tokio::test] +async fn proxy_starts_and_reaches_the_generator() { + if std::process::Command::new("python3").arg("--version").output().is_err() { + return; + } + use std::os::unix::fs::PermissionsExt; + let t = setup_with(|cfg, dir| { + let bin = dir.join("fake-vantage"); + std::fs::write( + &bin, + r#"#!/usr/bin/env python3 +import sys, os, http.server +a = sys.argv +port = int(a[a.index('--port') + 1]) +assert a[1] == 'server' and '--assets' in a and '--players-file' in a +token = os.environ['VANTAGE_SERVER_TOKEN'] +class H(http.server.BaseHTTPRequestHandler): + def do_GET(self): + if self.path == '/v1/health': + body = b'ok' + elif self.headers.get('Authorization') != 'Bearer ' + token: + self.send_response(401); self.end_headers(); return + else: + body = b'{"fake":true,"path":"' + self.path.encode() + b'"}' + self.send_response(200); self.send_header('ETag', '"v1"'); self.send_header('Content-Type', 'application/json'); self.end_headers(); self.wfile.write(body) + def log_message(self, *a): pass +http.server.ThreadingHTTPServer(('127.0.0.1', port), H).serve_forever() +"#, + ) + .unwrap(); + std::fs::set_permissions(&bin, std::fs::Permissions::from_mode(0o755)).unwrap(); + std::fs::create_dir_all(dir.join("assets")).unwrap(); + cfg.vantage_bin = Some(bin); + cfg.vantage_assets = Some(dir.join("assets")); + }) + .await; + let admin = t.login("admin", "supersecret").await; + let (_, v) = t.call("POST", "/api/admin/servers", Some(&admin), Some(json!({"name": "S", "live_map_enabled": true}))).await; + let id = v["server"]["id"].as_i64().unwrap(); + let token = v["token"].as_str().unwrap().to_string(); + + let uri = format!("/api/livemap/{id}/overworld/v1/worlds/default/manifest.json"); + let (s, _) = t.call("GET", &uri, Some(&admin), None).await; + assert_eq!(s, StatusCode::SERVICE_UNAVAILABLE, "no world uploaded yet"); + + post_bytes(&t, &format!("/api/server/v1/livemap/chunks?dim=overworld"), &token, record(0, 0, 5, b"abc")).await; + let mut gz = vec![0x1f, 0x8b]; + gz.extend([0u8; 30]); + post_bytes(&t, "/api/server/v1/livemap/level", &token, gz).await; + + let (_, info) = t.call("GET", &format!("/api/v1/servers/{id}/livemap"), Some(&admin), None).await; + assert_eq!(info["ready"], true, "{info}"); + let (s, body) = t.call("GET", &uri, Some(&admin), None).await; + assert_eq!(s, StatusCode::OK, "{body}"); + assert_eq!(body["fake"], true); + assert_eq!(body["path"], "/v1/worlds/default/manifest.json"); + let (_, st) = t.call("GET", &format!("/api/admin/servers/{id}/livemap"), Some(&admin), None).await; + assert_eq!(st["running"][0], "overworld"); + // Reset stops the generator and deletes the mirror. + t.call("DELETE", &format!("/api/admin/servers/{id}/livemap"), Some(&admin), None).await; + let (_, st) = t.call("GET", &format!("/api/admin/servers/{id}/livemap"), Some(&admin), None).await; + assert!(st["running"].as_array().unwrap().is_empty()); +}