use std::sync::Arc;
use std::time::Duration;
use tracing::{error, info, warn};
use super::logging;
use acme_proxy_core::config::Config;
use acme_proxy_net::tls;
use super::sockets::{Role, SocketPlans, check_metrics_config, plan_sockets};
use super::supervisor::Cells;
use super::{Assembly, GenerationParts};
use acme_proxy_protocol::profile::Profile;
use acme_proxy_protocol::router::build_app;
pub(crate) struct Generation {
pub(super) profiles: Vec<Arc<Profile>>,
pub(super) acme_app: axum::Router,
pub(super) admin_app: Option<axum::Router>,
pub(super) job_registry: acme_proxy_jobs::jobs::JobRegistry,
pub(super) tls: Option<tls::TlsSettings>,
pub(super) admin_tls: Option<tls::TlsSettings>,
pub(super) logins: Option<Arc<acme_proxy_admin::webadmin::LoginLimiter>>,
}
pub(crate) fn build_generation(
roles: crate::RoleSet,
config: &Arc<Config>,
resolved: &[acme_proxy_core::config::ProfileConfig],
assembly: &Assembly,
parts: &GenerationParts,
previous_logins: Option<&acme_proxy_admin::webadmin::LoginLimiter>,
) -> anyhow::Result<Generation> {
let admin_enabled = config.admin.enabled && roles.has(crate::ProcessRole::Admin);
let database = assembly.database.clone();
let profiles = crate::profile::build_all_with(config, resolved, parts)?;
let tls = tls::from_config(&config.server)
.inspect_err(|error| {
error!(event = "tls_init_failed", outcome = "failure", error = %error);
})?
.map(|acceptor| {
tls::TlsSettings::new(
acceptor,
Duration::from_millis(config.server.tls.handshake_timeout_ms),
)
});
let admin_tls = match admin_enabled {
false => None,
true => tls::admin_from_config(&config.admin)
.inspect_err(|error| {
error!(event = "admin_tls_init_failed", outcome = "failure", error = %error);
})?
.map(|acceptor| {
tls::TlsSettings::new(
acceptor,
Duration::from_millis(config.admin.tls.handshake_timeout_ms),
)
}),
};
let job_registry =
job_registry_for(config, resolved, assembly, parts, &profiles, admin_enabled)?;
let auditor = Arc::new(
acme_proxy_jobs::auditor::Auditor::from_config(
&config.audit,
&config.dns,
database.clone(),
assembly.metrics.clone(),
)
.inspect_err(|error| {
error!(event = "audit_init_failed", outcome = "failure", error = %error);
})?,
);
let (admin_app, logins) = match admin_enabled {
false => (None, None),
true => {
let policy =
acme_proxy_admin::webadmin::filter::build(config).inspect_err(|error| {
error!(event = "admin_filter_init_failed", outcome = "failure", error = %error);
})?;
let (router, logins) = acme_proxy_admin::webadmin::build_admin_app_with_logins(
database.clone(),
config.clone(),
&profiles,
auditor.clone(),
assembly.notifiers.clone(),
assembly.jobs.clone(),
policy,
previous_logins,
);
(Some(router), Some(logins))
}
};
let acme_app = build_app(
database,
config.clone(),
profiles.clone(),
auditor,
assembly.metrics.clone(),
assembly.jobs.clone(),
);
Ok(Generation {
profiles,
acme_app,
admin_app,
job_registry,
tls,
admin_tls,
logins,
})
}
fn job_registry_for(
config: &Config,
resolved: &[acme_proxy_core::config::ProfileConfig],
assembly: &Assembly,
parts: &GenerationParts,
profiles: &[Arc<Profile>],
admin_enabled: bool,
) -> anyhow::Result<acme_proxy_jobs::jobs::JobRegistry> {
let database = assembly.database.clone();
let mut job_registry = acme_proxy_jobs::jobs::JobRegistry::new();
let mut registered: Vec<usize> = Vec::new();
let mut refreshers: Vec<Arc<dyn acme_proxy_signer::CrlRefresher>> = Vec::new();
let mut relays: Vec<(String, acme_proxy_signer::relay::RelayState)> = Vec::new();
let mut backends = parts.signers.by_profile();
backends.sort_by(|a, b| a.0.cmp(&b.0));
for (profile, backend) in &backends {
relays.extend(backend.relay_state().map(|state| (profile.clone(), state)));
let identity = Arc::as_ptr(backend).cast::<()>() as usize;
if registered.contains(&identity) {
continue;
}
registered.push(identity);
refreshers.extend(backend.crl_refresher());
}
if !refreshers.is_empty() {
job_registry
.register(Arc::new(
acme_proxy_signer::local_ca::sweep::CrlRegenerateJob::new(refreshers.clone()),
))
.inspect_err(|error| {
error!(event = "job_registry_init_failed", outcome = "failure", error = %error);
})?;
job_registry
.register(Arc::new(
acme_proxy_signer::local_ca::sweep::CrlSweepJob::new(refreshers),
))
.inspect_err(|error| {
error!(event = "job_registry_init_failed", outcome = "failure", error = %error);
})?;
}
if !relays.is_empty() {
job_registry
.register(Arc::new(acme_proxy_signer::relay::flow::RelayJob::new(
database.clone(),
relays,
)))
.inspect_err(|error| {
error!(event = "job_registry_init_failed", outcome = "failure", error = %error);
})?;
}
job_registry
.register(Arc::new(
acme_proxy_protocol::acme::revoke::SignerRevokeJob::new(
database.clone(),
Arc::new(
acme_proxy_jobs::auditor::Auditor::offline(database.clone())
.with_metrics(assembly.metrics.clone()),
),
backends.clone(),
assembly.notifiers.clone(),
),
))
.inspect_err(|error| {
error!(event = "job_registry_init_failed", outcome = "failure", error = %error);
})?;
job_registry
.register(Arc::new(
acme_proxy_protocol::acme::issue::SignerIssueJob::new(
database.clone(),
Arc::new(
acme_proxy_jobs::auditor::Auditor::offline(database.clone())
.with_metrics(assembly.metrics.clone()),
),
backends.clone(),
assembly.notifiers.clone(),
),
))
.inspect_err(|error| {
error!(event = "job_registry_init_failed", outcome = "failure", error = %error);
})?;
job_registry
.register(Arc::new(
acme_proxy_protocol::acme::validate::ChallengeValidateJob::new(
database.clone(),
Arc::new(
acme_proxy_jobs::auditor::Auditor::offline(database.clone())
.with_metrics(assembly.metrics.clone()),
),
profiles
.iter()
.map(|profile| (profile.name.clone(), profile.clone()))
.collect(),
),
))
.inspect_err(|error| {
error!(event = "job_registry_init_failed", outcome = "failure", error = %error);
})?;
job_registry
.register(Arc::new(acme_proxy_jobs::notify::NotifyJob::new(
assembly.notifiers.clone(),
)))
.inspect_err(|error| {
error!(event = "job_registry_init_failed", outcome = "failure", error = %error);
})?;
if let Some(digest) = acme_proxy_jobs::notify::expiry::ExpiryDigestJob::from_profiles(
resolved,
assembly.notifiers.clone(),
database.clone(),
assembly.jobs.clone(),
) {
job_registry
.register(Arc::new(digest))
.inspect_err(|error| {
error!(event = "job_registry_init_failed", outcome = "failure", error = %error);
})?;
}
let ttl = Duration::from_secs(config.nonce.ttl_seconds);
let mut sweeps = vec![acme_proxy_jobs::jobs::SweepJob::nonces(
database.clone(),
ttl,
)];
if config.audit.retention_days > 0 {
sweeps.push(acme_proxy_jobs::jobs::SweepJob::audit(
database.clone(),
config.audit.retention_days,
));
}
if config.jobs.retention_days > 0 {
sweeps.push(acme_proxy_jobs::jobs::SweepJob::jobs(
database.clone(),
config.jobs.retention_days,
));
}
let order_retention: Vec<(String, u64)> = resolved
.iter()
.filter(|profile| profile.sections.order.retention_days > 0)
.map(|profile| (profile.name.clone(), profile.sections.order.retention_days))
.collect();
if !order_retention.is_empty() {
sweeps.push(acme_proxy_jobs::jobs::SweepJob::orders(
database.clone(),
order_retention,
));
}
if profiles
.iter()
.any(|profile| profile.signer_info.http01_tokens().is_some())
{
sweeps.push(acme_proxy_jobs::jobs::SweepJob::http01_tokens(
database.clone(),
));
}
if admin_enabled {
sweeps.push(acme_proxy_jobs::jobs::SweepJob::admin_sessions(
database.clone(),
Duration::from_secs(config.admin.session_idle_timeout_seconds),
config.admin.session_ttl_seconds,
));
}
for sweep in sweeps {
job_registry
.register(Arc::new(sweep))
.inspect_err(|error| {
error!(event = "job_registry_init_failed", outcome = "failure", error = %error);
})?;
}
Ok(job_registry)
}
pub(super) struct Reloaded {
pub(super) report: crate::reload::ReloadReport,
pub(super) config: Arc<Config>,
pub(super) resolved: Vec<acme_proxy_core::config::ProfileConfig>,
pub(super) logins: Option<Arc<acme_proxy_admin::webadmin::LoginLimiter>>,
pub(super) opened: Vec<(Role, String)>,
pub(super) mounted: Vec<Arc<Profile>>,
}
pub(super) struct Prepared {
config: Arc<Config>,
resolved: Vec<acme_proxy_core::config::ProfileConfig>,
parts: GenerationParts,
generation: Generation,
sockets: SocketPlans,
logging: logging::PreparedLogging,
logging_filter_source: logging::FilterSource,
mounted: Vec<Arc<Profile>>,
unmounted: Vec<String>,
}
impl Prepared {
pub(super) fn signers(&self) -> &acme_proxy_signer::SignerSet {
&self.parts.signers
}
}
pub(super) fn prepare_reload(
roles: crate::RoleSet,
config: &Arc<Config>,
resolved: &[acme_proxy_core::config::ProfileConfig],
assembly: &Assembly,
logins: Option<&acme_proxy_admin::webadmin::LoginLimiter>,
) -> Result<Prepared, crate::reload::ReloadError> {
use crate::reload::{Applied, ReloadError, check_frozen};
let next = Arc::new(Config::load().map_err(|error| ReloadError::Load(error.to_string()))?);
let next_resolved = next
.resolve_profiles()
.map_err(|error| ReloadError::Load(error.to_string()))?;
check_frozen(
&Applied {
config,
profiles: resolved,
},
&Applied {
config: &next,
profiles: &next_resolved,
},
)?;
let logging = logging::prepare_logging(&next.logging, logging::flag_override())
.map_err(ReloadError::Build)?;
let logging_filter_source = logging.filter_source;
if roles.has(crate::ProcessRole::Admin) {
acme_proxy_admin::webadmin::check_config(&next)
.map_err(|error| ReloadError::Build(error.to_string()))?;
}
check_metrics_config(&next).map_err(|error| ReloadError::Build(error.to_string()))?;
let sockets = plan_sockets(roles, config, &next)?;
let parts = assembly
.build_parts(&next_resolved, &next)
.map_err(|error| ReloadError::Build(error.to_string()))?;
let generation = build_generation(roles, &next, &next_resolved, assembly, &parts, logins)
.map_err(|error| ReloadError::Build(error.to_string()))?;
let running: std::collections::HashSet<&str> = resolved
.iter()
.map(|profile| profile.name.as_str())
.collect();
let mounted = generation
.profiles
.iter()
.filter(|profile| !running.contains(profile.name.as_str()))
.cloned()
.collect();
let next_names: std::collections::HashSet<&str> = next_resolved
.iter()
.map(|profile| profile.name.as_str())
.collect();
let unmounted = resolved
.iter()
.map(|profile| profile.name.clone())
.filter(|name| !next_names.contains(name.as_str()))
.collect();
Ok(Prepared {
config: next,
resolved: next_resolved,
parts,
generation,
sockets,
logging,
logging_filter_source,
mounted,
unmounted,
})
}
pub(super) fn publish_reload(
prepared: Prepared,
applied: &Arc<Config>,
assembly: &Assembly,
cells: &Cells,
generation: u64,
started: std::time::Instant,
) -> Reloaded {
use crate::reload::ReloadReport;
let Prepared {
config: next,
resolved: next_resolved,
parts,
generation: built,
sockets,
logging,
logging_filter_source,
mounted,
unmounted,
} = prepared;
let logging_reloaded = logging::publish_logging(logging);
let report = ReloadReport {
generation,
profiles: built
.profiles
.iter()
.map(|profile| profile.name.clone())
.collect(),
job_kinds: built.job_registry.kinds(),
tls_reloaded: built.tls.is_some(),
admin_tls_reloaded: built.admin_tls.is_some(),
listeners_rebound: sockets.rebound(),
logging_reloaded,
duration: started.elapsed(),
};
let next_logins = built.logins.clone();
assembly.publish_notifiers(parts.dispatchers);
assembly.publish_signers(parts.signers, parts.infos);
cells
.job_registry
.send_replace(Arc::new(built.job_registry));
cells.jobs.send_replace(Arc::new(next.jobs.clone()));
assembly.jobs.set_max_attempts(next.jobs.max_attempts);
cells.acme.set_tls(built.tls);
cells.admin.set_tls(built.admin_tls);
let opened = sockets.bound.clone();
sockets.publish(cells);
cells
.acme_router
.send_replace(built.acme_app.into_service::<axum::body::Body>());
cells.admin_router.send_replace(
built
.admin_app
.unwrap_or_default()
.into_service::<axum::body::Body>(),
);
for profile in unmounted {
warn!(
event = "profile_unmounted",
outcome = "advisory",
profile = %profile,
"the endpoint is no longer served: its accounts and orders stay in the \
database and come back if it is mounted again, but any issuance still in \
flight for it has no handler left to finish it"
);
}
if logging_filter_source.outranks_config() && applied.logging.filter != next.logging.filter {
warn!(
event = "server_logging_filter_overridden",
outcome = "advisory",
source = logging_filter_source.as_str(),
configured = %next.logging.filter,
);
}
Reloaded {
report,
config: next,
resolved: next_resolved,
logins: next_logins,
opened,
mounted,
}
}
pub(super) async fn announce_profile(profile: &Arc<Profile>) {
info!(
event = "profile_mounted",
outcome = "success",
profile = %profile.name,
directory = %profile.directory_url(),
challenge_bypass = profile.challenges.is_bypassed(),
eab_enabled = profile.eab.enabled
);
profile
.notify
.dispatch(acme_proxy_jobs::notify::NotifyEvent::ProfileMounted(
acme_proxy_jobs::notify::ProfileMountedData {
profile: profile.name.clone(),
},
))
.await;
}