From 226f75eaf2c71f277d41c7290f30da0701800fc8 Mon Sep 17 00:00:00 2001 From: Junyi Ou Date: Thu, 20 Aug 2026 15:49:56 -0400 Subject: [PATCH 1/9] test(dgw): cover credential injection reconnect Add process-level tests for DVLS-like preflight, first inject, jet_reuse reconnect, fail-closed missing mappings, and synthetic KDC generation reuse. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- Cargo.lock | 2 + testsuite/Cargo.toml | 2 + testsuite/src/dgw_config.rs | 7 +- testsuite/tests/cli/dgw/cred_injection.rs | 611 ++++++++++++++++++++++ testsuite/tests/cli/dgw/mod.rs | 1 + 5 files changed, 622 insertions(+), 1 deletion(-) create mode 100644 testsuite/tests/cli/dgw/cred_injection.rs diff --git a/Cargo.lock b/Cargo.lock index f2fa1f931..04dc60798 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7516,6 +7516,8 @@ dependencies = [ "escargot", "expect-test", "fastrand", + "ironrdp-core 0.2.1", + "ironrdp-pdu", "libsql", "mcp-proxy", "network-scanner", diff --git a/testsuite/Cargo.toml b/testsuite/Cargo.toml index 895e12531..74ea00174 100644 --- a/testsuite/Cargo.toml +++ b/testsuite/Cargo.toml @@ -32,6 +32,8 @@ tokio-tungstenite = { version = "0.29", features = ["rustls-tls-native-roots"] } [dev-dependencies] base64 = "0.23" +ironrdp-core = { version = "0.2", features = ["std"] } +ironrdp-pdu = { version = "0.9", features = ["std"] } proxy-socks = { path = "../crates/proxy-socks" } libsql = { version = "0.9", default-features = false, features = ["core"] } mcp-proxy.path = "../crates/mcp-proxy" diff --git a/testsuite/src/dgw_config.rs b/testsuite/src/dgw_config.rs index f5fc7a441..885e71e48 100644 --- a/testsuite/src/dgw_config.rs +++ b/testsuite/src/dgw_config.rs @@ -45,6 +45,9 @@ pub struct DgwConfig { /// Enable unstable features. #[builder(default = false)] enable_unstable: bool, + /// Enable Kerberos credential injection (also requires `enable_unstable`). + #[builder(default = false)] + kerberos_credential_injection: bool, /// Override the recording path in the gateway config. /// /// When `None`, the gateway uses its default (`/recordings`). @@ -84,6 +87,7 @@ impl DgwConfigHandle { disable_token_validation, verbosity_profile, enable_unstable, + kerberos_credential_injection, recording_path, agent_tunnel, } = config; @@ -137,7 +141,8 @@ impl DgwConfigHandle { "VerbosityProfile": "{verbosity_profile}", "__debug__": {{ "disable_token_validation": {disable_token_validation}, - "enable_unstable": {enable_unstable} + "enable_unstable": {enable_unstable}, + "kerberos_credential_injection": {kerberos_credential_injection} }}{recording_path_json}{agent_tunnel_json} }}"# ); diff --git a/testsuite/tests/cli/dgw/cred_injection.rs b/testsuite/tests/cli/dgw/cred_injection.rs new file mode 100644 index 000000000..01bcfed94 --- /dev/null +++ b/testsuite/tests/cli/dgw/cred_injection.rs @@ -0,0 +1,611 @@ +//! Process-level tests for RDP credential injection, reconnect, and fail-closed routing. +//! +//! These tests start a real Gateway, provision credentials over `/jet/preflight` the way DVLS +//! does, then connect to the TCP listener with an RDP preconnection blob. A loopback peer stands +//! in for the destination RDP server and records the X.224 Connection Request the proxy forwards. +//! Injection is observed from Gateway logs and from the rewritten mstshash cookie. CredSSP is not +//! completed: the contract under test is checkout, reconnect, and fail-closed routing. + +use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant}; + +use anyhow::Context as _; +use base64::Engine as _; +use testsuite::cli::{dgw_tokio_cmd, wait_for_tcp_port}; +use testsuite::dgw_config::{DgwConfig, DgwConfigHandle, VerbosityProfile}; +use tokio::io::{AsyncBufReadExt as _, AsyncReadExt as _, AsyncWriteExt as _, BufReader}; +use tokio::net::{TcpListener, TcpStream}; +use tokio::process::Child; + +const CLIENT_COOKIE: &str = "client-cookie-user"; +const TARGET_USER: &str = "injected-target-user"; +const PROXY_USER: &str = "injected-proxy-user"; +const KERBEROS_TARGET_USER: &str = "administrator@example.invalid"; +const INJECT_LOG: &str = "RDP-TLS forwarding with credential injection"; +const FORWARD_LOG: &str = "Upstream forwarding"; +const MISSING_LOG: &str = "missing or expired; re-provision to retry"; +const PUBLISHED_KDC_LOG: &str = "Published synthetic KDC"; +const REGISTERED_KDC_LOG: &str = "Registered synthetic KDC for credential-injection session"; + +fn next_id() -> String { + static COUNTER: AtomicU64 = AtomicU64::new(1); + let n = COUNTER.fetch_add(1, Ordering::Relaxed); + format!("00000000-0000-4000-a000-{n:012x}") +} + +fn unsigned_jws(header: serde_json::Value, payload: serde_json::Value) -> anyhow::Result { + let engine = base64::engine::general_purpose::URL_SAFE_NO_PAD; + let header = engine.encode(serde_json::to_vec(&header).context("serialize JWT header")?); + let payload = engine.encode(serde_json::to_vec(&payload).context("serialize JWT payload")?); + Ok(format!("{header}.{payload}.ZHVtbXlfc2lnbmF0dXJl")) +} + +fn preflight_scope_token() -> anyhow::Result { + unsigned_jws( + serde_json::json!({"alg":"RS256","typ":"JWT","cty":"SCOPE"}), + serde_json::json!({ + "scope": "gateway.preflight", + "exp": 9_999_999_999_i64, + "jti": next_id(), + }), + ) +} + +fn association_token(jti: &str, jet_aid: &str, dest_port: u16, jet_reuse: u32) -> anyhow::Result { + unsigned_jws( + serde_json::json!({"alg":"RS256","typ":"JWT","cty":"ASSOCIATION"}), + serde_json::json!({ + "dst_hst": format!("127.0.0.1:{dest_port}"), + "exp": 9_999_999_999_i64, + "jet_aid": jet_aid, + "jet_ap": "rdp", + "jet_cm": "fwd", + "jet_rec": "none", + "jet_reuse": jet_reuse, + "jti": jti, + "nbf": 0, + }), + ) +} + +fn encode_pcb(token: &str) -> anyhow::Result> { + let pcb = ironrdp_pdu::pcb::PreconnectionBlob { + version: ironrdp_pdu::pcb::PcbVersion::V2, + id: 0, + v2_payload: Some(token.to_owned()), + }; + ironrdp_core::encode_vec(&pcb).context("encode preconnection blob") +} + +fn encode_connection_request(cookie: &str) -> anyhow::Result> { + use ironrdp_pdu::nego::{ConnectionRequest, NegoRequestData, RequestFlags, SecurityProtocol}; + use ironrdp_pdu::x224::X224; + + let pdu = X224(ConnectionRequest { + nego_data: Some(NegoRequestData::cookie(cookie.to_owned())), + flags: RequestFlags::empty(), + protocol: SecurityProtocol::HYBRID | SecurityProtocol::HYBRID_EX | SecurityProtocol::SSL, + }); + ironrdp_core::encode_vec(&pdu).context("encode X.224 connection request") +} + +fn strip_ansi(input: &str) -> String { + let mut out = String::with_capacity(input.len()); + let mut chars = input.chars().peekable(); + while let Some(c) = chars.next() { + if c == '\u{1b}' && chars.peek() == Some(&'[') { + chars.next(); + for next in chars.by_ref() { + if next.is_ascii_alphabetic() { + break; + } + } + } else { + out.push(c); + } + } + out +} + +struct LogBuffer(Arc>); + +impl LogBuffer { + fn new() -> Self { + Self(Arc::new(Mutex::new(String::new()))) + } + + fn snapshot(&self) -> String { + strip_ansi(&self.0.lock().expect("log mutex")) + } + + async fn wait_contains(&self, needle: &str) -> anyhow::Result { + self.wait_count(needle, 1).await + } + + async fn wait_count(&self, needle: &str, count: usize) -> anyhow::Result { + let deadline = Instant::now() + Duration::from_secs(15); + loop { + let snapshot = self.snapshot(); + if snapshot.matches(needle).count() >= count { + return Ok(snapshot); + } + if Instant::now() >= deadline { + anyhow::bail!("timed out waiting for {count} occurrence(s) of {needle:?}; logs:\n{snapshot}"); + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + } +} + +struct FakeRdpTarget { + port: u16, + accepted: Arc, + payloads: Arc>>>, +} + +impl FakeRdpTarget { + async fn start() -> anyhow::Result { + let listener = TcpListener::bind("127.0.0.1:0").await.context("bind fake RDP target")?; + let port = listener.local_addr().context("fake RDP local_addr")?.port(); + let accepted = Arc::new(AtomicUsize::new(0)); + let payloads = Arc::new(Mutex::new(Vec::new())); + let accepted_task = Arc::clone(&accepted); + let payloads_task = Arc::clone(&payloads); + + tokio::spawn(async move { + loop { + let Ok((mut stream, _)) = listener.accept().await else { + break; + }; + accepted_task.fetch_add(1, Ordering::SeqCst); + let payloads = Arc::clone(&payloads_task); + tokio::spawn(async move { + let mut buf = vec![0_u8; 4096]; + // CredSSP cert generation can delay the rewritten X.224 CR. + if let Ok(Ok(n)) = tokio::time::timeout(Duration::from_secs(30), stream.read(&mut buf)).await + && n > 0 + { + payloads.lock().expect("payload mutex").push(buf[..n].to_vec()); + } + // Keep the accepted socket open so the proxy can finish writing the CR. + tokio::time::sleep(Duration::from_secs(30)).await; + }); + } + }); + + Ok(Self { + port, + accepted, + payloads, + }) + } + + fn accepted(&self) -> usize { + self.accepted.load(Ordering::SeqCst) + } + + async fn wait_payloads(&self, count: usize) -> anyhow::Result>> { + let deadline = Instant::now() + Duration::from_secs(30); + loop { + { + let payloads = self.payloads.lock().expect("payload mutex"); + if payloads.len() >= count { + return Ok(payloads.clone()); + } + } + if Instant::now() >= deadline { + anyhow::bail!( + "timed out waiting for {count} target payload(s); accepted={}", + self.accepted() + ); + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + } +} + +struct GatewayProc { + config: DgwConfigHandle, + process: Child, + logs: LogBuffer, +} + +impl GatewayProc { + async fn start(kerberos: bool) -> anyhow::Result { + let config = DgwConfig::builder() + .disable_token_validation(true) + .verbosity_profile(VerbosityProfile::DEBUG) + .enable_unstable(kerberos) + .kerberos_credential_injection(kerberos) + .build() + .init() + .context("init gateway config")?; + + let mut process = dgw_tokio_cmd() + .env("DGATEWAY_CONFIG_PATH", config.config_dir()) + .env("RUST_LOG", "devolutions_gateway=debug") + .env("NO_COLOR", "1") + .kill_on_drop(true) + .stdout(std::process::Stdio::piped()) + .stderr(std::process::Stdio::piped()) + .spawn() + .context("start Devolutions Gateway")?; + + let logs = LogBuffer::new(); + spawn_stdio_collector(process.stdout.take(), Arc::clone(&logs.0)); + spawn_stdio_collector(process.stderr.take(), Arc::clone(&logs.0)); + + wait_for_tcp_port(config.http_port()) + .await + .context("wait for gateway HTTP port")?; + + Ok(Self { config, process, logs }) + } +} + +fn spawn_stdio_collector(stream: Option, logs: Arc>) +where + R: tokio::io::AsyncRead + Unpin + Send + 'static, +{ + let Some(stream) = stream else { + return; + }; + tokio::spawn(async move { + let mut reader = BufReader::new(stream); + let mut line = String::new(); + loop { + line.clear(); + match reader.read_line(&mut line).await { + Ok(0) | Err(_) => break, + Ok(_) => logs.lock().expect("log mutex").push_str(&line), + } + } + }); +} + +async fn post_preflight(http_port: u16, operations: serde_json::Value) -> anyhow::Result { + let bearer = preflight_scope_token()?; + let body = serde_json::to_string(&operations).context("serialize preflight body")?; + let request = format!( + "POST /jet/preflight HTTP/1.1\r\n\ + Host: 127.0.0.1:{http_port}\r\n\ + Content-Type: application/json\r\n\ + Authorization: Bearer {bearer}\r\n\ + Content-Length: {}\r\n\ + Connection: close\r\n\ + \r\n\ + {body}", + body.len() + ); + + let mut stream = TcpStream::connect(("127.0.0.1", http_port)) + .await + .context("connect to gateway HTTP")?; + stream.write_all(request.as_bytes()).await.context("write preflight")?; + stream.flush().await.context("flush preflight")?; + + let mut reader = BufReader::new(stream); + let mut status_line = String::new(); + reader + .read_line(&mut status_line) + .await + .context("read preflight status")?; + anyhow::ensure!(status_line.contains("200"), "preflight HTTP status was {status_line:?}"); + + let mut content_length = None; + loop { + let mut line = String::new(); + reader.read_line(&mut line).await.context("read preflight header")?; + if line == "\r\n" || line.is_empty() { + break; + } + if let Some(value) = line + .split_once(':') + .filter(|(name, _)| name.eq_ignore_ascii_case("content-length")) + .map(|(_, value)| value.trim().to_owned()) + { + content_length = Some(value.parse::().context("parse Content-Length")?); + } + } + + let response_body = if let Some(len) = content_length { + let mut buf = vec![0_u8; len]; + reader.read_exact(&mut buf).await.context("read preflight body")?; + String::from_utf8(buf).context("preflight body utf-8")? + } else { + let mut buf = String::new(); + reader + .read_to_string(&mut buf) + .await + .context("read preflight eof body")?; + buf + }; + + let json: serde_json::Value = + serde_json::from_str(&response_body).with_context(|| format!("parse preflight JSON: {response_body}"))?; + let outputs = json.as_array().context("preflight response is not an array")?; + // Re-provisioning the same JTI emits an info alert, then still acks. + let mut acked = 0_usize; + for output in outputs { + match output["kind"].as_str() { + Some("ack") => acked += 1, + Some("alert") if output["alert_status"] == "info" => {} + _ => anyhow::bail!("preflight operation was not ack: {output}"), + } + } + anyhow::ensure!(acked > 0, "preflight returned no ack: {json}"); + Ok(json) +} + +async fn provision_credentials( + http_port: u16, + token: &str, + target_username: &str, + time_to_live: u32, + krb_kdc: Option<&str>, +) -> anyhow::Result<()> { + let mut operations = vec![serde_json::json!({ + "id": next_id(), + "kind": "provision-credentials", + "token": token, + "proxy_credential": { + "kind": "username-password", + "username": PROXY_USER, + "password": "proxy-secret" + }, + "target_credential": { + "kind": "username-password", + "username": target_username, + "password": "target-secret" + }, + "time_to_live": time_to_live + })]; + + if let Some(krb_kdc) = krb_kdc { + operations.push(serde_json::json!({ + "id": next_id(), + "kind": "provision-connection-options", + "token": token, + "connection_options": { "krb_kdc": krb_kdc }, + "time_to_live": time_to_live + })); + } + + post_preflight(http_port, serde_json::Value::Array(operations)).await?; + Ok(()) +} + +async fn connect_rdp_client(gateway_tcp: u16, association_jwt: &str) -> anyhow::Result { + let mut stream = TcpStream::connect(("127.0.0.1", gateway_tcp)) + .await + .context("connect to gateway TCP")?; + stream + .write_all(&encode_pcb(association_jwt)?) + .await + .context("write preconnection blob")?; + stream + .write_all(&encode_connection_request(CLIENT_COOKIE)?) + .await + .context("write connection request")?; + stream.flush().await.context("flush RDP client")?; + Ok(stream) +} + +fn cookie_line(username: &str) -> String { + format!("Cookie: mstshash={username}") +} + +fn payloads_contain(payloads: &[Vec], needle: &str) -> bool { + payloads + .iter() + .any(|payload| String::from_utf8_lossy(payload).contains(needle)) +} + +#[tokio::test] +async fn first_rdp_connection_injects_ntlm() -> anyhow::Result<()> { + let target = FakeRdpTarget::start().await?; + let mut gateway = GatewayProc::start(false).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token(&jti, &jet_aid, target.port, 60)?; + provision_credentials(gateway.config.http_port(), &token, TARGET_USER, 300, None).await?; + + let _client = connect_rdp_client(gateway.config.tcp_port(), &token).await?; + let logs = gateway.logs.wait_contains(INJECT_LOG).await?; + assert!( + logs.contains("kerberos=false"), + "expected NTLM injection; logs:\n{logs}" + ); + assert!( + !logs.contains(FORWARD_LOG), + "injection must not fall back to ordinary forward; logs:\n{logs}" + ); + + let payloads = target.wait_payloads(1).await?; + assert!( + payloads_contain(&payloads, &cookie_line(TARGET_USER)), + "target should see injected cookie; payloads={payloads:?}" + ); + assert!( + !payloads_contain(&payloads, &cookie_line(CLIENT_COOKIE)), + "target must not see the client cookie; payloads={payloads:?}" + ); + + let _ = gateway.process.start_kill(); + Ok(()) +} + +#[tokio::test] +async fn reconnect_same_jwt_still_injects() -> anyhow::Result<()> { + let target = FakeRdpTarget::start().await?; + let mut gateway = GatewayProc::start(false).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token(&jti, &jet_aid, target.port, 60)?; + provision_credentials(gateway.config.http_port(), &token, TARGET_USER, 300, None).await?; + + let first = connect_rdp_client(gateway.config.tcp_port(), &token).await?; + gateway.logs.wait_count(INJECT_LOG, 1).await?; + target.wait_payloads(1).await?; + drop(first); + + let _second = connect_rdp_client(gateway.config.tcp_port(), &token).await?; + let logs = gateway.logs.wait_count(INJECT_LOG, 2).await?; + assert!( + !logs.contains(FORWARD_LOG), + "reconnect must keep injecting, not ordinary-forward; logs:\n{logs}" + ); + + let payloads = target.wait_payloads(2).await?; + assert_eq!(payloads.len(), 2, "both connections should reach the fake RDP target"); + assert!( + payloads + .iter() + .all(|payload| String::from_utf8_lossy(payload).contains(&cookie_line(TARGET_USER))), + "both reconnects should inject the target cookie; payloads={payloads:?}" + ); + + let _ = gateway.process.start_kill(); + Ok(()) +} + +#[tokio::test] +async fn required_missing_fails_closed() -> anyhow::Result<()> { + let target = FakeRdpTarget::start().await?; + let mut gateway = GatewayProc::start(false).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token(&jti, &jet_aid, target.port, 60)?; + provision_credentials(gateway.config.http_port(), &token, TARGET_USER, 1, None).await?; + tokio::time::sleep(Duration::from_secs(2)).await; + + let _client = connect_rdp_client(gateway.config.tcp_port(), &token).await?; + let logs = gateway.logs.wait_contains(MISSING_LOG).await?; + assert!( + !logs.contains(INJECT_LOG), + "expired mapping must not inject; logs:\n{logs}" + ); + assert!( + !logs.contains(FORWARD_LOG), + "expired mapping must fail closed, never silent ordinary forward; logs:\n{logs}" + ); + + tokio::time::sleep(Duration::from_millis(500)).await; + assert_eq!(target.accepted(), 0, "fail-closed routing must not connect upstream"); + + let _ = gateway.process.start_kill(); + Ok(()) +} + +#[tokio::test] +async fn unprovisioned_rdp_uses_ordinary_forward() -> anyhow::Result<()> { + let target = FakeRdpTarget::start().await?; + let mut gateway = GatewayProc::start(false).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token(&jti, &jet_aid, target.port, 60)?; + + let _client = connect_rdp_client(gateway.config.tcp_port(), &token).await?; + let logs = gateway.logs.wait_contains(FORWARD_LOG).await?; + assert!( + !logs.contains(INJECT_LOG), + "absent mapping should ordinary-forward; logs:\n{logs}" + ); + + let payloads = target.wait_payloads(1).await?; + assert!( + payloads_contain(&payloads, &cookie_line(CLIENT_COOKIE)), + "ordinary forward should keep the client cookie; payloads={payloads:?}" + ); + assert!( + !payloads_contain(&payloads, &cookie_line(TARGET_USER)), + "ordinary forward must not invent an injection cookie; payloads={payloads:?}" + ); + + let _ = gateway.process.start_kill(); + Ok(()) +} + +#[tokio::test] +async fn kerberos_reconnect_reuses_generation_until_reprovision() -> anyhow::Result<()> { + let target = FakeRdpTarget::start().await?; + let mut gateway = GatewayProc::start(true).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token(&jti, &jet_aid, target.port, 60)?; + provision_credentials( + gateway.config.http_port(), + &token, + KERBEROS_TARGET_USER, + 300, + Some("tcp://127.0.0.1:88"), + ) + .await?; + + // Keep overlapping reconnect sockets so the same-generation KDC lease stays live. + let _first = connect_rdp_client(gateway.config.tcp_port(), &token).await?; + let logs = gateway.logs.wait_count(INJECT_LOG, 1).await?; + assert!( + logs.contains("kerberos=true"), + "expected Kerberos injection; logs:\n{logs}" + ); + gateway.logs.wait_count(PUBLISHED_KDC_LOG, 1).await?; + + let _second = connect_rdp_client(gateway.config.tcp_port(), &token).await?; + let logs = gateway.logs.wait_count(INJECT_LOG, 2).await?; + assert_eq!( + logs.matches("kerberos=true").count(), + 2, + "same generation reconnect should still inject Kerberos; logs:\n{logs}" + ); + gateway.logs.wait_count(REGISTERED_KDC_LOG, 2).await?; + assert_eq!( + logs.matches(PUBLISHED_KDC_LOG).count(), + 1, + "same provisioning generation must reuse the interned synthetic KDC; logs:\n{logs}" + ); + + provision_credentials( + gateway.config.http_port(), + &token, + KERBEROS_TARGET_USER, + 300, + Some("tcp://127.0.0.1:88"), + ) + .await?; + + let _third = connect_rdp_client(gateway.config.tcp_port(), &token).await?; + let logs = gateway.logs.wait_count(INJECT_LOG, 3).await?; + assert_eq!( + logs.matches("kerberos=true").count(), + 3, + "newer provisioning generation should replace and still inject; logs:\n{logs}" + ); + let logs = gateway.logs.wait_count(PUBLISHED_KDC_LOG, 2).await?; + assert_eq!( + logs.matches(REGISTERED_KDC_LOG).count(), + 3, + "each connection should register a synthetic KDC lease; logs:\n{logs}" + ); + assert!( + !logs.contains(FORWARD_LOG), + "Kerberos injection must not ordinary-forward; logs:\n{logs}" + ); + + let payloads = target.wait_payloads(3).await?; + assert!( + payloads + .iter() + .all(|payload| String::from_utf8_lossy(payload).contains(&cookie_line(KERBEROS_TARGET_USER))), + "each generation should inject the Kerberos target username; payloads={payloads:?}" + ); + + let _ = gateway.process.start_kill(); + Ok(()) +} diff --git a/testsuite/tests/cli/dgw/mod.rs b/testsuite/tests/cli/dgw/mod.rs index c6cc1226b..f6e88737a 100644 --- a/testsuite/tests/cli/dgw/mod.rs +++ b/testsuite/tests/cli/dgw/mod.rs @@ -1,5 +1,6 @@ mod benign_disconnect; mod cli_args; +mod cred_injection; mod heartbeat; mod preflight; mod tls_anchoring; From 976b5777e41fb875edbebe89815e96eaf6b1ce7c Mon Sep 17 00:00:00 2001 From: Junyi Ou Date: Thu, 20 Aug 2026 15:57:07 -0400 Subject: [PATCH 2/9] style(dgw): fix clippy literal suffixes in injection e2e Clippy separated_literal_suffix failed CI lints on the stacked reconnect tests. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- testsuite/tests/cli/dgw/cred_injection.rs | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/testsuite/tests/cli/dgw/cred_injection.rs b/testsuite/tests/cli/dgw/cred_injection.rs index 01bcfed94..b58403f98 100644 --- a/testsuite/tests/cli/dgw/cred_injection.rs +++ b/testsuite/tests/cli/dgw/cred_injection.rs @@ -46,7 +46,7 @@ fn preflight_scope_token() -> anyhow::Result { serde_json::json!({"alg":"RS256","typ":"JWT","cty":"SCOPE"}), serde_json::json!({ "scope": "gateway.preflight", - "exp": 9_999_999_999_i64, + "exp": 9_999_999_999i64, "jti": next_id(), }), ) @@ -57,7 +57,7 @@ fn association_token(jti: &str, jet_aid: &str, dest_port: u16, jet_reuse: u32) - serde_json::json!({"alg":"RS256","typ":"JWT","cty":"ASSOCIATION"}), serde_json::json!({ "dst_hst": format!("127.0.0.1:{dest_port}"), - "exp": 9_999_999_999_i64, + "exp": 9_999_999_999i64, "jet_aid": jet_aid, "jet_ap": "rdp", "jet_cm": "fwd", @@ -161,7 +161,7 @@ impl FakeRdpTarget { accepted_task.fetch_add(1, Ordering::SeqCst); let payloads = Arc::clone(&payloads_task); tokio::spawn(async move { - let mut buf = vec![0_u8; 4096]; + let mut buf = vec![0u8; 4096]; // CredSSP cert generation can delay the rewritten X.224 CR. if let Ok(Ok(n)) = tokio::time::timeout(Duration::from_secs(30), stream.read(&mut buf)).await && n > 0 @@ -310,7 +310,7 @@ async fn post_preflight(http_port: u16, operations: serde_json::Value) -> anyhow } let response_body = if let Some(len) = content_length { - let mut buf = vec![0_u8; len]; + let mut buf = vec![0u8; len]; reader.read_exact(&mut buf).await.context("read preflight body")?; String::from_utf8(buf).context("preflight body utf-8")? } else { @@ -326,7 +326,7 @@ async fn post_preflight(http_port: u16, operations: serde_json::Value) -> anyhow serde_json::from_str(&response_body).with_context(|| format!("parse preflight JSON: {response_body}"))?; let outputs = json.as_array().context("preflight response is not an array")?; // Re-provisioning the same JTI emits an info alert, then still acks. - let mut acked = 0_usize; + let mut acked = 0usize; for output in outputs { match output["kind"].as_str() { Some("ack") => acked += 1, From 70b21f69ae54d60a79ce4398938890211240c6b8 Mon Sep 17 00:00:00 2001 From: Junyi Ou Date: Thu, 20 Aug 2026 17:06:47 -0400 Subject: [PATCH 3/9] test(dgw): complete Kerberos injection against a mock KDC Drive CredSSP through a TCP kdc crate and IronRDP acceptor so target-leg Kerberos injection is proven, not just log-matched. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- Cargo.lock | 6 + testsuite/Cargo.toml | 6 + testsuite/tests/cli/dgw/cred_injection.rs | 44 +- testsuite/tests/cli/dgw/cred_injection_kdc.rs | 634 ++++++++++++++++++ testsuite/tests/cli/dgw/mod.rs | 1 + 5 files changed, 670 insertions(+), 21 deletions(-) create mode 100644 testsuite/tests/cli/dgw/cred_injection_kdc.rs diff --git a/Cargo.lock b/Cargo.lock index 04dc60798..b3f239ccc 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7516,12 +7516,17 @@ dependencies = [ "escargot", "expect-test", "fastrand", + "ironrdp-acceptor", + "ironrdp-connector", "ironrdp-core 0.2.1", "ironrdp-pdu", + "ironrdp-tokio", + "kdc", "libsql", "mcp-proxy", "network-scanner", "network-scanner-proto", + "picky-krb", "proxy-socks", "rstest", "serde", @@ -7536,6 +7541,7 @@ dependencies = [ "tokio-tungstenite", "tokio-util", "typed-builder", + "x509-cert 0.3.0", ] [[package]] diff --git a/testsuite/Cargo.toml b/testsuite/Cargo.toml index 74ea00174..5c023c919 100644 --- a/testsuite/Cargo.toml +++ b/testsuite/Cargo.toml @@ -32,8 +32,13 @@ tokio-tungstenite = { version = "0.29", features = ["rustls-tls-native-roots"] } [dev-dependencies] base64 = "0.23" +ironrdp-acceptor = "0.10" +ironrdp-connector = "0.10" ironrdp-core = { version = "0.2", features = ["std"] } ironrdp-pdu = { version = "0.9", features = ["std"] } +ironrdp-tokio = "0.10" +kdc = "0.1" +picky-krb = "0.12" proxy-socks = { path = "../crates/proxy-socks" } libsql = { version = "0.9", default-features = false, features = ["core"] } mcp-proxy.path = "../crates/mcp-proxy" @@ -45,6 +50,7 @@ sysevent.path = "../crates/sysevent" tempfile = "3" test-utils.path = "../crates/test-utils" tokio-rustls = { version = "0.26", features = ["ring"] } +x509-cert = { version = "0.3", default-features = false, features = ["std"] } [target.'cfg(unix)'.dev-dependencies] sysevent-syslog.path = "../crates/sysevent-syslog" diff --git a/testsuite/tests/cli/dgw/cred_injection.rs b/testsuite/tests/cli/dgw/cred_injection.rs index b58403f98..11fafd805 100644 --- a/testsuite/tests/cli/dgw/cred_injection.rs +++ b/testsuite/tests/cli/dgw/cred_injection.rs @@ -18,23 +18,25 @@ use tokio::io::{AsyncBufReadExt as _, AsyncReadExt as _, AsyncWriteExt as _, Buf use tokio::net::{TcpListener, TcpStream}; use tokio::process::Child; -const CLIENT_COOKIE: &str = "client-cookie-user"; -const TARGET_USER: &str = "injected-target-user"; -const PROXY_USER: &str = "injected-proxy-user"; -const KERBEROS_TARGET_USER: &str = "administrator@example.invalid"; -const INJECT_LOG: &str = "RDP-TLS forwarding with credential injection"; +pub(crate) const CLIENT_COOKIE: &str = "client-cookie-user"; +pub(crate) const TARGET_USER: &str = "injected-target-user"; +pub(crate) const PROXY_USER: &str = "injected-proxy-user"; +pub(crate) const PROXY_PASSWORD: &str = "proxy-secret"; +pub(crate) const TARGET_PASSWORD: &str = "target-secret"; +pub(crate) const KERBEROS_TARGET_USER: &str = "administrator@example.invalid"; +pub(crate) const INJECT_LOG: &str = "RDP-TLS forwarding with credential injection"; const FORWARD_LOG: &str = "Upstream forwarding"; const MISSING_LOG: &str = "missing or expired; re-provision to retry"; const PUBLISHED_KDC_LOG: &str = "Published synthetic KDC"; const REGISTERED_KDC_LOG: &str = "Registered synthetic KDC for credential-injection session"; -fn next_id() -> String { +pub(crate) fn next_id() -> String { static COUNTER: AtomicU64 = AtomicU64::new(1); let n = COUNTER.fetch_add(1, Ordering::Relaxed); format!("00000000-0000-4000-a000-{n:012x}") } -fn unsigned_jws(header: serde_json::Value, payload: serde_json::Value) -> anyhow::Result { +pub(crate) fn unsigned_jws(header: serde_json::Value, payload: serde_json::Value) -> anyhow::Result { let engine = base64::engine::general_purpose::URL_SAFE_NO_PAD; let header = engine.encode(serde_json::to_vec(&header).context("serialize JWT header")?); let payload = engine.encode(serde_json::to_vec(&payload).context("serialize JWT payload")?); @@ -52,7 +54,7 @@ fn preflight_scope_token() -> anyhow::Result { ) } -fn association_token(jti: &str, jet_aid: &str, dest_port: u16, jet_reuse: u32) -> anyhow::Result { +pub(crate) fn association_token(jti: &str, jet_aid: &str, dest_port: u16, jet_reuse: u32) -> anyhow::Result { unsigned_jws( serde_json::json!({"alg":"RS256","typ":"JWT","cty":"ASSOCIATION"}), serde_json::json!({ @@ -69,7 +71,7 @@ fn association_token(jti: &str, jet_aid: &str, dest_port: u16, jet_reuse: u32) - ) } -fn encode_pcb(token: &str) -> anyhow::Result> { +pub(crate) fn encode_pcb(token: &str) -> anyhow::Result> { let pcb = ironrdp_pdu::pcb::PreconnectionBlob { version: ironrdp_pdu::pcb::PcbVersion::V2, id: 0, @@ -108,22 +110,22 @@ fn strip_ansi(input: &str) -> String { out } -struct LogBuffer(Arc>); +pub(crate) struct LogBuffer(Arc>); impl LogBuffer { fn new() -> Self { Self(Arc::new(Mutex::new(String::new()))) } - fn snapshot(&self) -> String { + pub(crate) fn snapshot(&self) -> String { strip_ansi(&self.0.lock().expect("log mutex")) } - async fn wait_contains(&self, needle: &str) -> anyhow::Result { + pub(crate) async fn wait_contains(&self, needle: &str) -> anyhow::Result { self.wait_count(needle, 1).await } - async fn wait_count(&self, needle: &str, count: usize) -> anyhow::Result { + pub(crate) async fn wait_count(&self, needle: &str, count: usize) -> anyhow::Result { let deadline = Instant::now() + Duration::from_secs(15); loop { let snapshot = self.snapshot(); @@ -205,14 +207,14 @@ impl FakeRdpTarget { } } -struct GatewayProc { - config: DgwConfigHandle, - process: Child, - logs: LogBuffer, +pub(crate) struct GatewayProc { + pub(crate) config: DgwConfigHandle, + pub(crate) process: Child, + pub(crate) logs: LogBuffer, } impl GatewayProc { - async fn start(kerberos: bool) -> anyhow::Result { + pub(crate) async fn start(kerberos: bool) -> anyhow::Result { let config = DgwConfig::builder() .disable_token_validation(true) .verbosity_profile(VerbosityProfile::DEBUG) @@ -338,7 +340,7 @@ async fn post_preflight(http_port: u16, operations: serde_json::Value) -> anyhow Ok(json) } -async fn provision_credentials( +pub(crate) async fn provision_credentials( http_port: u16, token: &str, target_username: &str, @@ -352,12 +354,12 @@ async fn provision_credentials( "proxy_credential": { "kind": "username-password", "username": PROXY_USER, - "password": "proxy-secret" + "password": PROXY_PASSWORD }, "target_credential": { "kind": "username-password", "username": target_username, - "password": "target-secret" + "password": TARGET_PASSWORD }, "time_to_live": time_to_live })]; diff --git a/testsuite/tests/cli/dgw/cred_injection_kdc.rs b/testsuite/tests/cli/dgw/cred_injection_kdc.rs new file mode 100644 index 000000000..38846da56 --- /dev/null +++ b/testsuite/tests/cli/dgw/cred_injection_kdc.rs @@ -0,0 +1,634 @@ +//! Kerberos credential injection against a mock KDC and IronRDP CredSSP server. +//! +//! Proves the target-leg path: Gateway fetches tickets from a TCP KDC (`kdc` crate from +//! sspi-rs) and completes CredSSP with a fake RDP acceptor. The Gateway-facing client uses +//! NTLM so the test does not depend on the in-process synthetic KDC. + +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; +use std::time::{Duration, Instant}; + +use anyhow::Context as _; +use ironrdp_connector::sspi; +use ironrdp_connector::sspi::generator::GeneratorState; +use ironrdp_pdu::nego::{ + ConnectionConfirm, ConnectionRequest, NegoRequestData, RequestFlags, ResponseFlags, SecurityProtocol, +}; +use ironrdp_pdu::x224::X224; +use ironrdp_tokio::{FramedWrite as _, TokioFramed}; +use picky_krb::messages::KdcProxyMessage; +use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _}; +use tokio::net::{TcpListener, TcpStream}; +use tokio_rustls::rustls::pki_types::pem::PemObject as _; +use tokio_rustls::rustls::pki_types::{CertificateDer, PrivateKeyDer, ServerName}; +use tokio_rustls::rustls::{ClientConfig, ServerConfig}; +use x509_cert::der::Decode as _; + +use super::cred_injection::{ + GatewayProc, KERBEROS_TARGET_USER, PROXY_PASSWORD, PROXY_USER, TARGET_PASSWORD, encode_pcb, next_id, + provision_credentials, unsigned_jws, +}; + +const REALM: &str = "EXAMPLE.INVALID"; +// sspi-rs downgrades Negotiate to NTLM when the SPN host is an IP address. +const SERVICE_HOST: &str = "localhost"; +const KRBTGT_KEY: [u8; 32] = [0x11; 32]; +const TERMSRV_KEY: [u8; 32] = [0x22; 32]; + +struct MockKdc { + port: u16, + exchanges: Arc, +} + +impl MockKdc { + async fn start() -> anyhow::Result { + let listener = TcpListener::bind("127.0.0.1:0").await.context("bind mock KDC")?; + let port = listener.local_addr().context("mock KDC local_addr")?.port(); + let exchanges = Arc::new(AtomicUsize::new(0)); + let exchanges_task = Arc::clone(&exchanges); + let config = kdc_config(); + + tokio::spawn(async move { + loop { + let Ok((stream, _)) = listener.accept().await else { + break; + }; + let config = config.clone(); + let exchanges = Arc::clone(&exchanges_task); + tokio::spawn(async move { + match serve_kdc_exchange(stream, &config).await { + Ok(()) => { + exchanges.fetch_add(1, Ordering::SeqCst); + } + Err(error) => eprintln!("mock KDC exchange failed: {error:#}"), + } + }); + } + }); + + Ok(Self { port, exchanges }) + } + + fn url(&self) -> String { + format!("tcp://127.0.0.1:{}", self.port) + } + + fn exchanges(&self) -> usize { + self.exchanges.load(Ordering::SeqCst) + } +} + +fn kdc_config() -> kdc::config::KerberosServer { + let username = format!("administrator@{REALM}"); + kdc::config::KerberosServer { + realm: REALM.to_owned(), + users: vec![kdc::config::DomainUser { + username, + password: TARGET_PASSWORD.to_owned(), + salt: format!("{}administrator", REALM.to_ascii_uppercase()), + }], + max_time_skew: 300, + krbtgt_key: KRBTGT_KEY.to_vec(), + ticket_decryption_key: Some(TERMSRV_KEY.to_vec()), + service_user: None, + } +} + +async fn serve_kdc_exchange(mut stream: TcpStream, config: &kdc::config::KerberosServer) -> anyhow::Result<()> { + let mut len_buf = [0u8; 4]; + stream.read_exact(&mut len_buf).await.context("read KDC length")?; + let len = usize::try_from(u32::from_be_bytes(len_buf)).context("KDC length")?; + let mut body = vec![0u8; len]; + stream.read_exact(&mut body).await.context("read KDC body")?; + + let mut raw = Vec::with_capacity(4 + len); + raw.extend_from_slice(&len_buf); + raw.extend_from_slice(&body); + + let request = KdcProxyMessage::from_raw_kerb_message(&raw).context("wrap KDC TCP payload")?; + let reply = kdc::handle_kdc_proxy_message(request, config, SERVICE_HOST).context("handle KDC message")?; + stream + .write_all(&reply.kerb_message.0.0) + .await + .context("write KDC reply")?; + Ok(()) +} + +struct MockRdp { + port: u16, + credssp_ok: Arc, +} + +impl MockRdp { + async fn start(kdc_url: String) -> anyhow::Result { + install_crypto_provider(); + // Dual-stack so Windows `localhost` (IPv6 first) still hits the fake server. + let listener = match TcpListener::bind("[::]:0").await { + Ok(listener) => listener, + Err(_) => TcpListener::bind("127.0.0.1:0").await.context("bind mock RDP")?, + }; + let port = listener.local_addr().context("mock RDP local_addr")?.port(); + let credssp_ok = Arc::new(AtomicBool::new(false)); + let credssp_ok_task = Arc::clone(&credssp_ok); + let acceptor = tls_acceptor()?; + let public_key = server_public_key()?; + + tokio::spawn(async move { + loop { + let Ok((stream, peer)) = listener.accept().await else { + break; + }; + let acceptor = acceptor.clone(); + let public_key = public_key.clone(); + let credssp_ok = Arc::clone(&credssp_ok_task); + let kdc_url = kdc_url.clone(); + tokio::spawn(async move { + match accept_kerberos_rdp(stream, peer, acceptor, public_key, &kdc_url).await { + Ok(()) => credssp_ok.store(true, Ordering::SeqCst), + Err(error) => eprintln!("mock RDP CredSSP failed: {error:#}"), + } + }); + } + }); + + Ok(Self { port, credssp_ok }) + } + + fn credssp_ok(&self) -> bool { + self.credssp_ok.load(Ordering::SeqCst) + } + + async fn wait_credssp(&self) -> anyhow::Result<()> { + let deadline = Instant::now() + Duration::from_secs(30); + loop { + if self.credssp_ok() { + return Ok(()); + } + if Instant::now() >= deadline { + anyhow::bail!("timed out waiting for Kerberos CredSSP on mock RDP"); + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + } +} + +async fn accept_kerberos_rdp( + stream: TcpStream, + peer: std::net::SocketAddr, + acceptor: tokio_rustls::TlsAcceptor, + public_key: Vec, + kdc_url: &str, +) -> anyhow::Result<()> { + let mut framed = TokioFramed::new(stream); + let (_, request) = framed.read_pdu().await.context("read X.224 CR")?; + let _: X224 = ironrdp_core::decode(&request).context("decode X.224 CR")?; + + let confirm = X224(ConnectionConfirm::Response { + flags: ResponseFlags::empty(), + protocol: SecurityProtocol::HYBRID, + }); + framed + .write_all(&ironrdp_core::encode_vec(&confirm).context("encode X.224 CC")?) + .await + .context("write X.224 CC")?; + + let tcp = framed.into_inner_no_leftover(); + let tls = acceptor.accept(tcp).await.context("TLS accept")?; + let mut framed = TokioFramed::new(tls); + + let identity = sspi::AuthIdentity { + username: sspi::Username::parse(KERBEROS_TARGET_USER).context("parse target username")?, + password: TARGET_PASSWORD.to_owned().into(), + }; + let kerberos_config = sspi::KerberosServerConfig { + kerberos_config: sspi::KerberosConfig { + kdc_url: Some(kdc_url.parse().context("parse mock KDC URL")?), + client_computer_name: peer.to_string(), + }, + server_properties: sspi::kerberos::ServerProperties::new( + &["TERMSRV", SERVICE_HOST], + Some(sspi::CredentialsBuffers::AuthIdentity( + sspi::AuthIdentityBuffers::from_utf8(identity.username.account_name(), REALM, TARGET_PASSWORD), + )), + Duration::from_secs(300), + Some(sspi::Secret::new(TERMSRV_KEY.to_vec())), + ) + .context("Kerberos server properties")?, + }; + + let mut server = sspi::credssp::CredSspServer::new( + public_key, + IdentityProxy(identity), + sspi::credssp::ServerMode::Negotiate(sspi::NegotiateConfig::new( + Box::new(kerberos_config), + Some("kerberos,!ntlm".to_owned()), + peer.to_string(), + )), + ) + .context("init Kerberos-only CredSSP server")?; + + let hint = TsRequestHint; + let mut buf = ironrdp_pdu::WriteBuf::new(); + for _ in 0..6 { + let pdu = framed.read_by_hint(&hint).await.context("read CredSSP TSRequest")?; + let ts_request = sspi::credssp::TsRequest::from_buffer(&pdu).context("decode CredSSP")?; + let result = { + let mut generator = server.process(ts_request); + resolve_sspi_server(&mut generator) + .await + .map_err(|error| anyhow::anyhow!("mock RDP CredSSP: {error:?}"))? + }; + match result { + sspi::credssp::ServerState::ReplyNeeded(outbound) => { + buf.clear(); + let length = usize::from(outbound.buffer_len()); + outbound + .encode_ts_request(buf.unfilled_to(length)) + .context("encode server TSRequest")?; + buf.advance(length); + framed.write_all(&buf[..length]).await.context("write CredSSP")?; + } + sspi::credssp::ServerState::Finished(_) => return Ok(()), + } + } + anyhow::bail!("mock RDP CredSSP exceeded 6 round trips") +} + +async fn resolve_sspi_server( + generator: &mut sspi::generator::Generator< + '_, + sspi::generator::NetworkRequest, + sspi::Result>, + Result, + >, +) -> Result { + let mut state = generator.start(); + loop { + match state { + GeneratorState::Suspended(request) => { + let reply = send_kdc_tcp(&request) + .await + .map_err(|error| sspi::credssp::ServerError { + ts_request: None, + error: sspi::Error::new(sspi::ErrorKind::NoAuthenticatingAuthority, error), + })?; + state = generator.resume(Ok(reply)); + } + GeneratorState::Completed(result) => break result, + } + } +} + +async fn send_kdc_tcp(request: &sspi::generator::NetworkRequest) -> anyhow::Result> { + let host = request.url.host_str().context("KDC host")?; + let port = request.url.port().unwrap_or(88); + let mut stream = TcpStream::connect((host, port)).await.context("connect mock KDC")?; + stream.write_all(&request.data).await.context("write KDC request")?; + let mut len_buf = [0u8; 4]; + stream.read_exact(&mut len_buf).await.context("read KDC length")?; + let len = usize::try_from(u32::from_be_bytes(len_buf)).context("KDC length")?; + let mut body = vec![0u8; len]; + stream.read_exact(&mut body).await.context("read KDC body")?; + let mut reply = Vec::with_capacity(4 + len); + reply.extend_from_slice(&len_buf); + reply.extend_from_slice(&body); + Ok(reply) +} + +struct IdentityProxy(sspi::AuthIdentity); + +impl sspi::credssp::CredentialsProxy for IdentityProxy { + type AuthenticationData = sspi::AuthIdentity; + + fn auth_data_by_user(&mut self, username: &sspi::Username) -> std::io::Result { + if username.account_name() != self.0.username.account_name() { + return Err(std::io::Error::other("invalid username")); + } + let mut data = self.0.clone(); + data.username = username.clone(); + Ok(data) + } + + fn auth_data(&mut self) -> Result, std::io::Error> { + Ok(vec![self.0.clone()]) + } +} + +async fn connect_ntlm_client( + gateway_tcp: u16, + association_jwt: &str, +) -> anyhow::Result> { + let mut stream = TcpStream::connect(("127.0.0.1", gateway_tcp)) + .await + .context("connect gateway TCP")?; + stream + .write_all(&encode_pcb(association_jwt)?) + .await + .context("write PCB")?; + stream.write_all(&encode_hybrid_cr()?).await.context("write X.224 CR")?; + stream.flush().await.context("flush CR")?; + + let mut framed = TokioFramed::new(stream); + let (_, confirm) = framed.read_pdu().await.context("read X.224 CC")?; + let confirm: X224 = ironrdp_core::decode(&confirm).context("decode X.224 CC")?; + anyhow::ensure!( + matches!(confirm.0, ConnectionConfirm::Response { protocol, .. } if protocol.contains(SecurityProtocol::HYBRID)), + "gateway did not confirm CredSSP: {confirm:?}" + ); + + let tcp = framed.into_inner_no_leftover(); + let connector = dangerous_tls_connector(); + let server_name = ServerName::try_from("localhost").map_err(|error| anyhow::anyhow!("{error}"))?; + connector.connect(server_name, tcp).await.context("TLS to gateway") +} + +async fn complete_ntlm_credssp(tls: tokio_rustls::client::TlsStream) -> anyhow::Result<()> { + use sspi::credssp::{ClientMode, ClientState, CredSspClient, CredSspMode, TsRequest}; + use sspi::ntlm::NtlmConfig; + + let public_key = peer_public_key(&tls)?; + let mut framed = TokioFramed::new(tls); + let identity = sspi::AuthIdentity { + username: sspi::Username::parse(PROXY_USER).context("parse proxy username")?, + password: PROXY_PASSWORD.to_owned().into(), + }; + let mut client = CredSspClient::new( + public_key, + identity.into(), + CredSspMode::WithCredentials, + // Gateway's CredSSP server is Negotiate (Kerberos+NTLM). A raw NTLM client is rejected. + ClientMode::Negotiate(sspi::NegotiateConfig::new( + Box::new(NtlmConfig { + client_computer_name: Some("cred-injection-e2e".to_owned()), + }), + Some("ntlm,!kerberos,!pku2u".to_owned()), + "cred-injection-e2e".to_owned(), + )), + format!("TERMSRV/{SERVICE_HOST}"), + ) + .context("init Negotiate-NTLM CredSSP client")?; + + let mut ts_request = TsRequest::default(); + let mut buf = ironrdp_pdu::WriteBuf::new(); + let hint = TsRequestHint; + + for _ in 0..6 { + let client_state = { + let mut generator = client.process(std::mem::take(&mut ts_request)); + resolve_sspi_client(&mut generator)? + }; + let (outbound, finished) = match client_state { + ClientState::ReplyNeeded(request) => (request, false), + ClientState::FinalMessage(request) => (request, true), + }; + buf.clear(); + let length = usize::from(outbound.buffer_len()); + outbound + .encode_ts_request(buf.unfilled_to(length)) + .context("encode client TSRequest")?; + buf.advance(length); + framed.write_all(&buf[..length]).await.context("write client CredSSP")?; + if finished { + return Ok(()); + } + let pdu = framed.read_by_hint(&hint).await.context("read server CredSSP")?; + ts_request = TsRequest::from_buffer(&pdu).context("decode server TSRequest")?; + } + + anyhow::bail!("CredSSP exceeded 6 round trips") +} + +fn resolve_sspi_client( + generator: &mut sspi::generator::Generator< + '_, + sspi::generator::NetworkRequest, + sspi::Result>, + sspi::Result, + >, +) -> anyhow::Result { + let state = generator.start(); + match state { + GeneratorState::Suspended(request) => { + anyhow::bail!("NTLM CredSSP client issued a network request: {}", request.url); + } + GeneratorState::Completed(result) => result.map_err(|error| anyhow::anyhow!("client CredSSP: {error}")), + } +} + +#[derive(Debug)] +struct TsRequestHint; + +impl ironrdp_pdu::PduHint for TsRequestHint { + fn find_size(&self, bytes: &[u8]) -> ironrdp_core::DecodeResult> { + match sspi::credssp::TsRequest::read_length(bytes) { + Ok(length) => Ok(Some((true, length))), + Err(error) if error.kind() == std::io::ErrorKind::UnexpectedEof => Ok(None), + Err(error) => Err(ironrdp_core::other_err!("TsRequestHint", source: error)), + } + } +} + +fn association_token_for_host(jti: &str, jet_aid: &str, dest_port: u16, jet_reuse: u32) -> anyhow::Result { + unsigned_jws( + serde_json::json!({"alg":"RS256","typ":"JWT","cty":"ASSOCIATION"}), + serde_json::json!({ + "dst_hst": format!("{SERVICE_HOST}:{dest_port}"), + "exp": 9_999_999_999i64, + "jet_aid": jet_aid, + "jet_ap": "rdp", + "jet_cm": "fwd", + "jet_rec": "none", + "jet_reuse": jet_reuse, + "jti": jti, + "nbf": 0, + }), + ) +} + +fn encode_hybrid_cr() -> anyhow::Result> { + let pdu = X224(ConnectionRequest { + nego_data: Some(NegoRequestData::cookie(super::cred_injection::CLIENT_COOKIE.to_owned())), + flags: RequestFlags::empty(), + protocol: SecurityProtocol::HYBRID | SecurityProtocol::SSL, + }); + ironrdp_core::encode_vec(&pdu).context("encode hybrid CR") +} + +fn peer_public_key(tls: &tokio_rustls::client::TlsStream) -> anyhow::Result> { + let cert = tls + .get_ref() + .1 + .peer_certificates() + .and_then(|certs| certs.first()) + .context("gateway TLS certificate missing")?; + extract_public_key(cert) +} + +fn server_public_key() -> anyhow::Result> { + let cert = CertificateDer::from_pem_slice(CERT_PEM.as_bytes()).context("parse mock RDP cert")?; + extract_public_key(&cert) +} + +fn extract_public_key(cert: &CertificateDer<'_>) -> anyhow::Result> { + let cert = x509_cert::Certificate::from_der(cert.as_ref()).context("parse X509")?; + let public_key = cert + .tbs_certificate() + .subject_public_key_info() + .subject_public_key + .as_bytes() + .context("unaligned subject public key")? + .to_owned(); + Ok(public_key) +} + +fn tls_acceptor() -> anyhow::Result { + let cert = CertificateDer::from_pem_slice(CERT_PEM.as_bytes()).context("parse cert PEM")?; + let key = PrivateKeyDer::from_pem_slice(KEY_PEM.as_bytes()).context("parse key PEM")?; + let config = ServerConfig::builder() + .with_no_client_auth() + .with_single_cert(vec![cert], key) + .context("TLS server config")?; + Ok(tokio_rustls::TlsAcceptor::from(Arc::new(config))) +} + +fn dangerous_tls_connector() -> tokio_rustls::TlsConnector { + let mut config = ClientConfig::builder() + .dangerous() + .with_custom_certificate_verifier(Arc::new(NoCertificateVerification)) + .with_no_client_auth(); + config.resumption = tokio_rustls::rustls::client::Resumption::disabled(); + tokio_rustls::TlsConnector::from(Arc::new(config)) +} + +fn install_crypto_provider() { + let _ = tokio_rustls::rustls::crypto::ring::default_provider().install_default(); +} + +#[derive(Debug)] +struct NoCertificateVerification; + +impl tokio_rustls::rustls::client::danger::ServerCertVerifier for NoCertificateVerification { + fn verify_server_cert( + &self, + _: &CertificateDer<'_>, + _: &[CertificateDer<'_>], + _: &ServerName<'_>, + _: &[u8], + _: tokio_rustls::rustls::pki_types::UnixTime, + ) -> Result { + Ok(tokio_rustls::rustls::client::danger::ServerCertVerified::assertion()) + } + + fn verify_tls12_signature( + &self, + _: &[u8], + _: &CertificateDer<'_>, + _: &tokio_rustls::rustls::DigitallySignedStruct, + ) -> Result { + Ok(tokio_rustls::rustls::client::danger::HandshakeSignatureValid::assertion()) + } + + fn verify_tls13_signature( + &self, + _: &[u8], + _: &CertificateDer<'_>, + _: &tokio_rustls::rustls::DigitallySignedStruct, + ) -> Result { + Ok(tokio_rustls::rustls::client::danger::HandshakeSignatureValid::assertion()) + } + + fn supported_verify_schemes(&self) -> Vec { + vec![ + tokio_rustls::rustls::SignatureScheme::RSA_PKCS1_SHA256, + tokio_rustls::rustls::SignatureScheme::ECDSA_NISTP256_SHA256, + tokio_rustls::rustls::SignatureScheme::RSA_PSS_SHA256, + tokio_rustls::rustls::SignatureScheme::ED25519, + ] + } +} + +#[tokio::test] +async fn kerberos_injection_completes_credssp_against_mock_kdc() -> anyhow::Result<()> { + install_crypto_provider(); + let kdc = MockKdc::start().await?; + let rdp = MockRdp::start(kdc.url()).await?; + let mut gateway = GatewayProc::start(true).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token_for_host(&jti, &jet_aid, rdp.port, 60)?; + provision_credentials( + gateway.config.http_port(), + &token, + KERBEROS_TARGET_USER, + 300, + Some(&kdc.url()), + ) + .await?; + + let tls = connect_ntlm_client(gateway.config.tcp_port(), &token).await?; + complete_ntlm_credssp(tls) + .await + .context("Gateway-facing NTLM CredSSP")?; + rdp.wait_credssp() + .await + .with_context(|| format!("gateway logs:\n{}", gateway.logs.snapshot()))?; + anyhow::ensure!( + kdc.exchanges() >= 2, + "expected AS-REQ and TGS-REQ against the mock KDC; exchanges={}; gateway logs:\n{}", + kdc.exchanges(), + gateway.logs.snapshot() + ); + + let _ = gateway.process.start_kill(); + Ok(()) +} + +const CERT_PEM: &str = r#"-----BEGIN CERTIFICATE----- +MIIDCzCCAfOgAwIBAgIUPRJa8i280unV3/kW6TE2fSUw8PwwDQYJKoZIhvcNAQEL +BQAwFDESMBAGA1UEAwwJbG9jYWxob3N0MCAXDTI1MTEyNTA5NDAzMFoYDzIxMjUx +MTAxMDk0MDMwWjAUMRIwEAYDVQQDDAlsb2NhbGhvc3QwggEiMA0GCSqGSIb3DQEB +AQUAA4IBDwAwggEKAoIBAQDHpBlyRgUx/V9cQGw/eqDFc6odxB2hvnbudi67LvEj +cNIWOU79R1e/NswME4oecqT9W05n4UyxkABfm2qjODO0nDf47W0DsgbEA87qE715 +RWg8AtC529CZAazqTV3gqYyRMsCuVKzPVxgWa8rhPc7E6In1uDRak0lWKQPQSBbc +34nxMOVIusZNlkAEar8/aYPr/YWvdEqkobEvXp+g9WsuMaU913ecacWDjyWDkf80 +pPPtf+uet7WMysKMhzGQtpbgilT8XCo8uTsgUbK+TMWvkF9bcxAQDnJsrZRL7Jfh +ofsFfQbTIvbvpn+4J4kmHN36BTohlNL8TX1jrU3cPA7dAgMBAAGjUzBRMB0GA1Ud +DgQWBBTT+m6dyc/c3mXF3JAsZr9OqUwgWTAfBgNVHSMEGDAWgBTT+m6dyc/c3mXF +3JAsZr9OqUwgWTAPBgNVHRMBAf8EBTADAQH/MA0GCSqGSIb3DQEBCwUAA4IBAQBB +i/yonZY3ztaeGElzD8xkI+rJ+daJ5WzdfKnzudJllg/Ht8m7wO5SdQnMt2T44gbH +05uekc1zXnXb7fJKqs3R6DacctG0nQ3acuI+IMtTaBbbAcf3PJJlo0Pap0ypVC0R +IUiUhJGFNi4cCBOvJqsly0d3T5xqOXU1Q5j3mIwRBY68+m9btwwuZWvASRADtCyZ +RpisBzS4a6jSeHXa4iG/VhskbiZkcnfHNTw7yNJJdv125y2zQkWWF9wlLbYwWr40 +x9Ba6YbssOz6epATKhvt80yclO34AzUyimssvViIUpgFEyaPhZZTw46Q/6X3ixK4 +/v4eYM0cCHN0h+rynSor +-----END CERTIFICATE-----"#; + +const KEY_PEM: &str = r#"-----BEGIN PRIVATE KEY----- +MIIEvwIBADANBgkqhkiG9w0BAQEFAASCBKkwggSlAgEAAoIBAQDHpBlyRgUx/V9c +QGw/eqDFc6odxB2hvnbudi67LvEjcNIWOU79R1e/NswME4oecqT9W05n4UyxkABf +m2qjODO0nDf47W0DsgbEA87qE715RWg8AtC529CZAazqTV3gqYyRMsCuVKzPVxgW +a8rhPc7E6In1uDRak0lWKQPQSBbc34nxMOVIusZNlkAEar8/aYPr/YWvdEqkobEv +Xp+g9WsuMaU913ecacWDjyWDkf80pPPtf+uet7WMysKMhzGQtpbgilT8XCo8uTsg +UbK+TMWvkF9bcxAQDnJsrZRL7JfhofsFfQbTIvbvpn+4J4kmHN36BTohlNL8TX1j +rU3cPA7dAgMBAAECggEAKh7KK5zwTaq6atlAvWfe8anEk4EkC1MG/qq6k02FHMgZ +2wx+SNu7fKFQDaA1vNTNUJLqCOq05qWOHp3IsuURq6JmAMP/Aw+Vc9el2ScPC74E +Dt09MmlZKl77H3fxPYwoFx5RHrbIuvoSH/DgHgOPU2YIbWpOyWlXyLDgmBoNkM3N +fXYLXJONpStPHeQLhh7LcHO3CZgn6kycJyByEO2NtcchS5zITiJuwL+qR5/QIlvD +Yo7jdCjelJat38MZ9dE1us8xlIjQtsYF/acZZtcpYho+7ZpDCNcb+xF8KStKei+B +MMpWISsa+Zh9g7lPYTnG/i1dSMMT100XCEw8o4rBoQKBgQDnptz8acp7DB2wJH4L +c0xuw8IlrSl3BGUEj8H+RyFlpH3+//i6/fE9MrtF8b4FSYUp5AG4NVFGcRbwJVGW +jeL13YwIKMdXjmx8fDIylCgBB1tzBS9T/0ws3HS8avxhKvjgoXIZm6D3XDcBslrH +c9/LojT8YGI1wx7jWI2qKj8yeQKBgQDcn+kQ1QjzgIz6bAVWY3t1jr5uHHyaS+5G +ihY/mx4Mn3DURgPXZHz/HrN9rZkax0zuq9wuIlqgZ2KI37iCF49M4aZxC788LyDo +Hp0Cak3wt3g0Tj6J7SJiQe8h/6VBS4R5dRD2vhEc3xPAOf7WIFdlLYBOOvE/LmOt +N6ChkfgGhQKBgQDSiDqLRPJ7BjXtIh1T9sPeXxeR+mCXBG1yydx7ZtYZdHf2S1kZ +STX4cqT1GpGiaIEX41sUuZBWPu2j76bI98bvwRxFRhp1nsFGGfHdOf1pgfBBBtNO +udXXZ7zIiUs6XD24mcIDOAgBB9QOPLR4VP1uKsuRG1/mkKD/6jlGEANDsQKBgQDC +AoEygxQnBVFz2c/rwvnLS+Zb8AMGsGTtdPrRnjeThBX1JUi1fbGJq1bN2v27Fa2q +aEjr7NvjGGcG1C1tgQhL5Fa4LEtTwmHenSUW/aJiXwR+gpvuMDC/VRnTvPp2a9En ++XEcedGUoPq+XIGjjLctyxB8Osrw83tF1JgV3MXN/QKBgQC83B54rYDd4QmVH5nL +WLw834fgr+Z1hA6UqJIaahlD/bDwzbbJEv0pHCBxe01ywQFivqWBdVbuoy9YSeLS +KKEklzh+L0SorrYoBA5F63qx0zy05bba0ASplgDUEUNZn7oIFi7x5pVsNNaNxZpR +bQGM8UrNQvWQ+tutRmp7PM6VuQ== +-----END PRIVATE KEY-----"#; diff --git a/testsuite/tests/cli/dgw/mod.rs b/testsuite/tests/cli/dgw/mod.rs index f6e88737a..acb5a50a6 100644 --- a/testsuite/tests/cli/dgw/mod.rs +++ b/testsuite/tests/cli/dgw/mod.rs @@ -1,6 +1,7 @@ mod benign_disconnect; mod cli_args; mod cred_injection; +mod cred_injection_kdc; mod heartbeat; mod preflight; mod tls_anchoring; From 01c5545c1777ad309bac4b4a09da0b87139adf3a Mon Sep 17 00:00:00 2001 From: Junyi Ou Date: Thu, 20 Aug 2026 18:02:50 -0400 Subject: [PATCH 4/9] test(dgw): cover Kerberos client leg and fail-closed paths Prove both CredSSP hops against mock KDC/RDP, NTLM CredSSP both legs, and Kerberos fail-closed when the password, KDC, or krb_kdc is wrong. Token-cache jet_reuse still needs signed JWTs. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- testsuite/tests/cli/dgw/cred_injection.rs | 86 +++- testsuite/tests/cli/dgw/cred_injection_kdc.rs | 446 +++++++++++++++++- 2 files changed, 504 insertions(+), 28 deletions(-) diff --git a/testsuite/tests/cli/dgw/cred_injection.rs b/testsuite/tests/cli/dgw/cred_injection.rs index 11fafd805..6c5d236d4 100644 --- a/testsuite/tests/cli/dgw/cred_injection.rs +++ b/testsuite/tests/cli/dgw/cred_injection.rs @@ -25,8 +25,9 @@ pub(crate) const PROXY_PASSWORD: &str = "proxy-secret"; pub(crate) const TARGET_PASSWORD: &str = "target-secret"; pub(crate) const KERBEROS_TARGET_USER: &str = "administrator@example.invalid"; pub(crate) const INJECT_LOG: &str = "RDP-TLS forwarding with credential injection"; -const FORWARD_LOG: &str = "Upstream forwarding"; -const MISSING_LOG: &str = "missing or expired; re-provision to retry"; +pub(crate) const FORWARD_LOG: &str = "Upstream forwarding"; +pub(crate) const MISSING_LOG: &str = "missing or expired; re-provision to retry"; +pub(crate) const PROXY_KERBEROS_USER: &str = "injected-proxy-user@example.invalid"; const PUBLISHED_KDC_LOG: &str = "Published synthetic KDC"; const REGISTERED_KDC_LOG: &str = "Registered synthetic KDC for credential-injection session"; @@ -346,6 +347,27 @@ pub(crate) async fn provision_credentials( target_username: &str, time_to_live: u32, krb_kdc: Option<&str>, +) -> anyhow::Result<()> { + provision_mapping( + http_port, + token, + PROXY_USER, + target_username, + TARGET_PASSWORD, + time_to_live, + krb_kdc, + ) + .await +} + +pub(crate) async fn provision_mapping( + http_port: u16, + token: &str, + proxy_username: &str, + target_username: &str, + target_password: &str, + time_to_live: u32, + krb_kdc: Option<&str>, ) -> anyhow::Result<()> { let mut operations = vec![serde_json::json!({ "id": next_id(), @@ -353,13 +375,13 @@ pub(crate) async fn provision_credentials( "token": token, "proxy_credential": { "kind": "username-password", - "username": PROXY_USER, + "username": proxy_username, "password": PROXY_PASSWORD }, "target_credential": { "kind": "username-password", "username": target_username, - "password": TARGET_PASSWORD + "password": target_password }, "time_to_live": time_to_live })]; @@ -611,3 +633,59 @@ async fn kerberos_reconnect_reuses_generation_until_reprovision() -> anyhow::Res let _ = gateway.process.start_kill(); Ok(()) } + +#[tokio::test] +async fn domainless_target_stays_ntlm_even_with_krb_kdc() -> anyhow::Result<()> { + let target = FakeRdpTarget::start().await?; + let mut gateway = GatewayProc::start(true).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token(&jti, &jet_aid, target.port, 60)?; + provision_credentials( + gateway.config.http_port(), + &token, + TARGET_USER, + 300, + Some("tcp://127.0.0.1:88"), + ) + .await?; + + let _client = connect_rdp_client(gateway.config.tcp_port(), &token).await?; + let logs = gateway.logs.wait_contains(INJECT_LOG).await?; + assert!( + logs.contains("kerberos=false"), + "username without a realm must stay NTLM even if krb_kdc is provisioned; logs:\n{logs}" + ); + + let _ = gateway.process.start_kill(); + Ok(()) +} + +#[tokio::test] +async fn kerberos_opt_out_uses_ntlm_for_domain_user() -> anyhow::Result<()> { + let target = FakeRdpTarget::start().await?; + let mut gateway = GatewayProc::start(false).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token(&jti, &jet_aid, target.port, 60)?; + provision_credentials( + gateway.config.http_port(), + &token, + KERBEROS_TARGET_USER, + 300, + Some("tcp://127.0.0.1:88"), + ) + .await?; + + let _client = connect_rdp_client(gateway.config.tcp_port(), &token).await?; + let logs = gateway.logs.wait_contains(INJECT_LOG).await?; + assert!( + logs.contains("kerberos=false"), + "Kerberos injection opt-out must NTLM even with a domain username; logs:\n{logs}" + ); + + let _ = gateway.process.start_kill(); + Ok(()) +} diff --git a/testsuite/tests/cli/dgw/cred_injection_kdc.rs b/testsuite/tests/cli/dgw/cred_injection_kdc.rs index 38846da56..364097ca0 100644 --- a/testsuite/tests/cli/dgw/cred_injection_kdc.rs +++ b/testsuite/tests/cli/dgw/cred_injection_kdc.rs @@ -17,7 +17,7 @@ use ironrdp_pdu::nego::{ use ironrdp_pdu::x224::X224; use ironrdp_tokio::{FramedWrite as _, TokioFramed}; use picky_krb::messages::KdcProxyMessage; -use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _}; +use tokio::io::{AsyncBufReadExt as _, AsyncReadExt as _, AsyncWriteExt as _, BufReader}; use tokio::net::{TcpListener, TcpStream}; use tokio_rustls::rustls::pki_types::pem::PemObject as _; use tokio_rustls::rustls::pki_types::{CertificateDer, PrivateKeyDer, ServerName}; @@ -25,8 +25,9 @@ use tokio_rustls::rustls::{ClientConfig, ServerConfig}; use x509_cert::der::Decode as _; use super::cred_injection::{ - GatewayProc, KERBEROS_TARGET_USER, PROXY_PASSWORD, PROXY_USER, TARGET_PASSWORD, encode_pcb, next_id, - provision_credentials, unsigned_jws, + FORWARD_LOG, GatewayProc, INJECT_LOG, KERBEROS_TARGET_USER, MISSING_LOG, PROXY_KERBEROS_USER, PROXY_PASSWORD, + PROXY_USER, TARGET_PASSWORD, TARGET_USER, encode_pcb, next_id, provision_credentials, provision_mapping, + unsigned_jws, }; const REALM: &str = "EXAMPLE.INVALID"; @@ -114,13 +115,27 @@ async fn serve_kdc_exchange(mut stream: TcpStream, config: &kdc::config::Kerbero Ok(()) } +#[derive(Clone)] +enum MockRdpMode { + Kerberos { kdc_url: String }, + Ntlm, +} + struct MockRdp { port: u16, credssp_ok: Arc, } impl MockRdp { - async fn start(kdc_url: String) -> anyhow::Result { + async fn start_kerberos(kdc_url: String) -> anyhow::Result { + Self::start(MockRdpMode::Kerberos { kdc_url }).await + } + + async fn start_ntlm() -> anyhow::Result { + Self::start(MockRdpMode::Ntlm).await + } + + async fn start(mode: MockRdpMode) -> anyhow::Result { install_crypto_provider(); // Dual-stack so Windows `localhost` (IPv6 first) still hits the fake server. let listener = match TcpListener::bind("[::]:0").await { @@ -141,9 +156,15 @@ impl MockRdp { let acceptor = acceptor.clone(); let public_key = public_key.clone(); let credssp_ok = Arc::clone(&credssp_ok_task); - let kdc_url = kdc_url.clone(); + let mode = mode.clone(); tokio::spawn(async move { - match accept_kerberos_rdp(stream, peer, acceptor, public_key, &kdc_url).await { + let result = match &mode { + MockRdpMode::Kerberos { kdc_url } => { + accept_kerberos_rdp(stream, peer, acceptor, public_key, kdc_url).await + } + MockRdpMode::Ntlm => accept_ntlm_rdp(stream, peer, acceptor, public_key).await, + }; + match result { Ok(()) => credssp_ok.store(true, Ordering::SeqCst), Err(error) => eprintln!("mock RDP CredSSP failed: {error:#}"), } @@ -254,6 +275,69 @@ async fn accept_kerberos_rdp( anyhow::bail!("mock RDP CredSSP exceeded 6 round trips") } +async fn accept_ntlm_rdp( + stream: TcpStream, + peer: std::net::SocketAddr, + acceptor: tokio_rustls::TlsAcceptor, + public_key: Vec, +) -> anyhow::Result<()> { + let mut framed = TokioFramed::new(stream); + let (_, request) = framed.read_pdu().await.context("read X.224 CR")?; + let _: X224 = ironrdp_core::decode(&request).context("decode X.224 CR")?; + + let confirm = X224(ConnectionConfirm::Response { + flags: ResponseFlags::empty(), + protocol: SecurityProtocol::HYBRID, + }); + framed + .write_all(&ironrdp_core::encode_vec(&confirm).context("encode X.224 CC")?) + .await + .context("write X.224 CC")?; + + let tcp = framed.into_inner_no_leftover(); + let tls = acceptor.accept(tcp).await.context("TLS accept")?; + let mut framed = TokioFramed::new(tls); + + let identity = sspi::AuthIdentity { + username: sspi::Username::parse(TARGET_USER).context("parse NTLM target username")?, + password: TARGET_PASSWORD.to_owned().into(), + }; + let mut server = sspi::credssp::CredSspServer::new( + public_key, + IdentityProxy(identity), + sspi::credssp::ServerMode::Ntlm(sspi::ntlm::NtlmConfig { + client_computer_name: Some(peer.to_string()), + }), + ) + .context("init NTLM CredSSP server")?; + + let hint = TsRequestHint; + let mut buf = ironrdp_pdu::WriteBuf::new(); + for _ in 0..6 { + let pdu = framed.read_by_hint(&hint).await.context("read CredSSP TSRequest")?; + let ts_request = sspi::credssp::TsRequest::from_buffer(&pdu).context("decode CredSSP")?; + let result = { + let mut generator = server.process(ts_request); + resolve_sspi_server(&mut generator) + .await + .map_err(|error| anyhow::anyhow!("mock RDP NTLM CredSSP: {error:?}"))? + }; + match result { + sspi::credssp::ServerState::ReplyNeeded(outbound) => { + buf.clear(); + let length = usize::from(outbound.buffer_len()); + outbound + .encode_ts_request(buf.unfilled_to(length)) + .context("encode server TSRequest")?; + buf.advance(length); + framed.write_all(&buf[..length]).await.context("write CredSSP")?; + } + sspi::credssp::ServerState::Finished(_) => return Ok(()), + } + } + anyhow::bail!("mock RDP NTLM CredSSP exceeded 6 round trips") +} + async fn resolve_sspi_server( generator: &mut sspi::generator::Generator< '_, @@ -343,39 +427,68 @@ async fn connect_ntlm_client( } async fn complete_ntlm_credssp(tls: tokio_rustls::client::TlsStream) -> anyhow::Result<()> { + complete_client_credssp(tls, PROXY_USER, None, false).await +} + +async fn complete_raw_ntlm_credssp(tls: tokio_rustls::client::TlsStream) -> anyhow::Result<()> { + complete_client_credssp(tls, PROXY_USER, None, true).await +} + +async fn complete_client_credssp( + tls: tokio_rustls::client::TlsStream, + username: &str, + kdc_proxy_url: Option<&str>, + raw_ntlm: bool, +) -> anyhow::Result<()> { use sspi::credssp::{ClientMode, ClientState, CredSspClient, CredSspMode, TsRequest}; use sspi::ntlm::NtlmConfig; let public_key = peer_public_key(&tls)?; let mut framed = TokioFramed::new(tls); let identity = sspi::AuthIdentity { - username: sspi::Username::parse(PROXY_USER).context("parse proxy username")?, + username: sspi::Username::parse(username).context("parse client username")?, password: PROXY_PASSWORD.to_owned().into(), }; - let mut client = CredSspClient::new( - public_key, - identity.into(), - CredSspMode::WithCredentials, - // Gateway's CredSSP server is Negotiate (Kerberos+NTLM). A raw NTLM client is rejected. + let client_mode = if let Some(kdc_url) = kdc_proxy_url { + ClientMode::Negotiate(sspi::NegotiateConfig::new( + Box::new(sspi::KerberosConfig { + kdc_url: Some(kdc_url.parse().context("parse KDC proxy URL")?), + client_computer_name: "cred-injection-e2e".to_owned(), + }), + Some("kerberos,!ntlm".to_owned()), + "cred-injection-e2e".to_owned(), + )) + } else if raw_ntlm { + // Gateway NTLM injection uses ServerMode::Ntlm, which rejects SPNEGO. + ClientMode::Ntlm(NtlmConfig { + client_computer_name: Some("cred-injection-e2e".to_owned()), + }) + } else { ClientMode::Negotiate(sspi::NegotiateConfig::new( Box::new(NtlmConfig { client_computer_name: Some("cred-injection-e2e".to_owned()), }), Some("ntlm,!kerberos,!pku2u".to_owned()), "cred-injection-e2e".to_owned(), - )), + )) + }; + let mut client = CredSspClient::new( + public_key, + identity.into(), + CredSspMode::WithCredentials, + client_mode, format!("TERMSRV/{SERVICE_HOST}"), ) - .context("init Negotiate-NTLM CredSSP client")?; + .context("init CredSSP client")?; let mut ts_request = TsRequest::default(); let mut buf = ironrdp_pdu::WriteBuf::new(); let hint = TsRequestHint; - for _ in 0..6 { + for _ in 0..8 { let client_state = { let mut generator = client.process(std::mem::take(&mut ts_request)); - resolve_sspi_client(&mut generator)? + resolve_sspi_client(&mut generator).await? }; let (outbound, finished) = match client_state { ClientState::ReplyNeeded(request) => (request, false), @@ -395,10 +508,10 @@ async fn complete_ntlm_credssp(tls: tokio_rustls::client::TlsStream) ts_request = TsRequest::from_buffer(&pdu).context("decode server TSRequest")?; } - anyhow::bail!("CredSSP exceeded 6 round trips") + anyhow::bail!("CredSSP exceeded 8 round trips") } -fn resolve_sspi_client( +async fn resolve_sspi_client( generator: &mut sspi::generator::Generator< '_, sspi::generator::NetworkRequest, @@ -406,15 +519,104 @@ fn resolve_sspi_client( sspi::Result, >, ) -> anyhow::Result { - let state = generator.start(); - match state { - GeneratorState::Suspended(request) => { - anyhow::bail!("NTLM CredSSP client issued a network request: {}", request.url); + let mut state = generator.start(); + loop { + match state { + GeneratorState::Suspended(request) => { + let reply = match request.url.scheme() { + "tcp" | "udp" => send_kdc_tcp(&request).await?, + "http" | "https" => send_kdc_http(&request).await?, + other => anyhow::bail!("unsupported KDC scheme {other}: {}", request.url), + }; + state = generator.resume(Ok(reply)); + } + GeneratorState::Completed(result) => { + break result.map_err(|error| anyhow::anyhow!("client CredSSP: {error}")); + } } - GeneratorState::Completed(result) => result.map_err(|error| anyhow::anyhow!("client CredSSP: {error}")), } } +async fn send_kdc_http(request: &sspi::generator::NetworkRequest) -> anyhow::Result> { + let host = request.url.host_str().context("KDC proxy host")?; + let port = request.url.port_or_known_default().unwrap_or(80); + let path = if request.url.query().is_some() { + format!("{}?{}", request.url.path(), request.url.query().unwrap_or_default()) + } else { + request.url.path().to_owned() + }; + let mut stream = TcpStream::connect((host, port)).await.context("connect KDC proxy")?; + let header = format!( + "POST {path} HTTP/1.1\r\n\ + Host: {host}:{port}\r\n\ + Content-Type: application/octet-stream\r\n\ + Content-Length: {}\r\n\ + Connection: close\r\n\ + \r\n", + request.data.len() + ); + stream + .write_all(header.as_bytes()) + .await + .context("write KDC proxy headers")?; + stream.write_all(&request.data).await.context("write KDC proxy body")?; + stream.flush().await.context("flush KDC proxy")?; + + let mut reader = BufReader::new(stream); + let mut status_line = String::new(); + reader + .read_line(&mut status_line) + .await + .context("read KDC proxy status")?; + anyhow::ensure!(status_line.contains("200"), "KDC proxy HTTP status was {status_line:?}"); + + let mut content_length = None; + loop { + let mut line = String::new(); + reader.read_line(&mut line).await.context("read KDC proxy header")?; + if line == "\r\n" || line.is_empty() { + break; + } + if let Some(value) = line + .split_once(':') + .filter(|(name, _)| name.eq_ignore_ascii_case("content-length")) + .map(|(_, value)| value.trim().to_owned()) + { + content_length = Some(value.parse::().context("parse KDC proxy Content-Length")?); + } + } + + if let Some(len) = content_length { + let mut buf = vec![0u8; len]; + tokio::io::AsyncReadExt::read_exact(&mut reader, &mut buf) + .await + .context("read KDC proxy body")?; + Ok(buf) + } else { + let mut buf = Vec::new(); + tokio::io::AsyncReadExt::read_to_end(&mut reader, &mut buf) + .await + .context("read KDC proxy eof body")?; + Ok(buf) + } +} + +fn kdc_inject_token(association_jti: &str) -> anyhow::Result { + unsigned_jws( + serde_json::json!({"alg":"RS256","typ":"JWT","cty":"KDC"}), + serde_json::json!({ + "exp": 9_999_999_999i64, + "jet_cred_id": association_jti, + "jti": next_id(), + }), + ) +} + +fn kdc_proxy_url(http_port: u16, association_jti: &str) -> anyhow::Result { + let token = kdc_inject_token(association_jti)?; + Ok(format!("http://127.0.0.1:{http_port}/jet/KdcProxy/{token}")) +} + #[derive(Debug)] struct TsRequestHint; @@ -551,7 +753,7 @@ impl tokio_rustls::rustls::client::danger::ServerCertVerifier for NoCertificateV async fn kerberos_injection_completes_credssp_against_mock_kdc() -> anyhow::Result<()> { install_crypto_provider(); let kdc = MockKdc::start().await?; - let rdp = MockRdp::start(kdc.url()).await?; + let rdp = MockRdp::start_kerberos(kdc.url()).await?; let mut gateway = GatewayProc::start(true).await?; let jti = next_id(); @@ -584,6 +786,202 @@ async fn kerberos_injection_completes_credssp_against_mock_kdc() -> anyhow::Resu Ok(()) } +#[tokio::test] +async fn kerberos_client_and_target_legs_complete_credssp() -> anyhow::Result<()> { + install_crypto_provider(); + let kdc = MockKdc::start().await?; + let rdp = MockRdp::start_kerberos(kdc.url()).await?; + let mut gateway = GatewayProc::start(true).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token_for_host(&jti, &jet_aid, rdp.port, 60)?; + provision_mapping( + gateway.config.http_port(), + &token, + PROXY_KERBEROS_USER, + KERBEROS_TARGET_USER, + TARGET_PASSWORD, + 300, + Some(&kdc.url()), + ) + .await?; + + let kdc_proxy = kdc_proxy_url(gateway.config.http_port(), &jti)?; + let tls = connect_ntlm_client(gateway.config.tcp_port(), &token).await?; + complete_client_credssp(tls, PROXY_KERBEROS_USER, Some(&kdc_proxy), false) + .await + .with_context(|| { + format!( + "client-leg Kerberos CredSSP; gateway logs:\n{}", + gateway.logs.snapshot() + ) + })?; + rdp.wait_credssp().await.with_context(|| { + format!( + "target-leg Kerberos CredSSP; gateway logs:\n{}", + gateway.logs.snapshot() + ) + })?; + anyhow::ensure!( + kdc.exchanges() >= 2, + "target-leg must talk to the mock KDC; exchanges={}; logs:\n{}", + kdc.exchanges(), + gateway.logs.snapshot() + ); + + let _ = gateway.process.start_kill(); + Ok(()) +} + +#[tokio::test] +async fn kerberos_wrong_target_password_fails_closed() -> anyhow::Result<()> { + install_crypto_provider(); + let kdc = MockKdc::start().await?; + let rdp = MockRdp::start_kerberos(kdc.url()).await?; + let mut gateway = GatewayProc::start(true).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token_for_host(&jti, &jet_aid, rdp.port, 60)?; + provision_mapping( + gateway.config.http_port(), + &token, + PROXY_USER, + KERBEROS_TARGET_USER, + "wrong-target-password", + 300, + Some(&kdc.url()), + ) + .await?; + + let tls = connect_ntlm_client(gateway.config.tcp_port(), &token).await?; + let _ = complete_ntlm_credssp(tls).await; + tokio::time::sleep(Duration::from_secs(2)).await; + anyhow::ensure!( + !rdp.credssp_ok(), + "wrong target password must not complete Kerberos CredSSP; logs:\n{}", + gateway.logs.snapshot() + ); + let logs = gateway.logs.snapshot(); + anyhow::ensure!( + !logs.contains(FORWARD_LOG), + "wrong password must not ordinary-forward; logs:\n{logs}" + ); + + let _ = gateway.process.start_kill(); + Ok(()) +} + +#[tokio::test] +async fn kerberos_kdc_down_fails_closed() -> anyhow::Result<()> { + install_crypto_provider(); + let rdp = MockRdp::start_kerberos("tcp://127.0.0.1:1".to_owned()).await?; + let mut gateway = GatewayProc::start(true).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token_for_host(&jti, &jet_aid, rdp.port, 60)?; + provision_credentials( + gateway.config.http_port(), + &token, + KERBEROS_TARGET_USER, + 300, + Some("tcp://127.0.0.1:1"), + ) + .await?; + + let tls = connect_ntlm_client(gateway.config.tcp_port(), &token).await?; + let _ = complete_ntlm_credssp(tls).await; + tokio::time::sleep(Duration::from_secs(2)).await; + anyhow::ensure!( + !rdp.credssp_ok(), + "unreachable KDC must not complete CredSSP; logs:\n{}", + gateway.logs.snapshot() + ); + let logs = gateway.logs.snapshot(); + anyhow::ensure!( + !logs.contains(FORWARD_LOG), + "KDC down must not ordinary-forward; logs:\n{logs}" + ); + + let _ = gateway.process.start_kill(); + Ok(()) +} + +#[tokio::test] +async fn kerberos_missing_krb_kdc_fails_closed() -> anyhow::Result<()> { + let rdp = FakeClosedTarget::start().await?; + let mut gateway = GatewayProc::start(true).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token_for_host(&jti, &jet_aid, rdp.port, 60)?; + provision_credentials(gateway.config.http_port(), &token, KERBEROS_TARGET_USER, 300, None).await?; + + let mut stream = TcpStream::connect(("127.0.0.1", gateway.config.tcp_port())) + .await + .context("connect gateway TCP")?; + stream.write_all(&encode_pcb(&token)?).await.context("write PCB")?; + stream.write_all(&encode_hybrid_cr()?).await.context("write CR")?; + stream.flush().await.context("flush CR")?; + let logs = gateway.logs.wait_contains(MISSING_LOG).await?; + anyhow::ensure!( + !logs.contains(FORWARD_LOG), + "missing krb_kdc must fail closed; logs:\n{logs}" + ); + + let _ = gateway.process.start_kill(); + Ok(()) +} + +#[tokio::test] +async fn ntlm_injection_completes_credssp_both_legs() -> anyhow::Result<()> { + install_crypto_provider(); + let rdp = MockRdp::start_ntlm().await?; + let mut gateway = GatewayProc::start(false).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token_for_host(&jti, &jet_aid, rdp.port, 60)?; + provision_credentials(gateway.config.http_port(), &token, TARGET_USER, 300, None).await?; + + let tls = connect_ntlm_client(gateway.config.tcp_port(), &token).await?; + complete_raw_ntlm_credssp(tls) + .await + .with_context(|| format!("client-leg NTLM CredSSP; logs:\n{}", gateway.logs.snapshot()))?; + rdp.wait_credssp() + .await + .with_context(|| format!("target-leg NTLM CredSSP; logs:\n{}", gateway.logs.snapshot()))?; + let logs = gateway.logs.wait_contains(INJECT_LOG).await?; + anyhow::ensure!( + logs.contains("kerberos=false"), + "expected NTLM injection; logs:\n{logs}" + ); + + let _ = gateway.process.start_kill(); + Ok(()) +} + +struct FakeClosedTarget { + port: u16, +} + +impl FakeClosedTarget { + async fn start() -> anyhow::Result { + let listener = TcpListener::bind("127.0.0.1:0").await.context("bind closed target")?; + let port = listener.local_addr()?.port(); + tokio::spawn(async move { + loop { + let Ok((_stream, _)) = listener.accept().await else { + break; + }; + } + }); + Ok(Self { port }) + } +} + const CERT_PEM: &str = r#"-----BEGIN CERTIFICATE----- MIIDCzCCAfOgAwIBAgIUPRJa8i280unV3/kW6TE2fSUw8PwwDQYJKoZIhvcNAQEL BQAwFDESMBAGA1UEAwwJbG9jYWxob3N0MCAXDTI1MTEyNTA5NDAzMFoYDzIxMjUx From 7d22c4895f79dc1e2ff30213c0c72c18b7a029f8 Mon Sep 17 00:00:00 2001 From: Junyi Ou Date: Fri, 21 Aug 2026 11:22:43 -0400 Subject: [PATCH 5/9] test(dgw): assert KDC principals and CredSSP identities Decode AS-REQ/TGS-REQ on the mock KDC (cname, realm, TERMSRV/localhost), record CredSSP Finished account names and X.224 cookies, and require /jet/KdcProxy AS-REP plus TGS-REP. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- Cargo.lock | 1 + testsuite/Cargo.toml | 1 + testsuite/tests/cli/dgw/cred_injection_kdc.rs | 253 ++++++++++++++++-- 3 files changed, 230 insertions(+), 25 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index b3f239ccc..758047ea8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7526,6 +7526,7 @@ dependencies = [ "mcp-proxy", "network-scanner", "network-scanner-proto", + "picky-asn1-der", "picky-krb", "proxy-socks", "rstest", diff --git a/testsuite/Cargo.toml b/testsuite/Cargo.toml index 5c023c919..5c6614577 100644 --- a/testsuite/Cargo.toml +++ b/testsuite/Cargo.toml @@ -38,6 +38,7 @@ ironrdp-core = { version = "0.2", features = ["std"] } ironrdp-pdu = { version = "0.9", features = ["std"] } ironrdp-tokio = "0.10" kdc = "0.1" +picky-asn1-der = "0.5" picky-krb = "0.12" proxy-socks = { path = "../crates/proxy-socks" } libsql = { version = "0.9", default-features = false, features = ["core"] } diff --git a/testsuite/tests/cli/dgw/cred_injection_kdc.rs b/testsuite/tests/cli/dgw/cred_injection_kdc.rs index 364097ca0..e1c80d1be 100644 --- a/testsuite/tests/cli/dgw/cred_injection_kdc.rs +++ b/testsuite/tests/cli/dgw/cred_injection_kdc.rs @@ -4,8 +4,8 @@ //! sspi-rs) and completes CredSSP with a fake RDP acceptor. The Gateway-facing client uses //! NTLM so the test does not depend on the in-process synthetic KDC. -use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; use anyhow::Context as _; @@ -16,7 +16,8 @@ use ironrdp_pdu::nego::{ }; use ironrdp_pdu::x224::X224; use ironrdp_tokio::{FramedWrite as _, TokioFramed}; -use picky_krb::messages::KdcProxyMessage; +use picky_krb::data_types::PrincipalName; +use picky_krb::messages::{AsRep, AsReq, KdcProxyMessage, KrbError, TgsRep, TgsReq}; use tokio::io::{AsyncBufReadExt as _, AsyncReadExt as _, AsyncWriteExt as _, BufReader}; use tokio::net::{TcpListener, TcpStream}; use tokio_rustls::rustls::pki_types::pem::PemObject as _; @@ -36,9 +37,17 @@ const SERVICE_HOST: &str = "localhost"; const KRBTGT_KEY: [u8; 32] = [0x11; 32]; const TERMSRV_KEY: [u8; 32] = [0x22; 32]; +#[derive(Clone, Debug, PartialEq, Eq)] +enum ObservedKdcReq { + As { cname: String, realm: String }, + Tgs { sname: Vec, realm: String }, + Other, +} + struct MockKdc { port: u16, exchanges: Arc, + requests: Arc>>, } impl MockKdc { @@ -46,7 +55,9 @@ impl MockKdc { let listener = TcpListener::bind("127.0.0.1:0").await.context("bind mock KDC")?; let port = listener.local_addr().context("mock KDC local_addr")?.port(); let exchanges = Arc::new(AtomicUsize::new(0)); + let requests = Arc::new(Mutex::new(Vec::new())); let exchanges_task = Arc::clone(&exchanges); + let requests_task = Arc::clone(&requests); let config = kdc_config(); tokio::spawn(async move { @@ -56,8 +67,9 @@ impl MockKdc { }; let config = config.clone(); let exchanges = Arc::clone(&exchanges_task); + let requests = Arc::clone(&requests_task); tokio::spawn(async move { - match serve_kdc_exchange(stream, &config).await { + match serve_kdc_exchange(stream, &config, &requests).await { Ok(()) => { exchanges.fetch_add(1, Ordering::SeqCst); } @@ -67,7 +79,11 @@ impl MockKdc { } }); - Ok(Self { port, exchanges }) + Ok(Self { + port, + exchanges, + requests, + }) } fn url(&self) -> String { @@ -77,6 +93,10 @@ impl MockKdc { fn exchanges(&self) -> usize { self.exchanges.load(Ordering::SeqCst) } + + fn requests(&self) -> Vec { + self.requests.lock().expect("kdc request mutex").clone() + } } fn kdc_config() -> kdc::config::KerberosServer { @@ -95,13 +115,19 @@ fn kdc_config() -> kdc::config::KerberosServer { } } -async fn serve_kdc_exchange(mut stream: TcpStream, config: &kdc::config::KerberosServer) -> anyhow::Result<()> { +async fn serve_kdc_exchange( + mut stream: TcpStream, + config: &kdc::config::KerberosServer, + requests: &Mutex>, +) -> anyhow::Result<()> { let mut len_buf = [0u8; 4]; stream.read_exact(&mut len_buf).await.context("read KDC length")?; let len = usize::try_from(u32::from_be_bytes(len_buf)).context("KDC length")?; let mut body = vec![0u8; len]; stream.read_exact(&mut body).await.context("read KDC body")?; + requests.lock().expect("kdc request mutex").push(observe_kdc_req(&body)); + let mut raw = Vec::with_capacity(4 + len); raw.extend_from_slice(&len_buf); raw.extend_from_slice(&body); @@ -115,6 +141,61 @@ async fn serve_kdc_exchange(mut stream: TcpStream, config: &kdc::config::Kerbero Ok(()) } +fn principal_strings(name: &PrincipalName) -> Vec { + name.name_string.0.0.iter().map(|part| part.0.to_string()).collect() +} + +fn observe_kdc_req(body: &[u8]) -> ObservedKdcReq { + if let Ok(as_req) = picky_asn1_der::from_bytes::(body) { + let req = &as_req.0.req_body.0; + let cname = req + .cname + .0 + .as_ref() + .map(|name| principal_strings(&name.0).join("/")) + .unwrap_or_default(); + return ObservedKdcReq::As { + cname, + realm: req.realm.0.to_string(), + }; + } + if let Ok(tgs_req) = picky_asn1_der::from_bytes::(body) { + let req = &tgs_req.0.req_body.0; + let sname = req + .sname + .0 + .as_ref() + .map(|name| principal_strings(&name.0)) + .unwrap_or_default(); + return ObservedKdcReq::Tgs { + sname, + realm: req.realm.0.to_string(), + }; + } + ObservedKdcReq::Other +} + +fn observe_kdc_reply(body: &[u8]) -> ObservedKdcReply { + let krb = body.get(4..).unwrap_or(body); + if picky_asn1_der::from_bytes::(krb).is_ok() { + ObservedKdcReply::AsRep + } else if picky_asn1_der::from_bytes::(krb).is_ok() { + ObservedKdcReply::TgsRep + } else if picky_asn1_der::from_bytes::(krb).is_ok() { + ObservedKdcReply::KrbError + } else { + ObservedKdcReply::Other + } +} + +#[derive(Clone, Debug, PartialEq, Eq)] +enum ObservedKdcReply { + AsRep, + TgsRep, + KrbError, + Other, +} + #[derive(Clone)] enum MockRdpMode { Kerberos { kdc_url: String }, @@ -124,6 +205,8 @@ enum MockRdpMode { struct MockRdp { port: u16, credssp_ok: Arc, + finished_account: Arc>>, + cookies: Arc>>, } impl MockRdp { @@ -144,7 +227,11 @@ impl MockRdp { }; let port = listener.local_addr().context("mock RDP local_addr")?.port(); let credssp_ok = Arc::new(AtomicBool::new(false)); + let finished_account = Arc::new(Mutex::new(None)); + let cookies = Arc::new(Mutex::new(Vec::new())); let credssp_ok_task = Arc::clone(&credssp_ok); + let finished_account_task = Arc::clone(&finished_account); + let cookies_task = Arc::clone(&cookies); let acceptor = tls_acceptor()?; let public_key = server_public_key()?; @@ -156,13 +243,26 @@ impl MockRdp { let acceptor = acceptor.clone(); let public_key = public_key.clone(); let credssp_ok = Arc::clone(&credssp_ok_task); + let finished_account = Arc::clone(&finished_account_task); + let cookies = Arc::clone(&cookies_task); let mode = mode.clone(); tokio::spawn(async move { let result = match &mode { MockRdpMode::Kerberos { kdc_url } => { - accept_kerberos_rdp(stream, peer, acceptor, public_key, kdc_url).await + accept_kerberos_rdp( + stream, + peer, + acceptor, + public_key, + kdc_url, + &cookies, + &finished_account, + ) + .await + } + MockRdpMode::Ntlm => { + accept_ntlm_rdp(stream, peer, acceptor, public_key, &cookies, &finished_account).await } - MockRdpMode::Ntlm => accept_ntlm_rdp(stream, peer, acceptor, public_key).await, }; match result { Ok(()) => credssp_ok.store(true, Ordering::SeqCst), @@ -172,13 +272,26 @@ impl MockRdp { } }); - Ok(Self { port, credssp_ok }) + Ok(Self { + port, + credssp_ok, + finished_account, + cookies, + }) } fn credssp_ok(&self) -> bool { self.credssp_ok.load(Ordering::SeqCst) } + fn finished_account(&self) -> Option { + self.finished_account.lock().expect("finished account mutex").clone() + } + + fn cookies(&self) -> Vec { + self.cookies.lock().expect("cookie mutex").clone() + } + async fn wait_credssp(&self) -> anyhow::Result<()> { let deadline = Instant::now() + Duration::from_secs(30); loop { @@ -199,10 +312,13 @@ async fn accept_kerberos_rdp( acceptor: tokio_rustls::TlsAcceptor, public_key: Vec, kdc_url: &str, + cookies: &Mutex>, + finished_account: &Mutex>, ) -> anyhow::Result<()> { let mut framed = TokioFramed::new(stream); let (_, request) = framed.read_pdu().await.context("read X.224 CR")?; - let _: X224 = ironrdp_core::decode(&request).context("decode X.224 CR")?; + let cr: X224 = ironrdp_core::decode(&request).context("decode X.224 CR")?; + record_cookie(&cr, cookies); let confirm = X224(ConnectionConfirm::Response { flags: ResponseFlags::empty(), @@ -269,7 +385,11 @@ async fn accept_kerberos_rdp( buf.advance(length); framed.write_all(&buf[..length]).await.context("write CredSSP")?; } - sspi::credssp::ServerState::Finished(_) => return Ok(()), + sspi::credssp::ServerState::Finished(identity) => { + *finished_account.lock().expect("finished account mutex") = + Some(identity.username.account_name().to_owned()); + return Ok(()); + } } } anyhow::bail!("mock RDP CredSSP exceeded 6 round trips") @@ -280,10 +400,13 @@ async fn accept_ntlm_rdp( peer: std::net::SocketAddr, acceptor: tokio_rustls::TlsAcceptor, public_key: Vec, + cookies: &Mutex>, + finished_account: &Mutex>, ) -> anyhow::Result<()> { let mut framed = TokioFramed::new(stream); let (_, request) = framed.read_pdu().await.context("read X.224 CR")?; - let _: X224 = ironrdp_core::decode(&request).context("decode X.224 CR")?; + let cr: X224 = ironrdp_core::decode(&request).context("decode X.224 CR")?; + record_cookie(&cr, cookies); let confirm = X224(ConnectionConfirm::Response { flags: ResponseFlags::empty(), @@ -332,12 +455,22 @@ async fn accept_ntlm_rdp( buf.advance(length); framed.write_all(&buf[..length]).await.context("write CredSSP")?; } - sspi::credssp::ServerState::Finished(_) => return Ok(()), + sspi::credssp::ServerState::Finished(identity) => { + *finished_account.lock().expect("finished account mutex") = + Some(identity.username.account_name().to_owned()); + return Ok(()); + } } } anyhow::bail!("mock RDP NTLM CredSSP exceeded 6 round trips") } +fn record_cookie(cr: &X224, cookies: &Mutex>) { + if let Some(NegoRequestData::Cookie(cookie)) = &cr.0.nego_data { + cookies.lock().expect("cookie mutex").push(cookie.0.clone()); + } +} + async fn resolve_sspi_server( generator: &mut sspi::generator::Generator< '_, @@ -427,11 +560,11 @@ async fn connect_ntlm_client( } async fn complete_ntlm_credssp(tls: tokio_rustls::client::TlsStream) -> anyhow::Result<()> { - complete_client_credssp(tls, PROXY_USER, None, false).await + complete_client_credssp(tls, PROXY_USER, None, false, None).await } async fn complete_raw_ntlm_credssp(tls: tokio_rustls::client::TlsStream) -> anyhow::Result<()> { - complete_client_credssp(tls, PROXY_USER, None, true).await + complete_client_credssp(tls, PROXY_USER, None, true, None).await } async fn complete_client_credssp( @@ -439,6 +572,7 @@ async fn complete_client_credssp( username: &str, kdc_proxy_url: Option<&str>, raw_ntlm: bool, + proxy_replies: Option<&Mutex>>, ) -> anyhow::Result<()> { use sspi::credssp::{ClientMode, ClientState, CredSspClient, CredSspMode, TsRequest}; use sspi::ntlm::NtlmConfig; @@ -488,7 +622,7 @@ async fn complete_client_credssp( for _ in 0..8 { let client_state = { let mut generator = client.process(std::mem::take(&mut ts_request)); - resolve_sspi_client(&mut generator).await? + resolve_sspi_client(&mut generator, proxy_replies).await? }; let (outbound, finished) = match client_state { ClientState::ReplyNeeded(request) => (request, false), @@ -518,6 +652,7 @@ async fn resolve_sspi_client( sspi::Result>, sspi::Result, >, + proxy_replies: Option<&Mutex>>, ) -> anyhow::Result { let mut state = generator.start(); loop { @@ -525,7 +660,7 @@ async fn resolve_sspi_client( GeneratorState::Suspended(request) => { let reply = match request.url.scheme() { "tcp" | "udp" => send_kdc_tcp(&request).await?, - "http" | "https" => send_kdc_http(&request).await?, + "http" | "https" => send_kdc_http(&request, proxy_replies).await?, other => anyhow::bail!("unsupported KDC scheme {other}: {}", request.url), }; state = generator.resume(Ok(reply)); @@ -537,7 +672,10 @@ async fn resolve_sspi_client( } } -async fn send_kdc_http(request: &sspi::generator::NetworkRequest) -> anyhow::Result> { +async fn send_kdc_http( + request: &sspi::generator::NetworkRequest, + proxy_replies: Option<&Mutex>>, +) -> anyhow::Result> { let host = request.url.host_str().context("KDC proxy host")?; let port = request.url.port_or_known_default().unwrap_or(80); let path = if request.url.query().is_some() { @@ -586,19 +724,27 @@ async fn send_kdc_http(request: &sspi::generator::NetworkRequest) -> anyhow::Res } } - if let Some(len) = content_length { + let buf = if let Some(len) = content_length { let mut buf = vec![0u8; len]; tokio::io::AsyncReadExt::read_exact(&mut reader, &mut buf) .await .context("read KDC proxy body")?; - Ok(buf) + buf } else { let mut buf = Vec::new(); tokio::io::AsyncReadExt::read_to_end(&mut reader, &mut buf) .await .context("read KDC proxy eof body")?; - Ok(buf) + buf + }; + if let Ok(message) = KdcProxyMessage::from_raw(&buf) + && let Some(log) = proxy_replies + { + log.lock() + .expect("proxy reply mutex") + .push(observe_kdc_reply(&message.kerb_message.0.0)); } + Ok(buf) } fn kdc_inject_token(association_jti: &str) -> anyhow::Result { @@ -749,6 +895,27 @@ impl tokio_rustls::rustls::client::danger::ServerCertVerifier for NoCertificateV } } +fn assert_target_kdc_as_and_tgs(kdc: &MockKdc) -> anyhow::Result<()> { + let reqs = kdc.requests(); + anyhow::ensure!( + reqs.iter().any(|req| matches!( + req, + ObservedKdcReq::As { cname, realm } + if cname.eq_ignore_ascii_case("administrator") && realm.eq_ignore_ascii_case(REALM) + )), + "KDC must see AS-REQ cname=administrator realm={REALM}; requests={reqs:?}" + ); + anyhow::ensure!( + reqs.iter().any(|req| matches!( + req, + ObservedKdcReq::Tgs { sname, realm } + if *sname == ["TERMSRV", SERVICE_HOST] && realm.eq_ignore_ascii_case(REALM) + )), + "KDC must see TGS-REQ sname=TERMSRV/{SERVICE_HOST} realm={REALM}; requests={reqs:?}" + ); + Ok(()) +} + #[tokio::test] async fn kerberos_injection_completes_credssp_against_mock_kdc() -> anyhow::Result<()> { install_crypto_provider(); @@ -781,6 +948,18 @@ async fn kerberos_injection_completes_credssp_against_mock_kdc() -> anyhow::Resu kdc.exchanges(), gateway.logs.snapshot() ); + assert_target_kdc_as_and_tgs(&kdc)?; + anyhow::ensure!( + rdp.finished_account().as_deref() == Some("administrator"), + "RDP CredSSP Finished account must be administrator; got={:?}; cookies={:?}", + rdp.finished_account(), + rdp.cookies() + ); + anyhow::ensure!( + rdp.cookies().iter().any(|cookie| cookie == KERBEROS_TARGET_USER), + "RDP X.224 cookie must be {KERBEROS_TARGET_USER}; cookies={:?}", + rdp.cookies() + ); let _ = gateway.process.start_kill(); Ok(()) @@ -808,8 +987,9 @@ async fn kerberos_client_and_target_legs_complete_credssp() -> anyhow::Result<() .await?; let kdc_proxy = kdc_proxy_url(gateway.config.http_port(), &jti)?; + let proxy_replies = Mutex::new(Vec::new()); let tls = connect_ntlm_client(gateway.config.tcp_port(), &token).await?; - complete_client_credssp(tls, PROXY_KERBEROS_USER, Some(&kdc_proxy), false) + complete_client_credssp(tls, PROXY_KERBEROS_USER, Some(&kdc_proxy), false, Some(&proxy_replies)) .await .with_context(|| { format!( @@ -829,6 +1009,17 @@ async fn kerberos_client_and_target_legs_complete_credssp() -> anyhow::Result<() kdc.exchanges(), gateway.logs.snapshot() ); + assert_target_kdc_as_and_tgs(&kdc)?; + let replies = proxy_replies.lock().expect("proxy reply mutex").clone(); + anyhow::ensure!( + replies.contains(&ObservedKdcReply::AsRep) && replies.contains(&ObservedKdcReply::TgsRep), + "/jet/KdcProxy must return AS-REP and TGS-REP (PREAUTH KRB-ERROR is allowed first); replies={replies:?}" + ); + anyhow::ensure!( + rdp.finished_account().as_deref() == Some("administrator"), + "RDP CredSSP Finished account must be administrator; got={:?}", + rdp.finished_account() + ); let _ = gateway.process.start_kill(); Ok(()) @@ -859,8 +1050,9 @@ async fn kerberos_wrong_target_password_fails_closed() -> anyhow::Result<()> { let _ = complete_ntlm_credssp(tls).await; tokio::time::sleep(Duration::from_secs(2)).await; anyhow::ensure!( - !rdp.credssp_ok(), - "wrong target password must not complete Kerberos CredSSP; logs:\n{}", + !rdp.credssp_ok() && rdp.finished_account().is_none(), + "wrong target password must not complete Kerberos CredSSP; account={:?}; logs:\n{}", + rdp.finished_account(), gateway.logs.snapshot() ); let logs = gateway.logs.snapshot(); @@ -895,8 +1087,9 @@ async fn kerberos_kdc_down_fails_closed() -> anyhow::Result<()> { let _ = complete_ntlm_credssp(tls).await; tokio::time::sleep(Duration::from_secs(2)).await; anyhow::ensure!( - !rdp.credssp_ok(), - "unreachable KDC must not complete CredSSP; logs:\n{}", + !rdp.credssp_ok() && rdp.finished_account().is_none(), + "unreachable KDC must not complete CredSSP; account={:?}; logs:\n{}", + rdp.finished_account(), gateway.logs.snapshot() ); let logs = gateway.logs.snapshot(); @@ -958,6 +1151,16 @@ async fn ntlm_injection_completes_credssp_both_legs() -> anyhow::Result<()> { logs.contains("kerberos=false"), "expected NTLM injection; logs:\n{logs}" ); + anyhow::ensure!( + rdp.finished_account().as_deref() == Some(TARGET_USER), + "RDP NTLM CredSSP Finished account must be {TARGET_USER}; got={:?}", + rdp.finished_account() + ); + anyhow::ensure!( + rdp.cookies().iter().any(|cookie| cookie == TARGET_USER), + "RDP X.224 cookie must be {TARGET_USER}; cookies={:?}", + rdp.cookies() + ); let _ = gateway.process.start_kill(); Ok(()) From 4b2eca3a024934ac07d9baaa1a1471911909cfa8 Mon Sep 17 00:00:00 2001 From: Junyi Ou Date: Fri, 21 Aug 2026 11:32:56 -0400 Subject: [PATCH 6/9] test(dgw): decode cookies and KdcProxy AS-REQ principal Routing tests now decode X.224 Cookie instead of raw-byte search. Client-leg Kerberos asserts synthetic-KDC AS-REQ cname. Fail-closed Kerberos paths require inject-started then no Finished identity. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- testsuite/tests/cli/dgw/cred_injection.rs | 103 +++++++++--------- testsuite/tests/cli/dgw/cred_injection_kdc.rs | 62 ++++++++--- 2 files changed, 101 insertions(+), 64 deletions(-) diff --git a/testsuite/tests/cli/dgw/cred_injection.rs b/testsuite/tests/cli/dgw/cred_injection.rs index 6c5d236d4..c42d9430c 100644 --- a/testsuite/tests/cli/dgw/cred_injection.rs +++ b/testsuite/tests/cli/dgw/cred_injection.rs @@ -12,6 +12,8 @@ use std::time::{Duration, Instant}; use anyhow::Context as _; use base64::Engine as _; +use ironrdp_pdu::nego::{ConnectionRequest, NegoRequestData}; +use ironrdp_pdu::x224::X224; use testsuite::cli::{dgw_tokio_cmd, wait_for_tcp_port}; use testsuite::dgw_config::{DgwConfig, DgwConfigHandle, VerbosityProfile}; use tokio::io::{AsyncBufReadExt as _, AsyncReadExt as _, AsyncWriteExt as _, BufReader}; @@ -144,7 +146,7 @@ impl LogBuffer { struct FakeRdpTarget { port: u16, accepted: Arc, - payloads: Arc>>>, + cookies: Arc>>, } impl FakeRdpTarget { @@ -152,9 +154,9 @@ impl FakeRdpTarget { let listener = TcpListener::bind("127.0.0.1:0").await.context("bind fake RDP target")?; let port = listener.local_addr().context("fake RDP local_addr")?.port(); let accepted = Arc::new(AtomicUsize::new(0)); - let payloads = Arc::new(Mutex::new(Vec::new())); + let cookies = Arc::new(Mutex::new(Vec::new())); let accepted_task = Arc::clone(&accepted); - let payloads_task = Arc::clone(&payloads); + let cookies_task = Arc::clone(&cookies); tokio::spawn(async move { loop { @@ -162,14 +164,15 @@ impl FakeRdpTarget { break; }; accepted_task.fetch_add(1, Ordering::SeqCst); - let payloads = Arc::clone(&payloads_task); + let cookies = Arc::clone(&cookies_task); tokio::spawn(async move { let mut buf = vec![0u8; 4096]; // CredSSP cert generation can delay the rewritten X.224 CR. if let Ok(Ok(n)) = tokio::time::timeout(Duration::from_secs(30), stream.read(&mut buf)).await && n > 0 + && let Some(cookie) = decode_x224_cookie(&buf[..n]) { - payloads.lock().expect("payload mutex").push(buf[..n].to_vec()); + cookies.lock().expect("cookie mutex").push(cookie); } // Keep the accepted socket open so the proxy can finish writing the CR. tokio::time::sleep(Duration::from_secs(30)).await; @@ -180,7 +183,7 @@ impl FakeRdpTarget { Ok(Self { port, accepted, - payloads, + cookies, }) } @@ -188,18 +191,18 @@ impl FakeRdpTarget { self.accepted.load(Ordering::SeqCst) } - async fn wait_payloads(&self, count: usize) -> anyhow::Result>> { + async fn wait_cookies(&self, count: usize) -> anyhow::Result> { let deadline = Instant::now() + Duration::from_secs(30); loop { { - let payloads = self.payloads.lock().expect("payload mutex"); - if payloads.len() >= count { - return Ok(payloads.clone()); + let cookies = self.cookies.lock().expect("cookie mutex"); + if cookies.len() >= count { + return Ok(cookies.clone()); } } if Instant::now() >= deadline { anyhow::bail!( - "timed out waiting for {count} target payload(s); accepted={}", + "timed out waiting for {count} decoded X.224 cookie(s); accepted={}", self.accepted() ); } @@ -208,6 +211,14 @@ impl FakeRdpTarget { } } +fn decode_x224_cookie(payload: &[u8]) -> Option { + let cr: X224 = ironrdp_core::decode(payload).ok()?; + match cr.0.nego_data { + Some(NegoRequestData::Cookie(cookie)) => Some(cookie.0), + _ => None, + } +} + pub(crate) struct GatewayProc { pub(crate) config: DgwConfigHandle, pub(crate) process: Child, @@ -416,16 +427,6 @@ async fn connect_rdp_client(gateway_tcp: u16, association_jwt: &str) -> anyhow:: Ok(stream) } -fn cookie_line(username: &str) -> String { - format!("Cookie: mstshash={username}") -} - -fn payloads_contain(payloads: &[Vec], needle: &str) -> bool { - payloads - .iter() - .any(|payload| String::from_utf8_lossy(payload).contains(needle)) -} - #[tokio::test] async fn first_rdp_connection_injects_ntlm() -> anyhow::Result<()> { let target = FakeRdpTarget::start().await?; @@ -447,14 +448,11 @@ async fn first_rdp_connection_injects_ntlm() -> anyhow::Result<()> { "injection must not fall back to ordinary forward; logs:\n{logs}" ); - let payloads = target.wait_payloads(1).await?; - assert!( - payloads_contain(&payloads, &cookie_line(TARGET_USER)), - "target should see injected cookie; payloads={payloads:?}" - ); - assert!( - !payloads_contain(&payloads, &cookie_line(CLIENT_COOKIE)), - "target must not see the client cookie; payloads={payloads:?}" + let cookies = target.wait_cookies(1).await?; + assert_eq!( + cookies, + vec![TARGET_USER], + "decoded X.224 cookie must be the injected user" ); let _ = gateway.process.start_kill(); @@ -473,7 +471,7 @@ async fn reconnect_same_jwt_still_injects() -> anyhow::Result<()> { let first = connect_rdp_client(gateway.config.tcp_port(), &token).await?; gateway.logs.wait_count(INJECT_LOG, 1).await?; - target.wait_payloads(1).await?; + target.wait_cookies(1).await?; drop(first); let _second = connect_rdp_client(gateway.config.tcp_port(), &token).await?; @@ -483,13 +481,11 @@ async fn reconnect_same_jwt_still_injects() -> anyhow::Result<()> { "reconnect must keep injecting, not ordinary-forward; logs:\n{logs}" ); - let payloads = target.wait_payloads(2).await?; - assert_eq!(payloads.len(), 2, "both connections should reach the fake RDP target"); - assert!( - payloads - .iter() - .all(|payload| String::from_utf8_lossy(payload).contains(&cookie_line(TARGET_USER))), - "both reconnects should inject the target cookie; payloads={payloads:?}" + let cookies = target.wait_cookies(2).await?; + assert_eq!( + cookies, + vec![TARGET_USER, TARGET_USER], + "both decoded X.224 cookies must be the injected user" ); let _ = gateway.process.start_kill(); @@ -541,14 +537,11 @@ async fn unprovisioned_rdp_uses_ordinary_forward() -> anyhow::Result<()> { "absent mapping should ordinary-forward; logs:\n{logs}" ); - let payloads = target.wait_payloads(1).await?; - assert!( - payloads_contain(&payloads, &cookie_line(CLIENT_COOKIE)), - "ordinary forward should keep the client cookie; payloads={payloads:?}" - ); - assert!( - !payloads_contain(&payloads, &cookie_line(TARGET_USER)), - "ordinary forward must not invent an injection cookie; payloads={payloads:?}" + let cookies = target.wait_cookies(1).await?; + assert_eq!( + cookies, + vec![CLIENT_COOKIE], + "ordinary forward must keep the decoded client cookie" ); let _ = gateway.process.start_kill(); @@ -612,6 +605,11 @@ async fn kerberos_reconnect_reuses_generation_until_reprovision() -> anyhow::Res "newer provisioning generation should replace and still inject; logs:\n{logs}" ); let logs = gateway.logs.wait_count(PUBLISHED_KDC_LOG, 2).await?; + assert_eq!( + logs.matches(PUBLISHED_KDC_LOG).count(), + 2, + "re-provision must publish exactly one extra synthetic KDC; logs:\n{logs}" + ); assert_eq!( logs.matches(REGISTERED_KDC_LOG).count(), 3, @@ -622,12 +620,15 @@ async fn kerberos_reconnect_reuses_generation_until_reprovision() -> anyhow::Res "Kerberos injection must not ordinary-forward; logs:\n{logs}" ); - let payloads = target.wait_payloads(3).await?; - assert!( - payloads - .iter() - .all(|payload| String::from_utf8_lossy(payload).contains(&cookie_line(KERBEROS_TARGET_USER))), - "each generation should inject the Kerberos target username; payloads={payloads:?}" + let cookies = target.wait_cookies(3).await?; + assert_eq!( + cookies, + vec![ + KERBEROS_TARGET_USER.to_owned(), + KERBEROS_TARGET_USER.to_owned(), + KERBEROS_TARGET_USER.to_owned() + ], + "each decoded X.224 cookie must be the Kerberos target user" ); let _ = gateway.process.start_kill(); diff --git a/testsuite/tests/cli/dgw/cred_injection_kdc.rs b/testsuite/tests/cli/dgw/cred_injection_kdc.rs index e1c80d1be..010041a6c 100644 --- a/testsuite/tests/cli/dgw/cred_injection_kdc.rs +++ b/testsuite/tests/cli/dgw/cred_injection_kdc.rs @@ -560,11 +560,11 @@ async fn connect_ntlm_client( } async fn complete_ntlm_credssp(tls: tokio_rustls::client::TlsStream) -> anyhow::Result<()> { - complete_client_credssp(tls, PROXY_USER, None, false, None).await + complete_client_credssp(tls, PROXY_USER, None, false, None, None).await } async fn complete_raw_ntlm_credssp(tls: tokio_rustls::client::TlsStream) -> anyhow::Result<()> { - complete_client_credssp(tls, PROXY_USER, None, true, None).await + complete_client_credssp(tls, PROXY_USER, None, true, None, None).await } async fn complete_client_credssp( @@ -573,6 +573,7 @@ async fn complete_client_credssp( kdc_proxy_url: Option<&str>, raw_ntlm: bool, proxy_replies: Option<&Mutex>>, + proxy_requests: Option<&Mutex>>, ) -> anyhow::Result<()> { use sspi::credssp::{ClientMode, ClientState, CredSspClient, CredSspMode, TsRequest}; use sspi::ntlm::NtlmConfig; @@ -622,7 +623,7 @@ async fn complete_client_credssp( for _ in 0..8 { let client_state = { let mut generator = client.process(std::mem::take(&mut ts_request)); - resolve_sspi_client(&mut generator, proxy_replies).await? + resolve_sspi_client(&mut generator, proxy_replies, proxy_requests).await? }; let (outbound, finished) = match client_state { ClientState::ReplyNeeded(request) => (request, false), @@ -653,6 +654,7 @@ async fn resolve_sspi_client( sspi::Result, >, proxy_replies: Option<&Mutex>>, + proxy_requests: Option<&Mutex>>, ) -> anyhow::Result { let mut state = generator.start(); loop { @@ -660,7 +662,7 @@ async fn resolve_sspi_client( GeneratorState::Suspended(request) => { let reply = match request.url.scheme() { "tcp" | "udp" => send_kdc_tcp(&request).await?, - "http" | "https" => send_kdc_http(&request, proxy_replies).await?, + "http" | "https" => send_kdc_http(&request, proxy_replies, proxy_requests).await?, other => anyhow::bail!("unsupported KDC scheme {other}: {}", request.url), }; state = generator.resume(Ok(reply)); @@ -675,6 +677,7 @@ async fn resolve_sspi_client( async fn send_kdc_http( request: &sspi::generator::NetworkRequest, proxy_replies: Option<&Mutex>>, + proxy_requests: Option<&Mutex>>, ) -> anyhow::Result> { let host = request.url.host_str().context("KDC proxy host")?; let port = request.url.port_or_known_default().unwrap_or(80); @@ -699,6 +702,12 @@ async fn send_kdc_http( .context("write KDC proxy headers")?; stream.write_all(&request.data).await.context("write KDC proxy body")?; stream.flush().await.context("flush KDC proxy")?; + if let Ok(message) = KdcProxyMessage::from_raw(&request.data) + && let Some(log) = proxy_requests + { + let kerb = message.kerb_message.0.0.get(4..).unwrap_or(&message.kerb_message.0.0); + log.lock().expect("proxy request mutex").push(observe_kdc_req(kerb)); + } let mut reader = BufReader::new(stream); let mut status_line = String::new(); @@ -988,15 +997,23 @@ async fn kerberos_client_and_target_legs_complete_credssp() -> anyhow::Result<() let kdc_proxy = kdc_proxy_url(gateway.config.http_port(), &jti)?; let proxy_replies = Mutex::new(Vec::new()); + let proxy_requests = Mutex::new(Vec::new()); let tls = connect_ntlm_client(gateway.config.tcp_port(), &token).await?; - complete_client_credssp(tls, PROXY_KERBEROS_USER, Some(&kdc_proxy), false, Some(&proxy_replies)) - .await - .with_context(|| { - format!( - "client-leg Kerberos CredSSP; gateway logs:\n{}", - gateway.logs.snapshot() - ) - })?; + complete_client_credssp( + tls, + PROXY_KERBEROS_USER, + Some(&kdc_proxy), + false, + Some(&proxy_replies), + Some(&proxy_requests), + ) + .await + .with_context(|| { + format!( + "client-leg Kerberos CredSSP; gateway logs:\n{}", + gateway.logs.snapshot() + ) + })?; rdp.wait_credssp().await.with_context(|| { format!( "target-leg Kerberos CredSSP; gateway logs:\n{}", @@ -1015,6 +1032,15 @@ async fn kerberos_client_and_target_legs_complete_credssp() -> anyhow::Result<() replies.contains(&ObservedKdcReply::AsRep) && replies.contains(&ObservedKdcReply::TgsRep), "/jet/KdcProxy must return AS-REP and TGS-REP (PREAUTH KRB-ERROR is allowed first); replies={replies:?}" ); + let requests = proxy_requests.lock().expect("proxy request mutex").clone(); + anyhow::ensure!( + requests.iter().any(|req| matches!( + req, + ObservedKdcReq::As { cname, realm } + if cname.eq_ignore_ascii_case("injected-proxy-user") && realm.eq_ignore_ascii_case(REALM) + )), + "synthetic KDC AS-REQ must be proxy user injected-proxy-user@{REALM}; requests={requests:?}" + ); anyhow::ensure!( rdp.finished_account().as_deref() == Some("administrator"), "RDP CredSSP Finished account must be administrator; got={:?}", @@ -1048,6 +1074,11 @@ async fn kerberos_wrong_target_password_fails_closed() -> anyhow::Result<()> { let tls = connect_ntlm_client(gateway.config.tcp_port(), &token).await?; let _ = complete_ntlm_credssp(tls).await; + let logs = gateway.logs.wait_contains(INJECT_LOG).await?; + anyhow::ensure!( + logs.contains("kerberos=true"), + "wrong password must still start Kerberos injection; logs:\n{logs}" + ); tokio::time::sleep(Duration::from_secs(2)).await; anyhow::ensure!( !rdp.credssp_ok() && rdp.finished_account().is_none(), @@ -1085,6 +1116,11 @@ async fn kerberos_kdc_down_fails_closed() -> anyhow::Result<()> { let tls = connect_ntlm_client(gateway.config.tcp_port(), &token).await?; let _ = complete_ntlm_credssp(tls).await; + let logs = gateway.logs.wait_contains(INJECT_LOG).await?; + anyhow::ensure!( + logs.contains("kerberos=true"), + "KDC down must still start Kerberos injection; logs:\n{logs}" + ); tokio::time::sleep(Duration::from_secs(2)).await; anyhow::ensure!( !rdp.credssp_ok() && rdp.finished_account().is_none(), @@ -1120,7 +1156,7 @@ async fn kerberos_missing_krb_kdc_fails_closed() -> anyhow::Result<()> { stream.flush().await.context("flush CR")?; let logs = gateway.logs.wait_contains(MISSING_LOG).await?; anyhow::ensure!( - !logs.contains(FORWARD_LOG), + !logs.contains(FORWARD_LOG) && !logs.contains(INJECT_LOG), "missing krb_kdc must fail closed; logs:\n{logs}" ); From 79b95b132b3bbcffacafc6c18507c5379a3acdb5 Mon Sep 17 00:00:00 2001 From: Junyi Ou Date: Fri, 21 Aug 2026 12:34:21 -0400 Subject: [PATCH 7/9] test(dgw): tighten hop asserts without Gateway source changes Decode X.224 cookies, KdcProxy AS/TGS principals, and attribute KDC-down TCP to Gateway via a separate refusing listener. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- testsuite/tests/cli/dgw/cred_injection.rs | 11 +- testsuite/tests/cli/dgw/cred_injection_kdc.rs | 119 ++++++++++++++++-- 2 files changed, 119 insertions(+), 11 deletions(-) diff --git a/testsuite/tests/cli/dgw/cred_injection.rs b/testsuite/tests/cli/dgw/cred_injection.rs index c42d9430c..83810c432 100644 --- a/testsuite/tests/cli/dgw/cred_injection.rs +++ b/testsuite/tests/cli/dgw/cred_injection.rs @@ -476,6 +476,11 @@ async fn reconnect_same_jwt_still_injects() -> anyhow::Result<()> { let _second = connect_rdp_client(gateway.config.tcp_port(), &token).await?; let logs = gateway.logs.wait_count(INJECT_LOG, 2).await?; + assert_eq!( + logs.matches(INJECT_LOG).count(), + 2, + "reconnect must inject exactly twice; logs:\n{logs}" + ); assert!( !logs.contains(FORWARD_LOG), "reconnect must keep injecting, not ordinary-forward; logs:\n{logs}" @@ -514,7 +519,11 @@ async fn required_missing_fails_closed() -> anyhow::Result<()> { "expired mapping must fail closed, never silent ordinary forward; logs:\n{logs}" ); - tokio::time::sleep(Duration::from_millis(500)).await; + let deadline = Instant::now() + Duration::from_secs(2); + while Instant::now() < deadline { + anyhow::ensure!(target.accepted() == 0, "fail-closed routing must not connect upstream"); + tokio::time::sleep(Duration::from_millis(50)).await; + } assert_eq!(target.accepted(), 0, "fail-closed routing must not connect upstream"); let _ = gateway.process.start_kill(); diff --git a/testsuite/tests/cli/dgw/cred_injection_kdc.rs b/testsuite/tests/cli/dgw/cred_injection_kdc.rs index 010041a6c..dabb4e6fc 100644 --- a/testsuite/tests/cli/dgw/cred_injection_kdc.rs +++ b/testsuite/tests/cli/dgw/cred_injection_kdc.rs @@ -198,7 +198,7 @@ enum ObservedKdcReply { #[derive(Clone)] enum MockRdpMode { - Kerberos { kdc_url: String }, + Kerberos { kdc_url: Option }, Ntlm, } @@ -211,7 +211,7 @@ struct MockRdp { impl MockRdp { async fn start_kerberos(kdc_url: String) -> anyhow::Result { - Self::start(MockRdpMode::Kerberos { kdc_url }).await + Self::start(MockRdpMode::Kerberos { kdc_url: Some(kdc_url) }).await } async fn start_ntlm() -> anyhow::Result { @@ -254,7 +254,7 @@ impl MockRdp { peer, acceptor, public_key, - kdc_url, + kdc_url.as_deref(), &cookies, &finished_account, ) @@ -311,7 +311,7 @@ async fn accept_kerberos_rdp( peer: std::net::SocketAddr, acceptor: tokio_rustls::TlsAcceptor, public_key: Vec, - kdc_url: &str, + kdc_url: Option<&str>, cookies: &Mutex>, finished_account: &Mutex>, ) -> anyhow::Result<()> { @@ -339,7 +339,10 @@ async fn accept_kerberos_rdp( }; let kerberos_config = sspi::KerberosServerConfig { kerberos_config: sspi::KerberosConfig { - kdc_url: Some(kdc_url.parse().context("parse mock KDC URL")?), + kdc_url: kdc_url + .map(|url| url.parse()) + .transpose() + .context("parse mock KDC URL")?, client_computer_name: peer.to_string(), }, server_properties: sspi::kerberos::ServerProperties::new( @@ -715,7 +718,10 @@ async fn send_kdc_http( .read_line(&mut status_line) .await .context("read KDC proxy status")?; - anyhow::ensure!(status_line.contains("200"), "KDC proxy HTTP status was {status_line:?}"); + anyhow::ensure!( + status_line.starts_with("HTTP/1.1 200") || status_line.starts_with("HTTP/1.0 200"), + "KDC proxy HTTP status was {status_line:?}" + ); let mut content_length = None; loop { @@ -1041,11 +1047,24 @@ async fn kerberos_client_and_target_legs_complete_credssp() -> anyhow::Result<() )), "synthetic KDC AS-REQ must be proxy user injected-proxy-user@{REALM}; requests={requests:?}" ); + anyhow::ensure!( + requests.iter().any(|req| matches!( + req, + ObservedKdcReq::Tgs { sname, realm } + if *sname == ["TERMSRV", SERVICE_HOST] && realm.eq_ignore_ascii_case(REALM) + )), + "synthetic KDC TGS-REQ must be TERMSRV/{SERVICE_HOST}; requests={requests:?}" + ); anyhow::ensure!( rdp.finished_account().as_deref() == Some("administrator"), "RDP CredSSP Finished account must be administrator; got={:?}", rdp.finished_account() ); + anyhow::ensure!( + rdp.cookies().iter().any(|cookie| cookie == KERBEROS_TARGET_USER), + "RDP X.224 cookie must be {KERBEROS_TARGET_USER}; cookies={:?}", + rdp.cookies() + ); let _ = gateway.process.start_kill(); Ok(()) @@ -1079,6 +1098,25 @@ async fn kerberos_wrong_target_password_fails_closed() -> anyhow::Result<()> { logs.contains("kerberos=true"), "wrong password must still start Kerberos injection; logs:\n{logs}" ); + let deadline = Instant::now() + Duration::from_secs(5); + loop { + if kdc.requests().iter().any(|req| { + matches!( + req, + ObservedKdcReq::As { cname, realm } + if cname.eq_ignore_ascii_case("administrator") && realm.eq_ignore_ascii_case(REALM) + ) + }) { + break; + } + if Instant::now() >= deadline { + anyhow::bail!( + "timed out waiting for AS-REQ as administrator; requests={:?}", + kdc.requests() + ); + } + tokio::time::sleep(Duration::from_millis(50)).await; + } tokio::time::sleep(Duration::from_secs(2)).await; anyhow::ensure!( !rdp.credssp_ok() && rdp.finished_account().is_none(), @@ -1086,6 +1124,18 @@ async fn kerberos_wrong_target_password_fails_closed() -> anyhow::Result<()> { rdp.finished_account(), gateway.logs.snapshot() ); + anyhow::ensure!( + !kdc.requests() + .iter() + .any(|req| matches!(req, ObservedKdcReq::Tgs { .. })), + "wrong password must not obtain a TGS; requests={:?}", + kdc.requests() + ); + anyhow::ensure!( + rdp.cookies().iter().any(|cookie| cookie == KERBEROS_TARGET_USER), + "wrong password still rewrites the X.224 cookie; cookies={:?}", + rdp.cookies() + ); let logs = gateway.logs.snapshot(); anyhow::ensure!( !logs.contains(FORWARD_LOG), @@ -1096,10 +1146,43 @@ async fn kerberos_wrong_target_password_fails_closed() -> anyhow::Result<()> { Ok(()) } +struct RefusingKdc { + port: u16, + accepted: Arc, +} + +impl RefusingKdc { + async fn start() -> anyhow::Result { + let listener = TcpListener::bind("127.0.0.1:0").await.context("bind refusing KDC")?; + let port = listener.local_addr()?.port(); + let accepted = Arc::new(AtomicUsize::new(0)); + let accepted_task = Arc::clone(&accepted); + tokio::spawn(async move { + loop { + let Ok((_stream, _)) = listener.accept().await else { + break; + }; + accepted_task.fetch_add(1, Ordering::SeqCst); + } + }); + Ok(Self { port, accepted }) + } + + fn url(&self) -> String { + format!("tcp://127.0.0.1:{}", self.port) + } + + fn accepted(&self) -> usize { + self.accepted.load(Ordering::SeqCst) + } +} + #[tokio::test] async fn kerberos_kdc_down_fails_closed() -> anyhow::Result<()> { install_crypto_provider(); - let rdp = MockRdp::start_kerberos("tcp://127.0.0.1:1".to_owned()).await?; + let kdc = RefusingKdc::start().await?; + let rdp_kdc = MockKdc::start().await?; + let rdp = MockRdp::start_kerberos(rdp_kdc.url()).await?; let mut gateway = GatewayProc::start(true).await?; let jti = next_id(); @@ -1110,7 +1193,7 @@ async fn kerberos_kdc_down_fails_closed() -> anyhow::Result<()> { &token, KERBEROS_TARGET_USER, 300, - Some("tcp://127.0.0.1:1"), + Some(&kdc.url()), ) .await?; @@ -1121,13 +1204,30 @@ async fn kerberos_kdc_down_fails_closed() -> anyhow::Result<()> { logs.contains("kerberos=true"), "KDC down must still start Kerberos injection; logs:\n{logs}" ); - tokio::time::sleep(Duration::from_secs(2)).await; + let deadline = Instant::now() + Duration::from_secs(10); + while kdc.accepted() == 0 { + if Instant::now() >= deadline { + anyhow::bail!("Gateway never TCP-connected the provisioned KDC"); + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + tokio::time::sleep(Duration::from_secs(1)).await; anyhow::ensure!( !rdp.credssp_ok() && rdp.finished_account().is_none(), "unreachable KDC must not complete CredSSP; account={:?}; logs:\n{}", rdp.finished_account(), gateway.logs.snapshot() ); + anyhow::ensure!( + rdp.cookies().iter().any(|cookie| cookie == KERBEROS_TARGET_USER), + "KDC down still rewrites the X.224 cookie; cookies={:?}", + rdp.cookies() + ); + anyhow::ensure!( + kdc.accepted() >= 1, + "Gateway must TCP-connect the provisioned KDC; accepted={}", + kdc.accepted() + ); let logs = gateway.logs.snapshot(); anyhow::ensure!( !logs.contains(FORWARD_LOG), @@ -1159,7 +1259,6 @@ async fn kerberos_missing_krb_kdc_fails_closed() -> anyhow::Result<()> { !logs.contains(FORWARD_LOG) && !logs.contains(INJECT_LOG), "missing krb_kdc must fail closed; logs:\n{logs}" ); - let _ = gateway.process.start_kill(); Ok(()) } From 476b8c3e46b028c09c035b3fda6758ee0e4bc857 Mon Sep 17 00:00:00 2001 From: Junyi Ou Date: Fri, 21 Aug 2026 13:57:57 -0400 Subject: [PATCH 8/9] test(dgw): assert missing krb_kdc never dials the target Checkout-before-connect on the stack below lets the fail-closed path prove accepted==0. Issue: DGW-1900 Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- testsuite/tests/cli/dgw/cred_injection_kdc.rs | 16 +++++++++++++++- 1 file changed, 15 insertions(+), 1 deletion(-) diff --git a/testsuite/tests/cli/dgw/cred_injection_kdc.rs b/testsuite/tests/cli/dgw/cred_injection_kdc.rs index dabb4e6fc..8242923d4 100644 --- a/testsuite/tests/cli/dgw/cred_injection_kdc.rs +++ b/testsuite/tests/cli/dgw/cred_injection_kdc.rs @@ -1259,6 +1259,12 @@ async fn kerberos_missing_krb_kdc_fails_closed() -> anyhow::Result<()> { !logs.contains(FORWARD_LOG) && !logs.contains(INJECT_LOG), "missing krb_kdc must fail closed; logs:\n{logs}" ); + tokio::time::sleep(Duration::from_millis(250)).await; + anyhow::ensure!( + rdp.accepted() == 0, + "missing krb_kdc must not dial the target; accepted={}", + rdp.accepted() + ); let _ = gateway.process.start_kill(); Ok(()) } @@ -1303,20 +1309,28 @@ async fn ntlm_injection_completes_credssp_both_legs() -> anyhow::Result<()> { struct FakeClosedTarget { port: u16, + accepted: Arc, } impl FakeClosedTarget { async fn start() -> anyhow::Result { let listener = TcpListener::bind("127.0.0.1:0").await.context("bind closed target")?; let port = listener.local_addr()?.port(); + let accepted = Arc::new(AtomicUsize::new(0)); + let accepted_task = Arc::clone(&accepted); tokio::spawn(async move { loop { let Ok((_stream, _)) = listener.accept().await else { break; }; + accepted_task.fetch_add(1, Ordering::SeqCst); } }); - Ok(Self { port }) + Ok(Self { port, accepted }) + } + + fn accepted(&self) -> usize { + self.accepted.load(Ordering::SeqCst) } } From 2865edad02f60cad06fd1f6d832166e4ceb57e5a Mon Sep 17 00:00:00 2001 From: Junyi Ou Date: Fri, 21 Aug 2026 14:40:06 -0400 Subject: [PATCH 9/9] test(dgw): drive RDCleanPath injection with ironrdp-agent 0.1.0 Pin the public CLI release and complete NTLM and Kerberos target CredSSP over ws://127.0.0.1/jet/rdp. Issue: DGW-1900 Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- testsuite/tests/cli/dgw/cred_injection_kdc.rs | 267 ++++++++++++++++++ 1 file changed, 267 insertions(+) diff --git a/testsuite/tests/cli/dgw/cred_injection_kdc.rs b/testsuite/tests/cli/dgw/cred_injection_kdc.rs index 8242923d4..411bbbeaa 100644 --- a/testsuite/tests/cli/dgw/cred_injection_kdc.rs +++ b/testsuite/tests/cli/dgw/cred_injection_kdc.rs @@ -3,7 +3,12 @@ //! Proves the target-leg path: Gateway fetches tickets from a TCP KDC (`kdc` crate from //! sspi-rs) and completes CredSSP with a fake RDP acceptor. The Gateway-facing client uses //! NTLM so the test does not depend on the in-process synthetic KDC. +//! +//! RDCleanPath coverage drives the public `ironrdp-agent` 0.1.0 CLI (`cargo install +//! ironrdp-agent --version 0.1.0`) over `ws://127.0.0.1/jet/rdp`. +use std::path::{Path, PathBuf}; +use std::process::Stdio; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; @@ -1334,6 +1339,268 @@ impl FakeClosedTarget { } } +const IRONRDP_AGENT_VERSION: &str = "0.1.0"; +const RDCLEANPATH_INJECT_LOG: &str = "Switching to RdpProxy for credential injection (WebSocket)"; +const RDCLEANPATH_FORWARD_LOG: &str = "RDP-TLS forwarding (RDCleanPath)"; + +fn ironrdp_agent_bin() -> Option { + if let Ok(path) = std::env::var("IRONRDP_AGENT") { + return Some(PathBuf::from(path)); + } + let name = if cfg!(windows) { + "ironrdp-agent.exe" + } else { + "ironrdp-agent" + }; + if let Ok(home) = std::env::var("CARGO_HOME") { + let path = PathBuf::from(home).join("bin").join(name); + if path.is_file() { + return Some(path); + } + } + let cargo_home = std::env::var_os("USERPROFILE") + .or_else(|| std::env::var_os("HOME")) + .map(PathBuf::from) + .map(|home| home.join(".cargo").join("bin").join(name)); + if let Some(path) = cargo_home + && path.is_file() + { + return Some(path); + } + if let Ok(path) = std::env::var("PATH") { + for dir in std::env::split_paths(&path) { + let candidate = dir.join(name); + if candidate.is_file() { + return Some(candidate); + } + } + } + None +} + +fn require_ironrdp_agent() -> anyhow::Result> { + let Some(bin) = ironrdp_agent_bin() else { + eprintln!( + "skipping RDCleanPath ironrdp-agent test: cargo install ironrdp-agent --version {IRONRDP_AGENT_VERSION}" + ); + return Ok(None); + }; + let output = std::process::Command::new(&bin) + .arg("--version") + .output() + .with_context(|| format!("run {} --version", bin.display()))?; + let version = String::from_utf8_lossy(&output.stdout); + anyhow::ensure!( + version.contains(IRONRDP_AGENT_VERSION), + "expected ironrdp-agent {IRONRDP_AGENT_VERSION}, got {version:?} from {}", + bin.display() + ); + Ok(Some(bin)) +} + +fn ironrdp_agent_endpoint() -> String { + let name = format!("ironrdp-e2e-{}", next_id().replace('-', "")); + if cfg!(windows) { + format!(r"\\.\pipe\{name}") + } else { + std::env::temp_dir().join(format!("{name}.sock")).display().to_string() + } +} + +async fn start_ironrdp_daemon(bin: &Path, endpoint: &str) -> anyhow::Result { + let child = tokio::process::Command::new(bin) + .args(["--endpoint", endpoint, "daemon-start"]) + .kill_on_drop(true) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .context("start ironrdp-agent daemon")?; + let deadline = Instant::now() + Duration::from_secs(10); + loop { + let status = tokio::process::Command::new(bin) + .args(["--endpoint", endpoint, "status"]) + .output() + .await + .context("ironrdp-agent status")?; + if status.status.success() { + return Ok(child); + } + if Instant::now() >= deadline { + anyhow::bail!( + "ironrdp-agent daemon not ready at {endpoint}: {}", + String::from_utf8_lossy(&status.stderr) + ); + } + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + +async fn connect_ironrdp_rdcleanpath( + bin: &Path, + endpoint: &str, + server: &str, + token: &str, + http_port: u16, +) -> anyhow::Result { + let url = format!("ws://127.0.0.1:{http_port}/jet/rdp"); + tokio::process::Command::new(bin) + .args([ + "--endpoint", + endpoint, + "connect", + "--server", + server, + "--username", + PROXY_USER, + "--password", + PROXY_PASSWORD, + "--prop", + &format!("ironrdp_rdcleanpathurl:s:{url}"), + "--prop", + &format!("ironrdp_rdcleanpathtoken:s:{token}"), + ]) + .kill_on_drop(true) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .context("start ironrdp-agent connect") +} + +async fn agent_query_logs(bin: &Path, endpoint: &str) -> String { + tokio::process::Command::new(bin) + .args(["--endpoint", endpoint, "query-logs"]) + .output() + .await + .ok() + .map(|output| String::from_utf8_lossy(&output.stdout).into_owned()) + .unwrap_or_default() +} + +#[tokio::test] +async fn ironrdp_agent_rdcleanpath_ntlm_injection() -> anyhow::Result<()> { + let Some(bin) = require_ironrdp_agent()? else { + return Ok(()); + }; + install_crypto_provider(); + let rdp = MockRdp::start_ntlm().await?; + let mut gateway = GatewayProc::start(false).await?; + let endpoint = ironrdp_agent_endpoint(); + let mut daemon = start_ironrdp_daemon(&bin, &endpoint).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token_for_host(&jti, &jet_aid, rdp.port, 60)?; + provision_credentials(gateway.config.http_port(), &token, TARGET_USER, 300, None).await?; + + let mut connect = connect_ironrdp_rdcleanpath( + &bin, + &endpoint, + &format!("{SERVICE_HOST}:{}", rdp.port), + &token, + gateway.config.http_port(), + ) + .await?; + let wait = rdp.wait_credssp().await; + let agent_logs = agent_query_logs(&bin, &endpoint).await; + wait.with_context(|| { + format!( + "RDCleanPath NTLM target CredSSP; gateway logs:\n{}\nagent logs:\n{agent_logs}", + gateway.logs.snapshot() + ) + })?; + let logs = gateway.logs.snapshot(); + anyhow::ensure!( + logs.contains(RDCLEANPATH_INJECT_LOG), + "RDCleanPath must take the injection path; logs:\n{logs}" + ); + anyhow::ensure!( + !logs.contains(RDCLEANPATH_FORWARD_LOG), + "RDCleanPath injection must not ordinary-forward; logs:\n{logs}" + ); + anyhow::ensure!( + rdp.finished_account().as_deref() == Some(TARGET_USER), + "RDP NTLM CredSSP Finished account must be {TARGET_USER}; got={:?}", + rdp.finished_account() + ); + anyhow::ensure!( + rdp.cookies().iter().any(|cookie| cookie == PROXY_USER), + "RDCleanPath forwards the client X.224 cookie; cookies={:?}", + rdp.cookies() + ); + + let _ = connect.start_kill(); + let _ = daemon.start_kill(); + let _ = gateway.process.start_kill(); + Ok(()) +} + +#[tokio::test] +async fn ironrdp_agent_rdcleanpath_kerberos_injection() -> anyhow::Result<()> { + let Some(bin) = require_ironrdp_agent()? else { + return Ok(()); + }; + install_crypto_provider(); + let kdc = MockKdc::start().await?; + let rdp = MockRdp::start_kerberos(kdc.url()).await?; + let mut gateway = GatewayProc::start(true).await?; + let endpoint = ironrdp_agent_endpoint(); + let mut daemon = start_ironrdp_daemon(&bin, &endpoint).await?; + + let jti = next_id(); + let jet_aid = next_id(); + let token = association_token_for_host(&jti, &jet_aid, rdp.port, 60)?; + provision_credentials( + gateway.config.http_port(), + &token, + KERBEROS_TARGET_USER, + 300, + Some(&kdc.url()), + ) + .await?; + + let mut connect = connect_ironrdp_rdcleanpath( + &bin, + &endpoint, + &format!("{SERVICE_HOST}:{}", rdp.port), + &token, + gateway.config.http_port(), + ) + .await?; + let wait = rdp.wait_credssp().await; + let agent_logs = agent_query_logs(&bin, &endpoint).await; + wait.with_context(|| { + format!( + "RDCleanPath Kerberos target CredSSP; gateway logs:\n{}\nagent logs:\n{agent_logs}", + gateway.logs.snapshot() + ) + })?; + assert_target_kdc_as_and_tgs(&kdc)?; + let logs = gateway.logs.snapshot(); + anyhow::ensure!( + logs.contains(RDCLEANPATH_INJECT_LOG), + "RDCleanPath must take the injection path; logs:\n{logs}" + ); + anyhow::ensure!( + !logs.contains(RDCLEANPATH_FORWARD_LOG), + "RDCleanPath injection must not ordinary-forward; logs:\n{logs}" + ); + anyhow::ensure!( + rdp.finished_account().as_deref() == Some("administrator"), + "RDP CredSSP Finished account must be administrator; got={:?}", + rdp.finished_account() + ); + anyhow::ensure!( + rdp.cookies().iter().any(|cookie| cookie == PROXY_USER), + "RDCleanPath forwards the client X.224 cookie; cookies={:?}", + rdp.cookies() + ); + + let _ = connect.start_kill(); + let _ = daemon.start_kill(); + let _ = gateway.process.start_kill(); + Ok(()) +} + const CERT_PEM: &str = r#"-----BEGIN CERTIFICATE----- MIIDCzCCAfOgAwIBAgIUPRJa8i280unV3/kW6TE2fSUw8PwwDQYJKoZIhvcNAQEL BQAwFDESMBAGA1UEAwwJbG9jYWxob3N0MCAXDTI1MTEyNTA5NDAzMFoYDzIxMjUx