diff --git a/src/cli/init.rs b/src/cli/init.rs index f53fa5f..43282c3 100644 --- a/src/cli/init.rs +++ b/src/cli/init.rs @@ -3,6 +3,8 @@ use std::process::Command; use rand::RngExt; +use crate::util::trusted_command::resolve_trusted_helper; + /// Options for the fire-and-forget init command. #[derive(Debug, Clone)] pub struct InitOptions { @@ -165,11 +167,14 @@ pub fn run_init(opts: InitOptions) -> Result<(), Box> { eprintln!("[+] Service started"); std::thread::sleep(std::time::Duration::from_secs(1)); - let status = Command::new("systemctl") - .args(["is-active", "telemt.service"]) - .output(); + let status = resolve_trusted_helper("systemctl").and_then(|command_path| { + Command::new(command_path) + .args(["is-active", "telemt.service"]) + .output() + .ok() + }); match status { - Ok(out) if out.status.success() => { + Some(out) if out.status.success() => { eprintln!("[+] Service is running"); } _ => { @@ -329,7 +334,11 @@ weight = 10 } fn run_cmd(cmd: &str, args: &[&str]) { - match Command::new(cmd).args(args).output() { + let Some(command_path) = resolve_trusted_helper(cmd) else { + eprintln!("[!] Refusing unavailable or untrusted command: {}", cmd); + return; + }; + match Command::new(command_path).args(args).output() { Ok(output) => { if !output.status.success() { let stderr = String::from_utf8_lossy(&output.stderr); diff --git a/src/config/load/runtime_auth.rs b/src/config/load/runtime_auth.rs index 8dc3964..c7e5ec1 100644 --- a/src/config/load/runtime_auth.rs +++ b/src/config/load/runtime_auth.rs @@ -117,12 +117,14 @@ impl UserAuthSnapshot { self.entries.get(idx) } + /// Returns the stable credential identity for an exact configured username. pub(crate) fn credential_id_by_name(&self, user: &str) -> Option<[u8; 16]> { self.user_id_by_name(user) .and_then(|user_id| self.entry_by_id(user_id)) .map(|entry| entry.credential_id) } + /// Returns every bounded authentication candidate sharing a stable hint key. pub(crate) fn candidate_ids_by_hint_key(&self, hint_key: u64) -> Option<&[u32]> { self.by_hint_key.get(&hint_key).map(Vec::as_slice) } diff --git a/src/config/load/runtime_web/static_site_fallback.rs b/src/config/load/runtime_web/static_site_fallback.rs index 6a1d667..d248895 100644 --- a/src/config/load/runtime_web/static_site_fallback.rs +++ b/src/config/load/runtime_web/static_site_fallback.rs @@ -4,6 +4,7 @@ use std::path::Path; use super::*; +/// Builds a bounded static-site snapshot on platforms without directory descriptors. pub(super) fn load_static_site_by_path( root: &Path, limits: &WebLimitsConfig, diff --git a/src/conntrack_control/firewall.rs b/src/conntrack_control/firewall.rs index 73110cf..5632594 100644 --- a/src/conntrack_control/firewall.rs +++ b/src/conntrack_control/firewall.rs @@ -13,6 +13,7 @@ use crate::util::trusted_command::resolve_trusted_helper; use super::{ConntrackRuntimeSupport, NetfilterBackend}; +/// Reconciles kernel NOTRACK rules with the active listener policy. pub(super) async fn reconcile_rules( cfg: &ProxyConfig, runtime_support: ConntrackRuntimeSupport, @@ -45,6 +46,7 @@ pub(super) async fn reconcile_rules( } } +/// Probes the effective firewall backend and conntrack deletion capability. pub(super) fn probe_runtime_support( configured_backend: ConntrackBackend, ) -> ConntrackRuntimeSupport { @@ -55,6 +57,7 @@ pub(super) fn probe_runtime_support( } } +/// Resolves whether conntrack close publication is usable for this runtime. pub(super) fn effective_conntrack_enabled( cfg: &ProxyConfig, runtime_support: ConntrackRuntimeSupport, @@ -317,12 +320,17 @@ async fn clear_notrack_rules_all_backends() { let _ = run_command("ip6tables", &["-t", "raw", "-X", "TELEMT_NOTRACK"], None).await; } +/// Result of one best-effort kernel conntrack deletion. pub(super) enum DeleteOutcome { + /// The kernel reported successful deletion. Deleted, + /// No matching conntrack entry existed. NotFound, + /// The helper was unavailable or returned an unexpected failure. Error, } +/// Deletes the exact TCP tuple represented by one close event. pub(super) async fn delete_conntrack_entry(event: ConntrackCloseEvent) -> DeleteOutcome { if !command_exists("conntrack") { return DeleteOutcome::Error; diff --git a/src/daemon/pid_file.rs b/src/daemon/pid_file.rs index bc0acdd..b886576 100644 --- a/src/daemon/pid_file.rs +++ b/src/daemon/pid_file.rs @@ -1,7 +1,7 @@ use std::ffi::OsStr; use std::fs::{self, File}; use std::io::{self, ErrorKind, Read, Write}; -use std::os::unix::fs::MetadataExt; +use std::os::unix::fs::{MetadataExt, PermissionsExt}; #[cfg(target_os = "linux")] use std::os::fd::{FromRawFd, OwnedFd}; use std::path::{Path, PathBuf}; @@ -42,7 +42,7 @@ impl FileIdentity { impl PidFile { /// Creates a new PID file manager for the given path. pub fn new>(path: P) -> Self { - let path = path.as_ref().to_path_buf(); + let path = normalize_pid_path(path.as_ref()); let lock_path = sibling_lock_path(&path); Self { path, @@ -209,6 +209,33 @@ fn sibling_lock_path(path: &Path) -> PathBuf { lock_path.into() } +fn normalize_pid_path(path: &Path) -> PathBuf { + let legacy_run = Path::new("/var/run"); + let Ok(remainder) = path.strip_prefix(legacy_run) else { + return path.to_path_buf(); + }; + let Ok(var_metadata) = fs::metadata("/var") else { + return path.to_path_buf(); + }; + let Ok(link_metadata) = fs::symlink_metadata(legacy_run) else { + return path.to_path_buf(); + }; + let Ok(target) = fs::read_link(legacy_run) else { + return path.to_path_buf(); + }; + let trusted_var = var_metadata.is_dir() + && var_metadata.uid() == 0 + && var_metadata.permissions().mode() & 0o022 == 0; + let trusted_alias = link_metadata.file_type().is_symlink() + && link_metadata.uid() == 0 + && (target == Path::new("/run") || target == Path::new("../run")); + if trusted_var && trusted_alias { + Path::new("/run").join(remainder) + } else { + path.to_path_buf() + } +} + fn open_file_at( anchor: &AnchoredPath, name: &OsStr, @@ -343,8 +370,8 @@ fn validate_regular_single_link( /// Reads a PID from a PID file. #[allow(dead_code)] pub fn read_pid_file>(path: P) -> Result { - let path = path.as_ref(); - read_pid_file_if_exists(path)?.ok_or_else(|| { + let path = normalize_pid_path(path.as_ref()); + read_pid_file_if_exists(&path)?.ok_or_else(|| { DaemonError::PidFile(format!( "cannot read {}: file does not exist", path.display() @@ -358,11 +385,11 @@ pub fn signal_pid_file>( path: P, signal: nix::sys::signal::Signal, ) -> Result<(), DaemonError> { - let path = path.as_ref(); - let pid = read_pid_file(path)?; + let path = normalize_pid_path(path.as_ref()); + let pid = read_pid_file(&path)?; #[cfg(target_os = "linux")] let pidfd = open_pidfd(pid)?; - if !daemon_lock_is_held(path)? { + if !daemon_lock_is_held(&path)? { return Err(DaemonError::PidFile(format!( "refusing to signal unlocked or stale PID file {}", path.display() @@ -390,10 +417,10 @@ pub enum DaemonStatus { /// Checks daemon status without modifying the PID or lock file. #[allow(dead_code)] pub fn check_status>(path: P) -> DaemonStatus { - let path = path.as_ref(); - match read_pid_file_if_exists(path) { + let path = normalize_pid_path(path.as_ref()); + match read_pid_file_if_exists(&path) { Ok(Some(pid)) - if daemon_lock_is_held(path).unwrap_or(false) && is_process_running(pid) => + if daemon_lock_is_held(&path).unwrap_or(false) && is_process_running(pid) => { DaemonStatus::Running(pid) } diff --git a/src/daemon/pid_file/tests.rs b/src/daemon/pid_file/tests.rs index db23065..2203eca 100644 --- a/src/daemon/pid_file/tests.rs +++ b/src/daemon/pid_file/tests.rs @@ -38,6 +38,25 @@ fn pid_file_remains_send_and_sync() { assert_send_sync::(); } +#[test] +fn system_var_run_alias_keeps_the_default_pid_path_usable() { + let Ok(metadata) = fs::symlink_metadata("/var/run") else { + return; + }; + let Ok(target) = fs::read_link("/var/run") else { + return; + }; + if !metadata.file_type().is_symlink() + || (target != Path::new("/run") && target != Path::new("../run")) + { + return; + } + + let pid_file = PidFile::new("/var/run/telemt.pid"); + + assert_eq!(pid_file.path(), Path::new("/run/telemt.pid")); +} + #[test] fn lock_holder_subprocess() { let Some(pid_path) = std::env::var_os(HELPER_PID_PATH) else { diff --git a/src/maestro/helpers/runtime.rs b/src/maestro/helpers/runtime.rs index 245650f..725cb61 100644 --- a/src/maestro/helpers/runtime.rs +++ b/src/maestro/helpers/runtime.rs @@ -12,6 +12,7 @@ use crate::transport::middle_proxy::{ use super::print_maestro_line; +/// Prints configured MTProxy links through the direct MAESTRO output channel. pub(crate) fn print_proxy_links(host: &str, port: u16, config: &ProxyConfig) { print_maestro_line(format!("Proxy links ({host})")); for user_name in config @@ -94,6 +95,7 @@ pub(crate) fn print_web_proxy_links(config: &ProxyConfig) { } } +/// Durably replaces one Beobachten snapshot without following Unix symlinks. pub(crate) async fn write_beobachten_snapshot(path: &str, payload: &str) -> std::io::Result<()> { #[cfg(unix)] { @@ -115,10 +117,12 @@ pub(crate) async fn write_beobachten_snapshot(path: &str, payload: &str) -> std: } } +/// Selects a singular or plural display label for one integer value. pub(crate) fn unit_label(value: u64, singular: &'static str, plural: &'static str) -> &'static str { if value == 1 { singular } else { plural } } +/// Formats process uptime into bounded human-readable units and exact seconds. pub(crate) fn format_uptime(total_secs: u64) -> String { const SECS_PER_MINUTE: u64 = 60; const SECS_PER_HOUR: u64 = 60 * SECS_PER_MINUTE; @@ -172,6 +176,7 @@ pub(crate) fn format_uptime(total_secs: u64) -> String { } #[allow(dead_code)] +/// Waits until admission opens or its watch channel closes. pub(crate) async fn wait_until_admission_open(admission_rx: &mut watch::Receiver) -> bool { loop { if *admission_rx.borrow() { @@ -183,10 +188,12 @@ pub(crate) async fn wait_until_admission_open(admission_rx: &mut watch::Receiver } } +/// Classifies peer closure that is expected during an incomplete handshake. pub(crate) fn is_expected_handshake_eof(err: &crate::error::ProxyError) -> bool { expected_handshake_close_description(err).is_some() } +/// Returns a stable diagnostic description for transport-level peer closure. pub(crate) fn peer_close_description(err: &crate::error::ProxyError) -> Option<&'static str> { fn from_kind(kind: std::io::ErrorKind) -> Option<&'static str> { match kind { @@ -209,6 +216,7 @@ pub(crate) fn peer_close_description(err: &crate::error::ProxyError) -> Option<& } } +/// Returns a stable diagnostic description for expected handshake closure. pub(crate) fn expected_handshake_close_description( err: &crate::error::ProxyError, ) -> Option<&'static str> { @@ -243,6 +251,7 @@ pub(crate) fn expected_handshake_close_description( } } +/// Loads a non-empty startup endpoint snapshot with bounded cache fallback. pub(crate) async fn load_startup_proxy_config_snapshot( url: &str, cache_path: Option<&str>, diff --git a/src/metrics/render/me_hardswap.rs b/src/metrics/render/me_hardswap.rs index e3efe21..e8d19d8 100644 --- a/src/metrics/render/me_hardswap.rs +++ b/src/metrics/render/me_hardswap.rs @@ -2,6 +2,7 @@ use std::fmt::Write; use crate::transport::middle_proxy::MeApiHardswapSnapshot; +/// Renders fixed-cardinality hardswap and writer-replacement gauges. pub(super) fn render( out: &mut String, snapshot: Option<&MeApiHardswapSnapshot>, diff --git a/src/proxy/direct_relay.rs b/src/proxy/direct_relay.rs index ef65298..c22b209 100644 --- a/src/proxy/direct_relay.rs +++ b/src/proxy/direct_relay.rs @@ -1,5 +1,6 @@ use std::collections::HashSet; use std::ffi::OsString; +#[cfg(all(test, unix))] use std::fs::OpenOptions; use std::io::Write; use std::net::SocketAddr; @@ -33,7 +34,7 @@ use nix::fcntl::{Flock, FlockArg, OFlag, openat}; #[cfg(unix)] use nix::sys::stat::Mode; -#[cfg(unix)] +#[cfg(all(test, unix))] use std::os::unix::fs::OpenOptionsExt; // Direct relay lifecycle and conntrack publication. @@ -178,10 +179,7 @@ fn open_unknown_dc_log_append_anchored( ) -> std::io::Result { #[cfg(unix)] { - let parent = OpenOptions::new() - .read(true) - .custom_flags(libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC) - .open(&path.allowed_parent)?; + let parent = crate::util::secure_fs::open_dir_nofollow(&path.allowed_parent)?; let oflags = OFlag::O_CREAT | OFlag::O_APPEND diff --git a/src/proxy/handshake/auth_probe/testing.rs b/src/proxy/handshake/auth_probe/testing.rs index 18329f4..03b26ca 100644 --- a/src/proxy/handshake/auth_probe/testing.rs +++ b/src/proxy/handshake/auth_probe/testing.rs @@ -1,5 +1,6 @@ use super::*; +/// Records one deterministic authentication failure against an isolated shared state. pub(crate) fn auth_probe_record_failure_for_testing( shared: &ProxySharedState, peer_ip: IpAddr, @@ -8,6 +9,7 @@ pub(crate) fn auth_probe_record_failure_for_testing( auth_probe_record_failure_in(shared, peer_ip, now); } +/// Returns the normalized peer failure streak from an isolated shared state. pub(crate) fn auth_probe_fail_streak_for_testing_in_shared( shared: &ProxySharedState, peer_ip: IpAddr, @@ -20,6 +22,7 @@ pub(crate) fn auth_probe_fail_streak_for_testing_in_shared( .map(|entry| entry.fail_streak) } +/// Clears probe entries, exact capacity accounting, and saturation state together. pub(crate) fn clear_auth_probe_state_for_testing_in_shared(shared: &ProxySharedState) { let removed = shared.handshake.auth_probe.len(); assert_eq!(shared.handshake.auth_probe_slots.used(), removed); @@ -37,6 +40,7 @@ pub(crate) fn clear_auth_probe_state_for_testing_in_shared(shared: &ProxySharedS } } +/// Inserts one fixture entry while preserving exact registry capacity accounting. pub(crate) fn insert_auth_probe_state_for_testing_in_shared( shared: &ProxySharedState, peer_ip: IpAddr, @@ -59,22 +63,26 @@ pub(crate) fn insert_auth_probe_state_for_testing_in_shared( } } +/// Exposes the isolated probe registry to adversarial tests. pub(crate) fn auth_probe_state_for_testing_in_shared( shared: &ProxySharedState, ) -> &DashMap { &shared.handshake.auth_probe } +/// Returns exact committed probe slots for capacity assertions. pub(crate) fn auth_probe_slots_for_testing_in_shared(shared: &ProxySharedState) -> usize { shared.handshake.auth_probe_slots.used() } +/// Exposes the isolated saturation state mutex to tests. pub(crate) fn auth_probe_saturation_state_for_testing_in_shared( shared: &ProxySharedState, ) -> &Mutex> { &shared.handshake.auth_probe_saturation } +/// Locks isolated saturation state while recovering poisoned test fixtures. pub(crate) fn auth_probe_saturation_state_lock_for_testing_in_shared( shared: &ProxySharedState, ) -> std::sync::MutexGuard<'_, Option> { @@ -85,6 +93,7 @@ pub(crate) fn auth_probe_saturation_state_lock_for_testing_in_shared( .unwrap_or_else(|poisoned| poisoned.into_inner()) } +/// Resets the isolated unknown-SNI warning rate limiter. pub(crate) fn clear_unknown_sni_warn_state_for_testing_in_shared(shared: &ProxySharedState) { let mut guard = shared .handshake @@ -94,6 +103,7 @@ pub(crate) fn clear_unknown_sni_warn_state_for_testing_in_shared(shared: &ProxyS *guard = None; } +/// Evaluates unknown-SNI warning admission at a deterministic instant. pub(crate) fn should_emit_unknown_sni_warn_for_testing_in_shared( shared: &ProxySharedState, now: Instant, @@ -101,18 +111,21 @@ pub(crate) fn should_emit_unknown_sni_warn_for_testing_in_shared( should_emit_unknown_sni_warn_in(shared, now) } +/// Clears the isolated invalid-secret warning deduplication set. pub(crate) fn clear_warned_secrets_for_testing_in_shared(shared: &ProxySharedState) { if let Ok(mut guard) = shared.handshake.invalid_secret_warned.lock() { guard.clear(); } } +/// Exposes the isolated invalid-secret warning set to tests. pub(crate) fn warned_secrets_for_testing_in_shared( shared: &ProxySharedState, ) -> &Mutex> { &shared.handshake.invalid_secret_warned } +/// Evaluates peer throttling against the current test clock. pub(crate) fn auth_probe_is_throttled_for_testing_in_shared( shared: &ProxySharedState, peer_ip: IpAddr, @@ -120,12 +133,14 @@ pub(crate) fn auth_probe_is_throttled_for_testing_in_shared( auth_probe_is_throttled_in(shared, peer_ip, Instant::now()) } +/// Evaluates global saturation throttling against the current test clock. pub(crate) fn auth_probe_saturation_is_throttled_for_testing_in_shared( shared: &ProxySharedState, ) -> bool { auth_probe_saturation_is_throttled_in(shared, Instant::now()) } +/// Evaluates global saturation throttling at a deterministic instant. pub(crate) fn auth_probe_saturation_is_throttled_at_for_testing_in_shared( shared: &ProxySharedState, now: Instant, diff --git a/src/proxy/shared_state.rs b/src/proxy/shared_state.rs index 2749d7a..2cd9d6e 100644 --- a/src/proxy/shared_state.rs +++ b/src/proxy/shared_state.rs @@ -56,17 +56,25 @@ pub(crate) enum ConntrackClosePolicy { pub(crate) struct HandshakeSharedState { pub(crate) auth_probe: DashMap, + /// Exact capacity authority for the authentication probe registry. pub(crate) auth_probe_slots: SlotBudget, pub(crate) auth_probe_saturation: Mutex>, pub(crate) auth_probe_eviction_hasher: RandomState, pub(crate) invalid_secret_warned: Mutex>, pub(crate) unknown_sni_warn_next_allowed: Mutex>, + /// Stable credential hints keyed by exact peer IP. pub(crate) sticky_user_by_ip: DashMap, + /// Exact capacity authority for peer-IP credential hints. pub(crate) sticky_user_by_ip_slots: SlotBudget, + /// Stable credential hints keyed by bounded peer network prefix. pub(crate) sticky_user_by_ip_prefix: DashMap, + /// Exact capacity authority for peer-prefix credential hints. pub(crate) sticky_user_by_ip_prefix_slots: SlotBudget, + /// Stable credential hints keyed by normalized SNI hash. pub(crate) sticky_user_by_sni_hash: DashMap, + /// Exact capacity authority for SNI credential hints. pub(crate) sticky_user_by_sni_hash_slots: SlotBudget, + /// Bounded recent credential-hint ring used as an authentication fallback. pub(crate) recent_user_ring: Box<[AtomicU64]>, pub(crate) recent_user_ring_seq: AtomicU64, pub(crate) auth_expensive_checks_total: AtomicU64, diff --git a/src/proxy/tests/direct_relay_security_tests/anchored.rs b/src/proxy/tests/direct_relay_security_tests/anchored.rs index 6f0b242..02bc01e 100644 --- a/src/proxy/tests/direct_relay_security_tests/anchored.rs +++ b/src/proxy/tests/direct_relay_security_tests/anchored.rs @@ -72,6 +72,51 @@ fn adversarial_parent_swap_after_check_is_blocked_by_anchored_open() { ); } +#[cfg(unix)] +#[test] +fn adversarial_intermediate_parent_swap_is_blocked_by_component_walk() { + use std::os::unix::fs::symlink; + + let directory = tempfile::tempdir().expect("temporary directory must be creatable"); + let parent = directory.path().join("parent"); + let moved = directory.path().join("moved"); + let outside = directory.path().join("outside"); + fs::create_dir_all(parent.join("nested")) + .expect("original nested directory must be creatable"); + fs::create_dir_all(outside.join("nested")) + .expect("outside nested directory must be creatable"); + + let candidate = parent.join("nested/unknown-dc.log"); + let sanitized = sanitize_unknown_dc_log_path( + candidate + .to_str() + .expect("temporary path must be valid UTF-8"), + ) + .expect("candidate must sanitize before intermediate parent swap"); + assert!( + unknown_dc_log_path_is_still_safe(&sanitized), + "precondition: target should initially pass revalidation" + ); + + fs::rename(&parent, &moved).expect("intermediate parent must be movable"); + symlink(&outside, &parent).expect("intermediate parent symlink must be creatable"); + + let err = open_unknown_dc_log_append_anchored(&sanitized) + .expect_err("anchored open must reject a swapped intermediate component"); + let raw = err.raw_os_error(); + assert!( + matches!( + raw, + Some(libc::ELOOP) | Some(libc::ENOTDIR) | Some(libc::ENOENT) + ), + "component walk must fail closed on intermediate swap, got raw_os_error={raw:?}" + ); + assert!( + !outside.join("nested/unknown-dc.log").exists(), + "component walk must not create a log through a swapped intermediate directory" + ); +} + #[cfg(unix)] #[test] fn anchored_open_nix_path_writes_expected_lines() { diff --git a/src/proxy/tests/handshake_security_tests.rs b/src/proxy/tests/handshake_security_tests.rs index 010e56f..48a3f3d 100644 --- a/src/proxy/tests/handshake_security_tests.rs +++ b/src/proxy/tests/handshake_security_tests.rs @@ -1115,7 +1115,8 @@ async fn tls_unknown_sni_reject_handshake_policy_emits_unrecognized_name_alert() // Drain what the server wrote. We expect exactly one TLS alert record: // 0x15 0x03 0x03 0x00 0x02 0x02 0x70 // (ContentType.alert, TLS 1.2, length=2, fatal, unrecognized_name) - drop(result); // drops the server-side writer so peer_side sees EOF + // Drop the server-side writer so `peer_side` observes EOF. + drop(result); let mut buf = Vec::new(); peer_side.read_to_end(&mut buf).await.unwrap(); assert_eq!( diff --git a/src/transport/middle_proxy/pool.rs b/src/transport/middle_proxy/pool.rs index 8c551cc..d34fda7 100644 --- a/src/transport/middle_proxy/pool.rs +++ b/src/transport/middle_proxy/pool.rs @@ -266,16 +266,22 @@ pub struct RoutingCore { pub(super) writers: Arc, pub(super) rr: AtomicU64, pub(super) writer_epoch: watch::Sender, + /// Coherent immutable authority for endpoint maps and reverse indexes. pub(super) endpoint_snapshot: ArcSwap, } /// Immutable endpoint routing authority published as one coherent revision. #[derive(Clone, Debug)] pub(super) struct EndpointSnapshot { + /// Monotonic revision covering every endpoint-derived index in this snapshot. pub(super) revision: u64, + /// IPv4 endpoint map by Telegram DC. pub(super) map_v4: HashMap>, + /// IPv6 endpoint map by Telegram DC. pub(super) map_v6: HashMap>, + /// Reverse lookup from an endpoint to its optional Telegram DC. pub(super) endpoint_dc_map: HashMap>, + /// Ordered endpoint candidates used for per-DC writer selection. pub(super) preferred_endpoints_by_dc: HashMap>, } @@ -312,6 +318,7 @@ pub(super) struct ReinitPendingState { pub(super) generation: u64, pub(super) started_at_epoch_secs: u64, pub(super) map_hash: u64, + /// Endpoint authority revision targeted by the pending generation. pub(super) endpoint_revision: u64, } @@ -319,6 +326,7 @@ pub(super) struct ReinitPendingState { pub(super) struct ReinitAttemptState { pub(super) generation: u64, pub(super) map_hash: u64, + /// Endpoint authority revision captured by this attempt. pub(super) endpoint_revision: u64, pub(super) hardswap: bool, pub(super) committed: bool, @@ -328,6 +336,7 @@ pub(super) struct ReinitCoordinatorState { pub(super) next_attempt_id: u64, pub(super) active_generation: u64, pub(super) desired_map_hash: u64, + /// Latest endpoint authority revision accepted by the coordinator. pub(super) endpoint_revision: u64, pub(super) pending: Option, pub(super) attempts: HashMap, @@ -492,8 +501,10 @@ pub struct MePool { pub(super) next_writer_id: AtomicU64, pub(super) writer_connect_active_reserved: AtomicUsize, pub(super) writer_connect_warm_reserved: AtomicUsize, + /// Replacement connections opened but not yet committed to writer visibility. pub(super) writer_replacement_open_reserved: AtomicUsize, pub(super) rtt_stats: Arc>>, + /// Coalesced refill state keyed by exact generation and contour ownership. pub(super) refill_states: Arc>>, pub(super) refill_running: AtomicUsize, pub(super) refill_pending: AtomicUsize, diff --git a/src/transport/middle_proxy/pool/routing.rs b/src/transport/middle_proxy/pool/routing.rs index 50e9616..c183404 100644 --- a/src/transport/middle_proxy/pool/routing.rs +++ b/src/transport/middle_proxy/pool/routing.rs @@ -144,6 +144,7 @@ impl MePool { } } + /// Builds all endpoint-derived indexes under one immutable revision. pub(in crate::transport::middle_proxy) fn build_endpoint_snapshot( decision: &NetworkDecision, mut map_v4: HashMap>, @@ -272,6 +273,7 @@ impl MePool { endpoint_dc_map } + /// Removes runtime endpoint state absent from the current coherent snapshot. pub(in crate::transport::middle_proxy) async fn prune_endpoint_runtime_state(&self) { let configured_endpoints = self .endpoint_snapshot diff --git a/src/transport/middle_proxy/pool_reinit/coordination.rs b/src/transport/middle_proxy/pool_reinit/coordination.rs index 4b27b19..ab8e3c6 100644 --- a/src/transport/middle_proxy/pool_reinit/coordination.rs +++ b/src/transport/middle_proxy/pool_reinit/coordination.rs @@ -307,6 +307,7 @@ impl MePool { self.desired_dc_endpoints_from_snapshot(&endpoint_snapshot) } + /// Projects desired per-DC endpoint sets from one immutable endpoint revision. pub(super) fn desired_dc_endpoints_from_snapshot( &self, endpoint_snapshot: &EndpointSnapshot, diff --git a/src/transport/middle_proxy/pool_status/hardswap_snapshot.rs b/src/transport/middle_proxy/pool_status/hardswap_snapshot.rs index d332249..770791b 100644 --- a/src/transport/middle_proxy/pool_status/hardswap_snapshot.rs +++ b/src/transport/middle_proxy/pool_status/hardswap_snapshot.rs @@ -30,6 +30,7 @@ impl MePool { self.api_hardswap_snapshot_for_reinit(reinit.as_ref()).await } + /// Builds the bounded projection from one coherent reinitialization snapshot. pub(super) async fn api_hardswap_snapshot_for_reinit( &self, reinit: &ReinitStatusSnapshot, diff --git a/src/util/secure_fs/mod.rs b/src/util/secure_fs/mod.rs index e8641fb..916ee95 100644 --- a/src/util/secure_fs/mod.rs +++ b/src/util/secure_fs/mod.rs @@ -9,7 +9,7 @@ mod write; pub(crate) use path::{ AnchoredPath, chdir_nofollow_or_create, open_dir_nofollow, - open_dir_nofollow_or_create, open_trusted_dir_nofollow_or_create, + open_trusted_dir_nofollow_or_create, }; pub(crate) use write::{ atomic_replace, atomic_replace_async, open_append_regular, open_append_regular_at, diff --git a/src/util/secure_fs/path.rs b/src/util/secure_fs/path.rs index f395b44..b7e8183 100644 --- a/src/util/secure_fs/path.rs +++ b/src/util/secure_fs/path.rs @@ -182,10 +182,11 @@ fn validate_trusted_directory(descriptor: &OwnedFd, allow_sticky_parent: bool) - /// Creates missing components and changes cwd to the exact opened directory inode. pub(crate) fn chdir_nofollow_or_create(path: &Path, mode: u32) -> io::Result<()> { - let descriptor = open_or_create_dir_nofollow(path, mode)?; + let descriptor = open_dir_nofollow_or_create(path, mode)?; nix::unistd::fchdir(&descriptor).map_err(errno_to_io) } +/// Converts one `nix` errno without discarding its platform error code. pub(super) fn errno_to_io(error: nix::errno::Errno) -> io::Error { io::Error::from_raw_os_error(error as i32) } diff --git a/src/util/trusted_command.rs b/src/util/trusted_command.rs index 676350d..ab677be 100644 --- a/src/util/trusted_command.rs +++ b/src/util/trusted_command.rs @@ -2,7 +2,18 @@ use std::os::unix::fs::{MetadataExt, PermissionsExt}; use std::path::{Path, PathBuf}; const TRUSTED_HELPER_DIRS: [&str; 4] = ["/usr/sbin", "/usr/bin", "/sbin", "/bin"]; -const TRUSTED_HELPERS: [&str; 5] = ["nft", "iptables", "ip6tables", "conntrack", "pfctl"]; +const TRUSTED_HELPERS: [&str; 10] = [ + "nft", + "iptables", + "ip6tables", + "conntrack", + "pfctl", + "systemctl", + "rc-update", + "rc-service", + "sysrc", + "service", +]; /// Resolves a privileged helper only through the fixed system allowlist. pub(crate) fn resolve_trusted_helper(binary: &str) -> Option {