Skip to main content

fraiseql_server/server/
lifecycle.rs

1//! Server lifecycle: serve, `serve_with_shutdown`, and `shutdown_signal`.
2
3use std::net::SocketAddr;
4
5use axum::serve::ListenerExt;
6use tokio::net::TcpListener;
7use tracing::{error, info, warn};
8
9use super::{DatabaseAdapter, Result, Server, ServerError, TlsSetup};
10#[cfg(feature = "observers")]
11use crate::subscriptions::event_bridge::{EventBridge, EventBridgeConfig};
12
13impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<A> {
14    /// Start server and listen for requests.
15    ///
16    /// Uses SIGUSR1-aware shutdown signal when a schema path is configured,
17    /// enabling zero-downtime schema reloads via `kill -USR1 <pid>`.
18    ///
19    /// # Errors
20    ///
21    /// Returns error if server fails to bind or encounters runtime errors.
22    pub async fn serve(self) -> Result<()> {
23        self.serve_with_shutdown(Self::shutdown_signal()).await
24    }
25
26    /// Start server with a custom shutdown future.
27    ///
28    /// Enables programmatic shutdown (e.g., for `--watch` hot-reload) by accepting any
29    /// future that resolves when the server should stop.
30    ///
31    /// # Errors
32    ///
33    /// Returns error if server fails to bind or encounters runtime errors.
34    #[allow(clippy::cognitive_complexity)] // Reason: server lifecycle with TLS/non-TLS binding, signal handling, and graceful shutdown
35    pub async fn serve_with_shutdown<F>(mut self, shutdown: F) -> Result<()>
36    where
37        F: std::future::Future<Output = ()> + Send + 'static,
38    {
39        // Ensure RBAC schema exists before the router mounts RBAC endpoints.
40        // Must run here (async context) rather than inside build_router() (sync).
41        #[cfg(feature = "observers")]
42        if let Some(ref db_pool) = self.db_pool {
43            if self.config.admin_token.is_some() {
44                let rbac_backend =
45                    crate::api::rbac_management::db_backend::RbacDbBackend::new(db_pool.clone());
46                rbac_backend.ensure_schema().await.map_err(|e| {
47                    ServerError::ConfigError(format!("Failed to initialize RBAC schema: {e}"))
48                })?;
49            }
50        }
51
52        // Ensure the inbound-ingestion tables exist before the router mounts
53        // POST /webhooks/{provider}. Like RBAC, this must run here (async context)
54        // rather than in the sync build_router(). The spine holds normalized
55        // messages; the idempotency ledger backs the pipeline's atomic claim.
56        #[cfg(feature = "inbound")]
57        if let Some(ref db_pool) = self.db_pool {
58            if !self.config.webhooks.is_empty() {
59                crate::inbound::WebhookInboundState::init_spine(db_pool).await.map_err(|e| {
60                    ServerError::ConfigError(format!(
61                        "Failed to initialize inbound spine schema: {e}"
62                    ))
63                })?;
64                fraiseql_webhooks::PostgresIdempotencyStore::new(db_pool.clone())
65                    .init()
66                    .await
67                    .map_err(|e| {
68                        ServerError::ConfigError(format!(
69                            "Failed to initialize webhook idempotency schema: {e}"
70                        ))
71                    })?;
72                info!(
73                    routes = self.config.webhooks.len(),
74                    "Inbound ingestion schema ready (spine + idempotency ledger)"
75                );
76            }
77        }
78
79        // Initialize usage persistence backend if configured.
80        // Must run before build_router() so the aggregator is populated before
81        // serving requests, but after the DB pool is available (async context).
82        if let Some(ref usage_cfg) = self.config.usage.clone() {
83            use std::time::Duration;
84
85            use sqlx::postgres::PgPoolOptions;
86            use tokio::time::MissedTickBehavior;
87
88            use crate::usage::aggregator::{PostgresBackend, global_aggregator};
89
90            match PgPoolOptions::new()
91                .max_connections(2) // small dedicated pool — only used for periodic flushes
92                .connect(&self.config.database_url)
93                .await
94            {
95                Ok(pool) => {
96                    match PostgresBackend::new(pool).await {
97                        Ok(backend) => {
98                            let backend = std::sync::Arc::new(backend);
99                            // Upgrade global aggregator's backend from NoopBackend.
100                            global_aggregator().set_backend(backend.clone());
101                            // Restore persisted counters before serving requests.
102                            if let Err(e) = global_aggregator().load_from_backend().await {
103                                warn!(error = %e, "Usage persistence: startup load failed — continuing with in-memory counters");
104                            } else {
105                                info!("Usage persistence: loaded counters from PostgreSQL");
106                            }
107                            // Spawn background flush task on the server's JoinSet
108                            // so graceful shutdown can await its termination.
109                            let flush_interval = Duration::from_secs(usage_cfg.flush_interval_secs);
110                            let agg = std::sync::Arc::clone(global_aggregator());
111                            self.tasks.spawn(async move {
112                                let mut ticker = tokio::time::interval(flush_interval);
113                                ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
114                                ticker.tick().await; // skip immediate first tick
115                                loop {
116                                    ticker.tick().await;
117                                    if let Err(e) = agg.flush_to_backend().await {
118                                        warn!(error = %e, "Usage persistence: background flush failed");
119                                    }
120                                }
121                            });
122                            info!(
123                                flush_interval_secs = usage_cfg.flush_interval_secs,
124                                "Usage persistence: PostgreSQL backend active"
125                            );
126                        },
127                        Err(e) => {
128                            warn!(
129                                error = %e,
130                                "Usage persistence: PostgresBackend initialization failed — \
131                                 continuing with in-memory (NoopBackend)"
132                            );
133                        },
134                    }
135                },
136                Err(e) => {
137                    warn!(
138                        error = %e,
139                        "Usage persistence: failed to connect to PostgreSQL — \
140                         continuing with in-memory (NoopBackend)"
141                    );
142                },
143            }
144        }
145
146        // Prepare functions-runtime dispatch (load modules, register runtimes,
147        // attach the send_email wiring) before the router is built, so
148        // `build_app_state` mounts the before-mutation hooks. Async + fail-loud,
149        // like the RBAC/inbound schema init above; a no-op when no functions are
150        // declared or the feature is off.
151        #[cfg(feature = "functions-runtime")]
152        self.prepare_functions_runtime().await?;
153
154        let (app, app_state) = self.build_router();
155
156        // Start the poll-IMAP email workers.
157        // Each configured `[mailbox.<name>.imap]` half runs a background poll loop
158        // on the server's JoinSet, so graceful shutdown aborts them. The workers
159        // reuse the durable inbound spine and the `after:ingest` dispatch path;
160        // attachments stream into the legacy storage backend when one is
161        // configured. Must run here (async, after `build_router` supplies the
162        // function-dispatch hooks) rather than in the sync `build_router`.
163        #[cfg(feature = "inbound-email")]
164        if let Some(ref db_pool) = self.db_pool {
165            if self.config.mailbox.values().any(|mailbox| mailbox.imap.is_some()) {
166                use crate::inbound::{email, spine::PostgresInboundSpine};
167
168                // The email path shares the inbound spine; create it here too in
169                // case only `[mailbox.*.imap]` (no `[webhooks.*]`) is configured.
170                // Both DDLs are idempotent.
171                PostgresInboundSpine::new(db_pool.clone()).init().await.map_err(|e| {
172                    ServerError::ConfigError(format!(
173                        "Failed to initialize inbound spine schema: {e}"
174                    ))
175                })?;
176                email::init_cursor_store(db_pool).await.map_err(|e| {
177                    ServerError::ConfigError(format!(
178                        "Failed to initialize inbound email cursor schema: {e}"
179                    ))
180                })?;
181
182                // The correlator transitions send-status and suppression from
183                // inbound signals; its tables are also created by the send path,
184                // but a receive-only deployment needs them here too (idempotent).
185                let tracker = std::sync::Arc::new(email::PgSendTracker::new(db_pool.clone()));
186                tracker.init().await.map_err(|e| {
187                    ServerError::ConfigError(format!(
188                        "Failed to initialize send-tracking schema: {e}"
189                    ))
190                })?;
191                let correlator =
192                    std::sync::Arc::clone(&tracker) as std::sync::Arc<dyn email::SendCorrelator>;
193                let address_hash_key = self.build_address_hash_key();
194
195                let sink = self.storage_backend.as_ref().map(|backend| {
196                    std::sync::Arc::new(email::LegacyStorageSink::new(backend.clone()))
197                        as std::sync::Arc<
198                            dyn fraiseql_functions::host::live::storage::StorageBackend,
199                        >
200                });
201                // Opt-in Return-Path probe: verify the provider preserves
202                // plus-addressing before trusting VERP correlation. Blocks startup
203                // up to the probe window per eligible mailbox; off by default.
204                if self.config.send.verp_probe_on_start {
205                    email::run_startup_probes(&self.config.mailbox, |name| {
206                        std::env::var(name).ok()
207                    })
208                    .await;
209                }
210
211                let hooks = app_state.before_mutation_hooks.clone();
212                // #594: the after:ingest `fraiseql_query` bridge factory, built over the
213                // app's hot-reloadable executor (same one the route handlers use); only
214                // meaningful when function-dispatch hooks are present.
215                let query_executor_factory = hooks.as_ref().map(|_| {
216                    crate::routes::after_mutation::make_query_executor_factory(
217                        app_state.executor.clone(),
218                    )
219                });
220                let pollers = email::build_pollers(
221                    &self.config.mailbox,
222                    db_pool,
223                    hooks.as_ref(),
224                    query_executor_factory.as_ref(),
225                    sink.as_ref(),
226                    Some(&correlator),
227                    address_hash_key.as_ref(),
228                    self.config.send.challenge_suppress_after,
229                    |name| std::env::var(name).ok(),
230                );
231                let started = pollers.len();
232                for (poller, interval) in pollers {
233                    self.tasks.spawn(async move { poller.poll_forever(interval).await });
234                }
235                info!(mailboxes = started, "poll-IMAP email sources started");
236            }
237        }
238
239        // Start the scheduled-ingress source scheduler (#573): one poller per
240        // enabled Model B (Deno) source in the compiled schema, each firing its
241        // connector on its cron schedule under a single-firing lease with a host
242        // bound to the source's durable cursor and its `run_as` executor. Runs on
243        // the server's JoinSet so graceful shutdown drains it. Native pull sources
244        // (poll-IMAP email) run via their own pollers above. Must run here (async,
245        // after `build_router` supplies the function-dispatch hooks).
246        #[cfg(feature = "sources")]
247        if let Some(ref db_pool) = self.db_pool {
248            let sources = app_state.executor().schema().sources.clone();
249            if !sources.is_empty() {
250                let sources_config = self.config.sources.clone().unwrap_or_default();
251                if !crate::sources::sources_enabled(&sources_config) {
252                    info!(
253                        count = sources.len(),
254                        "source scheduler disabled by config — sources not started"
255                    );
256                } else if let Some(hooks) = app_state.before_mutation_hooks.as_ref() {
257                    // The shared durable source-cursor table (idempotent DDL).
258                    fraiseql_observers::PostgresSourceCursorStore::new(db_pool.clone())
259                        .init()
260                        .await
261                        .map_err(|e| {
262                            ServerError::ConfigError(format!(
263                                "Failed to initialize source cursor schema: {e}"
264                            ))
265                        })?;
266                    let host_config = crate::sources::source_host_config(&sources_config);
267                    let pollers = crate::sources::build_source_pollers(
268                        &sources,
269                        db_pool,
270                        &app_state.executor,
271                        hooks.as_ref(),
272                        &host_config,
273                        &fraiseql_functions::ResourceLimits::default(),
274                        sources_config.log_payloads,
275                    );
276                    let started = pollers.len();
277                    for poller in pollers {
278                        self.tasks.spawn(async move { poller.run_forever().await });
279                    }
280                    info!(sources = started, "source scheduler started");
281                } else {
282                    warn!(
283                        count = sources.len(),
284                        "compiled schema declares sources but the functions subsystem is not \
285                         configured — no source scheduler started"
286                    );
287                }
288            }
289        }
290
291        // Start the `cron:` function scheduler (#595): one leased poller per cron
292        // function in the compiled schema, each firing on its schedule under a
293        // single-firing advisory lease with the phase-02 `run_as` host. A cron
294        // function is a scheduled source without a cursor. Runs on the server's
295        // JoinSet so graceful shutdown drains it; requires a DB pool (the advisory
296        // lease + `_fraiseql_cron_state`) and the function-dispatch hooks.
297        #[cfg(feature = "functions-runtime")]
298        if let Some(ref db_pool) = self.db_pool {
299            if let Some(hooks) = app_state.before_mutation_hooks.as_ref() {
300                if !hooks.trigger_registry.cron_triggers.is_empty() {
301                    let cron_state = crate::cron::PgCronState::new(db_pool.clone());
302                    cron_state.init().await.map_err(|e| {
303                        ServerError::ConfigError(format!(
304                            "Failed to initialize cron state schema: {e}"
305                        ))
306                    })?;
307                    let host_config = crate::routes::after_mutation::host_context_config();
308                    let pollers = crate::cron::build_cron_pollers(
309                        db_pool,
310                        &app_state.executor,
311                        hooks.as_ref(),
312                        &host_config,
313                        &fraiseql_functions::ResourceLimits::default(),
314                    );
315                    let started = pollers.len();
316                    for poller in pollers {
317                        self.tasks.spawn(async move { poller.run_forever().await });
318                    }
319                    info!(cron_functions = started, "cron scheduler started");
320                }
321            }
322        }
323
324        // Spawn SIGUSR1 schema reload handler when running on Unix.
325        // The handler loops forever, reloading on each signal, until the
326        // server process exits — tracked on the server's JoinSet so graceful
327        // shutdown awaits its termination.
328        #[cfg(unix)]
329        if let Some(ref schema_path) = app_state.schema_path {
330            let reload_state = app_state.clone();
331            let reload_path = schema_path.clone();
332            self.tasks.spawn(async move {
333                let mut sigusr1 = match tokio::signal::unix::signal(
334                    tokio::signal::unix::SignalKind::user_defined1(),
335                ) {
336                    Ok(s) => s,
337                    Err(e) => {
338                        warn!(error = %e, "Failed to install SIGUSR1 handler — schema hot-reload disabled");
339                        return;
340                    },
341                };
342                loop {
343                    sigusr1.recv().await;
344                    info!(
345                        path = %reload_path.display(),
346                        "Received SIGUSR1 — reloading schema"
347                    );
348                    match reload_state.reload_schema(&reload_path).await {
349                        Ok(()) => {
350                            let hash = reload_state.executor().schema().content_hash();
351                            reload_state
352                                .metrics
353                                .schema_reloads_total
354                                .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
355                            info!(schema_hash = %hash, "Schema reloaded successfully via SIGUSR1");
356                        },
357                        Err(e) => {
358                            reload_state
359                                .metrics
360                                .schema_reload_errors_total
361                                .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
362                            error!(
363                                error = %e,
364                                path = %reload_path.display(),
365                                "Schema reload failed via SIGUSR1 — keeping previous schema"
366                            );
367                        },
368                    }
369                }
370            });
371            info!(
372                path = %schema_path.display(),
373                "SIGUSR1 schema reload handler installed"
374            );
375        }
376
377        // Initialize TLS setup (database connection TLS; server-side TLS is unsupported).
378        let tls_setup = TlsSetup::new(self.config.tls.clone(), self.config.database_tls.clone());
379
380        // Refuse to boot if server-side `[tls]` is enabled. FraiseQL does not terminate TLS
381        // itself — it serves plaintext and expects a reverse proxy / load balancer / service
382        // mesh to terminate TLS in front of it. Previously an enabled `[tls]` built a rustls
383        // config that was silently discarded while the server kept serving plaintext and
384        // logged `mtls_required = true` (M-tls-enforce); failing loud is honest.
385        if tls_setup.is_tls_enabled() {
386            return Err(ServerError::ConfigError(
387                "[tls] (server-side TLS termination) is enabled but not supported: FraiseQL \
388                 serves plaintext HTTP and expects TLS to be terminated by a reverse proxy, \
389                 load balancer, or service mesh. Remove the [tls] section (or set its \
390                 `enabled = false`) and terminate TLS in front of the server. Database \
391                 connection TLS ([database_tls]) is unaffected."
392                    .to_string(),
393            ));
394        }
395
396        info!(
397            bind_addr = %self.config.bind_addr,
398            graphql_path = %self.config.graphql_path,
399            tls_enabled = tls_setup.is_tls_enabled(),
400            "Starting FraiseQL server"
401        );
402
403        // Start observer runtime if configured, wiring CDC events to EventBridge
404        #[cfg(feature = "observers")]
405        #[allow(unused_variables)]
406        // Reason: _bridge_handle is kept alive to prevent task cancellation
407        let _bridge_handle = {
408            let mut handle: Option<tokio::task::JoinHandle<()>> = None;
409            if let Some(ref runtime) = self.observer_runtime {
410                info!("Starting observer runtime...");
411
412                // Create EventBridge to forward CDC events to GraphQL subscriptions
413                let bridge =
414                    EventBridge::new(self.subscription_manager.clone(), EventBridgeConfig::new());
415                let sender = bridge.sender();
416
417                let mut guard = runtime.write().await;
418                guard.set_event_bridge_sender(sender);
419
420                // #366: wire after:capture dispatch — externally-captured writes
421                // (from the change-log reader) drive `after:capture` functions on
422                // the phase-02 `run_as` host. The hook is a cheap no-op for
423                // FraiseQL's own (non-captured) rows, so mediated writes never loop.
424                #[cfg(feature = "functions-runtime")]
425                if let Some(hooks) = app_state.before_mutation_hooks.as_ref() {
426                    let hooks = std::sync::Arc::clone(hooks);
427                    let factory = crate::routes::after_mutation::make_query_executor_factory(
428                        app_state.executor.clone(),
429                    );
430                    guard.set_capture_dispatch(std::sync::Arc::new(move |event| {
431                        if let Some((fn_event, cdc)) =
432                            crate::routes::after_mutation::observer_event_to_capture(event)
433                        {
434                            let plans = crate::routes::after_mutation::plan_after_capture_dispatch(
435                                &hooks,
436                                &fn_event,
437                                cdc.as_deref(),
438                            );
439                            if !plans.is_empty() {
440                                crate::routes::after_mutation::spawn_after_capture(
441                                    &hooks,
442                                    plans,
443                                    Some(factory.clone()),
444                                );
445                            }
446                        }
447                    }));
448                }
449
450                match guard.start().await {
451                    Ok(()) => {
452                        info!("Observer runtime started");
453                        // Spawn EventBridge after observer runtime is running
454                        handle = Some(bridge.spawn());
455                        info!(
456                            "EventBridge started — CDC events will be forwarded to subscriptions"
457                        );
458                    },
459                    Err(e) => {
460                        // A broker-backed transport (NATS) was an explicit operator
461                        // choice; refusing to boot in production rather than silently
462                        // coming up without it is the #350 dead-broker contract. The
463                        // default PostgreSQL transport keeps the resilient
464                        // log-and-continue behaviour (and development downgrades the
465                        // NATS failure to the same warning).
466                        if guard.transport_requires_broker()
467                            && crate::ServerConfig::is_production_mode()
468                        {
469                            error!(
470                                error = %e,
471                                "Observer runtime failed to start on its configured \
472                                 transport; refusing to boot (set FRAISEQL_ENV=development \
473                                 to downgrade to a warning)"
474                            );
475                            return Err(e);
476                        }
477                        error!("Failed to start observer runtime: {}", e);
478                        warn!("Server will continue without observers");
479                    },
480                }
481                drop(guard);
482            }
483            handle
484        };
485
486        // Explicitly enable TCP_NODELAY (disable Nagle's algorithm) on every
487        // accepted connection to minimise latency for small GraphQL responses.
488        let listener = TcpListener::bind(self.config.bind_addr)
489            .await
490            .map_err(|e| ServerError::BindError(e.to_string()))?
491            .tap_io(|tcp_stream| {
492                if let Err(err) = tcp_stream.set_nodelay(true) {
493                    warn!("failed to set TCP_NODELAY: {err:#}");
494                }
495            });
496
497        // Warn if the process file descriptor limit is below the recommended minimum.
498        // A low limit causes "too many open files" errors under load.
499        #[cfg(target_os = "linux")]
500        {
501            if let Ok(limits) = std::fs::read_to_string("/proc/self/limits") {
502                for line in limits.lines() {
503                    if line.starts_with("Max open files") {
504                        let parts: Vec<&str> = line.split_whitespace().collect();
505                        if let Some(soft) = parts.get(3) {
506                            if let Ok(n) = soft.parse::<u64>() {
507                                if n < 65_536 {
508                                    warn!(
509                                        current_fd_limit = n,
510                                        recommended = 65_536,
511                                        "File descriptor limit is low; consider raising ulimit -n"
512                                    );
513                                }
514                            }
515                        }
516                        break;
517                    }
518                }
519            }
520        }
521
522        // Log database TLS configuration
523        info!(
524            postgres_ssl_mode = tls_setup.postgres_ssl_mode(),
525            redis_ssl = tls_setup.redis_ssl_enabled(),
526            clickhouse_https = tls_setup.clickhouse_https_enabled(),
527            elasticsearch_https = tls_setup.elasticsearch_https_enabled(),
528            "Database connection TLS configuration applied"
529        );
530
531        info!("Server listening on http://{}", self.config.bind_addr);
532
533        // Start both HTTP and gRPC servers concurrently if Arrow Flight is enabled
534        #[cfg(feature = "arrow")]
535        if let Some(flight_service) = self.flight_service.take() {
536            let flight_addr = self.config.flight_bind_addr;
537            info!("Arrow Flight server listening on grpc://{}", flight_addr);
538
539            // Spawn Flight server in background, registered on the server's
540            // JoinSet. The set's `shutdown` step abort-then-awaits the gRPC
541            // server when the HTTP server exits.
542            self.tasks.spawn(async move {
543                if let Err(e) = tonic::transport::Server::builder()
544                    .add_service(flight_service.into_server())
545                    .serve(flight_addr)
546                    .await
547                {
548                    error!(error = %e, "Arrow Flight server terminated with error");
549                }
550            });
551
552            // Wrap the user-supplied shutdown future so we can also stop observer runtime
553            #[cfg(feature = "observers")]
554            let observer_runtime = self.observer_runtime.clone();
555
556            let shutdown_with_cleanup = async move {
557                shutdown.await;
558                #[cfg(feature = "observers")]
559                if let Some(ref runtime) = observer_runtime {
560                    info!("Shutting down observer runtime");
561                    let mut guard = runtime.write().await;
562                    if let Err(e) = guard.stop().await {
563                        #[cfg(feature = "observers")]
564                        error!("Error stopping runtime: {}", e);
565                    } else {
566                        info!("Runtime stopped cleanly");
567                    }
568                }
569            };
570
571            // Run HTTP server with graceful shutdown
572            axum::serve(listener, app.into_make_service_with_connect_info::<SocketAddr>())
573                .with_graceful_shutdown(shutdown_with_cleanup)
574                .await
575                .map_err(|e| ServerError::IoError(std::io::Error::other(e)))?;
576
577            // Abort and await every lifecycle task (Flight server, SIGUSR1
578            // handler, PKCE cleanup, trusted-docs reload, usage flush, …).
579            drain_lifecycle_tasks(self.tasks, self.config.shutdown_timeout_secs).await;
580        }
581
582        // HTTP-only server (when arrow feature not enabled)
583        #[cfg(not(feature = "arrow"))]
584        {
585            axum::serve(listener, app.into_make_service_with_connect_info::<SocketAddr>())
586                .with_graceful_shutdown(shutdown)
587                .await
588                .map_err(|e| ServerError::IoError(std::io::Error::other(e)))?;
589
590            let shutdown_timeout =
591                std::time::Duration::from_secs(self.config.shutdown_timeout_secs);
592            info!(
593                timeout_secs = self.config.shutdown_timeout_secs,
594                "HTTP server stopped, draining remaining work"
595            );
596
597            let drain = tokio::time::timeout(shutdown_timeout, async {
598                #[cfg(feature = "observers")]
599                if let Some(ref runtime) = self.observer_runtime {
600                    let mut guard = runtime.write().await;
601                    match guard.stop().await {
602                        Ok(()) => info!("Observer runtime stopped cleanly"),
603                        Err(e) => warn!("Observer runtime shutdown error: {e}"),
604                    }
605                }
606            })
607            .await;
608
609            if drain.is_err() {
610                warn!(
611                    timeout_secs = self.config.shutdown_timeout_secs,
612                    "Shutdown drain timed out; forcing exit"
613                );
614            } else {
615                info!("Graceful shutdown complete");
616            }
617
618            // Abort and await every lifecycle task (SIGUSR1 handler, PKCE
619            // cleanup, trusted-docs reload, usage flush, …).
620            drain_lifecycle_tasks(self.tasks, self.config.shutdown_timeout_secs).await;
621        }
622
623        Ok(())
624    }
625
626    /// Start server on an externally created listener.
627    ///
628    /// Used in tests to discover the bound port before serving.
629    /// Skips TLS, Flight, and observer startup — suitable for unit/integration tests only.
630    ///
631    /// # Errors
632    ///
633    /// Returns error if the server encounters a runtime error.
634    pub async fn serve_on_listener<F>(self, listener: TcpListener, shutdown: F) -> Result<()>
635    where
636        F: std::future::Future<Output = ()> + Send + 'static,
637    {
638        let (app, _app_state) = self.build_router();
639        axum::serve(listener, app.into_make_service_with_connect_info::<SocketAddr>())
640            .with_graceful_shutdown(shutdown)
641            .await
642            .map_err(|e| ServerError::IoError(std::io::Error::other(e)))?;
643        // Abort and await any lifecycle tasks spawned during construction
644        // (e.g. PKCE cleanup, trusted-docs reload) so the test path doesn't
645        // leak background work into the next test.
646        drain_lifecycle_tasks(self.tasks, self.config.shutdown_timeout_secs).await;
647        Ok(())
648    }
649
650    /// Listen for shutdown signals (Ctrl+C or SIGTERM)
651    pub async fn shutdown_signal() {
652        use tokio::signal;
653
654        let ctrl_c = async {
655            match signal::ctrl_c().await {
656                Ok(()) => {},
657                Err(e) => {
658                    warn!(error = %e, "Failed to install Ctrl+C handler");
659                    std::future::pending::<()>().await;
660                },
661            }
662        };
663
664        #[cfg(unix)]
665        let terminate = async {
666            match signal::unix::signal(signal::unix::SignalKind::terminate()) {
667                Ok(mut s) => {
668                    s.recv().await;
669                },
670                Err(e) => {
671                    warn!(error = %e, "Failed to install SIGTERM handler");
672                    std::future::pending::<()>().await;
673                },
674            }
675        };
676
677        #[cfg(not(unix))]
678        let terminate = std::future::pending::<()>();
679
680        tokio::select! {
681            () = ctrl_c => info!("Received Ctrl+C"),
682            () = terminate => info!("Received SIGTERM"),
683        }
684    }
685}
686
687/// Abort every lifecycle task on the supplied [`tokio::task::JoinSet`] and await
688/// the resulting `JoinError`s so the runtime is fully drained before
689/// `serve_with_shutdown` returns.
690///
691/// Tasks are awaited under an outer timeout so a stuck task cannot prevent
692/// process exit. A `JoinError::is_cancelled()` after `JoinSet::abort_all` is
693/// the expected case — only unexpected panics are logged.
694pub(super) async fn drain_lifecycle_tasks(
695    mut tasks: tokio::task::JoinSet<()>,
696    shutdown_timeout_secs: u64,
697) {
698    if tasks.is_empty() {
699        return;
700    }
701
702    tasks.abort_all();
703    let timeout = std::time::Duration::from_secs(shutdown_timeout_secs);
704    let drained = tokio::time::timeout(timeout, async {
705        while let Some(res) = tasks.join_next().await {
706            if let Err(e) = res {
707                if !e.is_cancelled() {
708                    warn!(error = %e, "Lifecycle task terminated with a non-cancellation error");
709                }
710            }
711        }
712    })
713    .await;
714    if drained.is_err() {
715        warn!(
716            timeout_secs = shutdown_timeout_secs,
717            "Lifecycle task drain timed out; some background tasks did not stop in time"
718        );
719    } else {
720        info!("All lifecycle background tasks drained");
721    }
722}