//! Comprehensive Integration and Lifecycle Test Suite for NX9-WG Optional Upstream interfaces. //! //! Verifies: //! - ProtonVPN-style .conf import, parsing, validation, persistence, and kernel synchronization //! - wg0 overlay non-regression during all upstream operations //! - Upstream enable, disable, restart, and deletion lifecycles //! - Reconciliation engine drift detection, convergence, and orphan cleanup //! - Zero secret leakage across API preview, import, status, list, and CLI use axum::body::{Body, to_bytes}; use axum::http::{Request, StatusCode, header}; use chrono::Utc; use nx9_wg_api::auth::{BootstrapOptions, bootstrap_admin}; use nx9_wg_api::collect_managed_wg_subnets; use nx9_wg_api::reconciliation::ReconciliationEngine; use nx9_wg_api::routes::build_api_router; use nx9_wg_api::routes::cli::{ExecuteCliRequest, build_safe_argv, scrub_secrets}; use nx9_wg_api::state::AppState; use nx9_wg_core::config::AppConfig; use nx9_wg_core::crypto::generate_keypair; use nx9_wg_core::types::wireguard::{ Interface, InterfaceRole, Peer, PeerProfile, PeerState, PeerType, }; use nx9_wg_core::validation::validate_cidr; use nx9_wg_db::Store; use nx9_wg_network::SimulatedNetworkEngine; use nx9_wireguard::{LiveInterfaceStats, SimulatedWireGuardEngine, WireGuardEngine}; use serde_json::{Value, json}; use std::collections::HashMap; use std::sync::Arc; use tempfile::{TempDir, tempdir}; use tower::ServiceExt; use uuid::Uuid; struct TestHarness { _dir: TempDir, store: Store, _state: AppState, wg_engine: Arc, _net_engine: Arc, reconciler: Arc, app: axum::Router, session_cookie: String, } async fn setup_test_harness() -> TestHarness { let dir = tempdir().expect("create temp dir"); let db_path = dir.path().join("upstream_test.db"); let store = Store::connect(&db_path.to_string_lossy()) .await .expect("connect to db"); store.migrate().await.expect("run migrations"); let config = AppConfig::default(); let opts = BootstrapOptions { cli_password: Some("AdminSecret123!".to_string()), ..Default::default() }; bootstrap_admin(&store, &config, &opts) .await .expect("bootstrap admin"); let wg_engine = Arc::new(SimulatedWireGuardEngine::new()); let net_engine = Arc::new(SimulatedNetworkEngine::new()); let state = AppState::with_engines(store.clone(), wg_engine.clone(), net_engine.clone()); let reconciler = Arc::new(ReconciliationEngine::new( state.clone(), wg_engine.clone(), net_engine.clone(), )); let app = build_api_router(state.clone()); // Login to get session ID let login_req = Request::builder() .method("POST") .uri("/api/v1/auth/login") .header(header::CONTENT_TYPE, "application/json") .body(Body::from( json!({ "username": "admin", "password": "AdminSecret123!" }) .to_string(), )) .unwrap(); let resp = app.clone().oneshot(login_req).await.expect("login request"); assert_eq!(resp.status(), StatusCode::OK); let cookie_header = resp .headers() .get(header::SET_COOKIE) .expect("set-cookie") .to_str() .unwrap(); let session_cookie = cookie_header.split(';').next().unwrap().to_string(); let now = Utc::now().naive_utc(); let (wg0_priv, wg0_pub) = generate_keypair(); let wg0 = Interface { id: Uuid::new_v4(), name: "wg0".to_string(), role: InterfaceRole::Overlay, private_key: wg0_priv, public_key: wg0_pub.clone(), listen_port: Some(51820), address_v4: validate_cidr("10.100.0.1/24").unwrap(), address_v6: None, mtu: Some(1420), dns: Some("1.1.1.1".to_string()), enabled: true, pre_up: None, post_up: None, pre_down: None, post_down: None, created_at: now, updated_at: now, }; store.create_interface(&wg0).await.unwrap(); let (client_priv, client_pub) = generate_keypair(); let client_peer = Peer { id: Uuid::new_v4(), interface_id: wg0.id, name: "client-alice".to_string(), peer_type: PeerType::RoadWarrior, state: PeerState::Active, public_key: client_pub, private_key: Some(client_priv), preshared_key: None, endpoint: None, allowed_ips: "10.100.0.2/32".to_string(), server_allowed_ips: None, address_v4: Some(validate_cidr("10.100.0.2/32").unwrap()), address_v6: None, dns: None, mtu: None, persistent_keepalive: Some(25), profile: PeerProfile::FullTunnel, expires_at: None, last_handshake_at: None, created_at: now, updated_at: now, }; store.create_peer(&client_peer).await.unwrap(); // Baseline reconciliation to converge initial network/firewall/wg state reconciler.apply().await.unwrap(); TestHarness { _dir: dir, store, _state: state, wg_engine, _net_engine: net_engine, reconciler, app, session_cookie, } } fn sample_proton_conf(priv_k_str: &str, provider_pub_k_str: &str) -> String { format!( r#" # ProtonVPN WireGuard Configuration [Interface] PrivateKey = {} Address = 10.2.0.2/32 DNS = 10.2.0.1 MTU = 1420 [Peer] PublicKey = {} AllowedIPs = 0.0.0.0/0, ::/0 Endpoint = 37.19.199.155:51820 PersistentKeepalive = 25 "#, priv_k_str, provider_pub_k_str ) } #[tokio::test] async fn test_proton0_import_and_kernel_sync() { let harness = setup_test_harness().await; let (priv_k, pub_k) = generate_keypair(); let (_, provider_pub_k) = generate_keypair(); let conf = sample_proton_conf(priv_k.as_str(), provider_pub_k.as_str()); // 1. Preview API endpoint (read-only, no side effects) let preview_req = Request::builder() .method("POST") .uri("/api/v1/interfaces/upstreams/preview") .header(header::COOKIE, &harness.session_cookie) .header(header::CONTENT_TYPE, "application/json") .body(Body::from( json!({ "name": "proton0", "config": conf }) .to_string(), )) .unwrap(); let preview_resp = harness.app.clone().oneshot(preview_req).await.unwrap(); assert_eq!(preview_resp.status(), StatusCode::OK); let preview_body: Value = serde_json::from_slice( &to_bytes(preview_resp.into_body(), usize::MAX) .await .unwrap(), ) .unwrap(); assert_eq!(preview_body["name"], "proton0"); assert_eq!(preview_body["role"], "upstream"); assert_eq!(preview_body["address_v4"], "10.2.0.2/32"); assert_eq!(preview_body["dns"], "10.2.0.1"); assert_eq!(preview_body["provider_public_key"], provider_pub_k.as_str()); assert_eq!(preview_body["provider_endpoint"], "37.19.199.155:51820"); assert_eq!(preview_body["provider_allowed_ips"], "0.0.0.0/0, ::/0"); assert_eq!(preview_body["persistent_keepalive"], 25); // Ensure secrets are never in response assert!(preview_body.get("private_key").is_none()); assert!(preview_body.get("preshared_key").is_none()); // Verify DB still only has wg0 (preview didn't write to DB) assert_eq!(harness.store.list_interfaces().await.unwrap().len(), 1); // 2. Import API endpoint (transactional persistence + kernel sync) let import_req = Request::builder() .method("POST") .uri("/api/v1/interfaces/upstreams/import") .header(header::COOKIE, &harness.session_cookie) .header(header::CONTENT_TYPE, "application/json") .body(Body::from( json!({ "name": "proton0", "config": conf }) .to_string(), )) .unwrap(); let import_resp = harness.app.clone().oneshot(import_req).await.unwrap(); assert_eq!(import_resp.status(), StatusCode::OK); let import_body: Value = serde_json::from_slice(&to_bytes(import_resp.into_body(), usize::MAX).await.unwrap()) .unwrap(); let iface_id = import_body["interface_id"].as_str().unwrap(); let peer_id = import_body["peer_id"].as_str().unwrap(); assert_eq!(import_body["name"], "proton0"); assert_eq!(import_body["role"], "upstream"); assert!(import_body.get("private_key").is_none()); assert!(import_body.get("preshared_key").is_none()); // 3. Verify SQLite desired state let iface = harness .store .get_interface(Uuid::parse_str(iface_id).unwrap()) .await .unwrap() .expect("proton0 in db"); assert_eq!(iface.name, "proton0"); assert_eq!(iface.role, InterfaceRole::Upstream); assert_eq!(iface.public_key.as_str(), pub_k.as_str()); let peers = harness .store .list_peers_for_interface(iface.id) .await .unwrap(); assert_eq!(peers.len(), 1); assert_eq!(peers[0].id.to_string(), peer_id); assert_eq!(peers[0].public_key.as_str(), provider_pub_k.as_str()); assert_eq!(peers[0].allowed_ips, "0.0.0.0/0, ::/0"); // 4. Verify Kernel Simulation state let kernel_stats = harness .wg_engine .get_interface_stats("proton0") .await .unwrap() .expect("proton0 in kernel"); assert_eq!(kernel_stats.name, "proton0"); assert!(kernel_stats.is_up); assert_eq!(kernel_stats.peers.len(), 1); assert_eq!(kernel_stats.peers[0].public_key, provider_pub_k.as_str()); assert_eq!( kernel_stats.peers[0].endpoint, Some("37.19.199.155:51820".to_string()) ); assert_eq!( kernel_stats.peers[0].allowed_ips, vec!["0.0.0.0/0".to_string(), "::/0".to_string()] ); assert_eq!(kernel_stats.peers[0].persistent_keepalive, Some(25)); } #[tokio::test] async fn test_wg0_non_regression_during_upstream_operations() { let harness = setup_test_harness().await; let (priv_k, _) = generate_keypair(); let (_, provider_pub_k) = generate_keypair(); let conf = sample_proton_conf(priv_k.as_str(), provider_pub_k.as_str()); // Import proton0 let import_req = Request::builder() .method("POST") .uri("/api/v1/interfaces/upstreams/import") .header(header::COOKIE, &harness.session_cookie) .header(header::CONTENT_TYPE, "application/json") .body(Body::from( json!({ "name": "proton0", "config": conf }) .to_string(), )) .unwrap(); let resp = harness.app.clone().oneshot(import_req).await.unwrap(); assert_eq!(resp.status(), StatusCode::OK); // 1. wg0 remains Overlay let wg0 = harness .store .get_interface_by_name("wg0") .await .unwrap() .expect("wg0 exists"); assert_eq!(wg0.role, InterfaceRole::Overlay); assert_eq!(wg0.address_v4.to_string(), "10.100.0.1/24"); // 2. wg0 peers unchanged and RoadWarrior AllowedIPs remain strictly /32 let wg0_peers = harness .store .list_peers_for_interface(wg0.id) .await .unwrap(); assert_eq!(wg0_peers.len(), 1); assert_eq!(wg0_peers[0].name, "client-alice"); assert_eq!( wg0_peers[0].server_wireguard_allowed_ips_for_role(InterfaceRole::Overlay), "10.100.0.2/32" ); // 3. Managed subnets for client NAT masquerade only includes Overlay interfaces let subnets = collect_managed_wg_subnets(&harness.store).await.unwrap(); assert_eq!(subnets.len(), 1); assert_eq!(subnets[0].to_string(), "10.100.0.1/24"); // proton0 address (10.2.0.2/32) is NOT in client NAT subnets! assert!(!subnets.iter().any(|s| s.to_string().contains("10.2.0.2"))); // 4. Reconciliation plan reports zero drift let plan = harness.reconciler.plan().await.unwrap(); assert!(!plan.has_drift, "Plan must be clean and fully converged"); } #[tokio::test] async fn test_upstream_restart_lifecycle() { let harness = setup_test_harness().await; let (priv_k, _) = generate_keypair(); let (_, provider_pub_k) = generate_keypair(); let conf = sample_proton_conf(priv_k.as_str(), provider_pub_k.as_str()); // Import proton0 let import_req = Request::builder() .method("POST") .uri("/api/v1/interfaces/upstreams/import") .header(header::COOKIE, &harness.session_cookie) .header(header::CONTENT_TYPE, "application/json") .body(Body::from( json!({ "name": "proton0", "config": conf }) .to_string(), )) .unwrap(); let import_resp = harness.app.clone().oneshot(import_req).await.unwrap(); let import_body: Value = serde_json::from_slice(&to_bytes(import_resp.into_body(), usize::MAX).await.unwrap()) .unwrap(); let iface_id = import_body["interface_id"].as_str().unwrap(); // Restart proton0 let restart_req = Request::builder() .method("POST") .uri(format!("/api/v1/interfaces/{iface_id}/restart")) .header(header::COOKIE, &harness.session_cookie) .body(Body::empty()) .unwrap(); let restart_resp = harness.app.clone().oneshot(restart_req).await.unwrap(); assert_eq!(restart_resp.status(), StatusCode::OK); // Verify same interface ID in DB let iface_after = harness .store .get_interface(Uuid::parse_str(iface_id).unwrap()) .await .unwrap() .expect("iface exists"); assert_eq!(iface_after.name, "proton0"); assert_eq!(iface_after.role, InterfaceRole::Upstream); // Verify provider peer restored in kernel let kernel_stats = harness .wg_engine .get_interface_stats("proton0") .await .unwrap() .expect("proton0 live"); assert_eq!(kernel_stats.peers.len(), 1); assert_eq!(kernel_stats.peers[0].public_key, provider_pub_k.as_str()); assert_eq!( kernel_stats.peers[0].allowed_ips, vec!["0.0.0.0/0".to_string(), "::/0".to_string()] ); } #[tokio::test] async fn test_upstream_delete_lifecycle() { let harness = setup_test_harness().await; let (priv_k, _) = generate_keypair(); let (_, provider_pub_k) = generate_keypair(); let conf = sample_proton_conf(priv_k.as_str(), provider_pub_k.as_str()); // Import proton0 let import_req = Request::builder() .method("POST") .uri("/api/v1/interfaces/upstreams/import") .header(header::COOKIE, &harness.session_cookie) .header(header::CONTENT_TYPE, "application/json") .body(Body::from( json!({ "name": "proton0", "config": conf }) .to_string(), )) .unwrap(); let import_resp = harness.app.clone().oneshot(import_req).await.unwrap(); let import_body: Value = serde_json::from_slice(&to_bytes(import_resp.into_body(), usize::MAX).await.unwrap()) .unwrap(); let iface_id = import_body["interface_id"].as_str().unwrap(); // Verify present in kernel before delete assert!( harness .wg_engine .get_interface_stats("proton0") .await .unwrap() .is_some() ); // Delete proton0 let del_req = Request::builder() .method("DELETE") .uri(format!("/api/v1/interfaces/{iface_id}")) .header(header::COOKIE, &harness.session_cookie) .body(Body::empty()) .unwrap(); let del_resp = harness.app.clone().oneshot(del_req).await.unwrap(); assert_eq!(del_resp.status(), StatusCode::OK); // Verify absent from kernel assert!( harness .wg_engine .get_interface_stats("proton0") .await .unwrap() .is_none() ); // Verify absent from DB assert!( harness .store .get_interface(Uuid::parse_str(iface_id).unwrap()) .await .unwrap() .is_none() ); // Verify wg0 remains untouched assert!( harness .store .get_interface_by_name("wg0") .await .unwrap() .is_some() ); } #[tokio::test] async fn test_upstream_reconciliation_orphan_detection() { let harness = setup_test_harness().await; // Inject an orphan upstream interface into simulated kernel harness .wg_engine .inject_interface_stats(LiveInterfaceStats { name: "orphan_vpn0".to_string(), public_key: "orphanpubkey12345".to_string(), listen_port: 51830, fwmark: 0, peers: vec![], addresses: vec!["10.99.0.1/24".to_string()], mtu: Some(1420), is_up: true, }) .await; // Detect orphan in plan let plan = harness.reconciler.plan().await.unwrap(); assert!(plan.has_drift); let orphan_action = plan .actions .iter() .find(|a| a.resource_id == "orphan_vpn0") .expect("orphan action in plan"); assert_eq!(orphan_action.action_type, "delete_orphan_interface"); // Apply cleanup let report = harness.reconciler.apply().await.unwrap(); assert!( report .details .iter() .any(|d| d.contains("Removed orphan kernel interface 'orphan_vpn0'")) ); // Verify orphan was deleted from kernel assert!( harness .wg_engine .get_interface_stats("orphan_vpn0") .await .unwrap() .is_none() ); // Verify wg0 remains active assert!( harness .wg_engine .get_interface_stats("wg0") .await .unwrap() .is_some() ); } #[tokio::test] async fn test_upstream_secret_safety() { let harness = setup_test_harness().await; let (priv_k, _) = generate_keypair(); let (_, provider_pub_k) = generate_keypair(); let raw_priv = priv_k.as_str().to_string(); let conf = sample_proton_conf(&raw_priv, provider_pub_k.as_str()); // 1. Preview response secret check let preview_req = Request::builder() .method("POST") .uri("/api/v1/interfaces/upstreams/preview") .header(header::COOKIE, &harness.session_cookie) .header(header::CONTENT_TYPE, "application/json") .body(Body::from( json!({ "name": "proton0", "config": conf }) .to_string(), )) .unwrap(); let preview_resp = harness.app.clone().oneshot(preview_req).await.unwrap(); let preview_text = String::from_utf8( to_bytes(preview_resp.into_body(), usize::MAX) .await .unwrap() .to_vec(), ) .unwrap(); assert!( !preview_text.contains(&raw_priv), "PrivateKey leaked in preview response" ); // 2. Import response secret check let import_req = Request::builder() .method("POST") .uri("/api/v1/interfaces/upstreams/import") .header(header::COOKIE, &harness.session_cookie) .header(header::CONTENT_TYPE, "application/json") .body(Body::from( json!({ "name": "proton0", "config": conf }) .to_string(), )) .unwrap(); let import_resp = harness.app.clone().oneshot(import_req).await.unwrap(); let import_text = String::from_utf8( to_bytes(import_resp.into_body(), usize::MAX) .await .unwrap() .to_vec(), ) .unwrap(); assert!( !import_text.contains(&raw_priv), "PrivateKey leaked in import response" ); // 3. Read-only CLI output secret scrubber check let scrubbed = scrub_secrets(&format!( "private_key: {}\nPrivateKey = {}", raw_priv, raw_priv )); assert!( !scrubbed.contains(&raw_priv), "PrivateKey leaked past scrubber" ); // 4. Safe argv builder allows read-only Upstream queries let list_req = ExecuteCliRequest { command: "interface".to_string(), subcommand: Some("upstream".to_string()), sub_subcommand: Some("list".to_string()), target: None, parameters: HashMap::new(), }; let argv = build_safe_argv(&list_req).unwrap(); assert_eq!(argv, vec!["interface", "upstream", "list"]); // 5. Prohibited mutating commands rejected by CLI allowlist let import_cli_req = ExecuteCliRequest { command: "interface".to_string(), subcommand: Some("upstream".to_string()), sub_subcommand: Some("import".to_string()), target: Some("proton0".to_string()), parameters: HashMap::new(), }; assert!(build_safe_argv(&import_cli_req).is_err()); } #[tokio::test] async fn test_upstream_without_listen_port_does_not_conflict_with_wg0() { let harness = setup_test_harness().await; // 1. Verify wg0 already owns local UDP 51820 let wg0_initial = harness .wg_engine .get_interface_stats("wg0") .await .unwrap() .unwrap(); assert_eq!(wg0_initial.listen_port, 51820); // 2. Import proton0 from a configuration with no ListenPort let (proton_priv, _) = generate_keypair(); let (_, provider_pub) = generate_keypair(); let conf = format!( r#" [Interface] PrivateKey = {} Address = 10.2.0.2/32 DNS = 10.2.0.1 [Peer] PublicKey = {} AllowedIPs = 0.0.0.0/0, ::/0 Endpoint = 37.19.199.155:51820 PersistentKeepalive = 25 "#, proton_priv.as_str(), provider_pub.as_str() ); let import_req = Request::builder() .uri("/api/v1/interfaces/upstreams/import") .method("POST") .header(header::COOKIE, &harness.session_cookie) .header(header::CONTENT_TYPE, "application/json") .body(Body::from( json!({ "name": "proton0", "config": conf }) .to_string(), )) .unwrap(); let resp = harness.app.clone().oneshot(import_req).await.unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body_bytes = to_bytes(resp.into_body(), usize::MAX).await.unwrap(); let import_res: Value = serde_json::from_slice(&body_bytes).unwrap(); assert_eq!(import_res["name"], "proton0"); assert_eq!(import_res["role"], "upstream"); assert_eq!(import_res["listen_port"], Value::Null); assert_eq!(import_res["provider_endpoint"], "37.19.199.155:51820"); assert_eq!(import_res["provider_allowed_ips"], "0.0.0.0/0, ::/0"); // 3. Verify wg0 remains on UDP 51820 and unchanged let wg0_db = harness .store .get_interface_by_name("wg0") .await .unwrap() .unwrap(); assert_eq!(wg0_db.listen_port, Some(51820)); assert_eq!(wg0_db.role, InterfaceRole::Overlay); // 4. Verify proton0 desired state in DB has listen_port = None let proton_db = harness .store .get_interface_by_name("proton0") .await .unwrap() .unwrap(); assert_eq!(proton_db.listen_port, None); assert_eq!(proton_db.role, InterfaceRole::Upstream); // 5. Verify simulated kernel state has both wg0 (51820) and proton0 (dynamic/0) let live_wg0 = harness .wg_engine .get_interface_stats("wg0") .await .unwrap() .unwrap(); assert_eq!(live_wg0.listen_port, 51820); let live_proton = harness .wg_engine .get_interface_stats("proton0") .await .unwrap() .unwrap(); assert_eq!(live_proton.listen_port, 0); assert_eq!(live_proton.peers.len(), 1); assert_eq!( live_proton.peers[0].allowed_ips, vec!["0.0.0.0/0".to_string(), "::/0".to_string()] ); } #[tokio::test] async fn test_explicit_upstream_listen_port_is_preserved() { let harness = setup_test_harness().await; let (proton_priv, _) = generate_keypair(); let (_, provider_pub) = generate_keypair(); let conf = format!( r#" [Interface] PrivateKey = {} Address = 10.2.0.2/32 ListenPort = 45000 [Peer] PublicKey = {} AllowedIPs = 0.0.0.0/0, ::/0 Endpoint = 37.19.199.155:51820 "#, proton_priv.as_str(), provider_pub.as_str() ); let import_req = Request::builder() .uri("/api/v1/interfaces/upstreams/import") .method("POST") .header(header::COOKIE, &harness.session_cookie) .header(header::CONTENT_TYPE, "application/json") .body(Body::from( json!({ "name": "custom_vpn0", "config": conf }) .to_string(), )) .unwrap(); let resp = harness.app.clone().oneshot(import_req).await.unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body_bytes = to_bytes(resp.into_body(), usize::MAX).await.unwrap(); let import_res: Value = serde_json::from_slice(&body_bytes).unwrap(); assert_eq!(import_res["listen_port"], 45000); let iface_db = harness .store .get_interface_by_name("custom_vpn0") .await .unwrap() .unwrap(); assert_eq!(iface_db.listen_port, Some(45000)); let live_custom = harness .wg_engine .get_interface_stats("custom_vpn0") .await .unwrap() .unwrap(); assert_eq!(live_custom.listen_port, 45000); } #[tokio::test] async fn test_upstream_missing_listen_port_no_false_drift() { let harness = setup_test_harness().await; // 1. Create upstream interface proton0 in DB with listen_port = None let (priv_k, pub_k) = generate_keypair(); let (_, peer_pub) = generate_keypair(); let iface_id = Uuid::new_v4(); let iface = Interface { id: iface_id, name: "proton0".to_string(), role: InterfaceRole::Upstream, private_key: priv_k, public_key: pub_k.clone(), listen_port: None, address_v4: validate_cidr("10.2.0.2/32").unwrap(), address_v6: None, mtu: Some(1420), dns: None, enabled: true, pre_up: None, post_up: None, pre_down: None, post_down: None, created_at: Utc::now().naive_utc(), updated_at: Utc::now().naive_utc(), }; harness.store.create_interface(&iface).await.unwrap(); let peer = Peer { id: Uuid::new_v4(), interface_id: iface_id, name: "proton0-provider".to_string(), peer_type: PeerType::Server, state: PeerState::Active, public_key: peer_pub.clone(), private_key: None, preshared_key: None, endpoint: Some("37.19.199.155:51820".to_string()), allowed_ips: "0.0.0.0/0, ::/0".to_string(), server_allowed_ips: Some("0.0.0.0/0, ::/0".to_string()), address_v4: None, address_v6: None, dns: None, mtu: Some(1420), persistent_keepalive: Some(25), profile: PeerProfile::Custom, expires_at: None, last_handshake_at: None, created_at: Utc::now().naive_utc(), updated_at: Utc::now().naive_utc(), }; harness.store.create_peer(&peer).await.unwrap(); // 2. Inject live kernel stats where the kernel has allocated an ephemeral dynamic port 54321 harness .wg_engine .inject_interface_stats(LiveInterfaceStats { name: "proton0".to_string(), public_key: pub_k.as_str().to_string(), listen_port: 54321, // dynamic kernel-allocated port fwmark: 0, peers: vec![nx9_wireguard::LivePeerStats { public_key: peer_pub.as_str().to_string(), endpoint: Some("37.19.199.155:51820".to_string()), rx_bytes: 100, tx_bytes: 200, last_handshake_at: None, allowed_ips: vec!["0.0.0.0/0".to_string(), "::/0".to_string()], persistent_keepalive: Some(25), }], addresses: vec!["10.2.0.2/32".to_string()], mtu: Some(1420), is_up: true, }) .await; // 3. Run reconciliation plan — must NOT flag drift for the dynamic listen port let plan = harness.reconciler.plan().await.unwrap(); assert!( !plan.has_drift, "Expected zero drift for dynamic kernel listen port when desired listen_port is None, but got: {:?}", plan.actions ); assert_eq!(plan.interface_changes, 0); assert_eq!(plan.peer_changes, 0); }