//! opentracker server: one binary, one SQLite file. //! //! Wiring, in dependency order: //! //! ```text //! Db ──▶ writer task ──▶ Ingest ──▶ UDP workers //! │ │ //! └──▶ axum HTTP ─────────┘ (login mints tokens straight into Ingest) //! ``` //! //! The one operational trap worth repeating: **HTTP reverse proxies do not forward //! UDP.** The HTTP listener belongs behind nginx or Caddy; the UDP port needs its //! own firewall rule. mod api; mod auth; mod config; mod db; mod ingest; mod keys; mod limits; mod polyline; mod retention; mod simulate; mod tiles; mod udp; mod web; mod writer; use std::net::SocketAddr; use std::path::PathBuf; use std::sync::Arc; use anyhow::{Context, Result, bail}; use tower_sessions::cookie::SameSite; use tower_sessions::{Expiry, SessionManagerLayer}; use tower_sessions_sqlx_store::SqliteStore; use tracing::{info, warn}; use tracing_subscriber::EnvFilter; use crate::api::AppState; use crate::config::Config; use crate::db::Db; use crate::ingest::Ingest; use crate::keys::KeyVault; /// Rolling session lifetime. const SESSION_DAYS: i64 = 30; struct Args { config: Option, /// Relaxes the cookie's `Secure` requirement so the Vite dev server works /// over plain HTTP on localhost. Never for production. dev: bool, simulate_device: bool, /// `--create-admin `: bootstrap the first account. create_admin: Option<(String, String)>, } fn parse_args() -> Result { let mut args = Args { config: None, dev: false, simulate_device: false, create_admin: None, }; let mut it = std::env::args().skip(1); while let Some(arg) = it.next() { match arg.as_str() { "--config" | "-c" => { args.config = Some(PathBuf::from(it.next().context("--config needs a path")?)); } "--dev" => args.dev = true, "--simulate-device" => args.simulate_device = true, "--create-admin" => { let user = it.next().context("--create-admin needs a username")?; let pass = it.next().context("--create-admin needs a password")?; args.create_admin = Some((user, pass)); } "--help" | "-h" => { println!( "opentracker {}\n\n\ USAGE:\n \ otserver [--config FILE] [--dev] [--simulate-device]\n \ otserver --create-admin USERNAME PASSWORD\n\n\ ENVIRONMENT:\n \ OT_SECRET_KEY required; 32 random bytes, base64, or a path to them\n \ OT_SECRET_KEY_OLD accepted for unwrapping during a key rotation\n \ OT_* override any config field (see config.rs)\n", env!("CARGO_PKG_VERSION") ); std::process::exit(0); } other => bail!("unknown argument {other:?} (try --help)"), } } Ok(args) } #[tokio::main] async fn main() -> Result<()> { tracing_subscriber::fmt() .with_env_filter( EnvFilter::try_from_env("OT_LOG") .unwrap_or_else(|_| EnvFilter::new("info,otserver=debug")), ) .init(); let args = parse_args()?; let cfg = Config::load(args.config.as_deref())?; cfg.validate()?; let vault = KeyVault::from_env()?; let db = Db::open(&cfg.db_path).await?; info!(path = %cfg.db_path.display(), "database ready"); let (writer, writer_task) = writer::spawn(db.write.clone()); // Handing over the master is what enables replies to datagrams naming tokens // this server has no record of. Withheld unless config asks for it, so the // one reflective path cannot be switched on by accident. let revocation_master = cfg.revocation_notices.then(|| vault.revocation_master()); if revocation_master.is_some() { info!( "revocation_notices enabled: an unknown token draws a rate-limited 38-byte reply. \ See config.rs for the trade." ); } let ingest = Arc::new(Ingest::new( writer.clone(), cfg.ts_window_days * 86_400, revocation_master, )); // One-shot bootstrap, before anything starts listening. if let Some((username, password)) = args.create_admin { let hash = auth::hash_password(&cfg, password).await?; let at = db::now(); sqlx::query( "INSERT INTO users (username, pw_hash, display_name, is_admin, created_at, pw_changed_at) \ VALUES (?, ?, ?, 1, ?, ?)", ) .bind(&username) .bind(hash) .bind(&username) .bind(at) .bind(at) .execute(&db.write) .await .with_context(|| format!("creating admin {username}"))?; println!("created admin account {username}"); db.checkpoint().await?; db.close().await; return Ok(()); } let loaded = auth::load_tokens(&db.read, &vault, &ingest).await?; info!( tokens = loaded, "loaded device tokens into the ingest cache" ); // Sessions live in the same SQLite file, on the writer pool. They are // low-traffic enough not to disturb the batching that exists to protect the // position write path. let session_store = SqliteStore::new(db.write.clone()); session_store .migrate() .await .context("migrating the session store")?; if args.dev { warn!("--dev: session cookies will not require HTTPS. Never use this in production."); } let session_layer = SessionManagerLayer::new(session_store) .with_name(if args.dev { "otsid" } else { "__Host-otsid" }) .with_http_only(true) .with_secure(!args.dev) .with_same_site(SameSite::Lax) .with_expiry(Expiry::OnInactivity(time::Duration::days(SESSION_DAYS))); let http_addr = cfg.http_addr; let udp_addr = cfg.udp_addr; let workers = cfg.udp_worker_count(); let retention_days = cfg.retention_days; let token_stale_days = cfg.token_stale_days; let throttle = Arc::new(auth::LoginThrottle::default()); let state: api::Shared = Arc::new(AppState { db, cfg, vault, ingest: Arc::clone(&ingest), writer, throttle: Arc::clone(&throttle), }); let udp_tasks = udp::spawn(udp_addr, workers, Arc::clone(&ingest))?; let retention_task = retention::Retention { pool: state.db.write.clone(), retention_days, token_stale_days, ingest: Arc::clone(&ingest), } .spawn(); let limits_task = retention::spawn_limits_gc(Arc::clone(&ingest), throttle); let eviction_task = tiles::spawn_eviction( state.db.write.clone(), state.cfg.cache_dir.clone(), state.cfg.max_cache_bytes, ); let app = api::router(Arc::clone(&state)).layer(session_layer); let listener = tokio::net::TcpListener::bind(http_addr) .await .with_context(|| format!("binding {http_addr}"))?; info!(addr = %http_addr, "HTTP listener started (put a TLS-terminating proxy in front)"); if args.simulate_device { // After the listeners are up, so the login cannot race them. simulate::SimulatedDevice { base_url: format!("http://{http_addr}"), username: std::env::var("OT_SIM_USER").unwrap_or_else(|_| "sim".into()), password: std::env::var("OT_SIM_PASSWORD").unwrap_or_else(|_| "simsimsimsim".into()), udp_addr: loopback_of(udp_addr), } .spawn(); } axum::serve( listener, app.into_make_service_with_connect_info::(), ) .with_graceful_shutdown(shutdown_signal()) .await .context("HTTP server failed")?; // Stop the periodic tasks and the receive loops, then let the writer drain so // nothing already acknowledged is lost, then checkpoint so the .db file on // disk is self-contained. info!("shutting down"); retention_task.abort(); limits_task.abort(); eviction_task.abort(); for task in udp_tasks { task.abort(); } if let Err(e) = state.db.checkpoint().await { warn!(error = %e, "WAL checkpoint failed"); } drop(state); let _ = tokio::time::timeout(std::time::Duration::from_secs(5), writer_task).await; Ok(()) } /// The simulated device must talk to a routable address; the listener is usually /// bound to the wildcard, which cannot be a destination. fn loopback_of(addr: SocketAddr) -> String { if addr.ip().is_unspecified() { format!("127.0.0.1:{}", addr.port()) } else { addr.to_string() } } async fn shutdown_signal() { let ctrl_c = async { tokio::signal::ctrl_c().await.ok(); }; #[cfg(unix)] let terminate = async { if let Ok(mut sig) = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) { sig.recv().await; } }; #[cfg(not(unix))] let terminate = std::future::pending::<()>(); tokio::select! { () = ctrl_c => {}, () = terminate => {}, } } #[cfg(test)] mod tests { use super::*; #[test] fn the_wildcard_address_is_rewritten_for_the_simulated_device() { assert_eq!( loopback_of("0.0.0.0:7373".parse().expect("literal")), "127.0.0.1:7373" ); assert_eq!( loopback_of("192.0.2.1:7373".parse().expect("literal")), "192.0.2.1:7373" ); } }