Skip to main content

acme_proxy_server/
supervisor.rs

1//! The one task that serves reload requests for the life of the process.
2//!
3//! Reloads are serialized through it, and a generation is published into the
4//! `watch` cells with no `.await` between the sends, so no request ever sees
5//! half of one generation and half of the next. A reload that fails to build
6//! publishes nothing.
7
8use std::sync::Arc;
9
10use tracing::{error, info, warn};
11
12use acme_proxy_core::config::Config;
13
14use super::Assembly;
15use super::generation::{announce_profile, prepare_reload, publish_reload};
16use super::sockets::{Role, announce_admin_listener, announce_metrics_listener};
17
18/// The cells one generation is published into.
19///
20/// Held by the supervisor and by nothing else. Every field is a `watch::Sender`,
21/// and `send_replace` is synchronous — so publishing a generation is a run of
22/// sends with no `.await` between them, which no other task can interleave with.
23/// That is what makes a reload atomic without a lock.
24pub(super) struct Cells {
25    pub(super) acme_router:
26        tokio::sync::watch::Sender<axum::routing::RouterIntoService<axum::body::Body>>,
27    pub(super) admin_router:
28        tokio::sync::watch::Sender<axum::routing::RouterIntoService<axum::body::Body>>,
29    pub(super) job_registry: tokio::sync::watch::Sender<Arc<acme_proxy_jobs::jobs::JobRegistry>>,
30    /// The runner's own pacing. Separate from the registry above because the two
31    /// reach it by different routes: the registry carries what a *handler*
32    /// captured, this carries what the *loop* re-reads each pass.
33    pub(super) jobs: tokio::sync::watch::Sender<Arc<acme_proxy_core::config::JobsConfig>>,
34    /// The three sockets. Each carries its role's TLS mode as well, since both
35    /// are read by the same accept loop and both are published the same
36    /// synchronous way — see [`acme_proxy_net::listener::ListenerHandle`].
37    pub(super) acme: acme_proxy_net::listener::ListenerHandle,
38    pub(super) admin: acme_proxy_net::listener::ListenerHandle,
39    pub(super) metrics: acme_proxy_net::listener::ListenerHandle,
40}
41
42/// Serves reload requests for the life of the process.
43///
44/// One task, so reloads are serialised: two overlapping rebuilds could publish
45/// their cells interleaved, and the second-newest generation would win some of
46/// them. Ends when the last [`crate::reload::ReloadHandle`] is dropped, which is
47/// what makes [`crate::reload::Reloads::none`] cost nothing.
48pub(super) async fn supervise_reloads(
49    roles: crate::RoleSet,
50    mut reloads: crate::reload::Reloads,
51    mut config: Arc<Config>,
52    mut resolved: Vec<acme_proxy_core::config::ProfileConfig>,
53    assembly: Arc<Assembly>,
54    cells: Cells,
55    mut logins: Option<Arc<acme_proxy_admin::webadmin::LoginLimiter>>,
56) {
57    let mut generation: u64 = 1;
58
59    while let Some(request) = reloads.recv().await {
60        let started = std::time::Instant::now();
61        info!(
62            event = "server_config_reload_requested",
63            outcome = "progress",
64            generation = generation,
65        );
66
67        // The build phase runs on a blocking thread, and that is not a
68        // precaution: `RelaySigner::from_config` contacts the upstream the first
69        // time it is built for an account with no `kid` sidecar yet, on a scoped
70        // OS thread it then *joins*. Mounting a relay profile by `SIGHUP` would
71        // otherwise park a runtime worker for as long as
72        // `signer.relay.poll_timeout_secs` allows — five minutes by default —
73        // with every connection that worker was polling parked behind it. The
74        // publish phase stays on this task, where its lack of an await point is
75        // what makes a generation unobservable half-applied.
76        let outcome = {
77            let config = config.clone();
78            let resolved = resolved.clone();
79            let assembly = assembly.clone();
80            let logins = logins.clone();
81            tokio::task::spawn_blocking(move || {
82                prepare_reload(roles, &config, &resolved, &assembly, logins.as_deref())
83            })
84            .await
85            .unwrap_or_else(|error| {
86                Err(crate::reload::ReloadError::Build(format!(
87                    "the reload build task did not finish: {error}"
88                )))
89            })
90        };
91        // A CA this reload mounted has no stored CRL yet, and the read side
92        // never signs one. Stored here, before the routers that serve it are
93        // published, for startup's reason — see `store_first_crls`. Run over
94        // every CA of the new generation, once each: for one already serving, a
95        // refresh that finds a fresh CRL and nothing expired signs nothing.
96        if let Ok(prepared) = &outcome
97            && roles.has(super::ProcessRole::Worker)
98        {
99            super::store_first_crls(prepared.signers()).await;
100        }
101        let outcome = outcome.map(|prepared| {
102            publish_reload(
103                prepared,
104                &config,
105                &assembly,
106                &cells,
107                generation + 1,
108                started,
109            )
110        });
111
112        match outcome {
113            Ok(reloaded) => {
114                let report = reloaded.report;
115                config = reloaded.config;
116                resolved = reloaded.resolved;
117                logins = reloaded.logins;
118                generation = report.generation;
119                info!(
120                    event = "server_config_reloaded",
121                    outcome = "success",
122                    generation = report.generation,
123                    profiles = ?report.profiles,
124                    job_kinds = ?report.job_kinds,
125                    tls_reloaded = report.tls_reloaded,
126                    admin_tls_reloaded = report.admin_tls_reloaded,
127                    listeners_rebound = ?report.listeners_rebound,
128                    logging_reloaded = report.logging_reloaded,
129                    duration_ms = acme_proxy_core::logfields::millis(report.duration),
130                );
131                // After the reload's own line, and under the new configuration,
132                // since that is what these describe. Each is the same
133                // announcement startup makes for a listener that has just come
134                // up — including the panel's two warnings, which is why this is
135                // here rather than inside the synchronous publishing run.
136                for (role, address) in reloaded.opened {
137                    match role {
138                        Role::Acme => info!(
139                            event = "server_listening",
140                            outcome = "success",
141                            bind_address = %address,
142                            protocol = if config.server.tls.enabled { "https" } else { "http" }
143                        ),
144                        Role::Admin => {
145                            announce_admin_listener(&config, &assembly.database, &address).await;
146                        }
147                        Role::Metrics => announce_metrics_listener(&address),
148                    }
149                }
150                // An endpoint this reload mounted really did come up, so it gets
151                // the same announcement and the same notification startup makes
152                // for one. An endpoint that was *already* mounted stays silent:
153                // `profile_mounted` is a lifecycle event and not a heartbeat,
154                // and re-firing it per `SIGHUP` would make the notify surface
155                // noisiest in exactly the config-managed deployments that would
156                // least want it. Here rather than in the publishing run because
157                // dispatching reaches the database.
158                for profile in reloaded.mounted {
159                    announce_profile(&profile).await;
160                }
161                if let Some(respond) = request.respond {
162                    let _ = respond.send(Ok(report));
163                }
164            }
165            Err(error) => {
166                // Two names, because they are two different things for whoever
167                // is reading: a refusal is a configuration an operator must
168                // change, a failure is one the server could not build.
169                match &error {
170                    crate::reload::ReloadError::Frozen { .. } => warn!(
171                        event = "server_config_reload_refused",
172                        outcome = "failure",
173                        generation = generation,
174                        reason = error.kind(),
175                        error = %error,
176                    ),
177                    _ => error!(
178                        event = "server_config_reload_failed",
179                        outcome = "failure",
180                        generation = generation,
181                        reason = error.kind(),
182                        error = %error,
183                    ),
184                }
185                if let Some(respond) = request.respond {
186                    let _ = respond.send(Err(error));
187                }
188            }
189        }
190    }
191}