release: NX9-WG v1.0.0

This commit is contained in:
thakares committed 2026-08-18 17:32:56 +05:30
1 parent c8a9b7cde6
commit 4dfe42fe68
42 files changed
+4689 -336

No files matched your search

+193 -34
View File
@@ -3,6 +3,7 @@
use crate::error::{ApiError, ApiResult};
use crate::state::{AppState, SystemEvent};
use chrono::Utc;
use ipnet::IpNet;
use nx9_wg_core::types::audit::AuditEventType;
use nx9_wg_core::types::wireguard::PeerState;
use nx9_wg_network::NetworkEngine;
@@ -11,6 +12,36 @@ use serde::{Deserialize, Serialize};
use std::sync::Arc;
use std::time::Duration;
/// Check if a slice of live address strings contains the desired IpNet.
fn matches_ipnet(live_addrs: &[String], desired: &IpNet) -> bool {
live_addrs.iter().any(|s| {
if let Ok(net) = s.parse::<IpNet>() {
net.addr() == desired.addr() && net.prefix_len() == desired.prefix_len()
} else {
false
}
})
}
/// Check if live WireGuard peer allowed IPs match desired server-side allowed IPs.
fn matches_allowed_ips(live_allowed_ips: &[String], desired_str: &str) -> bool {
let desired_nets: std::collections::BTreeSet<IpNet> = desired_str
.split(',')
.map(|s| s.trim())
.filter(|s| !s.is_empty())
.filter_map(|s| s.parse::<IpNet>().ok())
.collect();
let live_nets: std::collections::BTreeSet<IpNet> = live_allowed_ips
.iter()
.map(|s| s.trim())
.filter(|s| !s.is_empty())
.filter_map(|s| s.parse::<IpNet>().ok())
.collect();
desired_nets == live_nets
}
/// Individual action proposed or taken by the reconciler.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ReconciliationAction {
@@ -53,6 +84,8 @@ pub struct ReconciliationReport {
#[serde(default)]
pub status: ReconciliationStatus,
pub executed_actions: usize,
#[serde(default)]
pub failed_actions: usize,
pub details: Vec<String>,
}
@@ -132,40 +165,66 @@ impl ReconciliationEngine {
.ok()
.flatten();
let live_peer_keys: Vec<String> = live_stats
.as_ref()
.map(|s| s.peers.iter().map(|p| p.public_key.clone()).collect())
.unwrap_or_default();
let iface_exists = live_interfaces.contains(&iface.name) || live_stats.is_some();
match live_stats.as_ref() {
Some(stats) => {
if stats.public_key != iface.public_key.as_str()
|| stats.listen_port != iface.listen_port
if iface_exists {
if let Some(stats) = live_stats.as_ref() {
let mut drift_reasons = Vec::new();
if !stats.public_key.is_empty()
&& stats.public_key != iface.public_key.as_str()
{
drift_reasons.push("public key mismatch".to_string());
}
if stats.listen_port != 0 && stats.listen_port != iface.listen_port {
drift_reasons.push("listen port mismatch".to_string());
}
if !matches_ipnet(&stats.addresses, &iface.address_v4) {
drift_reasons
.push(format!("missing IPv4 address '{}'", iface.address_v4));
}
if let Some(ref v6) = iface.address_v6
&& !matches_ipnet(&stats.addresses, v6)
{
drift_reasons.push(format!("missing IPv6 address '{v6}'"));
}
if let Some(desired_mtu) = iface.mtu
&& let Some(live_mtu) = stats.mtu
&& live_mtu != desired_mtu as u32
{
drift_reasons.push(format!(
"MTU mismatch (live: {live_mtu}, desired: {desired_mtu})"
));
}
if !stats.is_up {
drift_reasons.push("interface link is down".to_string());
}
if !drift_reasons.is_empty() {
plan.actions.push(ReconciliationAction {
subsystem: "wireguard".to_string(),
resource_id: iface.id.to_string(),
action_type: "update_interface".to_string(),
description: format!(
"Interface '{}' configuration drift detected; update listen port / keys",
iface.name
"Interface '{}' configuration drift detected ({}); synchronize link, address, port, or keys",
iface.name,
drift_reasons.join(", ")
),
});
plan.interface_changes += 1;
}
}
None => {
plan.actions.push(ReconciliationAction {
subsystem: "wireguard".to_string(),
resource_id: iface.id.to_string(),
action_type: "create_interface".to_string(),
description: format!(
"Interface '{}' missing in kernel; create and sync",
iface.name
),
});
plan.interface_changes += 1;
}
} else {
plan.actions.push(ReconciliationAction {
subsystem: "wireguard".to_string(),
resource_id: iface.id.to_string(),
action_type: "create_interface".to_string(),
description: format!(
"Interface '{}' missing in kernel; create and sync",
iface.name
),
});
plan.interface_changes += 1;
}
// Check peers (only Active desired peers should be live)
@@ -175,16 +234,86 @@ impl ReconciliationEngine {
.filter(|p| p.state == PeerState::Active)
.collect();
let live_peers_map: std::collections::HashMap<
String,
&nx9_wireguard::LivePeerStats,
> = if let Some(ref stats) = live_stats {
stats
.peers
.iter()
.map(|p| (p.public_key.clone(), p))
.collect()
} else {
std::collections::HashMap::new()
};
for p in &active_desired_peers {
if !live_peer_keys.contains(&p.public_key.as_str().to_string()) {
let pub_key_str = p.public_key.as_str();
let desired_server_allowed = p.server_wireguard_allowed_ips();
if let Some(live_p) = live_peers_map.get(pub_key_str) {
// Peer is present in live kernel interface. Verify semantic drift:
let mut peer_drifts = Vec::new();
if !matches_allowed_ips(&live_p.allowed_ips, &desired_server_allowed) {
peer_drifts.push(format!(
"AllowedIPs drift (live: [{:?}], desired: [{desired_server_allowed}])",
live_p.allowed_ips
));
}
if let (Some(desired_ka), Some(live_ka)) =
(p.persistent_keepalive, live_p.persistent_keepalive)
&& live_ka != desired_ka
{
peer_drifts.push(format!(
"persistent keepalive drift (live: {live_ka}s, desired: {desired_ka}s)"
));
}
if !peer_drifts.is_empty() {
plan.actions.push(ReconciliationAction {
subsystem: "wireguard".to_string(),
resource_id: p.id.to_string(),
action_type: "update_peer".to_string(),
description: format!(
"Peer '{}' ({}) drift detected: {}; re-sync in kernel",
p.name,
pub_key_str,
peer_drifts.join(", ")
),
});
plan.peer_changes += 1;
}
// Update operational telemetry (handshake timestamp and learned endpoint) from kernel
if live_p.last_handshake_at.is_some() || live_p.endpoint.is_some() {
let hs_newer = live_p.last_handshake_at.is_some()
&& live_p.last_handshake_at != p.last_handshake_at;
let ep_newer = live_p.endpoint.is_some()
&& live_p.endpoint.as_deref() != p.endpoint.as_deref();
if hs_newer || ep_newer {
let _ = self
.state
.store
.update_peer_learned_telemetry(
p.id,
live_p.last_handshake_at.or(p.last_handshake_at),
live_p.endpoint.as_deref().or(p.endpoint.as_deref()),
)
.await;
}
}
} else {
plan.actions.push(ReconciliationAction {
subsystem: "wireguard".to_string(),
resource_id: p.id.to_string(),
action_type: "add_peer".to_string(),
description: format!(
"Peer '{}' ({}) missing in live interface",
"Peer '{}' ({}) missing in live interface; add to kernel with AllowedIPs [{desired_server_allowed}]",
p.name,
p.public_key.as_str()
pub_key_str
),
});
plan.peer_changes += 1;
@@ -293,7 +422,7 @@ impl ReconciliationEngine {
.await
.unwrap_or_default();
if expected_ruleset.trim() != active_ruleset.trim() {
if nx9_wg_network::has_nftables_drift(&expected_ruleset, &active_ruleset) {
plan.actions.push(ReconciliationAction {
subsystem: "firewall".to_string(),
resource_id: "nftables".to_string(),
@@ -340,6 +469,17 @@ impl ReconciliationEngine {
// Sweep expired peers
let _ = self.sweep_expired_peers().await;
let initial_plan = self.plan().await.unwrap_or_default();
if !initial_plan.has_drift {
return Ok(ReconciliationReport {
success: true,
status: ReconciliationStatus::Converged,
executed_actions: 0,
failed_actions: 0,
details: vec!["System is already fully converged; zero drift detected".to_string()],
});
}
let desired_interfaces = self.state.store.list_interfaces().await?;
let mut details = Vec::new();
@@ -419,10 +559,28 @@ impl ReconciliationEngine {
// 4. Verify post-apply convergence
let post_plan = self.plan().await.unwrap_or_default();
let (success, status) = if !post_plan.has_drift {
(true, ReconciliationStatus::Converged)
let (success, status, executed_actions, failed_actions) = if !post_plan.has_drift {
(
true,
ReconciliationStatus::Converged,
initial_plan.actions.len(),
0,
)
} else {
(false, ReconciliationStatus::DriftRemains)
let remaining = post_plan.actions.len();
let completed = initial_plan.actions.len().saturating_sub(remaining);
for action in &post_plan.actions {
details.push(format!(
"Unresolved drift: [{}] {}",
action.subsystem, action.description
));
}
(
false,
ReconciliationStatus::DriftRemains,
completed,
remaining,
)
};
// 5. Audit reconciliation run
@@ -435,8 +593,8 @@ impl ReconciliationEngine {
Some("reconciliation"),
None,
Some(&format!(
"Reconciliation applied {} actions (status: {status:?})",
details.len()
"Reconciliation applied {} actions (status: {status:?}, failed: {failed_actions})",
executed_actions
)),
None,
None,
@@ -446,8 +604,8 @@ impl ReconciliationEngine {
self.state.broadcast(SystemEvent::AuditEvent {
event_type: AuditEventType::ReconciliationRun,
message: Some(format!(
"Reconciliation applied {} actions (status: {status:?})",
details.len()
"Reconciliation applied {} actions (status: {status:?}, failed: {failed_actions})",
executed_actions
)),
resource_type: Some("reconciliation".to_string()),
resource_id: None,
@@ -456,7 +614,8 @@ impl ReconciliationEngine {
Ok(ReconciliationReport {
success,
status,
executed_actions: details.len(),
executed_actions,
failed_actions,
details,
})
}
File diff suppressed because it is too large. Load diff
+25 -1
View File
@@ -10,7 +10,31 @@
</style>
</head>
<body>
<div id="app-layout">
<!-- Unauthenticated Login View -->
<div id="login-view" style="display: none;">
<div class="login-card">
<div class="login-header">
<span class="brand-mark">NX9</span>
<h2>Administrator Login</h2>
<p>Sign in to nx9-wg Native Linux Appliance</p>
</div>
<div id="login-error-msg" class="login-error" style="display: none;"></div>
<form id="login-form" onsubmit="event.preventDefault(); submitLogin();">
<div class="form-group">
<label class="form-label" for="login-username">Username</label>
<input type="text" id="login-username" class="form-input" value="admin" required autocomplete="username">
</div>
<div class="form-group">
<label class="form-label" for="login-password">Password</label>
<input type="password" id="login-password" class="form-input" required autofocus autocomplete="current-password">
</div>
<button type="submit" id="login-submit-btn" class="btn btn-primary" style="width: 100%; margin-top: 8px;">Sign In</button>
</form>
</div>
</div>
<!-- Authenticated Application Shell -->
<div id="app-layout" style="display: none;">
<!-- Top Application Bar -->
<header class="topbar">
<div class="topbar-left">
+35 -11
View File
@@ -3,9 +3,9 @@
use crate::auth::middleware::AuthenticatedAdmin;
use crate::error::{ApiError, ApiResult};
use crate::state::AppState;
use axum::extract::{Path, State};
use axum::extract::{Path, Request, State};
use axum::http::HeaderMap;
use axum::http::header::SET_COOKIE;
use axum::http::header::{AUTHORIZATION, COOKIE, SET_COOKIE};
use axum::response::{IntoResponse, Response};
use axum::{Extension, Json};
use chrono::NaiveDateTime;
@@ -67,9 +67,8 @@ pub async fn login_handler(
.await?;
let cookie_val = format!(
"nx9_session={}; Path=/; HttpOnly; SameSite=Lax; Max-Age={}",
session.id,
24 * 3600
"nx9_session={}; Path=/; HttpOnly; SameSite=Lax; Max-Age=86400",
session.id
);
let mut headers = HeaderMap::new();
@@ -89,12 +88,37 @@ pub async fn login_handler(
}
/// POST /api/v1/auth/logout
pub async fn logout_handler(
State(state): State<AppState>,
Extension(auth_user): Extension<AuthenticatedAdmin>,
) -> ApiResult<Response> {
if let Some(ref session_id) = auth_user.session_id {
state.auth.logout(session_id, None).await?;
pub async fn logout_handler(State(state): State<AppState>, req: Request) -> ApiResult<Response> {
let mut session_to_delete = None;
// 1. Try Bearer token in Authorization header
if let Some(token) = req
.headers()
.get(AUTHORIZATION)
.and_then(|v| v.to_str().ok())
.and_then(|h| h.strip_prefix("Bearer "))
{
let token = token.trim();
if !token.starts_with("nx9_") {
session_to_delete = Some(token.to_string());
}
}
// 2. Try session cookie (nx9_session=...)
if session_to_delete.is_none()
&& let Some(cookie_header) = req.headers().get(COOKIE).and_then(|v| v.to_str().ok())
{
for cookie in cookie_header.split(';') {
let cookie = cookie.trim();
if let Some(session_id) = cookie.strip_prefix("nx9_session=") {
session_to_delete = Some(session_id.trim().to_string());
break;
}
}
}
if let Some(ref session_id) = session_to_delete {
let _ = state.auth.logout(session_id, None).await;
}
let cookie_val = "nx9_session=; Path=/; HttpOnly; SameSite=Lax; Max-Age=0";
+28 -4
View File
@@ -38,6 +38,7 @@ pub struct UpdateInterfaceRequest {
pub address_v6: Option<String>,
pub mtu: Option<u16>,
pub dns: Option<String>,
pub enabled: Option<bool>,
pub pre_up: Option<String>,
pub post_up: Option<String>,
pub pre_down: Option<String>,
@@ -144,8 +145,18 @@ pub async fn update_interface_handler(
.ok_or_else(|| ApiError::NotFound(format!("Interface '{id}' not found")))?;
if let Some(ref name) = payload.name {
validate_interface_name(name)?;
iface.name = name.clone();
let trimmed = name.trim();
validate_interface_name(trimmed)?;
if iface.name != trimmed {
if let Ok(Some(existing)) = state.store.get_interface_by_name(trimmed).await
&& existing.id != iface.id
{
return Err(ApiError::Validation(format!(
"Interface with name '{trimmed}' already exists"
)));
}
iface.name = trimmed.to_string();
}
}
if let Some(port) = payload.listen_port {
validate_listen_port(port)?;
@@ -155,14 +166,25 @@ pub async fn update_interface_handler(
iface.address_v4 = validate_cidr(v4)?;
}
if let Some(ref v6) = payload.address_v6 {
iface.address_v6 = Some(validate_cidr(v6)?);
if v6.trim().is_empty() {
iface.address_v6 = None;
} else {
iface.address_v6 = Some(validate_cidr(v6.trim())?);
}
}
if let Some(m) = payload.mtu {
validate_mtu(m)?;
iface.mtu = Some(m);
}
if let Some(ref dns) = payload.dns {
iface.dns = Some(dns.clone());
if dns.trim().is_empty() {
iface.dns = None;
} else {
iface.dns = Some(dns.trim().to_string());
}
}
if let Some(en) = payload.enabled {
iface.enabled = en;
}
if payload.pre_up.is_some() {
iface.pre_up = payload.pre_up;
@@ -177,6 +199,8 @@ pub async fn update_interface_handler(
iface.post_down = payload.post_down;
}
iface.updated_at = Utc::now().naive_utc();
state.store.update_interface(&iface).await?;
state.broadcast(SystemEvent::InterfaceChanged {
+3 -1
View File
@@ -28,7 +28,6 @@ pub fn build_api_router(state: AppState) -> Router {
// 1. Protected routes (require authenticated admin via session or token)
let protected_router = Router::new()
// Auth management
.route("/auth/logout", post(auth::logout_handler))
.route("/auth/session", get(auth::session_handler))
.route("/auth/password", post(auth::change_password_handler))
.route("/auth/tokens", post(auth::create_token_handler))
@@ -36,6 +35,7 @@ pub fn build_api_router(state: AppState) -> Router {
.route("/auth/tokens/{id}", delete(auth::revoke_token_handler))
// System
.route("/system", get(system::system_overview_handler))
.route("/system/live-state", get(system::live_state_handler))
.route("/system/settings", get(system::list_settings_handler))
.route("/system/settings", put(system::upsert_setting_handler))
// Interfaces
@@ -68,6 +68,7 @@ pub fn build_api_router(state: AppState) -> Router {
)
.route("/interfaces/{id}/peers", post(peers::create_peer_handler))
// Peers
.route("/peers", get(peers::list_peers_handler))
.route("/peers/{id}", get(peers::get_peer_handler))
.route("/peers/{id}", put(peers::update_peer_handler))
.route("/peers/{id}", delete(peers::delete_peer_handler))
@@ -191,6 +192,7 @@ pub fn build_api_router(state: AppState) -> Router {
// 2. Public API routes (no authentication required)
let public_router = Router::new()
.route("/auth/login", post(auth::login_handler))
.route("/auth/logout", post(auth::logout_handler))
.route("/system/health", get(system::health_handler))
.route("/system/version", get(system::version_handler))
.route("/ws", get(ws::ws_handler));
+254 -16
View File
@@ -8,6 +8,7 @@ use axum::Json;
use axum::extract::{Path, Query, State};
use axum::response::{IntoResponse, Response};
use chrono::{NaiveDateTime, Utc};
use ipnet::IpNet;
use nx9_wg_core::crypto::{generate_keypair, generate_preshared_key};
use nx9_wg_core::types::network::Network;
use nx9_wg_core::types::wireguard::{
@@ -67,13 +68,200 @@ pub struct PeerLifecycleResponse {
pub updated_at: NaiveDateTime,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct PeerResponse {
pub id: Uuid,
pub interface_id: Uuid,
pub name: String,
pub peer_type: PeerType,
pub state: PeerState,
pub public_key: WireGuardPublicKey,
#[serde(skip_serializing_if = "Option::is_none")]
pub private_key: Option<WireGuardPrivateKey>,
#[serde(skip_serializing_if = "Option::is_none")]
pub preshared_key: Option<WireGuardPresharedKey>,
pub endpoint: Option<String>,
pub allowed_ips: String,
pub server_allowed_ips: Option<String>,
pub address_v4: Option<IpNet>,
pub address_v6: Option<IpNet>,
pub dns: Option<String>,
pub mtu: Option<u16>,
pub persistent_keepalive: Option<u16>,
pub profile: PeerProfile,
pub expires_at: Option<String>,
pub last_handshake_at: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub rx_bytes: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub tx_bytes: Option<u64>,
pub created_at: String,
pub updated_at: String,
}
fn format_utc_rfc3339(dt: NaiveDateTime) -> String {
let utc_dt = chrono::DateTime::<Utc>::from_naive_utc_and_offset(dt, Utc);
utc_dt.to_rfc3339()
}
fn to_peer_response(peer: Peer, live_stats: Option<&nx9_wireguard::LivePeerStats>) -> PeerResponse {
let (endpoint, last_handshake_at, rx_bytes, tx_bytes) = if let Some(live) = live_stats {
let ep = live.endpoint.clone().or(peer.endpoint.clone());
let hs = live.last_handshake_at.or(peer.last_handshake_at);
(ep, hs, Some(live.rx_bytes), Some(live.tx_bytes))
} else {
(peer.endpoint.clone(), peer.last_handshake_at, None, None)
};
PeerResponse {
id: peer.id,
interface_id: peer.interface_id,
name: peer.name,
peer_type: peer.peer_type,
state: peer.state,
public_key: peer.public_key,
private_key: peer.private_key,
preshared_key: peer.preshared_key,
endpoint,
allowed_ips: peer.allowed_ips,
server_allowed_ips: peer.server_allowed_ips,
address_v4: peer.address_v4,
address_v6: peer.address_v6,
dns: peer.dns,
mtu: peer.mtu,
persistent_keepalive: peer.persistent_keepalive,
profile: peer.profile,
expires_at: peer.expires_at.map(format_utc_rfc3339),
last_handshake_at: last_handshake_at.map(format_utc_rfc3339),
rx_bytes,
tx_bytes,
created_at: format_utc_rfc3339(peer.created_at),
updated_at: format_utc_rfc3339(peer.updated_at),
}
}
async fn enrich_peers_with_live_telemetry(state: &AppState, peers: Vec<Peer>) -> Vec<PeerResponse> {
if peers.is_empty() {
return Vec::new();
}
// 1. Gather all unique interface IDs from peers and find interface names
let mut iface_map = std::collections::HashMap::new();
for p in &peers {
if !iface_map.contains_key(&p.interface_id)
&& let Ok(Some(iface)) = state.store.get_interface(p.interface_id).await
{
iface_map.insert(p.interface_id, iface.name);
}
}
// 2. Query live interface stats for each interface
let mut live_map = std::collections::HashMap::new();
for iface_name in iface_map.values() {
if let Ok(Some(stats)) = state.wg_engine.get_interface_stats(iface_name).await {
for lp in stats.peers {
live_map.insert(lp.public_key.clone(), lp);
}
}
}
// 3. Construct enriched PeerResponse and update DB cache if newer
let mut responses = Vec::with_capacity(peers.len());
for p in peers {
let pub_key_str = p.public_key.to_string();
let live_stat = live_map.get(&pub_key_str);
if let Some(live) = live_stat {
let hs_newer =
live.last_handshake_at.is_some() && live.last_handshake_at != p.last_handshake_at;
let ep_newer =
live.endpoint.is_some() && live.endpoint.as_deref() != p.endpoint.as_deref();
if hs_newer || ep_newer {
let latest_hs = live.last_handshake_at.or(p.last_handshake_at);
let latest_ep = live.endpoint.as_deref().or(p.endpoint.as_deref());
let _ = state
.store
.update_peer_learned_telemetry(p.id, latest_hs, latest_ep)
.await;
}
}
responses.push(to_peer_response(p, live_stat));
}
responses
}
/// GET /api/v1/peers
pub async fn list_peers_handler(
State(state): State<AppState>,
) -> ApiResult<Json<Vec<PeerResponse>>> {
let peers = state.store.list_all_peers().await?;
let enriched = enrich_peers_with_live_telemetry(&state, peers).await;
Ok(Json(enriched))
}
/// GET /api/v1/interfaces/{id}/peers
pub async fn list_peers_for_interface_handler(
State(state): State<AppState>,
Path(interface_id): Path<Uuid>,
) -> ApiResult<Json<Vec<Peer>>> {
) -> ApiResult<Json<Vec<PeerResponse>>> {
let peers = state.store.list_peers_for_interface(interface_id).await?;
Ok(Json(peers))
let enriched = enrich_peers_with_live_telemetry(&state, peers).await;
Ok(Json(enriched))
}
async fn validate_no_server_allowed_ips_conflict(
store: &nx9_wg_db::Store,
interface_id: Uuid,
peer_id: Option<Uuid>,
candidate_server_allowed_ips: &str,
) -> ApiResult<()> {
if candidate_server_allowed_ips.trim().is_empty() {
return Ok(());
}
let candidate_nets: Vec<IpNet> = candidate_server_allowed_ips
.split(',')
.map(|s| s.trim())
.filter(|s| !s.is_empty())
.filter_map(|s| s.parse::<IpNet>().ok())
.collect();
if candidate_nets.is_empty() {
return Ok(());
}
let existing_peers = store.list_peers_for_interface(interface_id).await?;
for ep in existing_peers {
if ep.state != PeerState::Active {
continue;
}
if Some(ep.id) == peer_id {
continue;
}
let ep_server_allowed = ep.server_wireguard_allowed_ips();
let ep_nets: Vec<IpNet> = ep_server_allowed
.split(',')
.map(|s| s.trim())
.filter(|s| !s.is_empty())
.filter_map(|s| s.parse::<IpNet>().ok())
.collect();
for n1 in &candidate_nets {
for n2 in &ep_nets {
if n1.contains(n2) || n2.contains(n1) {
return Err(ApiError::Validation(format!(
"Server-side AllowedIP '{n1}' overlaps with active peer '{}' AllowedIP '{n2}'",
ep.name
)));
}
}
}
}
Ok(())
}
/// POST /api/v1/interfaces/{id}/peers
@@ -182,6 +370,15 @@ pub async fn create_peer_handler(
updated_at: now,
};
// Validate no overlapping server-side AllowedIPs with active peers on the same interface
validate_no_server_allowed_ips_conflict(
&state.store,
interface_id,
None,
&peer.server_wireguard_allowed_ips(),
)
.await?;
state.store.create_peer(&peer).await?;
state.broadcast(SystemEvent::PeerChanged {
@@ -196,13 +393,15 @@ pub async fn create_peer_handler(
pub async fn get_peer_handler(
State(state): State<AppState>,
Path(id): Path<Uuid>,
) -> ApiResult<Json<Peer>> {
) -> ApiResult<Json<PeerResponse>> {
let peer = state
.store
.get_peer(id)
.await?
.ok_or_else(|| ApiError::NotFound(format!("Peer '{id}' not found")))?;
Ok(Json(peer))
let mut enriched = enrich_peers_with_live_telemetry(&state, vec![peer]).await;
let peer_resp = enriched.pop().unwrap();
Ok(Json(peer_resp))
}
/// PUT /api/v1/peers/{id}
@@ -256,6 +455,15 @@ pub async fn update_peer_handler(
peer.expires_at = payload.expires_at;
}
// Validate no overlapping server-side AllowedIPs with active peers on the same interface
validate_no_server_allowed_ips_conflict(
&state.store,
peer.interface_id,
Some(peer.id),
&peer.server_wireguard_allowed_ips(),
)
.await?;
state.store.update_peer(&peer).await?;
state.broadcast(SystemEvent::PeerChanged {
@@ -391,6 +599,46 @@ pub struct ClientProfileQuery {
pub nat: Option<String>,
pub mtu: Option<u16>,
pub profile: Option<String>,
pub server_endpoint: Option<String>,
pub endpoint: Option<String>,
}
async fn resolve_server_endpoint(
state: &AppState,
query: &ClientProfileQuery,
) -> ApiResult<String> {
// 1. Explicit query parameter (server_endpoint or endpoint)
if let Some(ep) = query
.server_endpoint
.as_deref()
.or(query.endpoint.as_deref())
{
let trimmed = ep.trim();
if !trimmed.is_empty() {
return Ok(trimmed.to_string());
}
}
// 2. Persistent server_endpoint configuration from store
if let Some(setting) = state.store.get_setting("server_endpoint").await? {
let trimmed = setting.value.trim();
if !trimmed.is_empty() {
return Ok(trimmed.to_string());
}
}
// 3. Persistent public_endpoint configuration from store
if let Some(setting) = state.store.get_setting("public_endpoint").await? {
let trimmed = setting.value.trim();
if !trimmed.is_empty() {
return Ok(trimmed.to_string());
}
}
// Explicit actionable error if no reachable server endpoint is configured
Err(ApiError::Validation(
"No reachable WireGuard server endpoint is configured. Configure 'server_endpoint' in settings or provide --endpoint / query parameter.".to_string(),
))
}
#[derive(Debug, serde::Serialize)]
@@ -418,12 +666,7 @@ pub async fn download_peer_config_handler(
.await?
.ok_or_else(|| ApiError::NotFound("Associated interface not found".to_string()))?;
let host = state
.store
.get_setting("server_endpoint")
.await?
.map(|s| s.value)
.unwrap_or_else(|| "127.0.0.1".to_string());
let host = resolve_server_endpoint(&state, &query).await?;
let resolved_profile = if query.provider.is_some()
|| query.device.is_some()
@@ -509,12 +752,7 @@ pub async fn get_peer_qr_handler(
.await?
.ok_or_else(|| ApiError::NotFound("Associated interface not found".to_string()))?;
let host = state
.store
.get_setting("server_endpoint")
.await?
.map(|s| s.value)
.unwrap_or_else(|| "127.0.0.1".to_string());
let host = resolve_server_endpoint(&state, &query).await?;
let resolved_profile = if query.provider.is_some()
|| query.device.is_some()
+109 -4
View File
@@ -94,24 +94,129 @@ pub async fn upsert_setting_handler(
State(state): State<AppState>,
Json(payload): Json<UpsertSettingRequest>,
) -> ApiResult<Json<GenericSuccess>> {
if payload.key.trim().is_empty() {
let key_trimmed = payload.key.trim();
if key_trimmed.is_empty() {
return Err(ApiError::Validation(
"Setting key cannot be empty".to_string(),
));
}
let val_trimmed = payload.value.trim();
if (key_trimmed == "server_endpoint" || key_trimmed == "public_endpoint")
&& !val_trimmed.is_empty()
{
let has_valid_port = if let Some(last_colon) = val_trimmed.rfind(':') {
let port_str = &val_trimmed[last_colon + 1..];
if let Ok(port) = port_str.parse::<u16>() {
port > 0 && !val_trimmed[..last_colon].trim().is_empty()
} else {
false
}
} else {
false
};
if !has_valid_port {
return Err(ApiError::Validation(format!(
"Invalid server endpoint '{val_trimmed}'. Endpoint must be formatted as host:port (e.g. 192.168.1.8:51820 or vpn.domain.com:51820)"
)));
}
}
let is_secret = payload.is_secret.unwrap_or(false);
state
.store
.set_setting(&payload.key, &payload.value, is_secret)
.set_setting(key_trimmed, val_trimmed, is_secret)
.await?;
state.broadcast(SystemEvent::SettingsChanged {
key: payload.key.clone(),
key: key_trimmed.to_string(),
});
Ok(Json(GenericSuccess {
success: true,
message: format!("Setting '{}' saved successfully", payload.key),
message: format!("Setting '{key_trimmed}' saved"),
}))
}
#[derive(Debug, Serialize)]
pub struct LiveInterfaceTelemetry {
pub name: String,
pub public_key: String,
pub listen_port: u16,
pub fwmark: u32,
pub addresses: Vec<String>,
pub mtu: Option<u32>,
pub is_up: bool,
pub peer_count: usize,
pub peers: Vec<nx9_wireguard::LivePeerStats>,
}
#[derive(Debug, Serialize)]
pub struct LiveRouteTelemetry {
pub destination: String,
pub gateway: Option<String>,
pub metric: Option<u32>,
pub table: u32,
}
#[derive(Debug, Serialize)]
pub struct LiveSystemState {
pub interfaces: Vec<LiveInterfaceTelemetry>,
pub routes: Vec<LiveRouteTelemetry>,
pub ipv4_forwarding: bool,
pub ipv6_forwarding: bool,
pub active_nftables: Option<String>,
}
/// GET /api/v1/system/live-state
pub async fn live_state_handler(State(state): State<AppState>) -> ApiResult<Json<LiveSystemState>> {
let iface_names = state.wg_engine.list_interfaces().await.unwrap_or_default();
let mut interfaces = Vec::new();
for name in iface_names {
if let Ok(Some(stats)) = state.wg_engine.get_interface_stats(&name).await {
let peer_count = stats.peers.len();
interfaces.push(LiveInterfaceTelemetry {
name: stats.name,
public_key: stats.public_key,
listen_port: stats.listen_port,
fwmark: stats.fwmark,
addresses: stats.addresses,
mtu: stats.mtu,
is_up: stats.is_up,
peer_count,
peers: stats.peers,
});
}
}
let desired_routes = state.store.list_routes().await.unwrap_or_default();
let routes = desired_routes
.into_iter()
.filter(|r| r.enabled)
.map(|r| LiveRouteTelemetry {
destination: r.destination.to_string(),
gateway: r.gateway.map(|g| g.to_string()),
metric: r.metric,
table: 254,
})
.collect();
let fwd = state.net_engine.get_forwarding_status().await.unwrap_or(
nx9_wg_network::IpForwardingStatus {
ipv4_enabled: false,
ipv6_enabled: false,
},
);
let active_nftables = state.net_engine.get_active_nftables_ruleset().await.ok();
Ok(Json(LiveSystemState {
interfaces,
routes,
ipv4_forwarding: fwd.ipv4_enabled,
ipv6_forwarding: fwd.ipv6_enabled,
active_nftables,
}))
}
+32 -1
View File
@@ -3,7 +3,14 @@
use crate::auth::service::AuthService;
use nx9_wg_core::types::audit::AuditEventType;
use nx9_wg_db::Store;
#[allow(unused_imports)]
use nx9_wg_network::engine::{NativeLinuxNetworkEngine, NetworkEngine, SimulatedNetworkEngine};
#[allow(unused_imports)]
use nx9_wireguard::engine::{
NativeLinuxWireGuardEngine, SimulatedWireGuardEngine, WireGuardEngine,
};
use serde::{Deserialize, Serialize};
use std::sync::Arc;
use tokio::sync::broadcast;
/// Real-time system event broadcasted over WebSocket to connected clients.
@@ -39,17 +46,41 @@ pub struct AppState {
pub store: Store,
pub auth: AuthService,
pub event_tx: broadcast::Sender<SystemEvent>,
pub wg_engine: Arc<dyn WireGuardEngine>,
pub net_engine: Arc<dyn NetworkEngine>,
}
impl AppState {
/// Create a new AppState instance.
/// Create a new AppState instance with default engines.
pub fn new(store: Store) -> Self {
#[cfg(target_os = "linux")]
let (wg, net): (Arc<dyn WireGuardEngine>, Arc<dyn NetworkEngine>) = (
Arc::new(NativeLinuxWireGuardEngine::new()),
Arc::new(NativeLinuxNetworkEngine::new()),
);
#[cfg(not(target_os = "linux"))]
let (wg, net): (Arc<dyn WireGuardEngine>, Arc<dyn NetworkEngine>) = (
Arc::new(SimulatedWireGuardEngine::new()),
Arc::new(SimulatedNetworkEngine::new()),
);
Self::with_engines(store, wg, net)
}
/// Create a new AppState instance with custom engines.
pub fn with_engines(
store: Store,
wg_engine: Arc<dyn WireGuardEngine>,
net_engine: Arc<dyn NetworkEngine>,
) -> Self {
let (event_tx, _) = broadcast::channel(256);
let auth = AuthService::new(store.clone());
Self {
store,
auth,
event_tx,
wg_engine,
net_engine,
}
}