diff --git a/Cargo.toml b/Cargo.toml index eaa9638..c6127da 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -49,7 +49,9 @@ use-rustls = ["rustls"] use-rustls-ring = ["rustls-ring"] use-openssl = ["openssl"] -# optional dependencies only used in `jwt_dynamic_auth` example +# optional dependencies only used in examples and integration tests [dev-dependencies] tokio = { version = "1", features = ["full"] } bitreq = { version = "0.3.4", features = ["async-https", "json-using-serde"] } +bdk_testenv = "0.13.1" +testcontainers = { version = "0.17", features = ["blocking"] } diff --git a/tests/auth_integration.rs b/tests/auth_integration.rs new file mode 100644 index 0000000..6caab23 --- /dev/null +++ b/tests/auth_integration.rs @@ -0,0 +1,266 @@ +use std::io::{BufRead, BufReader, Write}; +use std::net::{TcpListener, TcpStream}; +use std::sync::{Arc, Mutex}; +use std::thread; +use std::time::Duration; + +use bdk_testenv::TestEnv; +use electrum_client::raw_client::RawClient; +use electrum_client::ElectrumApi; +mod helper; + +fn start_auth_proxy( + target_addr: String, + expected_auths: Option>>>, + recorded_auth: Arc>>, +) -> std::io::Result { + let listener = TcpListener::bind("127.0.0.1:0")?; + let listen_addr = listener.local_addr()?.to_string(); + + thread::spawn(move || { + let (client_stream, _) = match listener.accept() { + Ok(pair) => pair, + Err(_) => return, + }; + + let expected_auths = expected_auths.clone(); + let recorded_auth = recorded_auth.clone(); + + let server_stream = + TcpStream::connect(&target_addr).expect("failed to connect to electrum target"); + let mut client_reader = BufReader::new(client_stream.try_clone().unwrap()); + let mut client_writer_copy = client_stream.try_clone().unwrap(); + let mut client_writer_error = client_stream; + let mut server_writer = server_stream.try_clone().unwrap(); + + let mut server_reader = BufReader::new(server_stream.try_clone().unwrap()); + thread::spawn(move || { + let mut server_line = String::new(); + while let Ok(bytes_read) = server_reader.read_line(&mut server_line) { + if bytes_read == 0 { + break; + } + client_writer_copy + .write_all(server_line.as_bytes()) + .unwrap(); + client_writer_copy.flush().unwrap(); + server_line.clear(); + } + }); + + let mut line = String::new(); + let mut request_index = 0; + while let Ok(bytes_read) = client_reader.read_line(&mut line) { + if bytes_read == 0 { + break; + } + + let trimmed = line.trim_end(); + let parsed = serde_json::from_str::(trimmed); + let mut auth_value = None; + let mut forward_line = trimmed.to_string(); + if let Ok(mut json) = parsed { + if let Some(auth) = json.get("authorization").and_then(|v| v.as_str()) { + auth_value = Some(auth.to_string()); + } + + { + let mut recorded = recorded_auth.lock().unwrap(); + recorded.push(auth_value.clone().unwrap_or_default()); + } + + let auth_ok = if let Some(expected_auths) = expected_auths.as_ref() { + let expected = expected_auths.lock().unwrap(); + match expected.get(request_index) { + Some(expected) => auth_value.as_ref() == Some(expected), + None => auth_value.is_some(), + } + } else { + true + }; + + request_index += 1; + + if !auth_ok { + let resp = serde_json::json!({ + "jsonrpc": "2.0", + "id": json.get("id").cloned().unwrap_or(serde_json::json!(0)), + "error": {"code": -32600, "message": "authorization failed"} + }); + let mut out = serde_json::to_vec(&resp).unwrap(); + out.extend_from_slice(b"\n"); + client_writer_error.write_all(&out).unwrap(); + client_writer_error.flush().unwrap(); + line.clear(); + continue; + } + + json.as_object_mut().map(|obj| obj.remove("authorization")); + forward_line = serde_json::to_string(&json).unwrap(); + } + + server_writer.write_all(forward_line.as_bytes()).unwrap(); + server_writer.write_all(b"\n").unwrap(); + server_writer.flush().unwrap(); + line.clear(); + } + }); + + Ok(listen_addr) +} + +#[test] +fn test_auth_success_with_regtest_electrum() { + let env = TestEnv::new().expect("start regtest env"); + env.mine_blocks(1, None).expect("mine initial block"); + env.wait_until_electrum_sees_block(Duration::from_secs(6)) + .expect("electrum should see block"); + + let recorded_auth = Arc::new(Mutex::new(Vec::new())); + let expected_auths = Arc::new(Mutex::new(vec![ + "Bearer test-token-1".to_string(), + "Bearer test-token-1".to_string(), + ])); + let proxy_addr = start_auth_proxy( + env.electrsd.electrum_url.clone(), + Some(expected_auths), + recorded_auth.clone(), + ) + .expect("start auth proxy"); + + let token = Arc::new(Mutex::new(Some("Bearer test-token-1".to_string()))); + let auth_provider = { + let token = token.clone(); + Arc::new(move || token.lock().unwrap().clone()) + }; + + let client = RawClient::new( + proxy_addr, + Some(Duration::from_secs(5)), + Some(auth_provider), + ) + .expect("create client"); + + client.ping().expect("ping ok"); + + let auths = recorded_auth.lock().unwrap(); + assert_eq!( + auths.as_slice(), + ["Bearer test-token-1", "Bearer test-token-1"] + ); +} + +#[test] +fn test_auth_failure_when_missing_token_with_regtest_electrum() { + let env = TestEnv::new().expect("start regtest env"); + env.mine_blocks(1, None).expect("mine initial block"); + env.wait_until_electrum_sees_block(Duration::from_secs(6)) + .expect("electrum should see block"); + + let recorded_auth = Arc::new(Mutex::new(Vec::new())); + let expected_auths = Arc::new(Mutex::new(vec!["Bearer secret".to_string()])); + let proxy_addr = start_auth_proxy( + env.electrsd.electrum_url.clone(), + Some(expected_auths), + recorded_auth, + ) + .expect("start auth proxy"); + + let client_res = RawClient::new(proxy_addr, Some(Duration::from_secs(5)), None); + assert!(client_res.is_err()); +} + +#[test] +fn test_token_refresh_with_regtest_electrum() { + let env = TestEnv::new().expect("start regtest env"); + env.mine_blocks(1, None).expect("mine initial block"); + env.wait_until_electrum_sees_block(Duration::from_secs(6)) + .expect("electrum should see block"); + + let recorded_auth = Arc::new(Mutex::new(Vec::new())); + let expected_auths = Arc::new(Mutex::new(vec![ + "Bearer test-token-1".to_string(), + "Bearer test-token-1".to_string(), + "Bearer test-token-2".to_string(), + ])); + let proxy_addr = start_auth_proxy( + env.electrsd.electrum_url.clone(), + Some(expected_auths), + recorded_auth.clone(), + ) + .expect("start auth proxy"); + + let token = Arc::new(Mutex::new(Some("Bearer test-token-1".to_string()))); + let auth_provider = { + let token = token.clone(); + Arc::new(move || token.lock().unwrap().clone()) + }; + + let client = RawClient::new( + proxy_addr, + Some(Duration::from_secs(5)), + Some(auth_provider.clone()), + ) + .expect("create client"); + + client.ping().expect("first request ok"); + *token.lock().unwrap() = Some("Bearer test-token-2".to_string()); + client.ping().expect("second request ok"); + + let auths = recorded_auth.lock().unwrap(); + assert_eq!( + auths.as_slice(), + [ + "Bearer test-token-1", + "Bearer test-token-1", + "Bearer test-token-2" + ] + ); +} + +// Ignored by default since it requires Docker/Keycloak available. +#[test] +fn test_auth_with_keycloak() { + use helper::keycloak::TestKeycloak; + + // Start Keycloak container + let kc = TestKeycloak::start().unwrap(); + + // Setup realm/client/user + let token_url = kc + .setup_realm("test-realm", "test-client", "test-user", "password") + .expect("setup realm"); + + // Acquire token + let token = helper::keycloak::TestKeycloak::get_token( + &token_url, + "test-client", + "test-user", + "password", + ) + .expect("get token"); + + // Use token in auth provider + let recorded_auth = Arc::new(Mutex::new(Vec::new())); + let proxy_addr = start_auth_proxy( + TestEnv::new().unwrap().electrsd.electrum_url.clone(), + None, + recorded_auth.clone(), + ) + .expect("start auth proxy"); + + let token_val = Arc::new(Mutex::new(Some(format!("Bearer {}", token)))); + let auth_provider = { + let token_val = token_val.clone(); + Arc::new(move || token_val.lock().unwrap().clone()) + }; + + let client = RawClient::new( + proxy_addr, + Some(Duration::from_secs(5)), + Some(auth_provider), + ) + .expect("create client"); + + client.ping().expect("ping ok"); +} diff --git a/tests/helper/keycloak.rs b/tests/helper/keycloak.rs new file mode 100644 index 0000000..d0f3699 --- /dev/null +++ b/tests/helper/keycloak.rs @@ -0,0 +1,203 @@ +use std::path::Path; +use std::process::Command; +use std::thread::sleep; +use std::time::Duration; + +use serde::{Deserialize, Serialize}; +use testcontainers::core::RunnableImage; +use testcontainers::runners::SyncRunner; +use testcontainers::GenericImage; + +/// Minimal Keycloak test helper that starts a Keycloak container via testcontainers. +pub struct TestKeycloak { + pub base_url: String, +} + +#[derive(Serialize)] +struct CreateUser { + username: String, + enabled: bool, + credentials: Vec, +} + +#[derive(Serialize)] +struct Credential { + #[serde(rename = "type")] + typ: String, + value: String, + temporary: bool, +} + +#[derive(Deserialize)] +struct TokenResponse { + access_token: String, +} + +impl TestKeycloak { + /// Start Keycloak using testcontainers and return a helper bound to the mapped host port. + pub fn start() -> Result> { + if !docker_is_available() { + return Err("Docker is not available or the Docker daemon is not running".into()); + } + + let image = RunnableImage::from(GenericImage::new("quay.io/keycloak/keycloak", "20.0.2")) + .with_args(vec!["start-dev".to_string()]) + .with_env_var(("KEYCLOAK_ADMIN", "admin")) + .with_env_var(("KEYCLOAK_ADMIN_PASSWORD", "admin")) + .with_mapped_port((8080, 8080)); + + let container = image.start().map_err(|e| { + let hint = if docker_is_available() { + "Docker is available, but the container runtime could not create the Keycloak container." + } else { + "Docker does not appear to be available or the Docker daemon is not running." + }; + format!("{} Error: {}", hint, e) + })?; + sleep(Duration::from_secs(10)); + + println!("Loading image"); + let host_port = container.get_host_port_ipv4(8080)?; + let base_url = format!("http://127.0.0.1:{}", host_port); + + Ok(Self { base_url }) + } + + /// Create a realm, client, and user via the admin REST API. + pub fn setup_realm( + &self, + realm: &str, + client_id: &str, + username: &str, + password: &str, + ) -> Result> { + let admin_token = self.get_admin_token()?; + + let realm_payload = serde_json::json!({ "realm": realm, "enabled": true }); + let body = serde_json::to_vec(&realm_payload)?; + bitreq::post(format!("{}/admin/realms", self.base_url)) + .with_header("Authorization", format!("Bearer {}", admin_token)) + .with_header("Content-Type", "application/json") + .with_body(body) + .send()?; + + let client_payload = serde_json::json!({ + "clientId": client_id, + "enabled": true, + "publicClient": true, + "redirectUris": ["*"], + }); + let body = serde_json::to_vec(&client_payload)?; + bitreq::post(format!("{}/admin/realms/{}/clients", self.base_url, realm)) + .with_header("Authorization", format!("Bearer {}", admin_token)) + .with_header("Content-Type", "application/json") + .with_body(body) + .send()?; + + let cred = Credential { + typ: "password".to_string(), + value: password.to_string(), + temporary: false, + }; + let user_payload = CreateUser { + username: username.to_string(), + enabled: true, + credentials: vec![cred], + }; + let body = serde_json::to_vec(&user_payload)?; + bitreq::post(format!("{}/admin/realms/{}/users", self.base_url, realm)) + .with_header("Authorization", format!("Bearer {}", admin_token)) + .with_header("Content-Type", "application/json") + .with_body(body) + .send()?; + + sleep(Duration::from_secs(1)); + + Ok(format!( + "{}/realms/{}/protocol/openid-connect/token", + self.base_url, realm + )) + } + + fn get_admin_token(&self) -> Result> { + let body = build_form_body(&[ + ("grant_type", "password"), + ("client_id", "admin-cli"), + ("username", "admin"), + ("password", "admin"), + ]); + let resp = bitreq::post(format!( + "{}/realms/master/protocol/openid-connect/token", + self.base_url + )) + .with_header("Content-Type", "application/x-www-form-urlencoded") + .with_body(body) + .send()?; + + if !(200..300).contains(&resp.status_code) { + return Err(format!("admin token request failed: {}", resp.status_code).into()); + } + let tr: TokenResponse = resp.json()?; + Ok(tr.access_token) + } + + /// Request an access token using password grant for the configured realm token endpoint. + pub fn get_token( + token_url: &str, + client_id: &str, + username: &str, + password: &str, + ) -> Result> { + let body = build_form_body(&[ + ("grant_type", "password"), + ("client_id", client_id), + ("username", username), + ("password", password), + ]); + let resp = bitreq::post(token_url) + .with_header("Content-Type", "application/x-www-form-urlencoded") + .with_body(body) + .send()?; + + if !(200..300).contains(&resp.status_code) { + return Err(format!("token request failed: {}", resp.status_code).into()); + } + + let tr: TokenResponse = resp.json()?; + Ok(tr.access_token) + } +} + +fn build_form_body(pairs: &[(&str, &str)]) -> String { + pairs + .iter() + .enumerate() + .map(|(idx, (key, value))| { + let prefix = if idx == 0 { + String::new() + } else { + "&".to_string() + }; + format!("{prefix}{key}={value}") + }) + .collect() +} + +fn docker_is_available() -> bool { + if std::env::var_os("DOCKER_HOST").is_some() { + return true; + } + + if ["/var/run/docker.sock", "/run/docker.sock"] + .iter() + .any(|p| Path::new(p).exists()) + { + return true; + } + + Command::new("docker") + .arg("info") + .output() + .map(|output| output.status.success()) + .unwrap_or(false) +} diff --git a/tests/helper/mod.rs b/tests/helper/mod.rs new file mode 100644 index 0000000..94a7ca8 --- /dev/null +++ b/tests/helper/mod.rs @@ -0,0 +1 @@ +pub mod keycloak;