Skip to main content

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                let mut wal = WalWatch::default();
482                loop {
483                    let s = state.clone();
484                    match tokio::task::spawn_blocking(move || sweep_devices(&s)).await {
485                        Ok(Ok((0, 0))) => {}
486                        Ok(Ok((devices, enrollments))) => eprintln!(
487                            "removed {devices} idle ephemeral devices and {enrollments} expired enrolments"
488                        ),
489                        Ok(Err(e)) => eprintln!("device sweep failed: {e:#}"),
490                        Err(_) => {}
491                    }
492                    let s = state.clone();
493                    match tokio::task::spawn_blocking(move || jobs::prune(&s)).await {
494                        Ok(Ok(0)) | Err(_) => {}
495                        Ok(Ok(n)) => eprintln!("removed {n} finished merge jobs"),
496                        Ok(Err(e)) => eprintln!("job prune failed: {e:#}"),
497                    }
498                    // Jobs no worker is left to take: merged here, or
499                    // marked failed. The first sweep usually comes before
500                    // the first check of the CLI has finished, and leaves
501                    // them to the drain that check starts.
502                    if let Err(e) = jobs::drain_without_worker(&state).await {
503                        eprintln!("draining the merge queue failed: {e:#}");
504                    }
505                    // The WAL's commits copied back into recall.db, so the
506                    // file on its own is never more than a sweep behind:
507                    // see Store::checkpoint.
508                    let s = state.clone();
509                    match tokio::task::spawn_blocking(move || s.store.checkpoint()).await {
510                        Ok(Ok(complete)) => {
511                            if let Some(said) = wal.record(complete) {
512                                eprintln!("{said}");
513                            }
514                        }
515                        Ok(Err(e)) => eprintln!("checkpointing the WAL failed: {e:#}"),
516                        Err(_) => {}
517                    }
518                    tokio::time::sleep(SWEEP_EVERY).await;
519                }
520            }));
521        }
522        if self.state.cfg.merge_enabled {
523            let state = self.state.clone();
524            tasks.push(tokio::spawn(async move {
525                let every = state.cfg.claude_status_interval;
526                let mut first = true;
527                loop {
528                    let status = state.merger.check_status().await;
529                    state.write().claude_status = status;
530                    // Now that it is known whether this server can merge,
531                    // jobs a worker left before the last restart need not
532                    // wait for the next sweep.
533                    if std::mem::take(&mut first) {
534                        if let Err(e) = jobs::drain_without_worker(&state).await {
535                            eprintln!("draining the merge queue failed: {e:#}");
536                        }
537                    }
538                    tokio::time::sleep(every).await;
539                }
540            }));
541        }
542        if !self.state.cfg.backup_dir.is_empty() {
543            let state = self.state.clone();
544            tasks.push(tokio::spawn(async move {
545                let every = state.cfg.backup_interval;
546                loop {
547                    // VACUUM INTO can take a while on a large database and
548                    // holds the store lock, so it stays off the async
549                    // worker threads.
550                    let s = state.clone();
551                    let _ = tokio::task::spawn_blocking(move || run_backup(&s)).await;
552                    tokio::time::sleep(every).await;
553                }
554            }));
555        }
556        tasks
557    }
558
559    /// Binds `cfg.addr` and serves until SIGTERM or ctrl-c, then shuts down
560    /// gracefully so an in-flight merge isn't cut off mid-write, and empties
561    /// the WAL into the database file ([`Store::checkpoint_all`]).
562    pub async fn serve(&self) -> Result<()> {
563        let listener = TcpListener::bind(&self.state.cfg.addr)
564            .await
565            .with_context(|| format!("binding {}", self.state.cfg.addr))?;
566        self.serve_with_shutdown(listener, shutdown_signal()).await
567    }
568
569    /// Serves on an already-bound listener until `shutdown` resolves, in
570    /// whichever transport `cfg.tls` names (see `server/tls.rs`; plain HTTP,
571    /// the default, still goes through `axum::serve` directly, unchanged).
572    ///
573    /// Once the transport is ready, records a `start` leaf in the audit
574    /// log: the server's own doing, so its actor is
575    /// [`crate::audit::leaf::Actor::Server`], and its subject the version
576    /// that started. Here rather than in [`Server::new`] so that only a
577    /// server that got its port, and its certificate, records one: one that
578    /// failed to bind because another was still running, or to load its
579    /// key, leaves nothing. Best-effort — a store the audit table somehow
580    /// cannot be written to still serves sync, the same way a failed backup
581    /// does not take the server down.
582    pub async fn serve_with_shutdown<F>(&self, listener: TcpListener, shutdown: F) -> Result<()>
583    where
584        F: Future<Output = ()> + Send + 'static,
585    {
586        // The certificate is loaded (or the ACME state built) before
587        // anything claims the server is up, so a bad path or an unreadable
588        // key is the last line in the log rather than one after
589        // "listening".
590        let transport = match &self.state.cfg.tls {
591            TlsMode::Off => None,
592            mode => Some(tls::prepare(mode).await?),
593        };
594        eprintln!(
595            "recall server listening on {} ({}, db: {})",
596            listener
597                .local_addr()
598                .map_or_else(|_| self.state.cfg.addr.clone(), |a| a.to_string()),
599            transport
600                .as_ref()
601                .map_or("plain http", tls::Prepared::description),
602            self.state.cfg.db_path
603        );
604        if let Err(e) = record_start(&self.state.store) {
605            eprintln!("recording server start in the audit log: {e:#}");
606        }
607        let tasks = self.start_background();
608        let state = self.state.clone();
609        let shutdown = async move {
610            shutdown.await;
611            // A claim waiting for a job would otherwise hold the shutdown
612            // for up to its whole wait.
613            state.closing.store(true, Ordering::Relaxed);
614            state.jobs_ready.notify_waiters();
615        };
616        let result = match transport {
617            None => axum::serve(
618                listener,
619                self.router()
620                    .into_make_service_with_connect_info::<SocketAddr>(),
621            )
622            .with_graceful_shutdown(shutdown)
623            .await
624            .map_err(Into::into),
625            Some(prepared) => {
626                // axum-server runs its own accept loop rather than
627                // axum::serve's, so the listener crosses over to std here.
628                // It is already non-blocking (tokio bound it), which is
629                // exactly what tokio::net::TcpListener::from_std, which
630                // axum-server calls internally, requires.
631                let listener = listener.into_std().context("preparing the TLS listener")?;
632                let limits = tls::Limits::from_config(&self.state.cfg);
633                tls::serve(self.router(), listener, prepared, limits, shutdown).await
634            }
635        };
636        for task in tasks {
637            task.abort();
638        }
639        // Last, with every request answered: the WAL emptied into recall.db,
640        // so the file left behind is the whole database, even while
641        // sqlite-web has it open, which keeps closing the connection from
642        // doing the same. Best-effort, like the backups: a WAL it could not
643        // empty is still read on the next start.
644        match self.state.store.checkpoint_all() {
645            Ok(true) => {}
646            Ok(false) => eprintln!(
647                "stopping with commits still in the WAL: a reader held it past the busy \
648                 timeout. Nothing is lost; the next start reads them"
649            ),
650            Err(e) => eprintln!("checkpointing the WAL at shutdown failed: {e:#}"),
651        }
652        result
653    }
654}
655
656/// Appends the `start` leaf [`Server::serve_with_shutdown`] records.
657fn record_start(store: &Store) -> Result<()> {
658    let version = recall_wire::discovery::version();
659    store.audit_append(|seq, at| {
660        crate::audit::leaf::encode(
661            seq,
662            at,
663            crate::audit::leaf::action::START,
664            &crate::audit::leaf::Actor::Server,
665            crate::audit::leaf::subject_start(&version),
666            None,
667        )
668    })?;
669    Ok(())
670}
671
672/// The passkey routes: signing in, the bootstrap, and what a signed-in
673/// owner does with passkeys. None of them exists in a build without the
674/// feature, where each is the usual 404.
675#[cfg(feature = "passkeys")]
676fn sign_in_routes(state: &Arc<AppState>) -> Router<Arc<AppState>> {
677    use passkeys::{
678        handle_add_finish, handle_add_start, handle_bootstrap_finish, handle_bootstrap_start,
679        handle_list_passkeys, handle_remove, handle_sign_in_finish, handle_sign_in_start,
680        handle_sign_out, handle_sign_out_others, json_only,
681    };
682    // Every ceremony's start and finish takes JSON and says so, which a
683    // cross-site form or a blind `no-cors` fetch cannot. The last layer
684    // added runs first, so each router's own check comes before this one.
685    //
686    // Anyone may try to sign in; only a registered passkey finishes.
687    let sign_in = Router::new()
688        .route(
689            "/admin/login/start",
690            post(handle_sign_in_start).fallback(not_found),
691        )
692        .route(
693            "/admin/login/finish",
694            post(handle_sign_in_finish).fallback(not_found),
695        )
696        .route_layer(from_fn(json_only))
697        .route_layer(DefaultBodyLimit::max(SIGN_IN_BODY_BYTES))
698        .route_layer(from_fn_with_state(state.clone(), limited_sign_in));
699    // The first passkey: the operator's token, and the handlers refuse it
700    // once any passkey exists, whatever the token.
701    let bootstrap = Router::new()
702        .route(
703            "/admin/bootstrap/register",
704            post(handle_bootstrap_start).fallback(not_found),
705        )
706        .route(
707            "/admin/bootstrap/register/finish",
708            post(handle_bootstrap_finish).fallback(not_found),
709        )
710        .route_layer(from_fn(json_only))
711        .route_layer(DefaultBodyLimit::max(SIGN_IN_BODY_BYTES))
712        .route_layer(from_fn_with_state(state.clone(), guard));
713    // Only a signed-in owner: never the token, never a device.
714    let adding = Router::new()
715        .route(
716            "/admin/passkeys/register",
717            post(handle_add_start).fallback(not_found),
718        )
719        .route(
720            "/admin/passkeys/register/finish",
721            post(handle_add_finish).fallback(not_found),
722        )
723        .route_layer(from_fn(json_only))
724        .route_layer(DefaultBodyLimit::max(SIGN_IN_BODY_BYTES))
725        .route_layer(from_fn_with_state(state.clone(), admin::owner_only));
726    let owner = Router::new()
727        .route(
728            "/admin/passkeys",
729            get(handle_list_passkeys).fallback(not_found),
730        )
731        .route(
732            "/admin/passkeys/{id}/remove",
733            post(handle_remove).fallback(not_found),
734        )
735        .route("/admin/logout", post(handle_sign_out).fallback(not_found))
736        .route(
737            "/admin/logout/others",
738            post(handle_sign_out_others).fallback(not_found),
739        )
740        .route_layer(DefaultBodyLimit::max(SIGN_IN_BODY_BYTES))
741        .route_layer(from_fn_with_state(state.clone(), admin::owner_only));
742    sign_in.merge(bootstrap).merge(adding).merge(owner)
743}
744
745#[cfg(not(feature = "passkeys"))]
746fn sign_in_routes(_state: &Arc<AppState>) -> Router<Arc<AppState>> {
747    Router::new()
748}
749
750fn sweep_devices(state: &AppState) -> Result<(usize, usize)> {
751    let now = time::OffsetDateTime::now_utc();
752    // The admin page's leftovers go on the same round: ceremonies finished
753    // that would have expired by now, and sessions that have ended.
754    state.passkeys.prune(state.now());
755    let clock = state.clock();
756    if let Err(e) = state.store.sweep_admin_sessions(
757        &format_timestamp(clock),
758        &format_timestamp(clock - admin::SESSION_IDLE),
759    ) {
760        eprintln!("admin session sweep failed: {e:#}");
761    }
762    state.store.sweep_devices_audited(
763        &format_timestamp(now - state.cfg.ephemeral_device_ttl),
764        &format_timestamp(now - devices::EXPIRED_ENROLLMENT_KEPT),
765    )
766}
767
768/// Counts the sweeps in a row whose checkpoint could not copy the whole WAL
769/// back into `recall.db`, and says so once that has lasted long enough to
770/// be a reader holding a transaction open rather than a page being read.
771///
772/// Only a log line, not a `/health` field: `Health` is a frozen shape
773/// released clients read, and this is something for the owner looking at
774/// the server's own output, which is where the other background failures
775/// are reported too.
776#[derive(Debug, Default)]
777struct WalWatch {
778    behind: u32,
779}
780
781impl WalWatch {
782    /// How many sweeps in a row, of `SWEEP_EVERY` each, before it is said.
783    const SWEEPS: u32 = 3;
784
785    /// Records one sweep's checkpoint; answers what to log, if anything:
786    /// every [`Self::SWEEPS`] sweeps while it lasts, and once when it ends.
787    fn record(&mut self, complete: bool) -> Option<String> {
788        if complete {
789            let was = std::mem::take(&mut self.behind);
790            return (was >= Self::SWEEPS)
791                .then(|| "the WAL is fully checkpointed into recall.db again".to_string());
792        }
793        self.behind += 1;
794        self.behind.is_multiple_of(Self::SWEEPS).then(|| {
795            format!(
796                "the WAL has not been fully checkpointed into recall.db for {} sweeps in a \
797                 row: a reader is keeping a transaction open (sqlite-web on a page, a \
798                 `sqlite3` shell), and recall.db-wal grows until it ends. Nothing is lost",
799                self.behind
800            )
801        })
802    }
803}
804
805fn run_backup(state: &AppState) {
806    if state.cfg.backup_dir.is_empty() {
807        return;
808    }
809    match state
810        .store
811        .backup(&state.cfg.backup_dir, state.cfg.backup_keep)
812    {
813        Ok(dest) => {
814            state.write().last_backup_at = now();
815            eprintln!("backup written: {}", dest.display());
816        }
817        Err(e) => eprintln!("backup failed: {e:#}"),
818    }
819}
820
821async fn shutdown_signal() {
822    let ctrl_c = async {
823        let _ = tokio::signal::ctrl_c().await;
824    };
825    #[cfg(unix)]
826    let terminate = async {
827        match tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) {
828            Ok(mut sig) => {
829                sig.recv().await;
830            }
831            Err(_) => std::future::pending::<()>().await,
832        }
833    };
834    #[cfg(not(unix))]
835    let terminate = std::future::pending::<()>();
836
837    tokio::select! {
838        _ = ctrl_c => {}
839        _ = terminate => {}
840    }
841}
842
843#[cfg(test)]
844mod tests {
845    use super::{Server, WalWatch};
846    use crate::{Config, Store};
847    use std::sync::Arc;
848    use std::time::{Duration, Instant};
849
850    /// The background sweep checkpoints the WAL, so `recall.db` on its own,
851    /// which is all sqlite-web's single-file mount has, catches up within a
852    /// sweep rather than whenever SQLite's own threshold of 1000 pages comes
853    /// round. The file alone is copied out and opened with no WAL beside
854    /// it, the copy made holding the store's lock, which the checkpoint
855    /// takes too, so it is never of a file a checkpoint is half way through
856    /// writing.
857    #[tokio::test]
858    async fn the_sweep_checkpoints_the_wal_into_the_file() {
859        let dir = tempfile::tempdir().unwrap();
860        let db = dir.path().join("recall.db");
861        let store = Arc::new(Store::open(&db).unwrap());
862        let server = Server::new(
863            Config {
864                token: "sweep-token".into(),
865                merge_enabled: false,
866                ..Config::default()
867            },
868            store.clone(),
869        );
870        // The schema into the file first, so the file alone opens, and what
871        // it lacks below is only the rows.
872        assert!(store.checkpoint_all().unwrap());
873        for i in 0..5 {
874            store
875                .upsert_audited(
876                    "acme/app",
877                    &format!("f{i}.md"),
878                    "x",
879                    "",
880                    crate::store::test_leaf,
881                )
882                .unwrap();
883        }
884        let alone = || -> i64 {
885            let copy = tempfile::tempdir().unwrap();
886            let file = copy.path().join("recall.db");
887            store
888                .with_raw(|_| {
889                    std::fs::copy(&db, &file).unwrap();
890                    Ok(())
891                })
892                .unwrap();
893            rusqlite::Connection::open(&file)
894                .unwrap()
895                .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
896                .unwrap()
897        };
898        assert_eq!(alone(), 0, "the rows are only in the WAL before the sweep");
899
900        let tasks = server.start_background();
901        let started = Instant::now();
902        while alone() != 5 {
903            assert!(
904                started.elapsed() < Duration::from_secs(10),
905                "the first sweep did not checkpoint: the file alone has {}",
906                alone()
907            );
908            tokio::time::sleep(Duration::from_millis(20)).await;
909        }
910        for t in tasks {
911            t.abort();
912        }
913    }
914
915    /// Said after three sweeps behind in a row and every three after, said
916    /// once more when it catches up, and not at all for a sweep or two.
917    #[test]
918    fn a_wal_held_back_for_several_sweeps_is_reported_and_so_is_its_end() {
919        let mut w = WalWatch::default();
920        assert_eq!(w.record(false), None);
921        assert_eq!(w.record(true), None, "one sweep behind is nothing");
922        assert_eq!(w.record(false), None);
923        assert_eq!(w.record(false), None);
924        let said = w.record(false).expect("three in a row");
925        assert!(said.contains("for 3 sweeps"), "{said}");
926        assert_eq!(w.record(false), None);
927        assert_eq!(w.record(false), None);
928        assert!(w.record(false).unwrap().contains("for 6 sweeps"));
929        assert!(w.record(true).unwrap().contains("again"));
930        assert_eq!(w.record(true), None);
931    }
932}