fix: CI workflow, password reset flow, health/limits in services, session refactor, dead mixin stub cleanup

This commit is contained in:
loki5512344 2026-10-09 19:21:17 +02:00
parent 45dd592c40
commit f4e15b45c9
Signed by: boba
GPG key ID: 253067914055423B
95 changed files with 899 additions and 185 deletions

View file

@ -21,6 +21,10 @@ ACCOUNTS_GRPC_URL=http://127.0.0.1:50051
CONFIGS_HTTP_URL=http://127.0.0.1:8082
CHAT_HTTP_URL=http://127.0.0.1:8083
SITE_ORIGIN=http://localhost:5173
# true ONLY when the gateway is behind a reverse proxy that sets
# X-Forwarded-For (nginx in deploy/): per-IP rate limits then key on the
# real client. docker-compose.prod.yml overrides this to "true" itself;
# direct exposure must keep it false, or clients could spoof the header.
TRUST_PROXY=false
# Directory the gateway serves under GET /downloads/* (lovisual.jar lives here;
# docker-compose.prod.yml mounts it as /srv/downloads and sets the variable itself)

15
backend/Cargo.lock generated
View file

@ -872,6 +872,7 @@ dependencies = [
"serde",
"serde_json",
"tokio",
"tower-http 0.7.1",
"tracing",
"tracing-subscriber",
"uuid",
@ -941,6 +942,7 @@ dependencies = [
"sqlx",
"tokio",
"tonic",
"tower-http 0.7.1",
"tracing",
"tracing-subscriber",
"uuid",
@ -2330,6 +2332,15 @@ dependencies = [
"hashbrown 0.17.1",
]
[[package]]
name = "matchers"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d1525a2a28c7f4fa0fc98bb91ae755d1e2d1505079e05539e35bc876b5d65ae9"
dependencies = [
"regex-automata",
]
[[package]]
name = "matchit"
version = "0.8.4"
@ -4174,10 +4185,14 @@ version = "0.3.23"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "cb7f578e5945fb242538965c2d0b04418d38ec25c79d160cd279bf0731c8d319"
dependencies = [
"matchers",
"nu-ansi-term",
"once_cell",
"regex-automata",
"sharded-slab",
"smallvec",
"thread_local",
"tracing",
"tracing-core",
"tracing-log",
]

View file

@ -9,10 +9,10 @@ path = "src/lib.rs"
[dependencies]
axum = { version = "0.8", features = ["multipart", "macros"] }
tokio = { version = "1", features = ["rt-multi-thread", "macros"] }
tokio = { version = "1", features = ["rt-multi-thread", "macros", "sync"] }
tower-http = { version = "0.7", features = ["trace", "cors"] }
tracing = "0.1"
tracing-subscriber = "0.3"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
serde = { version = "1", features = ["derive"] }
serde_json = "1"
sqlx = { version = "0.9", default-features = false, features = ["runtime-tokio", "tls-rustls", "postgres", "uuid", "chrono", "macros", "migrate"] }
@ -36,4 +36,5 @@ tonic = "0.14"
[dev-dependencies]
axum-test = "21"
aws-sdk-s3 = "1"
tokio-stream = { version = "0.1", features = ["net"] }

View file

@ -0,0 +1,13 @@
-- Single-use password reset tokens (stored hashed, like refresh tokens).
-- A token is valid for 30 minutes; issuing a new request for the same
-- account supersedes any token still pending from an earlier request.
CREATE TABLE password_reset_tokens (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
account_id UUID NOT NULL REFERENCES accounts(id) ON DELETE CASCADE,
token_hash TEXT NOT NULL UNIQUE,
expires_at TIMESTAMPTZ NOT NULL,
used_at TIMESTAMPTZ,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE INDEX idx_password_reset_tokens_account_id ON password_reset_tokens(account_id);

View file

@ -24,6 +24,17 @@ pub fn validate_register(
password: &str,
nick: &str,
) -> Result<(String, String), AppError> {
let email = normalize_email(email)?;
let nick = validate_nick(nick)?;
validate_password(password)?;
Ok((email, nick))
}
/// Trims and shape-checks an email address the same way for every auth flow
/// that accepts one (register, password reset request). The lookup itself is
/// case-insensitive: `accounts_email_lower_idx` is a unique index on
/// `lower(email)`, and `repo::find_by_email` lowercases the query value.
pub fn normalize_email(email: &str) -> Result<String, AppError> {
let email = email.trim();
if email.is_empty() {
return Err(AppError::Validation("email must not be empty".into()));
@ -42,7 +53,10 @@ pub fn validate_register(
"email must have exactly one '@' with non-empty parts".into(),
));
}
Ok(email.to_string())
}
fn validate_nick(nick: &str) -> Result<String, AppError> {
let nick = nick.trim();
let nick_len = nick.chars().count();
if nick_len == 0 || nick_len > 32 {
@ -55,7 +69,11 @@ pub fn validate_register(
"nick must not contain control characters".into(),
));
}
Ok(nick.to_string())
}
/// Password policy shared by registration and password reset.
pub fn validate_password(password: &str) -> Result<(), AppError> {
if password.chars().count() < 8 {
return Err(AppError::Validation(
"password must be at least 8 characters".into(),
@ -66,8 +84,7 @@ pub fn validate_register(
"password must be at most {MAX_PASSWORD_BYTES} bytes"
)));
}
Ok((email.to_string(), nick.to_string()))
Ok(())
}
#[cfg(test)]

View file

@ -1,6 +1,6 @@
use crate::accounts::model::validate_register;
use crate::accounts::repo;
use crate::auth::{password, tokens};
use crate::auth::{password, reset, tokens};
use crate::error::{AppError, AppJson};
use axum::{Json, extract::State, http::StatusCode};
use axum_extra::extract::cookie::CookieJar;
@ -13,6 +13,9 @@ pub struct AuthState {
pub jwt_secret: String,
pub cookie_secure: bool,
pub hasher: password::PasswordHasher,
/// Reset-email delivery back-end: Log in production (until a real
/// provider lands), Queue in tests.
pub mail: reset::MailBox,
// A real Argon2id hash of a throwaway string. `login` verifies against
// it when the email is unknown so that "no such account" costs the same
// ~100ms as "wrong password" — otherwise response time leaks which
@ -22,6 +25,15 @@ pub struct AuthState {
impl AuthState {
pub fn new(pool: sqlx::PgPool, jwt_secret: String, cookie_secure: bool) -> Self {
Self::with_mail(pool, jwt_secret, cookie_secure, reset::MailBox::Log)
}
pub fn with_mail(
pool: sqlx::PgPool,
jwt_secret: String,
cookie_secure: bool,
mail: reset::MailBox,
) -> Self {
// Hashing a fixed, short constant with fixed valid params cannot
// fail; this is not user input, so the expect is a startup invariant.
let dummy_hash = password::hash_password("timing-equalizer-not-a-real-password")
@ -31,6 +43,7 @@ impl AuthState {
jwt_secret,
cookie_secure,
hasher: password::PasswordHasher::new(),
mail,
dummy_hash,
}
}

View file

@ -1,3 +1,4 @@
pub mod handlers;
pub mod password;
pub mod reset;
pub mod tokens;

View file

@ -0,0 +1,86 @@
use crate::accounts::{model, repo};
use crate::auth::reset;
use crate::error::{AppError, AppJson};
use axum::{Json, extract::State, http::StatusCode};
use serde::{Deserialize, Serialize};
use super::super::handlers::AuthState;
#[derive(Deserialize)]
pub struct ForgotPasswordRequest {
pub email: String,
}
#[derive(Serialize)]
pub struct AcceptedResponse {
pub status: &'static str,
}
/// POST /auth/forgot-password { email }
///
/// The response is identical whether or not the email is registered: no
/// oracle for account enumeration. Delivery happens out of band; the gateway
/// rate-limits this route per IP (3/hour) to keep the mailer from being
/// weaponised.
pub async fn forgot_password(
State(state): State<AuthState>,
AppJson(req): AppJson<ForgotPasswordRequest>,
) -> Result<(StatusCode, Json<AcceptedResponse>), AppError> {
let email = model::normalize_email(&req.email)?;
if let Some(account) = repo::find_by_email(&state.pool, &email).await? {
let token = reset::issue_reset_token(&state.pool, account.id).await?;
state.mail.send(&email, &token);
}
Ok((
StatusCode::ACCEPTED,
Json(AcceptedResponse {
status: "reset email sent if the account exists",
}),
))
}
#[derive(Deserialize)]
pub struct ResetPasswordRequest {
pub token: String,
pub new_password: String,
}
/// POST /auth/reset-password { token, new_password }
///
/// Consumes the token atomically, rewrites the password hash and revokes
/// every refresh session of the account. A bad token is a plain 400 with no
/// distinction between unknown, expired and already-used — all three are the
/// same "try again" situation from an attacker's point of view.
pub async fn reset_password(
State(state): State<AuthState>,
AppJson(req): AppJson<ResetPasswordRequest>,
) -> Result<Json<AcceptedResponse>, AppError> {
if !reset::is_well_formed_token(&req.token) {
return Err(AppError::Validation("malformed reset token".into()));
}
model::validate_password(&req.new_password)?;
let account_id = reset::consume_reset_token(&state.pool, &req.token)
.await?
.ok_or_else(|| AppError::Validation("reset token is invalid or expired".into()))?;
let hash = state
.hasher
.hash(req.new_password)
.await
.map_err(AppError::Internal)?;
let updated = sqlx::query("UPDATE accounts SET password_hash = $2 WHERE id = $1")
.bind(account_id)
.bind(&hash)
.execute(&state.pool)
.await?;
if updated.rows_affected() != 1 {
// The token row referenced a cascade-deleted account.
return Err(AppError::Validation(
"reset token is invalid or expired".into(),
));
}
reset::revoke_all_sessions(&state.pool, account_id).await?;
Ok(Json(AcceptedResponse {
status: "password updated",
}))
}

View file

@ -0,0 +1,129 @@
pub mod handlers;
use sqlx::PgPool;
use tokio::sync::mpsc::UnboundedSender;
use uuid::Uuid;
use crate::auth::tokens::{hash_token, new_opaque_token};
/// One outgoing reset email. `token` is the plaintext token: the only place
/// it ever exists outside the response of `issue_reset_token`.
#[derive(Debug, Clone)]
pub struct Mail {
pub to: String,
pub token: String,
}
/// Delivery back-end for reset emails.
#[derive(Clone)]
pub enum MailBox {
/// Development delivery: a structured log line carrying the token, which
/// local/dev stacks pick up from the container logs. When a real provider
/// is wired in (lettre/SES/anything), this variant is the single switch
/// point - the endpoint contract does not change.
Log,
/// In-process queue: tests (and a future in-process mail worker) receive
/// every mail exactly as the handler produced it.
Queue(UnboundedSender<Mail>),
}
impl MailBox {
pub fn send(&self, to: &str, token: &str) {
match self {
MailBox::Log => {
tracing::info!(
account = %to,
reset_token = %token,
"password reset requested; deliver the reset link to the account owner"
);
}
MailBox::Queue(tx) => {
let _ = tx.send(Mail {
to: to.to_string(),
token: token.to_string(),
});
}
}
}
}
/// Reset links must be used quickly: long windows turn a leaked email into
/// an account takeover. 30 minutes is the common industry compromise.
const TTL_MINUTES: i64 = 30;
pub const TOKEN_PREFIX: &str = "lvpr_";
/// A reset token is base64url like every other opaque token in this service
/// ("lvpr_" prefix + 43 chars), so the length check alone filters out most
/// junk before the database is ever touched.
pub fn is_well_formed_token(token: &str) -> bool {
token.len() == TOKEN_PREFIX.len() + 43 && token.starts_with(TOKEN_PREFIX)
}
/// Invalidates any token still pending for the account, then stores the hash
/// of a fresh one. Returns the plaintext token for the mailer only — it is
/// never persisted in clear form.
pub async fn issue_reset_token(pool: &PgPool, account_id: Uuid) -> Result<String, sqlx::Error> {
sqlx::query(
"UPDATE password_reset_tokens SET used_at = now()
WHERE account_id = $1 AND used_at IS NULL",
)
.bind(account_id)
.execute(pool)
.await?;
let token = new_opaque_token(TOKEN_PREFIX);
sqlx::query(
"INSERT INTO password_reset_tokens (account_id, token_hash, expires_at)
VALUES ($1, $2, now() + make_interval(mins => $3::int))",
)
.bind(account_id)
.bind(hash_token(&token))
.bind(TTL_MINUTES)
.execute(pool)
.await?;
Ok(token)
}
/// Atomically marks the token as used and returns its account. The single
/// conditional UPDATE makes double-spend impossible even for concurrent
/// callers: exactly one of them gets the row back.
pub async fn consume_reset_token(pool: &PgPool, token: &str) -> Result<Option<Uuid>, sqlx::Error> {
let row: Option<(Uuid,)> = sqlx::query_as(
"UPDATE password_reset_tokens SET used_at = now()
WHERE token_hash = $1 AND used_at IS NULL AND expires_at > now()
RETURNING account_id",
)
.bind(hash_token(token))
.fetch_optional(pool)
.await?;
Ok(row.map(|(id,)| id))
}
/// Kill every refresh session of the account: whoever holds a stolen session
/// cookie must not survive a password change.
pub async fn revoke_all_sessions(pool: &PgPool, account_id: Uuid) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE refresh_tokens SET revoked_at = now()
WHERE account_id = $1 AND revoked_at IS NULL",
)
.bind(account_id)
.execute(pool)
.await?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn well_formed_token_shape_matches_opaque_generator() {
let token = new_opaque_token(TOKEN_PREFIX);
assert!(is_well_formed_token(&token));
assert!(!is_well_formed_token(&format!("{token}x")));
assert!(!is_well_formed_token("lvpr_short"));
assert!(!is_well_formed_token(
"XWpr_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
));
}
}

View file

@ -16,9 +16,21 @@ use axum::{
};
use config::Config;
use device::{handlers::DeviceState, store::DeviceStore};
use tower_http::trace::TraceLayer;
pub fn build_app(pool: sqlx::PgPool, cfg: &Config) -> Router {
let auth_state = AuthState::new(pool.clone(), cfg.jwt_secret.clone(), cfg.cookie_secure);
build_app_with_mail(pool, cfg, auth::reset::MailBox::Log)
}
/// Same router with a custom reset-email back-end: production uses the log
/// mailer, tests capture tokens through the queue variant.
pub fn build_app_with_mail(pool: sqlx::PgPool, cfg: &Config, mail: auth::reset::MailBox) -> Router {
let auth_state = AuthState::with_mail(
pool.clone(),
cfg.jwt_secret.clone(),
cfg.cookie_secure,
mail,
);
let device_state = DeviceState {
store: DeviceStore::default(),
pool: pool.clone(),
@ -43,6 +55,14 @@ pub fn build_app(pool: sqlx::PgPool, cfg: &Config) -> Router {
.route("/auth/login", post(auth::handlers::login))
.route("/auth/refresh", post(auth::handlers::refresh))
.route("/auth/logout", post(auth::handlers::logout))
.route(
"/auth/forgot-password",
post(auth::reset::handlers::forgot_password),
)
.route(
"/auth/reset-password",
post(auth::reset::handlers::reset_password),
)
.with_state(auth_state);
let device_routes = Router::new()
@ -87,4 +107,5 @@ pub fn build_app(pool: sqlx::PgPool, cfg: &Config) -> Router {
Router::new()
.route("/health", get(|| async { "ok" }))
.merge(api)
.layer(TraceLayer::new_for_http())
}

View file

@ -2,7 +2,12 @@ use accounts_service::{build_app, config::Config};
#[tokio::main]
async fn main() -> anyhow::Result<()> {
tracing_subscriber::fmt::init();
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")),
)
.init();
dotenvy::dotenv().ok();
let cfg = Config::from_env()?;
cfg.validate()?;

View file

@ -20,6 +20,8 @@ fn form(bytes: Vec<u8>) -> MultipartForm {
#[tokio::test]
async fn valid_upload_stores_avatar_and_returns_url() {
common::init_test_logging();
common::ensure_avatar_bucket().await;
let pool = common::test_pool().await;
let server = common::test_server(accounts_service::build_app(
pool.clone(),

View file

@ -8,6 +8,10 @@ use uuid::Uuid;
// mounts as `mod common;`).
pub use ::common::internal::{ACCOUNT_ID_HEADER, INTERNAL_KEY_HEADER};
pub fn init_test_logging() {
let _ = tracing_subscriber::fmt::try_init();
}
pub async fn test_pool() -> sqlx::PgPool {
let url = std::env::var("DATABASE_URL")
.unwrap_or_else(|_| "postgres://lovisual:lovisual@localhost:5432/accounts_db".into());
@ -60,3 +64,65 @@ pub async fn register_account(server: &TestServer) -> (Uuid, String) {
email,
)
}
/// Creates the test avatar bucket if it does not exist yet, so the avatar
/// tests are self-contained: any S3-compatible backend (MinIO, SeaweedFS)
/// works without external `mc mb` bootstrap. Mirrors `S3Storage::from_config`.
/// Retries briefly so a just-started container does not race the test.
pub async fn ensure_avatar_bucket() {
let creds =
aws_sdk_s3::config::Credentials::new("minioadmin", "minioadmin", None, None, "static");
let config = aws_sdk_s3::config::Builder::new()
.endpoint_url("http://localhost:9000")
.credentials_provider(creds)
.region(aws_sdk_s3::config::Region::new("us-east-1"))
.force_path_style(true)
.behavior_version(aws_sdk_s3::config::BehaviorVersion::latest())
.build();
let client = aws_sdk_s3::Client::from_conf(config);
let mut last_err = String::new();
for _ in 0..30 {
// Probe writability, not just bucket existence: a freshly started
// S3 backend may still be electing volumes ("Not enough data nodes"
// on SeaweedFS) right after create_bucket succeeds.
match async {
// Already-exists is success (MinIO: BucketAlreadyOwnedByYou,
// SeaweedFS: BucketAlreadyExists); anything else aborts.
if let Err(e) = client
.create_bucket()
.bucket("lovisual-avatars-test")
.send()
.await
{
let msg = format!("{e:?}");
if !msg.contains("BucketAlready") {
return Err(msg);
}
}
client
.put_object()
.bucket("lovisual-avatars-test")
.key(".probe")
.body(aws_sdk_s3::primitives::ByteStream::from_static(b"probe"))
.send()
.await
.map_err(|e| format!("{e:?}"))?;
let _ = client
.delete_object()
.bucket("lovisual-avatars-test")
.key(".probe")
.send()
.await;
Ok::<(), String>(())
}
.await
{
Ok(_) => return,
Err(msg) => {
last_err = msg;
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
}
}
}
panic!("S3 backend is not writable: {last_err}");
}

View file

@ -0,0 +1,215 @@
mod common;
use accounts_service::auth::reset::{Mail, MailBox};
use axum::http::StatusCode;
use axum_test::TestServer;
use sqlx::PgPool;
use tokio::sync::mpsc::{UnboundedReceiver, unbounded_channel};
type App = (TestServer, PgPool, UnboundedReceiver<Mail>);
async fn app() -> App {
let (tx, rx) = unbounded_channel();
let pool = common::test_pool().await;
let server = common::test_server(accounts_service::build_app_with_mail(
pool.clone(),
&common::test_config(),
MailBox::Queue(tx),
));
(server, pool, rx)
}
async fn request_reset(server: &TestServer, email: &str) {
server
.post("/auth/forgot-password")
.json(&serde_json::json!({ "email": email }))
.await
.assert_status(StatusCode::ACCEPTED);
}
async fn reset_with(server: &TestServer, token: &str, password: &str) -> axum_test::TestResponse {
server
.post("/auth/reset-password")
.json(&serde_json::json!({ "token": token, "new_password": password }))
.await
}
#[tokio::test]
async fn forgot_password_never_reveals_account_existence() {
let (server, _pool, _mail) = app().await;
// Unknown email and known email must be indistinguishable: same status,
// same body. Enumeration is the first step of account takeover.
let unknown = server
.post("/auth/forgot-password")
.json(&serde_json::json!({ "email": "nobody-here@example.com" }))
.await;
unknown.assert_status(StatusCode::ACCEPTED);
let unknown_body: serde_json::Value = unknown.json();
assert_eq!(
unknown_body["status"],
"reset email sent if the account exists"
);
let (_, email) = common::register_account(&server).await;
let known = server
.post("/auth/forgot-password")
.json(&serde_json::json!({ "email": email }))
.await;
known.assert_status(StatusCode::ACCEPTED);
let known_body: serde_json::Value = known.json();
assert_eq!(unknown_body, known_body);
}
#[tokio::test]
async fn full_reset_flow_changes_password_and_kills_sessions() {
let (server, _pool, mut mail) = app().await;
let (_, email) = common::register_account(&server).await;
let old_password = "correct-horse-battery-staple";
// Log in to create a live refresh session the reset must kill.
let login = server
.post("/auth/login")
.json(&serde_json::json!({ "email": email, "password": old_password }))
.await;
login.assert_status_ok();
let old_refresh = refresh_cookie(&login).expect("login must set the refresh cookie");
request_reset(&server, &email).await;
let token = mail.recv().await.expect("queue mailer must deliver").token;
reset_with(&server, &token, "brand-new-password-1")
.await
.assert_status_ok();
// Old password is dead, new password works.
server
.post("/auth/login")
.json(&serde_json::json!({ "email": email, "password": old_password }))
.await
.assert_status(StatusCode::UNAUTHORIZED);
server
.post("/auth/login")
.json(&serde_json::json!({ "email": email, "password": "brand-new-password-1" }))
.await
.assert_status_ok();
// Every refresh session issued before the reset is revoked: replaying
// the pre-reset cookie must not survive (a stolen session cannot
// outlive a password change).
let replay = server
.post("/auth/refresh")
.add_cookie(old_refresh.as_str().into())
.await;
replay.assert_status(StatusCode::UNAUTHORIZED);
}
#[tokio::test]
async fn reset_token_is_single_use() {
let (server, _pool, mut mail) = app().await;
let (_, email) = common::register_account(&server).await;
request_reset(&server, &email).await;
let token = mail.recv().await.expect("mail").token;
reset_with(&server, &token, "brand-new-password-1")
.await
.assert_status_ok();
reset_with(&server, &token, "another-password-2")
.await
.assert_status(StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn new_request_supersedes_pending_token() {
let (server, _pool, mut mail) = app().await;
let (_, email) = common::register_account(&server).await;
request_reset(&server, &email).await;
let first = mail.recv().await.expect("mail").token;
request_reset(&server, &email).await;
let second = mail.recv().await.expect("mail").token;
assert_ne!(first, second);
// The superseded token no longer works, the fresh one does.
reset_with(&server, &first, "brand-new-password-1")
.await
.assert_status(StatusCode::BAD_REQUEST);
reset_with(&server, &second, "brand-new-password-1")
.await
.assert_status_ok();
}
#[tokio::test]
async fn bad_tokens_are_rejected_without_oracle() {
let (server, _pool, _mail) = app().await;
server
.post("/auth/reset-password")
.json(&serde_json::json!({ "token": "short", "new_password": "brand-new-password-1" }))
.await
.assert_status(StatusCode::BAD_REQUEST);
// Unknown but well-formed token: same 400 family, and unlike the
// shape-rejection above it must not leak which of unknown/expired/used
// it is.
let unknown = server
.post("/auth/reset-password")
.json(&serde_json::json!({
"token": "lvpr_AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA",
"new_password": "brand-new-password-1"
}))
.await;
unknown.assert_status(StatusCode::BAD_REQUEST);
let body: serde_json::Value = unknown.json();
assert_eq!(body["error"], "reset token is invalid or expired");
}
#[tokio::test]
async fn weak_password_is_rejected_before_token_consumption() {
let (server, _pool, mut mail) = app().await;
let (_, email) = common::register_account(&server).await;
request_reset(&server, &email).await;
let token = mail.recv().await.expect("mail").token;
reset_with(&server, &token, "short12")
.await
.assert_status(StatusCode::BAD_REQUEST);
// The token must still be usable afterwards: rejecting a weak password
// must not burn the user's one link.
reset_with(&server, &token, "brand-new-password-1")
.await
.assert_status_ok();
}
#[tokio::test]
async fn reset_mail_only_goes_to_known_accounts() {
let (server, _pool, mut mail) = app().await;
request_reset(&server, "ghost@example.com").await;
assert!(
mail.try_recv().is_err(),
"unknown account must not enqueue a reset email"
);
let (_, email) = common::register_account(&server).await;
request_reset(&server, &email).await;
let sent = mail.recv().await.expect("mail");
assert_eq!(sent.to, email);
assert!(sent.token.starts_with("lvpr_"));
}
/// The refresh cookie is httpOnly and scoped to /auth; axum-test exposes
/// response headers, so parse Set-Cookie directly.
fn refresh_cookie(res: &axum_test::TestResponse) -> Option<String> {
res.headers()
.get_all(axum::http::header::SET_COOKIE)
.iter()
.find_map(|v| {
let s = v.to_str().ok()?;
s.strip_prefix("lv_refresh=")
.map(|rest| format!("lv_refresh={}", rest.split(';').next().unwrap_or("")))
})
}

View file

@ -10,9 +10,10 @@ path = "src/lib.rs"
[dependencies]
common = { path = "../common" }
axum = { version = "0.8", features = ["macros"] }
tower-http = { version = "0.7", features = ["trace"] }
tokio = { version = "1", features = ["rt-multi-thread", "macros", "time"] }
tracing = "0.1"
tracing-subscriber = "0.3"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
serde = { version = "1", features = ["derive"] }
serde_json = "1"
uuid = { version = "1", features = ["v4", "serde"] }

View file

@ -1,11 +1,16 @@
pub mod config;
pub mod presence;
use axum::{Router, extract::DefaultBodyLimit, routing::{get, post}};
use axum::{
Router,
extract::DefaultBodyLimit,
routing::{get, post},
};
use common::internal::{InternalKey, require_internal_key};
use config::Config;
use presence::store::Store;
use std::sync::Arc;
use tower_http::trace::TraceLayer;
/// Presence requests are tiny JSON (a ≤64-char payload and ≤40 uuids).
const MAX_BODY_BYTES: usize = 8 * 1024;
@ -22,6 +27,7 @@ pub fn build_app(cfg: &Config) -> (Router, Arc<Store>) {
));
let app = Router::new()
.route("/health", get(|| async { "ok" }))
.merge(api);
.merge(api)
.layer(TraceLayer::new_for_http());
(app, store)
}

View file

@ -1,7 +1,12 @@
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
tracing_subscriber::fmt::init();
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")),
)
.init();
let cfg = chat_service::config::Config::from_env()?;
cfg.validate()?;
let (app, state) = chat_service::build_app(&cfg);

View file

@ -63,15 +63,23 @@ pub async fn sync(
}
let now = Instant::now();
if req.publish
&& let Err(PublishError::Claimed) = store.publish(id.account_id, &server, req.mc, req.gui, now)
&& let Err(PublishError::Claimed) =
store.publish(id.account_id, &server, req.mc, req.gui, now)
{
return error(StatusCode::CONFLICT, "player uuid is bound to another account");
return error(
StatusCode::CONFLICT,
"player uuid is bound to another account",
);
}
let players: Vec<PlayerState> = store
.lookup(&server, &req.want, now)
.into_iter()
.filter(|s| s.mc != req.mc)
.map(|s| PlayerState { mc: s.mc, gui: s.gui, age: s.age_ms })
.map(|s| PlayerState {
mc: s.mc,
gui: s.gui,
age: s.age_ms,
})
.collect();
Json(json!({ "players": players })).into_response()
}

View file

@ -6,7 +6,10 @@ pub mod handlers;
pub mod payload;
pub mod store;
use std::{sync::Arc, time::{Duration, Instant}};
use std::{
sync::Arc,
time::{Duration, Instant},
};
use store::Store;
pub fn spawn_purger(store: Arc<Store>) {

View file

@ -10,7 +10,8 @@ pub const MAX_WANT: usize = 40;
pub fn valid_payload(s: &str) -> bool {
!s.is_empty()
&& s.len() <= MAX_PAYLOAD_CHARS
&& s.bytes().all(|b| b.is_ascii_alphanumeric() || matches!(b, b'+' | b'/' | b'='))
&& s.bytes()
.all(|b| b.is_ascii_alphanumeric() || matches!(b, b'+' | b'/' | b'='))
}
/// Normalised server key: trimmed, lower-case, no control characters, bounded.
@ -38,7 +39,10 @@ mod tests {
#[test]
fn server_is_normalised() {
assert_eq!(normalize_server(" Play.Example.COM ").as_deref(), Some("play.example.com"));
assert_eq!(
normalize_server(" Play.Example.COM ").as_deref(),
Some("play.example.com")
);
assert_eq!(normalize_server(""), None);
assert_eq!(normalize_server("a\nb"), None);
assert_eq!(normalize_server(&"x".repeat(65)), None);

View file

@ -97,7 +97,9 @@ impl Store {
pub fn purge(&self, now: Instant) {
let mut inner = self.inner.lock().expect("presence lock");
inner.entries.retain(|_, e| now.duration_since(e.updated) < ENTRY_TTL);
inner
.entries
.retain(|_, e| now.duration_since(e.updated) < ENTRY_TTL);
let expired: Vec<Uuid> = inner
.claims
.iter()
@ -117,7 +119,12 @@ mod tests {
use super::*;
fn ids() -> (Uuid, Uuid, Uuid, Uuid) {
(Uuid::from_u128(1), Uuid::from_u128(2), Uuid::from_u128(10), Uuid::from_u128(20))
(
Uuid::from_u128(1),
Uuid::from_u128(2),
Uuid::from_u128(10),
Uuid::from_u128(20),
)
}
#[test]
@ -125,7 +132,9 @@ mod tests {
let (acc, _, mc, _) = ids();
let store = Store::default();
let t0 = Instant::now();
store.publish(acc, "srv", mc, Some("AAAA".into()), t0).unwrap();
store
.publish(acc, "srv", mc, Some("AAAA".into()), t0)
.unwrap();
let seen = store.lookup("srv", &[mc], t0 + Duration::from_millis(500));
assert_eq!(seen.len(), 1);
assert_eq!(seen[0].gui.as_deref(), Some("AAAA"));
@ -149,7 +158,10 @@ mod tests {
let store = Store::default();
let t0 = Instant::now();
store.publish(acc, "s", mc, None, t0).unwrap();
assert_eq!(store.publish(other, "s", mc, None, t0 + Duration::from_secs(1)), Err(PublishError::Claimed));
assert_eq!(
store.publish(other, "s", mc, None, t0 + Duration::from_secs(1)),
Err(PublishError::Claimed)
);
// after the claim goes idle it can be taken over
assert!(store.publish(other, "s", mc, None, t0 + CLAIM_TTL).is_ok());
}
@ -183,6 +195,10 @@ mod tests {
let t0 = Instant::now();
store.publish(acc, "s", mc, None, t0).unwrap();
store.purge(t0 + CLAIM_TTL);
assert!(store.publish(Uuid::from_u128(2), "s", mc, None, t0 + CLAIM_TTL).is_ok());
assert!(
store
.publish(Uuid::from_u128(2), "s", mc, None, t0 + CLAIM_TTL)
.is_ok()
);
}
}

View file

@ -7,7 +7,10 @@ use uuid::Uuid;
const KEY: &str = "kkkkkkkkkkkkkkkkkkkkkkkkkkkkkkkk";
fn server() -> TestServer {
let cfg = Config { port: 0, internal_key: KEY.into() };
let cfg = Config {
port: 0,
internal_key: KEY.into(),
};
let (app, _) = chat_service::build_app(&cfg);
TestServer::new(app)
}
@ -37,10 +40,20 @@ async fn requires_internal_key_and_identity() {
async fn two_players_see_each_other() {
let s = server();
let (a, b) = (Uuid::from_u128(10), Uuid::from_u128(20));
let r = sync(&s, 1, json!({"server": "Play.X", "mc": a, "gui": "AQEAAAA=", "want": [b]})).await;
let r = sync(
&s,
1,
json!({"server": "Play.X", "mc": a, "gui": "AQEAAAA=", "want": [b]}),
)
.await;
assert_eq!(r.status_code(), 200);
assert_eq!(r.json::<Value>()["players"].as_array().unwrap().len(), 0);
let r = sync(&s, 2, json!({"server": "play.x", "mc": b, "gui": null, "want": [a]})).await;
let r = sync(
&s,
2,
json!({"server": "play.x", "mc": b, "gui": null, "want": [a]}),
)
.await;
let body = r.json::<Value>();
let players = body["players"].as_array().unwrap();
assert_eq!(players.len(), 1);
@ -56,7 +69,12 @@ async fn rejects_bad_input_and_uuid_theft() {
assert_eq!(r.status_code(), 400);
let r = sync(&s, 1, json!({"server": "", "mc": mc})).await;
assert_eq!(r.status_code(), 400);
assert_eq!(sync(&s, 1, json!({"server": "x", "mc": mc})).await.status_code(), 200);
assert_eq!(
sync(&s, 1, json!({"server": "x", "mc": mc}))
.await
.status_code(),
200
);
let r = sync(&s, 2, json!({"server": "x", "mc": mc})).await;
assert_eq!(r.status_code(), 409);
}
@ -65,7 +83,12 @@ async fn rejects_bad_input_and_uuid_theft() {
async fn read_only_poll_does_not_publish() {
let s = server();
let (a, b) = (Uuid::from_u128(10), Uuid::from_u128(20));
sync(&s, 1, json!({"server": "x", "mc": a, "publish": false, "want": [b]})).await;
sync(
&s,
1,
json!({"server": "x", "mc": a, "publish": false, "want": [b]}),
)
.await;
let r = sync(&s, 2, json!({"server": "x", "mc": b, "want": [a]})).await;
assert_eq!(r.json::<Value>()["players"].as_array().unwrap().len(), 0);
}

View file

@ -10,9 +10,10 @@ path = "src/lib.rs"
[dependencies]
common = { path = "../common" }
axum = { version = "0.8", features = ["macros"] }
tower-http = { version = "0.7", features = ["trace"] }
tokio = { version = "1", features = ["rt-multi-thread", "macros"] }
tracing = "0.1"
tracing-subscriber = "0.3"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
serde = { version = "1", features = ["derive"] }
serde_json = "1"
sqlx = { version = "0.9", default-features = false, features = ["runtime-tokio", "tls-rustls", "postgres", "uuid", "chrono", "json", "macros", "migrate"] }

View file

@ -13,6 +13,7 @@ use common::internal::{InternalKey, require_internal_key};
use config::Config;
use showcase::profiles::ProfileSource;
use std::sync::Arc;
use tower_http::trace::TraceLayer;
pub fn build_app(pool: sqlx::PgPool, cfg: &Config, profiles: Arc<dyn ProfileSource>) -> Router {
let slots_state = slots::handlers::SlotsState { pool: pool.clone() };
@ -52,4 +53,5 @@ pub fn build_app(pool: sqlx::PgPool, cfg: &Config, profiles: Arc<dyn ProfileSour
Router::new()
.route("/health", get(|| async { "ok" }))
.merge(api)
.layer(TraceLayer::new_for_http())
}

View file

@ -1,7 +1,12 @@
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
tracing_subscriber::fmt::init();
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")),
)
.init();
let cfg = configs_service::config::Config::from_env()?;
cfg.validate()?;
let pool = sqlx::PgPool::connect(&cfg.database_url).await?;

View file

@ -126,6 +126,12 @@ services:
ACCOUNTS_GRPC_URL: http://accounts-service:50051
CONFIGS_HTTP_URL: http://configs-service:8082
CHAT_HTTP_URL: http://chat-service:8083
# This stack is only ever exposed through nginx (see deploy/), so the
# rate limiter must read the client IP from X-Forwarded-For. Without
# this every client shares the proxy's socket address and one bucket:
# the global 300/min and login 5/min would be site-wide self-DoS.
# Keep TRUST_PROXY=false in .env for direct-exposure dev runs.
TRUST_PROXY: "true"
# Read-only static files served under GET /downloads/* (lovisual.jar).
DOWNLOADS_DIR: /srv/downloads
volumes:

View file

@ -14,7 +14,7 @@ http-body-util = "0.1"
tokio = { version = "1", features = ["rt-multi-thread", "macros", "net", "time"] }
tower-http = { version = "0.7", features = ["cors", "trace", "fs"] }
tracing = "0.1"
tracing-subscriber = "0.3"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
serde_json = "1"
uuid = { version = "1", features = ["v4"] }
reqwest = { version = "0.13", default-features = false, features = ["stream"] }

View file

@ -15,7 +15,7 @@ use axum::{
use config::Config;
use identity::device::DeviceAuthenticator;
use std::sync::Arc;
use tower_http::{cors::CorsLayer, services::ServeDir};
use tower_http::{cors::CorsLayer, services::ServeDir, trace::TraceLayer};
pub fn build_app(cfg: &Config, devices: Arc<dyn DeviceAuthenticator>) -> Router {
let upstreams = Arc::new(proxy::forward::Upstreams::new(cfg).expect("valid upstream config"));
@ -69,6 +69,9 @@ pub fn build_app(cfg: &Config, devices: Arc<dyn DeviceAuthenticator>) -> Router
guard::reject_ambiguous_paths,
))
.layer(cors)
// Outermost so every route (health, downloads, proxied API) is
// traced; spans carry method/path/status for the fmt subscriber.
.layer(TraceLayer::new_for_http())
}
fn spawn_purger(limits: Arc<rate_limit::RateLimits>) {

View file

@ -4,7 +4,12 @@ use std::sync::Arc;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
tracing_subscriber::fmt::init();
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")),
)
.init();
let cfg = gateway::config::Config::from_env()?;
cfg.validate()?;
let devices = Arc::new(GrpcDevices::connect_lazy(

View file

@ -46,6 +46,20 @@ pub fn rules() -> Vec<Rule> {
vec![
rule(Method::POST, "/auth/login", Quota::per_minute(n(5)), Ip),
rule(Method::POST, "/auth/register", Quota::per_hour(n(3)), Ip),
// Same threat as register: each accepted request sends an email, so
// the mailer must not be usable as a spam cannon.
rule(
Method::POST,
"/auth/forgot-password",
Quota::per_hour(n(3)),
Ip,
),
rule(
Method::POST,
"/auth/reset-password",
Quota::per_minute(n(10)),
Ip,
),
rule(Method::POST, "/auth/refresh", Quota::per_minute(n(30)), Ip),
rule(Method::POST, "/device/code", Quota::per_minute(n(10)), Ip),
// The mod polls every 2–3 s for up to 10 min: 10/min would break linking.
@ -61,7 +75,12 @@ pub fn rules() -> Vec<Rule> {
rule(Method::POST, "/avatars", Quota::per_hour(n(5)), Account),
rule(Method::GET, "/showcase*", Quota::per_minute(n(60)), Ip),
// The mod syncs menu presence about 3 times a second while a menu is open.
rule(Method::POST, "/presence/*", Quota::per_second(n(6)), Account),
rule(
Method::POST,
"/presence/*",
Quota::per_second(n(6)),
Account,
),
]
}