main.rs
⎇
Raw
1//! opentracker server: one binary, one SQLite file.
2//!
3//! Wiring, in dependency order:
4//!
5//! ```text
6//! Db ──▶ writer task ──▶ Ingest ──▶ UDP workers
7//! │ │
8//! └──▶ axum HTTP ─────────┘ (login mints tokens straight into Ingest)
9//! ```
10//!
11//! The one operational trap worth repeating: **HTTP reverse proxies do not forward
12//! UDP.** The HTTP listener belongs behind nginx or Caddy; the UDP port needs its
13//! own firewall rule.
14
15mod api;
16mod auth;
17mod config;
18mod db;
19mod ingest;
20mod keys;
21mod limits;
22mod polyline;
23mod retention;
24mod simulate;
25mod udp;
26mod web;
27mod writer;
28
29use std::net::SocketAddr;
30use std::path::PathBuf;
31use std::sync::Arc;
32
33use anyhow::{Context, Result, bail};
34use tower_sessions::cookie::SameSite;
35use tower_sessions::{Expiry, SessionManagerLayer};
36use tower_sessions_sqlx_store::SqliteStore;
37use tracing::{info, warn};
38use tracing_subscriber::EnvFilter;
39
40use crate::api::AppState;
41use crate::config::Config;
42use crate::db::Db;
43use crate::ingest::Ingest;
44use crate::keys::KeyVault;
45
46/// Rolling session lifetime.
47const SESSION_DAYS: i64 = 30;
48
49struct Args {
50 config: Option<PathBuf>,
51 /// Relaxes the cookie's `Secure` requirement so the Vite dev server works
52 /// over plain HTTP on localhost. Never for production.
53 dev: bool,
54 simulate_device: bool,
55 /// `--create-admin <user> <password>`: bootstrap the first account.
56 create_admin: Option<(String, String)>,
57}
58
59fn parse_args() -> Result<Args> {
60 let mut args = Args {
61 config: None,
62 dev: false,
63 simulate_device: false,
64 create_admin: None,
65 };
66 let mut it = std::env::args().skip(1);
67 while let Some(arg) = it.next() {
68 match arg.as_str() {
69 "--config" | "-c" => {
70 args.config = Some(PathBuf::from(it.next().context("--config needs a path")?));
71 }
72 "--dev" => args.dev = true,
73 "--simulate-device" => args.simulate_device = true,
74 "--create-admin" => {
75 let user = it.next().context("--create-admin needs a username")?;
76 let pass = it.next().context("--create-admin needs a password")?;
77 args.create_admin = Some((user, pass));
78 }
79 "--help" | "-h" => {
80 println!(
81 "opentracker {}\n\n\
82 USAGE:\n \
83 otserver [--config FILE] [--dev] [--simulate-device]\n \
84 otserver --create-admin USERNAME PASSWORD\n\n\
85 ENVIRONMENT:\n \
86 OT_SECRET_KEY required; 32 random bytes, base64, or a path to them\n \
87 OT_SECRET_KEY_OLD accepted for unwrapping during a key rotation\n \
88 OT_* override any config field (see config.rs)\n",
89 env!("CARGO_PKG_VERSION")
90 );
91 std::process::exit(0);
92 }
93 other => bail!("unknown argument {other:?} (try --help)"),
94 }
95 }
96 Ok(args)
97}
98
99#[tokio::main]
100async fn main() -> Result<()> {
101 tracing_subscriber::fmt()
102 .with_env_filter(
103 EnvFilter::try_from_env("OT_LOG")
104 .unwrap_or_else(|_| EnvFilter::new("info,otserver=debug")),
105 )
106 .init();
107
108 let args = parse_args()?;
109 let cfg = Config::load(args.config.as_deref())?;
110 cfg.validate()?;
111
112 let vault = KeyVault::from_env()?;
113 let db = Db::open(&cfg.db_path).await?;
114 info!(path = %cfg.db_path.display(), "database ready");
115
116 let (writer, writer_task) = writer::spawn(db.write.clone());
117 // Handing over the master is what enables replies to datagrams naming tokens
118 // this server has no record of. Withheld unless config asks for it, so the
119 // one reflective path cannot be switched on by accident.
120 let revocation_master = cfg.revocation_notices.then(|| vault.revocation_master());
121 if revocation_master.is_some() {
122 info!(
123 "revocation_notices enabled: an unknown token draws a rate-limited 38-byte reply. \
124 See config.rs for the trade."
125 );
126 }
127 let ingest = Arc::new(Ingest::new(
128 writer.clone(),
129 cfg.ts_window_days * 86_400,
130 revocation_master,
131 ));
132
133 // One-shot bootstrap, before anything starts listening.
134 if let Some((username, password)) = args.create_admin {
135 let hash = auth::hash_password(&cfg, password).await?;
136 let at = db::now();
137 sqlx::query(
138 "INSERT INTO users (username, pw_hash, display_name, is_admin, created_at, pw_changed_at) \
139 VALUES (?, ?, ?, 1, ?, ?)",
140 )
141 .bind(&username)
142 .bind(hash)
143 .bind(&username)
144 .bind(at)
145 .bind(at)
146 .execute(&db.write)
147 .await
148 .with_context(|| format!("creating admin {username}"))?;
149 println!("created admin account {username}");
150 db.checkpoint().await?;
151 db.close().await;
152 return Ok(());
153 }
154
155 let loaded = auth::load_tokens(&db.read, &vault, &ingest).await?;
156 info!(
157 tokens = loaded,
158 "loaded device tokens into the ingest cache"
159 );
160
161 // Sessions live in the same SQLite file, on the writer pool. They are
162 // low-traffic enough not to disturb the batching that exists to protect the
163 // position write path.
164 let session_store = SqliteStore::new(db.write.clone());
165 session_store
166 .migrate()
167 .await
168 .context("migrating the session store")?;
169
170 if args.dev {
171 warn!("--dev: session cookies will not require HTTPS. Never use this in production.");
172 }
173 let session_layer = SessionManagerLayer::new(session_store)
174 .with_name(if args.dev { "otsid" } else { "__Host-otsid" })
175 .with_http_only(true)
176 .with_secure(!args.dev)
177 .with_same_site(SameSite::Lax)
178 .with_expiry(Expiry::OnInactivity(time::Duration::days(SESSION_DAYS)));
179
180 let http_addr = cfg.http_addr;
181 let udp_addr = cfg.udp_addr;
182 let workers = cfg.udp_worker_count();
183 let retention_days = cfg.retention_days;
184 let token_stale_days = cfg.token_stale_days;
185
186 let throttle = Arc::new(auth::LoginThrottle::default());
187 let state: api::Shared = Arc::new(AppState {
188 db,
189 cfg,
190 vault,
191 ingest: Arc::clone(&ingest),
192 writer,
193 throttle: Arc::clone(&throttle),
194 });
195
196 let udp_tasks = udp::spawn(udp_addr, workers, Arc::clone(&ingest))?;
197
198 let retention_task = retention::Retention {
199 pool: state.db.write.clone(),
200 retention_days,
201 token_stale_days,
202 ingest: Arc::clone(&ingest),
203 }
204 .spawn();
205 let limits_task = retention::spawn_limits_gc(Arc::clone(&ingest), throttle);
206
207 let app = api::router(Arc::clone(&state)).layer(session_layer);
208
209 let listener = tokio::net::TcpListener::bind(http_addr)
210 .await
211 .with_context(|| format!("binding {http_addr}"))?;
212 info!(addr = %http_addr, "HTTP listener started (put a TLS-terminating proxy in front)");
213
214 if args.simulate_device {
215 // After the listeners are up, so the login cannot race them.
216 simulate::SimulatedDevice {
217 base_url: format!("http://{http_addr}"),
218 username: std::env::var("OT_SIM_USER").unwrap_or_else(|_| "sim".into()),
219 password: std::env::var("OT_SIM_PASSWORD").unwrap_or_else(|_| "simsimsimsim".into()),
220 udp_addr: loopback_of(udp_addr),
221 }
222 .spawn();
223 }
224
225 axum::serve(
226 listener,
227 app.into_make_service_with_connect_info::<SocketAddr>(),
228 )
229 .with_graceful_shutdown(shutdown_signal())
230 .await
231 .context("HTTP server failed")?;
232
233 // Stop the periodic tasks and the receive loops, then let the writer drain so
234 // nothing already acknowledged is lost, then checkpoint so the .db file on
235 // disk is self-contained.
236 info!("shutting down");
237 retention_task.abort();
238 limits_task.abort();
239 for task in udp_tasks {
240 task.abort();
241 }
242 if let Err(e) = state.db.checkpoint().await {
243 warn!(error = %e, "WAL checkpoint failed");
244 }
245 drop(state);
246 let _ = tokio::time::timeout(std::time::Duration::from_secs(5), writer_task).await;
247 Ok(())
248}
249
250/// The simulated device must talk to a routable address; the listener is usually
251/// bound to the wildcard, which cannot be a destination.
252fn loopback_of(addr: SocketAddr) -> String {
253 if addr.ip().is_unspecified() {
254 format!("127.0.0.1:{}", addr.port())
255 } else {
256 addr.to_string()
257 }
258}
259
260async fn shutdown_signal() {
261 let ctrl_c = async {
262 tokio::signal::ctrl_c().await.ok();
263 };
264 #[cfg(unix)]
265 let terminate = async {
266 if let Ok(mut sig) =
267 tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
268 {
269 sig.recv().await;
270 }
271 };
272 #[cfg(not(unix))]
273 let terminate = std::future::pending::<()>();
274
275 tokio::select! {
276 () = ctrl_c => {},
277 () = terminate => {},
278 }
279}
280
281#[cfg(test)]
282mod tests {
283 use super::*;
284
285 #[test]
286 fn the_wildcard_address_is_rewritten_for_the_simulated_device() {
287 assert_eq!(
288 loopback_of("0.0.0.0:7373".parse().expect("literal")),
289 "127.0.0.1:7373"
290 );
291 assert_eq!(
292 loopback_of("192.0.2.1:7373".parse().expect("literal")),
293 "192.0.2.1:7373"
294 );
295 }
296}
297