Files
nx9-wg/crates/nx9-wg-api/tests/test_reconciliation_drift.rs
2026-09-02 15:19:19 +05:30

593 lines
20 KiB
Rust

//! 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<SimulatedWireGuardEngine>,
Arc<SimulatedNetworkEngine>,
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::<IpNet>().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);
}