Compare commits
8
Commits
222567b735
...
master
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7c52146d8e | ||
|
|
b72279a7fd | ||
|
|
399aced29d | ||
|
|
a5a937694d | ||
|
|
0179da4bbf | ||
|
|
56f511d5f6 | ||
|
|
c4ff921a7c | ||
|
|
c9abf5df77 |
Generated
+2
@@ -1463,6 +1463,7 @@ dependencies = [
|
||||
"rustls-platform-verifier",
|
||||
"serde",
|
||||
"serde_json",
|
||||
"serde_urlencoded",
|
||||
"sync_wrapper",
|
||||
"tokio",
|
||||
"tokio-rustls",
|
||||
@@ -1746,6 +1747,7 @@ name = "server"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"axum",
|
||||
"rand 0.9.2",
|
||||
"reqwest",
|
||||
"serde",
|
||||
"serde_json",
|
||||
|
||||
+2
-1
@@ -5,7 +5,8 @@ edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
axum = { version = "0.8.9", features = ["ws"] }
|
||||
reqwest = { version = "0.13.2", features = ["json"] }
|
||||
rand = "0.9"
|
||||
reqwest = { version = "0.13.2", features = ["json", "query"] }
|
||||
serde = { version = "1.0.228", features = ["serde_derive"] }
|
||||
serde_json = "1.0.149"
|
||||
session-rs = { version = "0.2.1", default-features = false, features = ["axum"] }
|
||||
|
||||
@@ -6,6 +6,10 @@ services:
|
||||
# Coolify's proxy; set the domain in Coolify as https://<domain>:8080.
|
||||
ports:
|
||||
- "${HOST_PORT:-8080}:8080"
|
||||
environment:
|
||||
# Comma-separated client versions served at GET /versions.
|
||||
SUPPORTED_VERSIONS: ${SUPPORTED_VERSIONS:-0.1.3-beta}
|
||||
DEPRECATED_VERSIONS: ${DEPRECATED_VERSIONS:-0.1.2-beta}
|
||||
volumes:
|
||||
- server-data:/data
|
||||
|
||||
|
||||
+2
-4
@@ -44,8 +44,7 @@ pub async fn buy(
|
||||
item_id: String,
|
||||
pool: Arc<SqlitePool>,
|
||||
) -> Result<String, String> {
|
||||
// Lock once
|
||||
let uuid = uuid.lock().await.clone();
|
||||
let uuid = crate::methods::auth::require(&uuid).await?;
|
||||
|
||||
let mut user = user::get(&uuid, &pool).await?;
|
||||
|
||||
@@ -88,8 +87,7 @@ pub async fn equip(
|
||||
item_id: String,
|
||||
pool: Arc<SqlitePool>,
|
||||
) -> Result<String, String> {
|
||||
// Lock once
|
||||
let uuid = uuid.lock().await.clone();
|
||||
let uuid = crate::methods::auth::require(&uuid).await?;
|
||||
|
||||
let mut user = user::get(&uuid, &pool).await?;
|
||||
|
||||
|
||||
@@ -0,0 +1,79 @@
|
||||
use std::{
|
||||
collections::HashSet,
|
||||
time::{Duration, Instant},
|
||||
};
|
||||
|
||||
use tokio::sync::Mutex;
|
||||
|
||||
/// Most players a single `emote` or `send_player` request may target.
|
||||
pub const MAX_TARGETS: usize = 256;
|
||||
|
||||
/// Validates and de-duplicates a request's target UUIDs.
|
||||
pub fn targets(targets: Vec<String>) -> Result<HashSet<String>, String> {
|
||||
if targets.len() > MAX_TARGETS {
|
||||
return Err(format!("Too many targets (max {MAX_TARGETS})"));
|
||||
}
|
||||
|
||||
Ok(targets.into_iter().collect())
|
||||
}
|
||||
|
||||
/// Token bucket: allows bursts of up to `burst` requests, refilled at one
|
||||
/// token per `refill`.
|
||||
pub struct RateLimiter {
|
||||
state: Mutex<(f64, Instant)>,
|
||||
burst: f64,
|
||||
per_sec: f64,
|
||||
}
|
||||
|
||||
impl RateLimiter {
|
||||
pub fn new(burst: u32, refill: Duration) -> Self {
|
||||
Self {
|
||||
state: Mutex::new((burst as f64, Instant::now())),
|
||||
burst: burst as f64,
|
||||
per_sec: 1.0 / refill.as_secs_f64(),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn check(&self) -> Result<(), String> {
|
||||
let mut state = self.state.lock().await;
|
||||
let (tokens, last) = &mut *state;
|
||||
|
||||
let now = Instant::now();
|
||||
*tokens = (*tokens + now.duration_since(*last).as_secs_f64() * self.per_sec).min(self.burst);
|
||||
*last = now;
|
||||
|
||||
if *tokens < 1.0 {
|
||||
return Err("Rate limited, try again shortly".to_string());
|
||||
}
|
||||
|
||||
*tokens -= 1.0;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[tokio::test]
|
||||
async fn allows_a_burst_then_refills() {
|
||||
let limiter = RateLimiter::new(2, Duration::from_millis(50));
|
||||
|
||||
assert!(limiter.check().await.is_ok());
|
||||
assert!(limiter.check().await.is_ok());
|
||||
assert!(limiter.check().await.is_err());
|
||||
|
||||
tokio::time::sleep(Duration::from_millis(60)).await;
|
||||
assert!(limiter.check().await.is_ok());
|
||||
assert!(limiter.check().await.is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn caps_and_dedupes_targets() {
|
||||
let many = (0..=MAX_TARGETS).map(|i| i.to_string()).collect();
|
||||
assert!(targets(many).is_err());
|
||||
|
||||
let dupes = vec!["a".to_string(), "a".to_string(), "b".to_string()];
|
||||
assert_eq!(targets(dupes).unwrap().len(), 2);
|
||||
}
|
||||
}
|
||||
+73
-8
@@ -1,4 +1,5 @@
|
||||
mod cosmetics;
|
||||
mod limits;
|
||||
mod methods;
|
||||
mod types;
|
||||
mod user;
|
||||
@@ -6,12 +7,13 @@ mod user;
|
||||
use std::{collections::HashMap, sync::Arc, time::Duration};
|
||||
|
||||
use axum::{
|
||||
Router,
|
||||
Json, Router,
|
||||
extract::{State, WebSocketUpgrade},
|
||||
http::StatusCode,
|
||||
response::Response,
|
||||
routing::get,
|
||||
};
|
||||
use serde::Serialize;
|
||||
use session_rs::Session;
|
||||
use sqlx::SqlitePool;
|
||||
use tokio::sync::Mutex;
|
||||
@@ -25,6 +27,35 @@ const MAX_MESSAGE_SIZE: usize = 1 << 20;
|
||||
struct AppState {
|
||||
pool: Arc<SqlitePool>,
|
||||
sessions: SessionMap,
|
||||
versions: Arc<VersionManifest>,
|
||||
}
|
||||
|
||||
/// Client versions the server accepts. Clients on a `deprecated` version are
|
||||
/// warned but still connect; clients on neither list refuse to connect.
|
||||
#[derive(Debug, Serialize)]
|
||||
struct VersionManifest {
|
||||
supported: Vec<String>,
|
||||
deprecated: Vec<String>,
|
||||
}
|
||||
|
||||
impl VersionManifest {
|
||||
/// Reads comma-separated `SUPPORTED_VERSIONS` and `DEPRECATED_VERSIONS`.
|
||||
fn from_env() -> Self {
|
||||
let list = |key: &str, default: &str| -> Vec<String> {
|
||||
std::env::var(key)
|
||||
.unwrap_or_else(|_| default.to_string())
|
||||
.split(',')
|
||||
.map(str::trim)
|
||||
.filter(|v| !v.is_empty())
|
||||
.map(String::from)
|
||||
.collect()
|
||||
};
|
||||
|
||||
Self {
|
||||
supported: list("SUPPORTED_VERSIONS", "0.1.3-beta"),
|
||||
deprecated: list("DEPRECATED_VERSIONS", "0.1.2-beta"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
@@ -36,14 +67,19 @@ async fn main() -> std::io::Result<()> {
|
||||
std::process::exit(healthcheck(&bind_addr).await);
|
||||
}
|
||||
|
||||
let versions = VersionManifest::from_env();
|
||||
println!("Client versions: {versions:?}");
|
||||
|
||||
let state = AppState {
|
||||
pool: Arc::new(user::init_db().await),
|
||||
sessions: Arc::new(Mutex::new(HashMap::new())),
|
||||
versions: Arc::new(versions),
|
||||
};
|
||||
|
||||
let app = Router::new()
|
||||
.route("/", get(ws))
|
||||
.route("/health", get(health))
|
||||
.route("/versions", get(versions_manifest))
|
||||
.with_state(state);
|
||||
|
||||
let listener = tokio::net::TcpListener::bind(&bind_addr).await?;
|
||||
@@ -76,10 +112,19 @@ async fn health(State(state): State<AppState>) -> (StatusCode, &'static str) {
|
||||
}
|
||||
}
|
||||
|
||||
async fn versions_manifest(State(state): State<AppState>) -> Json<Arc<VersionManifest>> {
|
||||
Json(state.versions)
|
||||
}
|
||||
|
||||
async fn register_handlers(session: &Session, state: AppState) {
|
||||
let AppState { pool, sessions } = state;
|
||||
let AppState { pool, sessions, .. } = state;
|
||||
let uuid = Arc::new(Mutex::new(String::new()));
|
||||
let name = Arc::new(Mutex::new(String::new()));
|
||||
let pending: methods::auth::PendingChallenge = Arc::new(Mutex::new(None));
|
||||
|
||||
// Per-connection limits on the methods that fan out to other players.
|
||||
let emote_limit = Arc::new(limits::RateLimiter::new(5, Duration::from_secs(1)));
|
||||
let send_player_limit = Arc::new(limits::RateLimiter::new(10, Duration::from_secs(1)));
|
||||
|
||||
session
|
||||
.on_close({
|
||||
@@ -106,20 +151,30 @@ async fn register_handlers(session: &Session, state: AppState) {
|
||||
.await;
|
||||
|
||||
session
|
||||
.on_request::<methods::Auth, _>({
|
||||
.on_request::<methods::AuthChallenge, _>({
|
||||
let uuid = Arc::clone(&uuid);
|
||||
let pending = Arc::clone(&pending);
|
||||
|
||||
move |_, username| methods::auth::challenge(uuid.clone(), pending.clone(), username)
|
||||
})
|
||||
.await;
|
||||
|
||||
session
|
||||
.on_request::<methods::AuthVerify, _>({
|
||||
let pool = Arc::clone(&pool);
|
||||
let uuid = Arc::clone(&uuid);
|
||||
let sessions = Arc::clone(&sessions);
|
||||
let name = Arc::clone(&name);
|
||||
let pending = Arc::clone(&pending);
|
||||
let session = session.clone();
|
||||
|
||||
move |_, token| {
|
||||
methods::auth::authenticate(
|
||||
move |_, ()| {
|
||||
methods::auth::verify(
|
||||
sessions.clone(),
|
||||
session.clone(),
|
||||
name.clone(),
|
||||
uuid.clone(),
|
||||
token,
|
||||
pending.clone(),
|
||||
pool.clone(),
|
||||
)
|
||||
}
|
||||
@@ -195,7 +250,14 @@ async fn register_handlers(session: &Session, state: AppState) {
|
||||
let sessions = Arc::clone(&sessions);
|
||||
let uuid = Arc::clone(&uuid);
|
||||
|
||||
move |_, emote| methods::emote::send_emote(Arc::clone(&sessions), Arc::clone(&uuid), emote)
|
||||
move |_, emote| {
|
||||
methods::emote::send_emote(
|
||||
Arc::clone(&sessions),
|
||||
Arc::clone(&uuid),
|
||||
Arc::clone(&emote_limit),
|
||||
emote,
|
||||
)
|
||||
}
|
||||
})
|
||||
.await;
|
||||
|
||||
@@ -204,7 +266,9 @@ async fn register_handlers(session: &Session, state: AppState) {
|
||||
let sessions = Arc::clone(&sessions);
|
||||
let pool = Arc::clone(&pool);
|
||||
|
||||
move |_, uuid| methods::user::get_user(sessions.clone(), uuid, pool.clone())
|
||||
let uuid = Arc::clone(&uuid);
|
||||
|
||||
move |_, target| methods::user::get_user(sessions.clone(), uuid.clone(), target, pool.clone())
|
||||
})
|
||||
.await;
|
||||
|
||||
@@ -220,6 +284,7 @@ async fn register_handlers(session: &Session, state: AppState) {
|
||||
sessions.clone(),
|
||||
name.clone(),
|
||||
uuid.clone(),
|
||||
Arc::clone(&send_player_limit),
|
||||
targets.targets,
|
||||
pool.clone(),
|
||||
)
|
||||
|
||||
+79
-22
@@ -2,63 +2,120 @@ use std::sync::Arc;
|
||||
|
||||
use session_rs::session::Session;
|
||||
use sqlx::SqlitePool;
|
||||
use tokio::sync::Mutex;
|
||||
|
||||
use crate::{
|
||||
types::{SessionMap, UUID},
|
||||
types::{MinecraftAuthResponse, SessionMap, UUID},
|
||||
user::User,
|
||||
};
|
||||
|
||||
pub async fn authenticate(
|
||||
const HAS_JOINED_URL: &str = "https://sessionserver.mojang.com/session/minecraft/hasJoined";
|
||||
|
||||
/// A challenge issued to a connection that hasn't authenticated yet.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct Challenge {
|
||||
username: String,
|
||||
server_id: String,
|
||||
}
|
||||
|
||||
pub type PendingChallenge = Arc<Mutex<Option<Challenge>>>;
|
||||
|
||||
/// Returns the connection's authenticated UUID, or an error before `auth_verify`.
|
||||
pub async fn require(uuid: &UUID) -> Result<String, String> {
|
||||
let uuid = uuid.lock().await;
|
||||
|
||||
if uuid.is_empty() {
|
||||
return Err("Not authenticated".to_string());
|
||||
}
|
||||
|
||||
Ok(uuid.clone())
|
||||
}
|
||||
|
||||
/// Minecraft usernames: 1–16 characters of letters, digits and underscores.
|
||||
fn is_valid_username(name: &str) -> bool {
|
||||
(1..=16).contains(&name.len()) && name.bytes().all(|b| b.is_ascii_alphanumeric() || b == b'_')
|
||||
}
|
||||
|
||||
/// Issues a random server id for the client to `join` with at Mojang.
|
||||
pub async fn challenge(uuid: UUID, pending: PendingChallenge, username: String) -> Result<String, String> {
|
||||
if !uuid.lock().await.is_empty() {
|
||||
return Err("Already authenticated".to_string());
|
||||
}
|
||||
|
||||
if !is_valid_username(&username) {
|
||||
return Err("Invalid username".to_string());
|
||||
}
|
||||
|
||||
let server_id: String = rand::random::<[u8; 20]>()
|
||||
.iter()
|
||||
.map(|b| format!("{b:02x}"))
|
||||
.collect();
|
||||
|
||||
*pending.lock().await = Some(Challenge {
|
||||
username,
|
||||
server_id: server_id.clone(),
|
||||
});
|
||||
|
||||
Ok(server_id)
|
||||
}
|
||||
|
||||
/// Confirms with Mojang that the challenged player joined with our server id.
|
||||
pub async fn verify(
|
||||
sessions: SessionMap,
|
||||
session: Session,
|
||||
name: UUID,
|
||||
uuid: UUID,
|
||||
session_token: String,
|
||||
pending: PendingChallenge,
|
||||
pool: Arc<SqlitePool>,
|
||||
) -> Result<User, String> {
|
||||
if !uuid.lock().await.is_empty() {
|
||||
return Err(format!("Already authenticated"));
|
||||
return Err("Already authenticated".to_string());
|
||||
}
|
||||
|
||||
let client = reqwest::Client::new();
|
||||
// Single use: a failed verification needs a fresh challenge.
|
||||
let Some(challenge) = pending.lock().await.take() else {
|
||||
return Err("No pending challenge, call auth_challenge first".to_string());
|
||||
};
|
||||
|
||||
let response = client
|
||||
.get("https://api.minecraftservices.com/minecraft/profile")
|
||||
.bearer_auth(&session_token)
|
||||
let response = reqwest::Client::new()
|
||||
.get(HAS_JOINED_URL)
|
||||
.query(&[
|
||||
("username", challenge.username.as_str()),
|
||||
("serverId", challenge.server_id.as_str()),
|
||||
])
|
||||
.send()
|
||||
.await
|
||||
.map_err(|_| "Failed to validate session".to_string())?;
|
||||
.map_err(|_| "Failed to reach the Mojang session server".to_string())?;
|
||||
|
||||
// 204 No Content: the player didn't join with this server id.
|
||||
if response.status() == reqwest::StatusCode::NO_CONTENT {
|
||||
return Err("Session not verified by Mojang".to_string());
|
||||
}
|
||||
|
||||
if !response.status().is_success() {
|
||||
return Err(format!(
|
||||
"Authentication failed with code {}",
|
||||
"Mojang session server returned {}",
|
||||
response.status()
|
||||
));
|
||||
}
|
||||
|
||||
let auth = response
|
||||
.json::<crate::types::MinecraftAuthResponse>()
|
||||
.json::<MinecraftAuthResponse>()
|
||||
.await
|
||||
.map_err(|_| "Unable to parse auth response".to_string())
|
||||
.and_then(|auth| {
|
||||
.map_err(|_| "Unable to parse Mojang response".to_string())?;
|
||||
let id = crate::types::format_uuid(&auth.id)?;
|
||||
Ok(crate::types::MinecraftAuthResponse {
|
||||
name: auth.name,
|
||||
id,
|
||||
})
|
||||
})?;
|
||||
|
||||
*uuid.lock().await = auth.id.clone();
|
||||
*uuid.lock().await = id.clone();
|
||||
*name.lock().await = auth.name.clone();
|
||||
|
||||
sessions
|
||||
.lock()
|
||||
.await
|
||||
.entry(auth.id.clone())
|
||||
.entry(id.clone())
|
||||
.or_default()
|
||||
.insert(session);
|
||||
|
||||
println!("{:?}", auth);
|
||||
println!("Authenticated {} ({id})", auth.name);
|
||||
|
||||
crate::user::get_put(&auth.id, &pool).await
|
||||
crate::user::get_put(&id, &pool).await
|
||||
}
|
||||
|
||||
+13
-3
@@ -1,16 +1,26 @@
|
||||
use crate::types::{SessionMap, UUID};
|
||||
use std::sync::Arc;
|
||||
|
||||
use crate::{
|
||||
limits::{self, RateLimiter},
|
||||
types::{SessionMap, UUID},
|
||||
};
|
||||
|
||||
pub async fn send_emote(
|
||||
sessions: SessionMap,
|
||||
uuid: UUID,
|
||||
limit: Arc<RateLimiter>,
|
||||
emote: crate::types::EmoteRequest,
|
||||
) -> Result<(), String> {
|
||||
let from = crate::methods::auth::require(&uuid).await?;
|
||||
let targets = limits::targets(emote.targets)?;
|
||||
limit.check().await?;
|
||||
|
||||
let emote_event = crate::types::EventEmote {
|
||||
from: uuid.lock().await.clone(),
|
||||
from,
|
||||
emote: emote.emote,
|
||||
};
|
||||
|
||||
for target in emote.targets {
|
||||
for target in targets {
|
||||
if let Some(sessions) = sessions.lock().await.get_mut(&target) {
|
||||
let emote_event = emote_event.clone();
|
||||
let mut bad_sessions = Vec::new();
|
||||
|
||||
+17
-3
@@ -10,12 +10,26 @@ use crate::{
|
||||
user::User,
|
||||
};
|
||||
|
||||
/// Step 1 of authentication: the client sends its username and gets a random
|
||||
/// server id to pass to Mojang's session server `join` endpoint.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct Auth;
|
||||
pub struct AuthChallenge;
|
||||
|
||||
impl Method for Auth {
|
||||
const NAME: &'static str = "auth";
|
||||
impl Method for AuthChallenge {
|
||||
const NAME: &'static str = "auth_challenge";
|
||||
type Request = String;
|
||||
type Response = String;
|
||||
type Error = String;
|
||||
}
|
||||
|
||||
/// Step 2 of authentication: after joining, the server confirms the player
|
||||
/// with Mojang's `hasJoined` endpoint. The access token never reaches us.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct AuthVerify;
|
||||
|
||||
impl Method for AuthVerify {
|
||||
const NAME: &'static str = "auth_verify";
|
||||
type Request = ();
|
||||
type Response = User;
|
||||
type Error = String;
|
||||
}
|
||||
|
||||
+9
-1
@@ -3,6 +3,7 @@ use std::sync::Arc;
|
||||
use sqlx::SqlitePool;
|
||||
|
||||
use crate::{
|
||||
limits::{self, RateLimiter},
|
||||
methods,
|
||||
types::{PlayerStream, SessionMap, UUID},
|
||||
user::User,
|
||||
@@ -10,9 +11,12 @@ use crate::{
|
||||
|
||||
pub async fn get_user(
|
||||
sessions: SessionMap,
|
||||
caller: UUID,
|
||||
uuid: String,
|
||||
pool: Arc<SqlitePool>,
|
||||
) -> Result<Option<User>, String> {
|
||||
crate::methods::auth::require(&caller).await?;
|
||||
|
||||
if !sessions.lock().await.contains_key(&uuid) {
|
||||
return Ok(None);
|
||||
}
|
||||
@@ -24,11 +28,15 @@ pub async fn send_user(
|
||||
sessions: SessionMap,
|
||||
name: UUID,
|
||||
uuid: UUID,
|
||||
limit: Arc<RateLimiter>,
|
||||
targets: Vec<String>,
|
||||
pool: Arc<SqlitePool>,
|
||||
) -> Result<(), String> {
|
||||
let uuid = crate::methods::auth::require(&uuid).await?;
|
||||
let targets = limits::targets(targets)?;
|
||||
limit.check().await?;
|
||||
|
||||
println!("send player {targets:?}");
|
||||
let uuid = uuid.lock().await.to_string();
|
||||
|
||||
let user = PlayerStream {
|
||||
player: crate::user::get(&uuid, pool.as_ref()).await?,
|
||||
|
||||
Reference in New Issue
Block a user