diff --git a/CHANGELOG.md b/CHANGELOG.md index 2a6dab00..3d23dd54 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,7 @@ All Nullnet releases with the relative changes are documented in this file. ## [UNRELEASED] ### Added +- Add an opt-in “Pausable when idle” checkbox for Docker services, persisted in configuration with a false default ([#187](https://github.com/NullNet-ai/nullnet/pull/187) — fixes [#180](https://github.com/NullNet-ai/nullnet/issues/180)) - Optional service host pinning and automatic proxy TCP/UDP listen-port firewall allowances ([#184](https://github.com/NullNet-ai/nullnet/pull/184) — fixes [#177](https://github.com/NullNet-ai/nullnet/issues/177)) - Per-service egress/ingress traffic filters: arbitrary AND/OR/group combinations of Country, Organization, Src IP (ingress), and Dst IP (egress) conditions, evaluated via `rpn-predicate-interpreter` — replaces the country-only egress/ingress policy ([#171](https://github.com/NullNet-ai/nullnet/pull/171) — fixes [#143](https://github.com/NullNet-ai/nullnet/issues/143)) - Persist ingress and egress sessions to SQLite and show the full history on the Sessions page, filterable by status, service, direction, and policy verdict, with its own retention window ([#170](https://github.com/NullNet-ai/nullnet/pull/170) — fixes [#156](https://github.com/NullNet-ai/nullnet/issues/156)) diff --git a/SETUP.md b/SETUP.md index 6bfd1949..b6575f47 100644 --- a/SETUP.md +++ b/SETUP.md @@ -220,6 +220,11 @@ The repository should be cloned under `/root` so the provided `setup-*.sh` scrip (non-Docker) service - `host_ip` optionally limits either match to one node's control-channel IPv4 address (for example, `host_ip = "192.168.1.103"` for that host's SSH service). Omit it to match all hosts. +- `pausable = true` opts a Docker service into pausing when idle (the Config checkbox). + It defaults to `false`, including existing DB rows. Backend services follow the same setting. + Live chains and egress sessions keep containers running; disabling pause starts an asynchronous resume of a paused replica. + If several declarations share a container, all must opt in. Paused initiators cannot start work + themselves; use this only when incoming traffic can wake them. - `timeout` controls proxy-reachability: when present the service is a proxy-reachable entry point with that per-client idle timeout in seconds (`0` disables the timeout); omit it to keep the service off the proxy (backend-only) diff --git a/members/nullnet-client/src/main.rs b/members/nullnet-client/src/main.rs index f88f48cc..de80112a 100644 --- a/members/nullnet-client/src/main.rs +++ b/members/nullnet-client/src/main.rs @@ -345,16 +345,7 @@ async fn declare_services( // Report raw local observations; the server joins them against its // per-stack config to decide what this node hosts. Running containers: // logical (Swarm label / name) -> real container name(s). - let containers: Vec = get_running_docker_containers() - .await - .into_iter() - .flat_map(|(match_key, real_names)| { - real_names.into_iter().map(move |real_name| Container { - match_key: match_key.clone(), - real_name, - }) - }) - .collect(); + let containers = get_running_docker_containers().await; // One Listener per distinct listening process, keyed by exe path. let mut paths: Vec = listeners::get_all() @@ -465,43 +456,40 @@ async fn declare_services( } } -/// Returns a map of logical name -> real container names for all running Docker containers. -/// -/// Supports both standalone Docker (name -> [name]) and Swarm mode (swarm service label -> [replicas]). -async fn get_running_docker_containers() -> HashMap> { - let mut map: HashMap> = HashMap::new(); - - // Query container name and Swarm service label together +/// Report running and paused containers, using Swarm labels when present. +async fn get_running_docker_containers() -> Vec { let output = tokio::process::Command::new("docker") .args([ "ps", "--format", - "{{.Names}}\t{{.Label \"com.docker.swarm.service.name\"}}", + "{{.Names}}\t{{.Label \"com.docker.swarm.service.name\"}}\t{{.State}}", ]) .output() .await; + let Ok(out) = output else { + return Vec::new(); + }; + parse_container_report(&String::from_utf8_lossy(&out.stdout)) +} - if let Ok(out) = output { - for line in String::from_utf8_lossy(&out.stdout).lines() { - if line.is_empty() { - continue; - } - let parts: Vec<&str> = line.split('\t').collect(); - let real_name = parts[0].to_string(); - let swarm_label = parts.get(1).unwrap_or(&"").trim(); - if swarm_label.is_empty() { - // standalone: logical name = container name - map.entry(real_name.clone()).or_default().push(real_name); - } else { - // Swarm: logical name = swarm service label, may have multiple replicas - map.entry(swarm_label.to_string()) - .or_default() - .push(real_name); +fn parse_container_report(output: &str) -> Vec { + output + .lines() + .filter_map(|line| { + let mut parts = line.split('\t'); + let real_name = parts.next()?.trim(); + if real_name.is_empty() { + return None; } - } - } - - map + let label = parts.next()?.trim(); + let state = parts.next()?.trim(); + Some(Container { + match_key: if label.is_empty() { real_name } else { label }.to_string(), + real_name: real_name.to_string(), + paused: state == "paused", + }) + }) + .collect() } async fn setup_tap( @@ -567,3 +555,17 @@ async fn setup_tap( // } // None // } + +#[cfg(test)] +mod container_report_tests { + #[test] + fn reports_pause_state_for_standalone_and_swarm() { + let report = + super::parse_container_report("web\t\trunning\nworker.1.x\tstack_worker\tpaused\n"); + assert_eq!(report.len(), 2); + assert_eq!(report[0].match_key, "web"); + assert!(!report[0].paused); + assert_eq!(report[1].match_key, "stack_worker"); + assert!(report[1].paused); + } +} diff --git a/members/nullnet-grpc-lib/proto/nullnet_grpc.proto b/members/nullnet-grpc-lib/proto/nullnet_grpc.proto index 143b08b4..b9885185 100644 --- a/members/nullnet-grpc-lib/proto/nullnet_grpc.proto +++ b/members/nullnet-grpc-lib/proto/nullnet_grpc.proto @@ -252,6 +252,7 @@ message ServiceReport { message Container { string match_key = 1; // matched against a service's docker_container (Swarm label / container name) string real_name = 2; // actual container name; stored as the replica identity + bool paused = 3; } message Listener { diff --git a/members/nullnet-grpc-lib/src/proto/nullnet_grpc.rs b/members/nullnet-grpc-lib/src/proto/nullnet_grpc.rs index c73a786b..da0cdaba 100644 --- a/members/nullnet-grpc-lib/src/proto/nullnet_grpc.rs +++ b/members/nullnet-grpc-lib/src/proto/nullnet_grpc.rs @@ -214,6 +214,8 @@ pub struct Container { /// actual container name; stored as the replica identity #[prost(string, tag = "2")] pub real_name: ::prost::alloc::string::String, + #[prost(bool, tag = "3")] + pub paused: bool, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct Listener { diff --git a/members/nullnet-server/src/db/migrations/2026-09-09-000001_service_pausable/down.sql b/members/nullnet-server/src/db/migrations/2026-09-09-000001_service_pausable/down.sql new file mode 100644 index 00000000..869d8cb3 --- /dev/null +++ b/members/nullnet-server/src/db/migrations/2026-09-09-000001_service_pausable/down.sql @@ -0,0 +1 @@ +ALTER TABLE services DROP COLUMN pausable; diff --git a/members/nullnet-server/src/db/migrations/2026-09-09-000001_service_pausable/up.sql b/members/nullnet-server/src/db/migrations/2026-09-09-000001_service_pausable/up.sql new file mode 100644 index 00000000..1683dc9b --- /dev/null +++ b/members/nullnet-server/src/db/migrations/2026-09-09-000001_service_pausable/up.sql @@ -0,0 +1 @@ +ALTER TABLE services ADD COLUMN pausable BOOLEAN NOT NULL DEFAULT FALSE; diff --git a/members/nullnet-server/src/db/models.rs b/members/nullnet-server/src/db/models.rs index 2abf5f96..3562cd8d 100644 --- a/members/nullnet-server/src/db/models.rs +++ b/members/nullnet-server/src/db/models.rs @@ -72,6 +72,7 @@ pub(crate) struct ServiceRow { pub(crate) docker_container: Option, pub(crate) process_path: Option, pub(crate) host_ip: Option, + pub(crate) pausable: bool, pub(crate) port: Option, pub(crate) timeout: Option, pub(crate) max_networks: Option, @@ -91,6 +92,7 @@ pub(crate) struct NewServiceRow<'a> { pub(crate) docker_container: Option<&'a str>, pub(crate) process_path: Option<&'a str>, pub(crate) host_ip: Option<&'a str>, + pub(crate) pausable: bool, pub(crate) port: Option, pub(crate) timeout: Option, pub(crate) max_networks: Option, diff --git a/members/nullnet-server/src/db/schema.rs b/members/nullnet-server/src/db/schema.rs index a678f0c7..5b63e55f 100644 --- a/members/nullnet-server/src/db/schema.rs +++ b/members/nullnet-server/src/db/schema.rs @@ -33,6 +33,7 @@ diesel::table! { docker_container -> Nullable, process_path -> Nullable, host_ip -> Nullable, + pausable -> Bool, port -> Nullable, timeout -> Nullable, max_networks -> Nullable, diff --git a/members/nullnet-server/src/db/stacks.rs b/members/nullnet-server/src/db/stacks.rs index 46ec27c3..779f31e8 100644 --- a/members/nullnet-server/src/db/stacks.rs +++ b/members/nullnet-server/src/db/stacks.rs @@ -19,6 +19,7 @@ pub(crate) struct ServiceInsert<'a> { pub(crate) docker_container: Option<&'a str>, pub(crate) process_path: Option<&'a str>, pub(crate) host_ip: Option<&'a str>, + pub(crate) pausable: bool, pub(crate) port: Option, pub(crate) timeout: Option, pub(crate) max_networks: Option, @@ -170,6 +171,7 @@ impl StackRepository { docker_container: s.docker_container, process_path: s.process_path, host_ip: s.host_ip, + pausable: s.pausable, port: s.port, timeout: s.timeout, max_networks: s.max_networks, @@ -286,6 +288,7 @@ mod tests { docker_container: Some("my-app_web"), process_path: None, host_ip: Some("192.0.2.1"), + pausable: true, port: Some(8080), timeout: Some(0), max_networks: None, @@ -310,6 +313,7 @@ mod tests { assert_eq!(services.len(), 1); assert_eq!(services[0].name, "web"); assert_eq!(services[0].host_ip.as_deref(), Some("192.0.2.1")); + assert!(services[0].pausable); assert_eq!(services[0].docker_container.as_deref(), Some("my-app_web")); let ids: Vec = services.iter().map(|s| s.id).collect(); diff --git a/members/nullnet-server/src/nullnet_grpc_impl.rs b/members/nullnet-server/src/nullnet_grpc_impl.rs index 023090e4..dab2aecc 100644 --- a/members/nullnet-server/src/nullnet_grpc_impl.rs +++ b/members/nullnet-server/src/nullnet_grpc_impl.rs @@ -15,9 +15,7 @@ use crate::services::clients::{Client, ClientInfo}; use crate::services::edge::RegisteredEdge; use crate::services::firewall::{FilterContext, FilterPolicy}; use crate::services::input::{MatchIndex, RouteMap, RouteTarget, ServicesToml, StackMap}; -use crate::services::service_info::{ - RegisteredServiceInfo, ServiceInfo, backend_involved_services, -}; +use crate::services::service_info::{RegisteredServiceInfo, ServiceInfo}; use crate::timeout::check_timeouts; use nullnet_grpc_lib::nullnet_grpc::nullnet_grpc_server::NullnetGrpc; use nullnet_grpc_lib::nullnet_grpc::{ @@ -808,7 +806,13 @@ impl NullnetGrpcImpl { } } - self.apply_services_list_by_stack(sender_ip, &service_list_by_stack) + let paused: HashSet = report + .containers + .iter() + .filter(|c| c.paused) + .map(|c| c.real_name.clone()) + .collect(); + self.apply_services_report(sender_ip, &service_list_by_stack, &paused) .await?; drop(index); @@ -1217,8 +1221,15 @@ impl NullnetGrpcImpl { let this = self.clone(); let initiator_name = initiator_name.to_string(); detached(async move { - let result = this - .setup_backend_chain_owned( + let result = async { + this.resume_trigger_source( + &stack, + &initiator_name, + initiator_ip, + initiator_docker.as_deref(), + ) + .await?; + this.setup_backend_chain_owned( &stack, &initiator_name, initiator_ip, @@ -1226,7 +1237,9 @@ impl NullnetGrpcImpl { port, Some((key.clone(), generation)), ) - .await; + .await + } + .await; if !matches!(result, Ok(true)) { this.orchestrator .cancel_backend_session(&key, generation) @@ -1381,17 +1394,25 @@ impl NullnetGrpcImpl { // would drop the line precisely when the trigger timed out — leaving a // trigger with no edge up, which is the signature of the stranded // reservation this detaching exists to prevent. - let orchestrator = self.orchestrator.clone(); + let this = self.clone(); detached(async move { - let built = orchestrator + let built = this + .orchestrator .ensure_egress_edge( &initiator_stack, &initiator_name, initiator_ip, - initiator_docker, + initiator_docker.clone(), proxy_ip, ) .await?; + this.resume_trigger_source( + &initiator_stack, + &initiator_name, + initiator_ip, + initiator_docker.as_deref(), + ) + .await?; if built { println!( "[egress] edge up for '{initiator_name}' ({initiator_ip}) -> proxy {proxy_ip}" @@ -1403,6 +1424,37 @@ impl NullnetGrpcImpl { Ok(()) } + async fn resume_trigger_source( + &self, + stack: &str, + service: &str, + ip: IpAddr, + docker: Option<&str>, + ) -> Result<(), Error> { + let pending = { + let services = self.services.read().await; + let replica = services + .get(stack) + .and_then(|s| s.get(service)) + .and_then(|info| match info { + ServiceInfo::Registered(reg) => reg + .replicas() + .iter() + .find(|r| r.matches_identity(ip, docker)), + ServiceInfo::Unregistered(_) => None, + }) + .ok_or("trigger source disappeared") + .handle_err(location!())?; + replica.resume(&self.orchestrator) + }; + if let Some(pending) = pending + && !pending.wait().await + { + return Err("trigger source could not be resumed").handle_err(location!()); + } + Ok(()) + } + /// Resolve an egress sender `(sender_ip, container)` to the *registered* /// replica identity `(stack, service_name, ip, docker)` — scanning every stack, /// since the client sends no logical service name. The returned `(ip, docker)` @@ -1652,11 +1704,23 @@ impl NullnetGrpcImpl { &self.http_routes_changed } + #[cfg(test)] #[allow(clippy::type_complexity)] pub(crate) async fn apply_services_list_by_stack( &self, sender_ip: IpAddr, service_list_by_stack: &HashMap)>>, + ) -> Result<(), Error> { + self.apply_services_report(sender_ip, service_list_by_stack, &HashSet::new()) + .await + } + + #[allow(clippy::type_complexity)] + async fn apply_services_report( + &self, + sender_ip: IpAddr, + service_list_by_stack: &HashMap)>>, + paused: &HashSet, ) -> Result<(), Error> { let mut services_mut = self.services.write().await; @@ -1686,6 +1750,14 @@ impl NullnetGrpcImpl { .unwrap_or(false); stack_map.entry(name.clone()).and_modify(|si| { si.add_replica(sender_ip, *port, docker_container.clone()); + if is_new + && docker_container + .as_ref() + .is_some_and(|c| paused.contains(c)) + && let ServiceInfo::Registered(reg) = si + { + reg.mark_replica_suspended(sender_ip, docker_container.as_deref()); + } }); if is_new { self.orchestrator @@ -1696,21 +1768,11 @@ impl NullnetGrpcImpl { } } - // Enforce the invariant: any Docker-backed replica that is idle (e.g. a - // freshly declared, never-requested container at startup) must be paused. - // Backend-involved services are pinned and never paused. - for stack in service_list_by_stack.keys() { - let Some(stack_map) = services_mut.get_mut(stack) else { - continue; - }; - let pinned = backend_involved_services(stack_map); - for (name, si) in stack_map.iter_mut() { - if let ServiceInfo::Registered(reg) = si { - reg.reconcile_suspends(&self.orchestrator, pinned.contains(name)) - .await; - } - } - } + crate::services::service_info::reconcile_container_pauses( + &mut services_mut, + &self.orchestrator, + ) + .await; Ok(()) } @@ -1997,17 +2059,11 @@ async fn run_net_chain_setup( if any_failure || orphaned { if let Some(stack_map) = services_mut.get_mut(&stack) { - let pinned = backend_involved_services(stack_map); for edge in &successful { if let Some(ServiceInfo::Registered(reg)) = stack_map.get_mut(&edge.server_name) && reg.client_net_id(&edge.client) == Some(edge.net_id) { - reg.decrement_chain( - &edge.client, - &orchestrator, - pinned.contains(&edge.server_name), - ) - .await; + reg.decrement_chain(&edge.client, &orchestrator).await; } } } @@ -2096,7 +2152,7 @@ async fn setup_edge( // by (client, proxy) and dep clients by source replica, and a source // replica can only route to one replica of a given dep. let deadline = std::time::Instant::now() + EDGE_CLAIM_TIMEOUT; - let (server_ethernet, server_docker, server_suspended, reservation) = loop { + let (server_ethernet, server_docker, server_resume, reservation) = loop { let waiting_on = { let mut services_guard = services.write().await; let Some(stack_map) = services_guard.get_mut(stack) else { @@ -2132,9 +2188,13 @@ async fn setup_edge( return EdgeOutcome::Failed; } // Does the target replica need unpausing first? - let suspended = reg.replica_suspended(ip, docker.as_deref()); + let resume = reg + .replicas() + .iter() + .find(|r| r.matches_identity(ip, docker.as_deref())) + .and_then(|r| r.resume(orchestrator)); let reservation = reg.pending_notify(&client).expect("just reserved edge"); - break (ip, docker, suspended, reservation); + break (ip, docker, resume, reservation); } } }; @@ -2152,37 +2212,18 @@ async fn setup_edge( let _ = tokio::time::timeout(EDGE_CLAIM_POLL, waiting_on.notified()).await; }; - // Resume the target container before bringing up the link, so it is - // serving by the time traffic arrives. This covers the proxy entry, - // proxy dependencies, and every hop of a backend-triggered chain - // uniformly (it mirrors the per-edge suspend in `decrement_chain`). - if server_suspended && let Some(container) = server_docker.clone() { - if orchestrator - .send_container_resume(server_ethernet, container.clone()) - .await + // The reservation protects the target while its shared resume completes. + if let Some(resume) = server_resume + && !resume.wait().await + { + // roll back the reserved placeholder; the idle replica stays + // suspended (consistent) and the request fails fast. + if let Some(stack_map) = services.write().await.get_mut(stack) + && let Some(ServiceInfo::Registered(reg)) = stack_map.get_mut(&server_name) { - if let Some(stack_map) = services.write().await.get_mut(stack) - && let Some(ServiceInfo::Registered(reg)) = stack_map.get_mut(&server_name) - { - reg.mark_replica_resumed(server_ethernet, server_docker.as_deref()); - } - } else { - orchestrator - .events - .emit(Event::container_resume_failed( - container, - format!("no ack from {server_ethernet} within timeout"), - )) - .await; - // roll back the reserved placeholder; the idle replica stays - // suspended (consistent) and the request fails fast. - if let Some(stack_map) = services.write().await.get_mut(stack) - && let Some(ServiceInfo::Registered(reg)) = stack_map.get_mut(&server_name) - { - reg.remove_pending_client_if(&client, &reservation); - } - return EdgeOutcome::Failed; + reg.remove_pending_client_if(&client, &reservation); } + return EdgeOutcome::Failed; } let Some(net_id) = orchestrator.allocate_net_id().await else { @@ -3006,6 +3047,54 @@ mod host_pin_tests { use nullnet_grpc_lib::nullnet_grpc::Listener; use tonic::transport::server::TcpConnectInfo; + #[tokio::test] + async fn pausable_report_recovers_paused_container_after_restart() { + let (services, entries, _) = validate_stack_toml( + r#" +[[services]] +name = "restarted" +docker_container = "restarted_c" +port = 80 +"#, + ) + .unwrap(); + let server = NullnetGrpcImpl::new_for_test(HashMap::from([("restart".into(), services)])); + *server.match_index.write().await = HashMap::from([("restart".into(), entries)]); + let ip: IpAddr = "192.0.2.1".parse().unwrap(); + let log = server.orchestrator.register_recording_client(ip).await; + let mut request = Request::new(ServiceReport { + containers: vec![nullnet_grpc_lib::nullnet_grpc::Container { + match_key: "restarted_c".into(), + real_name: "restarted_c".into(), + paused: true, + }], + listeners: vec![], + }); + request.extensions_mut().insert(TcpConnectInfo { + local_addr: None, + remote_addr: Some("192.0.2.1:50000".parse().unwrap()), + }); + server.services_list_impl(request).await.unwrap(); + let guard = server.services.read().await; + let ServiceInfo::Registered(reg) = &guard["restart"]["restarted"] else { + panic!("not registered"); + }; + let pending = reg.replicas()[0].resume(&server.orchestrator); + drop(guard); + if let Some(pending) = pending { + assert!(pending.wait().await); + } + let guard = server.services.read().await; + let ServiceInfo::Registered(reg) = &guard["restart"]["restarted"] else { + unreachable!() + }; + assert!(!reg.replica_suspended(ip, Some("restarted_c"))); + assert!(log.lock().await.iter().any(|msg| matches!( + &msg.message, + Some(nullnet_grpc_lib::nullnet_grpc::net_message::Message::ContainerResume(_)) + ))); + } + #[tokio::test] async fn ssh_host_pins_separate_identical_listeners() { let (services, entries, _) = validate_stack_toml( diff --git a/members/nullnet-server/src/orchestrator.rs b/members/nullnet-server/src/orchestrator.rs index 9c9eb535..e655bde1 100644 --- a/members/nullnet-server/src/orchestrator.rs +++ b/members/nullnet-server/src/orchestrator.rs @@ -736,6 +736,25 @@ impl Orchestrator { .min() } + /// Outbound reservations and live sessions also keep their source running. + pub(crate) async fn has_outgoing_work(&self, ip: IpAddr, container: &str) -> bool { + if self + .backend_sessions + .read() + .await + .keys() + .any(|(_, source_ip, docker, _)| { + *source_ip == ip && docker.as_deref() == Some(container) + }) + { + return true; + } + self.egress_edges + .read() + .await + .contains_key(&(ip, Some(container.to_string()))) + } + /// Called under the services lock, before any edge can be reserved. pub(crate) async fn claim_backend_session(&self, key: BackendKey, stack: &str) -> Option { let mut sessions = self.backend_sessions.write().await; diff --git a/members/nullnet-server/src/services/changes.rs b/members/nullnet-server/src/services/changes.rs index a9e8c62b..99d1455a 100644 --- a/members/nullnet-server/src/services/changes.rs +++ b/members/nullnet-server/src/services/changes.rs @@ -1,7 +1,7 @@ use crate::events::Event; use crate::orchestrator::Orchestrator; use crate::services::clients::Client; -use crate::services::service_info::{ServiceInfo, backend_involved_services}; +use crate::services::service_info::ServiceInfo; use std::collections::{HashMap, HashSet}; use std::net::IpAddr; use std::time::Duration; @@ -233,6 +233,7 @@ async fn teardown_invalidated_service( let listen_port = si.listen_port(); let egress_policy = si.egress_policy().clone(); let ingress_policy = si.ingress_policy().clone(); + let pausable = si.pausable(); services.insert( invalidated_service.to_string(), ServiceInfo::new( @@ -244,7 +245,8 @@ async fn teardown_invalidated_service( listen_port, egress_policy, ingress_policy, - ), + ) + .with_pausable(pausable), ); } } @@ -542,12 +544,9 @@ async fn teardown_backend_chain( )); } } - let pinned = backend_involved_services(services); for (client, dep_name) in edges { if let Some(ServiceInfo::Registered(dep_reg)) = services.get_mut(&dep_name) { - dep_reg - .decrement_chain(&client, orchestrator, pinned.contains(&dep_name)) - .await; + dep_reg.decrement_chain(&client, orchestrator).await; } } // This teardown consumed whatever refcount the trigger sessions on these @@ -647,12 +646,9 @@ async fn teardown_dep_chain( orchestrator: &Orchestrator, ) { let edges = collect_dep_chain_edges(service_name, replica_ip, replica_docker, services); - let pinned = backend_involved_services(services); for (client, dep_name) in edges { if let Some(ServiceInfo::Registered(dep_reg)) = services.get_mut(&dep_name) { - dep_reg - .decrement_chain(&client, orchestrator, pinned.contains(&dep_name)) - .await; + dep_reg.decrement_chain(&client, orchestrator).await; } } } diff --git a/members/nullnet-server/src/services/input.rs b/members/nullnet-server/src/services/input.rs index 5293cd58..3d7d1765 100644 --- a/members/nullnet-server/src/services/input.rs +++ b/members/nullnet-server/src/services/input.rs @@ -481,6 +481,13 @@ impl ServicesToml { // proxy-reachable entry point, `None` (omitted) leaves it backend-only. // A declared service can thus host triggers without being reachable. for s in self.services { + if s.pausable && s.docker_container.is_none() { + return Err(format!( + "service '{}': pausable requires docker_container", + s.name + )) + .handle_err(location!()); + } let protocol = s .protocol .map_or(ServiceProtocol::Http, ServiceProtocol::from); @@ -527,7 +534,8 @@ impl ServicesToml { s.listen_port, egress_policy, ingress_policy, - ), + ) + .with_pausable(s.pausable), ); } @@ -716,6 +724,7 @@ fn services_from_rows( docker_container: row.docker_container, process_path: row.process_path, host_ip: row.host_ip, + pausable: row.pausable, port: row.port.and_then(|p| u16::try_from(p).ok()), timeout: row.timeout.and_then(|t| u64::try_from(t).ok()), max_networks: row.max_networks.and_then(|m| u32::try_from(m).ok()), @@ -743,6 +752,7 @@ pub(crate) fn services_to_inserts(services: &[ServiceToml]) -> Vec, /// Optional node IPv4 address, as seen by the control channel. host_ip: Option, + /// Allow idle Docker replicas to be paused. Opt-in only. + #[serde(default)] + pausable: bool, /// Backend port the overlay/proxy connects to on this service's replicas. /// Required when any host-match key is set. Distinct from `listen_port` /// (the proxy's external tcp/udp front port). @@ -1790,12 +1805,39 @@ proxy_dependencies = [["api"]] assert!(detect_route_conflicts(&routes).is_empty()); } + #[test] + fn pausable_defaults_off_and_requires_docker() { + let (services, _, _) = validate_stack_toml( + r#" +[[services]] +name = "plain" +docker_container = "plain" +port = 80 +"#, + ) + .unwrap(); + assert!(!services["plain"].pausable()); + assert!( + validate_stack_toml( + r#" +[[services]] +name = "process" +process_path = "/usr/sbin/sshd" +port = 22 +pausable = true +"# + ) + .is_err() + ); + } + fn empty_service(name: &str) -> ServiceToml { ServiceToml { name: name.to_string(), docker_container: None, process_path: None, host_ip: None, + pausable: false, port: None, timeout: None, proxy_dependencies: Vec::new(), @@ -1820,6 +1862,7 @@ proxy_dependencies = [["api"]] let services = vec![ServiceToml { docker_container: Some("my-app_color".to_string()), host_ip: Some("192.0.2.1".to_string()), + pausable: true, port: Some(3001), timeout: Some(0), proxy_dependencies: vec![vec!["a.dep".to_string(), "b.dep".to_string()]], @@ -1860,6 +1903,7 @@ proxy_dependencies = [["api"]] docker_container: insert.docker_container.map(str::to_string), process_path: insert.process_path.map(str::to_string), host_ip: insert.host_ip.map(str::to_string), + pausable: insert.pausable, port: insert.port, timeout: insert.timeout, max_networks: insert.max_networks, diff --git a/members/nullnet-server/src/services/service_info.rs b/members/nullnet-server/src/services/service_info.rs index f957de0c..17243c5b 100644 --- a/members/nullnet-server/src/services/service_info.rs +++ b/members/nullnet-server/src/services/service_info.rs @@ -3,30 +3,39 @@ use crate::orchestrator::Orchestrator; use crate::services::clients::{Client, ClientInfo, Clients}; use crate::services::firewall::FilterPolicy; use nullnet_grpc_lib::nullnet_grpc::{ServiceProtocol, Upstream}; -use std::collections::{HashMap, HashSet}; +use std::collections::HashMap; use std::net::{IpAddr, Ipv4Addr}; +use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; - -/// Service names that participate in backend trigger chains: those that declare -/// triggers ("have backend deps") and every service named in a trigger chain -/// ("is a backend dep"). These are never paused — backend-dep networks aren't -/// torn down on idle (see `decrement_chain`), so pausing them would strand their -/// networks and leave the replica unreachable. -pub(crate) fn backend_involved_services( - services: &HashMap, -) -> HashSet { - let mut pinned = HashSet::new(); - for (name, info) in services { - let triggers = info.triggers(); - if triggers.is_empty() { - continue; - } - pinned.insert(name.clone()); - for dep in triggers.values().flatten() { - pinned.insert(dep.clone()); - } - } - pinned +use tokio::sync::watch; + +/// A container is idle only when every declaration permits pausing and is idle. +pub(crate) async fn reconcile_container_pauses( + services: &mut crate::services::input::StackMap, + orchestrator: &Orchestrator, +) { + let mut blocked = std::collections::HashSet::new(); + for info in services.values().flat_map(|stack| stack.values()) { + if let ServiceInfo::Registered(reg) = info { + for replica in ®.replicas { + if !reg.pausable + || !replica.clients.clients().is_empty() + || replica.pause.lock().unwrap().resume.is_some() + { + blocked.insert((replica.ip, replica.docker_container.clone())); + } + } + } + } + for info in services.values_mut().flat_map(|stack| stack.values_mut()) { + if let ServiceInfo::Registered(reg) = info { + for replica in &mut reg.replicas { + let allowed = reg.pausable + && !blocked.contains(&(replica.ip, replica.docker_container.clone())); + replica.reconcile_suspend(orchestrator, allowed).await; + } + } + } } #[derive(Clone, Debug)] @@ -56,9 +65,25 @@ impl ServiceInfo { listen_port, egress_policy, ingress_policy, + false, )) } + pub(crate) fn with_pausable(mut self, pausable: bool) -> Self { + match &mut self { + Self::Unregistered(info) => info.pausable = pausable, + Self::Registered(info) => info.pausable = pausable, + } + self + } + + pub(crate) fn pausable(&self) -> bool { + match self { + Self::Unregistered(info) => info.pausable, + Self::Registered(info) => info.pausable, + } + } + pub(crate) fn add_replica(&mut self, ip: IpAddr, port: u16, docker_container: Option) { match self { ServiceInfo::Unregistered(unreg) => { @@ -71,6 +96,7 @@ impl ServiceInfo { listen_port: unreg.listen_port, egress_policy: unreg.egress_policy.clone(), ingress_policy: unreg.ingress_policy.clone(), + pausable: unreg.pausable, replicas: vec![Replica::new(ip, port, docker_container)], }); } @@ -113,6 +139,7 @@ impl ServiceInfo { reg.listen_port, reg.egress_policy.clone(), reg.ingress_policy.clone(), + reg.pausable, )); } } @@ -134,6 +161,7 @@ impl ServiceInfo { reg.listen_port, reg.egress_policy.clone(), reg.ingress_policy.clone(), + reg.pausable, )); } } @@ -158,6 +186,7 @@ impl ServiceInfo { unreg.proxy_deps = loaded.proxy_deps().to_vec(); unreg.triggers.clone_from(loaded.triggers()); unreg.timeout = loaded_timeout; + unreg.pausable = loaded.pausable(); unreg.max_networks = loaded_max_networks; unreg.protocol = loaded_protocol; unreg.listen_port = loaded_listen_port; @@ -168,6 +197,7 @@ impl ServiceInfo { reg.proxy_deps = loaded.proxy_deps().to_vec(); reg.triggers.clone_from(loaded.triggers()); reg.timeout = loaded_timeout; + reg.pausable = loaded.pausable(); reg.max_networks = loaded_max_networks; reg.protocol = loaded_protocol; reg.listen_port = loaded_listen_port; @@ -261,6 +291,7 @@ pub(crate) struct UnregisteredServiceInfo { egress_policy: FilterPolicy, /// Ingress traffic filter for external clients reaching this service via the proxy. ingress_policy: FilterPolicy, + pausable: bool, } impl UnregisteredServiceInfo { @@ -274,6 +305,7 @@ impl UnregisteredServiceInfo { listen_port: Option, egress_policy: FilterPolicy, ingress_policy: FilterPolicy, + pausable: bool, ) -> Self { Self { proxy_deps, @@ -284,21 +316,39 @@ impl UnregisteredServiceInfo { listen_port, egress_policy, ingress_policy, + pausable, } } } +#[derive(Debug, Default)] +struct PauseState { + suspended: bool, + resume: Option>>, +} + +/// Callers share one detached resume and wait without holding service state. +pub(crate) struct Resume(watch::Receiver>); + +impl Resume { + pub(crate) async fn wait(mut self) -> bool { + self.0 + .wait_for(|result| result.is_some()) + .await + .ok() + .and_then(|result| *result) + .unwrap_or(false) + } +} + #[derive(Clone, Debug)] pub(crate) struct Replica { ip: IpAddr, port: u16, docker_container: Option, clients: Clients, - /// Whether the backing Docker container is currently paused. Invariant for - /// non-pinned Docker-backed replicas: `suspended ⟺ clients is empty` - /// (backend-involved services are pinned and never paused — see - /// [`backend_involved_services`]). - suspended: bool, + /// Actual pause state, independent of permission to pause. + pause: Arc>, } impl Replica { @@ -308,7 +358,7 @@ impl Replica { port, docker_container, clients: Clients::default(), - suspended: false, + pause: Arc::default(), } } @@ -334,24 +384,64 @@ impl Replica { } pub(crate) fn suspended(&self) -> bool { - self.suspended + self.pause.lock().unwrap().suspended } - /// Pause this replica's Docker container when it is Docker-backed, has no - /// clients, isn't already suspended, and the service isn't `pinned` - /// (backend-involved — see [`backend_involved_services`]). Fire-and-forget; - /// flips the flag so the invariant `suspended ⟺ no clients` holds. - async fn reconcile_suspend(&mut self, orchestrator: &Orchestrator, pinned: bool) { - if pinned || self.suspended || !self.clients.clients().is_empty() { - return; + pub(crate) fn resume(&self, orchestrator: &Orchestrator) -> Option { + let container = self.docker_container.clone()?; + let mut state = self.pause.lock().unwrap(); + if !state.suspended { + return None; } + if let Some(pending) = &state.resume { + return Some(Resume(pending.clone())); + } + let (tx, rx) = watch::channel(None); + state.resume = Some(rx.clone()); + let pause = self.pause.clone(); + let ip = self.ip; + let orchestrator = orchestrator.clone(); + tokio::spawn(async move { + let ok = orchestrator + .send_container_resume(ip, container.clone()) + .await; + if !ok { + orchestrator + .events + .emit(crate::events::Event::container_resume_failed( + container, + format!("no ack from {ip} within timeout"), + )) + .await; + } + // The captured state belongs to this incarnation, even if it was removed. + let mut state = pause.lock().unwrap(); + state.suspended = !ok; + state.resume = None; + tx.send_replace(Some(ok)); + }); + Some(Resume(rx)) + } + + /// Apply the opt-in policy without suspending live incoming or outgoing work. + async fn reconcile_suspend(&mut self, orchestrator: &Orchestrator, pausable: bool) { let Some(container) = self.docker_container.clone() else { return; }; + if !pausable { + self.resume(orchestrator); + return; + } + if self.suspended() + || !self.clients.clients().is_empty() + || orchestrator.has_outgoing_work(self.ip, &container).await + { + return; + } orchestrator .send_container_suspend(self.ip, container) .await; - self.suspended = true; + self.pause.lock().unwrap().suspended = true; } } @@ -375,6 +465,7 @@ pub(crate) struct RegisteredServiceInfo { egress_policy: FilterPolicy, /// Ingress traffic filter for external clients reaching this service via the proxy. ingress_policy: FilterPolicy, + pausable: bool, /// Replicas of this service. replicas: Vec, } @@ -480,12 +571,7 @@ impl RegisteredServiceInfo { /// Decrement `active_chains` for a specific client entry. /// If it reaches 0, the VXLAN is torn down and the entry is removed. - pub(crate) async fn decrement_chain( - &mut self, - client: &Client, - orchestrator: &Orchestrator, - pinned: bool, - ) { + pub(crate) async fn decrement_chain(&mut self, client: &Client, orchestrator: &Orchestrator) { for replica in &mut self.replicas { if let Some(ci) = replica.clients.clients_mut().get_mut(client) { // A reservation's single refcount belongs to the task building @@ -506,30 +592,14 @@ impl RegisteredServiceInfo { ci.net_id(), ) .await; - // Invariant: a non-pinned Docker-backed replica with no clients is paused. - // Sent immediately, not deferred behind the teardown ack: the flag flips - // here, and `replica_suspended` gates the resume, so any gap between flag - // and pause is a window where an arriving request "resumes" a container - // that was never paused and the late pause then freezes a live client out. - // Ordering the pause is unnecessary anyway — every teardown step works on - // a paused container (hosts edits go through the host-side file). - replica.reconcile_suspend(orchestrator, pinned).await; } return; } } } - /// Pause every Docker-backed replica that is idle and not yet suspended. - /// Used as the registration-time and periodic safety net that enforces the - /// invariant for replicas missed by the per-event hooks. - pub(crate) async fn reconcile_suspends(&mut self, orchestrator: &Orchestrator, pinned: bool) { - for replica in &mut self.replicas { - replica.reconcile_suspend(orchestrator, pinned).await; - } - } - /// Whether the replica identified by `(ip, docker)` is currently suspended. + #[cfg(test)] pub(crate) fn replica_suspended(&self, ip: IpAddr, docker: Option<&str>) -> bool { self.replicas .iter() @@ -537,14 +607,14 @@ impl RegisteredServiceInfo { .is_some_and(Replica::suspended) } - /// Clear the suspended flag for `(ip, docker)` after a successful unpause. - pub(crate) fn mark_replica_resumed(&mut self, ip: IpAddr, docker: Option<&str>) { + /// Initialize a newly discovered replica from the host's pause observation. + pub(crate) fn mark_replica_suspended(&mut self, ip: IpAddr, docker: Option<&str>) { if let Some(replica) = self .replicas .iter_mut() .find(|r| r.matches_identity(ip, docker)) { - replica.suspended = false; + replica.pause.lock().unwrap().suspended = true; } } @@ -847,3 +917,39 @@ impl RegisteredServiceInfo { .collect() } } + +#[cfg(test)] +mod resume_tests { + use super::*; + + #[tokio::test] + async fn resume_survives_a_dropped_waiter_and_preserves_replacement_state() { + let orchestrator = Orchestrator::new(); + let ip = "192.0.2.1".parse().unwrap(); + let (log, gate) = orchestrator.register_gated_client(ip).await; + let old = Replica::new(ip, 80, Some("container".into())); + old.pause.lock().unwrap().suspended = true; + drop(old.resume(&orchestrator).unwrap()); + let pending = old.resume(&orchestrator).unwrap(); + let replacement = Replica::new(ip, 80, Some("container".into())); + replacement.pause.lock().unwrap().suspended = true; + drop(old); + gate.add_permits(1); + assert!(pending.wait().await); + assert!(replacement.suspended()); + assert_eq!(log.lock().await.len(), 1); + } + + #[tokio::test] + async fn failed_resume_stays_paused_and_can_retry() { + let orchestrator = Orchestrator::new(); + let ip = "192.0.2.1".parse().unwrap(); + let replica = Replica::new(ip, 80, Some("container".into())); + replica.pause.lock().unwrap().suspended = true; + assert!(!replica.resume(&orchestrator).unwrap().wait().await); + assert!(replica.suspended()); + orchestrator.register_fake_client(ip).await; + assert!(replica.resume(&orchestrator).unwrap().wait().await); + assert!(!replica.suspended()); + } +} diff --git a/members/nullnet-server/src/tests.rs b/members/nullnet-server/src/tests.rs index e548f090..d8f12fde 100644 --- a/members/nullnet-server/src/tests.rs +++ b/members/nullnet-server/src/tests.rs @@ -2854,7 +2854,8 @@ fn suspend_test_server() -> NullnetGrpcImpl { None, FilterPolicy::None, FilterPolicy::None, - ), + ) + .with_pausable(true), ); inner.insert( "host_svc".to_string(), @@ -2867,7 +2868,8 @@ fn suspend_test_server() -> NullnetGrpcImpl { None, FilterPolicy::None, FilterPolicy::None, - ), + ) + .with_pausable(true), ); NullnetGrpcImpl::new_for_test(into_stack_map(inner)) } @@ -2920,12 +2922,10 @@ async fn docker_replica_suspended_on_registration() { assert_eq!(total_suspends(&log.lock().await), 1); } -/// Services that participate in backend trigger chains are pinned: the -/// initiator ("has backend deps") and the dep ("is a backend dep") are never -/// paused, even idle and Docker-backed, because backend-dep networks aren't -/// torn down on idle. A plain Docker replica in the same list is still paused. +/// Backend involvement does not override the explicit pause policy. + #[tokio::test] -async fn backend_involved_replicas_never_suspended() { +async fn backend_services_follow_pause_checkbox() { let mut inner = HashMap::new(); // initiator declares a trigger chain → dep (so it "has backend deps") inner.insert( @@ -2969,6 +2969,10 @@ async fn backend_involved_replicas_never_suspended() { FilterPolicy::None, ), ); + for name in ["dep", "plain"] { + let info = inner.remove(name).unwrap(); + inner.insert(name.to_string(), info.with_pausable(true)); + } let server = NullnetGrpcImpl::new_for_test(into_stack_map(inner)); let svc_ip = ip(1, 1, 1, 1); @@ -2995,16 +2999,16 @@ async fn backend_involved_replicas_never_suspended() { svc_ip, "initiator_c1" )); - assert!(!replica_is_suspended(&guard, "dep", svc_ip, "dep_c1")); + assert!(replica_is_suspended(&guard, "dep", svc_ip, "dep_c1")); // control: the unrelated replica is still paused while idle assert!(replica_is_suspended(&guard, "plain", svc_ip, "plain_c1")); } - // only the plain replica was ever suspended + // Both opted-in replicas pause; the unchecked initiator stays running. wait_for_log(&log, |l| count_suspends(l, "plain_c1") == 1).await; assert_eq!(count_suspends(&log.lock().await, "initiator_c1"), 0); - assert_eq!(count_suspends(&log.lock().await, "dep_c1"), 0); - assert_eq!(total_suspends(&log.lock().await), 1); + assert_eq!(count_suspends(&log.lock().await, "dep_c1"), 1); + assert_eq!(total_suspends(&log.lock().await), 2); } #[tokio::test] @@ -3116,7 +3120,7 @@ async fn stale_session_broken_chain_is_evicted_and_rebuilt() { panic!("C is not registered"); }; c_reg - .decrement_chain(&b_client, server.orchestrator(), false) + .decrement_chain(&b_client, server.orchestrator()) .await; } assert!( @@ -4241,3 +4245,319 @@ async fn clients_sharing_a_net_id_each_close_their_own_row() { "every client's row must close: {rows:?}" ); } + +#[tokio::test] +async fn pausable_defaults_off_and_live_toggle_resumes() { + let server = suspend_test_server(); + let svc_ip = ip(1, 1, 1, 1); + let log = server + .orchestrator() + .register_recording_client(svc_ip) + .await; + let list = declared_list(vec![("svc", 8080, Some("svc_c1"))]); + server + .apply_services_list_by_stack(svc_ip, &list) + .await + .unwrap(); + assert!(replica_is_suspended( + &*server.services().read().await, + "svc", + svc_ip, + "svc_c1" + )); + { + let mut guard = server.services().write().await; + let info = guard.get_mut(TEST_STACK).unwrap().remove("svc").unwrap(); + guard + .get_mut(TEST_STACK) + .unwrap() + .insert("svc".into(), info.with_pausable(false)); + crate::services::service_info::reconcile_container_pauses( + &mut guard, + server.orchestrator(), + ) + .await; + let ServiceInfo::Registered(reg) = &guard[TEST_STACK]["svc"] else { + unreachable!() + }; + let pending = reg.replicas()[0].resume(server.orchestrator()); + drop(guard); + if let Some(pending) = pending { + assert!(pending.wait().await); + } + assert!(!replica_is_suspended( + &*server.services().read().await, + "svc", + svc_ip, + "svc_c1" + )); + } + assert_eq!(count_suspends(&log.lock().await, "svc_c1"), 1); + server + .apply_services_list_by_stack(svc_ip, &list) + .await + .unwrap(); + assert_eq!(count_suspends(&log.lock().await, "svc_c1"), 1); +} + +#[tokio::test] +async fn pausable_backend_source_resumes_and_stays_running_while_held() { + let (server, node, log) = paused_backend_test_server().await; + server + .handle_backend_trigger("source", 8080, node, Some("source_c")) + .await + .unwrap(); + let mut guard = server.services().write().await; + crate::services::service_info::reconcile_container_pauses(&mut guard, server.orchestrator()) + .await; + assert!(!replica_is_suspended(&guard, "source", node, "source_c")); + assert!(!replica_is_suspended(&guard, "dep", node, "dep_c")); + assert_eq!(count_suspends(&log.lock().await, "source_c"), 1); +} + +#[tokio::test] +async fn pausable_shared_container_requires_all_declarations_to_opt_in() { + let (inner, _, _) = crate::services::input::validate_stack_toml( + r#" +[[services]] +name = "opted" +docker_container = "shared" +port = 80 +pausable = true +[[services]] +name = "unchecked" +docker_container = "shared" +port = 81 +"#, + ) + .unwrap(); + let server = NullnetGrpcImpl::new_for_test(into_stack_map(inner)); + let node = ip(1, 1, 1, 1); + let log = server.orchestrator().register_recording_client(node).await; + server + .apply_services_list_by_stack( + node, + &declared_list(vec![ + ("opted", 80, Some("shared")), + ("unchecked", 81, Some("shared")), + ]), + ) + .await + .unwrap(); + assert_eq!(total_suspends(&log.lock().await), 0); +} + +#[tokio::test] +async fn pausable_survives_last_replica_disappearing() { + let server = suspend_test_server(); + let node = ip(1, 1, 1, 1); + let log = server.orchestrator().register_recording_client(node).await; + let list = declared_list(vec![("svc", 8080, Some("svc_c1"))]); + server + .apply_services_list_by_stack(node, &list) + .await + .unwrap(); + server + .apply_services_list_by_stack(node, &HashMap::new()) + .await + .unwrap(); + assert!(stack_view(&*server.services().read().await)["svc"].pausable()); + server + .apply_services_list_by_stack(node, &list) + .await + .unwrap(); + assert!(replica_is_suspended( + &*server.services().read().await, + "svc", + node, + "svc_c1" + )); + wait_for_log(&log, |l| count_suspends(l, "svc_c1") == 2).await; +} + +#[tokio::test] +async fn pausable_unacknowledged_policy_resume_does_not_lock_services() { + let server = suspend_test_server(); + let node = ip(1, 1, 1, 1); + server.orchestrator().register_fake_client(node).await; + server + .apply_services_list_by_stack(node, &declared_list(vec![("svc", 8080, Some("svc_c1"))])) + .await + .unwrap(); + let (log, gate) = server.orchestrator().register_gated_client(node).await; + let worker = server.clone(); + let task = tokio::spawn(async move { + let mut guard = worker.services().write().await; + let info = guard.get_mut(TEST_STACK).unwrap().remove("svc").unwrap(); + guard + .get_mut(TEST_STACK) + .unwrap() + .insert("svc".into(), info.with_pausable(false)); + crate::services::service_info::reconcile_container_pauses( + &mut guard, + worker.orchestrator(), + ) + .await; + }); + wait_for_log(&log, |l| count_resumes(l, "svc_c1") == 1).await; + let unlocked = tokio::time::timeout( + std::time::Duration::from_millis(100), + server.services().write(), + ) + .await + .is_ok(); + task.await.unwrap(); + assert!( + unlocked, + "a withheld resume ack blocked unrelated service operations" + ); + let pending = { + let mut guard = server.services().write().await; + crate::services::service_info::reconcile_container_pauses( + &mut guard, + server.orchestrator(), + ) + .await; + let ServiceInfo::Registered(reg) = &guard[TEST_STACK]["svc"] else { + unreachable!() + }; + assert!(reg.replicas()[0].suspended()); + reg.replicas()[0].resume(server.orchestrator()).unwrap() + }; + gate.add_permits(10); + assert!(pending.wait().await); + assert!(!replica_is_suspended( + &*server.services().read().await, + "svc", + node, + "svc_c1" + )); + assert_eq!(count_resumes(&log.lock().await, "svc_c1"), 1); +} + +async fn paused_backend_test_server() -> ( + NullnetGrpcImpl, + IpAddr, + std::sync::Arc>>, +) { + let (inner, _, _) = crate::services::input::validate_stack_toml( + r#" +[[services]] +name = "source" +docker_container = "source_c" +port = 80 +pausable = true +triggers = [{port = 8080, chain = ["dep"]}] +[[services]] +name = "dep" +docker_container = "dep_c" +port = 80 +pausable = true +"#, + ) + .unwrap(); + let server = NullnetGrpcImpl::new_for_test(into_stack_map(inner)); + let node = ip(1, 1, 1, 1); + let log = server.orchestrator().register_recording_client(node).await; + server + .apply_services_list_by_stack( + node, + &declared_list(vec![ + ("source", 80, Some("source_c")), + ("dep", 80, Some("dep_c")), + ]), + ) + .await + .unwrap(); + (server, node, log) +} + +#[tokio::test] +async fn pausable_unacknowledged_trigger_resume_does_not_lock_services() { + let (server, node, _) = paused_backend_test_server().await; + let (log, gate) = server.orchestrator().register_gated_client(node).await; + let worker = server.clone(); + let task = tokio::spawn(async move { + worker + .handle_backend_trigger("source", 8080, node, Some("source_c")) + .await + }); + wait_for_log(&log, |l| count_resumes(l, "source_c") == 1).await; + let unlocked = tokio::time::timeout( + std::time::Duration::from_millis(100), + server.services().write(), + ) + .await + .is_ok(); + gate.add_permits(20); + task.await.unwrap().unwrap(); + assert!( + unlocked, + "a withheld trigger resume ack blocked unrelated service operations" + ); + assert_eq!(count_resumes(&log.lock().await, "source_c"), 1); +} + +#[tokio::test] +async fn pausable_shared_container_does_not_pause_during_alias_resume() { + let (inner, _, _) = crate::services::input::validate_stack_toml( + r#" +[[services]] +name = "first" +docker_container = "shared" +port = 80 +pausable = true +[[services]] +name = "second" +docker_container = "shared" +port = 81 +pausable = true +"#, + ) + .unwrap(); + let server = NullnetGrpcImpl::new_for_test(into_stack_map(inner)); + let node = ip(1, 1, 1, 1); + server.orchestrator().register_fake_client(node).await; + server + .apply_services_list_by_stack( + node, + &declared_list(vec![ + ("first", 80, Some("shared")), + ("second", 81, Some("shared")), + ]), + ) + .await + .unwrap(); + let pending = { + let guard = server.services().read().await; + let ServiceInfo::Registered(reg) = &guard[TEST_STACK]["first"] else { + unreachable!() + }; + reg.replicas()[0].resume(server.orchestrator()).unwrap() + }; + assert!(pending.wait().await); + let (log, gate) = server.orchestrator().register_gated_client(node).await; + let pending = { + let guard = server.services().read().await; + let ServiceInfo::Registered(reg) = &guard[TEST_STACK]["second"] else { + unreachable!() + }; + reg.replicas()[0].resume(server.orchestrator()).unwrap() + }; + wait_for_log(&log, |l| count_resumes(l, "shared") == 1).await; + let remained_running = { + let mut guard = server.services().write().await; + crate::services::service_info::reconcile_container_pauses( + &mut guard, + server.orchestrator(), + ) + .await; + !replica_is_suspended(&guard, "first", node, "shared") + }; + gate.add_permits(20); + assert!(pending.wait().await); + assert!( + remained_running, + "an idle alias paused a container while its resume was pending" + ); +} diff --git a/members/nullnet-server/src/timeout.rs b/members/nullnet-server/src/timeout.rs index 8ae910ad..8c43642d 100644 --- a/members/nullnet-server/src/timeout.rs +++ b/members/nullnet-server/src/timeout.rs @@ -3,7 +3,7 @@ use crate::services::changes::{ ServiceChange, apply_changes, dep_chain_has_pending, release_backend_chain, }; use crate::services::input::StackMap; -use crate::services::service_info::{ServiceInfo, backend_involved_services}; +use crate::services::service_info::ServiceInfo; use std::collections::HashMap; use std::sync::Arc; use std::time::Duration; @@ -73,6 +73,8 @@ pub(crate) async fn check_timeouts( apply_timeouts(stack_map, &orchestrator, &stack).await; } } + crate::services::service_info::reconcile_container_pauses(&mut services_mut, &orchestrator) + .await; } } @@ -116,19 +118,6 @@ pub(crate) async fn apply_timeouts( if !changes.is_empty() { apply_changes(changes, services, None, orchestrator, stack).await; } - - // Safety net: enforce the invariant that every idle Docker-backed replica is - // paused, catching any missed by the per-event hooks (startup, races, - // restarts). Cheap when nothing is pending — `reconcile_suspends` skips - // replicas that are already suspended or still have clients. Backend-involved - // services are pinned and never paused. - let pinned = backend_involved_services(services); - for (name, si) in services.iter_mut() { - if let ServiceInfo::Registered(reg) = si { - reg.reconcile_suspends(orchestrator, pinned.contains(name)) - .await; - } - } } fn collect_timed_out_clients(services: &HashMap) -> Vec { diff --git a/members/nullnet-server/tests/fixtures/multi_replica/after_b_same_ip_disconnect.dot b/members/nullnet-server/tests/fixtures/multi_replica/after_b_same_ip_disconnect.dot index 5ff746db..cd6b673e 100644 --- a/members/nullnet-server/tests/fixtures/multi_replica/after_b_same_ip_disconnect.dot +++ b/members/nullnet-server/tests/fixtures/multi_replica/after_b_same_ip_disconnect.dot @@ -3,7 +3,7 @@ digraph G { node [color=white, fontcolor=white]; edge [color=white, fontcolor=white, fontsize=9, labelangle=180, labeldistance=0.8]; - "A" [label=1 active / 1 paused / 2 total>] [style=solid, color=green]; + "A" [label=1 active / 0 paused / 2 total>] [style=solid, color=green]; "10.0.0.4 (via 7.7.7.7)" -> "A" [label="VXLAN 106 [0ms]"]; "B" [label=2 active / 0 paused / 2 total>] [style=dashed, color=green]; diff --git a/members/nullnet-server/ui/src/pages/Config.tsx b/members/nullnet-server/ui/src/pages/Config.tsx index dd139e94..183c9843 100644 --- a/members/nullnet-server/ui/src/pages/Config.tsx +++ b/members/nullnet-server/ui/src/pages/Config.tsx @@ -241,6 +241,7 @@ interface ServiceFormState { matchKind: MatchKind; matchValue: string; hostIp: string; + pausable: boolean; port: string; reachable: boolean; timeout: string; @@ -258,6 +259,7 @@ const EMPTY_FORM: ServiceFormState = { matchKind: 'docker', matchValue: '', hostIp: '', + pausable: false, port: '', reachable: false, timeout: '0', @@ -361,6 +363,7 @@ function serviceToForm(s: ServiceConfigJson): ServiceFormState { matchKind: s.process_path ? 'process' : 'docker', matchValue: s.docker_container ?? s.process_path ?? '', hostIp: s.host_ip ?? '', + pausable: s.pausable ?? false, port: s.port != null ? String(s.port) : '', reachable: s.timeout != null, timeout: s.timeout != null ? String(s.timeout) : '0', @@ -384,6 +387,7 @@ function formToService(f: ServiceFormState): ServiceConfigJson { docker_container: f.matchKind === 'docker' ? f.matchValue.trim() : null, process_path: f.matchKind === 'process' ? f.matchValue.trim() : null, host_ip: f.hostIp.trim() || null, + pausable: f.matchKind === 'docker' && f.pausable, port: f.port.trim() !== '' ? Number(f.port) : null, timeout: f.reachable ? Number(f.timeout || '0') : null, proxy_dependencies: f.dependencies.map(chain).filter(branch => branch.length > 0), @@ -833,6 +837,17 @@ export default function Config() { /> + {form.matchKind === 'docker' && ( + + )} +