diff --git a/rust/Cargo.lock b/rust/Cargo.lock index ef67753f..6c8ac4a4 100644 --- a/rust/Cargo.lock +++ b/rust/Cargo.lock @@ -83,24 +83,37 @@ version = "0.30.0" dependencies = [ "adc-backend-api7", "adc-backend-apisix", + "adc-backend-apisix-standalone", "adc-backend-core", "adc-converter-openapi", "adc-differ", "adc-sdk", + "axum", + "axum-server", "chrono", "clap", "dotenvy", "glob", "http", + "http-body-util", "humantime", "indicatif", + "libc", + "lru", + "reqwest 0.12.28", + "rustls", + "rustls-pemfile", + "serde", "serde_json", "serde_yaml_ng", "thiserror", "tokio", + "tower", "tracing", "tracing-indicatif", "tracing-subscriber", + "url", + "uuid", ] [[package]] @@ -236,6 +249,15 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "arc-swap" +version = "1.9.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c049c0be4daef0b145cb3555416b3b8ef5b7888a38aea1a3a155801fe7b0810b" +dependencies = [ + "rustversion", +] + [[package]] name = "arrayvec" version = "0.7.8" @@ -340,6 +362,28 @@ dependencies = [ "tracing", ] +[[package]] +name = "axum-server" +version = "0.7.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c1ab4a3ec9ea8a657c72d99a03a824af695bd0fb5ec639ccbd9cd3543b41a5f9" +dependencies = [ + "arc-swap", + "bytes", + "fs-err", + "http", + "http-body", + "hyper", + "hyper-util", + "pin-project-lite", + "rustls", + "rustls-pemfile", + "rustls-pki-types", + "tokio", + "tokio-rustls", + "tower-service", +] + [[package]] name = "base64" version = "0.22.1" @@ -835,6 +879,12 @@ version = "1.0.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3f9eec918d3f24069decb9af1554cad7c880e2da24a9afd88aca000531ab82c1" +[[package]] +name = "foldhash" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d9c4f5dac5e15c24eb999c26181a6ca40b39fe946cbe4c263c7209467bc83af2" + [[package]] name = "foldhash" version = "0.2.0" @@ -860,6 +910,16 @@ dependencies = [ "num", ] +[[package]] +name = "fs-err" +version = "3.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b91aa448ca50d7e79433bdf3ee8d99215430d2ec02ade5aefab2a073a1822e8a" +dependencies = [ + "autocfg", + "tokio", +] + [[package]] name = "fs_extra" version = "1.3.0" @@ -1043,6 +1103,17 @@ version = "0.14.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e5274423e17b7c9fc20b6e7e208532f9b19825d82dfd615708b70edd83df41f1" +[[package]] +name = "hashbrown" +version = "0.15.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9229cfe53dfd69f0609a49f65461bd93001ea1ef889cd5529dd176593f5338a1" +dependencies = [ + "allocator-api2", + "equivalent", + "foldhash 0.1.5", +] + [[package]] name = "hashbrown" version = "0.17.1" @@ -1051,7 +1122,7 @@ checksum = "ed5909b6e89a2db4456e54cd5f673791d7eca6732202bbf2a9cc504fe2f9b84a" dependencies = [ "allocator-api2", "equivalent", - "foldhash", + "foldhash 0.2.0", ] [[package]] @@ -1522,6 +1593,15 @@ version = "0.4.33" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0ceec5bc11778974d1bcb055b18002eba7f4b3518b6a0081b3af5f21666da9ad" +[[package]] +name = "lru" +version = "0.12.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "234cf4f4a04dc1f57e24b96cc0cd600cf2af460d4161ac5ecdd0af8e1f3b2a38" +dependencies = [ + "hashbrown 0.15.5", +] + [[package]] name = "lru-slab" version = "0.1.2" @@ -2168,6 +2248,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0283386ce02abc0151e1761d08802dfe86c173b0b494af5cbc086574e453da06" dependencies = [ "aws-lc-rs", + "log", "once_cell", "ring", "rustls-pki-types", @@ -2188,6 +2269,15 @@ dependencies = [ "security-framework", ] +[[package]] +name = "rustls-pemfile" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dce314e5fee3f39953d46bb63bb8a46d40c2f8fb7cc5a3b6cab2bde9721d6e50" +dependencies = [ + "rustls-pki-types", +] + [[package]] name = "rustls-pki-types" version = "1.15.1" @@ -2850,6 +2940,16 @@ dependencies = [ "tracing-core", ] +[[package]] +name = "tracing-serde" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "704b1aeb7be0d0a84fc9828cae51dab5970fee5088f83d1dd7ee6f6246fc6ff1" +dependencies = [ + "serde", + "tracing-core", +] + [[package]] name = "tracing-subscriber" version = "0.3.23" @@ -2860,12 +2960,15 @@ dependencies = [ "nu-ansi-term", "once_cell", "regex-automata", + "serde", + "serde_json", "sharded-slab", "smallvec", "thread_local", "tracing", "tracing-core", "tracing-log", + "tracing-serde", ] [[package]] @@ -2949,6 +3052,17 @@ version = "0.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" +[[package]] +name = "uuid" +version = "1.24.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2cefc03fd367c0c6d4305de1b312cf00248c4114f4a0418ce6a6af769e3b0bd9" +dependencies = [ + "getrandom 0.4.3", + "js-sys", + "wasm-bindgen", +] + [[package]] name = "uuid-simd" version = "0.8.0" diff --git a/rust/Cargo.toml b/rust/Cargo.toml index 1699900d..083a157b 100644 --- a/rust/Cargo.toml +++ b/rust/Cargo.toml @@ -31,6 +31,8 @@ url = "2" regex = "1" schemars = "1" jsonschema = "0.49" +axum = "0.8" +rustls = "0.23" [profile.release] lto = "fat" diff --git a/rust/crates/adc-backend-api7/Cargo.toml b/rust/crates/adc-backend-api7/Cargo.toml index 0613153d..03d6b8c0 100644 --- a/rust/crates/adc-backend-api7/Cargo.toml +++ b/rust/crates/adc-backend-api7/Cargo.toml @@ -28,4 +28,4 @@ adc-backend-api7 = { path = ".", features = ["test-utils"] } adc-differ = { path = "../adc-differ" } tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "time"] } reqwest = { version = "0.12", default-features = false, features = ["rustls-tls", "json", "cookies"] } -axum = "0.8" +axum = { workspace = true } diff --git a/rust/crates/adc-backend-api7/src/backend.rs b/rust/crates/adc-backend-api7/src/backend.rs index 7c8b3ee7..8ace7e41 100644 --- a/rust/crates/adc-backend-api7/src/backend.rs +++ b/rust/crates/adc-backend-api7/src/backend.rs @@ -21,7 +21,7 @@ pub struct Backend { filter: ResourceFilter, version: OnceCell, default_value: OnceCell, - fetch_concurrency: usize, + concurrency: usize, } impl Backend { @@ -30,7 +30,7 @@ impl Backend { gateway_group_name: String, token: &str, filter: ResourceFilter, - fetch_concurrency: usize, + concurrency: usize, ) -> Self { let client = client.with_log_scope(vec!["API7".to_string()]); let gateway_group = GatewayGroupResolver::new(client.clone(), gateway_group_name, token); @@ -40,7 +40,7 @@ impl Backend { filter, version: OnceCell::new(), default_value: OnceCell::new(), - fetch_concurrency, + concurrency, } } @@ -104,7 +104,7 @@ impl adc_sdk::Backend for Backend { version, gateway_group_id, self.filter.clone(), - self.fetch_concurrency, + self.concurrency, ) .dump() .await diff --git a/rust/crates/adc-backend-apisix/Cargo.toml b/rust/crates/adc-backend-apisix/Cargo.toml index 4870a053..aa16e7bf 100644 --- a/rust/crates/adc-backend-apisix/Cargo.toml +++ b/rust/crates/adc-backend-apisix/Cargo.toml @@ -29,4 +29,4 @@ url = { workspace = true } [dev-dependencies] adc-backend-apisix = { path = ".", features = ["test-utils"] } tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net"] } -axum = "0.8" +axum = { workspace = true } diff --git a/rust/crates/adc-backend-apisix/src/backend.rs b/rust/crates/adc-backend-apisix/src/backend.rs index af9fce6e..64b9818b 100644 --- a/rust/crates/adc-backend-apisix/src/backend.rs +++ b/rust/crates/adc-backend-apisix/src/backend.rs @@ -19,16 +19,16 @@ pub struct Backend { client: HttpClient, filter: ResourceFilter, version: OnceCell, - fetch_concurrency: usize, + concurrency: usize, } impl Backend { - pub fn new(client: HttpClient, filter: ResourceFilter, fetch_concurrency: usize) -> Self { + pub fn new(client: HttpClient, filter: ResourceFilter, concurrency: usize) -> Self { Self { client: client.with_log_scope(vec!["APISIX".to_string()]), filter, version: OnceCell::new(), - fetch_concurrency, + concurrency, } } @@ -89,7 +89,7 @@ impl adc_sdk::Backend for Backend { self.client.clone(), version, self.filter.clone(), - self.fetch_concurrency, + self.concurrency, ) .dump() .await diff --git a/rust/crates/adc-backend-core/Cargo.toml b/rust/crates/adc-backend-core/Cargo.toml index 316c7ba7..42a32f7a 100644 --- a/rust/crates/adc-backend-core/Cargo.toml +++ b/rust/crates/adc-backend-core/Cargo.toml @@ -19,4 +19,4 @@ serde_json = { workspace = true } [dev-dependencies] tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "time"] } -axum = "0.8" +axum = { workspace = true } diff --git a/rust/crates/adc-backend-core/src/client.rs b/rust/crates/adc-backend-core/src/client.rs index 8c78645d..c7eee14c 100644 --- a/rust/crates/adc-backend-core/src/client.rs +++ b/rust/crates/adc-backend-core/src/client.rs @@ -58,37 +58,55 @@ pub struct HttpClient { /// (`"APISIX"`), not the generic `"ADC"`. Set via `with_log_scope` /// since `HttpClient` itself doesn't know which backend owns it. log_scope: Vec, + /// Applied per-request in [`HttpClient::request`], not just baked into + /// `inner` at build time — a caller-owned pooled client (see + /// [`HttpClient::with_shared_client`]) may have no client-level timeout + /// of its own, so this is the only way its callers still get one. + timeout: Option, } impl HttpClient { pub fn new(config: HttpClientConfig) -> Result { - let base_url = Url::parse(&config.server).map_err(|e| { - BackendError::Other(format!("invalid server URL {:?}: {e}", config.server).into()) - })?; - - let mut headers = HeaderMap::new(); - headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/json")); - let mut token = HeaderValue::from_str(&config.token) - .map_err(|e| BackendError::Other(format!("invalid token: {e}").into()))?; - token.set_sensitive(true); - headers.insert("X-API-KEY", token); - let default_headers = headers.clone(); - - let mut builder = reqwest::Client::builder().default_headers(headers); + let mut builder = reqwest::Client::builder(); if let Some(timeout) = config.timeout { builder = builder.timeout(timeout); } builder = config.tls.apply(builder)?; - let inner = builder.build().map_err(|e| { BackendError::Other(format!("failed to build HTTP client: {}", with_source(&e)).into()) })?; + Self::with_shared_client(inner, config.server, config.token, config.timeout) + } + + /// Like [`HttpClient::new`], but wraps an already-built `reqwest::Client` + /// — for callers pooling clients across many `HttpClient`s by TLS + /// material. That pooled client is shared across callers that may want + /// different timeouts, so `timeout` (if given) is applied per-request + /// here rather than relying on whatever the pooled client itself was + /// built with. + pub fn with_shared_client( + client: reqwest::Client, + server: String, + token: String, + timeout: Option, + ) -> Result { + let base_url = Url::parse(&server) + .map_err(|e| BackendError::Other(format!("invalid server URL {server:?}: {e}").into()))?; + + let mut headers = HeaderMap::new(); + headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/json")); + let mut token = HeaderValue::from_str(&token) + .map_err(|e| BackendError::Other(format!("invalid token: {e}").into()))?; + token.set_sensitive(true); + headers.insert("X-API-KEY", token); + Ok(Self { - inner, + inner: client, base_url, - default_headers, + default_headers: headers, log_scope: vec!["ADC".to_string()], + timeout, }) } @@ -102,6 +120,8 @@ impl HttpClient { /// Plain string concatenation, not `Url::join` — a root-anchored `path` /// via `join` would replace the base URL's path outright, dropping any /// prefix the server URL carries (e.g. behind a reverse-proxy). + /// `default_headers` are attached here, not baked into the underlying + /// client, since that client may be shared across tokens. pub fn request(&self, method: Method, path: &str) -> Result { let base = self.base_url.as_str().trim_end_matches('/'); let path = path.trim_start_matches('/'); @@ -109,7 +129,11 @@ impl HttpClient { let url = Url::parse(&combined).map_err(|e| { BackendError::Other(format!("invalid request path {path:?}: {e}").into()) })?; - Ok(self.inner.request(method, url)) + let mut builder = self.inner.request(method, url).headers(self.default_headers.clone()); + if let Some(timeout) = self.timeout { + builder = builder.timeout(timeout); + } + Ok(builder) } /// Classifies only transport-level failures into `BackendError::Transport`. @@ -135,8 +159,7 @@ impl HttpClient { let server_address = url.host_str().unwrap_or("").to_string(); let server_port = url.port_or_known_default().unwrap_or(0); - // `request.headers()` alone misses `default_headers` (merged only - // at send time); request-set headers win on conflict. + // Defensive fallback for a `RequestBuilder` not built via `request()`. let mut request_headers = self.default_headers.clone(); for (name, value) in request.headers() { request_headers.insert(name.clone(), value.clone()); diff --git a/rust/crates/adc-backend-core/src/tls.rs b/rust/crates/adc-backend-core/src/tls.rs index 79599d8c..ef0d3441 100644 --- a/rust/crates/adc-backend-core/src/tls.rs +++ b/rust/crates/adc-backend-core/src/tls.rs @@ -18,6 +18,14 @@ pub struct TlsConfig { } impl TlsConfig { + /// Builds a bare `reqwest::Client` from this config alone, with no + /// server/token/timeout attached. + pub fn build_client(&self) -> Result { + self.apply(reqwest::Client::builder())? + .build() + .map_err(|e| BackendError::Other(format!("failed to build HTTP client: {e}").into())) + } + pub(crate) fn apply(&self, mut builder: reqwest::ClientBuilder) -> Result { if self.skip_verify { builder = builder.danger_accept_invalid_certs(true); diff --git a/rust/crates/adc-cli/Cargo.toml b/rust/crates/adc-cli/Cargo.toml index 00ef4897..f2701f0f 100644 --- a/rust/crates/adc-cli/Cargo.toml +++ b/rust/crates/adc-cli/Cargo.toml @@ -15,18 +15,33 @@ adc-differ = { path = "../adc-differ" } adc-backend-core = { path = "../adc-backend-core" } adc-backend-apisix = { path = "../adc-backend-apisix" } adc-backend-api7 = { path = "../adc-backend-api7" } +adc-backend-apisix-standalone = { path = "../adc-backend-apisix-standalone" } adc-converter-openapi = { path = "../adc-converter-openapi" } clap = { version = "4", features = ["derive", "env"] } humantime = "2" +serde = { workspace = true } serde_json = { workspace = true } serde_yaml_ng = "0.10" thiserror = { workspace = true } -tokio = { workspace = true, features = ["rt-multi-thread", "macros", "fs"] } +tokio = { workspace = true, features = ["rt-multi-thread", "macros", "fs", "net", "signal", "sync", "process"] } glob = "0.3" -tracing-subscriber = { version = "0.3", features = ["env-filter"] } +tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] } dotenvy = "0.15" tracing-indicatif = "0.3" tracing = { workspace = true } indicatif = "0.18" chrono = { version = "0.4", default-features = false, features = ["clock"] } http = "1" +axum = { workspace = true } +axum-server = { version = "0.7", features = ["tls-rustls"] } +rustls = { workspace = true } +rustls-pemfile = "2" +reqwest = { version = "0.12", default-features = false, features = ["rustls-tls", "json"] } +uuid = { version = "1", features = ["v4"] } +lru = "0.12" +url = { workspace = true } +http-body-util = "0.1" + +[dev-dependencies] +tower = { version = "0.5", features = ["util"] } +libc = "0.2" diff --git a/rust/crates/adc-cli/src/cli.rs b/rust/crates/adc-cli/src/cli.rs index eab4c4e8..770fc755 100644 --- a/rust/crates/adc-cli/src/cli.rs +++ b/rust/crates/adc-cli/src/cli.rs @@ -39,7 +39,41 @@ pub enum Command { IngressSync, /// run the local ingress server #[command(hide = true)] - IngressServer, + IngressServer(IngressServerArgs), +} + +#[derive(Args, Debug)] +pub struct IngressServerArgs { + /// listen address of the ADC server, in the form scheme://host:port + /// (http/https/unix are supported) + #[arg(long, default_value = "http://127.0.0.1:3000", value_parser = parse_listen_url)] + pub listen: url::Url, + + /// status listen port (exposes GET /healthz/ready) + #[arg(long, default_value_t = 3001, value_parser = clap::value_parser!(u16).range(1..=65535))] + pub listen_status: u16, + + /// path to the CA certificate used to verify client certificates (enables mTLS) + #[arg(long, value_parser = existing_file)] + pub ca_cert_file: Option, + + /// path to the TLS server certificate (required for https:// listen addresses) + #[arg(long, value_parser = existing_file)] + pub tls_cert_file: Option, + + /// path to the TLS server key (required for https:// listen addresses) + #[arg(long, value_parser = existing_file)] + pub tls_key_file: Option, +} + +fn parse_listen_url(raw: &str) -> Result { + let url = url::Url::parse(raw).map_err(|e| e.to_string())?; + match url.scheme() { + "http" | "https" | "unix" => Ok(url), + other => Err(format!( + "unsupported --listen scheme \"{other}\": expected http, https, or unix" + )), + } } #[derive(ValueEnum, Debug, Clone, Copy, PartialEq, Eq)] @@ -141,16 +175,16 @@ pub struct BackendArgs { pub request_concurrent: usize, /// path to the CA certificate to verify the backend - #[arg(long, env = "ADC_CA_CERT_FILE", value_parser = existing_file)] - pub ca_cert_file: Option, + #[arg(long = "ca-cert-file", env = "ADC_CA_CERT_FILE", value_parser = read_file_bytes)] + pub ca_cert_pem: Option>, /// path to the mutual TLS client certificate to verify the backend - #[arg(long, env = "ADC_TLS_CLIENT_CERT_FILE", value_parser = existing_file, requires = "tls_client_key_file")] - pub tls_client_cert_file: Option, + #[arg(long = "tls-client-cert-file", env = "ADC_TLS_CLIENT_CERT_FILE", value_parser = read_file_bytes, requires = "tls_client_key_pem")] + pub tls_client_cert_pem: Option>, /// path to the mutual TLS client key to verify the backend - #[arg(long, env = "ADC_TLS_CLIENT_KEY_FILE", value_parser = existing_file, requires = "tls_client_cert_file")] - pub tls_client_key_file: Option, + #[arg(long = "tls-client-key-file", env = "ADC_TLS_CLIENT_KEY_FILE", value_parser = read_file_bytes, requires = "tls_client_cert_pem")] + pub tls_client_key_pem: Option>, /// disable the verification of the backend TLS certificate #[arg(long, env = "ADC_TLS_SKIP_VERIFY")] @@ -259,3 +293,9 @@ fn existing_file(raw: &str) -> Result { } Ok(path) } + +/// Reads the file's content at parse time — so downstream code (`BackendSpec`) +/// only ever deals with bytes already in memory, never a path. +fn read_file_bytes(raw: &str) -> Result, String> { + std::fs::read(raw).map_err(|e| format!("{raw}: {e}")) +} diff --git a/rust/crates/adc-cli/src/main.rs b/rust/crates/adc-cli/src/main.rs index 0af9746f..565d71b6 100644 --- a/rust/crates/adc-cli/src/main.rs +++ b/rust/crates/adc-cli/src/main.rs @@ -4,14 +4,23 @@ mod error; mod logging; mod pipeline; mod progress; +mod server; use std::collections::{HashMap, HashSet}; use adc_sdk::{BackendSyncOptions, Event, EventType, ResourceType}; use clap::Parser; -use cli::{BackendArgs, Cli, Command, ConvertArgs, ConvertFormat, DiffArgs, DumpArgs, LintArgs, SyncArgs, ValidateArgs}; +use cli::{ + BackendArgs, Cli, Command, ConvertArgs, ConvertFormat, DiffArgs, DumpArgs, LintArgs, SyncArgs, + ValidateArgs, +}; use error::CliError; +use pipeline::BackendSpec; + +pub(crate) fn install_crypto_provider() { + let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); +} #[tokio::main] async fn main() { @@ -20,8 +29,19 @@ async fn main() { let _ = dotenvy::from_path(".env"); let cli = Cli::parse(); - logging::init(cli.verbose); - progress::set_verbose(cli.verbose); + + install_crypto_provider(); + + // The global tracing subscriber can only be installed once, so pick + // which logging setup to use before dispatching below. + let is_ingress_server = matches!(cli.command, Command::IngressServer(_)); + match &cli.command { + Command::IngressServer(_) => server::logging::init(), + _ => { + logging::init(cli.verbose); + progress::set_verbose(cli.verbose); + } + } let result = match cli.command { Command::Ping(args) => cmd_ping(args).await, @@ -32,10 +52,13 @@ async fn main() { Command::Validate(args) => cmd_validate(args).await, Command::Convert(args) => cmd_convert(args).await, Command::IngressSync => Err(CliError::msg("adc ingress-sync: not yet implemented")), - Command::IngressServer => Err(CliError::msg("adc ingress-server: not yet implemented")), + Command::IngressServer(args) => server::run(args).await, }; match result { + // The ingress-server daemon has its own "Stopping..." log line on + // shutdown — skip the one-shot-command "All is well" line here. + Ok(()) if is_ingress_server => {} Ok(()) => progress::finish_ok(), Err(CliError::AlreadyReported) => std::process::exit(1), Err(err) => { @@ -46,7 +69,7 @@ async fn main() { } async fn cmd_ping(args: BackendArgs) -> Result<(), CliError> { - let backend = pipeline::init_backend(&args).await?; + let backend = pipeline::init_backend(BackendSpec::try_from(&args)?, None)?; progress::stage("Connecting to backend...", backend.ping()).await?; println!( "Connected to the \"{}\" backend successfully!", @@ -56,7 +79,7 @@ async fn cmd_ping(args: BackendArgs) -> Result<(), CliError> { } async fn cmd_dump(args: DumpArgs) -> Result<(), CliError> { - let backend = pipeline::init_backend(&args.backend).await?; + let backend = pipeline::init_backend(BackendSpec::try_from(&args.backend)?, None)?; let (include, exclude) = pipeline::resource_type_sets(&args.backend); let label_selector = pipeline::label_selector_map(&args.backend)?; let remote = progress::stage( @@ -79,7 +102,7 @@ async fn cmd_dump(args: DumpArgs) -> Result<(), CliError> { } async fn cmd_diff(args: DiffArgs) -> Result<(), CliError> { - let backend = pipeline::init_backend(&args.backend).await?; + let backend = pipeline::init_backend(BackendSpec::try_from(&args.backend)?, None)?; let (include, exclude) = pipeline::resource_type_sets(&args.backend); let label_selector = pipeline::label_selector_map(&args.backend)?; @@ -113,7 +136,7 @@ async fn cmd_diff(args: DiffArgs) -> Result<(), CliError> { } async fn cmd_sync(args: SyncArgs) -> Result<(), CliError> { - let backend = pipeline::init_backend(&args.backend).await?; + let backend = pipeline::init_backend(BackendSpec::try_from(&args.backend)?, None)?; let (include, exclude) = pipeline::resource_type_sets(&args.backend); let label_selector = pipeline::label_selector_map(&args.backend)?; @@ -179,7 +202,11 @@ async fn cmd_sync(args: SyncArgs) -> Result<(), CliError> { // blame for the failure — report which server instead. None => println!( "[FAILED] sync{}", - result.server.as_deref().map(|server| format!(" to {server}")).unwrap_or_default() + result + .server + .as_deref() + .map(|server| format!(" to {server}")) + .unwrap_or_default() ), } if let Some(err) = &result.error { @@ -237,7 +264,7 @@ async fn cmd_lint(args: LintArgs) -> Result<(), CliError> { } async fn cmd_validate(args: ValidateArgs) -> Result<(), CliError> { - let backend = pipeline::init_backend(&args.backend).await?; + let backend = pipeline::init_backend(BackendSpec::try_from(&args.backend)?, None)?; let (include, exclude) = pipeline::resource_type_sets(&args.backend); let label_selector = pipeline::label_selector_map(&args.backend)?; diff --git a/rust/crates/adc-cli/src/pipeline.rs b/rust/crates/adc-cli/src/pipeline.rs index 5bf22a5c..3c5bcd53 100644 --- a/rust/crates/adc-cli/src/pipeline.rs +++ b/rust/crates/adc-cli/src/pipeline.rs @@ -6,75 +6,137 @@ use std::collections::{HashMap, HashSet}; use std::path::PathBuf; +use std::time::Duration; use adc_backend_core::{HttpClient, HttpClientConfig, ResourceFilter, TlsConfig}; use adc_differ::DifferV4; use adc_sdk::resources::Configuration; use adc_sdk::{Backend, Converter, Event, InternalConfiguration, ResourceType}; -use crate::cli::{BackendArgs, BackendKind}; +use crate::cli::BackendArgs; use crate::config; use crate::error::CliError; -pub async fn init_backend(args: &BackendArgs) -> Result, CliError> { - let filter = resource_filter(args)?; - match args.backend { - BackendKind::Apisix => { - let (client, _token) = build_client(args).await?; - Ok(Box::new(adc_backend_apisix::Backend::new( +/// Everything [`init_backend`] needs to pick and construct a backend, +/// independent of where the caller sourced it from (CLI args or an +/// ingress-server request body). +pub struct BackendSpec { + pub kind: String, + pub servers: Vec, + pub tokens: Vec, + pub gateway_group: Option, + pub filter: ResourceFilter, + pub concurrency: usize, + pub cache_key: String, + pub bypass_cache: bool, + pub timeout: Option, + pub tls: TlsConfig, +} + +/// CLI args resolve into a `BackendSpec` with no I/O — `--ca-cert-file` and +/// friends already read their file's bytes at argument-parse time (see +/// `cli::read_file_bytes`), so this is a plain, synchronous field mapping. +/// Fallible only because `label_selector_map` rejects `managed-by` as a +/// selector key. +impl TryFrom<&BackendArgs> for BackendSpec { + type Error = CliError; + + fn try_from(args: &BackendArgs) -> Result { + let (include, exclude) = resource_type_sets(args); + let label_selector = label_selector_map(args)?; + let token = args.token.as_deref().ok_or_else(|| { + CliError::msg("a backend token is required: pass --token or set ADC_TOKEN") + })?; + + Ok(BackendSpec { + kind: args.backend.as_str().to_string(), + servers: args.server.split(',').map(str::to_string).collect(), + tokens: token.split(',').map(str::to_string).collect(), + gateway_group: Some(args.gateway_group.clone()), + filter: ResourceFilter { + include, + exclude, + label_selector, + }, + concurrency: args.request_concurrent, + cache_key: "default".to_string(), + bypass_cache: false, + timeout: Some(args.timeout), + tls: TlsConfig { + ca_cert_pem: args.ca_cert_pem.clone(), + client_cert_pem: args.tls_client_cert_pem.clone(), + client_key_pem: args.tls_client_key_pem.clone(), + skip_verify: args.tls_skip_verify, + }, + }) + } +} + +/// `shared_client`: `None` builds a fresh `HttpClient` from `spec.tls` (the +/// one-shot CLI); `Some` reuses an already-built, pooled `reqwest::Client` +/// (the ingress-server daemon). Only `apisix`/`api7ee` look at it — +/// `apisix-standalone` always builds its own clients from `spec.tls`. +pub fn init_backend( + spec: BackendSpec, + shared_client: Option, +) -> Result, CliError> { + let server = spec.servers.first().cloned().unwrap_or_default(); + let token = spec.tokens.first().cloned().unwrap_or_default(); + + match spec.kind.as_str() { + "api7ee" => { + let client = http_client(shared_client, &server, &token, spec.timeout, &spec.tls)?; + let gateway_group = spec.gateway_group.unwrap_or_else(|| "default".to_string()); + Ok(Box::new(adc_backend_api7::Backend::new( client, - filter, - args.request_concurrent, + gateway_group, + &token, + spec.filter, + spec.concurrency, ))) } - BackendKind::Api7Ee => { - let (client, token) = build_client(args).await?; - Ok(Box::new(adc_backend_api7::Backend::new( + "apisix-standalone" => Ok(Box::new(adc_backend_apisix_standalone::Backend::new( + adc_backend_apisix_standalone::BackendOptions { + servers: spec.servers, + tokens: spec.tokens, + cache_key: spec.cache_key, + bypass_cache: spec.bypass_cache, + timeout: spec.timeout, + tls: spec.tls, + }, + )?)), + "" | "apisix" => { + let client = http_client(shared_client, &server, &token, spec.timeout, &spec.tls)?; + Ok(Box::new(adc_backend_apisix::Backend::new( client, - args.gateway_group.clone(), - &token, - filter, - args.request_concurrent, + spec.filter, + spec.concurrency, ))) } - BackendKind::ApisixStandalone => Err(CliError::msg(format!( - "backend \"{}\" is not yet implemented (only \"apisix\"/\"api7ee\" are supported so far)", - args.backend.as_str() + other => Err(CliError::msg(format!( + "unrecognized backend kind: \"{other}\"" ))), } } -/// Shared by every backend: the `X-API-KEY`/TLS-configured `HttpClient` -/// every one of them wraps. Returns the raw token alongside it — `api7ee` -/// needs it separately (to recognize an `a7adm-` admin token, which skips -/// gateway_group resolution entirely), not just baked into the client's -/// headers. -async fn build_client(args: &BackendArgs) -> Result<(HttpClient, String), CliError> { - let token = args.token.clone().ok_or_else(|| { - CliError::msg("a backend token is required: pass --token or set ADC_TOKEN") - })?; - let ca_cert_pem = read_optional(&args.ca_cert_file).await?; - let client_cert_pem = read_optional(&args.tls_client_cert_file).await?; - let client_key_pem = read_optional(&args.tls_client_key_file).await?; - let client = HttpClient::new(HttpClientConfig { - server: args.server.clone(), - token: token.clone(), - timeout: Some(args.timeout), - tls: TlsConfig { - ca_cert_pem, - client_cert_pem, - client_key_pem, - skip_verify: args.tls_skip_verify, - }, - })?; - Ok((client, token)) -} - -async fn read_optional(path: &Option) -> Result>, CliError> { - match path { - Some(path) => Ok(Some(tokio::fs::read(path).await?)), - None => Ok(None), - } +fn http_client( + shared_client: Option, + server: &str, + token: &str, + timeout: Option, + tls: &TlsConfig, +) -> Result { + Ok(match shared_client { + None => HttpClient::new(HttpClientConfig { + server: server.to_string(), + token: token.to_string(), + timeout, + tls: tls.clone(), + })?, + Some(client) => { + HttpClient::with_shared_client(client, server.to_string(), token.to_string(), timeout)? + } + }) } /// Resource types nested under a service/consumer rather than a top-level @@ -123,7 +185,11 @@ pub fn resource_type_sets(args: &BackendArgs) -> (HashSet, HashSet // only ever check top-level types, so a nested one left in `include` // would make every real top-level type look excluded instead of doing // nothing. - let is_unfilterable = |rt: &ResourceType| UNFILTERABLE_RESOURCE_TYPES.iter().any(|(_, unfilterable)| unfilterable == rt); + let is_unfilterable = |rt: &ResourceType| { + UNFILTERABLE_RESOURCE_TYPES + .iter() + .any(|(_, unfilterable)| unfilterable == rt) + }; include.retain(|rt| !is_unfilterable(rt)); exclude.retain(|rt| !is_unfilterable(rt)); (include, exclude) @@ -160,21 +226,6 @@ fn parse_label_selector(entries: &[String]) -> Result, C .collect() } -/// The filter a backend applies at fetch time: skipping whole resource -/// types the request never needed, and (where the admin API supports it) -/// asking the server itself to narrow results by label. This is an -/// optimization only — `config::filter_resource_types`/`filter_by_labels` -/// still run afterward and are what actually guarantee the result matches. -fn resource_filter(args: &BackendArgs) -> Result { - let (include, exclude) = resource_type_sets(args); - let label_selector = label_selector_map(args)?; - Ok(ResourceFilter { - include, - exclude, - label_selector, - }) -} - /// Loads, merges, and structurally parses the local configuration file(s), /// then (unless `lint` is `false`, i.e. `--no-lint`) runs semantic /// validation on top. Deserializing into `Configuration` is the @@ -209,11 +260,10 @@ pub async fn load_local( Ok(configuration) } -/// Collects every lint violation into one multi-line message — mirrors the -/// TS CLI wrapping `z.prettifyError`'s multi-issue output into a single -/// thrown `Error`. +/// Collects every lint violation into one multi-line message. fn format_lint_issues(issues: &[adc_sdk::lint::LintIssue]) -> String { - let mut message = "Lint configuration\nThe following errors were found in configuration:\n".to_string(); + let mut message = + "Lint configuration\nThe following errors were found in configuration:\n".to_string(); for issue in issues { message.push_str(&format!(" - {issue}\n")); } @@ -235,7 +285,12 @@ pub async fn convert_openapi(files: &[PathBuf]) -> Result>(), vec!["svc-a", "svc-b"]); + assert_eq!( + merged.iter().map(|s| s.name.as_str()).collect::>(), + vec!["svc-a", "svc-b"] + ); } #[test] @@ -369,13 +428,19 @@ mod tests { let err = merge_openapi_services(per_file).unwrap_err(); let message = err.to_string(); assert!(message.contains("b.yaml"), "{message}"); - assert!(message.contains("a.yaml"), "{message}: should name the file that first produced this service"); + assert!( + message.contains("a.yaml"), + "{message}: should name the file that first produced this service" + ); assert!(message.contains("shared"), "{message}"); } #[test] fn merge_openapi_services_rejects_a_duplicate_name_within_the_same_file() { - let per_file = vec![(PathBuf::from("a.yaml"), vec![service("shared"), service("shared")])]; + let per_file = vec![( + PathBuf::from("a.yaml"), + vec![service("shared"), service("shared")], + )]; assert!(merge_openapi_services(per_file).is_err()); } @@ -390,14 +455,39 @@ mod tests { exclude_resource_type: vec![], timeout: Duration::from_secs(10), request_concurrent: 10, - ca_cert_file: None, - tls_client_cert_file: None, - tls_client_key_file: None, + ca_cert_pem: None, + tls_client_cert_pem: None, + tls_client_key_pem: None, tls_skip_verify: false, managed_by_label: true, } } + #[test] + fn backend_spec_splits_comma_joined_servers_and_tokens() { + let mut args = backend_args(vec![]); + args.server = "http://a:9180,http://b:9180".to_string(); + args.token = Some("tok-a,tok-b".to_string()); + let spec = BackendSpec::try_from(&args).unwrap(); + assert_eq!(spec.servers, vec!["http://a:9180", "http://b:9180"]); + assert_eq!(spec.tokens, vec!["tok-a", "tok-b"]); + } + + #[test] + fn backend_spec_accepts_a_single_server_and_token() { + let mut args = backend_args(vec![]); + args.token = Some("tok".to_string()); + let spec = BackendSpec::try_from(&args).unwrap(); + assert_eq!(spec.servers, vec!["http://localhost:9180"]); + assert_eq!(spec.tokens, vec!["tok"]); + } + + #[test] + fn backend_spec_rejects_a_missing_token() { + let args = backend_args(vec![]); + assert!(BackendSpec::try_from(&args).is_err()); + } + #[test] fn rejects_managed_by_as_a_selector_key_regardless_of_the_value_supplied() { let args = backend_args(vec!["managed-by=custom".to_string()]); @@ -431,7 +521,9 @@ mod tests { fn captured_warning(set: &HashSet) -> String { let buffer = SharedBuffer::default(); let writer = buffer.clone(); - let subscriber = tracing_subscriber::fmt().with_writer(move || writer.clone()).finish(); + let subscriber = tracing_subscriber::fmt() + .with_writer(move || writer.clone()) + .finish(); tracing::subscriber::with_default(subscriber, || { warn_on_unfilterable_resource_types("--include-resource-type", set); }); @@ -447,7 +539,10 @@ mod tests { #[test] fn no_warning_when_only_top_level_types_are_named() { - let output = captured_warning(&HashSet::from([ResourceType::Service, ResourceType::Consumer])); + let output = captured_warning(&HashSet::from([ + ResourceType::Service, + ResourceType::Consumer, + ])); assert!(output.is_empty(), "{output}"); } @@ -465,6 +560,9 @@ mod tests { let mut args = backend_args(vec![]); args.include_resource_type = vec![ResourceTypeArg::Route, ResourceTypeArg::Upstream]; let (include, _) = resource_type_sets(&args); - assert!(include.is_empty(), "an include set with nothing but unfilterable types must behave like no --include-resource-type was passed at all"); + assert!( + include.is_empty(), + "an include set with nothing but unfilterable types must behave like no --include-resource-type was passed at all" + ); } } diff --git a/rust/crates/adc-cli/src/server/agent_pool.rs b/rust/crates/adc-cli/src/server/agent_pool.rs new file mode 100644 index 00000000..07bcf5c7 --- /dev/null +++ b/rust/crates/adc-cli/src/server/agent_pool.rs @@ -0,0 +1,118 @@ +//! A process-wide, LRU-bounded pool of `reqwest::Client`s keyed by TLS +//! material, so requests to the same backend gateway share a connection +//! pool. Eviction is plain `Arc` drop — no manual refcounting needed. + +use std::num::NonZeroUsize; +use std::sync::{Arc, LazyLock, Mutex}; + +use adc_backend_core::TlsConfig; +use adc_sdk::BackendError; +use lru::LruCache; + +/// Distinguishes one pooled client from another — `Hash`/`Eq` let it double as the pool key. +#[derive(Debug, Clone, Default, PartialEq, Eq, Hash)] +pub struct TlsMaterial { + pub skip_verify: bool, + pub ca_cert: Option, + pub client_cert: Option, + pub client_key: Option, +} + +const DEFAULT_MAX_ENTRIES: usize = 16; + +fn env_max_entries() -> NonZeroUsize { + std::env::var("ADC_INGRESS_TLS_AGENT_POOL_MAX") + .ok() + .and_then(|value| value.parse::().ok()) + .and_then(NonZeroUsize::new) + .unwrap_or(NonZeroUsize::new(DEFAULT_MAX_ENTRIES).expect("16 is nonzero")) +} + +pub struct AgentPool { + entries: Mutex>>, +} + +impl AgentPool { + pub fn with_capacity(max_entries: usize) -> Self { + let capacity = NonZeroUsize::new(max_entries).unwrap_or(NonZeroUsize::new(1).expect("1 is nonzero")); + Self { entries: Mutex::new(LruCache::new(capacity)) } + } + + pub fn global() -> &'static AgentPool { + static GLOBAL: LazyLock = LazyLock::new(|| AgentPool::with_capacity(env_max_entries().get())); + &GLOBAL + } + + pub fn get_client(&self, tls: &TlsMaterial) -> Result, BackendError> { + let mut entries = self.entries.lock().expect("agent pool mutex poisoned"); + if let Some(client) = entries.get(tls) { + return Ok(client.clone()); + } + let client = Arc::new(build_client(tls)?); + entries.put(tls.clone(), client.clone()); + Ok(client) + } +} + +fn build_client(tls: &TlsMaterial) -> Result { + TlsConfig { + ca_cert_pem: tls.ca_cert.clone().map(String::into_bytes), + client_cert_pem: tls.client_cert.clone().map(String::into_bytes), + client_key_pem: tls.client_key.clone().map(String::into_bytes), + skip_verify: tls.skip_verify, + } + .build_client() +} + +pub fn get_client(tls: &TlsMaterial) -> Result, BackendError> { + AgentPool::global().get_client(tls) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn material(ca_cert: &str) -> TlsMaterial { + TlsMaterial { ca_cert: Some(ca_cert.to_string()), ..Default::default() } + } + + #[test] + fn reuses_the_same_client_for_identical_material() { + let pool = AgentPool::with_capacity(16); + let a = pool.get_client(&material("shared")).unwrap(); + let b = pool.get_client(&material("shared")).unwrap(); + assert!(Arc::ptr_eq(&a, &b)); + } + + #[test] + fn different_material_gets_isolated_clients() { + let pool = AgentPool::with_capacity(16); + let a = pool.get_client(&material("a")).unwrap(); + let b = pool.get_client(&material("b")).unwrap(); + assert!(!Arc::ptr_eq(&a, &b)); + } + + #[test] + fn no_tls_material_and_some_tls_material_are_isolated_too() { + let pool = AgentPool::with_capacity(16); + let default_client = pool.get_client(&TlsMaterial::default()).unwrap(); + let skip_verify = pool.get_client(&TlsMaterial { skip_verify: true, ..Default::default() }).unwrap(); + assert!(!Arc::ptr_eq(&default_client, &skip_verify)); + } + + #[test] + fn evicts_the_least_recently_used_entry_not_merely_the_first_inserted() { + let pool = AgentPool::with_capacity(2); + let a = pool.get_client(&material("a")).unwrap(); + let b = pool.get_client(&material("b")).unwrap(); + let _ = pool.get_client(&material("a")).unwrap(); // refreshes "a"'s recency + + let _c = pool.get_client(&material("c")).unwrap(); // exceeds capacity(2), evicts "b" + + let a_again = pool.get_client(&material("a")).unwrap(); + assert!(Arc::ptr_eq(&a, &a_again), "\"a\" should not have been evicted"); + + let b_again = pool.get_client(&material("b")).unwrap(); + assert!(!Arc::ptr_eq(&b, &b_again), "\"b\" should have been evicted and rebuilt"); + } +} diff --git a/rust/crates/adc-cli/src/server/backend.rs b/rust/crates/adc-cli/src/server/backend.rs new file mode 100644 index 00000000..28a1ef3b --- /dev/null +++ b/rust/crates/adc-cli/src/server/backend.rs @@ -0,0 +1,176 @@ +//! Translates a request's `opts` into `pipeline::BackendSpec` and delegates +//! to `pipeline::init_backend` — the only server-specific pieces are +//! sourcing TLS material from inline PEM (not files) and getting the +//! `HttpClient` from the shared pool (see `agent_pool`) instead of building +//! a fresh one. + +use adc_backend_core::{ResourceFilter, TlsConfig}; +use adc_sdk::Backend; + +use super::agent_pool::{self, TlsMaterial}; +use super::schema::Opts; +use crate::error::CliError; +use crate::pipeline::{self, BackendSpec}; + +/// No fallible parts, unlike `BackendSpec`'s `TryFrom<&BackendArgs>` +/// (which rejects `managed-by` as a label selector) — nothing here can fail. +impl From<&Opts> for BackendSpec { + fn from(opts: &Opts) -> Self { + let (include, exclude) = opts.resource_type_sets(); + BackendSpec { + kind: opts.backend.clone(), + servers: opts.server.as_list(), + tokens: opts.token.split(',').map(str::to_string).collect(), + gateway_group: opts.gateway_group.clone(), + filter: ResourceFilter { + include, + exclude, + label_selector: opts.label_selector_or_default(), + }, + concurrency: opts.request_concurrent, + cache_key: opts.cache_key.clone(), + bypass_cache: opts.bypass_cache, + timeout: Some(std::time::Duration::from_millis(opts.timeout)), + tls: TlsConfig { + ca_cert_pem: opts.ca_cert.clone().map(String::into_bytes), + client_cert_pem: opts.tls_client_cert.clone().map(String::into_bytes), + client_key_pem: opts.tls_client_key.clone().map(String::into_bytes), + skip_verify: opts.tls_skip_verify, + }, + } + } +} + +pub fn build_backend(opts: &Opts) -> Result, CliError> { + let tls_material = TlsMaterial { + skip_verify: opts.tls_skip_verify, + ca_cert: opts.ca_cert.clone(), + client_cert: opts.tls_client_cert.clone(), + client_key: opts.tls_client_key.clone(), + }; + let shared_client = agent_pool::get_client(&tls_material)?; + pipeline::init_backend(opts.into(), Some((*shared_client).clone())) +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use super::super::schema::ServerAddr; + use super::*; + + fn tls_asset(name: &str) -> std::path::PathBuf { + std::path::PathBuf::from(concat!(env!("CARGO_MANIFEST_DIR"), "/tests/assets/tls")) + .join(name) + } + + fn opts(server: &str, ca_cert: Option) -> Opts { + Opts { + backend: "apisix".to_string(), + server: ServerAddr::Single(server.to_string()), + token: "token".to_string(), + lint: true, + include_resource_type: None, + exclude_resource_type: None, + label_selector: None, + cache_key: "default".to_string(), + bypass_cache: false, + gateway_group: None, + request_concurrent: 10, + timeout: 30_000, + tls_skip_verify: false, + ca_cert, + tls_client_cert: None, + tls_client_key: None, + } + } + + /// A bare HTTPS stub answering every request with `200 {}`. + async fn spawn_https_stub() -> (std::net::SocketAddr, tokio::task::JoinHandle<()>) { + let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); + + let cert = super::super::load_certs(&tls_asset("server.cer")).unwrap(); + let key = super::super::load_key(&tls_asset("server.key")).unwrap(); + let config = rustls::ServerConfig::builder() + .with_no_client_auth() + .with_single_cert(cert, key) + .unwrap(); + let tls_config = axum_server::tls_rustls::RustlsConfig::from_config(Arc::new(config)); + + let app = axum::Router::new().fallback(axum::routing::any(|| async { + (axum::http::StatusCode::OK, "{}") + })); + let addr: std::net::SocketAddr = "127.0.0.1:0".parse().unwrap(); + let listener = std::net::TcpListener::bind(addr).unwrap(); + let addr = listener.local_addr().unwrap(); + let handle = axum_server::Handle::new(); + let server_handle = handle.clone(); + let task = tokio::spawn(async move { + axum_server::from_tcp_rustls(listener, tls_config) + .handle(server_handle) + .serve(app.into_make_service()) + .await + .unwrap(); + }); + handle.listening().await; + (addr, task) + } + + #[tokio::test] + async fn ping_fails_certificate_verification_without_a_ca_cert() { + let (addr, task) = spawn_https_stub().await; + let backend = build_backend(&opts(&format!("https://{addr}"), None)).unwrap(); + let error = backend.ping().await.unwrap_err(); + assert!( + matches!(error, adc_sdk::BackendError::Transport(_)), + "expected a Transport error (certificate verification is a transport-level failure), got {error:?}" + ); + assert!( + error.to_string().to_lowercase().contains("certificate") + || error.to_string().to_lowercase().contains("unknownissuer"), + "{error}" + ); + task.abort(); + } + + #[tokio::test] + async fn ping_succeeds_once_the_signing_ca_is_provided() { + let (addr, task) = spawn_https_stub().await; + let ca = std::fs::read_to_string(tls_asset("ca.cer")).unwrap(); + let backend = build_backend(&opts(&format!("https://{addr}"), Some(ca))).unwrap(); + backend + .ping() + .await + .expect("certificate verification should succeed with the correct CA"); + task.abort(); + } + + #[test] + fn build_backend_constructs_an_api7ee_backend() { + let mut o = opts("http://127.0.0.1:9180", None); + o.backend = "api7ee".to_string(); + o.gateway_group = Some("prod".to_string()); + assert!(build_backend(&o).is_ok()); + } + + #[test] + fn build_backend_constructs_an_apisix_standalone_backend_from_multiple_servers_and_tokens() { + let mut o = opts("http://127.0.0.1:9180", None); + o.backend = "apisix-standalone".to_string(); + o.server = ServerAddr::Multiple(vec![ + "http://127.0.0.1:9180".to_string(), + "http://127.0.0.1:9181".to_string(), + ]); + o.token = "t1,t2".to_string(); + o.cache_key = "test-standalone-key".to_string(); + assert!(build_backend(&o).is_ok()); + } + + #[test] + fn build_backend_rejects_an_apisix_standalone_backend_with_no_servers() { + let mut o = opts("http://127.0.0.1:9180", None); + o.backend = "apisix-standalone".to_string(); + o.server = ServerAddr::Multiple(vec![]); + assert!(build_backend(&o).is_err()); + } +} diff --git a/rust/crates/adc-cli/src/server/logging.rs b/rust/crates/adc-cli/src/server/logging.rs new file mode 100644 index 00000000..531f0b54 --- /dev/null +++ b/rust/crates/adc-cli/src/server/logging.rs @@ -0,0 +1,108 @@ +//! JSON structured logging for the ingress-server daemon. Level is +//! controlled by `ADC_INGRESS_LOG_LEVEL` (default `info`), not `--verbose`. + +use axum::body::Body; +use axum::extract::Request; +use axum::middleware::Next; +use axum::response::{IntoResponse, Response}; +use serde_json::Value; + +pub fn init() { + let level = std::env::var("ADC_INGRESS_LOG_LEVEL").unwrap_or_else(|_| "info".to_string()); + let env_filter = + tracing_subscriber::EnvFilter::try_new(&level).unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")); + tracing_subscriber::fmt().json().with_env_filter(env_filter).init(); +} + +pub async fn request_logger(request: Request, next: Next) -> Response { + let request_id = uuid::Uuid::new_v4().to_string(); + let method = request.method().clone(); + let path = request.uri().path().to_string(); + tracing::info!(request_id = %request_id, "{method} {path}"); + + let (parts, body) = request.into_parts(); + let bytes = match axum::body::to_bytes(body, 100 * 1024 * 1024).await { + Ok(bytes) => bytes, + Err(error) => { + tracing::warn!(request_id = %request_id, %error, "failed to read request body"); + let status = if std::error::Error::source(&error).is_some_and(|source| source.is::()) { + axum::http::StatusCode::PAYLOAD_TOO_LARGE + } else { + axum::http::StatusCode::BAD_REQUEST + }; + return status.into_response(); + } + }; + if !bytes.is_empty() { + tracing::debug!(request_id = %request_id, request_body = %redacted_body_text(&bytes)); + } + + let request = Request::from_parts(parts, Body::from(bytes)); + next.run(request).await +} + +fn redacted_body_text(bytes: &[u8]) -> String { + match serde_json::from_slice::(bytes) { + Ok(value) => redact_request_body(&value).to_string(), + // Non-JSON bytes could still contain a raw token/key substring — + // logging them verbatim would defeat the redaction above. + Err(_) => format!("", bytes.len()), + } +} + +/// Never let a raw mTLS private key or backend token reach a debug log. +pub fn redact_request_body(body: &Value) -> Value { + let mut redacted = body.clone(); + for pointer in ["/task/opts/tlsClientKey", "/task/opts/token"] { + if let Some(field) = redacted.pointer_mut(pointer) { + *field = Value::String("***".to_string()); + } + } + redacted +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + #[test] + fn redacts_tls_client_key_while_preserving_other_fields() { + let body = json!({"task": {"opts": {"backend": "apisix", "tlsClientKey": "SECRET", "tlsClientCert": "cert"}, "config": {}}}); + let redacted = redact_request_body(&body); + assert_eq!( + redacted, + json!({"task": {"opts": {"backend": "apisix", "tlsClientKey": "***", "tlsClientCert": "cert"}, "config": {}}}) + ); + } + + #[test] + fn redacts_the_backend_token() { + let body = json!({"task": {"opts": {"backend": "apisix", "token": "SECRET"}, "config": {}}}); + let redacted = redact_request_body(&body); + assert_eq!( + redacted, + json!({"task": {"opts": {"backend": "apisix", "token": "***"}, "config": {}}}) + ); + } + + #[test] + fn returns_the_body_unchanged_when_tls_client_key_is_absent() { + let body = json!({"task": {"opts": {"backend": "apisix"}, "config": {}}}); + assert_eq!(redact_request_body(&body), body); + } + + #[test] + fn does_not_panic_on_malformed_bodies() { + for body in [ + json!({"task": {"opts": 1}}), + json!({"task": {"opts": "not-an-object"}}), + json!({"task": {"opts": null}}), + json!({"task": {}}), + json!({}), + Value::Null, + ] { + let _ = redact_request_body(&body); + } + } +} diff --git a/rust/crates/adc-cli/src/server/mod.rs b/rust/crates/adc-cli/src/server/mod.rs new file mode 100644 index 00000000..fdf57816 --- /dev/null +++ b/rust/crates/adc-cli/src/server/mod.rs @@ -0,0 +1,574 @@ +//! The ingress-server daemon: a long-running HTTP(S)/Unix-socket sidecar +//! exposing `PUT /sync`/`PUT /validate` on an adc listener, +//! `GET /healthz/ready` on a separate status listener. + +pub mod agent_pool; +mod backend; +pub mod logging; +mod schema; +mod sync; +mod validate; + +use std::net::{IpAddr, Ipv4Addr, SocketAddr, ToSocketAddrs}; +use std::os::unix::fs::{FileTypeExt, PermissionsExt}; +use std::path::Path; +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::time::Duration; + +use axum::extract::State; +use axum::http::StatusCode; +use axum::response::{IntoResponse, Response}; +use axum::routing::{get, put}; +use axum::{Json, Router}; +use rustls::pki_types::{CertificateDer, PrivateKeyDer}; +use serde_json::Value; +use tokio::sync::watch; + +use crate::cli::IngressServerArgs; +use crate::error::CliError; + +const MAX_BODY_BYTES: usize = 100 * 1024 * 1024; + +/// Bound on how long the HTTPS listener waits for in-flight (e.g. +/// keep-alive) connections to finish once a shutdown starts — without this, +/// an idle keep-alive connection could block process exit indefinitely. +const HTTPS_SHUTDOWN_DEADLINE: Duration = Duration::from_secs(10); + +/// Checks the current value first so an already-sent `true` isn't missed. +async fn wait_for_shutdown(mut rx: watch::Receiver) { + if *rx.borrow() { + return; + } + let _ = rx.changed().await; +} + +fn adc_router() -> Router { + Router::new() + .route("/sync", put(sync::sync_handler)) + .route("/validate", put(validate::validate_handler)) + .layer(axum::middleware::from_fn(logging::request_logger)) + .layer(axum::extract::DefaultBodyLimit::max(MAX_BODY_BYTES)) +} + +fn status_router(ready: Arc) -> Router { + Router::new() + .route("/healthz/ready", get(healthz)) + .with_state(ready) +} + +pub async fn run(args: IngressServerArgs) -> Result<(), CliError> { + // `signal()` registers synchronously, unlike `ctrl_c()` (an async fn + // that only registers on first poll) — called first to close the race + // where an early SIGINT hits before a spawned ctrl_c task gets polled. + let mut sigint = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::interrupt()) + .map_err(|e| CliError::msg(format!("failed to install SIGINT handler: {e}")))?; + + let app = adc_router(); + + let ready = Arc::new(AtomicBool::new(false)); + let status_app = status_router(ready.clone()); + + let status_addr = SocketAddr::from((IpAddr::V4(Ipv4Addr::UNSPECIFIED), args.listen_status)); + let status_listener = tokio::net::TcpListener::bind(status_addr) + .await + .map_err(|e| { + CliError::msg(format!( + "failed to bind status listener on {status_addr}: {e}" + )) + })?; + + let (shutdown_tx, shutdown_rx) = watch::channel(false); + + let status_shutdown = shutdown_rx.clone(); + let status_server = async move { + axum::serve(status_listener, status_app) + .with_graceful_shutdown(wait_for_shutdown(status_shutdown)) + .await + .map_err(|e| CliError::msg(format!("status server error: {e}"))) + }; + + tracing::info!( + "ADC server is running on: {}", + display_listen_address(&args.listen) + ); + let adc_server = serve_adc(&args, app, shutdown_rx, ready); + + tokio::spawn(async move { + sigint.recv().await; + tracing::info!("Stopping, see you next time!"); + let _ = shutdown_tx.send(true); + }); + + tokio::try_join!(adc_server, status_server)?; + Ok(()) +} + +async fn healthz(State(ready): State>) -> Response { + if ready.load(Ordering::Acquire) { + (StatusCode::OK, "ok").into_response() + } else { + (StatusCode::SERVICE_UNAVAILABLE, "not ready").into_response() + } +} + +async fn serve_adc( + args: &IngressServerArgs, + app: Router, + shutdown: watch::Receiver, + ready: Arc, +) -> Result<(), CliError> { + match args.listen.scheme() { + "unix" => serve_unix(args, app, shutdown, ready).await, + "https" => serve_https(args, app, shutdown, ready).await, + _ => serve_http(args, app, shutdown, ready).await, + } +} + +async fn serve_http( + args: &IngressServerArgs, + app: Router, + shutdown: watch::Receiver, + ready: Arc, +) -> Result<(), CliError> { + let addr = tcp_addr(&args.listen)?; + let listener = tokio::net::TcpListener::bind(addr) + .await + .map_err(|e| CliError::msg(format!("failed to bind {addr}: {e}")))?; + ready.store(true, Ordering::Release); + axum::serve(listener, app) + .with_graceful_shutdown(wait_for_shutdown(shutdown)) + .await + .map_err(|e| CliError::msg(format!("server error: {e}"))) +} + +/// An existing stale socket file is removed before binding, and the fresh +/// one gets `0o660` permissions once bound. The socket file itself is +/// removed again on shutdown so a normal restart doesn't depend on this +/// stale-file cleanup happening next time. +async fn serve_unix( + args: &IngressServerArgs, + app: Router, + shutdown: watch::Receiver, + ready: Arc, +) -> Result<(), CliError> { + let path = args.listen.path(); + match std::fs::symlink_metadata(path) { + Ok(metadata) if metadata.file_type().is_socket() => { + std::fs::remove_file(path) + .map_err(|e| CliError::msg(format!("failed to remove stale socket {path}: {e}")))?; + } + Ok(_) => { + return Err(CliError::msg(format!( + "refusing to bind unix socket: {path} already exists and is not a socket" + ))); + } + Err(e) if e.kind() == std::io::ErrorKind::NotFound => {} + Err(e) => return Err(CliError::msg(format!("failed to inspect {path}: {e}"))), + } + let listener = tokio::net::UnixListener::bind(path) + .map_err(|e| CliError::msg(format!("failed to bind unix socket {path}: {e}")))?; + std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o660)) + .map_err(|e| CliError::msg(format!("failed to chmod unix socket {path}: {e}")))?; + ready.store(true, Ordering::Release); + let result = axum::serve(listener, app) + .with_graceful_shutdown(wait_for_shutdown(shutdown)) + .await + .map_err(|e| CliError::msg(format!("server error: {e}"))); + std::fs::remove_file(path).ok(); + result +} + +async fn serve_https( + args: &IngressServerArgs, + app: Router, + shutdown: watch::Receiver, + ready: Arc, +) -> Result<(), CliError> { + let cert_path = args + .tls_cert_file + .as_deref() + .ok_or_else(|| CliError::msg("--tls-cert-file is required when --listen uses https"))?; + let key_path = args + .tls_key_file + .as_deref() + .ok_or_else(|| CliError::msg("--tls-key-file is required when --listen uses https"))?; + let server_config = build_rustls_config(cert_path, key_path, args.ca_cert_file.as_deref())?; + let addr = tcp_addr(&args.listen)?; + + let handle = axum_server::Handle::new(); + let handle_for_shutdown = handle.clone(); + tokio::spawn(async move { + wait_for_shutdown(shutdown).await; + handle_for_shutdown.graceful_shutdown(Some(HTTPS_SHUTDOWN_DEADLINE)); + }); + + let handle_for_ready = handle.clone(); + tokio::spawn(async move { + if handle_for_ready.listening().await.is_some() { + ready.store(true, Ordering::Release); + } + }); + + let tls_config = axum_server::tls_rustls::RustlsConfig::from_config(Arc::new(server_config)); + axum_server::tls_rustls::bind_rustls(addr, tls_config) + .handle(handle) + .serve(app.into_make_service()) + .await + .map_err(|e| CliError::msg(format!("server error: {e}"))) +} + +/// `requestCert: true, rejectUnauthorized: true` when a CA is given (mTLS); +/// plain server-only TLS otherwise. +fn build_rustls_config( + cert_path: &Path, + key_path: &Path, + ca_path: Option<&Path>, +) -> Result { + crate::install_crypto_provider(); + + let certs = load_certs(cert_path)?; + let key = load_key(key_path)?; + + let builder = rustls::ServerConfig::builder(); + let builder = match ca_path { + Some(ca_path) => { + let mut roots = rustls::RootCertStore::empty(); + for cert in load_certs(ca_path)? { + roots + .add(cert) + .map_err(|e| CliError::msg(format!("invalid CA certificate: {e}")))?; + } + let verifier = rustls::server::WebPkiClientVerifier::builder(Arc::new(roots)) + .build() + .map_err(|e| CliError::msg(format!("invalid CA certificate: {e}")))?; + builder.with_client_cert_verifier(verifier) + } + None => builder.with_no_client_auth(), + }; + builder + .with_single_cert(certs, key) + .map_err(|e| CliError::msg(format!("invalid TLS certificate/key: {e}"))) +} + +fn load_certs(path: &Path) -> Result>, CliError> { + let file = + std::fs::File::open(path).map_err(|e| CliError::msg(format!("{}: {e}", path.display())))?; + let mut reader = std::io::BufReader::new(file); + rustls_pemfile::certs(&mut reader) + .collect::, _>>() + .map_err(|e| CliError::msg(format!("{}: {e}", path.display()))) +} + +fn load_key(path: &Path) -> Result, CliError> { + let file = + std::fs::File::open(path).map_err(|e| CliError::msg(format!("{}: {e}", path.display())))?; + let mut reader = std::io::BufReader::new(file); + rustls_pemfile::private_key(&mut reader) + .map_err(|e| CliError::msg(format!("{}: {e}", path.display())))? + .ok_or_else(|| CliError::msg(format!("{}: no private key found", path.display()))) +} + +fn tcp_addr(listen: &url::Url) -> Result { + let host = listen + .host_str() + .ok_or_else(|| CliError::msg("--listen must include a host"))?; + let port = listen + .port_or_known_default() + .ok_or_else(|| CliError::msg("--listen must include a port"))?; + (host, port) + .to_socket_addrs() + .ok() + .and_then(|mut addrs| addrs.next()) + .ok_or_else(|| CliError::msg(format!("could not resolve listen address {host}:{port}"))) +} + +fn display_listen_address(listen: &url::Url) -> String { + if listen.scheme() == "unix" { + listen.path().to_string() + } else { + listen.as_str().trim_end_matches('/').to_string() + } +} + +fn bad_request(body: Value) -> Response { + (StatusCode::BAD_REQUEST, Json(body)).into_response() +} + +fn internal_error(body: Value) -> Response { + (StatusCode::INTERNAL_SERVER_ERROR, Json(body)).into_response() +} + +#[cfg(test)] +mod tests { + use axum::body::Body; + use axum::http::Request; + use tower::ServiceExt; + + use super::*; + + fn tls_asset(name: &str) -> std::path::PathBuf { + std::path::PathBuf::from(concat!(env!("CARGO_MANIFEST_DIR"), "/tests/assets/tls")) + .join(name) + } + + async fn send(router: Router, method: &str, path: &str, body: &str) -> (StatusCode, Value) { + let request = Request::builder() + .method(method) + .uri(path) + .header("content-type", "application/json") + .body(Body::from(body.to_string())) + .unwrap(); + let response = router.oneshot(request).await.unwrap(); + let status = response.status(); + let bytes = axum::body::to_bytes(response.into_body(), usize::MAX) + .await + .unwrap(); + let json = serde_json::from_slice(&bytes).unwrap_or(Value::Null); + (status, json) + } + + #[tokio::test] + async fn healthz_ready_returns_ok_once_marked_ready() { + let router = status_router(Arc::new(AtomicBool::new(true))); + let (status, _) = send(router, "GET", "/healthz/ready", "").await; + assert_eq!(status, StatusCode::OK); + } + + #[tokio::test] + async fn healthz_ready_returns_service_unavailable_before_ready() { + let router = status_router(Arc::new(AtomicBool::new(false))); + let (status, _) = send(router, "GET", "/healthz/ready", "").await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); + } + + #[tokio::test] + async fn sync_rejects_malformed_input_with_400() { + let (status, body) = send(adc_router(), "PUT", "/sync", r#"{"not":"valid"}"#).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert!(body["message"].is_string(), "{body}"); + } + + #[tokio::test] + async fn sync_rejects_unpaired_tls_client_material_and_names_the_field() { + let body = serde_json::json!({ + "task": { + "opts": { + "backend": "apisix", "server": "http://1.1.1.1:9180", "token": "t", "cacheKey": "default", + "tlsClientCert": "-----BEGIN CERTIFICATE-----\nx", + }, + "config": {}, + } + }) + .to_string(); + let (status, json) = send(adc_router(), "PUT", "/sync", &body).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + let errors = json["errors"].as_array().unwrap(); + assert!( + errors + .iter() + .any(|e| e["path"] == serde_json::json!(["tlsClientKey"])), + "{json}" + ); + } + + #[tokio::test] + async fn sync_against_an_unreachable_backend_returns_500() { + let body = serde_json::json!({ + "task": {"opts": {"backend": "apisix", "server": "http://127.0.0.1:1", "token": "t", "cacheKey": "default"}, "config": {}} + }) + .to_string(); + let (status, _) = send(adc_router(), "PUT", "/sync", &body).await; + assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR); + } + + #[tokio::test] + async fn sync_rejects_a_configuration_that_fails_lint() { + let body = serde_json::json!({ + "task": { + "opts": {"backend": "apisix", "server": "http://1.1.1.1:9180", "token": "t", "cacheKey": "default"}, + "config": {"services": [{"name": ""}]}, + } + }) + .to_string(); + let (status, json) = send(adc_router(), "PUT", "/sync", &body).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert!(json["message"].as_str().unwrap().contains("Lint"), "{json}"); + } + + #[tokio::test] + async fn validate_reports_source_input_on_a_malformed_body() { + let (status, json) = send(adc_router(), "PUT", "/validate", "{}").await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert_eq!(json["source"], "input"); + assert_eq!(json["success"], false); + } + + #[tokio::test] + async fn validate_reports_source_lint_on_a_lint_failure() { + let body = serde_json::json!({ + "task": { + "opts": {"backend": "apisix", "server": "http://1.1.1.1:9180", "token": "t", "cacheKey": "default"}, + "config": {"services": [{"name": ""}]}, + } + }) + .to_string(); + let (status, json) = send(adc_router(), "PUT", "/validate", &body).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert_eq!(json["source"], "lint"); + } + + #[tokio::test] + async fn http_listener_serves_real_requests_end_to_end() { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let (shutdown_tx, shutdown_rx) = watch::channel(false); + let server = tokio::spawn(async move { + axum::serve(listener, adc_router()) + .with_graceful_shutdown(wait_for_shutdown(shutdown_rx)) + .await + .unwrap(); + }); + + let response = reqwest::Client::new() + .put(format!("http://{addr}/sync")) + .json(&serde_json::json!({"not": "valid"})) + .send() + .await + .unwrap(); + assert_eq!(response.status(), reqwest::StatusCode::BAD_REQUEST); + + let _ = shutdown_tx.send(true); + server.await.unwrap(); + } + + /// A real `https://` listener requiring a client certificate. + async fn spawn_https_mtls_listener(app: Router) -> (SocketAddr, tokio::task::JoinHandle<()>) { + let config = build_rustls_config( + &tls_asset("server.cer"), + &tls_asset("server.key"), + Some(&tls_asset("ca.cer")), + ) + .unwrap(); + let tls_config = axum_server::tls_rustls::RustlsConfig::from_config(Arc::new(config)); + let handle = axum_server::Handle::new(); + let server_handle = handle.clone(); + let task = tokio::spawn(async move { + axum_server::tls_rustls::bind_rustls("127.0.0.1:0".parse().unwrap(), tls_config) + .handle(server_handle) + .serve(app.into_make_service()) + .await + .unwrap(); + }); + let addr = handle.listening().await.expect("server should have bound"); + (addr, task) + } + + #[tokio::test] + async fn https_listener_rejects_a_connection_without_a_client_certificate() { + let (addr, task) = spawn_https_mtls_listener(adc_router()).await; + let client = reqwest::Client::builder() + .danger_accept_invalid_certs(true) + .build() + .unwrap(); + let result = client + .put(format!("https://{addr}/sync")) + .body("{}") + .send() + .await; + assert!( + result.is_err(), + "expected the TLS handshake to fail without a client certificate" + ); + task.abort(); + } + + #[tokio::test] + async fn https_listener_accepts_a_connection_with_a_valid_client_certificate() { + let (addr, task) = spawn_https_mtls_listener(adc_router()).await; + + let mut identity_pem = std::fs::read_to_string(tls_asset("client.cer")).unwrap(); + identity_pem.push('\n'); + identity_pem.push_str(&std::fs::read_to_string(tls_asset("client.key")).unwrap()); + let identity = reqwest::Identity::from_pem(identity_pem.as_bytes()).unwrap(); + let client = reqwest::Client::builder() + .danger_accept_invalid_certs(true) + .identity(identity) + .build() + .unwrap(); + + // A real HTTP response (not a TLS error) proves the handshake succeeded. + let response = client + .put(format!("https://{addr}/sync")) + .header("content-type", "application/json") + .body("{}") + .send() + .await + .unwrap(); + assert_eq!(response.status(), reqwest::StatusCode::BAD_REQUEST); + task.abort(); + } + + #[tokio::test] + async fn unix_listener_removes_a_stale_socket_and_sets_0o660_permissions() { + let dir = std::env::temp_dir().join(format!("adc-server-test-{}", uuid::Uuid::new_v4())); + std::fs::create_dir_all(&dir).unwrap(); + let path = dir.join("adc.sock"); + // A real leftover socket from a crashed previous run, not just any file. + drop(std::os::unix::net::UnixListener::bind(&path).unwrap()); + + let args = IngressServerArgs { + listen: url::Url::parse(&format!("unix://{}", path.display())).unwrap(), + listen_status: 0, + ca_cert_file: None, + tls_cert_file: None, + tls_key_file: None, + }; + let (shutdown_tx, shutdown_rx) = watch::channel(false); + let ready = Arc::new(AtomicBool::new(false)); + let server = tokio::spawn({ + let ready = ready.clone(); + async move { serve_unix(&args, adc_router(), shutdown_rx, ready).await } + }); + + // `serve_unix` only flips `ready` after `set_permissions` succeeds — + // polling the socket's mere existence races the chmod that follows it. + let mut became_ready = false; + for _ in 0..100 { + if ready.load(Ordering::Acquire) { + became_ready = true; + break; + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + assert!(became_ready, "server never became ready"); + let mode = std::fs::metadata(&path).unwrap().permissions().mode() & 0o777; + assert_eq!(mode, 0o660); + + let _ = shutdown_tx.send(true); + server.await.unwrap().unwrap(); + std::fs::remove_dir_all(&dir).ok(); + } + + #[test] + fn build_rustls_config_accepts_a_matching_cert_and_key_with_a_ca() { + let config = build_rustls_config( + &tls_asset("server.cer"), + &tls_asset("server.key"), + Some(&tls_asset("ca.cer")), + ); + assert!(config.is_ok(), "{config:?}"); + } + + #[test] + fn build_rustls_config_accepts_a_matching_cert_and_key_without_a_ca() { + let config = build_rustls_config(&tls_asset("server.cer"), &tls_asset("server.key"), None); + assert!(config.is_ok(), "{config:?}"); + } + + #[test] + fn build_rustls_config_rejects_a_key_that_does_not_match_the_certificate() { + let config = build_rustls_config(&tls_asset("server.cer"), &tls_asset("client.key"), None); + assert!(config.is_err()); + } +} diff --git a/rust/crates/adc-cli/src/server/schema.rs b/rust/crates/adc-cli/src/server/schema.rs new file mode 100644 index 00000000..f300dc9c --- /dev/null +++ b/rust/crates/adc-cli/src/server/schema.rs @@ -0,0 +1,452 @@ +//! Request bodies for `PUT /sync` and `PUT /validate`. + +use std::collections::HashMap; + +use serde::{Deserialize, Serialize}; +use serde_json::Value; + +#[derive(Debug, Deserialize)] +pub struct SyncInput { + pub task: SyncTask, +} + +#[derive(Debug, Deserialize)] +pub struct SyncTask { + pub opts: Opts, + pub config: Value, +} + +#[derive(Debug, Deserialize)] +pub struct ValidateInput { + pub task: ValidateTask, +} + +#[derive(Debug, Deserialize)] +pub struct ValidateTask { + pub opts: Opts, + pub config: Value, +} + +/// Shared by both endpoints — `/validate` just ignores `bypass_cache`. +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct Opts { + pub backend: String, + pub server: ServerAddr, + pub token: String, + #[serde(default = "default_true")] + pub lint: bool, + pub include_resource_type: Option>, + pub exclude_resource_type: Option>, + pub label_selector: Option>, + pub cache_key: String, + #[serde(default)] + pub bypass_cache: bool, + pub gateway_group: Option, + #[serde(default = "default_request_concurrent")] + pub request_concurrent: usize, + #[serde(default = "default_timeout_ms")] + pub timeout: u64, + + // TLS/mTLS to the backend gateway — raw PEM, not a file path. + #[serde(default)] + pub tls_skip_verify: bool, + pub ca_cert: Option, + pub tls_client_cert: Option, + pub tls_client_key: Option, +} + +fn default_true() -> bool { + true +} + +fn default_request_concurrent() -> usize { + 10 +} + +fn default_timeout_ms() -> u64 { + 30_000 +} + +impl Opts { + pub fn resource_type_sets( + &self, + ) -> ( + std::collections::HashSet, + std::collections::HashSet, + ) { + let include = self + .include_resource_type + .iter() + .flatten() + .map(|t| (*t).into()) + .collect(); + let exclude = self + .exclude_resource_type + .iter() + .flatten() + .map(|t| (*t).into()) + .collect(); + (include, exclude) + } + + pub fn label_selector_or_default(&self) -> HashMap { + self.label_selector.clone().unwrap_or_default() + } +} + +/// `z.union([z.url(), z.array(z.url())])` — a single backend takes one +/// server URL, `apisix-standalone` addresses a cluster of them. +#[derive(Debug, Deserialize)] +#[serde(untagged)] +pub enum ServerAddr { + Single(String), + Multiple(Vec), +} + +impl ServerAddr { + pub fn as_list(&self) -> Vec { + match self { + ServerAddr::Single(server) => vec![server.clone()], + ServerAddr::Multiple(servers) => servers.clone(), + } + } +} + +/// Mirrors `cli::ResourceTypeArg`, kept separate so `adc_sdk::ResourceType` stays derive-free. +#[derive(Debug, Deserialize, Clone, Copy, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum ServerResourceType { + Route, + Service, + Upstream, + Ssl, + GlobalRule, + PluginConfig, + PluginMetadata, + Consumer, + ConsumerGroup, + ConsumerCredential, + StreamRoute, +} + +impl From for adc_sdk::ResourceType { + fn from(value: ServerResourceType) -> Self { + match value { + ServerResourceType::Route => adc_sdk::ResourceType::Route, + ServerResourceType::Service => adc_sdk::ResourceType::Service, + ServerResourceType::Upstream => adc_sdk::ResourceType::Upstream, + ServerResourceType::Ssl => adc_sdk::ResourceType::Ssl, + ServerResourceType::GlobalRule => adc_sdk::ResourceType::GlobalRule, + ServerResourceType::PluginConfig => adc_sdk::ResourceType::PluginConfig, + ServerResourceType::PluginMetadata => adc_sdk::ResourceType::PluginMetadata, + ServerResourceType::Consumer => adc_sdk::ResourceType::Consumer, + ServerResourceType::ConsumerGroup => adc_sdk::ResourceType::ConsumerGroup, + ServerResourceType::ConsumerCredential => adc_sdk::ResourceType::ConsumerCredential, + ServerResourceType::StreamRoute => adc_sdk::ResourceType::StreamRoute, + } + } +} + +/// One input-validation failure — `path` names the offending field. +#[derive(Debug, Serialize, PartialEq, Eq)] +pub struct ValidationIssue { + pub path: Vec, + pub message: String, +} + +impl ValidationIssue { + fn new(field: &str, message: impl Into) -> Self { + Self { + path: vec![field.to_string()], + message: message.into(), + } + } +} + +/// cert/key must be provided together; any of the three, if given, must be PEM. +pub fn validate_tls_material(opts: &Opts) -> Vec { + let mut issues = Vec::new(); + if opts.tls_client_cert.is_some() != opts.tls_client_key.is_some() { + issues.push(ValidationIssue::new( + "tlsClientKey", + "tlsClientCert and tlsClientKey must be provided together", + )); + } + for (field, value, noun, is_valid) in [ + ( + "caCert", + &opts.ca_cert, + "certificate", + is_valid_pem_certificate as fn(&str) -> bool, + ), + ( + "tlsClientCert", + &opts.tls_client_cert, + "certificate", + is_valid_pem_certificate, + ), + ( + "tlsClientKey", + &opts.tls_client_key, + "key", + is_valid_pem_private_key, + ), + ] { + if let Some(value) = value + && !is_valid(value) + { + issues.push(ValidationIssue::new( + field, + format!("{field} does not look like a PEM-encoded {noun}"), + )); + } + } + issues +} + +/// Parses the decoded DER via `RootCertStore::add`, not just the PEM +/// framing — garbage base64 between valid `-----BEGIN/END CERTIFICATE-----` +/// markers would otherwise pass and only fail later, deep inside backend +/// construction. +fn is_valid_pem_certificate(value: &str) -> bool { + let mut reader = std::io::BufReader::new(value.as_bytes()); + let Ok(certs) = rustls_pemfile::certs(&mut reader).collect::, _>>() else { + return false; + }; + if certs.is_empty() { + return false; + } + let mut store = rustls::RootCertStore::empty(); + certs.into_iter().all(|cert| store.add(cert).is_ok()) +} + +/// Goes through the installed `CryptoProvider`, not a specific backend +/// module, so this stays agnostic to which one `server::run` installs. +fn is_valid_pem_private_key(value: &str) -> bool { + let mut reader = std::io::BufReader::new(value.as_bytes()); + let Ok(Some(key)) = rustls_pemfile::private_key(&mut reader) else { + return false; + }; + let Some(provider) = rustls::crypto::CryptoProvider::get_default() else { + return false; + }; + provider.key_provider.load_private_key(key).is_ok() +} + +#[cfg(test)] +mod tests { + use super::*; + + fn opts(overrides: impl FnOnce(&mut Opts)) -> Opts { + let mut opts = Opts { + backend: "apisix".to_string(), + server: ServerAddr::Single("http://127.0.0.1:9180".to_string()), + token: "token".to_string(), + lint: true, + include_resource_type: None, + exclude_resource_type: None, + label_selector: None, + cache_key: "default".to_string(), + bypass_cache: false, + gateway_group: None, + request_concurrent: 10, + timeout: 30_000, + tls_skip_verify: false, + ca_cert: None, + tls_client_cert: None, + tls_client_key: None, + }; + overrides(&mut opts); + opts + } + + #[test] + fn deserializes_a_single_server_string() { + let input: SyncInput = serde_json::from_value(serde_json::json!({ + "task": { + "opts": {"backend": "apisix", "server": "http://a:9180", "token": "t", "cacheKey": "default"}, + "config": {} + } + })) + .unwrap(); + assert_eq!(input.task.opts.server.as_list(), vec!["http://a:9180"]); + } + + #[test] + fn deserializes_an_array_of_servers() { + let input: SyncInput = serde_json::from_value(serde_json::json!({ + "task": { + "opts": {"backend": "apisix-standalone", "server": ["http://a:9180", "http://b:9180"], "token": "t", "cacheKey": "default"}, + "config": {} + } + })) + .unwrap(); + assert_eq!( + input.task.opts.server.as_list(), + vec!["http://a:9180", "http://b:9180"] + ); + } + + #[test] + fn lint_defaults_to_true_and_bypass_cache_to_false() { + let input: SyncInput = serde_json::from_value(serde_json::json!({ + "task": { + "opts": {"backend": "apisix", "server": "http://a:9180", "token": "t", "cacheKey": "default"}, + "config": {} + } + })) + .unwrap(); + assert!(input.task.opts.lint); + assert!(!input.task.opts.bypass_cache); + } + + #[test] + fn resource_type_is_case_matched_snake_case() { + let input: SyncInput = serde_json::from_value(serde_json::json!({ + "task": { + "opts": { + "backend": "apisix", "server": "http://a:9180", "token": "t", "cacheKey": "default", + "includeResourceType": ["stream_route", "consumer_credential"] + }, + "config": {} + } + })) + .unwrap(); + assert_eq!( + input.task.opts.include_resource_type.unwrap(), + vec![ + ServerResourceType::StreamRoute, + ServerResourceType::ConsumerCredential + ] + ); + } + + #[test] + fn a_lone_tls_client_cert_without_a_key_is_rejected() { + let opts = + opts(|o| o.tls_client_cert = Some("-----BEGIN CERTIFICATE-----\n...".to_string())); + let issues = validate_tls_material(&opts); + assert!( + issues.iter().any(|i| i.path == vec!["tlsClientKey"]), + "{issues:?}" + ); + } + + #[test] + fn a_lone_tls_client_key_without_a_cert_is_rejected() { + let opts = + opts(|o| o.tls_client_key = Some("-----BEGIN PRIVATE KEY-----\n...".to_string())); + let issues = validate_tls_material(&opts); + assert!( + issues.iter().any(|i| i.path == vec!["tlsClientKey"]), + "{issues:?}" + ); + } + + // A real self-signed EC cert/key pair (generated once via `openssl req + // -x509 -newkey ec ...`, not fetched at test time) — validation parses + // this for real, not just checking for a `-----BEGIN` prefix. + const CERT_PEM: &str = "-----BEGIN CERTIFICATE-----\n\ +MIIBcjCCARmgAwIBAgIUWp+abBNKuUPdUIeouYDaDgHPIO4wCgYIKoZIzj0EAwIw\n\ +DzENMAsGA1UEAwwEdGVzdDAeFw0yNjA4MTgwMzIwMDJaFw0yNjA4MTkwMzIwMDJa\n\ +MA8xDTALBgNVBAMMBHRlc3QwWTATBgcqhkjOPQIBBggqhkjOPQMBBwNCAARVo3/X\n\ +uhOYfghuoLbag2VJvGofvgPYtXcdh4oFCmXB1MOupxSI3DqCFvMJc/QeH92Nz/qW\n\ +vLW7TEWRCo2/Bay1o1MwUTAdBgNVHQ4EFgQUy80qZFI7+wryg4UyeI+YsHfSqgow\n\ +HwYDVR0jBBgwFoAUy80qZFI7+wryg4UyeI+YsHfSqgowDwYDVR0TAQH/BAUwAwEB\n\ +/zAKBggqhkjOPQQDAgNHADBEAiB+ddl9S2GSo8/NF37M47JI1HtxOzQQTizSoAQd\n\ +tx5+SQIgKX3ASSnC8rrNGSFda+y79MOudxia/iQouBhv8Fb/hnE=\n\ +-----END CERTIFICATE-----\n"; + const KEY_PEM: &str = "-----BEGIN PRIVATE KEY-----\n\ +MIGHAgEAMBMGByqGSM49AgEGCCqGSM49AwEHBG0wawIBAQQgO3caXwd/kpykMTTw\n\ ++IoxA9NXadu2yhvQXw/rxkgjhZChRANCAARVo3/XuhOYfghuoLbag2VJvGofvgPY\n\ +tXcdh4oFCmXB1MOupxSI3DqCFvMJc/QeH92Nz/qWvLW7TEWRCo2/Bay1\n\ +-----END PRIVATE KEY-----\n"; + + #[test] + fn a_paired_cert_and_key_is_accepted_at_the_pairing_check() { + let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); + let opts = opts(|o| { + o.tls_client_cert = Some(CERT_PEM.to_string()); + o.tls_client_key = Some(KEY_PEM.to_string()); + }); + let issues = validate_tls_material(&opts); + assert!(issues.is_empty(), "{issues:?}"); + } + + #[test] + fn a_lone_tls_client_key_missing_its_cert_is_flagged_on_tls_client_key_only() { + // A lone tlsClientKey should be reported once, on the pairing + // check — not duplicated by the per-field PEM-validity check, + // since the key itself is valid PEM. + let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); + let opts = opts(|o| o.tls_client_key = Some(KEY_PEM.to_string())); + let issues = validate_tls_material(&opts); + assert_eq!( + issues, + vec![ValidationIssue::new( + "tlsClientKey", + "tlsClientCert and tlsClientKey must be provided together", + )], + "{issues:?}" + ); + } + + #[test] + fn an_incomplete_pem_certificate_is_rejected() { + let opts = opts(|o| o.ca_cert = Some("-----BEGIN CERTIFICATE-----\n...".to_string())); + let issues = validate_tls_material(&opts); + assert!( + issues.iter().any(|i| i.path == vec!["caCert"]), + "{issues:?}" + ); + } + + #[test] + fn a_well_formed_pem_wrapping_garbage_der_is_rejected() { + // Valid PEM framing, but the base64 payload isn't a real X.509 + // certificate — must fail at the DER-parsing check, not just the + // PEM-extraction one. + let garbage_cert = "-----BEGIN CERTIFICATE-----\n\ +dGhpcyBpcyBub3QgYSB2YWxpZCB4NTA5IGNlcnRpZmljYXRlLCBqdXN0IHNvbWUgcGFkZGluZyBieXRlcyB0byBtYWtlIGl0IGxvbmcgZW5vdWdo\n\ +-----END CERTIFICATE-----\n"; + let opts = opts(|o| o.ca_cert = Some(garbage_cert.to_string())); + let issues = validate_tls_material(&opts); + assert!( + issues.iter().any(|i| i.path == vec!["caCert"]), + "{issues:?}" + ); + } + + #[test] + fn a_well_formed_pem_wrapping_garbage_der_key_is_rejected() { + let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); + let garbage_key = "-----BEGIN PRIVATE KEY-----\n\ +dGhpcyBpcyBub3QgYSB2YWxpZCBwcml2YXRlIGtleSBkZXIsIGp1c3QgcGFkZGluZyBieXRlcyB0byBtYWtlIGl0IGxvbmcgZW5vdWdoIHRvIGxvb2sgcmVhbA==\n\ +-----END PRIVATE KEY-----\n"; + let opts = opts(|o| { + o.tls_client_cert = Some(CERT_PEM.to_string()); + o.tls_client_key = Some(garbage_key.to_string()); + }); + let issues = validate_tls_material(&opts); + assert!( + issues.iter().any(|i| i.path == vec!["tlsClientKey"]), + "{issues:?}" + ); + } + + #[test] + fn a_non_pem_ca_cert_is_rejected() { + let opts = opts(|o| o.ca_cert = Some("not-a-pem".to_string())); + let issues = validate_tls_material(&opts); + assert!( + issues.iter().any(|i| i.path == vec!["caCert"]), + "{issues:?}" + ); + } + + #[test] + fn no_tls_material_at_all_is_fine() { + assert!(validate_tls_material(&opts(|_| {})).is_empty()); + } +} diff --git a/rust/crates/adc-cli/src/server/sync.rs b/rust/crates/adc-cli/src/server/sync.rs new file mode 100644 index 00000000..fba852d3 --- /dev/null +++ b/rust/crates/adc-cli/src/server/sync.rs @@ -0,0 +1,177 @@ +//! `PUT /sync`: lint (optional) + diff against the remote backend + apply. + +use std::collections::{HashMap, HashSet}; + +use adc_sdk::resources::Configuration; +use adc_sdk::{BackendSyncOptions, BackendSyncResult, Event, ResourceType}; +use axum::Json; +use axum::body::Bytes; +use axum::http::StatusCode; +use axum::response::{IntoResponse, Response}; +use serde_json::{Value, json}; + +use super::schema::{self, SyncInput}; +use super::{backend, bad_request, internal_error}; +use crate::config; +use crate::error::CliError; +use crate::pipeline; + +pub async fn sync_handler(body: Bytes) -> Response { + let input: SyncInput = match serde_json::from_slice(&body) { + Ok(input) => input, + Err(error) => return bad_request(json!({"message": error.to_string(), "errors": []})), + }; + let opts = input.task.opts; + + let tls_issues = schema::validate_tls_material(&opts); + if !tls_issues.is_empty() { + return bad_request(json!({"message": "invalid TLS material", "errors": tls_issues})); + } + + let label_selector = opts.label_selector_or_default(); + let mut config_value = input.task.config; + config::fill_labels(&mut config_value, &label_selector); + + let mut configuration: Configuration = match serde_json::from_value(config_value) { + Ok(configuration) => configuration, + Err(error) => { + return bad_request( + json!({"message": format!("invalid configuration: {error}"), "errors": []}), + ); + } + }; + + let (include, exclude) = opts.resource_type_sets(); + config::filter_resource_types(&mut configuration, &include, &exclude); + + if opts.lint { + let issues = adc_sdk::lint::lint(&configuration); + if !issues.is_empty() { + return bad_request(json!({ + "message": "Lint configuration\nThe following errors were found in configuration:", + "errors": issues.iter().map(lint_issue_json).collect::>(), + })); + } + } + + match run( + &opts.backend, + &opts, + configuration, + &include, + &exclude, + &label_selector, + ) + .await + { + Ok(output) => (StatusCode::ACCEPTED, Json(output)).into_response(), + Err(error) => internal_error(json!({"message": error.to_string()})), + } +} + +async fn run( + backend_kind: &str, + opts: &schema::Opts, + local: Configuration, + include: &HashSet, + exclude: &HashSet, + label_selector: &HashMap, +) -> Result { + let gateway = backend::build_backend(opts)?; + let remote = pipeline::load_remote(gateway.as_ref(), include, exclude, label_selector).await?; + let events = pipeline::diff(gateway.as_ref(), &local, &remote).await?; + let sync_opts = BackendSyncOptions { + concurrent: Some(opts.request_concurrent), + exit_on_failure: Some(false), + }; + + let is_apisix_standalone = backend_kind == "apisix-standalone"; + let events_for_output = is_apisix_standalone.then(|| events.clone()); + let results = gateway.sync(events, sync_opts).await?; + Ok(match events_for_output { + Some(events) => output_for_apisix_standalone(&events, &results), + None => output(&results), + }) +} + +fn status_of(total: usize, successes: usize, failures: usize) -> &'static str { + if total == successes { + "success" + } else if total == failures { + "all_failed" + } else { + "partial_failure" + } +} + +fn output(results: &[BackendSyncResult]) -> Value { + let now = chrono::Utc::now().to_rfc3339(); + let (successes, failures): (Vec<_>, Vec<_>) = results.iter().partition(|r| r.success); + + json!({ + "status": status_of(results.len(), successes.len(), failures.len()), + "total_resources": results.len(), + "success_count": successes.len(), + "failed_count": failures.len(), + "success": successes.iter().map(|r| json!({ + "server": r.server, + "event": r.event.as_ref().map(simplify_event), + "synced_at": now, + })).collect::>(), + "failed": failures.iter().map(|r| json!({ + "server": r.server, + "event": r.event.as_ref().map(simplify_event), + "failed_at": now, + "reason": r.error.as_ref().map(|e| e.to_string()).unwrap_or_default(), + })).collect::>(), + }) +} + +/// One `BackendSyncResult` per *server* here, not per event — `success`/ +/// `failed` describe `events` directly, `endpoint_status` carries the +/// per-server detail. +fn output_for_apisix_standalone(events: &[Event], results: &[BackendSyncResult]) -> Value { + let now = chrono::Utc::now().to_rfc3339(); + let (successes, failures): (Vec<_>, Vec<_>) = results.iter().partition(|r| r.success); + + json!({ + "status": status_of(results.len(), successes.len(), failures.len()), + "total_resources": 0, + "success_count": successes.len(), + "failed_count": failures.len(), + "success": events.iter().map(|event| json!({ + "event": simplify_event(event), + "synced_at": now, + })).collect::>(), + "failed": Vec::::new(), + "endpoint_status": results.iter().map(|r| json!({ + "server": r.server, + "success": r.success, + "reason": r.error.as_ref().map(|e| e.to_string()), + "requested_at": now, + })).collect::>(), + }) +} + +fn simplify_event(event: &Event) -> Value { + let mut value = match serde_json::to_value(event) { + Ok(value) => value, + // The sync itself already happened by the time this runs — a + // response with a placeholder event beats panicking and losing the + // result entirely. + Err(error) => return json!({"error": format!("failed to serialize event: {error}")}), + }; + if let Value::Object(map) = &mut value { + map.remove("old_value"); + map.remove("new_value"); + map.remove("diff"); + } + value +} + +fn lint_issue_json(issue: &adc_sdk::lint::LintIssue) -> Value { + json!({ + "path": issue.path.iter().map(ToString::to_string).collect::>(), + "message": issue.message, + }) +} diff --git a/rust/crates/adc-cli/src/server/validate.rs b/rust/crates/adc-cli/src/server/validate.rs new file mode 100644 index 00000000..885fba95 --- /dev/null +++ b/rust/crates/adc-cli/src/server/validate.rs @@ -0,0 +1,115 @@ +//! `PUT /validate`: lint (optional) + backend-side validation of the +//! events that would be produced against an *empty* remote config. + +use adc_sdk::resources::Configuration; +use adc_sdk::BackendValidateResult; +use axum::Json; +use axum::body::Bytes; +use axum::http::StatusCode; +use axum::response::{IntoResponse, Response}; +use serde_json::json; + +use super::schema::{self, ValidateInput}; +use super::{backend, bad_request, internal_error}; +use crate::config; +use crate::pipeline; + +fn empty_configuration() -> Configuration { + Configuration { + services: None, + ssls: None, + consumers: None, + consumer_groups: None, + global_rules: None, + plugin_metadata: None, + } +} + +pub async fn validate_handler(body: Bytes) -> Response { + let input: ValidateInput = match serde_json::from_slice(&body) { + Ok(input) => input, + Err(error) => { + return bad_request(json!({"success": false, "source": "input", "message": error.to_string(), "errors": []})); + } + }; + let opts = input.task.opts; + + let tls_issues = schema::validate_tls_material(&opts); + if !tls_issues.is_empty() { + return bad_request(json!({ + "success": false, "source": "input", + "message": "invalid TLS material", "errors": tls_issues, + })); + } + + let label_selector = opts.label_selector_or_default(); + let mut config_value = input.task.config; + config::fill_labels(&mut config_value, &label_selector); + + let mut configuration: Configuration = match serde_json::from_value(config_value) { + Ok(configuration) => configuration, + Err(error) => { + return bad_request(json!({ + "success": false, "source": "input", + "message": format!("invalid configuration: {error}"), "errors": [], + })); + } + }; + + let (include, exclude) = opts.resource_type_sets(); + config::filter_resource_types(&mut configuration, &include, &exclude); + + if opts.lint { + let issues = adc_sdk::lint::lint(&configuration); + if !issues.is_empty() { + return bad_request(json!({ + "success": false, "source": "lint", + "message": "Lint configuration\nThe following errors were found in configuration:", + "errors": issues.iter().map(|issue| json!({ + "path": issue.path.iter().map(ToString::to_string).collect::>(), + "message": issue.message, + })).collect::>(), + })); + } + } + + let gateway = match backend::build_backend(&opts) { + Ok(gateway) => gateway, + Err(error) => return internal_error(json!({"success": false, "message": error.to_string(), "errors": []})), + }; + + let events = match pipeline::diff(gateway.as_ref(), &configuration, &empty_configuration()).await { + Ok(events) => events, + Err(error) => return internal_error(json!({"success": false, "message": error.to_string(), "errors": []})), + }; + + match gateway.validate(&events).await { + Ok(BackendValidateResult { success, error_message, errors }) => { + let mut body = json!({"success": success, "source": "validate", "errors": errors_json(&errors)}); + if let Some(message) = error_message { + body["message"] = json!(message); + } + (StatusCode::OK, Json(body)).into_response() + } + Err(adc_sdk::BackendError::Unsupported(_)) => bad_request(json!({ + "success": false, "source": "validate", + "message": "Validate is not supported by the current backend.", "errors": [], + })), + Err(error) => internal_error(json!({"success": false, "message": error.to_string(), "errors": []})), + } +} + +fn errors_json(errors: &[adc_sdk::BackendValidationError]) -> serde_json::Value { + json!( + errors + .iter() + .map(|e| json!({ + "resource_type": e.resource_type, + "resource_id": e.resource_id, + "resource_name": e.resource_name, + "index": e.index, + "error": e.error, + })) + .collect::>() + ) +} diff --git a/rust/crates/adc-cli/tests/assets/tls/ca.cer b/rust/crates/adc-cli/tests/assets/tls/ca.cer new file mode 100644 index 00000000..26008a9f --- /dev/null +++ b/rust/crates/adc-cli/tests/assets/tls/ca.cer @@ -0,0 +1,19 @@ +-----BEGIN CERTIFICATE----- +MIIDFTCCAf2gAwIBAgIUd8RVHd2+mZHqrYof1Ihf3l7H3pswDQYJKoZIhvcNAQEL +BQAwETEPMA0GA1UEAwwGUk9PVENBMCAXDTI2MDgxODE2NDQ0N1oYDzIxMjYwNzI1 +MTY0NDQ3WjARMQ8wDQYDVQQDDAZST09UQ0EwggEiMA0GCSqGSIb3DQEBAQUAA4IB +DwAwggEKAoIBAQDI3w1Bpe0S1Al6EJ7OEzdnVWsqZR5sD8At/hr94RDRpRte0+uw +NWPP6Bsi9XbUuC7tqGps2Cl5QMxpuxKHiJbhIca5lpIQvPV/BhbAER4CZ/dpGxJn +FN20F0NaxNghl83SjWngzc1rn4w9rwWBrdDZFPMgqNphj0PIlOR0pyE9UW6x8j8g +cqp0t+2hgcJ8zcX1abY1dNS1u7Kb6KFLNFkTkF7KLj4QgJ+WgYutYtCmPaGHc90q +F81VzeUTt6VoNRG1uemM9kPjaxthJkRQGV37+DaDYDPfroaZufNIjHeGkgWa9pvI +w0u0ZzyL73LGLj+uVEgngqUbDNqxWta0so4rAgMBAAGjYzBhMB0GA1UdDgQWBBR9 +a5lP7QVT4asmF9SC3AQwK04KqDAfBgNVHSMEGDAWgBR9a5lP7QVT4asmF9SC3AQw +K04KqDAPBgNVHRMBAf8EBTADAQH/MA4GA1UdDwEB/wQEAwIBBjANBgkqhkiG9w0B +AQsFAAOCAQEAeRQNdtITxb9p2Yo0+L36bZqYVp+3/v8hp548Lc/eEq2HLztlA3yL +hdukchmSgyHkDRenzBE1+m04V0FcoRqzm9gNueaVCflRSx276j+q8+9SdpC99jUG +jz91pNaReEkHi/MJPlT+gN078DhIqVq1tsuZ0oDyKoztQrHPGN7e9ypeYToQMWFQ +MOSA0mujILK7xb/dR92f9ExZNW1dGkEz9ts73NwfMjnnsn1MQXp3UDxsACPeTFee +EnY1QI/JeUwHkoE84bzJmJfOqXPYqeXdc1J1dR45jEy9W+/PtiacrqM9uTVzasSY +vAXrTehhOWBtK4QtBk/RBODndtHn++UIpw== +-----END CERTIFICATE----- diff --git a/rust/crates/adc-cli/tests/assets/tls/ca.key b/rust/crates/adc-cli/tests/assets/tls/ca.key new file mode 100644 index 00000000..06bbf36d --- /dev/null +++ b/rust/crates/adc-cli/tests/assets/tls/ca.key @@ -0,0 +1,28 @@ +-----BEGIN PRIVATE KEY----- +MIIEvQIBADANBgkqhkiG9w0BAQEFAASCBKcwggSjAgEAAoIBAQDI3w1Bpe0S1Al6 +EJ7OEzdnVWsqZR5sD8At/hr94RDRpRte0+uwNWPP6Bsi9XbUuC7tqGps2Cl5QMxp +uxKHiJbhIca5lpIQvPV/BhbAER4CZ/dpGxJnFN20F0NaxNghl83SjWngzc1rn4w9 +rwWBrdDZFPMgqNphj0PIlOR0pyE9UW6x8j8gcqp0t+2hgcJ8zcX1abY1dNS1u7Kb +6KFLNFkTkF7KLj4QgJ+WgYutYtCmPaGHc90qF81VzeUTt6VoNRG1uemM9kPjaxth +JkRQGV37+DaDYDPfroaZufNIjHeGkgWa9pvIw0u0ZzyL73LGLj+uVEgngqUbDNqx +Wta0so4rAgMBAAECggEAIZrPSPBNXR0ECNvG9YrZdfwgVZNdJ47rA8bDFT4V5jzM ++2xQvcXw0NNv1sVh/+xgTXojc9ol9hcVG4skanA7baaM7Hd4MDyshXerTq6OarCh +/3978KrY/Ev4BLNxxQz0bgkicW18tEiY2ajyLuO5UNfkZM5a2n9xQ5lFLw7WzL8K ++ZrS89rgnO9sPvD0aD6CUWkXUFP+qY4qskl9ylHprkIRCRXtk3Icw+lqCcQd4LYF +kdyjmHfhc0PQp3ididur3bZzhHXaKW9HcJ71BkMRtrSDGL1sjsdaP6oJIvSLvq57 +KY3AvvMNrsTH+Y/5ZAvaqCxydWO6nViJ3CqGBbOXQQKBgQDj6aI8NxZQIuP++JNO +Q8aOsixmzfE9wH21wOQdMvGsPGwaJD/g87tvWZkYOEhsZ1smO4QvQwEOE/DCnpyL +Uur8dlxCSQ0VYLeEwsCEuKtszcSU5zSHluI4verIgygi2Iqp6hTfXr3ooYjxQf7x +QVA05LnGKYGqUw9gDk24lozSCQKBgQDhoEv7eH62OerhV6X+VAuLse2gQK1cjIZ1 +I8TpP2aKF8crE8ooN8TnP3jj1ShozYDsH1WiY5lFTSKfoCHNiJscaCBDWzpgHuub +Gqb+gykEGTqYDepbxZP2Ro4NDD2Z/4MQyTymiv6lz0j7eTMHWxcLWOUR94IloAlD ++8Te+fcbkwKBgQDdhxkHOHA6wj8kdM8RkrUrvCmGX4SuBizqfiv76amYRT66BiQE +7kNwjwFM1mAm5itltRHdsl4TJfSt5ue4UIdRj2ZLk5/g+JpIs9fW6XzOjA8YwMaB +SHpotsi/zyQzApF9aKaTGw6yUFjAT+qS624fi3a7E1sSiBt4vU50LfmAqQKBgA/A +H+3DIJ1Z97KZasYRWej7l8oLGc8PJEfDInjh6yeSt12jeQZLtlwqSyckdzixt+FD +4rd+WnHDC7q29AUkFyfpgO8SzEVvgyUFvEiiIVfe5v88YXLcnRKhJEN26kn40051 +rd02cMZkbQTZFh3aVwZ8wyj47UXxIRR02+5w5rYvAoGAfl1ma8nILMmkdfwDUJpk +h6qn/6i7aXN2yn0ITCdmve0/ksTulxhVo1d9uIO/Z/NKeoEVS7u+YGahCDgpPBG0 +SEtesthdxRcjmsqyS+2oU0otzyG5VwFpS/O53lQE3LNlz7f97FZvCzFHonvw0nMb +8NKiS3s4ITYl5P0sLpyBlaE= +-----END PRIVATE KEY----- diff --git a/rust/crates/adc-cli/tests/assets/tls/client.cer b/rust/crates/adc-cli/tests/assets/tls/client.cer new file mode 100644 index 00000000..bafa5831 --- /dev/null +++ b/rust/crates/adc-cli/tests/assets/tls/client.cer @@ -0,0 +1,19 @@ +-----BEGIN CERTIFICATE----- +MIIDFDCCAfygAwIBAgIUIsGGu2NXIwAMe77j5l2wBmx+DAowDQYJKoZIhvcNAQEL +BQAwETEPMA0GA1UEAwwGUk9PVENBMCAXDTI2MDgxODE2NDQ0N1oYDzIxMjYwNzI1 +MTY0NDQ3WjARMQ8wDQYDVQQDDAZDTElFTlQwggEiMA0GCSqGSIb3DQEBAQUAA4IB +DwAwggEKAoIBAQDQpiPqAtPaalQo4pj76xc8duCxsnGwHe5OfXi/ltzXNogM9/p7 +dAhpGW5SEcYXyakA5hGnsiAlISfcx7W5YUXb0ML33AuTTylDRHLU0AnNyEoZjW5s +edFKu5wjQrM6C5AvcrvZzCkpsIsHv5pj4El7nil+LZohimnHD+ggLk9tJIAWyZdY +2G+IUjKN9RscN/PnR9l5yXGlC1kldixidqLerWaBQgoC2DfGJ8J7Q0es+j8kAZJL +scwz6TQy1hljqXZGtZvpZSEM+YPBT5NOW+vDUwJAFBUrUQnmeFneB/RAcPA1D4bg +t2yMN/kEdPT7+GyVS2hnsMKhPadiP+1CrltrAgMBAAGjYjBgMAkGA1UdEwQCMAAw +EwYDVR0lBAwwCgYIKwYBBQUHAwIwHQYDVR0OBBYEFI1CQzQgwxu+Ma9rvtGL+2T3 +Z2mCMB8GA1UdIwQYMBaAFH1rmU/tBVPhqyYX1ILcBDArTgqoMA0GCSqGSIb3DQEB +CwUAA4IBAQCGYTMJX+RhCN6BMRSU7PJLRxBPAPGeBL5DdNmkR7RgBV3juPjaYmBL +ZaqMDs5FOik47PZt+OmHd1UQ3SfTzfAe+1HwywYFtGu7apvYaSCTdZCWD6bYTgne +9khfU9aqThUEmUkp7dzWneWV+ovhbuap5GNaNtaR5DypCcrMAtE+agsUxn0Abvqv +ZV4RH6YsaHa80dKCFhZi0O0xBrClOqSJQL5GVTg4ai/ywMJQlN78eSR7aeqIcl7K +NHeMMi4LtwhewUVpRZAv3d9WocvQxcS/Yi+IYc1ByNZ8XLjKn8ad1Ytt3X6jYbJs +CpWwR7Oafnu+GjctORFUwP5YRq490hi+ +-----END CERTIFICATE----- diff --git a/rust/crates/adc-cli/tests/assets/tls/client.csr b/rust/crates/adc-cli/tests/assets/tls/client.csr new file mode 100644 index 00000000..183e3b6f --- /dev/null +++ b/rust/crates/adc-cli/tests/assets/tls/client.csr @@ -0,0 +1,15 @@ +-----BEGIN CERTIFICATE REQUEST----- +MIICVjCCAT4CAQAwETEPMA0GA1UEAwwGQ0xJRU5UMIIBIjANBgkqhkiG9w0BAQEF +AAOCAQ8AMIIBCgKCAQEA0KYj6gLT2mpUKOKY++sXPHbgsbJxsB3uTn14v5bc1zaI +DPf6e3QIaRluUhHGF8mpAOYRp7IgJSEn3Me1uWFF29DC99wLk08pQ0Ry1NAJzchK +GY1ubHnRSrucI0KzOguQL3K72cwpKbCLB7+aY+BJe54pfi2aIYppxw/oIC5PbSSA +FsmXWNhviFIyjfUbHDfz50fZeclxpQtZJXYsYnai3q1mgUIKAtg3xifCe0NHrPo/ +JAGSS7HMM+k0MtYZY6l2RrWb6WUhDPmDwU+TTlvrw1MCQBQVK1EJ5nhZ3gf0QHDw +NQ+G4LdsjDf5BHT0+/hslUtoZ7DCoT2nYj/tQq5bawIDAQABoAAwDQYJKoZIhvcN +AQELBQADggEBABsaa46LkaYY22rguLJ8osuXeGvEVo7wwQ0Y4dZmvC/L6pUUDlLU +QHx33sEeESrpgz9rbVZgVAgHCM543z1UJT+CRQZ9f3XrO9nin09REWY3dP0204tf +ZOZ0iynfJNGfQAfBQKhzBBTAfITsRrSMTYSXz4kAw0qwA5+wGOr+iokZlVs0yy0V +S16JblyJcZNs4gYF1STv2d1NQPV8gIeUkmP/tYRcpvQjZ1XVKCUvasq5gkjtLyRL +6HYO11MHE6VyZJShdQoDRR4m53eBZp571Y98rZqxBII18RVtGU2KRuO5azSBfF20 +bZdPS8095UU90amQZSqapsTtk2zh180jx1c= +-----END CERTIFICATE REQUEST----- diff --git a/rust/crates/adc-cli/tests/assets/tls/client.key b/rust/crates/adc-cli/tests/assets/tls/client.key new file mode 100644 index 00000000..be67b5fe --- /dev/null +++ b/rust/crates/adc-cli/tests/assets/tls/client.key @@ -0,0 +1,28 @@ +-----BEGIN PRIVATE KEY----- +MIIEvgIBADANBgkqhkiG9w0BAQEFAASCBKgwggSkAgEAAoIBAQDQpiPqAtPaalQo +4pj76xc8duCxsnGwHe5OfXi/ltzXNogM9/p7dAhpGW5SEcYXyakA5hGnsiAlISfc +x7W5YUXb0ML33AuTTylDRHLU0AnNyEoZjW5sedFKu5wjQrM6C5AvcrvZzCkpsIsH +v5pj4El7nil+LZohimnHD+ggLk9tJIAWyZdY2G+IUjKN9RscN/PnR9l5yXGlC1kl +dixidqLerWaBQgoC2DfGJ8J7Q0es+j8kAZJLscwz6TQy1hljqXZGtZvpZSEM+YPB +T5NOW+vDUwJAFBUrUQnmeFneB/RAcPA1D4bgt2yMN/kEdPT7+GyVS2hnsMKhPadi +P+1CrltrAgMBAAECggEAPH4+4WUKeUPkvKneAwQJC53HzZ1X+uDiq90S+jFKPBdy +YJgxBkQBAD/ATYkbrt/n4PvTWJR7X2h6fzdjx6idMXsYW/ZvYLlN1FPvGyZqAUC1 +wyzPPCIhfRJh1ZNMFWMu3aLdNetMb+rglFGH+LcZdv7HNu8PxfO0cWN6QIJMwu6Q +0Icdn/mKLqRVLUxdnDlswO0OOypOxMwKEZJ2kgWYaMBtNLlVDsRon90u31AKNppZ +Eu4d5cX37J1CXczFwSbTg65LJk68aR3+97DlOrCvRKv6BfWEHGrCgBZWBzji7jvG +c3KtKK3yv5vqwNpwJXKxXk6h1PUiLckF6nAT+Kuk4QKBgQD5dGcRaj0FfnZqjOJW +PqDyaepB2bXFh/D6fzJEiZia3yPmd00IrUiqoGSyYYGJmt8uVA9YjW+z2okKq+lc +GTrll7n9D46L9n95M0Z0BBMZGRHz+AJisYBz1qzKIQJXcOtiMERJLLNgp2vYimi1 +IDYORAU3dlc/AVPg+jjajHifsQKBgQDWH6TgIeV//zwhAvqqxiIIo54xbgfsvfyD +8qnQNYFYR2rqxQ/+mQ3XfNqXG0s9lXI6CS1ENpLbL/vzUxC+xheutTyoJxSP1ddk +8wEDN9Pgx6mOaQ6jSu46U4aovJeJSWODRxsudeCDMH44owslCQfWamNHB89ArbTw +I7FaPkBv2wKBgQDv3tG5OkpBPTDLFnwSaJjFYbmD5sBWmHjNt3/zzcfzrHxOAgwO +Ouq0QBV0PjScyFKxrt0uzppJ/OtoWpTEHfK3kaWjxNDSn45GUlr99mkS6juMOMC6 +fGrDePugRguFX6zINxeCsbwvRe57Q+SZvsacAyZtBZuxlyo8HQCMjyTykQKBgB1c +y4Q8wbbyrjEssmkWsHYU0c2fdBC/4M/LSAQYQjtz17KIAXB9VouVQHh2MrQoOTjC +J2XyQeMyyk8MtgAjM/4uNjos2cH7pgTe2eWyEykA2DyCJZK45MA00gNzkSgvWykW +aCDP41C6JqTnntCeU2fQwPptlLse1vATRO/GF5n/AoGBAN8iVTEcki/dLhTkq11l +r8vuBWhVq6HUlL0rwncEJ/naw1vfI5UQLzamJbCWXSC2+Puu5MdkkphCkXLmMBgK +mqVr0R2Sf2yy6Xl6cfNIzz1gC6648oFz/D3lnGeRAjWdYMFxZUgclL9O1vdQ3E2g +CyU/a9AtoOVawFQKqkfOTwGM +-----END PRIVATE KEY----- diff --git a/rust/crates/adc-cli/tests/assets/tls/generate-mtls.sh b/rust/crates/adc-cli/tests/assets/tls/generate-mtls.sh new file mode 100755 index 00000000..c4a1ceb5 --- /dev/null +++ b/rust/crates/adc-cli/tests/assets/tls/generate-mtls.sh @@ -0,0 +1,21 @@ +#!/bin/bash +# Generates the self-signed CA/server/client cert chain these tests use. +# Uses -addext/-extfile to produce real X.509v3 certs (rustls/webpki reject v1). +set -euo pipefail +cd "$(dirname "$0")" + +openssl genrsa -out ca.key 2048 +openssl req -new -x509 -sha256 -days 36500 -key ca.key -out ca.cer -subj "/CN=ROOTCA" \ + -addext "basicConstraints=critical,CA:TRUE" -addext "keyUsage=critical,keyCertSign,cRLSign" + +openssl genrsa -out server.key 2048 +openssl req -new -sha256 -key server.key -out server.csr -subj "/CN=localhost" +openssl x509 -req -days 36500 -sha256 -CA ca.cer -CAkey ca.key -CAcreateserial -in server.csr -out server.cer \ + -extfile <(printf "basicConstraints=CA:FALSE\nsubjectAltName=DNS:localhost,IP:127.0.0.1\nextendedKeyUsage=serverAuth") + +openssl genrsa -out client.key 2048 +openssl req -new -sha256 -key client.key -out client.csr -subj "/CN=CLIENT" +openssl x509 -req -days 36500 -sha256 -CA ca.cer -CAkey ca.key -CAcreateserial -in client.csr -out client.cer \ + -extfile <(printf "basicConstraints=CA:FALSE\nextendedKeyUsage=clientAuth") + +rm -f ca.srl diff --git a/rust/crates/adc-cli/tests/assets/tls/server.cer b/rust/crates/adc-cli/tests/assets/tls/server.cer new file mode 100644 index 00000000..94767f14 --- /dev/null +++ b/rust/crates/adc-cli/tests/assets/tls/server.cer @@ -0,0 +1,20 @@ +-----BEGIN CERTIFICATE----- +MIIDMzCCAhugAwIBAgIUIsGGu2NXIwAMe77j5l2wBmx+DAkwDQYJKoZIhvcNAQEL +BQAwETEPMA0GA1UEAwwGUk9PVENBMCAXDTI2MDgxODE2NDQ0N1oYDzIxMjYwNzI1 +MTY0NDQ3WjAUMRIwEAYDVQQDDAlsb2NhbGhvc3QwggEiMA0GCSqGSIb3DQEBAQUA +A4IBDwAwggEKAoIBAQDpPH5LPCO0dfl87t+y8/iJ5wnCJu/hZvYJha4H2CV+59rk +pw2q71BuEgieTTT49flprvhhOY6VnsLdR88vDey8H7GRTY/PaqxBliewETdTK5gZ +IppFOF2WwzY1GbMqfDmNB81JZEI2pTKagznHi8IOhfNzZAjCQwsbiRUE87R9BSo8 +cLv51aS+2VCaV3dlJE6n38+CHOfr/QC97YjefdGXYstrtKdEsQ758CtJ5XPezGYL +UMMm277ZD3oJ6Ce/kEnBjGZfcRUFa/n+IwyfTL2gI/MZ2I1pEyD4F7I4Z5Xw7xwB +tNenfpr3ckVMvDW78nml7SfteCdheKPHXTn9IkR7AgMBAAGjfjB8MAkGA1UdEwQC +MAAwGgYDVR0RBBMwEYIJbG9jYWxob3N0hwR/AAABMBMGA1UdJQQMMAoGCCsGAQUF +BwMBMB0GA1UdDgQWBBRjqD6zQxypM77JKmD1ck/6Ux75OjAfBgNVHSMEGDAWgBR9 +a5lP7QVT4asmF9SC3AQwK04KqDANBgkqhkiG9w0BAQsFAAOCAQEAe8+hcXQJq0vA +HkLX7/nZ5Fqu9sVxtF2K/RGVrKnR6xDimNfpVMhDovc5kBxuUoJ6LjAaeAaVMYEc +oxCuZLYGfpcazX2DH/hoGav/y4YkDeOAVxeB730A1Hhr1y07CVaJwyussoRSQcnO +2nK6VZlCrbVN2xUt5i/Qnb/Am+S1iZpzqu1pgyAaq06FOnu5MZ9himCkTJ//wzrs +KPCaGe5s490UkWGXzVZ7QsSWtKG4qXn3srR4+ksne8sOSfHjI0L7CqX52/4QZXsX +kcz5v4ET3MrrJWPt7RB8Sv4pgCC8wAaddwnRTR5e0/AgGcPozwhQKcAA3qeqpe9n +ZEWCuukGRQ== +-----END CERTIFICATE----- diff --git a/rust/crates/adc-cli/tests/assets/tls/server.csr b/rust/crates/adc-cli/tests/assets/tls/server.csr new file mode 100644 index 00000000..44e0f7f0 --- /dev/null +++ b/rust/crates/adc-cli/tests/assets/tls/server.csr @@ -0,0 +1,15 @@ +-----BEGIN CERTIFICATE REQUEST----- +MIICWTCCAUECAQAwFDESMBAGA1UEAwwJbG9jYWxob3N0MIIBIjANBgkqhkiG9w0B +AQEFAAOCAQ8AMIIBCgKCAQEA6Tx+SzwjtHX5fO7fsvP4iecJwibv4Wb2CYWuB9gl +fufa5KcNqu9QbhIInk00+PX5aa74YTmOlZ7C3UfPLw3svB+xkU2Pz2qsQZYnsBE3 +UyuYGSKaRThdlsM2NRmzKnw5jQfNSWRCNqUymoM5x4vCDoXzc2QIwkMLG4kVBPO0 +fQUqPHC7+dWkvtlQmld3ZSROp9/Pghzn6/0Ave2I3n3Rl2LLa7SnRLEO+fArSeVz +3sxmC1DDJtu+2Q96Cegnv5BJwYxmX3EVBWv5/iMMn0y9oCPzGdiNaRMg+BeyOGeV +8O8cAbTXp36a93JFTLw1u/J5pe0n7XgnYXijx105/SJEewIDAQABoAAwDQYJKoZI +hvcNAQELBQADggEBAMdlBHBIStN3yvkyjnP6Xyx0WCmI262XEZizMs7cZzaT372R +IBAypIWMdoRt8vLPS0mgD7dOhgv4dlorfaT+HH/ir4GDFr+EGyRLIuLfD37a/XOX +AQqUAOhL2z1V0fg8IQoHtZibaL3Zug1AkPQfk8Dk7qeFoPbItw5EKEp7iPXCrLlv +7EUx+MSm5J6AHjnTwmMGMyBdL2W85C0I51VT1m+fMgIHY96/6jesorZC5P3kCN1q +6mIn2heAxG/zdcSgYqS+W/dx10+fbLUa0rrmNRNe6UPe/S6yTmZmy/gv1TI+8LMR +BkIlPyUquwmzgv3ALecTsaOU4L6ahVjb/d+4WAU= +-----END CERTIFICATE REQUEST----- diff --git a/rust/crates/adc-cli/tests/assets/tls/server.key b/rust/crates/adc-cli/tests/assets/tls/server.key new file mode 100644 index 00000000..291469e8 --- /dev/null +++ b/rust/crates/adc-cli/tests/assets/tls/server.key @@ -0,0 +1,28 @@ +-----BEGIN PRIVATE KEY----- +MIIEvwIBADANBgkqhkiG9w0BAQEFAASCBKkwggSlAgEAAoIBAQDpPH5LPCO0dfl8 +7t+y8/iJ5wnCJu/hZvYJha4H2CV+59rkpw2q71BuEgieTTT49flprvhhOY6VnsLd +R88vDey8H7GRTY/PaqxBliewETdTK5gZIppFOF2WwzY1GbMqfDmNB81JZEI2pTKa +gznHi8IOhfNzZAjCQwsbiRUE87R9BSo8cLv51aS+2VCaV3dlJE6n38+CHOfr/QC9 +7YjefdGXYstrtKdEsQ758CtJ5XPezGYLUMMm277ZD3oJ6Ce/kEnBjGZfcRUFa/n+ +IwyfTL2gI/MZ2I1pEyD4F7I4Z5Xw7xwBtNenfpr3ckVMvDW78nml7SfteCdheKPH +XTn9IkR7AgMBAAECggEAV3T6DH8QCmakdz7hPeqy2w75v0Y3c+dWQcrRL5rSsIwD +LfMgMmULXULA3Y8o2mPtsr3L4DUjbKI8ApqfK09G4mHmBQy27LlcvzktR52lB7hU +j7REcclJerNXe8DXyIoNUH9I8Ii6NWBrobmsLFGRIj4DRFUR3bojC5+y9IjnuGrE +oQyAcQxmkImnw/fwlXkjx0ViaHzxvs3APwGeonB2o7ZUl9t6sevCzuGF1wzC6UHv +cti0Qh0K9bNlXyS+97ugW0bVDIsUw88dVB0L5aXJkJBvNPDA++Lire3B5qc5/dnT +fS3+boxIrLMqR8hDRmIZZUlkzNoSGVgC1/94q93zCQKBgQD2s8dmSs4tHReSxW5H +u2gva8jd29BGio/2jywSA2JvEgnznUQdx8Pyqr7wL3u/BH4xUe1w+793yy0c7E07 +d4EjgVMkkus9TXvIxvnCYYnhVPzAeYtmo/72GVVTR2jwHy0/x6YdlgYulE7F/rjJ +5tF+3CYWayPeigLMYzk2zBBFgwKBgQDyBsr05rdsvxWdABFQJnl5hbz+btbS6t8y +gAhs7CC5D5BqhuXkg/KryQ5Aok453FVQSREomvKONQNcGtE7gAUbME8r8INkMRe9 +d2lLNpeXDSI7KMS41vzLqcRgwpg/1UpmyYpZCS0NSFzW64XEYh+EsDyQclMLnrtT ++hAubi1LqQKBgQCCkMFejQazb6szPZRRGIlaV6Q2bwi63Mi2iC2d1va4rAZiTYBo +dnppKx7kxWyruvgCqEaPPl2mS/yzSwjRCT1qih5zw+IGTsTNjSlQTAkKHc2rHGi/ +yNm+a8fxzGBofUeYctSi4eyhqFJMjbRE/wkvJ9pskQWp2McEXxs/uh5+ewKBgQCu +6+Xn1pAfUoPGcvQQX55QDC6qHWW6DvK9xvdP8eE8n1kbBOBGpm7PZYKdiDDNdMdc +PVLfbA1+ZiZFfURXopEOM34lHbF4ylqEHzfEmnI5Q87HvxFfHlKax9ocrMfo6rjZ +TTRmYVFkVjZzRsnpQ5nQBqffJiGLNm/ho8vqIsst8QKBgQCW2M/2KIeYS4Z/ODMj +NQ/l7oYuSR+WZf03TlAetuG2E6wjUL6Z2mFOQqZ+0qeI4H8LINYLOYW1ltposQrj +rdBhgRoNMYe9cIhuilu26QaaibQuMPSl0nFnQIThbrNSEFPx+wQ4zqlwzq8kMr65 +N9Vj/MaZN6kprStDBEDhz4awqw== +-----END PRIVATE KEY----- diff --git a/rust/crates/adc-cli/tests/ingress_server_sigint.rs b/rust/crates/adc-cli/tests/ingress_server_sigint.rs new file mode 100644 index 00000000..2da92cbc --- /dev/null +++ b/rust/crates/adc-cli/tests/ingress_server_sigint.rs @@ -0,0 +1,98 @@ +//! Black-box test of the ingress-server daemon's graceful shutdown — +//! spawns the real `adc` binary as a subprocess and sends it a real SIGINT. + +use std::process::Stdio; +use std::time::Duration; + +use tokio::io::{AsyncBufReadExt, BufReader}; +use tokio::process::Command; +use tokio::time::timeout; + +fn free_port() -> u16 { + std::net::TcpListener::bind("127.0.0.1:0").unwrap().local_addr().unwrap().port() +} + +#[tokio::test] +async fn sigint_shuts_down_both_listeners_gracefully_and_exits_zero() { + let listen_port = free_port(); + let status_port = free_port(); + + let mut child = Command::new(env!("CARGO_BIN_EXE_adc")) + .args([ + "ingress-server", + "--listen", + &format!("http://127.0.0.1:{listen_port}"), + "--listen-status", + &format!("{status_port}"), + ]) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .kill_on_drop(true) + .spawn() + .expect("failed to spawn the adc binary"); + let pid = child.id().expect("child should have a pid right after spawning") as i32; + + let ready = timeout(Duration::from_secs(10), async { + loop { + if let Ok(response) = reqwest::get(format!("http://127.0.0.1:{status_port}/healthz/ready")).await + && response.status().is_success() + { + return; + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + }) + .await; + assert!(ready.is_ok(), "server never became ready"); + + // SAFETY: `pid` is this test's own freshly-spawned, still-alive child. + let rc = unsafe { libc::kill(pid, libc::SIGINT) }; + assert_eq!(rc, 0, "failed to send SIGINT: {}", std::io::Error::last_os_error()); + + let exit = timeout(Duration::from_secs(10), child.wait()) + .await + .expect("process did not exit within the timeout — SIGINT was not handled gracefully") + .expect("failed to wait on child"); + assert!(exit.success(), "expected exit code 0, got {exit:?}"); +} + +/// `Ctrl+C` fired before the server finishes starting up should still shut it down promptly. +#[tokio::test] +async fn sigint_before_the_server_is_ready_still_lets_it_exit() { + let listen_port = free_port(); + let status_port = free_port(); + + let mut child = Command::new(env!("CARGO_BIN_EXE_adc")) + .args([ + "ingress-server", + "--listen", + &format!("http://127.0.0.1:{listen_port}"), + "--listen-status", + &format!("{status_port}"), + ]) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .kill_on_drop(true) + .spawn() + .expect("failed to spawn the adc binary"); + let pid = child.id().expect("child should have a pid right after spawning") as i32; + + let mut lines = BufReader::new(child.stdout.take().unwrap()).lines(); + match timeout(Duration::from_secs(10), lines.next_line()).await { + Ok(Ok(Some(_))) => {} + Ok(Ok(None)) => panic!("child exited before printing its startup line"), + Ok(Err(error)) => panic!("failed to read the child's stdout: {error}"), + Err(_) => panic!("child did not print its startup line within the timeout"), + } + + // SAFETY: `pid` is this test's own freshly-spawned, still-alive child — + // proven by having just read its startup line above. + let rc = unsafe { libc::kill(pid, libc::SIGINT) }; + assert_eq!(rc, 0, "failed to send SIGINT: {}", std::io::Error::last_os_error()); + + let exit = timeout(Duration::from_secs(10), child.wait()) + .await + .expect("process did not exit within the timeout") + .expect("failed to wait on child"); + assert!(exit.success(), "expected exit code 0, got {exit:?}"); +}