mod common;
use std::sync::Arc;
use acme_proxy_core::config::Config;
use acme_proxy_server::ProcessRole;
use acme_proxy_server::RoleSet;
use acme_proxy_server::serve_on_with_reloads;
use acme_proxy_server::sockets::Sockets as ServerSockets;
use acme_proxy_store::db::Database;
use common::TempDir;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::{TcpListener, TcpStream};
struct Process {
acme: Option<std::net::SocketAddr>,
admin: Option<std::net::SocketAddr>,
shutdown: Option<tokio::sync::oneshot::Sender<()>>,
handle: tokio::task::JoinHandle<anyhow::Result<()>>,
}
impl Process {
async fn stop(mut self) {
if let Some(shutdown) = self.shutdown.take() {
let _ = shutdown.send(());
}
let _ = tokio::time::timeout(std::time::Duration::from_secs(10), self.handle).await;
}
}
fn write_config(dir: &TempDir) {
let ca = dir.join("ca");
write_config_at(dir, dir, &format!("{}.key", ca.display()));
}
fn write_config_at(dir: &TempDir, into: &TempDir, key_path: &str) {
let ca = dir.join("ca");
let body = format!(
r#"
[database]
url = "sqlite://{dir}/roles.db"
[server]
base_url = "http://localhost:3000"
[admin]
enabled = true
[profiles.default]
signer.local_ca.cert_path = "{ca}.pem"
signer.local_ca.key_path = "{key_path}"
signer.local_ca.crl_path = "{ca}.crl"
"#,
dir = dir.path().display(),
ca = ca.display(),
);
std::fs::write(into.join("config.toml"), body).unwrap();
}
async fn initialise(config: &Config) {
let database = Arc::new(
Database::open(&config.database.url)
.await
.expect("the database must open"),
);
acme_proxy::cli::init(
acme_proxy_core::palette::Palette::plain(),
&Arc::new(config.clone()),
database,
)
.await
.expect("init must succeed");
}
fn load_from(dir: &TempDir) -> Config {
unsafe {
std::env::set_var("ACME_PROXY_CONFIG", dir.join("config").to_str().unwrap());
}
Config::load().expect("the configuration must load")
}
async fn start(config: &Config, roles: RoleSet) -> Process {
let database = Arc::new(
Database::open(&config.database.url)
.await
.expect("the database must open"),
);
let (acme_listener, acme) = match roles.has(ProcessRole::Acme) {
false => (None, None),
true => {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
(Some(listener), Some(addr))
}
};
let (admin_listener, admin) = match roles.has(ProcessRole::Admin) {
false => (None, None),
true => {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
(Some(listener), Some(addr))
}
};
let (shutdown, rx) = tokio::sync::oneshot::channel::<()>();
let handle = tokio::spawn(serve_on_with_reloads(
roles,
Arc::new(config.clone()),
database,
ServerSockets {
acme: acme_listener,
admin: admin_listener,
metrics: None,
},
async {
let _ = rx.await;
},
acme_proxy_server::reload::Reloads::none(),
));
Process {
acme,
admin,
shutdown: Some(shutdown),
handle,
}
}
async fn status_of(addr: std::net::SocketAddr, path: &str) -> u16 {
let mut stream = TcpStream::connect(addr)
.await
.expect("the port must accept");
let request = format!("GET {path} HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n");
stream.write_all(request.as_bytes()).await.unwrap();
let mut response = Vec::new();
stream.read_to_end(&mut response).await.unwrap();
let head = String::from_utf8_lossy(&response);
head.split_whitespace()
.nth(1)
.and_then(|code| code.parse().ok())
.unwrap_or_else(|| panic!("no status line in: {head}"))
}
#[tokio::test]
async fn each_role_serves_only_its_own_surface() {
let dir = TempDir::new("roles-surfaces");
write_config(&dir);
let config = load_from(&dir);
initialise(&config).await;
let worker = start(&config, RoleSet::parse(Some("worker")).unwrap()).await;
let acme = start(&config, RoleSet::parse(Some("acme")).unwrap()).await;
let admin = start(&config, RoleSet::parse(Some("admin")).unwrap()).await;
let acme_addr = acme.acme.expect("the acme role binds the ACME socket");
let admin_addr = admin.admin.expect("the admin role binds the admin socket");
assert!(
acme.admin.is_none(),
"an acme process holds no admin socket"
);
assert!(
admin.acme.is_none(),
"an admin process holds no ACME socket"
);
assert_eq!(
status_of(acme_addr, "/profile/default/directory").await,
200
);
assert_eq!(status_of(acme_addr, "/api/accounts").await, 404);
assert_eq!(status_of(admin_addr, "/api/accounts").await, 401);
assert_eq!(
status_of(admin_addr, "/profile/default/directory").await,
404
);
admin.stop().await;
acme.stop().await;
worker.stop().await;
}
#[tokio::test]
async fn an_acme_process_queues_work_a_worker_performs() {
let dir = TempDir::new("roles-queue");
write_config(&dir);
let config = load_from(&dir);
initialise(&config).await;
let acme = start(&config, RoleSet::parse(Some("acme")).unwrap()).await;
let acme_addr = acme.acme.expect("the acme role binds the ACME socket");
assert_eq!(
status_of(acme_addr, "/profile/default/directory").await,
200
);
let database = Arc::new(
Database::open(&config.database.url)
.await
.expect("the database must open"),
);
let queue = acme_proxy_jobs::jobs::JobQueue::new(database.clone(), &config.jobs);
let sweep = acme_proxy_jobs::jobs::SweepJob::nonces(
database.clone(),
std::time::Duration::from_secs(1),
);
assert!(
queue
.enqueue(acme_proxy_jobs::jobs::JobSpec::now(
acme_proxy_jobs::jobs::JobHandler::kind(&sweep),
"nonces",
))
.await
.expect("the enqueue must land")
);
tokio::time::sleep(std::time::Duration::from_millis(300)).await;
assert_eq!(
acme_proxy_store::job::Job::count_live("nonce_sweep", &database)
.await
.unwrap(),
1,
"an acme-only process must not have run the job"
);
let worker = start(&config, RoleSet::parse(Some("worker")).unwrap()).await;
let mut ran = false;
for _ in 0..200 {
let row =
acme_proxy_store::job::Job::find_latest_by_dedup("nonce_sweep", "nonces", &database)
.await
.unwrap()
.expect("the row exists");
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs() as i64;
if row.run_at > now {
ran = true;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
assert!(ran, "the worker must have performed the queued sweep");
worker.stop().await;
acme.stop().await;
}
#[tokio::test]
async fn a_process_without_the_worker_role_refuses_an_unmigrated_database() {
let dir = TempDir::new("roles-schema");
write_config(&dir);
let config = load_from(&dir);
let database = Arc::new(
Database::open(&config.database.url)
.await
.expect("the database must open"),
);
assert!(
!database.pending_migrations().await.unwrap().is_empty(),
"the fixture must start behind, or this proves nothing"
);
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let error = serve_on_with_reloads(
RoleSet::parse(Some("acme")).unwrap(),
Arc::new(config.clone()),
database.clone(),
ServerSockets {
acme: Some(listener),
admin: None,
metrics: None,
},
std::future::pending(),
acme_proxy_server::reload::Reloads::none(),
)
.await
.expect_err("an acme process must refuse a schema it does not own");
let message = error.to_string();
assert!(message.contains("worker"), "{message}");
assert!(message.contains("acme-proxy migrate"), "{message}");
assert!(!database.pending_migrations().await.unwrap().is_empty());
}
#[tokio::test]
async fn the_three_roles_start_concurrently_over_one_database() {
let dir = TempDir::new("roles-concurrent");
write_config(&dir);
let config = load_from(&dir);
initialise(&config).await;
let (worker, acme, admin) = tokio::join!(
start(&config, RoleSet::parse(Some("worker")).unwrap()),
start(&config, RoleSet::parse(Some("acme")).unwrap()),
start(&config, RoleSet::parse(Some("admin")).unwrap()),
);
assert_eq!(
status_of(acme.acme.unwrap(), "/profile/default/directory").await,
200
);
assert_eq!(status_of(admin.admin.unwrap(), "/api/accounts").await, 401);
assert!(!worker.handle.is_finished());
assert!(!acme.handle.is_finished());
assert!(!admin.handle.is_finished());
admin.stop().await;
acme.stop().await;
worker.stop().await;
}
async fn fetch(addr: std::net::SocketAddr, path: &str) -> (u16, Vec<u8>) {
let mut stream = TcpStream::connect(addr)
.await
.expect("the port must accept");
let request = format!("GET {path} HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n");
stream.write_all(request.as_bytes()).await.unwrap();
let mut response = Vec::new();
stream.read_to_end(&mut response).await.unwrap();
let split = response
.windows(4)
.position(|window| window == b"\r\n\r\n")
.expect("a response has a head");
let head = String::from_utf8_lossy(&response[..split]);
let status = head
.split_whitespace()
.nth(1)
.and_then(|code| code.parse().ok())
.unwrap_or_else(|| panic!("no status line in: {head}"));
(status, response[split + 4..].to_vec())
}
#[tokio::test]
async fn the_acme_and_admin_processes_never_touch_the_ca_key() {
let dir = TempDir::new("roles-keyless");
write_config(&dir);
let config = load_from(&dir);
initialise(&config).await;
let ca_pem = std::fs::read_to_string(dir.join("ca.pem")).unwrap();
let keyless_dir = TempDir::new("roles-keyless-conf");
let absent_key = keyless_dir.join("no-such.key");
write_config_at(&dir, &keyless_dir, absent_key.to_str().unwrap());
let keyless = load_from(&keyless_dir);
let acme = start(&keyless, RoleSet::parse(Some("acme")).unwrap()).await;
let admin = start(&keyless, RoleSet::parse(Some("admin")).unwrap()).await;
let acme_addr = acme.acme.expect("the acme role binds the ACME socket");
let (status, body) = fetch(acme_addr, "/profile/default/ca.pem").await;
assert_eq!(status, 200);
assert_eq!(String::from_utf8(body).unwrap(), ca_pem);
assert_eq!(status_of(admin.admin.unwrap(), "/api/accounts").await, 401);
assert!(!acme.handle.is_finished() && !admin.handle.is_finished());
let database = Arc::new(Database::open(&config.database.url).await.unwrap());
let queue = acme_proxy_jobs::jobs::JobQueue::new(database.clone(), &config.jobs);
let order_id = claimed_order(&database, &queue).await;
tokio::time::sleep(std::time::Duration::from_millis(300)).await;
let order = acme_proxy_store::order::Order::find_by_id(&order_id, &database)
.await
.unwrap()
.unwrap();
assert!(order.certificate.is_none(), "nothing without the key signs");
let worker = start(&config, RoleSet::parse(Some("worker")).unwrap()).await;
let order = settled(&database, &order_id).await;
assert!(order.certificate.is_some(), "the worker must have signed");
let route = acme_proxy_signer::revocation_route(
&keyless.resolve_profiles().unwrap()[0].sections.signer,
)
.unwrap();
let audit = acme_proxy_jobs::auditor::Auditor::offline(database.clone());
acme_proxy_protocol::acme::revoke::Revocations {
database: &database,
audit: &audit,
notify: None,
revoker: acme_proxy_protocol::acme::revoke::Revoker::for_route(
&route,
&queue,
std::time::Duration::ZERO,
),
}
.revoke_order(
&order_id,
Some(1),
acme_proxy_core::audit::Actor::cli(),
acme_proxy_core::audit::ClientContext::default(),
)
.await
.expect("a local CA's revocation needs no key");
let serial = order.cert_serial.clone().unwrap();
let mut listed = false;
for _ in 0..300 {
let (status, der) = fetch(acme_addr, "/profile/default/crl").await;
assert_eq!(status, 200);
use x509_parser::prelude::FromDer;
let (_, crl) =
x509_parser::revocation_list::CertificateRevocationList::from_der(&der).unwrap();
if crl
.iter_revoked_certificates()
.any(|entry| hex::encode(entry.raw_serial()).eq_ignore_ascii_case(&serial))
{
listed = true;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
assert!(listed, "the acme process must serve the worker's CRL");
assert!(
!absent_key.exists(),
"a process without the worker role must never create a CA key"
);
assert_eq!(
std::fs::read_to_string(dir.join("ca.pem")).unwrap(),
ca_pem,
"nor replace the CA"
);
worker.stop().await;
admin.stop().await;
acme.stop().await;
}
#[tokio::test]
async fn a_process_without_the_worker_role_refuses_a_missing_ca() {
let dir = TempDir::new("roles-no-ca");
write_config(&dir);
let config = load_from(&dir);
Database::connect_and_migrate(&config.database.url)
.await
.expect("the schema must apply");
let database = Arc::new(Database::open(&config.database.url).await.unwrap());
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let error = serve_on_with_reloads(
RoleSet::parse(Some("acme")).unwrap(),
Arc::new(config.clone()),
database,
ServerSockets {
acme: Some(listener),
admin: None,
metrics: None,
},
std::future::pending(),
acme_proxy_server::reload::Reloads::none(),
)
.await
.expect_err("an acme process must refuse to start without a CA");
let message = error.to_string();
assert!(message.contains("acme-proxy init"), "{message}");
assert!(message.contains("worker"), "{message}");
assert!(!dir.join("ca.pem").exists());
assert!(!dir.join("ca.key").exists());
}
async fn claimed_order(
database: &Arc<Database>,
queue: &acme_proxy_jobs::jobs::JobQueue,
) -> String {
use acme_proxy_core::identifier::Identifier;
use acme_proxy_store::account::Account;
use acme_proxy_store::order::Order;
use base64::prelude::*;
let (account, _) = Account::find_or_create(
"default",
&[1, 2, 3],
vec![],
&acme_proxy_core::audit::ClientContext::default(),
database,
)
.await
.unwrap();
let mut order = Order::create(
"default",
account.id,
vec![Identifier::dns("a.example.com")],
2_000_000_000,
None,
None,
database,
)
.await
.unwrap();
let mut authz = acme_proxy_store::authz::Authorization::create(
order.id,
Identifier::dns("a.example.com"),
order.expires,
database,
)
.await
.unwrap();
assert!(authz.mark_valid(database).await.unwrap());
assert!(order.mark_ready(database).await.unwrap());
assert!(order.claim_for_finalize(database).await.unwrap());
let csr = BASE64_URL_SAFE_NO_PAD
.decode(common::make_csr("a.example.com"))
.unwrap();
assert!(
queue
.enqueue(acme_proxy_protocol::acme::issue::signer_issue_spec(
&order,
&csr,
&acme_proxy_core::audit::ClientContext::default(),
None,
))
.await
.unwrap()
);
order.id.to_string()
}
async fn settled(database: &Database, id: &str) -> acme_proxy_store::order::Order {
for _ in 0..500 {
let order = acme_proxy_store::order::Order::find_by_id(id, database)
.await
.unwrap()
.unwrap();
if order.status != acme_proxy_store::status::OrderStatus::Processing {
return order;
}
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
panic!("order {id} never left `processing`");
}