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 let pollers = email::build_pollers(
213 &self.config.mailbox,
214 db_pool,
215 hooks.as_ref(),
216 sink.as_ref(),
217 Some(&correlator),
218 address_hash_key.as_ref(),
219 self.config.send.challenge_suppress_after,
220 |name| std::env::var(name).ok(),
221 );
222 let started = pollers.len();
223 for (poller, interval) in pollers {
224 self.tasks.spawn(async move { poller.poll_forever(interval).await });
225 }
226 info!(mailboxes = started, "poll-IMAP email sources started");
227 }
228 }
229
230 // Start the scheduled-ingress source scheduler (#573): one poller per
231 // enabled Model B (Deno) source in the compiled schema, each firing its
232 // connector on its cron schedule under a single-firing lease with a host
233 // bound to the source's durable cursor and its `run_as` executor. Runs on
234 // the server's JoinSet so graceful shutdown drains it. Native pull sources
235 // (poll-IMAP email) run via their own pollers above. Must run here (async,
236 // after `build_router` supplies the function-dispatch hooks).
237 #[cfg(feature = "sources")]
238 if let Some(ref db_pool) = self.db_pool {
239 let sources = app_state.executor().schema().sources.clone();
240 if !sources.is_empty() {
241 let sources_config = self.config.sources.clone().unwrap_or_default();
242 if !crate::sources::sources_enabled(&sources_config) {
243 info!(
244 count = sources.len(),
245 "source scheduler disabled by config — sources not started"
246 );
247 } else if let Some(hooks) = app_state.before_mutation_hooks.as_ref() {
248 // The shared durable source-cursor table (idempotent DDL).
249 fraiseql_observers::PostgresSourceCursorStore::new(db_pool.clone())
250 .init()
251 .await
252 .map_err(|e| {
253 ServerError::ConfigError(format!(
254 "Failed to initialize source cursor schema: {e}"
255 ))
256 })?;
257 let host_config = crate::sources::source_host_config(&sources_config);
258 let pollers = crate::sources::build_source_pollers(
259 &sources,
260 db_pool,
261 &app_state.executor,
262 hooks.as_ref(),
263 &host_config,
264 &fraiseql_functions::ResourceLimits::default(),
265 sources_config.log_payloads,
266 );
267 let started = pollers.len();
268 for poller in pollers {
269 self.tasks.spawn(async move { poller.run_forever().await });
270 }
271 info!(sources = started, "source scheduler started");
272 } else {
273 warn!(
274 count = sources.len(),
275 "compiled schema declares sources but the functions subsystem is not \
276 configured — no source scheduler started"
277 );
278 }
279 }
280 }
281
282 // Spawn SIGUSR1 schema reload handler when running on Unix.
283 // The handler loops forever, reloading on each signal, until the
284 // server process exits — tracked on the server's JoinSet so graceful
285 // shutdown awaits its termination.
286 #[cfg(unix)]
287 if let Some(ref schema_path) = app_state.schema_path {
288 let reload_state = app_state.clone();
289 let reload_path = schema_path.clone();
290 self.tasks.spawn(async move {
291 let mut sigusr1 = match tokio::signal::unix::signal(
292 tokio::signal::unix::SignalKind::user_defined1(),
293 ) {
294 Ok(s) => s,
295 Err(e) => {
296 warn!(error = %e, "Failed to install SIGUSR1 handler — schema hot-reload disabled");
297 return;
298 },
299 };
300 loop {
301 sigusr1.recv().await;
302 info!(
303 path = %reload_path.display(),
304 "Received SIGUSR1 — reloading schema"
305 );
306 match reload_state.reload_schema(&reload_path).await {
307 Ok(()) => {
308 let hash = reload_state.executor().schema().content_hash();
309 reload_state
310 .metrics
311 .schema_reloads_total
312 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
313 info!(schema_hash = %hash, "Schema reloaded successfully via SIGUSR1");
314 },
315 Err(e) => {
316 reload_state
317 .metrics
318 .schema_reload_errors_total
319 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
320 error!(
321 error = %e,
322 path = %reload_path.display(),
323 "Schema reload failed via SIGUSR1 — keeping previous schema"
324 );
325 },
326 }
327 }
328 });
329 info!(
330 path = %schema_path.display(),
331 "SIGUSR1 schema reload handler installed"
332 );
333 }
334
335 // Initialize TLS setup (database connection TLS; server-side TLS is unsupported).
336 let tls_setup = TlsSetup::new(self.config.tls.clone(), self.config.database_tls.clone());
337
338 // Refuse to boot if server-side `[tls]` is enabled. FraiseQL does not terminate TLS
339 // itself — it serves plaintext and expects a reverse proxy / load balancer / service
340 // mesh to terminate TLS in front of it. Previously an enabled `[tls]` built a rustls
341 // config that was silently discarded while the server kept serving plaintext and
342 // logged `mtls_required = true` (M-tls-enforce); failing loud is honest.
343 if tls_setup.is_tls_enabled() {
344 return Err(ServerError::ConfigError(
345 "[tls] (server-side TLS termination) is enabled but not supported: FraiseQL \
346 serves plaintext HTTP and expects TLS to be terminated by a reverse proxy, \
347 load balancer, or service mesh. Remove the [tls] section (or set its \
348 `enabled = false`) and terminate TLS in front of the server. Database \
349 connection TLS ([database_tls]) is unaffected."
350 .to_string(),
351 ));
352 }
353
354 info!(
355 bind_addr = %self.config.bind_addr,
356 graphql_path = %self.config.graphql_path,
357 tls_enabled = tls_setup.is_tls_enabled(),
358 "Starting FraiseQL server"
359 );
360
361 // Start observer runtime if configured, wiring CDC events to EventBridge
362 #[cfg(feature = "observers")]
363 #[allow(unused_variables)]
364 // Reason: _bridge_handle is kept alive to prevent task cancellation
365 let _bridge_handle = {
366 let mut handle: Option<tokio::task::JoinHandle<()>> = None;
367 if let Some(ref runtime) = self.observer_runtime {
368 info!("Starting observer runtime...");
369
370 // Create EventBridge to forward CDC events to GraphQL subscriptions
371 let bridge =
372 EventBridge::new(self.subscription_manager.clone(), EventBridgeConfig::new());
373 let sender = bridge.sender();
374
375 let mut guard = runtime.write().await;
376 guard.set_event_bridge_sender(sender);
377
378 match guard.start().await {
379 Ok(()) => {
380 info!("Observer runtime started");
381 // Spawn EventBridge after observer runtime is running
382 handle = Some(bridge.spawn());
383 info!(
384 "EventBridge started — CDC events will be forwarded to subscriptions"
385 );
386 },
387 Err(e) => {
388 // A broker-backed transport (NATS) was an explicit operator
389 // choice; refusing to boot in production rather than silently
390 // coming up without it is the #350 dead-broker contract. The
391 // default PostgreSQL transport keeps the resilient
392 // log-and-continue behaviour (and development downgrades the
393 // NATS failure to the same warning).
394 if guard.transport_requires_broker()
395 && crate::ServerConfig::is_production_mode()
396 {
397 error!(
398 error = %e,
399 "Observer runtime failed to start on its configured \
400 transport; refusing to boot (set FRAISEQL_ENV=development \
401 to downgrade to a warning)"
402 );
403 return Err(e);
404 }
405 error!("Failed to start observer runtime: {}", e);
406 warn!("Server will continue without observers");
407 },
408 }
409 drop(guard);
410 }
411 handle
412 };
413
414 // Explicitly enable TCP_NODELAY (disable Nagle's algorithm) on every
415 // accepted connection to minimise latency for small GraphQL responses.
416 let listener = TcpListener::bind(self.config.bind_addr)
417 .await
418 .map_err(|e| ServerError::BindError(e.to_string()))?
419 .tap_io(|tcp_stream| {
420 if let Err(err) = tcp_stream.set_nodelay(true) {
421 warn!("failed to set TCP_NODELAY: {err:#}");
422 }
423 });
424
425 // Warn if the process file descriptor limit is below the recommended minimum.
426 // A low limit causes "too many open files" errors under load.
427 #[cfg(target_os = "linux")]
428 {
429 if let Ok(limits) = std::fs::read_to_string("/proc/self/limits") {
430 for line in limits.lines() {
431 if line.starts_with("Max open files") {
432 let parts: Vec<&str> = line.split_whitespace().collect();
433 if let Some(soft) = parts.get(3) {
434 if let Ok(n) = soft.parse::<u64>() {
435 if n < 65_536 {
436 warn!(
437 current_fd_limit = n,
438 recommended = 65_536,
439 "File descriptor limit is low; consider raising ulimit -n"
440 );
441 }
442 }
443 }
444 break;
445 }
446 }
447 }
448 }
449
450 // Log database TLS configuration
451 info!(
452 postgres_ssl_mode = tls_setup.postgres_ssl_mode(),
453 redis_ssl = tls_setup.redis_ssl_enabled(),
454 clickhouse_https = tls_setup.clickhouse_https_enabled(),
455 elasticsearch_https = tls_setup.elasticsearch_https_enabled(),
456 "Database connection TLS configuration applied"
457 );
458
459 info!("Server listening on http://{}", self.config.bind_addr);
460
461 // Start both HTTP and gRPC servers concurrently if Arrow Flight is enabled
462 #[cfg(feature = "arrow")]
463 if let Some(flight_service) = self.flight_service.take() {
464 let flight_addr = self.config.flight_bind_addr;
465 info!("Arrow Flight server listening on grpc://{}", flight_addr);
466
467 // Spawn Flight server in background, registered on the server's
468 // JoinSet. The set's `shutdown` step abort-then-awaits the gRPC
469 // server when the HTTP server exits.
470 self.tasks.spawn(async move {
471 if let Err(e) = tonic::transport::Server::builder()
472 .add_service(flight_service.into_server())
473 .serve(flight_addr)
474 .await
475 {
476 error!(error = %e, "Arrow Flight server terminated with error");
477 }
478 });
479
480 // Wrap the user-supplied shutdown future so we can also stop observer runtime
481 #[cfg(feature = "observers")]
482 let observer_runtime = self.observer_runtime.clone();
483
484 let shutdown_with_cleanup = async move {
485 shutdown.await;
486 #[cfg(feature = "observers")]
487 if let Some(ref runtime) = observer_runtime {
488 info!("Shutting down observer runtime");
489 let mut guard = runtime.write().await;
490 if let Err(e) = guard.stop().await {
491 #[cfg(feature = "observers")]
492 error!("Error stopping runtime: {}", e);
493 } else {
494 info!("Runtime stopped cleanly");
495 }
496 }
497 };
498
499 // Run HTTP server with graceful shutdown
500 axum::serve(listener, app.into_make_service_with_connect_info::<SocketAddr>())
501 .with_graceful_shutdown(shutdown_with_cleanup)
502 .await
503 .map_err(|e| ServerError::IoError(std::io::Error::other(e)))?;
504
505 // Abort and await every lifecycle task (Flight server, SIGUSR1
506 // handler, PKCE cleanup, trusted-docs reload, usage flush, …).
507 drain_lifecycle_tasks(self.tasks, self.config.shutdown_timeout_secs).await;
508 }
509
510 // HTTP-only server (when arrow feature not enabled)
511 #[cfg(not(feature = "arrow"))]
512 {
513 axum::serve(listener, app.into_make_service_with_connect_info::<SocketAddr>())
514 .with_graceful_shutdown(shutdown)
515 .await
516 .map_err(|e| ServerError::IoError(std::io::Error::other(e)))?;
517
518 let shutdown_timeout =
519 std::time::Duration::from_secs(self.config.shutdown_timeout_secs);
520 info!(
521 timeout_secs = self.config.shutdown_timeout_secs,
522 "HTTP server stopped, draining remaining work"
523 );
524
525 let drain = tokio::time::timeout(shutdown_timeout, async {
526 #[cfg(feature = "observers")]
527 if let Some(ref runtime) = self.observer_runtime {
528 let mut guard = runtime.write().await;
529 match guard.stop().await {
530 Ok(()) => info!("Observer runtime stopped cleanly"),
531 Err(e) => warn!("Observer runtime shutdown error: {e}"),
532 }
533 }
534 })
535 .await;
536
537 if drain.is_err() {
538 warn!(
539 timeout_secs = self.config.shutdown_timeout_secs,
540 "Shutdown drain timed out; forcing exit"
541 );
542 } else {
543 info!("Graceful shutdown complete");
544 }
545
546 // Abort and await every lifecycle task (SIGUSR1 handler, PKCE
547 // cleanup, trusted-docs reload, usage flush, …).
548 drain_lifecycle_tasks(self.tasks, self.config.shutdown_timeout_secs).await;
549 }
550
551 Ok(())
552 }
553
554 /// Start server on an externally created listener.
555 ///
556 /// Used in tests to discover the bound port before serving.
557 /// Skips TLS, Flight, and observer startup — suitable for unit/integration tests only.
558 ///
559 /// # Errors
560 ///
561 /// Returns error if the server encounters a runtime error.
562 pub async fn serve_on_listener<F>(self, listener: TcpListener, shutdown: F) -> Result<()>
563 where
564 F: std::future::Future<Output = ()> + Send + 'static,
565 {
566 let (app, _app_state) = self.build_router();
567 axum::serve(listener, app.into_make_service_with_connect_info::<SocketAddr>())
568 .with_graceful_shutdown(shutdown)
569 .await
570 .map_err(|e| ServerError::IoError(std::io::Error::other(e)))?;
571 // Abort and await any lifecycle tasks spawned during construction
572 // (e.g. PKCE cleanup, trusted-docs reload) so the test path doesn't
573 // leak background work into the next test.
574 drain_lifecycle_tasks(self.tasks, self.config.shutdown_timeout_secs).await;
575 Ok(())
576 }
577
578 /// Listen for shutdown signals (Ctrl+C or SIGTERM)
579 pub async fn shutdown_signal() {
580 use tokio::signal;
581
582 let ctrl_c = async {
583 match signal::ctrl_c().await {
584 Ok(()) => {},
585 Err(e) => {
586 warn!(error = %e, "Failed to install Ctrl+C handler");
587 std::future::pending::<()>().await;
588 },
589 }
590 };
591
592 #[cfg(unix)]
593 let terminate = async {
594 match signal::unix::signal(signal::unix::SignalKind::terminate()) {
595 Ok(mut s) => {
596 s.recv().await;
597 },
598 Err(e) => {
599 warn!(error = %e, "Failed to install SIGTERM handler");
600 std::future::pending::<()>().await;
601 },
602 }
603 };
604
605 #[cfg(not(unix))]
606 let terminate = std::future::pending::<()>();
607
608 tokio::select! {
609 () = ctrl_c => info!("Received Ctrl+C"),
610 () = terminate => info!("Received SIGTERM"),
611 }
612 }
613}
614
615/// Abort every lifecycle task on the supplied [`tokio::task::JoinSet`] and await
616/// the resulting `JoinError`s so the runtime is fully drained before
617/// `serve_with_shutdown` returns.
618///
619/// Tasks are awaited under an outer timeout so a stuck task cannot prevent
620/// process exit. A `JoinError::is_cancelled()` after `JoinSet::abort_all` is
621/// the expected case — only unexpected panics are logged.
622pub(super) async fn drain_lifecycle_tasks(
623 mut tasks: tokio::task::JoinSet<()>,
624 shutdown_timeout_secs: u64,
625) {
626 if tasks.is_empty() {
627 return;
628 }
629
630 tasks.abort_all();
631 let timeout = std::time::Duration::from_secs(shutdown_timeout_secs);
632 let drained = tokio::time::timeout(timeout, async {
633 while let Some(res) = tasks.join_next().await {
634 if let Err(e) = res {
635 if !e.is_cancelled() {
636 warn!(error = %e, "Lifecycle task terminated with a non-cancellation error");
637 }
638 }
639 }
640 })
641 .await;
642 if drained.is_err() {
643 warn!(
644 timeout_secs = shutdown_timeout_secs,
645 "Lifecycle task drain timed out; some background tasks did not stop in time"
646 );
647 } else {
648 info!("All lifecycle background tasks drained");
649 }
650}