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}