recall_server/server/mod.rs
1//! The HTTP surface.
2//!
3//! Routes and their auth posture are frozen:
4//!
5//! | route | auth |
6//! |---|---|
7//! | `GET /health` | none — uptime tooling holds no secret |
8//! | `GET /admin` | none — static markup, no data |
9//! | `GET /.well-known/recall` | none — a client asks before it can authenticate |
10//! | `POST /sync`, `GET /sync`, `GET /v1/devices/me` | bearer token, or the signature of any device but a worker |
11//! | `GET /v1/audit/checkpoint`, `GET /v1/audit/consistency` | bearer token, or any device's signature |
12//! | `POST /v1/devices/enroll`, `POST /v1/devices/enroll/poll` | none, but rate limited, and small bodies only |
13//! | `GET /admin/stats`, the rest of `/v1/devices`, `/v1/authkeys`, and `GET /v1/audit/entries` | bearer token, an admin device's signature, or the admin page's passkey session (with its CSRF header on a POST); small bodies only |
14//! | `GET /v1/jobs`, `POST /v1/jobs/{id}/retry` | bearer token, or an admin device's signature |
15//! | `POST /v1/jobs/claim`, `POST /v1/jobs/{id}/result` | a worker device's signature, and nothing else |
16//! | `GET /admin/session`, `POST /admin/login/start`, `POST /admin/login/finish` | none, but rate limited |
17//! | `POST /admin/bootstrap/register` and `…/finish` | bearer token only, with the one-time bootstrap code, and only while no passkey exists |
18//! | `GET /admin/passkeys`, `POST /admin/passkeys/…`, `POST /admin/logout`, `POST /admin/logout/others` | the passkey session only, with its CSRF header on a POST; adding or removing a passkey and signing out the others also need a sign-in in the last five minutes |
19//! | anything else | 404 JSON |
20//!
21//! This module owns the shared state, the router, and the background jobs.
22//! What it wires together are private submodules, each living next to its
23//! own tests: `middleware.rs` (rate limiting, then the protocol check, then
24//! auth), `auth.rs` (device signatures and the replay cache),
25//! `handlers.rs` (one function per route), `devices.rs` (the device
26//! routes), `jobs.rs` (the merge queue's routes and its drain), `admin.rs`
27//! (the admin page and its session), `passkeys.rs` (the WebAuthn ceremonies
28//! that start a session), `respond.rs` (the JSON shape of every reply,
29//! errors included), `limit.rs` (the per-IP window the middleware consults)
30//! and `tls.rs` (the direct-TLS accept loop, used only when `Config::tls`
31//! is on; plain HTTP, the default, never touches it). Both transports serve
32//! the one router [`Server::router`] builds, every route group and layer
33//! included.
34
35use std::future::Future;
36use std::net::SocketAddr;
37use std::sync::atomic::{AtomicBool, AtomicI64, Ordering};
38use std::sync::{Arc, PoisonError, RwLock};
39use std::time::{Instant, SystemTime, UNIX_EPOCH};
40
41use anyhow::{Context, Result};
42use axum::extract::DefaultBodyLimit;
43// Imported by name because this module has a `middleware` of its own, and
44// an unqualified `middleware::` would resolve to that one.
45use axum::middleware::{from_fn, from_fn_with_state};
46use axum::routing::{get, post};
47use axum::Router;
48use recall_wire::devices as paths;
49use recall_wire::{ClaudeCliStatus, MergeError};
50use tokio::net::TcpListener;
51use tokio::sync::Notify;
52use tokio::task::JoinHandle;
53
54use crate::config::TlsMode;
55use crate::merge::{Merger, Status};
56use crate::{format_timestamp, now, Config, Store};
57
58// Without passkeys nothing starts a session, so the half of this module
59// that makes one goes unused; the half that checks one still runs.
60#[cfg_attr(not(feature = "passkeys"), allow(dead_code))]
61mod admin;
62mod audit;
63mod auth;
64mod devices;
65mod handlers;
66mod jobs;
67mod limit;
68mod middleware;
69#[cfg(feature = "passkeys")]
70mod passkeys;
71mod respond;
72mod tls;
73
74use audit::{handle_checkpoint, handle_consistency, handle_entries};
75#[cfg(feature = "passkeys")]
76use passkeys::Passkeys;
77#[cfg(not(feature = "passkeys"))]
78use without_passkeys::Passkeys;
79
80/// What stands in for `passkeys.rs` in a build without the feature: sign-in
81/// is off, and the page says so.
82#[cfg(not(feature = "passkeys"))]
83mod without_passkeys {
84 pub(super) struct Passkeys;
85
86 impl Passkeys {
87 pub(super) fn new(_public_url: &str) -> Self {
88 Self
89 }
90
91 pub(super) fn status(&self) -> super::admin::PasskeyStatus {
92 super::admin::PasskeyStatus {
93 enabled: false,
94 origin: None,
95 reason: Some("this recall-server was built without passkey support".to_string()),
96 }
97 }
98
99 pub(super) fn prune(&self, _now: i64) -> usize {
100 0
101 }
102 }
103}
104
105use auth::ReplayCache;
106use devices::{
107 handle_approve, handle_create_authkey, handle_deny, handle_enroll, handle_list_authkeys,
108 handle_list_devices, handle_me, handle_pending, handle_poll, handle_revoke_authkey,
109 handle_revoke_device,
110};
111use handlers::{
112 handle_admin_stats, handle_discovery, handle_health, handle_pull, handle_push, not_found,
113};
114use limit::RateLimiter;
115use middleware::{admin_guard, admin_only, guard, limited, limited_sign_in, not_worker};
116
117/// How often idle ephemeral devices and long-expired enrolments are swept
118/// away. Removal is at most this late, which against a TTL counted in
119/// hours is nothing.
120const SWEEP_EVERY: std::time::Duration = std::time::Duration::from_secs(10 * 60);
121
122/// Bounds a single push. Memory files are prose; anything this large is a
123/// bug or an attack, not a note.
124const MAX_BODY_BYTES: usize = 5 << 20;
125
126/// Bounds a request to the routes anyone may call. An enrolment is a name,
127/// a key and an agent, and a poll is an id; nobody who has not proved
128/// anything gets to make the server hold megabytes.
129const ENROLL_BODY_BYTES: usize = 8 << 10;
130
131/// Bounds a request to the admin routes. Each body is a code, a scope, a
132/// fingerprint or a tag, and a signed one is kept whole in its audit leaf,
133/// so none may be more than a few kilobytes.
134const ADMIN_BODY_BYTES: usize = 8 << 10;
135
136/// Bounds a request to the admin page's sign-in and passkey routes. A
137/// passkey's answer is a few kilobytes at most, attestation certificates
138/// included.
139const SIGN_IN_BODY_BYTES: usize = 64 << 10;
140
141struct Runtime {
142 last_backup_at: String,
143 last_merge_at: String,
144 last_merge_error: Option<MergeError>,
145 claude_status: Status,
146 /// When a worker last claimed, since this process started.
147 worker_last_claim_at: Option<String>,
148 /// The same, as an instant, or when this process started if no worker
149 /// has claimed since: what says whether a worker's claims have stopped.
150 worker_last_claim: Instant,
151 /// The `User-Agent` of that claim.
152 worker_agent: String,
153 /// The CLI check that claim carried: what `/health` reports as
154 /// `claude_cli` while a worker is enrolled.
155 worker_cli: Option<ClaudeCliStatus>,
156}
157
158struct AppState {
159 cfg: Config,
160 store: Arc<Store>,
161 merger: Merger,
162 started_at: String,
163 /// When this process started, as a UNIX time: signatures the process
164 /// before could have accepted are refused, since the nonces that would
165 /// catch their replay were in its memory (see `auth.rs`). Moved only
166 /// by tests, through [`Server::backdate_start`].
167 started_unix: AtomicI64,
168 /// Added to the clock signatures are judged by. Zero except in tests.
169 clock_offset: AtomicI64,
170 runtime: RwLock<Runtime>,
171 limiter: RateLimiter,
172 replay: ReplayCache,
173 /// Passkey sign-in for the admin page, and the ceremonies it finished.
174 passkeys: Passkeys,
175 /// Wakes waiting claims when a job is queued.
176 jobs_ready: Notify,
177 /// Set once shutdown begins, so waiting claims answer at once rather
178 /// than holding the graceful shutdown for their whole wait.
179 closing: AtomicBool,
180 /// Set while the queue is being drained here, so two drains never run
181 /// at once.
182 draining: AtomicBool,
183}
184
185impl AppState {
186 fn read(&self) -> std::sync::RwLockReadGuard<'_, Runtime> {
187 self.runtime.read().unwrap_or_else(PoisonError::into_inner)
188 }
189 fn write(&self) -> std::sync::RwLockWriteGuard<'_, Runtime> {
190 self.runtime.write().unwrap_or_else(PoisonError::into_inner)
191 }
192
193 /// The UNIX time signatures are judged by.
194 fn now(&self) -> i64 {
195 unix_now() + self.clock_offset.load(Ordering::Relaxed)
196 }
197
198 /// The same clock, as the moment admin sessions are judged by.
199 fn clock(&self) -> time::OffsetDateTime {
200 time::OffsetDateTime::now_utc()
201 + time::Duration::seconds(self.clock_offset.load(Ordering::Relaxed))
202 }
203
204 /// When this process started, as a UNIX time.
205 fn started(&self) -> i64 {
206 self.started_unix.load(Ordering::Relaxed)
207 }
208}
209
210fn unix_now() -> i64 {
211 SystemTime::now()
212 .duration_since(UNIX_EPOCH)
213 .map(|d| d.as_secs() as i64)
214 .unwrap_or(0)
215}
216
217/// How many live nonces one device may have: as many requests as one
218/// address may make while a nonce stays live, which is up to the window
219/// and the few seconds `created` may be ahead of the clock.
220fn nonces_per_device(cfg: &Config) -> usize {
221 let live_ms = auth::NONCE_LIFETIME as u128 * 1000;
222 let window_ms = cfg.rate_limit_window.as_millis().max(1);
223 let windows = live_ms.div_ceil(window_ms);
224 (cfg.rate_limit_max as u128 * windows).min(usize::MAX as u128) as usize
225}
226
227/// The HTTP API, its background jobs, and the state they share.
228pub struct Server {
229 state: Arc<AppState>,
230}
231
232impl Server {
233 /// Builds a server around an already-open store. Nothing is recorded
234 /// yet: the `start` leaf waits until the server is serving (see
235 /// [`Server::serve_with_shutdown`]).
236 pub fn new(cfg: Config, store: Arc<Store>) -> Self {
237 let limiter = RateLimiter::new(cfg.rate_limit_window, cfg.rate_limit_max);
238 let merger = Merger::new(cfg.claude_bin.clone(), cfg.merge_timeout);
239 let replay = ReplayCache::new(auth::WINDOW, nonces_per_device(&cfg));
240 let passkeys = Passkeys::new(&cfg.public_url);
241 Self {
242 state: Arc::new(AppState {
243 cfg,
244 store,
245 merger,
246 started_at: now(),
247 started_unix: AtomicI64::new(unix_now()),
248 clock_offset: AtomicI64::new(0),
249 runtime: RwLock::new(Runtime {
250 last_backup_at: String::new(),
251 last_merge_at: String::new(),
252 last_merge_error: None,
253 claude_status: Status::default(),
254 worker_last_claim_at: None,
255 worker_last_claim: Instant::now(),
256 worker_agent: String::new(),
257 worker_cli: None,
258 }),
259 limiter,
260 replay,
261 passkeys,
262 jobs_ready: Notify::new(),
263 closing: AtomicBool::new(false),
264 draining: AtomicBool::new(false),
265 }),
266 }
267 }
268
269 /// The router, built separately from binding a port so tests can drive
270 /// it without real sockets.
271 pub fn router(&self) -> Router {
272 let state = self.state.clone();
273 // Managing devices: authenticated, then held to the admin scope.
274 // The last layer added runs first, so `guard` has put the caller
275 // in place by the time `admin_only` looks for it.
276 let admin = Router::new()
277 .route("/admin/stats", get(handle_admin_stats).fallback(not_found))
278 .route(
279 paths::DEVICES_PATH,
280 get(handle_list_devices).fallback(not_found),
281 )
282 .route(
283 paths::APPROVE_PATH,
284 post(handle_approve).fallback(not_found),
285 )
286 .route(paths::DENY_PATH, post(handle_deny).fallback(not_found))
287 .route(
288 "/v1/devices/pending/{user_code}",
289 get(handle_pending).fallback(not_found),
290 )
291 .route(
292 "/v1/devices/{id}/revoke",
293 post(handle_revoke_device).fallback(not_found),
294 )
295 .route(
296 paths::AUTHKEYS_PATH,
297 get(handle_list_authkeys)
298 .post(handle_create_authkey)
299 .fallback(not_found),
300 )
301 .route(
302 "/v1/authkeys/{id}/revoke",
303 post(handle_revoke_authkey).fallback(not_found),
304 )
305 // The leaves themselves: every project, file, device and
306 // authkey the log names, which is what the device list and the
307 // stats already keep to this scope. The checkpoint and the
308 // proofs, hashes only, stay open to any credential below.
309 .route(
310 recall_wire::audit::ENTRIES_PATH,
311 get(handle_entries).fallback(not_found),
312 )
313 .route_layer(DefaultBodyLimit::max(ADMIN_BODY_BYTES))
314 .route_layer(from_fn(admin_only))
315 .route_layer(from_fn_with_state(state.clone(), admin_guard));
316 // The audit log's hashes: any credential, a worker's included, since
317 // a checkpoint and a proof name nothing (the leaves themselves are
318 // admin, above).
319 let audit_routes = Router::new()
320 .route(
321 recall_wire::audit::CHECKPOINT_PATH,
322 get(handle_checkpoint).fallback(not_found),
323 )
324 .route(
325 recall_wire::audit::CONSISTENCY_PATH,
326 get(handle_consistency).fallback(not_found),
327 )
328 .route_layer(from_fn_with_state(state.clone(), guard));
329 // Enrolling: a machine has no credential yet, so no auth, but the
330 // same rate limit and protocol check as everything else, and a
331 // body limit sized for what an enrolment is. The inner limit wins
332 // over the router-wide one below.
333 let enrolment = Router::new()
334 .route(paths::ENROLL_PATH, post(handle_enroll).fallback(not_found))
335 .route(
336 paths::ENROLL_POLL_PATH,
337 post(handle_poll).fallback(not_found),
338 )
339 .route_layer(DefaultBodyLimit::max(ENROLL_BODY_BYTES))
340 .route_layer(from_fn_with_state(state.clone(), limited));
341 // What the admin page asks before it knows who is looking: open to
342 // anyone, rate limited.
343 let page = Router::new()
344 .route(
345 "/admin/session",
346 get(admin::handle_session_status).fallback(not_found),
347 )
348 .route_layer(from_fn_with_state(state.clone(), limited_sign_in));
349 Router::new()
350 // Go's mux dispatched every method through one guarded handler
351 // and 404'd the ones it didn't implement; the method fallbacks
352 // keep that shape (and its JSON body) instead of axum's bare
353 // 405.
354 .route(
355 "/sync",
356 get(handle_pull).post(handle_push).fallback(not_found),
357 )
358 .route(paths::DEVICES_ME_PATH, get(handle_me).fallback(not_found))
359 // Registered before the layers, so only these routes are rate
360 // limited and authenticated here. The last layer added runs
361 // first: `guard` puts the caller in place, then `not_worker`
362 // keeps a worker device out of memory.
363 .route_layer(from_fn(not_worker))
364 .route_layer(from_fn_with_state(state.clone(), guard))
365 .merge(admin)
366 .merge(audit_routes)
367 .merge(enrolment)
368 .merge(page)
369 .merge(sign_in_routes(&state))
370 .merge(jobs::routes(state.clone()))
371 .route("/health", get(handle_health).fallback(not_found))
372 .route(
373 recall_wire::DISCOVERY_PATH,
374 get(handle_discovery).fallback(not_found),
375 )
376 .route("/admin", get(admin::handle_admin_page).fallback(not_found))
377 .fallback(not_found)
378 .layer(DefaultBodyLimit::max(MAX_BODY_BYTES))
379 .with_state(state)
380 }
381
382 /// Re-runs the local `claude auth status` probe.
383 pub async fn refresh_claude_status(&self) {
384 let status = self.state.merger.check_status().await;
385 self.state.write().claude_status = status;
386 }
387
388 /// The last known state of the `claude` CLI.
389 pub fn claude_status(&self) -> Status {
390 self.state.read().claude_status.clone()
391 }
392
393 /// Overrides the cached CLI status.
394 ///
395 /// Exposed so tests can exercise the merge-failure path — reaching it
396 /// otherwise needs a real, logged-in CLI on the machine running them.
397 pub fn set_claude_status(&self, status: Status) {
398 self.state.write().claude_status = status;
399 }
400
401 /// Moves the clock device signatures are judged by, in seconds.
402 ///
403 /// Exposed so tests can reach the edges of the signature window, and
404 /// what happens after it, without waiting a minute for each.
405 pub fn set_clock_offset(&self, seconds: i64) {
406 self.state.clock_offset.store(seconds, Ordering::Relaxed);
407 }
408
409 /// Makes this server act as if it had started `seconds` earlier than
410 /// it did.
411 ///
412 /// Exposed so tests can sign requests at once, rather than waiting out
413 /// the few seconds after a start in which every signature is refused.
414 pub fn backdate_start(&self, seconds: i64) {
415 self.state
416 .started_unix
417 .fetch_sub(seconds, Ordering::Relaxed);
418 }
419
420 /// Makes the last claim by a worker, or this server's start if there
421 /// has been none, `seconds` older than it is.
422 ///
423 /// Exposed so tests can reach a worker whose claims have stopped
424 /// without waiting minutes for it.
425 pub fn backdate_last_claim(&self, seconds: u64) {
426 let mut rt = self.state.write();
427 if let Some(earlier) = rt
428 .worker_last_claim
429 .checked_sub(std::time::Duration::from_secs(seconds))
430 {
431 rt.worker_last_claim = earlier;
432 }
433 }
434
435 /// Merges the jobs left in the queue here, or marks them failed when
436 /// this server cannot merge, if no worker is enrolled; otherwise does
437 /// nothing. Run by itself when the last worker is revoked and on every
438 /// sweep; exposed so tests can wait for it.
439 pub async fn drain_jobs(&self) -> Result<()> {
440 jobs::drain_without_worker(&self.state).await
441 }
442
443 /// Writes a backup now. Failure is logged, never propagated: it becomes
444 /// visible through `/health`'s `last_backup_at` going stale.
445 pub fn run_backup(&self) {
446 run_backup(&self.state);
447 }
448
449 /// Issues a new bootstrap code, when passkey sign-in is on and no
450 /// passkey is registered: the one-time code that, with `RECALL_TOKEN`,
451 /// registers the first. [`None`] otherwise. `recall-server` calls it on
452 /// start and prints what it gets; see [`crate::bootstrap`].
453 pub fn issue_bootstrap_code(&self) -> Result<Option<crate::bootstrap::BootstrapCode>> {
454 if !self.state.passkeys.status().enabled || self.state.store.has_admin_credentials()? {
455 return Ok(None);
456 }
457 crate::bootstrap::issue(&self.state.store, self.state.clock()).map(Some)
458 }
459
460 /// Where the admin page is, when `RECALL_PUBLIC_URL` says.
461 pub fn public_url(&self) -> Option<&str> {
462 Some(self.state.cfg.public_url.as_str()).filter(|u| !u.is_empty())
463 }
464
465 /// Removes ephemeral devices idle for longer than
466 /// [`Config::ephemeral_device_ttl`], and enrolments that expired over
467 /// an hour ago. Answers how many of each went.
468 pub fn sweep_devices(&self) -> Result<(usize, usize)> {
469 sweep_devices(&self.state)
470 }
471
472 /// Starts background work: the first Claude CLI status check, its
473 /// refresh loop, backups, the device sweep, and draining a queue no
474 /// worker is left to take. All of it is best-effort — none of it may
475 /// take the sync API down.
476 pub fn start_background(&self) -> Vec<JoinHandle<()>> {
477 let mut tasks = Vec::new();
478 {
479 let state = self.state.clone();
480 tasks.push(tokio::spawn(async move {
481 loop {
482 let s = state.clone();
483 match tokio::task::spawn_blocking(move || sweep_devices(&s)).await {
484 Ok(Ok((0, 0))) => {}
485 Ok(Ok((devices, enrollments))) => eprintln!(
486 "removed {devices} idle ephemeral devices and {enrollments} expired enrolments"
487 ),
488 Ok(Err(e)) => eprintln!("device sweep failed: {e:#}"),
489 Err(_) => {}
490 }
491 let s = state.clone();
492 match tokio::task::spawn_blocking(move || jobs::prune(&s)).await {
493 Ok(Ok(0)) | Err(_) => {}
494 Ok(Ok(n)) => eprintln!("removed {n} finished merge jobs"),
495 Ok(Err(e)) => eprintln!("job prune failed: {e:#}"),
496 }
497 // Jobs no worker is left to take: merged here, or
498 // marked failed. The first sweep usually comes before
499 // the first check of the CLI has finished, and leaves
500 // them to the drain that check starts.
501 if let Err(e) = jobs::drain_without_worker(&state).await {
502 eprintln!("draining the merge queue failed: {e:#}");
503 }
504 tokio::time::sleep(SWEEP_EVERY).await;
505 }
506 }));
507 }
508 if self.state.cfg.merge_enabled {
509 let state = self.state.clone();
510 tasks.push(tokio::spawn(async move {
511 let every = state.cfg.claude_status_interval;
512 let mut first = true;
513 loop {
514 let status = state.merger.check_status().await;
515 state.write().claude_status = status;
516 // Now that it is known whether this server can merge,
517 // jobs a worker left before the last restart need not
518 // wait for the next sweep.
519 if std::mem::take(&mut first) {
520 if let Err(e) = jobs::drain_without_worker(&state).await {
521 eprintln!("draining the merge queue failed: {e:#}");
522 }
523 }
524 tokio::time::sleep(every).await;
525 }
526 }));
527 }
528 if !self.state.cfg.backup_dir.is_empty() {
529 let state = self.state.clone();
530 tasks.push(tokio::spawn(async move {
531 let every = state.cfg.backup_interval;
532 loop {
533 // VACUUM INTO can take a while on a large database and
534 // holds the store lock, so it stays off the async
535 // worker threads.
536 let s = state.clone();
537 let _ = tokio::task::spawn_blocking(move || run_backup(&s)).await;
538 tokio::time::sleep(every).await;
539 }
540 }));
541 }
542 tasks
543 }
544
545 /// Binds `cfg.addr` and serves until SIGTERM or ctrl-c, then shuts down
546 /// gracefully so an in-flight merge isn't cut off mid-write.
547 pub async fn serve(&self) -> Result<()> {
548 let listener = TcpListener::bind(&self.state.cfg.addr)
549 .await
550 .with_context(|| format!("binding {}", self.state.cfg.addr))?;
551 self.serve_with_shutdown(listener, shutdown_signal()).await
552 }
553
554 /// Serves on an already-bound listener until `shutdown` resolves, in
555 /// whichever transport `cfg.tls` names (see `server/tls.rs`; plain HTTP,
556 /// the default, still goes through `axum::serve` directly, unchanged).
557 ///
558 /// Once the transport is ready, records a `start` leaf in the audit
559 /// log: the server's own doing, so its actor is
560 /// [`crate::audit::leaf::Actor::Server`], and its subject the version
561 /// that started. Here rather than in [`Server::new`] so that only a
562 /// server that got its port, and its certificate, records one: one that
563 /// failed to bind because another was still running, or to load its
564 /// key, leaves nothing. Best-effort — a store the audit table somehow
565 /// cannot be written to still serves sync, the same way a failed backup
566 /// does not take the server down.
567 pub async fn serve_with_shutdown<F>(&self, listener: TcpListener, shutdown: F) -> Result<()>
568 where
569 F: Future<Output = ()> + Send + 'static,
570 {
571 // The certificate is loaded (or the ACME state built) before
572 // anything claims the server is up, so a bad path or an unreadable
573 // key is the last line in the log rather than one after
574 // "listening".
575 let transport = match &self.state.cfg.tls {
576 TlsMode::Off => None,
577 mode => Some(tls::prepare(mode).await?),
578 };
579 eprintln!(
580 "recall server listening on {} ({}, db: {})",
581 listener
582 .local_addr()
583 .map_or_else(|_| self.state.cfg.addr.clone(), |a| a.to_string()),
584 transport
585 .as_ref()
586 .map_or("plain http", tls::Prepared::description),
587 self.state.cfg.db_path
588 );
589 if let Err(e) = record_start(&self.state.store) {
590 eprintln!("recording server start in the audit log: {e:#}");
591 }
592 let tasks = self.start_background();
593 let state = self.state.clone();
594 let shutdown = async move {
595 shutdown.await;
596 // A claim waiting for a job would otherwise hold the shutdown
597 // for up to its whole wait.
598 state.closing.store(true, Ordering::Relaxed);
599 state.jobs_ready.notify_waiters();
600 };
601 let result = match transport {
602 None => axum::serve(
603 listener,
604 self.router()
605 .into_make_service_with_connect_info::<SocketAddr>(),
606 )
607 .with_graceful_shutdown(shutdown)
608 .await
609 .map_err(Into::into),
610 Some(prepared) => {
611 // axum-server runs its own accept loop rather than
612 // axum::serve's, so the listener crosses over to std here.
613 // It is already non-blocking (tokio bound it), which is
614 // exactly what tokio::net::TcpListener::from_std, which
615 // axum-server calls internally, requires.
616 let listener = listener.into_std().context("preparing the TLS listener")?;
617 let limits = tls::Limits::from_config(&self.state.cfg);
618 tls::serve(self.router(), listener, prepared, limits, shutdown).await
619 }
620 };
621 for task in tasks {
622 task.abort();
623 }
624 result
625 }
626}
627
628/// Appends the `start` leaf [`Server::serve_with_shutdown`] records.
629fn record_start(store: &Store) -> Result<()> {
630 let version = recall_wire::discovery::version();
631 store.audit_append(|seq, at| {
632 crate::audit::leaf::encode(
633 seq,
634 at,
635 crate::audit::leaf::action::START,
636 &crate::audit::leaf::Actor::Server,
637 crate::audit::leaf::subject_start(&version),
638 None,
639 )
640 })?;
641 Ok(())
642}
643
644/// The passkey routes: signing in, the bootstrap, and what a signed-in
645/// owner does with passkeys. None of them exists in a build without the
646/// feature, where each is the usual 404.
647#[cfg(feature = "passkeys")]
648fn sign_in_routes(state: &Arc<AppState>) -> Router<Arc<AppState>> {
649 use passkeys::{
650 handle_add_finish, handle_add_start, handle_bootstrap_finish, handle_bootstrap_start,
651 handle_list_passkeys, handle_remove, handle_sign_in_finish, handle_sign_in_start,
652 handle_sign_out, handle_sign_out_others, json_only,
653 };
654 // Every ceremony's start and finish takes JSON and says so, which a
655 // cross-site form or a blind `no-cors` fetch cannot. The last layer
656 // added runs first, so each router's own check comes before this one.
657 //
658 // Anyone may try to sign in; only a registered passkey finishes.
659 let sign_in = Router::new()
660 .route(
661 "/admin/login/start",
662 post(handle_sign_in_start).fallback(not_found),
663 )
664 .route(
665 "/admin/login/finish",
666 post(handle_sign_in_finish).fallback(not_found),
667 )
668 .route_layer(from_fn(json_only))
669 .route_layer(DefaultBodyLimit::max(SIGN_IN_BODY_BYTES))
670 .route_layer(from_fn_with_state(state.clone(), limited_sign_in));
671 // The first passkey: the operator's token, and the handlers refuse it
672 // once any passkey exists, whatever the token.
673 let bootstrap = Router::new()
674 .route(
675 "/admin/bootstrap/register",
676 post(handle_bootstrap_start).fallback(not_found),
677 )
678 .route(
679 "/admin/bootstrap/register/finish",
680 post(handle_bootstrap_finish).fallback(not_found),
681 )
682 .route_layer(from_fn(json_only))
683 .route_layer(DefaultBodyLimit::max(SIGN_IN_BODY_BYTES))
684 .route_layer(from_fn_with_state(state.clone(), guard));
685 // Only a signed-in owner: never the token, never a device.
686 let adding = Router::new()
687 .route(
688 "/admin/passkeys/register",
689 post(handle_add_start).fallback(not_found),
690 )
691 .route(
692 "/admin/passkeys/register/finish",
693 post(handle_add_finish).fallback(not_found),
694 )
695 .route_layer(from_fn(json_only))
696 .route_layer(DefaultBodyLimit::max(SIGN_IN_BODY_BYTES))
697 .route_layer(from_fn_with_state(state.clone(), admin::owner_only));
698 let owner = Router::new()
699 .route(
700 "/admin/passkeys",
701 get(handle_list_passkeys).fallback(not_found),
702 )
703 .route(
704 "/admin/passkeys/{id}/remove",
705 post(handle_remove).fallback(not_found),
706 )
707 .route("/admin/logout", post(handle_sign_out).fallback(not_found))
708 .route(
709 "/admin/logout/others",
710 post(handle_sign_out_others).fallback(not_found),
711 )
712 .route_layer(DefaultBodyLimit::max(SIGN_IN_BODY_BYTES))
713 .route_layer(from_fn_with_state(state.clone(), admin::owner_only));
714 sign_in.merge(bootstrap).merge(adding).merge(owner)
715}
716
717#[cfg(not(feature = "passkeys"))]
718fn sign_in_routes(_state: &Arc<AppState>) -> Router<Arc<AppState>> {
719 Router::new()
720}
721
722fn sweep_devices(state: &AppState) -> Result<(usize, usize)> {
723 let now = time::OffsetDateTime::now_utc();
724 // The admin page's leftovers go on the same round: ceremonies finished
725 // that would have expired by now, and sessions that have ended.
726 state.passkeys.prune(state.now());
727 let clock = state.clock();
728 if let Err(e) = state.store.sweep_admin_sessions(
729 &format_timestamp(clock),
730 &format_timestamp(clock - admin::SESSION_IDLE),
731 ) {
732 eprintln!("admin session sweep failed: {e:#}");
733 }
734 state.store.sweep_devices_audited(
735 &format_timestamp(now - state.cfg.ephemeral_device_ttl),
736 &format_timestamp(now - devices::EXPIRED_ENROLLMENT_KEPT),
737 )
738}
739
740fn run_backup(state: &AppState) {
741 if state.cfg.backup_dir.is_empty() {
742 return;
743 }
744 match state
745 .store
746 .backup(&state.cfg.backup_dir, state.cfg.backup_keep)
747 {
748 Ok(dest) => {
749 state.write().last_backup_at = now();
750 eprintln!("backup written: {}", dest.display());
751 }
752 Err(e) => eprintln!("backup failed: {e:#}"),
753 }
754}
755
756async fn shutdown_signal() {
757 let ctrl_c = async {
758 let _ = tokio::signal::ctrl_c().await;
759 };
760 #[cfg(unix)]
761 let terminate = async {
762 match tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) {
763 Ok(mut sig) => {
764 sig.recv().await;
765 }
766 Err(_) => std::future::pending::<()>().await,
767 }
768 };
769 #[cfg(not(unix))]
770 let terminate = std::future::pending::<()>();
771
772 tokio::select! {
773 _ = ctrl_c => {}
774 _ = terminate => {}
775 }
776}