Skip to main content

acme_proxy_server/
lib.rs

1//! The server runtime: how a configuration becomes routers, and how those are
2//! served and rebuilt.
3//!
4//! Startup is split on the socket boundary, which is what lets a test drive the
5//! whole path on an ephemeral port with its own shutdown future instead of a
6//! process signal:
7//!
8//! - [`run`] binds `server.bind_address`, installs the `SIGHUP` handler, and
9//!   hands the socket on. It is the whole of `acme-proxy serve`.
10//! - [`serve_on`] validates the admin configuration and binds that socket too,
11//!   when `[admin]` is enabled.
12//! - [`serve_on_with`] does everything else — profile resolution, deduplicated
13//!   signer backends, per-profile filters and validators, TLS, the job registry
14//!   (every signer's handlers, notification delivery and the six table sweeps,
15//!   built by `generation::job_registry_for`), the runner draining it, and
16//!   `axum::serve` with connect info attached.
17//!
18//! That assembly is [`generation::build_generation`], and it is called again on
19//! every reload rather than only at startup — so the two cannot drift, and a
20//! subsystem added to one is added to the other by construction. What a reload
21//! may change, and what it refuses by name, is [`crate::reload`]'s to say;
22//! [`serve_on_with_reloads`] is where the two meet.
23//!
24//! - [`profile`] — how a generation builds each ACME endpoint.
25//! - [`assembly`] — what a generation hands its profiles, and what outlives it.
26//! - [`generation`] — one generation built, then published: the reload policy.
27//! - [`supervisor`] — the task that serialises reloads.
28//! - [`sockets`] — the three listeners' binds, plans and announcements.
29//! - [`roles`] — which of `acme`, `admin` and `worker` this process runs.
30//! - [`logging`] — the subscriber `serve` installs, and its reloadable filter.
31//!
32//! What it serves sits below it: the endpoint itself is
33//! [`acme_proxy_protocol::profile::Profile`], and the routers each listener serves, with
34//! their shared layers, are [`acme_proxy_protocol::router`].
35//!
36//! [`reload`] beside it is what a reload may change and what it refuses by
37//! name. An internal crate of the `acme-proxy` binary, published in lockstep
38//! with it and with no semver promise of its own.
39use std::future::Future;
40use std::net::SocketAddr;
41use std::sync::Arc;
42
43use tokio::net::TcpListener;
44use tracing::{error, info, warn};
45
46use acme_proxy_core::config::Config;
47use acme_proxy_store::db::Database;
48
49pub mod assembly;
50pub mod generation;
51pub mod logging;
52pub mod profile;
53pub mod reload;
54pub mod roles;
55pub mod sockets;
56pub mod supervisor;
57#[cfg(test)]
58mod tests;
59
60pub use assembly::{Assembly, GenerationParts};
61pub use roles::{ProcessRole, RoleSet};
62
63pub use sockets::check_metrics_config;
64
65use generation::{Generation, announce_profile, build_generation};
66use sockets::{
67    Sockets, announce_admin_listener, announce_metrics_listener, bind_admin, bind_metrics,
68    bound_address,
69};
70use supervisor::{Cells, supervise_reloads};
71
72/// Binds the configured sockets and runs the server until a shutdown signal
73/// arrives: the whole of `acme-proxy serve`.
74///
75/// Every failure is logged here — `server_socket_bind_failed` for the ACME
76/// socket, `server_fatal_error` for anything after it — and returned for the
77/// caller to turn into an exit status.
78pub async fn run(
79    roles: RoleSet,
80    config: Arc<Config>,
81    database: Arc<Database>,
82) -> anyhow::Result<()> {
83    // Only the roles this process runs get a socket. The ACME listener is the
84    // one that used to be unconditional — there is no `server.enabled`, a CA
85    // serving no ACME having been a process with nothing to do — and a role
86    // split is exactly the case where that stops being true.
87    let listener = match roles.has(ProcessRole::Acme) {
88        false => None,
89        true => Some(
90            TcpListener::bind(&config.server.bind_address)
91                .await
92                .map_err(|error| {
93                    error!(event = "server_socket_bind_failed",
94                           outcome = "failure",
95                           bind_address = %config.server.bind_address,
96                           error = %error);
97                    anyhow::anyhow!("cannot bind {}: {error}", config.server.bind_address)
98                })?,
99        ),
100    };
101
102    // Installed here, before anything slow: `SIGHUP`'s default disposition is
103    // *terminate*, so until the handler exists a reload signal kills the
104    // process. `serve_on_with_reloads` does profile assembly and the relay's
105    // first upstream contact before it binds anything, which is exactly the
106    // window an operator's `systemctl reload` could land in.
107    let (reload_handle, reloads) = crate::reload::channel();
108    let _hangups = AbortOnDrop(tokio::spawn(watch_for_hangup(reload_handle)));
109
110    let fatal = |error: &anyhow::Error| {
111        error!(event = "server_fatal_error", outcome = "failure", error = %error);
112    };
113    let admin_listener = match roles.has(ProcessRole::Admin) {
114        false => None,
115        true => bind_admin(&config).await.inspect_err(fatal)?,
116    };
117    let metrics_listener = bind_metrics(&config).await.inspect_err(fatal)?;
118
119    serve_on_with_reloads(
120        roles,
121        config,
122        database,
123        Sockets {
124            acme: listener,
125            admin: admin_listener,
126            metrics: metrics_listener,
127        },
128        shutdown_signal(),
129        reloads,
130    )
131    .await
132    .inspect_err(fatal)
133}
134
135/// Turns every `SIGHUP` into a reload request, for the life of the process.
136///
137/// Unlike the shutdown signal, this one does not consume its stream: an operator
138/// reloads repeatedly, and a handler that fired once would leave the second
139/// `SIGHUP` back at its default disposition — killing the server.
140#[cfg(unix)]
141async fn watch_for_hangup(handle: crate::reload::ReloadHandle) {
142    let mut hangups = match tokio::signal::unix::signal(tokio::signal::unix::SignalKind::hangup()) {
143        Ok(stream) => stream,
144        Err(error) => {
145            error!(event = "server_signal_handler_failed", outcome = "failure", signal = "SIGHUP", error = %error);
146            return;
147        }
148    };
149    while hangups.recv().await.is_some() {
150        // `false` means the supervisor is gone — it panicked, or the server is
151        // already shutting down — so this signal, and every later one, does
152        // nothing. Silence there is an operator whose `systemctl reload`
153        // reports success and changes nothing.
154        if !handle.trigger() {
155            error!(
156                event = "server_reload_supervisor_gone",
157                outcome = "failure",
158                signal = "SIGHUP",
159                "the reload supervisor is not accepting requests; restart the process to \
160                 apply a configuration change"
161            );
162        }
163    }
164}
165
166/// No `SIGHUP` off Unix, so there is nothing to watch for.
167#[cfg(not(unix))]
168async fn watch_for_hangup(_handle: crate::reload::ReloadHandle) {
169    std::future::pending::<()>().await;
170}
171
172/// Assembles and serves the application over an already-bound socket.
173///
174/// Split from [`run`] on the socket boundary: a caller supplying its own
175/// listener and its own `shutdown` future can drive the whole startup path —
176/// profile assembly, TLS, backend resume, the nonce reaper, both `axum::serve`
177/// arms — without owning a fixed port or a process signal.
178///
179/// Binds the web admin socket itself when `[admin]` is enabled. The signature
180/// is unchanged, and `admin.enabled` is false by default, so every existing
181/// caller is untouched; a test that wants to drive *both* listeners supplies
182/// its own pair through [`serve_on_with`].
183pub async fn serve_on(
184    config: Arc<Config>,
185    database: Arc<Database>,
186    listener: TcpListener,
187    shutdown: impl Future<Output = ()> + Send + 'static,
188) -> anyhow::Result<()> {
189    let admin_listener = bind_admin(&config).await?;
190    let metrics_listener = bind_metrics(&config).await?;
191    serve_on_with(
192        config,
193        database,
194        listener,
195        admin_listener,
196        metrics_listener,
197        shutdown,
198    )
199    .await
200}
201
202/// [`serve_on`] with all three sockets supplied.
203///
204/// The full version, split on the same boundary and for the same reason: a
205/// caller handing in three ephemeral ports can drive the whole startup path —
206/// including that one shutdown signal stops all of them — without owning a
207/// fixed port or a process signal.
208pub async fn serve_on_with(
209    config: Arc<Config>,
210    database: Arc<Database>,
211    listener: TcpListener,
212    admin_listener: Option<TcpListener>,
213    metrics_listener: Option<TcpListener>,
214    shutdown: impl Future<Output = ()> + Send + 'static,
215) -> anyhow::Result<()> {
216    serve_on_with_reloads(
217        RoleSet::default(),
218        config,
219        database,
220        Sockets {
221            acme: Some(listener),
222            admin: admin_listener,
223            metrics: metrics_listener,
224        },
225        shutdown,
226        crate::reload::Reloads::none(),
227    )
228    .await
229}
230
231/// [`serve_on_with`], serving configuration reloads as well as requests.
232///
233/// The variant [`run`] uses, so a `SIGHUP` rebuilds both routers, the job
234/// registry, the notifier map and both TLS acceptors behind the sockets that are
235/// already bound. Every other caller goes through [`serve_on_with`] and gets a
236/// source that never fires, which costs one task that ends immediately.
237///
238/// See [`crate::reload`] for what a reload may change and what it refuses.
239pub async fn serve_on_with_reloads(
240    roles: RoleSet,
241    config: Arc<Config>,
242    database: Arc<Database>,
243    sockets: Sockets,
244    shutdown: impl Future<Output = ()> + Send + 'static,
245    reloads: crate::reload::Reloads,
246) -> anyhow::Result<()> {
247    let Sockets {
248        acme: listener,
249        admin: admin_listener,
250        metrics: metrics_listener,
251    } = sockets;
252
253    info!(
254        event = "server_startup",
255        outcome = "success",
256        roles = %roles.labels().join(","),
257        bind_address = %config.server.bind_address,
258        base_url = %config.server.base_url,
259        tls = config.server.tls.enabled,
260        // A PostgreSQL DSN carries `user:password@`; this line is the one
261        // place the whole value is printed at INFO on every start.
262        database_database_url = %acme_proxy_core::logfields::redact_url(&config.database.url)
263    );
264
265    // The schema, before anything reads a row. **One owner**: the `worker` role
266    // applies the migrations, and every other role checks and refuses by name.
267    // Opening the database used to do this as a side effect, which made every
268    // subcommand an upgrade step and let two processes starting together race
269    // `MIGRATOR::run` with no lock between them.
270    apply_or_require_schema(roles, &database).await?;
271
272    // One `shutdown` future, several consumers: both listeners and the job
273    // runner. Created here rather than beside `axum::serve` below so a signal
274    // arriving *during* startup is not ignored — profile assembly and the
275    // relay's first upstream contact both happen before anything binds. The
276    // relay task is held under `AbortOnDrop` so an error path below does not
277    // leak a task parked on a signal that will never arrive.
278    let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(false);
279    let _shutdown_relay = AbortOnDrop(tokio::spawn(async move {
280        shutdown.await;
281        let _ = shutdown_tx.send(true);
282    }));
283
284    // The enqueue side of the durable queue. Built before the profiles because
285    // a signer backend that defers issuance is handed one at construction, and
286    // process-wide for the reason `[audit]` is: one table, one runner, and a
287    // per-endpoint retry budget would make a job's pacing depend on which
288    // profile happened to queue it.
289    let job_queue = acme_proxy_jobs::jobs::JobQueue::new(database.clone(), &config.jobs);
290
291    let resolved = config.resolve_profiles().inspect_err(|error| {
292        error!(event = "profile_init_failed", outcome = "failure", error = %error);
293    })?;
294    // Everything that outlives a configuration generation — the signer backends
295    // above all, which are carried rather than rebuilt. See
296    // `crate::Assembly`.
297    let (assembly, parts) = Assembly::new(
298        roles,
299        &resolved,
300        database.clone(),
301        job_queue.clone(),
302        &config,
303    )
304    .inspect_err(|error| {
305        error!(event = "profile_init_failed", outcome = "failure", error = %error);
306    })?;
307    let assembly = Arc::new(assembly);
308    if roles.has(ProcessRole::Worker) {
309        store_first_crls(&parts.signers).await;
310    }
311
312    let generation = build_generation(roles, &config, &resolved, &assembly, &parts, None)
313        .inspect_err(|error| {
314            error!(event = "profile_init_failed", outcome = "failure", error = %error);
315        })?;
316
317    // Every endpoint in the first generation came up, so every one is announced.
318    // A reload announces only the endpoints it *added* — see
319    // `supervise_reloads`, which shares this function for exactly that reason.
320    for profile in &generation.profiles {
321        announce_profile(profile).await;
322    }
323
324    let Generation {
325        profiles: _,
326        acme_app,
327        admin_app,
328        job_registry,
329        tls,
330        admin_tls,
331        logins,
332    } = generation;
333
334    // The one task that drains the queue. Held under `AbortOnDrop` for the same
335    // reason the reapers are — an un-cancelled loop holding an `Arc<Database>`
336    // per `serve_on` call — and, unlike the relay tasks this replaced, it does
337    // *not* need to outlive this function: the work is durable now, so a job cut
338    // short is re-claimed from its own row rather than lost. It takes the same
339    // shutdown signal as both listeners, so a stop is graceful rather than an
340    // abort: it releases its leases on the way out, and a restart therefore
341    // re-claims its own work immediately instead of waiting one out. The wait
342    // for that stop is at the end of this function — `AbortOnDrop` is the
343    // backstop for the paths that return before it, where an abort is the only
344    // thing on offer.
345    //
346    // Neither the registry it drains nor the `[jobs]` section it paces itself
347    // from is a value: both are cells a reload republishes, so a changed
348    // retention, a rebuilt notify map and a retuned lease or concurrency all
349    // reach the runner without restarting it. `jobs.max_attempts` is the third
350    // piece and does not come through here — it belongs to the enqueue side, so
351    // it is published onto `job_queue` itself.
352    //
353    // Only the `worker` role runs it. The cells are created either way, so the
354    // supervisor keeps one shape and a reload republishes into them whether or
355    // not anything is draining here.
356    let (registry_tx, registry_rx) = tokio::sync::watch::channel(Arc::new(job_registry));
357    let (jobs_tx, jobs_rx) = tokio::sync::watch::channel(Arc::new(config.jobs.clone()));
358    let mut job_runner = roles.has(ProcessRole::Worker).then(|| {
359        AbortOnDrop(acme_proxy_jobs::jobs::spawn_runner_watching(
360            job_queue,
361            registry_rx,
362            jobs_rx,
363            shutdown_rx.clone(),
364        ))
365    });
366    if !roles.has(ProcessRole::Worker) {
367        // Advisory rather than a refusal: an operator may legitimately start an
368        // `acme` process before the worker, and the rows it queues are durable.
369        // But a deployment that never runs one issues nothing — every relayed
370        // order, notification, sweep and challenge validation waits for ever —
371        // so this has to be visible.
372        warn!(
373            event = "server_role_no_worker",
374            outcome = "advisory",
375            roles = %roles.labels().join(","),
376            "this process runs no worker, so nothing here drains the job queue: \
377             challenge validation, notifications and the periodic sweeps all wait for a \
378             process started with `--role worker`"
379        );
380    }
381
382    // Only when this process actually holds the ACME socket. A worker-only
383    // process announcing an address it never bound would send an operator
384    // looking for a listener that is somebody else's.
385    if roles.has(ProcessRole::Acme) {
386        info!(
387            event = "server_listening",
388            outcome = "success",
389            bind_address = %config.server.bind_address,
390            protocol = if tls.is_some() { "https" } else { "http" }
391        );
392    }
393
394    // One accept loop per role, each owning a socket a reload can replace and a
395    // TLS mode it can switch — see `acme_proxy_net::listener`. `axum::serve` below is
396    // handed one of these instead of a `TcpListener` and therefore outlives
397    // every rebind, which is what removes the listener from the list of things
398    // only a restart can change.
399    let admin_bound = bound_address(admin_listener.as_ref(), &config.admin.bind_address);
400    let metrics_bound = bound_address(metrics_listener.as_ref(), &config.metrics.bind_address);
401    let (acme_socket, acme_handle) = acme_proxy_net::listener::spawn("acme", listener, tls);
402    let (admin_socket, admin_handle) =
403        acme_proxy_net::listener::spawn("admin", admin_listener, admin_tls);
404    let (metrics_socket, metrics_handle) =
405        acme_proxy_net::listener::spawn("metrics", metrics_listener, None);
406
407    // Behind a swap cell rather than served directly, so a configuration reload
408    // can replace the whole router without the socket moving. The cell is what
409    // `axum::serve` holds; `acme_app` itself is only ever generation one.
410    let (acme_router_tx, acme_router_rx) = crate::reload::router_channel(acme_app);
411    let acme = serve_role(
412        crate::reload::swappable(acme_router_rx),
413        acme_socket,
414        shutdown_rx.clone(),
415    );
416
417    // Opened whether or not the panel is on, unlike the app inside it: with
418    // `admin.enabled` reloadable, a cell created only when the panel starts
419    // would be the one thing a reload turning it on could not reach. An empty
420    // `Router` answers `404` to everything, which is also what the panel being
421    // switched off later publishes here.
422    let (admin_router_tx, admin_router_rx) =
423        crate::reload::router_channel(admin_app.unwrap_or_default());
424    let admin = serve_role(
425        crate::reload::swappable(admin_router_rx),
426        admin_socket,
427        shutdown_rx.clone(),
428    );
429    if roles.has(ProcessRole::Admin) && config.admin.enabled {
430        announce_admin_listener(&config, &database, &admin_bound).await;
431    }
432
433    // The third listener. Served directly rather than through a
434    // `reload::router_channel` like the other two, and the asymmetry is
435    // deliberate: this router has one route whose only state is the registry,
436    // and the registry is carried across generations rather than rebuilt (see
437    // `Assembly`), so a new generation could put nothing new in it.
438    if config.metrics.enabled {
439        announce_metrics_listener(&metrics_bound);
440    }
441    let metrics = serve_role(
442        acme_proxy_protocol::router::metrics_app(assembly.metrics.clone()),
443        metrics_socket,
444        shutdown_rx,
445    );
446
447    // The supervisor owns every cell sender from here on, which is what makes it
448    // the only writer: a generation is published by one task or by nobody.
449    // Aborted on drop, so an error return below does not leave it parked on a
450    // channel nothing will ever send to.
451    let _reload_supervisor = AbortOnDrop(tokio::spawn(supervise_reloads(
452        roles,
453        reloads,
454        config.clone(),
455        resolved,
456        assembly,
457        Cells {
458            acme_router: acme_router_tx,
459            admin_router: admin_router_tx,
460            job_registry: registry_tx,
461            jobs: jobs_tx,
462            acme: acme_handle,
463            admin: admin_handle,
464            metrics: metrics_handle,
465        },
466        logins,
467    )));
468
469    // Nothing is drained here any more. A notification in flight at shutdown is
470    // a `notify_deliver` row, not a spawned task: the runner releases its lease
471    // on the way out and whoever starts next claims it. That is what replaced a
472    // best-effort five-second drain which still lost anything slower than it.
473    tokio::try_join!(acme, admin, metrics)?;
474
475    // The listeners are done; the runner may not be. It takes the same shutdown
476    // signal, but returning here would drop its `AbortOnDrop` and cut its stop
477    // short — and in a worker-only process, where the three futures above have
478    // nothing to finish, that happens immediately. Waiting for it is what makes
479    // the release of its leases real: without it a job in flight waits out its
480    // whole lease before another process may claim it, and `job_runner_stopped`
481    // is never logged. Bounded, because a handler that ignores the signal must
482    // not hold the process open: the runner's own budget plus a margin.
483    if let Some(mut runner) = job_runner.take()
484        && tokio::time::timeout(RUNNER_SHUTDOWN_BUDGET, runner.take())
485            .await
486            .is_err()
487    {
488        warn!(
489            event = "job_runner_shutdown_timed_out",
490            outcome = "failure",
491            budget_ms = acme_proxy_core::logfields::millis(RUNNER_SHUTDOWN_BUDGET),
492            "the job runner did not stop in time; its leases expire on their own"
493        );
494    }
495    Ok(())
496}
497
498/// How long [`serve_on`] waits for the job runner to stop once both listeners
499/// have. The runner's own drain budget is five seconds, so this is that plus
500/// room for the writes that release its leases.
501const RUNNER_SHUTDOWN_BUDGET: std::time::Duration = std::time::Duration::from_secs(10);
502
503/// Stores each local CA's first CRL before this process serves anything.
504///
505/// Every role serves `GET /crl` from the stored row, and the read side never
506/// signs, so a CA nothing has met yet has no CRL to serve. The daily
507/// `CrlSweepJob` stores one on its first pass, but that pass runs on the
508/// runner's first tick — after the listeners are up — so an all-in-one server
509/// would answer `500` for the first moments of its life. Doing it here, once,
510/// in the process that holds the keys, closes that for every topology with a
511/// worker. A failure is logged where it happened
512/// (`local_ca_crl_initialization_failed`) and the sweep's first pass tries
513/// again, so it never stops the process.
514/// Once per CA, not once per profile: two profiles over one `[signer]` section
515/// share a backend, and a third naming the same CA by another path is still
516/// the same issuer. Refreshing twice would sign the same CRL twice and export
517/// it twice.
518pub(crate) async fn store_first_crls(signers: &acme_proxy_signer::SignerSet) {
519    let mut done: std::collections::HashSet<String> = std::collections::HashSet::new();
520    for (_, backend) in signers.by_profile() {
521        if let Some(refresher) = backend.crl_refresher()
522            && done.insert(refresher.issuer().to_string())
523        {
524            let _ = refresher.refresh().await;
525        }
526    }
527}
528
529/// **The `worker` role owns the schema.** Every other role checks and stops by
530/// name, which is what removes the startup race: on `SQLite` `sqlx` takes no
531/// migration lock, so two processes that both ran `MIGRATOR::run` could
532/// interleave. (`PostgreSQL`'s advisory lock would serialize them, but a
533/// process that is not the worker still has no business changing the schema.) Naming `acme-proxy migrate` in the refusal also means a split
534/// deployment fails at the process that started too early rather than later, as
535/// a missing table in a request.
536///
537/// All-in-one is unaffected: a default `serve` runs `worker`, so a fresh
538/// database is migrated exactly as it always was.
539async fn apply_or_require_schema(roles: RoleSet, database: &Arc<Database>) -> anyhow::Result<()> {
540    if roles.has(ProcessRole::Worker) {
541        return database.migrate().await.map_err(|error| {
542            error!(event = "db_migration_failed", outcome = "failure", error = %error);
543            anyhow::anyhow!("cannot apply the database migrations: {error}")
544        });
545    }
546
547    let pending = database.pending_migrations().await.map_err(|error| {
548        error!(event = "server_schema_check_failed", outcome = "failure", error = %error);
549        anyhow::anyhow!("cannot read the database schema version: {error}")
550    })?;
551    if pending.is_empty() {
552        return Ok(());
553    }
554
555    error!(
556        event = "server_schema_behind",
557        outcome = "failure",
558        pending = pending.len(),
559        roles = %roles.labels().join(","),
560    );
561    anyhow::bail!(
562        "the database is {} migration(s) behind and this process does not run the `worker` \
563         role, which owns the schema: run `acme-proxy migrate` (or start the worker) first",
564        pending.len()
565    )
566}
567
568/// Serves `app` on one role's socket until the process shuts down.
569///
570/// One shape for all three roles, where there used to be a boxed future per
571/// listener per TLS arm: [`acme_proxy_net::listener::RoleSocket`] is the same type
572/// whether the role is speaking TLS, speaking cleartext or — a socket having
573/// been closed by a reload — not serving at all, so the four cases collapse into
574/// this one call. Its future lives for the process: a rebind replaces what is
575/// underneath it, never the `axum::serve` above.
576fn serve_role(
577    app: axum::Router,
578    socket: acme_proxy_net::listener::RoleSocket,
579    shutdown: tokio::sync::watch::Receiver<bool>,
580) -> impl Future<Output = std::io::Result<()>> + Send {
581    axum::serve(
582        socket,
583        app.into_make_service_with_connect_info::<SocketAddr>(),
584    )
585    .with_graceful_shutdown(on_shutdown(shutdown))
586    .into_future()
587}
588
589/// A future that completes when the shutdown relay fires.
590async fn on_shutdown(mut receiver: tokio::sync::watch::Receiver<bool>) {
591    // An error means the sender was dropped, which only happens when the relay
592    // task itself is gone — treat it as "shut down" rather than parking
593    // forever.
594    let _ = receiver.wait_for(|ready| *ready).await;
595}
596
597async fn shutdown_signal() {
598    let ctrl_c = async {
599        tokio::signal::ctrl_c()
600            .await
601            .expect("failed to install Ctrl+C handler");
602    };
603
604    #[cfg(unix)]
605    let terminate = async {
606        tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
607            .expect("failed to install signal handler")
608            .recv()
609            .await;
610    };
611
612    #[cfg(not(unix))]
613    let terminate = std::future::pending::<()>();
614
615    tokio::select! {
616        _ = ctrl_c => {},
617        _ = terminate => {},
618    }
619}
620
621/// Aborts a background task when it goes out of scope.
622struct AbortOnDrop(tokio::task::JoinHandle<()>);
623
624impl AbortOnDrop {
625    /// Takes the handle back, so the caller can await the task's own stop
626    /// instead of aborting it. What is left behind is an already-finished
627    /// handle, whose abort on drop is a no-op.
628    fn take(&mut self) -> tokio::task::JoinHandle<()> {
629        std::mem::replace(&mut self.0, tokio::spawn(std::future::ready(())))
630    }
631}
632
633impl Drop for AbortOnDrop {
634    fn drop(&mut self) {
635        self.0.abort();
636    }
637}