Add central live map backend: chunk/player ingest, Vantage sidecar supervisor and proxy

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01DjMbLQujBHunCCu5GpsHaT
This commit is contained in:
Claude committed 2026-09-30 17:09:31 +00:00
1 parent 004862ea88
commit 8acb4a64e6
13 files changed
+1255 -7

No files matched your search

+1 -1
View File
@@ -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"] }
+9
View File
@@ -14,6 +14,12 @@ pub struct Config {
pub max_upload_mb: usize,
pub public_url: Option<String>,
pub trusted_proxies: Vec<std::net::IpAddr>,
/// Vantage generator binary (defaults to `vantage` on PATH or /opt/vantage).
pub vantage_bin: Option<PathBuf>,
/// Minecraft client assets directory handed to Vantage.
pub vantage_assets: Option<PathBuf>,
/// Extra arguments appended to `vantage server`.
pub vantage_args: Vec<String>,
}
fn var(name: &str) -> Option<String> {
@@ -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(),
}
}
+3
View File
@@ -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<SqlitePool> {
+2
View File
@@ -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),
})
}
+686
View File
@@ -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 `<data>/livemap/<server>/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<Dim> {
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<u8>,
}
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<usize> {
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<u8> = 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<Vec<u8>> {
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<Vec<Chunk>, &'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<LivePlayer>,
}
// ---------------------------------------------------------------------------
// Sidecar supervisor
// ---------------------------------------------------------------------------
struct Sidecar {
child: tokio::process::Child,
port: u16,
token: String,
last_used: Instant,
log: std::sync::Arc<Mutex<VecDeque<String>>>,
}
#[derive(Default, Clone, Serialize)]
pub struct ServerStats {
pub chunks_received: u64,
pub last_chunk_at: Option<String>,
pub last_players_at: Option<String>,
}
pub struct LiveMap {
root: PathBuf,
bin: Option<PathBuf>,
assets: Option<PathBuf>,
extra_args: Vec<String>,
players: Mutex<HashMap<i64, Snapshot>>,
stats: Mutex<HashMap<i64, ServerStats>>,
sidecars: tokio::sync::Mutex<HashMap<(i64, Dim), Sidecar>>,
last_error: Mutex<HashMap<(i64, Dim), String>>,
write_lock: tokio::sync::Mutex<()>,
players_written: Mutex<HashMap<i64, Instant>>,
}
fn find_bin(configured: Option<&Path>) -> Option<PathBuf> {
let candidates: Vec<PathBuf> = match configured {
Some(p) => vec![p.to_path_buf()],
None => {
let mut v: Vec<PathBuf> = 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<PathBuf>, assets: Option<PathBuf>, extra_args: Vec<String>) -> 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<String> {
self.bin.as_ref().map(|p| p.display().to_string())
}
pub fn assets_path(&self) -> Option<String> {
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<Chunk>) -> std::io::Result<usize> {
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<Chunk>> = 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(&region_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<u8>) -> 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<u8>)> {
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::<i32>(), b.parse::<i32>()) 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<LivePlayer>) {
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<LivePlayer> {
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<String> {
self.sidecars.lock().await.keys().filter(|(s, _)| *s == id).map(|(_, d)| d.slug().to_string()).collect()
}
pub fn last_error(&self, id: i64) -> Option<String> {
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<Mutex<VecDeque<String>>> = 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::<Vec<_>>().into_iter().rev().collect::<Vec<_>>().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<Self>) {
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: tokio::io::AsyncRead + Unpin + Send + 'static>(r: R, log: std::sync::Arc<Mutex<VecDeque<String>>>) {
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<Dim>) -> 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<String> {
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<u8>)> {
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);
}
}
+1
View File
@@ -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"));
+268
View File
@@ -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> {
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<Value> {
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<AppState>, Query(q): Query<DimQuery>) -> AppResult<Json<Value>> {
require_enabled(server.live_map_enabled)?;
let dim = dim_of(&q.dim)?;
let b64 = base64::engine::general_purpose::STANDARD;
let regions: Vec<Value> = 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<AppState>, Query(q): Query<DimQuery>, body: Bytes) -> AppResult<Json<Value>> {
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<AppState>, body: Bytes) -> AppResult<Json<Value>> {
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<LivePlayer>,
}
fn sanitize(mut players: Vec<LivePlayer>) -> Vec<LivePlayer> {
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<AppState>, Json(b): Json<PlayersBody>) -> AppResult<Json<Value>> {
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<AppState>, ws: WebSocketUpgrade) -> AppResult<Response> {
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<LivePlayer>,
}
if let Ok(f) = serde_json::from_str::<Frame>(&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<Value> = 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<AppState>, Path(id): Path<i64>) -> AppResult<Json<Value>> {
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<AppState>, Path(id): Path<i64>) -> AppResult<Json<Value>> {
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<AppState>, Path(id): Path<i64>) -> AppResult<Json<Value>> {
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<AppState>, Path((id, dim, path)): Path<(i64, String, String)>, headers: HeaderMap) -> AppResult<Response> {
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()))?)
}
+15 -1
View File
@@ -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<AppState> {
.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<AppState> {
.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<AppState> {
.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))
+10 -3
View File
@@ -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<ServerRow> {
pub async fn get_server(state: &AppState, id: i64) -> AppResult<ServerRow> {
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<AppState>, 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<AppState>, 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<Json<ServerView>> {
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<AppState>) -> AppResult<Json<Val
"name": s.name,
"instance_id": s.instance_id,
"map_url": s.map_url,
"live_map": s.live_map_enabled,
"online": is_online,
"players_online": if is_online { s.online_count } else { 0 },
"players_max": s.max_players,
+1
View File
@@ -13,4 +13,5 @@ pub struct AppState {
pub ygg: Arc<crate::yggdrasil::keys::Keys>,
pub http: reqwest::Client,
pub login_guard: Arc<LoginGuard>,
pub livemap: Arc<crate::livemap::LiveMap>,
}
+10 -1
View File
@@ -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();
+182
View File
@@ -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<u8> {
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<u8>) -> (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());
}