diff --git a/members/nullnet-client/src/commands/mod.rs b/members/nullnet-client/src/commands/mod.rs index 00c7df2..d134910 100644 --- a/members/nullnet-client/src/commands/mod.rs +++ b/members/nullnet-client/src/commands/mod.rs @@ -124,11 +124,14 @@ impl RtNetLinkHandle { } } -pub(crate) async fn cleanup_network(rtnetlink_handle: &RtNetLinkHandle) { +/// Returns the MSS-clamp install error, if there was one. It can't be reported +/// from here — this runs before the control connection exists — so the caller +/// emits the event once it does. +pub(crate) async fn cleanup_network(rtnetlink_handle: &RtNetLinkHandle) -> Option { dnat::init(); nfqueue::init(); egress::init(); - install_mss_clamp(); + let mss_error = install_mss_clamp(); vxlan_cleanup_network(); vlan_cleanup_network(rtnetlink_handle).await; // State a killed process never got to tear down, and that no `VxlanTeardown` @@ -139,6 +142,7 @@ pub(crate) async fn cleanup_network(rtnetlink_handle: &RtNetLinkHandle) { purge_stale_xfrm(); crate::host_mappings::purge_stale_mappings(); egress::purge_stale_steers(); + mss_error } /// SPI range `vxlan-setup.sh` can install: it offsets the net id by 1000 to @@ -346,7 +350,7 @@ src 10.20.30.1/32 dst 10.20.30.2/32 /// removed first — `-C` only matches a rule verbatim, so without that an /// upgrade would leave two clamps installed and the older one would win by /// position. -fn install_mss_clamp() { +fn install_mss_clamp() -> Option { prune_superseded_mss_rules(); // Must match OVERLAY_MTU in vxlan_scripts/vxlan-setup.sh (1080) minus the // 40-byte IP+TCP headers. The previous 1400 came from a theoretical @@ -371,16 +375,23 @@ fn install_mss_clamp() { let mut check = vec!["iptables", "-t", "mangle", "-C", "FORWARD"]; check.extend_from_slice(&rule); if sudo(&check).map(|s| s.success()).unwrap_or(false) { - return; + return None; } let mut add = vec!["iptables", "-t", "mangle", "-A", "FORWARD"]; add.extend_from_slice(&rule); match sudo(&add) { Ok(s) if s.success() => { - println!("[mss] clamp installed on mangle/FORWARD: --set-mss {MSS}") + println!("[mss] clamp installed on mangle/FORWARD: --set-mss {MSS}"); + None + } + Ok(s) => { + eprintln!("[mss] clamp install exited {s}"); + Some(format!("iptables exited {s}")) + } + Err(e) => { + eprintln!("[mss] clamp install failed: {e}"); + Some(e.to_string()) } - Ok(s) => eprintln!("[mss] clamp install exited {s}"), - Err(e) => eprintln!("[mss] clamp install failed: {e}"), } } diff --git a/members/nullnet-client/src/control_channel.rs b/members/nullnet-client/src/control_channel.rs index c203fce..6f58d23 100644 --- a/members/nullnet-client/src/control_channel.rs +++ b/members/nullnet-client/src/control_channel.rs @@ -13,8 +13,8 @@ use nullnet_grpc_lib::NullnetGrpcInterface; use nullnet_grpc_lib::nullnet_grpc::{ AgentContainerResumeFailed, AgentContainerSuspendFailed, AgentControlChannelAckFailed, AgentControlChannelClosed, AgentControlChannelEstablished, AgentDnatInstallFailed, - AgentDnatRemovalFailed, AgentGatewayForwardInstallFailed, AgentHostMappingFailed, - AgentVlanSetupCompleted, AgentVlanSetupFailed, AgentVlanTeardownFailed, + AgentDnatRemovalFailed, AgentEgressSteerInstallFailed, AgentGatewayForwardInstallFailed, + AgentHostMappingFailed, AgentVlanSetupCompleted, AgentVlanSetupFailed, AgentVlanTeardownFailed, AgentVxlanSetupCompleted, AgentVxlanSetupFailed, AgentVxlanTeardownFailed, }; use nullnet_grpc_lib::nullnet_grpc::{ @@ -174,10 +174,11 @@ pub(crate) async fn control_channel( // re-enter the NFQUEUE as NEW — newly-denied ones die there. let verdicts = policy_verdicts.clone(); let cache = bridge_cache.clone(); + let grpc = server.clone(); tokio::spawn(async move { println!("[egress-policy] policy changed on server; re-verdicting flows"); verdicts.clear(); - flush_container_conntrack(cache.ips()).await; + flush_container_conntrack(&grpc, cache.ips()).await; }); } None => {} @@ -537,6 +538,17 @@ async fn handle_vxlan_setup( cip, ); } + } else { + // Mirrors the DNAT path below: without steering the held + // packet is never woken and drops at ACTIVE_TIMEOUT. + fire_event( + &grpc, + AgentEventKind::EgressSteerInstallFailed(AgentEgressSteerInstallFailed { + vxlan_id, + docker_container: message.docker_container.clone(), + error_message: "steer rules failed to install".to_string(), + }), + ); } } _ => { @@ -544,6 +556,16 @@ async fn handle_vxlan_setup( "[vxlan_setup] egress steer missing gateway/container_ip \ (gw={proxy_gw:?}, cip={container_ip:?}); steering not installed" ); + fire_event( + &grpc, + AgentEventKind::EgressSteerInstallFailed(AgentEgressSteerInstallFailed { + vxlan_id, + docker_container: message.docker_container.clone(), + error_message: format!( + "missing gateway/container_ip (gw={proxy_gw:?}, cip={container_ip:?})" + ), + }), + ); } } } else if egress_intercept { diff --git a/members/nullnet-client/src/egress_policy.rs b/members/nullnet-client/src/egress_policy.rs index e57ed5b..b08c130 100644 --- a/members/nullnet-client/src/egress_policy.rs +++ b/members/nullnet-client/src/egress_policy.rs @@ -7,6 +7,10 @@ //! so live flows re-enter the queue as NEW and get re-verdicted — flows the //! new policy denies die on their next packet. +use nullnet_grpc_lib::NullnetGrpcInterface; +use nullnet_grpc_lib::nullnet_grpc::{ + AgentConntrackFlushFailed, AgentEvent, agent_event::Event as AgentEventKind, +}; use std::collections::HashMap; use std::net::Ipv4Addr; use std::sync::Mutex; @@ -56,22 +60,40 @@ impl PolicyVerdicts { /// Delete the conntrack entries originating from each container bridge IP so /// every live flow re-enters the NFQUEUE as NEW and is re-verdicted. Exit /// code 1 just means "no entries matched" — only real failures are logged. -pub async fn flush_container_conntrack(ips: Vec) { +pub async fn flush_container_conntrack(grpc: &NullnetGrpcInterface, ips: Vec) { for ip in ips { let out = tokio::process::Command::new("conntrack") .args(["-D", "-s", &ip.to_string()]) .output() .await; - match out { - Ok(o) if o.status.code() == Some(0) || o.status.code() == Some(1) => {} - Ok(o) => eprintln!( - "[egress-policy] conntrack -D -s {ip} exited {}: {}", - o.status, - String::from_utf8_lossy(&o.stderr).trim() - ), + // A failed flush leaves flows the new policy denies running until they + // close on their own, so the policy change is only partly in force. + let error_message = match out { + Ok(o) if o.status.code() == Some(0) || o.status.code() == Some(1) => continue, + Ok(o) => { + let stderr = String::from_utf8_lossy(&o.stderr).trim().to_string(); + eprintln!( + "[egress-policy] conntrack -D -s {ip} exited {}: {stderr}", + o.status + ); + format!("conntrack exited {}: {stderr}", o.status) + } Err(e) => { eprintln!("[egress-policy] conntrack flush {ip}: {e} (is conntrack installed?)"); + format!("{e} (is conntrack installed?)") } - } + }; + let grpc = grpc.clone(); + let event = AgentEvent { + event: Some(AgentEventKind::ConntrackFlushFailed( + AgentConntrackFlushFailed { + ip: ip.to_string(), + error_message, + }, + )), + }; + tokio::spawn(async move { + let _ = grpc.report_event(event).await; + }); } } diff --git a/members/nullnet-client/src/main.rs b/members/nullnet-client/src/main.rs index 36d3e8f..e39ad00 100644 --- a/members/nullnet-client/src/main.rs +++ b/members/nullnet-client/src/main.rs @@ -16,8 +16,9 @@ use crate::triggers::TriggersState; use clap::Parser; use nullnet_grpc_lib::NullnetGrpcInterface; use nullnet_grpc_lib::nullnet_grpc::{ - AgentEvent, AgentServicesListUpdateFailed, AgentServicesListUpdated, Container, Listener, Net, - NetType, ServiceReport, agent_event::Event as AgentEventKind, + AgentEvent, AgentFirewallRulesLoadFailed, AgentMssClampInstallFailed, + AgentServicesListUpdateFailed, AgentServicesListUpdated, Container, Listener, Net, NetType, + ServiceReport, agent_event::Event as AgentEventKind, }; use nullnet_liberror::{Error, ErrorHandler, Location, location}; use std::collections::HashMap; @@ -77,7 +78,7 @@ async fn main() -> Result<(), Error> { let rtnetlink_handle = RtNetLinkHandle::new()?; // cleanup existing VLANs and VXLANs material - cleanup_network(&rtnetlink_handle).await; + let mss_error = cleanup_network(&rtnetlink_handle).await; // maps of all the peers let peers = Arc::new(RwLock::new(Peers::default())); @@ -88,6 +89,22 @@ async fn main() -> Result<(), Error> { let grpc_server2 = grpc_server.clone(); let grpc_server3 = grpc_server.clone(); + // Deferred from `cleanup_network`, which runs before this connection exists. + // Without the clamp, oversized segments are silently black-holed once they + // enter a tunnel, which is near-impossible to trace from the symptom. + if let Some(error_message) = mss_error { + let grpc = grpc_server.clone(); + tokio::spawn(async move { + let _ = grpc + .report_event(AgentEvent { + event: Some(AgentEventKind::MssClampInstallFailed( + AgentMssClampInstallFailed { error_message }, + )), + }) + .await; + }); + } + let net_type = grpc_server.network_type().await.handle_err(location!())?; if net_type.net() == Net::Vlan { @@ -108,6 +125,22 @@ async fn main() -> Result<(), Error> { } Err(e) => { eprintln!("Failed to enable eBPF firewall: {e:?}"); + // Awaited, not fire-and-forget: the exit below would kill a spawned + // task before it ever reached the server. Bounded, because + // `report_event` has no request timeout of its own and a hung one + // would keep us from exiting at all. + let _ = tokio::time::timeout( + Duration::from_secs(5), + grpc_server.report_event(AgentEvent { + event: Some(AgentEventKind::FirewallRulesLoadFailed( + AgentFirewallRulesLoadFailed { + path: "ebpf host firewall".to_string(), + error_message: format!("{e:?}"), + }, + )), + }), + ) + .await; process::exit(1); } }; diff --git a/members/nullnet-client/src/nfqueue/egress_listener.rs b/members/nullnet-client/src/nfqueue/egress_listener.rs index cf98db4..5c8b0ad 100644 --- a/members/nullnet-client/src/nfqueue/egress_listener.rs +++ b/members/nullnet-client/src/nfqueue/egress_listener.rs @@ -23,8 +23,8 @@ use crate::triggers::{EGRESS_TRIGGER_PORT, TriggerState, TriggersState}; use nfq::{Message, Verdict}; use nullnet_grpc_lib::NullnetGrpcInterface; use nullnet_grpc_lib::nullnet_grpc::{ - AgentEgressTriggerSendFailed, AgentEvent, EgressDestinationEntry, - agent_event::Event as AgentEventKind, + AgentEgressPolicyCheckFailed, AgentEgressSteerSetupTimedOut, AgentEgressTriggerSendFailed, + AgentEvent, EgressDestinationEntry, agent_event::Event as AgentEventKind, }; use std::collections::HashMap; use std::net::Ipv4Addr; @@ -105,7 +105,9 @@ pub fn spawn_egress_recv_thread( pending, verdicts, }; + let grpc = ctx.grpc.clone(); spawn_queue_loop( + &grpc, QUEUE_ID, COPY_RANGE, QUEUE_MAX_LEN, @@ -159,7 +161,9 @@ async fn decide_verdict(ctx: &EgressCtx, flow: Option<(Ipv4Addr, Ipv4Addr, u16)> match ctx.triggers_state.state(&container, EGRESS_TRIGGER_PORT) { TriggerState::Active => Verdict::Accept, - TriggerState::Pending(notify) => wait_for_steer(ctx, &container, notify).await, + TriggerState::Pending(notify) => { + wait_for_steer(ctx, &container, dst_ip, dst_port, notify).await + } TriggerState::Fresh => { let notify = ctx .triggers_state @@ -196,6 +200,13 @@ async fn decide_verdict(ctx: &EgressCtx, flow: Option<(Ipv4Addr, Ipv4Addr, u16)> // Pending (ages out at PENDING_TIMEOUT) so a retransmit // re-waits on the same notify; drop this held SYN. eprintln!("[egress-nfq] no egress steer for container {container}"); + report_steer_timed_out( + &ctx.grpc, + &container, + dst_ip, + dst_port, + format!("no egress steer within {STEER_TIMEOUT:?}"), + ); Verdict::Drop } }, @@ -242,21 +253,36 @@ async fn policy_allows(ctx: &EgressCtx, container: &str, dst_ip: Ipv4Addr) -> bo ctx.verdicts.put(container, dst_ip, allowed); allowed } + // Fail-closed means the flow is dropped, so the operator should see it: + // a control-plane blip blackholes egress with no other symptom. Ok(Err(e)) => { eprintln!("[egress-nfq] policy check {container} -> {dst_ip}: {e}; failing closed"); + report_policy_check_failed(&ctx.grpc, container, dst_ip, e); false } Err(_) => { eprintln!( "[egress-nfq] policy check timeout for {container} -> {dst_ip}; failing closed" ); + report_policy_check_failed( + &ctx.grpc, + container, + dst_ip, + format!("policy check timed out after {POLICY_TIMEOUT:?}"), + ); false } } } /// Hold the packet until steering is marked active (or time out and drop it). -async fn wait_for_steer(ctx: &EgressCtx, container: &str, notify: Arc) -> Verdict { +async fn wait_for_steer( + ctx: &EgressCtx, + container: &str, + dst_ip: Ipv4Addr, + dst_port: u16, + notify: Arc, +) -> Verdict { let notified = notify.notified(); tokio::pin!(notified); if notified.as_mut().enable() @@ -271,6 +297,13 @@ async fn wait_for_steer(ctx: &EgressCtx, container: &str, notify: Arc) - Ok(_) => Verdict::Accept, Err(_) => { eprintln!("[egress-nfq] timeout waiting for egress steer, container {container}"); + report_steer_timed_out( + &ctx.grpc, + container, + dst_ip, + dst_port, + format!("steering not active after {STEER_TIMEOUT:?}"), + ); Verdict::Drop } } @@ -344,6 +377,55 @@ fn spawn_flush_task(grpc: NullnetGrpcInterface, pending: PendingDsts) { }); } +/// Fire-and-forget: report a held packet dropped because steering never went +/// live. The trigger was accepted — see [`report_trigger_send_failed`] for the +/// case where the RPC itself failed. +fn report_steer_timed_out( + grpc: &NullnetGrpcInterface, + container: &str, + dst_ip: Ipv4Addr, + dst_port: u16, + error_message: String, +) { + let grpc = grpc.clone(); + let event = AgentEvent { + event: Some(AgentEventKind::EgressSteerSetupTimedOut( + AgentEgressSteerSetupTimedOut { + docker_container: container.to_string(), + dst_ip: dst_ip.to_string(), + dst_port: u32::from(dst_port), + error_message, + }, + )), + }; + tokio::spawn(async move { + let _ = grpc.report_event(event).await; + }); +} + +/// Fire-and-forget: report an unresolvable egress policy check. The flow was +/// denied by the fail-closed rule, so this is a drop, not just a warning. +fn report_policy_check_failed( + grpc: &NullnetGrpcInterface, + container: &str, + dst_ip: Ipv4Addr, + error_message: String, +) { + let grpc = grpc.clone(); + let event = AgentEvent { + event: Some(AgentEventKind::EgressPolicyCheckFailed( + AgentEgressPolicyCheckFailed { + docker_container: container.to_string(), + dst_ip: dst_ip.to_string(), + error_message, + }, + )), + }; + tokio::spawn(async move { + let _ = grpc.report_event(event).await; + }); +} + /// Fire-and-forget: report a failed `egress_trigger` to the server's event stream. fn report_trigger_send_failed( grpc: &NullnetGrpcInterface, diff --git a/members/nullnet-client/src/nfqueue/listener.rs b/members/nullnet-client/src/nfqueue/listener.rs index 03c6c4b..0d0ea7d 100644 --- a/members/nullnet-client/src/nfqueue/listener.rs +++ b/members/nullnet-client/src/nfqueue/listener.rs @@ -5,7 +5,8 @@ use crate::triggers::{TriggerState, TriggersState}; use nfq::{Message, Verdict}; use nullnet_grpc_lib::NullnetGrpcInterface; use nullnet_grpc_lib::nullnet_grpc::{ - AgentBackendTriggerSendFailed, AgentEvent, agent_event::Event as AgentEventKind, + AgentBackendTriggerSendFailed, AgentBackendTriggerSetupTimedOut, AgentEvent, + agent_event::Event as AgentEventKind, }; use std::collections::HashMap; use std::sync::mpsc::Sender; @@ -75,7 +76,9 @@ pub struct ListenerCtx { /// Spawn the backend-trigger recv loop (queue 0). Each packet is held until /// `handle_packet` resolves a verdict — see `recv_loop::spawn_queue_loop`. pub fn spawn_recv_thread(ctx: ListenerCtx) { + let grpc = ctx.grpc.clone(); spawn_queue_loop( + &grpc, QUEUE_ID, COPY_RANGE, QUEUE_MAX_LEN, @@ -168,6 +171,13 @@ async fn decide_verdict( eprintln!( "[nfqueue] timeout waiting for active state on '{service}' port {dst_port} container {container}" ); + report_setup_timed_out( + &ctx.grpc, + service, + dst_port, + container, + format!("chain not active after {ACTIVE_TIMEOUT:?}"), + ); Verdict::Drop } } @@ -211,6 +221,13 @@ async fn decide_verdict( eprintln!( "[nfqueue] no VxlanSetup for '{service}' port {dst_port} container {container}" ); + report_setup_timed_out( + &ctx.grpc, + service, + dst_port, + container, + format!("no VxlanSetup within {ACTIVE_TIMEOUT:?}"), + ); Verdict::Drop } }, @@ -263,6 +280,32 @@ fn report_trigger_send_failed( }); } +/// Fire-and-forget: report a held packet dropped because the chain never went +/// active. Distinct from [`report_trigger_send_failed`] — the trigger itself was +/// accepted here, so the failure is the setup not landing, not the RPC. +fn report_setup_timed_out( + grpc: &NullnetGrpcInterface, + service: &str, + port: u16, + container: &str, + error_message: String, +) { + let grpc = grpc.clone(); + let event = AgentEvent { + event: Some(AgentEventKind::BackendTriggerSetupTimedOut( + AgentBackendTriggerSetupTimedOut { + service_name: service.to_string(), + port: u32::from(port), + docker_container: container.to_string(), + error_message, + }, + )), + }; + tokio::spawn(async move { + let _ = grpc.report_event(event).await; + }); +} + #[cfg(test)] mod tests { use super::{TriggerMap, TriggerOwner, owner_for}; diff --git a/members/nullnet-client/src/nfqueue/recv_loop.rs b/members/nullnet-client/src/nfqueue/recv_loop.rs index 7f8a207..cf65193 100644 --- a/members/nullnet-client/src/nfqueue/recv_loop.rs +++ b/members/nullnet-client/src/nfqueue/recv_loop.rs @@ -9,6 +9,10 @@ //! the datapath (DNAT / egress steer) is ready, so the original packet isn't lost. use nfq::{Message, Queue}; +use nullnet_grpc_lib::NullnetGrpcInterface; +use nullnet_grpc_lib::nullnet_grpc::{ + AgentEvent, AgentNfqueueBindFailed, agent_event::Event as AgentEventKind, +}; use std::future::Future; use std::sync::mpsc::{Sender, TryRecvError}; use std::time::Duration; @@ -25,22 +29,30 @@ const IDLE_SLEEP: Duration = Duration::from_millis(1); /// Failure to open/bind the queue is logged and the thread exits; the rest of the /// client keeps running. With `--queue-bypass` on the iptables rule, the absence /// of a consumer fail-opens, so traffic flows unaltered. -pub fn spawn_queue_loop(queue_id: u16, copy_range: u16, queue_max_len: u32, handler: H) -where +pub fn spawn_queue_loop( + grpc: &NullnetGrpcInterface, + queue_id: u16, + copy_range: u16, + queue_max_len: u32, + handler: H, +) where H: Fn(Message, Sender) -> Fut + Send + 'static, Fut: Future + Send + 'static, { let runtime = tokio::runtime::Handle::current(); + let grpc = grpc.clone(); std::thread::spawn(move || { let mut queue = match Queue::open() { Ok(q) => q, Err(e) => { eprintln!("[nfqueue] open queue {queue_id} failed: {e} (need CAP_NET_ADMIN)"); + report_bind_failed(&runtime, &grpc, queue_id, format!("open: {e}")); return; } }; if let Err(e) = queue.bind(queue_id) { eprintln!("[nfqueue] bind queue {queue_id} failed: {e}"); + report_bind_failed(&runtime, &grpc, queue_id, format!("bind: {e}")); return; } if let Err(e) = queue.set_copy_range(queue_id, copy_range) { @@ -95,3 +107,25 @@ where } }); } + +/// Report a queue that never came up. Spawned onto the runtime because the +/// caller is a plain OS thread with no reactor of its own. Worth an event: with +/// no consumer, `--queue-bypass` fail-opens, so trigger detection and egress +/// policy are both off on this queue until the client restarts. +fn report_bind_failed( + runtime: &tokio::runtime::Handle, + grpc: &NullnetGrpcInterface, + queue_id: u16, + error_message: String, +) { + let grpc = grpc.clone(); + let event = AgentEvent { + event: Some(AgentEventKind::NfqueueBindFailed(AgentNfqueueBindFailed { + queue_id: u32::from(queue_id), + error_message, + })), + }; + runtime.spawn(async move { + let _ = grpc.report_event(event).await; + }); +} diff --git a/members/nullnet-grpc-lib/proto/nullnet_grpc.proto b/members/nullnet-grpc-lib/proto/nullnet_grpc.proto index 9c5fcba..35d7178 100644 --- a/members/nullnet-grpc-lib/proto/nullnet_grpc.proto +++ b/members/nullnet-grpc-lib/proto/nullnet_grpc.proto @@ -455,6 +455,13 @@ message AgentEvent { AgentContainerResumeFailed container_resume_failed = 25; AgentEgressTriggerSendFailed egress_trigger_send_failed = 30; AgentGatewayForwardInstallFailed gateway_forward_install_failed = 31; + AgentBackendTriggerSetupTimedOut backend_trigger_setup_timed_out = 32; + AgentEgressSteerSetupTimedOut egress_steer_setup_timed_out = 33; + AgentEgressSteerInstallFailed egress_steer_install_failed = 34; + AgentNfqueueBindFailed nfqueue_bind_failed = 35; + AgentMssClampInstallFailed mss_clamp_install_failed = 36; + AgentEgressPolicyCheckFailed egress_policy_check_failed = 37; + AgentConntrackFlushFailed conntrack_flush_failed = 38; // Client info events AgentVxlanSetupCompleted vxlan_setup_completed = 13; AgentVlanSetupCompleted vlan_setup_completed = 14; @@ -492,6 +499,13 @@ message AgentContainerSuspendFailed { string docker_container = 1; string error message AgentContainerResumeFailed { string docker_container = 1; string error_message = 2; } message AgentEgressTriggerSendFailed { string service_name = 1; string dst_ip = 2; uint32 dst_port = 3; string error_message = 4; } message AgentGatewayForwardInstallFailed { uint32 vxlan_id = 1; string br_net = 2; } +message AgentBackendTriggerSetupTimedOut { string service_name = 1; uint32 port = 2; string docker_container = 3; string error_message = 4; } +message AgentEgressSteerSetupTimedOut { string docker_container = 1; string dst_ip = 2; uint32 dst_port = 3; string error_message = 4; } +message AgentEgressSteerInstallFailed { uint32 vxlan_id = 1; optional string docker_container = 2; string error_message = 3; } +message AgentNfqueueBindFailed { uint32 queue_id = 1; string error_message = 2; } +message AgentMssClampInstallFailed { string error_message = 1; } +message AgentEgressPolicyCheckFailed { string docker_container = 1; string dst_ip = 2; string error_message = 3; } +message AgentConntrackFlushFailed { string ip = 1; string error_message = 2; } message AgentVxlanSetupCompleted { uint32 vxlan_id = 1; string ns_name = 2; } message AgentVlanSetupCompleted { uint32 vlan_id = 1; } message AgentControlChannelEstablished {} diff --git a/members/nullnet-grpc-lib/src/proto/nullnet_grpc.rs b/members/nullnet-grpc-lib/src/proto/nullnet_grpc.rs index cda8be5..2f142f1 100644 --- a/members/nullnet-grpc-lib/src/proto/nullnet_grpc.rs +++ b/members/nullnet-grpc-lib/src/proto/nullnet_grpc.rs @@ -457,7 +457,7 @@ pub struct Empty {} pub struct AgentEvent { #[prost( oneof = "agent_event::Event", - tags = "1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 24, 25, 30, 31, 13, 14, 15, 16, 17, 18, 19, 20, 21, 23, 26, 27, 28, 29, 22" + tags = "1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 24, 25, 30, 31, 32, 33, 34, 35, 36, 37, 38, 13, 14, 15, 16, 17, 18, 19, 20, 21, 23, 26, 27, 28, 29, 22" )] pub event: ::core::option::Option, } @@ -498,6 +498,20 @@ pub mod agent_event { EgressTriggerSendFailed(super::AgentEgressTriggerSendFailed), #[prost(message, tag = "31")] GatewayForwardInstallFailed(super::AgentGatewayForwardInstallFailed), + #[prost(message, tag = "32")] + BackendTriggerSetupTimedOut(super::AgentBackendTriggerSetupTimedOut), + #[prost(message, tag = "33")] + EgressSteerSetupTimedOut(super::AgentEgressSteerSetupTimedOut), + #[prost(message, tag = "34")] + EgressSteerInstallFailed(super::AgentEgressSteerInstallFailed), + #[prost(message, tag = "35")] + NfqueueBindFailed(super::AgentNfqueueBindFailed), + #[prost(message, tag = "36")] + MssClampInstallFailed(super::AgentMssClampInstallFailed), + #[prost(message, tag = "37")] + EgressPolicyCheckFailed(super::AgentEgressPolicyCheckFailed), + #[prost(message, tag = "38")] + ConntrackFlushFailed(super::AgentConntrackFlushFailed), /// Client info events #[prost(message, tag = "13")] VxlanSetupCompleted(super::AgentVxlanSetupCompleted), @@ -655,6 +669,65 @@ pub struct AgentGatewayForwardInstallFailed { pub br_net: ::prost::alloc::string::String, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct AgentBackendTriggerSetupTimedOut { + #[prost(string, tag = "1")] + pub service_name: ::prost::alloc::string::String, + #[prost(uint32, tag = "2")] + pub port: u32, + #[prost(string, tag = "3")] + pub docker_container: ::prost::alloc::string::String, + #[prost(string, tag = "4")] + pub error_message: ::prost::alloc::string::String, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct AgentEgressSteerSetupTimedOut { + #[prost(string, tag = "1")] + pub docker_container: ::prost::alloc::string::String, + #[prost(string, tag = "2")] + pub dst_ip: ::prost::alloc::string::String, + #[prost(uint32, tag = "3")] + pub dst_port: u32, + #[prost(string, tag = "4")] + pub error_message: ::prost::alloc::string::String, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct AgentEgressSteerInstallFailed { + #[prost(uint32, tag = "1")] + pub vxlan_id: u32, + #[prost(string, optional, tag = "2")] + pub docker_container: ::core::option::Option<::prost::alloc::string::String>, + #[prost(string, tag = "3")] + pub error_message: ::prost::alloc::string::String, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct AgentNfqueueBindFailed { + #[prost(uint32, tag = "1")] + pub queue_id: u32, + #[prost(string, tag = "2")] + pub error_message: ::prost::alloc::string::String, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct AgentMssClampInstallFailed { + #[prost(string, tag = "1")] + pub error_message: ::prost::alloc::string::String, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct AgentEgressPolicyCheckFailed { + #[prost(string, tag = "1")] + pub docker_container: ::prost::alloc::string::String, + #[prost(string, tag = "2")] + pub dst_ip: ::prost::alloc::string::String, + #[prost(string, tag = "3")] + pub error_message: ::prost::alloc::string::String, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct AgentConntrackFlushFailed { + #[prost(string, tag = "1")] + pub ip: ::prost::alloc::string::String, + #[prost(string, tag = "2")] + pub error_message: ::prost::alloc::string::String, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct AgentVxlanSetupCompleted { #[prost(uint32, tag = "1")] pub vxlan_id: u32, diff --git a/members/nullnet-server/src/cert_renewal.rs b/members/nullnet-server/src/cert_renewal.rs index 72c74d3..36b1a5c 100644 --- a/members/nullnet-server/src/cert_renewal.rs +++ b/members/nullnet-server/src/cert_renewal.rs @@ -66,6 +66,15 @@ async fn run_pass(events: &EventStore, config: &RenewalConfig) { continue; }; let Some(expiry) = crate::certs::read_expiry(&domain).await else { + // An ACME cert whose expiry can't be read is never renewed, so this + // silently opts it out of the guarantee the loop exists to provide. + eprintln!("Cert renewal: cannot read expiry of '{domain}'; skipping"); + events + .emit(Event::certificate_renewal_failed( + domain.clone(), + "certificate expiry could not be read".to_string(), + )) + .await; continue; }; if expiry > threshold { @@ -82,7 +91,15 @@ async fn run_pass(events: &EventStore, config: &RenewalConfig) { .await; println!("Cert renewal: '{domain}' renewed successfully"); } - Err(e) => eprintln!("Cert renewal: failed to renew '{domain}': {e:#}"), + Err(e) => { + eprintln!("Cert renewal: failed to renew '{domain}': {e:#}"); + events + .emit(Event::certificate_renewal_failed( + domain.clone(), + format!("{e:#}"), + )) + .await; + } } } } diff --git a/members/nullnet-server/src/events.rs b/members/nullnet-server/src/events.rs index 6d6703c..15aa254 100644 --- a/members/nullnet-server/src/events.rs +++ b/members/nullnet-server/src/events.rs @@ -168,6 +168,26 @@ pub(crate) enum Event { port: u16, timestamp: u64, }, + /// A `services/*.toml` change was picked up but could not be parsed, so the + /// previous configuration is still in force. + ConfigReloadFailed { + error_message: String, + timestamp: u64, + }, + /// A file watcher never started (or died), so changes to what it watched are + /// no longer picked up for the lifetime of this process. + FileWatchFailed { + target: String, + error_message: String, + timestamp: u64, + }, + /// The dedicated-UDP-port pool for encrypted cross-host tunnels ran dry; the + /// edge that needed one failed. Mirrors [`Self::NetIdPoolExhausted`]. + UdpPortPoolExhausted { + service: String, + client_ip: String, + timestamp: u64, + }, // --- Client error events --- VxlanSetupFailed { @@ -255,6 +275,61 @@ pub(crate) enum Event { error_message: String, timestamp: u64, }, + /// A held first packet was dropped because the chain never became active in + /// time. The trigger itself was accepted — this is the setup not landing, + /// not the RPC failing (see [`Self::BackendTriggerSendFailed`]). + BackendTriggerSetupTimedOut { + service_name: String, + port: u16, + docker_container: String, + error_message: String, + timestamp: u64, + }, + /// Egress counterpart of [`Self::BackendTriggerSetupTimedOut`]: the held + /// packet was dropped because steering never went live. + EgressSteerSetupTimedOut { + docker_container: String, + dst_ip: String, + dst_port: u32, + error_message: String, + timestamp: u64, + }, + /// Egress steering rules could not be installed for a new edge, so the + /// initiator's held packet will time out and drop. + EgressSteerInstallFailed { + vxlan_id: u32, + docker_container: Option, + error_message: String, + timestamp: u64, + }, + /// An NFQUEUE consumer never started. Trigger detection and egress policy + /// enforcement are both off on that queue for the rest of the process. + NfqueueBindFailed { + queue_id: u32, + error_message: String, + timestamp: u64, + }, + /// The TCP MSS clamp could not be installed, so oversized segments can be + /// silently black-holed once they enter an overlay tunnel. + MssClampInstallFailed { + error_message: String, + timestamp: u64, + }, + /// An egress country-policy check could not be resolved. The flow is denied + /// (fail-closed), so this is a drop the operator should see. + EgressPolicyCheckFailed { + docker_container: String, + dst_ip: String, + error_message: String, + timestamp: u64, + }, + /// Conntrack could not be flushed after a policy change, so flows the new + /// policy denies may keep running until they close on their own. + ConntrackFlushFailed { + ip: String, + error_message: String, + timestamp: u64, + }, // --- Client info events --- VxlanSetupCompleted { @@ -328,6 +403,17 @@ pub(crate) enum Event { timestamp: u64, }, + /// A proxy opened its certificate stream — i.e. a proxy came up. Paired + /// with [`Self::ProxyDisconnected`], mirroring the node events. + ProxyConnected { + ip: String, + timestamp: u64, + }, + ProxyDisconnected { + ip: String, + timestamp: u64, + }, + // --- Proxy info events --- ProxyRequestRouted { service_name: String, @@ -350,6 +436,21 @@ pub(crate) enum Event { domain: String, timestamp: u64, }, + /// Unattended renewal did not produce a usable certificate. Left unattended + /// this ends in an expired cert, so it is an error even though the current + /// one is still serving. + CertificateRenewalFailed { + domain: String, + error_message: String, + timestamp: u64, + }, + /// A certificate was issued but its DNS credentials could not be stored, so + /// unattended renewal will never run for it. + CertificateCredentialsStoreFailed { + domain: String, + error_message: String, + timestamp: u64, + }, } impl Event { @@ -379,6 +480,22 @@ impl Event { Self::NetIdPoolExhausted { .. } => "net_id_pool_exhausted", Self::ProxyChainSetupFailed { .. } => "proxy_chain_setup_failed", Self::BackendTriggerSetupBailed { .. } => "backend_trigger_setup_bailed", + Self::ConfigReloadFailed { .. } => "config_reload_failed", + Self::FileWatchFailed { .. } => "file_watch_failed", + Self::UdpPortPoolExhausted { .. } => "udp_port_pool_exhausted", + Self::BackendTriggerSetupTimedOut { .. } => "backend_trigger_setup_timed_out", + Self::EgressSteerSetupTimedOut { .. } => "egress_steer_setup_timed_out", + Self::EgressSteerInstallFailed { .. } => "egress_steer_install_failed", + Self::NfqueueBindFailed { .. } => "nfqueue_bind_failed", + Self::MssClampInstallFailed { .. } => "mss_clamp_install_failed", + Self::EgressPolicyCheckFailed { .. } => "egress_policy_check_failed", + Self::ConntrackFlushFailed { .. } => "conntrack_flush_failed", + Self::ProxyConnected { .. } => "proxy_connected", + Self::ProxyDisconnected { .. } => "proxy_disconnected", + Self::CertificateRenewalFailed { .. } => "certificate_renewal_failed", + Self::CertificateCredentialsStoreFailed { .. } => { + "certificate_credentials_store_failed" + } Self::VxlanSetupFailed { .. } => "vxlan_setup_failed", Self::VlanSetupFailed { .. } => "vlan_setup_failed", Self::VxlanTeardownFailed { .. } => "vxlan_teardown_failed", @@ -431,6 +548,7 @@ impl Event { | Self::ControlChannelEstablished { .. } | Self::ServicesListUpdated { .. } | Self::ProxyRequestRouted { .. } + | Self::ProxyConnected { .. } | Self::CertificateInstalled { .. } | Self::CertificateRenewed { .. } => Severity::Info, @@ -446,10 +564,23 @@ impl Event { | Self::BackendTriggerSetupBailed { .. } | Self::ControlChannelClosed { .. } | Self::NetTeardownUnconfirmed { .. } + | Self::ConntrackFlushFailed { .. } + | Self::ProxyDisconnected { .. } + | Self::CertificateCredentialsStoreFailed { .. } | Self::CertificateRemoved { .. } => Severity::Warning, Self::SetupTimeout { .. } | Self::NetIdPoolExhausted { .. } + | Self::UdpPortPoolExhausted { .. } + | Self::ConfigReloadFailed { .. } + | Self::FileWatchFailed { .. } + | Self::BackendTriggerSetupTimedOut { .. } + | Self::EgressSteerSetupTimedOut { .. } + | Self::EgressSteerInstallFailed { .. } + | Self::NfqueueBindFailed { .. } + | Self::MssClampInstallFailed { .. } + | Self::EgressPolicyCheckFailed { .. } + | Self::CertificateRenewalFailed { .. } | Self::ProxyChainSetupFailed { .. } | Self::VxlanSetupFailed { .. } | Self::VlanSetupFailed { .. } @@ -1038,6 +1169,141 @@ impl Event { timestamp: now_secs(), } } + + pub(crate) fn certificate_renewal_failed(domain: String, error_message: String) -> Self { + Self::CertificateRenewalFailed { + domain, + error_message, + timestamp: now_secs(), + } + } + + pub(crate) fn certificate_credentials_store_failed( + domain: String, + error_message: String, + ) -> Self { + Self::CertificateCredentialsStoreFailed { + domain, + error_message, + timestamp: now_secs(), + } + } + + pub(crate) fn config_reload_failed(error_message: String) -> Self { + Self::ConfigReloadFailed { + error_message, + timestamp: now_secs(), + } + } + + pub(crate) fn file_watch_failed(target: String, error_message: String) -> Self { + Self::FileWatchFailed { + target, + error_message, + timestamp: now_secs(), + } + } + + pub(crate) fn udp_port_pool_exhausted(service: String, client_ip: String) -> Self { + Self::UdpPortPoolExhausted { + service, + client_ip, + timestamp: now_secs(), + } + } + + pub(crate) fn proxy_connected(ip: String) -> Self { + Self::ProxyConnected { + ip, + timestamp: now_secs(), + } + } + + pub(crate) fn proxy_disconnected(ip: String) -> Self { + Self::ProxyDisconnected { + ip, + timestamp: now_secs(), + } + } + + pub(crate) fn backend_trigger_setup_timed_out( + service_name: String, + port: u16, + docker_container: String, + error_message: String, + ) -> Self { + Self::BackendTriggerSetupTimedOut { + service_name, + port, + docker_container, + error_message, + timestamp: now_secs(), + } + } + + pub(crate) fn egress_steer_setup_timed_out( + docker_container: String, + dst_ip: String, + dst_port: u32, + error_message: String, + ) -> Self { + Self::EgressSteerSetupTimedOut { + docker_container, + dst_ip, + dst_port, + error_message, + timestamp: now_secs(), + } + } + + pub(crate) fn egress_steer_install_failed( + vxlan_id: u32, + docker_container: Option, + error_message: String, + ) -> Self { + Self::EgressSteerInstallFailed { + vxlan_id, + docker_container, + error_message, + timestamp: now_secs(), + } + } + + pub(crate) fn nfqueue_bind_failed(queue_id: u32, error_message: String) -> Self { + Self::NfqueueBindFailed { + queue_id, + error_message, + timestamp: now_secs(), + } + } + + pub(crate) fn mss_clamp_install_failed(error_message: String) -> Self { + Self::MssClampInstallFailed { + error_message, + timestamp: now_secs(), + } + } + + pub(crate) fn egress_policy_check_failed( + docker_container: String, + dst_ip: String, + error_message: String, + ) -> Self { + Self::EgressPolicyCheckFailed { + docker_container, + dst_ip, + error_message, + timestamp: now_secs(), + } + } + + pub(crate) fn conntrack_flush_failed(ip: String, error_message: String) -> Self { + Self::ConntrackFlushFailed { + ip, + error_message, + timestamp: now_secs(), + } + } } /// Shared event store: ring buffer + broadcast channel for SSE subscribers. diff --git a/members/nullnet-server/src/http_server/certificates.rs b/members/nullnet-server/src/http_server/certificates.rs index 3096287..63c53d9 100644 --- a/members/nullnet-server/src/http_server/certificates.rs +++ b/members/nullnet-server/src/http_server/certificates.rs @@ -131,14 +131,30 @@ pub(super) async fn request_handler( // store creds for auto-renew, but only if the cert actually persisted // (otherwise we'd orphan a creds file with no matching cert) if resp.status() == StatusCode::NO_CONTENT { + // Without stored credentials the cert installs fine but can + // never auto-renew, and the operator sees only a success here. match serde_json::to_string(&credentials) { Ok(json) => { if let Err(e) = crate::certs::store_dns_credentials(&domain, &json).await { eprintln!("Failed to store DNS credentials for '{domain}': {e:?}"); + state + .events + .emit(Event::certificate_credentials_store_failed( + domain.clone(), + format!("{e:?}"), + )) + .await; } } Err(e) => { - eprintln!("Failed to serialize DNS credentials for '{domain}': {e:?}") + eprintln!("Failed to serialize DNS credentials for '{domain}': {e:?}"); + state + .events + .emit(Event::certificate_credentials_store_failed( + domain.clone(), + format!("{e:?}"), + )) + .await; } } } diff --git a/members/nullnet-server/src/nullnet_grpc_impl.rs b/members/nullnet-server/src/nullnet_grpc_impl.rs index 6a3f2c3..2c5224b 100644 --- a/members/nullnet-server/src/nullnet_grpc_impl.rs +++ b/members/nullnet-server/src/nullnet_grpc_impl.rs @@ -229,7 +229,7 @@ fn find_service_stack<'a>(services: &'a StackMap, service_name: &str) -> Option< impl NullnetGrpcImpl { pub async fn new() -> Result { - let (stacks, index, route_map) = ServicesToml::load_validated().await?; + let (stacks, index, route_map, startup_conflicts) = ServicesToml::load_validated().await?; let services = Arc::new(RwLock::new(stacks)); let match_index = Arc::new(RwLock::new(index)); let routes = Arc::new(RwLock::new(route_map)); @@ -241,6 +241,29 @@ impl NullnetGrpcImpl { }); let orchestrator = Orchestrator::new(); + + // Conflicts detected before the event store existed: the offending + // stacks were dropped, so report them now that we can. + for c in startup_conflicts.ports { + orchestrator + .events + .emit(Event::port_mapping_conflict( + c.stack_a, + c.service_a, + c.stack_b, + c.service_b, + format!("{:?}", c.protocol), + c.listen_port, + )) + .await; + } + for c in startup_conflicts.routes { + orchestrator + .events + .emit(Event::route_conflict(c.stack_a, c.stack_b, c.host, c.path)) + .await; + } + let config_changed = Arc::new(Notify::new()); // Separate from `config_changed`: `Notify::notify_one` wakes at most // one waiter, so each consumer needs its own `Notify` rather than @@ -256,6 +279,7 @@ impl NullnetGrpcImpl { let config_changed_2 = config_changed.clone(); let port_mappings_changed_2 = port_mappings_changed.clone(); let http_routes_changed_2 = http_routes_changed.clone(); + let events_2 = orchestrator.events.clone(); tokio::spawn(async move { if let Err(e) = ServicesToml::watch( &services_2, @@ -268,7 +292,15 @@ impl NullnetGrpcImpl { ) .await { + // Config hot-reload is dead for the rest of this process: every + // later edit is silently ignored. eprintln!("failed to watch services.toml for changes: {e:?}"); + events_2 + .emit(Event::file_watch_failed( + "services.toml".to_string(), + format!("{e:?}"), + )) + .await; } }); @@ -313,9 +345,18 @@ impl NullnetGrpcImpl { // load TLS certificates and keep them in sync with the ./certs dir let (certs_tx, certs_rx) = watch::channel(crate::certs::load_certificates().await); + let events_3 = orchestrator.events.clone(); tokio::spawn(async move { if let Err(e) = crate::certs::watch(certs_tx).await { + // Renewals still write to disk but never reach the proxies, so + // they keep serving the old cert until it expires. eprintln!("failed to watch certs for changes: {e:?}"); + events_3 + .emit(Event::file_watch_failed( + "certs".to_string(), + format!("{e:?}"), + )) + .await; } }); @@ -1449,6 +1490,13 @@ impl NullnetGrpcImpl { Some(port) => Some(u32::from(port)), None => { eprintln!("UDP port pool exhausted"); + orchestrator + .events + .emit(Event::udp_port_pool_exhausted( + server.name().to_string(), + client_ethernet.to_string(), + )) + .await; orchestrator.free_net_id(net_id).await; if let Some(stack_map) = services.write().await.get_mut(&stack) && let Some(ServiceInfo::Registered(reg)) = @@ -1795,22 +1843,29 @@ impl NullnetGrpc for NullnetGrpcImpl { async fn watch_certificates( &self, - _: Request, + req: Request, ) -> Result, Status> { let mut certs = self.certs.clone(); let (tx, rx) = mpsc::channel(4); + // Every proxy opens this stream once at startup and exits when it drops, + // so its lifetime is the proxy's — the node events' counterpart. + let proxy_ip = req + .remote_addr() + .map_or_else(|| "unknown".to_string(), |a| a.ip().to_string()); + let events = self.orchestrator.events.clone(); + events.emit(Event::proxy_connected(proxy_ip.clone())).await; tokio::spawn(async move { // send the current set immediately, then one snapshot per change let initial = certs.borrow_and_update().clone(); - if tx.send(Ok(initial)).await.is_err() { - return; - } - while certs.changed().await.is_ok() { - let snapshot = certs.borrow_and_update().clone(); - if tx.send(Ok(snapshot)).await.is_err() { - break; + if tx.send(Ok(initial)).await.is_ok() { + while certs.changed().await.is_ok() { + let snapshot = certs.borrow_and_update().clone(); + if tx.send(Ok(snapshot)).await.is_err() { + break; + } } } + events.emit(Event::proxy_disconnected(proxy_ip)).await; }); Ok(Response::new(ReceiverStream::new(rx))) } @@ -1908,6 +1963,35 @@ impl NullnetGrpc for NullnetGrpcImpl { AgentEventKind::GatewayForwardInstallFailed(e) => { Event::gateway_forward_install_failed(e.vxlan_id, e.br_net) } + AgentEventKind::BackendTriggerSetupTimedOut(e) => { + Event::backend_trigger_setup_timed_out( + e.service_name, + e.port as u16, + e.docker_container, + e.error_message, + ) + } + AgentEventKind::EgressSteerSetupTimedOut(e) => Event::egress_steer_setup_timed_out( + e.docker_container, + e.dst_ip, + e.dst_port, + e.error_message, + ), + AgentEventKind::EgressSteerInstallFailed(e) => { + Event::egress_steer_install_failed(e.vxlan_id, e.docker_container, e.error_message) + } + AgentEventKind::NfqueueBindFailed(e) => { + Event::nfqueue_bind_failed(e.queue_id, e.error_message) + } + AgentEventKind::MssClampInstallFailed(e) => { + Event::mss_clamp_install_failed(e.error_message) + } + AgentEventKind::EgressPolicyCheckFailed(e) => { + Event::egress_policy_check_failed(e.docker_container, e.dst_ip, e.error_message) + } + AgentEventKind::ConntrackFlushFailed(e) => { + Event::conntrack_flush_failed(e.ip, e.error_message) + } AgentEventKind::FirewallRulesLoadFailed(e) => { Event::firewall_rules_load_failed(e.path, e.error_message) } diff --git a/members/nullnet-server/src/services/input.rs b/members/nullnet-server/src/services/input.rs index da2e8e0..edb3ed2 100644 --- a/members/nullnet-server/src/services/input.rs +++ b/members/nullnet-server/src/services/input.rs @@ -186,6 +186,13 @@ pub(crate) struct RouteConflict { pub(crate) path: String, } +/// Both conflict kinds found by the initial load, for the caller to report +/// once the event store exists. +pub(crate) struct StartupConflicts { + pub(crate) ports: Vec, + pub(crate) routes: Vec, +} + /// Scan every stack for `(host, path)` pairs claimed by more than one route, /// including two routes within the same stack. pub(crate) fn detect_route_conflicts(routes: &RouteMap) -> Vec { @@ -291,7 +298,12 @@ impl ServicesToml { /// or `(host, path)` route pair is claimed by more than one service/route /// — both are global on the proxy, unlike service names which only need /// to be unique within a stack. - pub(crate) async fn load_validated() -> Result<(StackMap, MatchIndex, RouteMap), Error> { + /// + /// The conflicts are returned rather than reported here: the event store + /// doesn't exist yet this early in startup, so the caller emits them once + /// the orchestrator is up (the reload path emits its own directly). + pub(crate) async fn load_validated() + -> Result<(StackMap, MatchIndex, RouteMap, StartupConflicts), Error> { let (mut stacks, mut index, mut routes) = Self::load().await?; // Don't brick the control plane on a conflict (e.g. a bad UI edit or // hand-edited file left two stacks claiming the same listen_port/route). @@ -324,7 +336,15 @@ impl ServicesToml { index.remove(stack); routes.remove(stack); } - Ok((stacks, index, routes)) + Ok(( + stacks, + index, + routes, + StartupConflicts { + ports: conflicts, + routes: route_conflicts, + }, + )) } /// `config_changed`, `port_mappings_changed`, and `http_routes_changed` @@ -418,7 +438,15 @@ impl ServicesToml { } } } - Err(e) => eprintln!("Failed to reload services.toml: {e:?}"), + // Unparseable file: the previous config stays in force, + // which looks identical to "nothing changed" from the UI. + Err(e) => { + eprintln!("Failed to reload services.toml: {e:?}"); + orchestrator + .events + .emit(ServerEvent::config_reload_failed(format!("{e:?}"))) + .await; + } } last_update_time = Instant::now(); } diff --git a/members/nullnet-server/ui/src/pages/Events.tsx b/members/nullnet-server/ui/src/pages/Events.tsx index d3452c1..b07464e 100644 --- a/members/nullnet-server/ui/src/pages/Events.tsx +++ b/members/nullnet-server/ui/src/pages/Events.tsx @@ -34,6 +34,10 @@ const KIND_LABELS: Record = { net_id_pool_exhausted: 'net_id_pool_exhausted', proxy_chain_setup_failed: 'proxy_chain_setup_failed', backend_trigger_setup_bailed: 'backend_trigger_setup_bailed', + udp_port_pool_exhausted: 'udp_port_pool_exhausted', + config_reload_failed: 'config_reload_failed', + file_watch_failed: 'file_watch_failed', + port_mapping_conflict: 'port_mapping_conflict', // Client error vxlan_setup_failed: 'vxlan_setup_failed', vlan_setup_failed: 'vlan_setup_failed', @@ -51,6 +55,13 @@ const KIND_LABELS: Record = { firewall_rules_load_failed: 'firewall_rules_load_failed', container_suspend_failed: 'container_suspend_failed', container_resume_failed: 'container_resume_failed', + backend_trigger_setup_timed_out: 'backend_trigger_setup_timed_out', + egress_steer_setup_timed_out: 'egress_steer_setup_timed_out', + egress_steer_install_failed: 'egress_steer_install_failed', + nfqueue_bind_failed: 'nfqueue_bind_failed', + mss_clamp_install_failed: 'mss_clamp_install_failed', + egress_policy_check_failed: 'egress_policy_check_failed', + conntrack_flush_failed: 'conntrack_flush_failed', // Client info vxlan_setup_completed: 'vxlan_setup_completed', vlan_setup_completed: 'vlan_setup_completed', @@ -63,12 +74,20 @@ const KIND_LABELS: Record = { upstream_ip_parse_failed: 'upstream_ip_parse_failed', proxy_client_not_inet: 'proxy_client_not_inet', tls_certificate_invalid: 'tls_certificate_invalid', + tcp_listener_bind_failed: 'tcp_listener_bind_failed', + udp_listener_bind_failed: 'udp_listener_bind_failed', + tcp_upstream_connect_failed: 'tcp_upstream_connect_failed', + udp_upstream_connect_failed: 'udp_upstream_connect_failed', // Proxy info proxy_request_routed: 'proxy_request_routed', + proxy_connected: 'proxy_connected', + proxy_disconnected: 'proxy_disconnected', // Certificate certificate_installed: 'certificate_installed', certificate_renewed: 'certificate_renewed', certificate_removed: 'certificate_removed', + certificate_renewal_failed: 'certificate_renewal_failed', + certificate_credentials_store_failed: 'certificate_credentials_store_failed', }; const ALL_KINDS = Object.keys(KIND_LABELS); @@ -112,10 +131,17 @@ function eventDetail(e: EventJson): string { case 'max_networks_limit_enforced': return `${e.service} · proxy ${e.proxy_ip} · net ${e.net_id} · limit ${e.limit}`; case 'net_id_pool_exhausted': + case 'udp_port_pool_exhausted': case 'proxy_chain_setup_failed': return `${e.service} · ${e.client_ip}`; case 'backend_trigger_setup_bailed': return `${e.service} · port ${e.port}`; + case 'config_reload_failed': + return e.error_message; + case 'file_watch_failed': + return `${e.target} · ${e.error_message}`; + case 'port_mapping_conflict': + return `${e.protocol}/${e.listen_port} · ${e.stack_a}/${e.service_a} vs ${e.stack_b}/${e.service_b}`; // Client error case 'vxlan_setup_failed': case 'vxlan_teardown_failed': @@ -146,6 +172,20 @@ function eventDetail(e: EventJson): string { case 'container_suspend_failed': case 'container_resume_failed': return `${e.docker_container} · ${e.error_message}`; + case 'backend_trigger_setup_timed_out': + return `${e.service_name}:${e.port} · ${e.docker_container} · ${e.error_message}`; + case 'egress_steer_setup_timed_out': + return `${e.docker_container} → ${e.dst_ip}:${e.dst_port} · ${e.error_message}`; + case 'egress_steer_install_failed': + return `vxlan ${e.vxlan_id} · ${e.docker_container ?? '—'} · ${e.error_message}`; + case 'nfqueue_bind_failed': + return `queue ${e.queue_id} · ${e.error_message}`; + case 'mss_clamp_install_failed': + return e.error_message; + case 'egress_policy_check_failed': + return `${e.docker_container} → ${e.dst_ip} · ${e.error_message}`; + case 'conntrack_flush_failed': + return `${e.ip} · ${e.error_message}`; // Client info case 'vxlan_setup_completed': return `vxlan ${e.vxlan_id} · ${e.ns_name}`; @@ -167,14 +207,26 @@ function eventDetail(e: EventJson): string { return e.address_family; case 'tls_certificate_invalid': return `${e.domain} · ${e.reason}`; + case 'tcp_listener_bind_failed': + case 'udp_listener_bind_failed': + return `:${e.listen_port} · ${e.service_name} · ${e.error_message}`; + case 'tcp_upstream_connect_failed': + case 'udp_upstream_connect_failed': + return `${e.service_name} · ${e.client_ip} · ${e.error_message}`; // Certificate events case 'certificate_installed': case 'certificate_renewed': case 'certificate_removed': return e.domain; + case 'certificate_renewal_failed': + case 'certificate_credentials_store_failed': + return `${e.domain} · ${e.error_message}`; // Proxy info case 'proxy_request_routed': return `${e.service_name} · ${e.client_ip} → ${e.upstream_ip} · ${e.latency_ms}ms`; + case 'proxy_connected': + case 'proxy_disconnected': + return e.ip; } } diff --git a/members/nullnet-server/ui/src/types.ts b/members/nullnet-server/ui/src/types.ts index d3f6c9c..574397a 100644 --- a/members/nullnet-server/ui/src/types.ts +++ b/members/nullnet-server/ui/src/types.ts @@ -99,6 +99,10 @@ export type EventJson = | WithSeverity & { type: 'net_id_pool_exhausted'; service: string; client_ip: string } | WithSeverity & { type: 'proxy_chain_setup_failed'; service: string; client_ip: string } | WithSeverity & { type: 'backend_trigger_setup_bailed'; service: string; port: number } + | WithSeverity & { type: 'udp_port_pool_exhausted'; service: string; client_ip: string } + | WithSeverity & { type: 'config_reload_failed'; error_message: string } + | WithSeverity & { type: 'file_watch_failed'; target: string; error_message: string } + | WithSeverity & { type: 'port_mapping_conflict'; stack_a: string; service_a: string; stack_b: string; service_b: string; protocol: string; listen_port: number } // Client error events | WithSeverity & { type: 'vxlan_setup_failed'; vxlan_id: number; ns_name: string; error_code: number } | WithSeverity & { type: 'vlan_setup_failed'; vlan_id: number; local_veth: string; error_reason: string } @@ -116,6 +120,13 @@ export type EventJson = | WithSeverity & { type: 'firewall_rules_load_failed'; path: string; error_message: string } | WithSeverity & { type: 'container_suspend_failed'; docker_container: string; error_message: string } | WithSeverity & { type: 'container_resume_failed'; docker_container: string; error_message: string } + | WithSeverity & { type: 'backend_trigger_setup_timed_out'; service_name: string; port: number; docker_container: string; error_message: string } + | WithSeverity & { type: 'egress_steer_setup_timed_out'; docker_container: string; dst_ip: string; dst_port: number; error_message: string } + | WithSeverity & { type: 'egress_steer_install_failed'; vxlan_id: number; docker_container?: string; error_message: string } + | WithSeverity & { type: 'nfqueue_bind_failed'; queue_id: number; error_message: string } + | WithSeverity & { type: 'mss_clamp_install_failed'; error_message: string } + | WithSeverity & { type: 'egress_policy_check_failed'; docker_container: string; dst_ip: string; error_message: string } + | WithSeverity & { type: 'conntrack_flush_failed'; ip: string; error_message: string } // Client info events | WithSeverity & { type: 'vxlan_setup_completed'; vxlan_id: number; ns_name: string } | WithSeverity & { type: 'vlan_setup_completed'; vlan_id: number } @@ -128,12 +139,20 @@ export type EventJson = | WithSeverity & { type: 'upstream_ip_parse_failed'; raw_ip: string; service_name: string } | WithSeverity & { type: 'proxy_client_not_inet'; address_family: string } | WithSeverity & { type: 'tls_certificate_invalid'; domain: string; reason: string } + | WithSeverity & { type: 'tcp_listener_bind_failed'; listen_port: number; service_name: string; error_message: string } + | WithSeverity & { type: 'udp_listener_bind_failed'; listen_port: number; service_name: string; error_message: string } + | WithSeverity & { type: 'tcp_upstream_connect_failed'; service_name: string; client_ip: string; error_message: string } + | WithSeverity & { type: 'udp_upstream_connect_failed'; service_name: string; client_ip: string; error_message: string } // Proxy info events | WithSeverity & { type: 'proxy_request_routed'; service_name: string; client_ip: string; upstream_ip: string; latency_ms: number } + | WithSeverity & { type: 'proxy_connected'; ip: string } + | WithSeverity & { type: 'proxy_disconnected'; ip: string } // Certificate events | WithSeverity & { type: 'certificate_installed'; domain: string } | WithSeverity & { type: 'certificate_renewed'; domain: string } - | WithSeverity & { type: 'certificate_removed'; domain: string }; + | WithSeverity & { type: 'certificate_removed'; domain: string } + | WithSeverity & { type: 'certificate_renewal_failed'; domain: string; error_message: string } + | WithSeverity & { type: 'certificate_credentials_store_failed'; domain: string; error_message: string }; export interface GraphNodeJson { id: string;