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