use aether_contracts::tunnel_security::TUNNEL_SECURITY_NON_TLS_REQUIRED; use async_trait::async_trait; use serde_json::Value; const PROXY_NODE_BOUND_TUNNEL_SECRET_PREFIX: &str = "aether-proxy-node-secret-v2:aether-runtime-secret-v1:"; #[derive(Clone, PartialEq, serde::Serialize, serde::Deserialize)] pub struct StoredProxyNode { pub id: String, #[serde(default = "new_proxy_node_tunnel_generation")] pub tunnel_generation: String, pub name: String, pub ip: String, pub port: i32, pub region: Option, pub is_manual: bool, pub proxy_url: Option, pub proxy_username: Option, pub proxy_password: Option, pub status: String, pub registered_by: Option, pub last_heartbeat_at_unix_secs: Option, pub heartbeat_interval: i32, pub active_connections: i32, pub total_requests: i64, pub avg_latency_ms: Option, pub failed_requests: i64, pub dns_failures: i64, pub stream_errors: i64, pub proxy_metadata: Option, pub hardware_info: Option, pub estimated_max_concurrency: Option, pub tunnel_mode: bool, pub tunnel_connected: bool, pub tunnel_connected_at_unix_secs: Option, pub remote_config: Option, pub config_version: i32, pub created_at_unix_ms: Option, pub updated_at_unix_secs: Option, } impl std::fmt::Debug for StoredProxyNode { fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { formatter .debug_struct("StoredProxyNode") .field("id", &self.id) .field("tunnel_generation", &self.tunnel_generation) .field("name", &self.name) .field("ip", &self.ip) .field("port", &self.port) .field("region", &self.region) .field("is_manual", &self.is_manual) .field("proxy_url", &self.proxy_url.as_ref().map(|_| "[REDACTED]")) .field("proxy_username", &self.proxy_username) .field( "proxy_password", &self.proxy_password.as_ref().map(|_| "[REDACTED]"), ) .field("status", &self.status) .field("tunnel_mode", &self.tunnel_mode) .field("tunnel_connected", &self.tunnel_connected) .field( "proxy_metadata", &self.proxy_metadata.as_ref().map(|_| "[REDACTED]"), ) .field( "hardware_info", &self.hardware_info.as_ref().map(|_| "[REDACTED]"), ) .field( "remote_config", &self.remote_config.as_ref().map(|_| "[REDACTED]"), ) .finish_non_exhaustive() } } impl StoredProxyNode { #[allow(clippy::too_many_arguments)] pub fn new( id: String, name: String, ip: String, port: i32, is_manual: bool, status: String, heartbeat_interval: i32, active_connections: i32, total_requests: i64, failed_requests: i64, dns_failures: i64, stream_errors: i64, tunnel_mode: bool, tunnel_connected: bool, config_version: i32, ) -> Result { if id.trim().is_empty() { return Err(crate::DataLayerError::UnexpectedValue( "proxy_nodes.id is empty".to_string(), )); } if name.trim().is_empty() { return Err(crate::DataLayerError::UnexpectedValue( "proxy_nodes.name is empty".to_string(), )); } if ip.trim().is_empty() { return Err(crate::DataLayerError::UnexpectedValue( "proxy_nodes.ip is empty".to_string(), )); } if status.trim().is_empty() { return Err(crate::DataLayerError::UnexpectedValue( "proxy_nodes.status is empty".to_string(), )); } Ok(Self { id, tunnel_generation: new_proxy_node_tunnel_generation(), name, ip, port, region: None, is_manual, proxy_url: None, proxy_username: None, proxy_password: None, status, registered_by: None, last_heartbeat_at_unix_secs: None, heartbeat_interval, active_connections, total_requests, avg_latency_ms: None, failed_requests, dns_failures, stream_errors, proxy_metadata: None, hardware_info: None, estimated_max_concurrency: None, tunnel_mode, tunnel_connected, tunnel_connected_at_unix_secs: None, remote_config: None, config_version, created_at_unix_ms: None, updated_at_unix_secs: None, }) } #[allow(clippy::too_many_arguments)] pub fn with_runtime_fields( mut self, region: Option, registered_by: Option, last_heartbeat_at_unix_secs: Option, avg_latency_ms: Option, proxy_metadata: Option, hardware_info: Option, estimated_max_concurrency: Option, tunnel_connected_at_unix_secs: Option, remote_config: Option, created_at_unix_ms: Option, updated_at_unix_secs: Option, ) -> Self { self.region = region; self.registered_by = registered_by; self.last_heartbeat_at_unix_secs = last_heartbeat_at_unix_secs; self.avg_latency_ms = avg_latency_ms; self.proxy_metadata = proxy_metadata; self.hardware_info = hardware_info; self.estimated_max_concurrency = estimated_max_concurrency; self.tunnel_connected_at_unix_secs = tunnel_connected_at_unix_secs; self.remote_config = remote_config; self.created_at_unix_ms = created_at_unix_ms; self.updated_at_unix_secs = updated_at_unix_secs; self } pub fn with_manual_proxy_fields( mut self, proxy_url: Option, proxy_username: Option, proxy_password: Option, ) -> Self { self.proxy_url = proxy_url; self.proxy_username = proxy_username; self.proxy_password = proxy_password; self } pub fn with_tunnel_generation(mut self, tunnel_generation: String) -> Self { self.tunnel_generation = tunnel_generation; self } } pub fn new_proxy_node_tunnel_generation() -> String { uuid::Uuid::new_v4().to_string() } #[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] pub struct ProxyNodeHeartbeatMutation { pub node_id: String, #[serde(default)] pub expected_tunnel_generation: Option, pub heartbeat_interval: Option, pub active_connections: Option, pub total_requests_delta: Option, pub avg_latency_ms: Option, pub failed_requests_delta: Option, pub dns_failures_delta: Option, pub stream_errors_delta: Option, pub proxy_metadata: Option, pub proxy_version: Option, } #[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] pub struct ProxyNodeTrafficMutation { pub node_id: String, /// Incarnation fence captured when the request plan selected this node. /// Missing fences are rejected by the gateway path so a stale plan cannot /// update a node recreated under the same id. #[serde(default)] pub expected_tunnel_generation: Option, pub total_requests_delta: i64, pub failed_requests_delta: i64, pub dns_failures_delta: i64, pub stream_errors_delta: i64, } #[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] pub struct ProxyNodeRegistrationMutation { /// The stable identity selected by the caller before any secret is /// protected. Re-registration of an existing endpoint must use its /// existing id; repositories reject attempts to replace it. #[serde(default)] pub node_id: Option, pub name: String, pub ip: String, pub port: i32, pub region: Option, pub heartbeat_interval: i32, pub active_connections: Option, pub total_requests: Option, pub avg_latency_ms: Option, pub hardware_info: Option, pub estimated_max_concurrency: Option, pub proxy_metadata: Option, pub proxy_version: Option, pub registered_by: Option, pub tunnel_mode: bool, } #[derive(Clone, PartialEq, serde::Serialize, serde::Deserialize)] pub struct ProxyNodeManualCreateMutation { /// Optional caller-selected id used to bind credentials before the row is /// inserted. Repositories generate one only for legacy callers that do /// not provide it. #[serde(default)] pub node_id: Option, pub name: String, pub ip: String, pub port: i32, pub region: Option, pub proxy_url: String, pub proxy_username: Option, pub proxy_password: Option, pub registered_by: Option, } impl std::fmt::Debug for ProxyNodeManualCreateMutation { fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { formatter .debug_struct("ProxyNodeManualCreateMutation") .field("node_id", &self.node_id) .field("name", &self.name) .field("ip", &self.ip) .field("port", &self.port) .field("region", &self.region) .field("proxy_url", &"[REDACTED]") .field("proxy_username", &self.proxy_username) .field( "proxy_password", &self.proxy_password.as_ref().map(|_| "[REDACTED]"), ) .field("registered_by", &self.registered_by) .finish() } } #[derive(Clone, PartialEq, serde::Serialize, serde::Deserialize)] pub struct ProxyNodeManualUpdateMutation { pub node_id: String, pub name: Option, pub ip: Option, pub port: Option, pub region: Option, pub proxy_url: Option, pub proxy_username: Option, pub proxy_password: Option, } impl std::fmt::Debug for ProxyNodeManualUpdateMutation { fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { formatter .debug_struct("ProxyNodeManualUpdateMutation") .field("node_id", &self.node_id) .field("name", &self.name) .field("ip", &self.ip) .field("port", &self.port) .field("region", &self.region) .field("proxy_url", &self.proxy_url.as_ref().map(|_| "[REDACTED]")) .field("proxy_username", &self.proxy_username) .field( "proxy_password", &self.proxy_password.as_ref().map(|_| "[REDACTED]"), ) .finish() } } #[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] pub struct ProxyNodeTunnelStatusMutation { pub node_id: String, #[serde(default)] pub expected_tunnel_generation: Option, pub connected: bool, pub conn_count: i32, pub detail: Option, pub observed_at_unix_secs: Option, } #[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] pub struct ProxyNodeRemoteConfigMutation { pub node_id: String, #[serde(default)] pub expected_tunnel_generation: Option, pub node_name: Option, pub allowed_ports: Option>, pub log_level: Option, pub heartbeat_interval: Option, pub scheduling_state: Option>, pub upgrade_to: Option>, } #[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] pub struct StoredProxyNodeEvent { pub id: i64, pub node_id: String, pub event_type: String, pub detail: Option, pub event_metadata: Option, pub created_at_unix_ms: Option, } #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] pub struct ProxyNodeEventQuery { pub limit: usize, pub from_unix_secs: Option, pub to_unix_secs: Option, pub event_type: Option, } #[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] pub enum ProxyNodeMetricsStep { OneMinute, OneHour, } impl ProxyNodeMetricsStep { pub fn bucket_size_secs(self) -> u64 { match self { Self::OneMinute => 60, Self::OneHour => 3_600, } } pub fn as_api_value(self) -> &'static str { match self { Self::OneMinute => "1m", Self::OneHour => "1h", } } } #[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] pub struct StoredProxyNodeMetricsBucket { pub node_id: String, pub bucket_start_unix_secs: u64, pub samples: i64, pub uptime_samples: i64, pub active_connections_sum: i64, pub active_connections_max: i64, pub heartbeat_rtt_ms_sum: i64, pub heartbeat_rtt_ms_max: i64, pub connect_errors_delta: i64, pub disconnects_delta: i64, pub error_events_delta: i64, pub ws_in_bytes_delta: i64, pub ws_out_bytes_delta: i64, pub ws_in_frames_delta: i64, pub ws_out_frames_delta: i64, } #[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] pub struct StoredProxyFleetMetricsBucket { pub bucket_start_unix_secs: u64, pub samples: i64, pub uptime_samples: i64, pub active_connections_sum: i64, pub active_connections_max: i64, pub heartbeat_rtt_ms_sum: i64, pub heartbeat_rtt_ms_max: i64, pub connect_errors_delta: i64, pub disconnects_delta: i64, pub error_events_delta: i64, pub ws_in_bytes_delta: i64, pub ws_out_bytes_delta: i64, pub ws_in_frames_delta: i64, pub ws_out_frames_delta: i64, } #[derive(Debug, Clone, Copy, PartialEq, Eq, Default, serde::Serialize, serde::Deserialize)] pub struct ProxyNodeMetricsCleanupSummary { pub deleted_1m_rows: usize, pub deleted_1h_rows: usize, } pub const PROXY_NODE_EVENT_TYPE_TUNNEL_ERROR: &str = "tunnel_err"; #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] pub struct TunnelErrorEventRecord { pub timestamp_unix_secs: u64, pub timestamp_unix_ms: Option, pub category: String, pub message: String, pub severity: Option, pub component: Option, pub summary: Option, pub operator_action: Option, } #[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] pub struct TunnelMetricsCounters { pub connect_errors: u64, pub disconnects: u64, pub error_events_total: u64, pub ws_in_bytes: u64, pub ws_out_bytes: u64, pub ws_in_frames: u64, pub ws_out_frames: u64, pub heartbeat_rtt_last_ms: u64, } #[derive(Debug, Clone, PartialEq, Eq)] pub struct TunnelMetricsSample { pub samples: i64, pub uptime_samples: i64, pub active_connections_sum: i64, pub active_connections_max: i64, pub heartbeat_rtt_ms_sum: i64, pub heartbeat_rtt_ms_max: i64, pub connect_errors_delta: i64, pub disconnects_delta: i64, pub error_events_delta: i64, pub ws_in_bytes_delta: i64, pub ws_out_bytes_delta: i64, pub ws_in_frames_delta: i64, pub ws_out_frames_delta: i64, pub recent_error_events: Vec, } pub fn bucket_start_unix_secs(timestamp_unix_secs: u64, step: ProxyNodeMetricsStep) -> u64 { let size = step.bucket_size_secs(); timestamp_unix_secs / size * size } pub fn build_tunnel_metrics_sample( previous_proxy_metadata: Option<&Value>, current_proxy_metadata: Option<&Value>, active_connections: i32, tunnel_connected: bool, ) -> Option { let current = extract_tunnel_metrics_counters(current_proxy_metadata)?; let previous = extract_tunnel_metrics_counters(previous_proxy_metadata); let current_recent_errors = extract_recent_tunnel_errors(current_proxy_metadata); let connect_errors_delta = counter_delta_u64(previous.map(|v| v.connect_errors), current.connect_errors); let disconnects_delta = counter_delta_u64(previous.map(|v| v.disconnects), current.disconnects); let error_events_delta = counter_delta_u64( previous.map(|v| v.error_events_total), current.error_events_total, ); let ws_in_bytes_delta = counter_delta_u64(previous.map(|v| v.ws_in_bytes), current.ws_in_bytes); let ws_out_bytes_delta = counter_delta_u64(previous.map(|v| v.ws_out_bytes), current.ws_out_bytes); let ws_in_frames_delta = counter_delta_u64(previous.map(|v| v.ws_in_frames), current.ws_in_frames); let ws_out_frames_delta = counter_delta_u64(previous.map(|v| v.ws_out_frames), current.ws_out_frames); let take_recent = usize::try_from(error_events_delta).unwrap_or(usize::MAX); let recent_error_events = if take_recent == 0 { Vec::new() } else { let capture = take_recent.min(current_recent_errors.len()); let from = current_recent_errors.len().saturating_sub(capture); current_recent_errors[from..].to_vec() }; let active_connections = i64::from(active_connections.max(0)); let heartbeat_rtt_last_ms = i64::try_from(current.heartbeat_rtt_last_ms).unwrap_or(i64::MAX); Some(TunnelMetricsSample { samples: 1, uptime_samples: if tunnel_connected { 1 } else { 0 }, active_connections_sum: active_connections, active_connections_max: active_connections, heartbeat_rtt_ms_sum: heartbeat_rtt_last_ms, heartbeat_rtt_ms_max: heartbeat_rtt_last_ms, connect_errors_delta: i64::try_from(connect_errors_delta).unwrap_or(i64::MAX), disconnects_delta: i64::try_from(disconnects_delta).unwrap_or(i64::MAX), error_events_delta: i64::try_from(error_events_delta).unwrap_or(i64::MAX), ws_in_bytes_delta: i64::try_from(ws_in_bytes_delta).unwrap_or(i64::MAX), ws_out_bytes_delta: i64::try_from(ws_out_bytes_delta).unwrap_or(i64::MAX), ws_in_frames_delta: i64::try_from(ws_in_frames_delta).unwrap_or(i64::MAX), ws_out_frames_delta: i64::try_from(ws_out_frames_delta).unwrap_or(i64::MAX), recent_error_events, }) } pub fn build_tunnel_error_event_detail(event: &TunnelErrorEventRecord) -> String { format!( "[{}] {}", event.category, event.summary.as_deref().unwrap_or(event.message.as_str()) ) } pub fn normalize_proxy_metadata( proxy_metadata: Option<&serde_json::Value>, proxy_version: Option<&str>, ) -> Option { let mut normalized = match proxy_metadata { Some(serde_json::Value::Object(map)) => map.clone(), Some(_) | None => serde_json::Map::new(), }; let raw_version = normalized .remove("version") .and_then(|value| value.as_str().map(str::to_string)); let version = proxy_version .map(str::trim) .filter(|value| !value.is_empty()) .map(|value| value.chars().take(20).collect::()) .or_else(|| { raw_version .as_deref() .map(str::trim) .filter(|value| !value.is_empty()) .map(|value| value.chars().take(20).collect::()) }); if let Some(version) = version { normalized.insert("version".to_string(), serde_json::Value::String(version)); } if normalized.is_empty() { None } else { Some(serde_json::Value::Object(normalized)) } } pub fn normalize_heartbeat_proxy_metadata( previous_proxy_metadata: Option<&Value>, proxy_metadata: Option<&Value>, proxy_version: Option<&str>, ) -> Option { let Some(Value::Object(mut normalized)) = normalize_proxy_metadata(proxy_metadata, proxy_version) else { return None; }; // Tunnel security is control-plane state. A heartbeat may refresh runtime // metadata, but it must never introduce or replace this trusted field. normalized.remove("tunnel_security"); let merged = preserve_proxy_metadata_tunnel_security( previous_proxy_metadata, Some(Value::Object(normalized)), ); merged.filter(|value| { value .as_object() .is_some_and(|metadata| !metadata.is_empty()) }) } pub fn preserve_proxy_metadata_tunnel_security( previous_proxy_metadata: Option<&Value>, next_proxy_metadata: Option, ) -> Option { let Some(tunnel_security) = previous_proxy_metadata .and_then(|value| value.get("tunnel_security")) .filter(|value| value.is_object()) .cloned() else { return next_proxy_metadata; }; match next_proxy_metadata { Some(Value::Object(mut metadata)) => { metadata.insert("tunnel_security".to_string(), tunnel_security); Some(Value::Object(metadata)) } Some(value) => Some(value), None => { let mut metadata = serde_json::Map::new(); metadata.insert("tunnel_security".to_string(), tunnel_security); Some(Value::Object(metadata)) } } } /// Merge metadata received during a trusted registration/re-registration. /// /// Registration is the control-plane path that may rotate a tunnel PSK. A /// registration payload that omits `tunnel_security` is therefore a partial /// metadata refresh and must not clear the previously trusted security state. /// Only a non-empty, gateway-bound v2 ciphertext proves that the registration /// passed through the gateway credential-binding path, and that ciphertext is /// accepted only with the required non-TLS security mode. Mode-only, /// plaintext, malformed, null, scalar, empty, and disabled security values are /// treated as omission and cannot clear a previously trusted object. pub fn merge_proxy_metadata_for_registration( previous_proxy_metadata: Option<&Value>, next_proxy_metadata: Option, ) -> Option { let Some(next_proxy_metadata) = next_proxy_metadata else { return previous_proxy_metadata.cloned(); }; let incoming_security_is_explicit = proxy_metadata_has_explicit_tunnel_security(Some(&next_proxy_metadata)); let Value::Object(mut metadata) = next_proxy_metadata else { // `normalize_proxy_metadata` normally prevents this branch. Keep a // malformed replacement from erasing trusted control-plane state. return previous_proxy_metadata.cloned(); }; if incoming_security_is_explicit { return Some(Value::Object(metadata)); } // Null, scalar, and empty security objects are not valid replacements. // Remove them before restoring the previous trusted object so malformed // input cannot mask or downgrade the registered security policy. metadata.remove("tunnel_security"); if let Some(previous_security) = previous_proxy_metadata .and_then(|value| value.get("tunnel_security")) .filter(|value| value.is_object()) .cloned() { metadata.insert("tunnel_security".to_string(), previous_security); } (!metadata.is_empty()).then_some(Value::Object(metadata)) } /// Return whether metadata contains a complete gateway-bound tunnel security /// replacement that a trusted registration may persist. pub fn proxy_metadata_has_explicit_tunnel_security(proxy_metadata: Option<&Value>) -> bool { proxy_metadata .and_then(Value::as_object) .and_then(|metadata| metadata.get("tunnel_security")) .and_then(Value::as_object) .is_some_and(|security| { security.get("mode").and_then(Value::as_str) == Some(TUNNEL_SECURITY_NON_TLS_REQUIRED) && security .get("encryption_key_encrypted") .and_then(Value::as_str) .and_then(|encrypted| { encrypted.strip_prefix(PROXY_NODE_BOUND_TUNNEL_SECRET_PREFIX) }) .is_some_and(|ciphertext| !ciphertext.is_empty()) }) } fn extract_tunnel_metrics_counters( proxy_metadata: Option<&Value>, ) -> Option { let tunnel_metrics = proxy_metadata .and_then(Value::as_object) .and_then(|metadata| metadata.get("tunnel_metrics")) .and_then(Value::as_object)?; Some(TunnelMetricsCounters { connect_errors: json_u64(tunnel_metrics.get("connect_errors")).unwrap_or(0), disconnects: json_u64(tunnel_metrics.get("disconnects")).unwrap_or(0), error_events_total: json_u64(tunnel_metrics.get("error_events_total")).unwrap_or(0), ws_in_bytes: json_u64(tunnel_metrics.get("ws_in_bytes")).unwrap_or(0), ws_out_bytes: json_u64(tunnel_metrics.get("ws_out_bytes")).unwrap_or(0), ws_in_frames: json_u64(tunnel_metrics.get("ws_in_frames")).unwrap_or(0), ws_out_frames: json_u64(tunnel_metrics.get("ws_out_frames")).unwrap_or(0), heartbeat_rtt_last_ms: json_u64(tunnel_metrics.get("heartbeat_rtt_last_ms")).unwrap_or(0), }) } fn extract_recent_tunnel_errors(proxy_metadata: Option<&Value>) -> Vec { proxy_metadata .and_then(Value::as_object) .and_then(|metadata| metadata.get("recent_tunnel_errors")) .and_then(Value::as_array) .map(|items| { items .iter() .filter_map(|item| { let item = item.as_object()?; Some(TunnelErrorEventRecord { timestamp_unix_secs: json_u64(item.get("timestamp_unix_secs")) .unwrap_or_default(), timestamp_unix_ms: json_u64(item.get("timestamp_unix_ms")), category: item .get("category") .and_then(Value::as_str) .unwrap_or("unknown") .to_string(), message: item .get("message") .and_then(Value::as_str) .unwrap_or("n/a") .to_string(), severity: json_string(item.get("severity")), component: json_string(item.get("component")), summary: json_string(item.get("summary")), operator_action: json_string(item.get("operator_action")), }) }) .collect::>() }) .unwrap_or_default() } fn json_u64(value: Option<&Value>) -> Option { value.and_then(|value| { value .as_u64() .or_else(|| value.as_i64().and_then(|n| (n >= 0).then_some(n as u64))) }) } fn json_string(value: Option<&Value>) -> Option { value .and_then(Value::as_str) .map(str::trim) .filter(|value| !value.is_empty()) .map(ToString::to_string) } fn counter_delta_u64(previous: Option, current: u64) -> u64 { match previous { Some(previous) if current >= previous => current - previous, Some(_) => current, None => 0, } } fn normalize_proxy_version_label(value: &str) -> Option { let trimmed = value.trim(); if trimmed.is_empty() { return None; } Some( trimmed .strip_prefix("tunnel-v") .or_else(|| trimmed.strip_prefix("proxy-v")) .unwrap_or(trimmed) .to_ascii_lowercase(), ) } pub const PROXY_NODE_SCHEDULING_STATE_DRAINING: &str = "draining"; pub const PROXY_NODE_SCHEDULING_STATE_CORDONED: &str = "cordoned"; pub fn normalize_proxy_node_scheduling_state(value: &str) -> Option<&'static str> { let trimmed = value.trim(); if trimmed.eq_ignore_ascii_case(PROXY_NODE_SCHEDULING_STATE_DRAINING) { return Some(PROXY_NODE_SCHEDULING_STATE_DRAINING); } if trimmed.eq_ignore_ascii_case(PROXY_NODE_SCHEDULING_STATE_CORDONED) { return Some(PROXY_NODE_SCHEDULING_STATE_CORDONED); } None } pub fn remote_config_scheduling_state( remote_config: Option<&serde_json::Value>, ) -> Option<&'static str> { remote_config .and_then(serde_json::Value::as_object) .and_then(|value| value.get("scheduling_state")) .and_then(serde_json::Value::as_str) .and_then(normalize_proxy_node_scheduling_state) } pub fn proxy_node_accepts_new_tunnels(node: &StoredProxyNode) -> bool { remote_config_scheduling_state(node.remote_config.as_ref()).is_none() } pub fn proxy_reported_version(proxy_metadata: Option<&serde_json::Value>) -> Option { proxy_metadata .and_then(serde_json::Value::as_object) .and_then(|value| value.get("version")) .and_then(serde_json::Value::as_str) .and_then(normalize_proxy_version_label) } pub fn remote_config_upgrade_target(remote_config: Option<&serde_json::Value>) -> Option { remote_config .and_then(serde_json::Value::as_object) .and_then(|value| value.get("upgrade_to")) .and_then(serde_json::Value::as_str) .and_then(normalize_proxy_version_label) } pub fn reconcile_remote_config_after_heartbeat( remote_config: Option<&serde_json::Value>, proxy_version: Option<&str>, ) -> Option { let Some(mut config) = remote_config .and_then(serde_json::Value::as_object) .cloned() else { return remote_config.cloned(); }; let Some(target_version) = config .get("upgrade_to") .and_then(serde_json::Value::as_str) .and_then(normalize_proxy_version_label) else { return Some(serde_json::Value::Object(config)); }; let Some(reported_version) = proxy_version.and_then(normalize_proxy_version_label) else { return Some(serde_json::Value::Object(config)); }; if reported_version == target_version { config.remove("upgrade_to"); } (!config.is_empty()).then_some(serde_json::Value::Object(config)) } #[async_trait] pub trait ProxyNodeReadRepository: Send + Sync { async fn list_proxy_nodes(&self) -> Result, crate::DataLayerError>; async fn find_proxy_node( &self, node_id: &str, ) -> Result, crate::DataLayerError>; async fn list_proxy_node_events( &self, node_id: &str, limit: usize, ) -> Result, crate::DataLayerError>; async fn list_proxy_node_events_filtered( &self, node_id: &str, query: &ProxyNodeEventQuery, ) -> Result, crate::DataLayerError> { let mut items = self.list_proxy_node_events(node_id, query.limit).await?; if let Some(from_unix_secs) = query.from_unix_secs { items.retain(|item| item.created_at_unix_ms.unwrap_or(0) >= from_unix_secs); } if let Some(to_unix_secs) = query.to_unix_secs { items.retain(|item| item.created_at_unix_ms.unwrap_or(u64::MAX) <= to_unix_secs); } if let Some(event_type) = query.event_type.as_deref() { items.retain(|item| item.event_type.eq_ignore_ascii_case(event_type)); } items.truncate(query.limit); Ok(items) } async fn list_proxy_node_metrics( &self, node_id: &str, step: ProxyNodeMetricsStep, from_unix_secs: u64, to_unix_secs: u64, limit: usize, ) -> Result, crate::DataLayerError>; async fn list_proxy_fleet_metrics( &self, step: ProxyNodeMetricsStep, from_unix_secs: u64, to_unix_secs: u64, limit: usize, ) -> Result, crate::DataLayerError>; } #[async_trait] pub trait ProxyNodeWriteRepository: Send + Sync { async fn reset_stale_tunnel_statuses(&self) -> Result; async fn compare_and_set_proxy_password( &self, node_id: &str, expected: &str, replacement: &str, ) -> Result; async fn compare_and_set_proxy_metadata( &self, node_id: &str, expected: &serde_json::Value, replacement: &serde_json::Value, ) -> Result; async fn create_manual_node( &self, mutation: &ProxyNodeManualCreateMutation, ) -> Result; async fn update_manual_node( &self, mutation: &ProxyNodeManualUpdateMutation, ) -> Result, crate::DataLayerError>; async fn register_node( &self, mutation: &ProxyNodeRegistrationMutation, ) -> Result; async fn apply_heartbeat( &self, mutation: &ProxyNodeHeartbeatMutation, ) -> Result, crate::DataLayerError>; async fn record_traffic( &self, mutation: &ProxyNodeTrafficMutation, ) -> Result; async fn update_tunnel_status( &self, mutation: &ProxyNodeTunnelStatusMutation, ) -> Result, crate::DataLayerError>; async fn unregister_node( &self, node_id: &str, ) -> Result, crate::DataLayerError>; async fn delete_node( &self, node_id: &str, ) -> Result, crate::DataLayerError>; async fn update_remote_config( &self, mutation: &ProxyNodeRemoteConfigMutation, ) -> Result, crate::DataLayerError>; async fn increment_manual_node_requests( &self, node_id: &str, total_delta: i64, failed_delta: i64, latency_ms: Option, ) -> Result<(), crate::DataLayerError>; async fn cleanup_proxy_node_metrics( &self, retain_1m_from_unix_secs: u64, retain_1h_from_unix_secs: u64, delete_limit: usize, ) -> Result; } #[cfg(test)] mod tests { use serde_json::json; use super::{ bucket_start_unix_secs, build_tunnel_error_event_detail, build_tunnel_metrics_sample, merge_proxy_metadata_for_registration, normalize_heartbeat_proxy_metadata, normalize_proxy_node_scheduling_state, preserve_proxy_metadata_tunnel_security, proxy_node_accepts_new_tunnels, proxy_reported_version, reconcile_remote_config_after_heartbeat, remote_config_scheduling_state, remote_config_upgrade_target, ProxyNodeManualCreateMutation, ProxyNodeManualUpdateMutation, ProxyNodeMetricsStep, StoredProxyNode, }; #[test] fn proxy_node_debug_output_redacts_credentials_and_untrusted_metadata() { let password = "debug-secret-proxy-password"; let proxy_url = "https://user:debug-secret-url@example.com"; let metadata_secret = "debug-secret-proxy-metadata"; let mut stored = StoredProxyNode::new( "node-1".to_string(), "node".to_string(), "127.0.0.1".to_string(), 8080, true, "online".to_string(), 30, 0, 0, 0, 0, 0, false, false, 1, ) .expect("proxy node should build") .with_manual_proxy_fields( Some(proxy_url.to_string()), Some("user".to_string()), Some(password.to_string()), ); stored.proxy_metadata = Some(json!({"secret": metadata_secret})); let create = ProxyNodeManualCreateMutation { node_id: Some("node-1".to_string()), name: "node".to_string(), ip: "127.0.0.1".to_string(), port: 8080, region: None, proxy_url: proxy_url.to_string(), proxy_username: Some("user".to_string()), proxy_password: Some(password.to_string()), registered_by: None, }; let update = ProxyNodeManualUpdateMutation { node_id: "node-1".to_string(), name: None, ip: None, port: None, region: None, proxy_url: Some(proxy_url.to_string()), proxy_username: Some("user".to_string()), proxy_password: Some(password.to_string()), }; for rendered in [ format!("{stored:?}"), format!("{create:?}"), format!("{update:?}"), ] { for secret in [password, proxy_url, metadata_secret] { assert!(!rendered.contains(secret), "Debug output leaked {secret}"); } assert!(rendered.contains("[REDACTED]")); } } #[test] fn normalizes_reported_versions_and_clears_completed_upgrade_targets() { let remote_config = json!({ "node_name": "edge-1", "upgrade_to": "tunnel-v2.0.0", }); let proxy_metadata = json!({ "version": "2.0.0", "arch": "arm64", }); assert_eq!( proxy_reported_version(Some(&proxy_metadata)).as_deref(), Some("2.0.0") ); assert_eq!( remote_config_upgrade_target(Some(&remote_config)).as_deref(), Some("2.0.0") ); let reconciled = reconcile_remote_config_after_heartbeat(Some(&remote_config), Some("tunnel-v2.0.0")) .expect("reconciled config should remain an object"); assert_eq!(reconciled.get("upgrade_to"), None); assert_eq!(reconciled.get("node_name"), Some(&json!("edge-1"))); } #[test] fn normalizes_proxy_node_scheduling_state_and_detects_unschedulable_nodes() { assert_eq!( normalize_proxy_node_scheduling_state("draining"), Some("draining") ); assert_eq!( normalize_proxy_node_scheduling_state(" CORDONED "), Some("cordoned") ); assert_eq!(normalize_proxy_node_scheduling_state("active"), None); let remote_config = json!({ "node_name": "edge-1", "scheduling_state": "draining", }); assert_eq!( remote_config_scheduling_state(Some(&remote_config)), Some("draining") ); let node = StoredProxyNode::new( "node-1".to_string(), "edge-1".to_string(), "127.0.0.1".to_string(), 0, false, "online".to_string(), 30, 0, 0, 0, 0, 0, true, true, 0, ) .expect("node should build") .with_runtime_fields( None, None, Some(1_800_000_000), None, None, None, None, Some(1_800_000_001), Some(remote_config), Some(1_800_000_000), Some(1_800_000_001), ); assert!(!proxy_node_accepts_new_tunnels(&node)); } #[test] fn builds_tunnel_metrics_sample_with_reset_safe_counter_deltas() { let previous = json!({ "tunnel_metrics": { "connect_errors": 10, "disconnects": 5, "error_events_total": 7, "ws_in_bytes": 1_000, "ws_out_bytes": 2_000, "ws_in_frames": 10, "ws_out_frames": 20, "heartbeat_rtt_last_ms": 30 } }); let current = json!({ "tunnel_metrics": { "connect_errors": 12, "disconnects": 2, "error_events_total": 9, "ws_in_bytes": 1_500, "ws_out_bytes": 100, "ws_in_frames": 11, "ws_out_frames": 3, "heartbeat_rtt_last_ms": 44 }, "recent_tunnel_errors": [ {"timestamp_unix_secs": 100, "category": "older", "message": "old"}, { "timestamp_unix_secs": 101, "timestamp_unix_ms": 101_999, "category": "newer", "message": "new", "severity": "error", "component": "tunnel_write", "summary": "WebSocket write failed because the peer closed or reset the connection", "operator_action": "Check gateway restarts and network resets." } ] }); let sample = build_tunnel_metrics_sample(Some(&previous), Some(¤t), 4, true) .expect("sample should build"); assert_eq!(sample.samples, 1); assert_eq!(sample.uptime_samples, 1); assert_eq!(sample.active_connections_sum, 4); assert_eq!(sample.heartbeat_rtt_ms_sum, 44); assert_eq!(sample.connect_errors_delta, 2); assert_eq!(sample.disconnects_delta, 2); assert_eq!(sample.error_events_delta, 2); assert_eq!(sample.ws_in_bytes_delta, 500); assert_eq!(sample.ws_out_bytes_delta, 100); assert_eq!(sample.ws_out_frames_delta, 3); assert_eq!(sample.recent_error_events.len(), 2); assert_eq!(sample.recent_error_events[0].category, "older"); assert_eq!(sample.recent_error_events[0].summary, None); assert_eq!( sample.recent_error_events[1].severity.as_deref(), Some("error") ); assert_eq!( sample.recent_error_events[1].component.as_deref(), Some("tunnel_write") ); assert_eq!( sample.recent_error_events[1].timestamp_unix_ms, Some(101_999) ); assert_eq!( build_tunnel_error_event_detail(&sample.recent_error_events[1]), "[newer] WebSocket write failed because the peer closed or reset the connection" ); } #[test] fn builds_tunnel_metrics_sample_uses_first_counter_report_as_baseline() { let current = json!({ "tunnel_metrics": { "connect_errors": 12, "disconnects": 5, "error_events_total": 7, "ws_in_bytes": 1_500, "ws_out_bytes": 2_500, "ws_in_frames": 15, "ws_out_frames": 25, "heartbeat_rtt_last_ms": 44 }, "recent_tunnel_errors": [ {"timestamp_unix_secs": 101, "category": "newer", "message": "new"} ] }); let sample = build_tunnel_metrics_sample(None, Some(¤t), 4, true) .expect("sample should build"); assert_eq!(sample.samples, 1); assert_eq!(sample.heartbeat_rtt_ms_sum, 44); assert_eq!(sample.connect_errors_delta, 0); assert_eq!(sample.disconnects_delta, 0); assert_eq!(sample.error_events_delta, 0); assert_eq!(sample.ws_in_bytes_delta, 0); assert_eq!(sample.ws_out_bytes_delta, 0); assert_eq!(sample.ws_in_frames_delta, 0); assert_eq!(sample.ws_out_frames_delta, 0); assert!(sample.recent_error_events.is_empty()); } #[test] fn preserves_secure_tunnel_metadata_across_heartbeat_metadata_refresh() { let previous = json!({ "version": "1.0.0", "tunnel_security": { "mode": "non_tls_required", "encryption_key": "BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc=" } }); let next = json!({ "version": "1.0.1", "tunnel_metrics": {"connect_successes": 1}, "tunnel_security": { "mode": "disabled", "encryption_key": "attacker-controlled" } }); let merged = preserve_proxy_metadata_tunnel_security(Some(&previous), Some(next)) .expect("metadata should remain present"); assert_eq!( merged .pointer("/tunnel_security/mode") .and_then(|v| v.as_str()), Some("non_tls_required") ); assert_eq!( merged .pointer("/tunnel_security/encryption_key") .and_then(|v| v.as_str()), Some("BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc=") ); assert_eq!( merged.pointer("/tunnel_metrics/connect_successes"), Some(&json!(1)) ); } #[test] fn registration_metadata_preserves_omitted_tunnel_security() { let previous = json!({ "version": "1.0.0", "tunnel_security": { "mode": "non_tls_required", "encryption_key_encrypted": "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-old" } }); let next = json!({ "version": "1.1.0", "tunnel_metrics": {"connect_successes": 2} }); let merged = merge_proxy_metadata_for_registration(Some(&previous), Some(next)) .expect("registration metadata should remain present"); assert_eq!(merged.get("version"), Some(&json!("1.1.0"))); assert_eq!( merged.pointer("/tunnel_security/encryption_key_encrypted"), Some(&json!( "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-old" )) ); assert_eq!( merged.pointer("/tunnel_metrics/connect_successes"), Some(&json!(2)) ); } #[test] fn registration_metadata_accepts_explicit_tunnel_security_rotation() { let previous = json!({ "tunnel_security": { "mode": "non_tls_required", "encryption_key_encrypted": "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-old" } }); let next = json!({ "tunnel_security": { "mode": "non_tls_required", "encryption_key_encrypted": "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-new" } }); let merged = merge_proxy_metadata_for_registration(Some(&previous), Some(next)) .expect("rotated registration metadata should remain present"); assert_eq!( merged.pointer("/tunnel_security/encryption_key_encrypted"), Some(&json!( "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-new" )) ); } #[test] fn registration_metadata_rejects_invalid_security_replacement() { let previous = json!({ "tunnel_security": { "mode": "non_tls_required", "encryption_key_encrypted": "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-old" } }); let next = json!({ "version": "1.2.0", "tunnel_security": null }); let merged = merge_proxy_metadata_for_registration(Some(&previous), Some(next)) .expect("previous security should be retained"); assert_eq!(merged.get("version"), Some(&json!("1.2.0"))); assert_eq!( merged.pointer("/tunnel_security/encryption_key_encrypted"), Some(&json!( "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-old" )) ); } #[test] fn registration_metadata_rejects_mode_only_security_downgrade() { let previous = json!({ "tunnel_security": { "mode": "non_tls_required", "encryption_key_encrypted": "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-old" } }); let next = json!({ "version": "1.3.0", "tunnel_security": {"mode": "disabled"} }); let merged = merge_proxy_metadata_for_registration(Some(&previous), Some(next)) .expect("mode-only security must not replace the registered credential"); assert_eq!(merged.get("version"), Some(&json!("1.3.0"))); assert_eq!( merged.pointer("/tunnel_security/encryption_key_encrypted"), Some(&json!( "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-old" )) ); assert_eq!( merged.pointer("/tunnel_security/mode"), Some(&json!("non_tls_required")) ); } #[test] fn registration_metadata_rejects_bound_ciphertext_with_disabled_mode() { let previous = json!({ "tunnel_security": { "mode": "non_tls_required", "encryption_key_encrypted": "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-old" } }); let next = json!({ "version": "1.3.1", "tunnel_security": { "mode": "disabled", "encryption_key_encrypted": "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-attacker" } }); let merged = merge_proxy_metadata_for_registration(Some(&previous), Some(next)) .expect("disabled security must not replace the registered credential"); assert_eq!(merged.get("version"), Some(&json!("1.3.1"))); assert_eq!( merged.pointer("/tunnel_security/encryption_key_encrypted"), Some(&json!( "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-old" )) ); assert_eq!( merged.pointer("/tunnel_security/mode"), Some(&json!("non_tls_required")) ); } #[test] fn new_registration_drops_invalid_tunnel_security_fields() { for invalid_security in [ json!(null), json!("disabled"), json!({}), json!({"mode": "disabled"}), json!({"mode": "disabled", "encryption_key_encrypted": "aether-proxy-node-secret-v2:aether-runtime-secret-v1:sealed-attacker"}), json!({"mode": "non_tls_required", "encryption_key_encrypted": ""}), json!({"mode": "non_tls_required", "encryption_key_encrypted": "not-gateway-bound"}), ] { let next = json!({ "version": "1.4.0", "tunnel_security": invalid_security }); let merged = merge_proxy_metadata_for_registration(None, Some(next)) .expect("valid non-security metadata should remain"); assert_eq!(merged.get("version"), Some(&json!("1.4.0"))); assert!(merged.get("tunnel_security").is_none()); } assert_eq!( merge_proxy_metadata_for_registration( None, Some(json!({"tunnel_security": {"mode": "disabled"}})), ), None ); } #[test] fn heartbeat_metadata_cannot_introduce_tunnel_security() { let injected = json!({ "tunnel_security": { "mode": "disabled", "encryption_key": "attacker-controlled" } }); assert_eq!( normalize_heartbeat_proxy_metadata(None, Some(&injected), None), None ); let injected_with_runtime_metadata = json!({ "version": "1.2.3", "arch": "arm64", "tunnel_security": { "mode": "disabled", "encryption_key": "attacker-controlled" } }); let normalized = normalize_heartbeat_proxy_metadata(None, Some(&injected_with_runtime_metadata), None) .expect("runtime metadata should remain present"); assert_eq!(normalized.get("version"), Some(&json!("1.2.3"))); assert_eq!(normalized.get("arch"), Some(&json!("arm64"))); assert!(normalized.get("tunnel_security").is_none()); } #[test] fn maps_timestamps_to_metric_buckets() { assert_eq!( bucket_start_unix_secs(1_710_000_119, ProxyNodeMetricsStep::OneMinute), 1_710_000_060 ); assert_eq!( bucket_start_unix_secs(1_710_003_999, ProxyNodeMetricsStep::OneHour), 1_710_003_600 ); } }