//! Comprehensive Integration and Drift Matrix test suite for ReconciliationEngine. //! Covers: //! - WireGuard interface & peer drift (CREATE, UPDATE, DELETE, NOOP) //! - Route, Firewall, NAT, and Forwarding drift detection //! - Plan dry-run read-only determinism and idempotency //! - Concurrent apply serialization //! - Restart recovery //! - Secret safety across plans and reports use chrono::Utc; use ipnet::IpNet; use nx9_wg_api::reconciliation::ReconciliationEngine; use nx9_wg_api::state::AppState; use nx9_wg_core::crypto::generate_keypair; use nx9_wg_core::types::firewall::{ FirewallAction, FirewallDirection, FirewallProtocol, FirewallRule, }; use nx9_wg_core::types::network::Route; 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::{NetworkEngine, SimulatedNetworkEngine}; use nx9_wireguard::{SimulatedWireGuardEngine, WireGuardEngine}; use std::sync::Arc; use tempfile::{TempDir, tempdir}; use uuid::Uuid; async fn setup_test_env() -> ( TempDir, Store, AppState, Arc, Arc, ReconciliationEngine, ) { let dir = tempdir().expect("create temp dir"); let db_path = dir.path().join("drift_test.db"); let store = Store::connect(&db_path.to_string_lossy()) .await .expect("connect to db"); store.migrate().await.expect("run migrations"); let state = AppState::new(store.clone()); let wg_engine = Arc::new(SimulatedWireGuardEngine::new()); let net_engine = Arc::new(SimulatedNetworkEngine::new()); let reconciler = ReconciliationEngine::new(state.clone(), wg_engine.clone(), net_engine.clone()); (dir, store, state, wg_engine, net_engine, reconciler) } #[tokio::test] async fn test_drift_matrix_peer_lifecycle() { let (_dir, store, _state, wg_engine, _net_engine, reconciler) = setup_test_env().await; // 1. Create interface & active peer in SQLite let (priv_key, pub_key) = generate_keypair(); let iface_id = Uuid::new_v4(); let iface = Interface { id: iface_id, name: "nx9_test0".to_string(), role: InterfaceRole::Overlay, private_key: priv_key, public_key: pub_key, listen_port: Some(51820), address_v4: validate_cidr("10.10.0.1/24").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(), }; store.create_interface(&iface).await.unwrap(); let (_p_priv, p_pub) = generate_keypair(); let peer_id = Uuid::new_v4(); let peer = Peer { id: peer_id, interface_id: iface_id, name: "peer-alice".to_string(), public_key: p_pub.clone(), private_key: None, preshared_key: None, address_v4: Some(validate_cidr("10.10.0.2/32").unwrap()), address_v6: None, allowed_ips: "10.10.0.2/32".to_string(), server_allowed_ips: None, endpoint: Some("203.0.113.5:51820".to_string()), persistent_keepalive: Some(25), dns: None, mtu: None, profile: PeerProfile::FullTunnel, state: PeerState::Active, peer_type: PeerType::RoadWarrior, expires_at: None, last_handshake_at: None, created_at: Utc::now().naive_utc(), updated_at: Utc::now().naive_utc(), }; store.create_peer(&peer).await.unwrap(); // 2. Plan: Detect interface missing & peer missing let plan = reconciler.plan().await.unwrap(); assert!(plan.has_drift); assert_eq!(plan.interface_changes, 1); assert_eq!(plan.peer_changes, 1); // 3. Apply: Converges state to kernel let report = reconciler.apply().await.unwrap(); assert!(report.success); // 4. Verify live stats let stats = wg_engine .get_interface_stats("nx9_test0") .await .unwrap() .unwrap(); assert_eq!(stats.peers.len(), 1); assert_eq!(stats.peers[0].public_key, p_pub.as_str()); // 5. Post-apply verify: zero drift let plan2 = reconciler.verify().await.unwrap(); assert_eq!(plan2.interface_changes, 0); assert_eq!(plan2.peer_changes, 0); // 6. Drift injection: Mark peer Expired in SQLite store.mark_peer_expired(peer_id).await.unwrap(); // Plan should detect active peer in kernel is no longer active in DB -> remove_inactive_peer let plan3 = reconciler.plan().await.unwrap(); assert!(plan3.has_drift); assert_eq!(plan3.peer_changes, 1); assert!( plan3 .actions .iter() .any(|a| a.action_type == "remove_inactive_peer") ); // Apply removal let report2 = reconciler.apply().await.unwrap(); assert!(report2.success); // Live interface now has 0 peers let stats2 = wg_engine .get_interface_stats("nx9_test0") .await .unwrap() .unwrap(); assert_eq!(stats2.peers.len(), 0); } #[tokio::test] async fn test_drift_matrix_routes_and_firewall() { let (_dir, store, _state, _wg_engine, _net_engine, reconciler) = setup_test_env().await; // 1. Add route in SQLite let route = Route { id: Uuid::new_v4(), network_id: None, interface_id: None, destination: "192.168.50.0/24".parse::().unwrap(), gateway: Some("10.10.0.1".parse().unwrap()), interface_name: Some("nx9_test0".to_string()), metric: Some(100), description: Some("Test route".to_string()), enabled: true, created_at: Utc::now().naive_utc(), updated_at: Utc::now().naive_utc(), }; store.create_route(&route).await.unwrap(); // 2. Add firewall rule in SQLite let fw = FirewallRule { id: Uuid::new_v4(), name: "allow-http".to_string(), interface_id: None, peer_id: None, direction: FirewallDirection::In, source: None, destination: None, protocol: FirewallProtocol::Tcp, source_port: None, destination_port: Some(80), port_range: None, action: FirewallAction::Accept, priority: 100, enabled: true, description: None, created_at: Utc::now().naive_utc(), updated_at: Utc::now().naive_utc(), }; store.create_firewall_rule(&fw).await.unwrap(); // 3. Plan should detect route changes and firewall changes let plan = reconciler.plan().await.unwrap(); assert!(plan.has_drift); assert_eq!(plan.route_changes, 1); assert_eq!(plan.firewall_changes, 1); // 4. Apply let report = reconciler.apply().await.unwrap(); assert!(report.success); assert!(report.executed_actions >= 2); } #[tokio::test] async fn test_reconciliation_dry_run_idempotency_and_read_only() { let (_dir, store, _state, _wg_engine, _net_engine, reconciler) = setup_test_env().await; let audit_count_before = store .list_audit_events(&nx9_wg_db::AuditFilter::default(), 100, 0) .await .unwrap() .len(); // Run plan multiple times let plan1 = reconciler.plan().await.unwrap(); let plan2 = reconciler.plan().await.unwrap(); let plan3 = reconciler.verify().await.unwrap(); assert_eq!(plan1.has_drift, plan2.has_drift); assert_eq!(plan1.actions.len(), plan2.actions.len()); assert_eq!(plan1.actions.len(), plan3.actions.len()); // Audit logs must not increase during plan/verify dry-runs let audit_count_after = store .list_audit_events(&nx9_wg_db::AuditFilter::default(), 100, 0) .await .unwrap() .len(); assert_eq!(audit_count_before, audit_count_after); } #[tokio::test] async fn test_reconciliation_concurrent_apply_serialization() { let (_dir, _store, _state, _wg_engine, _net_engine, reconciler) = setup_test_env().await; let reconciler_arc = Arc::new(reconciler); let mut handles = Vec::new(); for _ in 0..5 { let r = Arc::clone(&reconciler_arc); handles.push(tokio::spawn(async move { r.apply().await })); } for handle in handles { let res = handle.await.unwrap(); assert!(res.is_ok()); } } #[tokio::test] async fn test_restart_recovery_simulation() { let dir = tempdir().expect("create temp dir"); let db_path = dir.path().join("restart_test.db"); let store = Store::connect(&db_path.to_string_lossy()) .await .expect("connect to db"); store.migrate().await.expect("run migrations"); // 1. Initial run with interface let (priv_key, pub_key) = generate_keypair(); let iface = Interface { id: Uuid::new_v4(), name: "nx9_boot".to_string(), role: InterfaceRole::Overlay, private_key: priv_key, public_key: pub_key, listen_port: Some(51820), address_v4: validate_cidr("10.20.0.1/24").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(), }; store.create_interface(&iface).await.unwrap(); let wg1 = Arc::new(SimulatedWireGuardEngine::new()); let net1 = Arc::new(SimulatedNetworkEngine::new()); let r1 = ReconciliationEngine::new(AppState::new(store.clone()), wg1.clone(), net1.clone()); r1.apply().await.unwrap(); assert!(wg1.get_interface_stats("nx9_boot").await.unwrap().is_some()); // 2. Simulate machine reboot / app restart: // Create new live engine instance (empty kernel state), but reconnect same store let wg2 = Arc::new(SimulatedWireGuardEngine::new()); let net2 = Arc::new(SimulatedNetworkEngine::new()); let r2 = ReconciliationEngine::new(AppState::new(store.clone()), wg2.clone(), net2.clone()); // Before reconcile, new engine is empty assert!(wg2.get_interface_stats("nx9_boot").await.unwrap().is_none()); // Compute plan: detects missing interface let plan = r2.plan().await.unwrap(); assert!(plan.has_drift); assert_eq!(plan.interface_changes, 1); // Apply reconciliation r2.apply().await.unwrap(); // Kernel converged assert!(wg2.get_interface_stats("nx9_boot").await.unwrap().is_some()); } #[tokio::test] async fn test_secret_redaction_in_reconciliation_plan_and_report() { let (_dir, store, _state, _wg_engine, _net_engine, reconciler) = setup_test_env().await; let (priv_key, pub_key) = generate_keypair(); let raw_secret = priv_key.as_str().to_string(); let iface = Interface { id: Uuid::new_v4(), name: "nx9_sec".to_string(), role: InterfaceRole::Overlay, private_key: priv_key, public_key: pub_key, listen_port: Some(51820), address_v4: validate_cidr("10.30.0.1/24").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(), }; store.create_interface(&iface).await.unwrap(); let plan = reconciler.plan().await.unwrap(); let plan_json = serde_json::to_string(&plan).unwrap(); assert!( !plan_json.contains(&raw_secret), "Private key must NOT leak into plan JSON" ); let report = reconciler.apply().await.unwrap(); let report_json = serde_json::to_string(&report).unwrap(); assert!( !report_json.contains(&raw_secret), "Private key must NOT leak into report JSON" ); } #[tokio::test] async fn test_reconciliation_status_lifecycle_and_multi_cycle_idempotency() { use nx9_wg_api::reconciliation::ReconciliationStatus; let (_dir, store, _state, _wg_engine, _net_engine, reconciler) = setup_test_env().await; let (priv_key, pub_key) = generate_keypair(); let iface = Interface { id: Uuid::new_v4(), name: "nx9_idem".to_string(), role: InterfaceRole::Overlay, private_key: priv_key, public_key: pub_key, listen_port: Some(51820), address_v4: validate_cidr("10.50.0.1/24").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(), }; store.create_interface(&iface).await.unwrap(); // 1. First apply converges let report1 = reconciler.apply().await.unwrap(); assert!(report1.success); assert_eq!(report1.status, ReconciliationStatus::Converged); // 2. Run 5 consecutive apply cycles: all must succeed with Converged status for cycle in 2..=6 { let report = reconciler.apply().await.unwrap(); assert!(report.success, "Cycle {cycle} must succeed"); assert_eq!( report.status, ReconciliationStatus::Converged, "Cycle {cycle} must report Converged" ); let plan = reconciler.plan().await.unwrap(); assert!(!plan.has_drift, "Cycle {cycle} plan must show zero drift"); } } #[tokio::test] async fn test_reconciliation_report_schema_and_json_contract() { use nx9_wg_api::reconciliation::{ReconciliationReport, ReconciliationStatus}; let report = ReconciliationReport { success: true, status: ReconciliationStatus::Converged, executed_actions: 3, failed_actions: 0, details: vec![ "Synchronized interface 'wg0' with 5 peers".to_string(), "Synchronized 1 routing entries".to_string(), "Synchronized 0 firewall rules into table inet nx9_wg (NAT: true)".to_string(), ], }; let json_val = serde_json::to_value(&report).unwrap(); assert_eq!(json_val["success"], true); assert_eq!(json_val["status"], "converged"); assert_eq!(json_val["executed_actions"], 3); assert_eq!(json_val["failed_actions"], 0); assert!(json_val["details"].is_array()); assert_eq!(json_val["details"].as_array().unwrap().len(), 3); } #[tokio::test] async fn test_reconciliation_nftables_canonical_drift_and_kernel_handle_tolerance() { let (_dir, store, _state, _wg_engine, net_engine, reconciler) = setup_test_env().await; // Add firewall rule in SQLite let fw = FirewallRule { id: Uuid::new_v4(), name: "allow-https".to_string(), interface_id: None, peer_id: None, direction: FirewallDirection::In, source: None, destination: None, protocol: FirewallProtocol::Tcp, source_port: None, destination_port: Some(443), port_range: None, action: FirewallAction::Accept, priority: 50, enabled: true, description: None, created_at: Utc::now().naive_utc(), updated_at: Utc::now().naive_utc(), }; store.create_firewall_rule(&fw).await.unwrap(); // 1. Initial Plan should detect drift let plan = reconciler.plan().await.unwrap(); assert!(plan.has_drift); assert_eq!(plan.firewall_changes, 1); // 2. Apply should converge let report = reconciler.apply().await.unwrap(); assert!(report.success); assert_eq!( report.status, nx9_wg_api::reconciliation::ReconciliationStatus::Converged ); assert_eq!(report.failed_actions, 0); // 3. Post-apply verify: exactly 0 drift let plan_after = reconciler.plan().await.unwrap(); assert!(!plan_after.has_drift); assert_eq!(plan_after.firewall_changes, 0); // 4. Simulate kernel returning ruleset with handles and tabs let simulated_kernel_output_with_handles = r#"table inet nx9_wg { chain input { type filter hook input priority filter; policy accept; ct state established,related accept # handle 46 iifname "lo" accept # handle 1 tcp dport 443 accept # handle 10 } chain forward { type filter hook forward priority filter; policy accept; ct state established,related accept # handle 4 } chain postrouting { type nat hook postrouting priority srcnat; policy accept; } } "#; // Set simulated ruleset to text containing kernel handles net_engine .sync_firewall(std::slice::from_ref(&fw), false, &[]) .await .unwrap(); // Directly test drift function against simulated kernel handles let expected = nx9_wg_network::NftablesRulesetBuilder::build(&[fw], false, &[]); assert!( !nx9_wg_network::has_nftables_drift(&expected, simulated_kernel_output_with_handles), "Ruleset with handles must not trigger false drift" ); } #[tokio::test] async fn test_interface_address_and_mtu_drift_lifecycle() { let (_dir, store, _state, wg_engine, _net_engine, reconciler) = setup_test_env().await; let (priv_key, pub_key) = generate_keypair(); let iface_id = Uuid::new_v4(); let iface = Interface { id: iface_id, name: "wg0".to_string(), role: InterfaceRole::Overlay, private_key: priv_key, public_key: pub_key.clone(), listen_port: Some(51820), address_v4: validate_cidr("10.100.0.1/24").unwrap(), address_v6: Some(validate_cidr("fd00::1/64").unwrap()), 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(), }; store.create_interface(&iface).await.unwrap(); // 1. Initially, interface does not exist in wg_engine -> plan reports create_interface drift let plan = reconciler.plan().await.unwrap(); assert!(plan.has_drift); assert_eq!(plan.interface_changes, 1); assert_eq!(plan.actions[0].action_type, "create_interface"); // 2. Apply initial sync -> interface is created and synchronized let report = reconciler.apply().await.unwrap(); assert!(report.success); assert_eq!( report.status, nx9_wg_api::reconciliation::ReconciliationStatus::Converged ); // 3. Post-apply plan must have 0 drift let plan_after = reconciler.plan().await.unwrap(); assert!(!plan_after.has_drift); assert_eq!(plan_after.interface_changes, 0); // 4. Manually strip IPv4 address from live interface to simulate kernel address drop let mut stats = wg_engine.get_interface_stats("wg0").await.unwrap().unwrap(); stats.addresses = vec!["fd00::1/64".to_string()]; // IPv4 missing // Sync altered stats wg_engine .sync_interface( &Interface { address_v4: validate_cidr("10.99.99.99/24").unwrap(), // different ..iface.clone() }, &[], ) .await .unwrap(); // 5. Plan MUST detect the missing/mismatched IPv4 address as drift let plan_drift = reconciler.plan().await.unwrap(); assert!(plan_drift.has_drift); assert_eq!(plan_drift.interface_changes, 1); assert_eq!(plan_drift.actions[0].action_type, "update_interface"); assert!(plan_drift.actions[0].description.contains("IPv4 address")); // 6. Apply reconciliation -> restores correct addresses let report2 = reconciler.apply().await.unwrap(); assert!(report2.success); assert_eq!( report2.status, nx9_wg_api::reconciliation::ReconciliationStatus::Converged ); // 7. Final plan reports 0 drift let final_plan = reconciler.plan().await.unwrap(); assert!(!final_plan.has_drift); assert_eq!(final_plan.interface_changes, 0); }