diff --git a/crates/subc-core/src/bin/ck.rs b/crates/subc-core/src/bin/ck.rs index 93bb6470..e1060b99 100644 --- a/crates/subc-core/src/bin/ck.rs +++ b/crates/subc-core/src/bin/ck.rs @@ -1799,6 +1799,10 @@ fn provenance_value(value: Option<&Value>) -> String { } } +fn terminal_safe_string(value: &str) -> String { + provenance_value(Some(&Value::String(value.to_string()))) +} + fn provenance_image(value: Option<&Value>) -> String { let Some(value) = value else { return "is unknown".to_string(); @@ -2739,27 +2743,55 @@ fn daemon_frame_drop_summary(describe: &Value) -> String { .and_then(|counters| counters.get("module_frames_dropped_no_route_last_10m")) .and_then(Value::as_u64) .unwrap_or(0); - if drops == 0 { - return "no frame drops in the last 10 minutes".to_string(); + let frame_drop_summary = if drops == 0 { + "no frame drops in the last 10 minutes".to_string() + } else { + let top = describe + .get("counters") + .and_then(|counters| counters.get("module_frames_dropped_no_route_by_module")) + .and_then(Value::as_object) + .and_then(|modules| { + modules + .iter() + .filter_map(|(module_id, count)| count.as_u64().map(|count| (module_id, count))) + .max_by(|(left_id, left_count), (right_id, right_count)| { + left_count + .cmp(right_count) + .then_with(|| right_id.cmp(left_id)) + }) + .map(|(module_id, _)| module_id.as_str()) + }) + .unwrap_or("unknown"); + let noun = if drops == 1 { "drop" } else { "drops" }; + format!("{drops} frame {noun} in the last 10 minutes, top: {top}") + }; + let refusals = describe + .get("counters") + .and_then(|counters| counters.get("route_open_refused_by_code")) + .and_then(Value::as_object) + .map(|codes| codes.values().filter_map(Value::as_u64).sum::()) + .unwrap_or(0); + if refusals == 0 { + return frame_drop_summary; } let top = describe .get("counters") - .and_then(|counters| counters.get("module_frames_dropped_no_route_by_module")) + .and_then(|counters| counters.get("route_open_refused_by_code")) .and_then(Value::as_object) - .and_then(|modules| { - modules + .and_then(|codes| { + codes .iter() - .filter_map(|(module_id, count)| count.as_u64().map(|count| (module_id, count))) - .max_by(|(left_id, left_count), (right_id, right_count)| { + .filter_map(|(code, count)| count.as_u64().map(|count| (code, count))) + .max_by(|(left_code, left_count), (right_code, right_count)| { left_count .cmp(right_count) - .then_with(|| right_id.cmp(left_id)) + .then_with(|| right_code.cmp(left_code)) }) - .map(|(module_id, _)| module_id.as_str()) + .map(|(code, _)| terminal_safe_string(code)) }) - .unwrap_or("unknown"); - let noun = if drops == 1 { "drop" } else { "drops" }; - format!("{drops} frame {noun} in the last 10 minutes, top: {top}") + .unwrap_or_else(|| "unknown".to_string()); + let noun = if refusals == 1 { "refusal" } else { "refusals" }; + format!("{frame_drop_summary}; {refusals} route.open {noun}, top: {top}") } /// Compare the daemon's embedded build provenance against this CLI's own. @@ -4655,11 +4687,25 @@ fn display_json_value(value: &Value) -> String { Value::Bool(value) => value.to_string(), Value::Number(value) => value.to_string(), Value::Array(_) | Value::Object(_) => { - serde_json::to_string(value).unwrap_or_else(|_| value.to_string()) + serde_json::to_string(&terminal_safe_json(value)).unwrap_or_else(|_| value.to_string()) } } } +fn terminal_safe_json(value: &Value) -> Value { + match value { + Value::String(value) => Value::String(terminal_safe_string(value)), + Value::Array(values) => Value::Array(values.iter().map(terminal_safe_json).collect()), + Value::Object(values) => Value::Object( + values + .iter() + .map(|(key, value)| (terminal_safe_string(key), terminal_safe_json(value))) + .collect(), + ), + other => other.clone(), + } +} + fn connection_file_age(path: &Path) -> Option { let modified = fs::metadata(path).ok()?.modified().ok()?; SystemTime::now().duration_since(modified).ok() @@ -5905,6 +5951,23 @@ mod tests { } } + #[test] + fn route_open_refusal_counter_renderers_escape_terminal_controls() { + let hostile = "\u{1b}]52;c;AAAA\u{07}"; + let summary = daemon_frame_drop_summary(&serde_json::json!({ + "counters": { + "module_frames_dropped_no_route_last_10m": 0, + "route_open_refused_by_code": { hostile: 1 } + } + })); + let verbose = display_json_value(&serde_json::json!({ hostile: 1 })); + + for rendered in [summary, verbose] { + assert!(!rendered.bytes().any(|byte| byte < 0x20)); + assert!(rendered.contains(r"\x1b"), "rendered: {rendered:?}"); + } + } + #[test] fn provenance_image_renders_unknown_reason_escaped() { let value = serde_json::json!({ diff --git a/crates/subc-core/src/control.rs b/crates/subc-core/src/control.rs index 0358bc60..4981e634 100644 --- a/crates/subc-core/src/control.rs +++ b/crates/subc-core/src/control.rs @@ -1526,6 +1526,58 @@ impl ControlHandler { )?)) } + fn route_open_refusal_frame( + &self, + ctx: &RouteCtx, + frame: &Frame, + module_id: &str, + code: &'static str, + message: impl Into, + ) -> Result { + self.observe_route_open_refusal(ctx, module_id, code); + control_error_frame(frame, code, message.into()) + } + + fn observe_route_open_refusal(&self, ctx: &RouteCtx, module_id: &str, code: &'static str) { + self.counters.increment_route_open_refused(code); + info!( + target: "subc_core::control", + code, + module_id = ?module_id, + connection_id = ctx.connection_id.get(), + "route.open refused" + ); + } + + fn supervised_absent_route_open_refusal_frame( + &self, + ctx: &RouteCtx, + frame: &Frame, + module_id: &str, + code: &'static str, + status: &crate::supervise::ModuleStatus, + ) -> Result { + self.counters.increment_route_open_refused(code); + info!( + target: "subc_core::control", + code, + module_id = ?module_id, + connection_id = ctx.connection_id.get(), + state = %status.state, + enabled = status.enabled, + live = status.live, + "route.open refused" + ); + control_error_frame( + frame, + code, + format!( + "module_id '{module_id}' is supervised but not available (state={}, enabled={}, live={})", + status.state, status.enabled, status.live + ), + ) + } + async fn handle_route_open( &self, ctx: &RouteCtx, @@ -1573,46 +1625,54 @@ impl ControlHandler { if let Some((status, warming)) = self.supervisor_status(&target_module_id, frame.header.corr)? { - return Ok(vec![control_error_frame( + let code = if warming { + "module_warming" + } else { + "target_unavailable" + }; + return Ok(vec![self.supervised_absent_route_open_refusal_frame( + ctx, &frame, - if warming { - "module_warming" - } else { - "target_unavailable" - }, - format!( - "module_id '{target_module_id}' is supervised but not available (state={}, enabled={}, live={})", - status.state, status.enabled, status.live - ), + &target_module_id, + code, + &status, )?]); } if let Some(removed_ago_ms) = self.supervisor.removal_tombstone_age_ms(&target_module_id) { - return Ok(vec![control_error_frame( + return Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, error_codes::MODULE_REMOVED, format!("module_id '{target_module_id}' was removed {removed_ago_ms} ms ago"), )?]); } - return Ok(vec![control_error_frame( + return Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, error_codes::UNKNOWN_MODULE, format!("module_id '{target_module_id}' is not registered"), )?]); }; if !target_has_required_role(&target, ®istration.manifest.provides) { - return Ok(vec![control_error_frame( + return Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, "target_unavailable", format!("module_id '{target_module_id}' does not provide the requested target"), )?]); } if registration.state != ChannelState::Active { - return Ok(vec![control_error_frame( + return Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, "target_unavailable", format!("module_id '{target_module_id}' is not active"), )?]); @@ -1623,8 +1683,10 @@ impl ControlHandler { .module_is_draining(&target_module_id) .map_err(RouterError::Forwarding)? { - return Ok(vec![control_error_frame( + return Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, "module_reloading", format!("module_id '{target_module_id}' is reloading"), )?]); @@ -1636,8 +1698,10 @@ impl ControlHandler { .and_then(|process_liveness| process_liveness.process_live(&target_module_id)) == Some(false) { - return Ok(vec![control_error_frame( + return Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, "target_unavailable", format!("module_id '{target_module_id}' is not live"), )?]); @@ -1648,8 +1712,10 @@ impl ControlHandler { .has_live_module_connection(&target_module_id) .map_err(RouterError::Forwarding)? { - return Ok(vec![control_error_frame( + return Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, "target_unavailable", format!("module_id '{target_module_id}' has no live forwarding connection"), )?]); @@ -1658,12 +1724,16 @@ impl ControlHandler { if let Some(error) = self.guard_module_control_op(&frame, &target_module_id, "route.bind")? { + self.observe_route_open_refusal(ctx, &target_module_id, "op_not_allowed"); return Ok(vec![error]); } let principal = match self.route_open_principal(&frame, consumer_identity)? { Ok(principal) => principal, - Err(error) => return Ok(vec![error]), + Err(error) => { + self.observe_route_open_refusal(ctx, &target_module_id, "bad_consumer_identity"); + return Ok(vec![error]); + } }; // This is attested, control-plane policy for supervised module origins. @@ -1687,8 +1757,10 @@ impl ControlHandler { capability, "refusing route.open because an attested capability deny edge matches" ); - return Ok(vec![control_error_frame( + return Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, "capability_forbidden", format!( "module_id '{opening_module_id}' must never reach capability '{capability}' provided by '{target_module_id}'" @@ -1705,8 +1777,10 @@ impl ControlHandler { if self.admission_facts_carrier_module_id.as_deref() == Some(module_id) ); if !carrier_matches { - return Ok(vec![control_error_frame( + return Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, "admission_facts_not_permitted", "admission facts may only be carried by the configured reserved module", )?]); @@ -1717,8 +1791,10 @@ impl ControlHandler { .as_ref() .is_some_and(|targets| targets.iter().any(|id| id == &target_module_id)); if !target_allowed { - return Ok(vec![control_error_frame( + return Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, "admission_facts_target_not_allowed", format!( "admission facts are not permitted for target module_id '{target_module_id}'" @@ -1785,8 +1861,10 @@ impl ControlHandler { { Ok(pending) => pending, Err(err) => { - return Ok(vec![control_error_frame( + return Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, forwarding_error_code(&err), err.to_string(), )?]) @@ -1843,8 +1921,10 @@ impl ControlHandler { if let Err(err) = module_sink.send(relay_frame).await { reservation.release_and_disarm(); - return Ok(vec![control_error_frame( + return Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, "target_unavailable", err.to_string(), )?]); @@ -1870,6 +1950,16 @@ impl ControlHandler { } Ok(Ok(RouteBindRelayOutcome::Rejected(body))) => { reservation.release_and_disarm(); + self.counters + .increment_route_open_refused("module_rejected"); + info!( + target: "subc_core::control", + code = "module_rejected", + module_code = ?body.code, + module_id = ?target_module_id, + connection_id = ctx.connection_id.get(), + "route.open refused" + ); Ok(vec![control_error_body_frame(&frame, body)?]) } Ok(Ok(RouteBindRelayOutcome::ModuleGone(message))) => { @@ -1884,16 +1974,20 @@ impl ControlHandler { module_id = %target_module_id, "route.bind relay abandoned: {message}" ); - Ok(vec![control_error_frame( + Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, "target_unavailable", message, )?]) } Ok(Err(_)) => { reservation.release_and_disarm(); - Ok(vec![control_error_frame( + Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, "target_unavailable", "route.bind relay waiter was canceled before the module responded", )?]) @@ -1912,8 +2006,10 @@ impl ControlHandler { timeout_ms = route_bind_relay_timeout.as_millis() as u64, "route.bind relay timed out: module did not ack within budget" ); - Ok(vec![control_error_frame( + Ok(vec![self.route_open_refusal_frame( + ctx, &frame, + &target_module_id, "module_timeout", format!( "module_id '{target_module_id}' did not answer route.bind within {:?}", @@ -3995,7 +4091,13 @@ pub(crate) fn send_route_control_pushes( #[cfg(test)] mod tests { - use std::{path::PathBuf, sync::Arc, time::Duration}; + use std::{ + collections::BTreeMap, + fmt, + path::PathBuf, + sync::{Arc, Mutex}, + time::Duration, + }; use serde_json::{json, Value}; use subc_protocol::{ @@ -4021,6 +4123,11 @@ mod tests { sync::mpsc, time::{sleep, Instant}, }; + use tracing::{ + field::{Field, Visit}, + Event, Subscriber, + }; + use tracing_subscriber::{layer::Context, prelude::*, Layer}; /// Locates the `fake-aft-stub` binary from a `src/lib.rs` unit test. /// @@ -4442,6 +4549,49 @@ mod tests { Frame::build(FrameType::Request, control_flags(), 0, 0, corr, body).unwrap() } + #[derive(Clone, Default)] + struct EventCapture { + events: Arc>>, + } + + #[derive(Clone, Debug)] + struct CapturedEvent { + target: String, + fields: BTreeMap, + } + + impl EventCapture { + fn events(&self) -> Vec { + self.events.lock().unwrap().clone() + } + } + + impl Layer for EventCapture + where + S: Subscriber, + { + fn on_event(&self, event: &Event<'_>, _context: Context<'_, S>) { + let mut visitor = EventFieldVisitor::default(); + event.record(&mut visitor); + self.events.lock().unwrap().push(CapturedEvent { + target: event.metadata().target().to_string(), + fields: visitor.fields, + }); + } + } + + #[derive(Default)] + struct EventFieldVisitor { + fields: BTreeMap, + } + + impl Visit for EventFieldVisitor { + fn record_debug(&mut self, field: &Field, value: &dyn fmt::Debug) { + self.fields + .insert(field.name().to_string(), format!("{value:?}")); + } + } + fn health_response(corr: u64, status: HealthStatus) -> Frame { let body = serde_json::to_vec(&ModuleControlResponse::HealthCheck { status, @@ -5368,7 +5518,9 @@ mod tests { let handler = handler.clone(); let project_root = unique_project_root(project_root_label); let module_id = module_id.to_string(); + let dispatch = tracing::dispatcher::get_default(|dispatch| dispatch.clone()); let task = tokio::spawn(async move { + let _guard = tracing::dispatcher::set_default(&dispatch); handler .handle_control_frame(&ctx, route_open_frame(corr, &module_id, project_root)) .await @@ -6519,6 +6671,175 @@ mod tests { .contains("state=running, enabled=true, live=false")); } + #[tokio::test] + async fn route_open_supervised_absence_emits_refusal_fields_and_counts_code() { + let registry = Arc::new(Registry::default()); + let supervisor_handle = SupervisorHandle::new(); + let supervisor = + Supervisor::new(Arc::clone(®istry), RestartPolicy::new(0, Duration::ZERO)) + .with_handle(supervisor_handle.clone()) + .with_connection_file_path(std::env::temp_dir().join(format!( + "subc-route-open-refusal-info-{}", + std::process::id() + ))); + let module = supervisor + .supervise_configured( + ModuleSpec { + module_id: "warming".to_string(), + program: fake_aft_stub_path(), + args: Vec::new(), + env: Vec::new(), + reserved: false, + reserved_prefixes: Vec::new(), + }, + true, + ) + .unwrap(); + assert_eq!(module.state().unwrap(), ModuleState::Running); + + let handler = ControlHandler::new(Arc::clone(®istry)).with_supervisor(supervisor_handle); + assert!(handler + .counters() + .snapshot() + .get("route_open_refused_by_code") + .is_none()); + let capture = EventCapture::default(); + let _subscriber = + tracing::subscriber::set_default(tracing_subscriber::registry().with(capture.clone())); + let (ctx, _rx) = route_ctx(ConnectionId::new(94)); + let response = handler + .handle_control_frame( + &ctx, + route_open_frame(394, "warming", unique_project_root("refusal-info")), + ) + .await + .unwrap(); + module.stop().await.unwrap(); + + assert_eq!(parse_error(&response[0])["code"], "module_warming"); + let event = capture + .events() + .into_iter() + .find(|event| { + event.target == "subc_core::control" + && event.fields.get("code") == Some(&"\"module_warming\"".to_string()) + }) + .expect("route.open refusal event"); + assert_eq!( + event.fields.get("module_id"), + Some(&"\"warming\"".to_string()) + ); + assert_eq!(event.fields.get("connection_id"), Some(&"94".to_string())); + assert_eq!(event.fields.get("state"), Some(&"running".to_string())); + assert_eq!(event.fields.get("enabled"), Some(&"true".to_string())); + assert_eq!(event.fields.get("live"), Some(&"false".to_string())); + assert_eq!( + handler.counters().snapshot()["route_open_refused_by_code"], + json!({ "module_warming": 1 }) + ); + } + + #[tokio::test(flavor = "current_thread")] + async fn route_open_unknown_module_escapes_target_module_id() { + let handler = ControlHandler::new(Arc::new(Registry::default())); + let capture = EventCapture::default(); + let _subscriber = + tracing::subscriber::set_default(tracing_subscriber::registry().with(capture.clone())); + let hostile_module_id = "\u{1b}]52;c;AAAA\u{07}"; + let (ctx, _rx) = route_ctx(ConnectionId::new(95)); + let response = handler + .handle_control_frame( + &ctx, + route_open_frame( + 395, + hostile_module_id, + unique_project_root("hostile-target-module-id"), + ), + ) + .await + .unwrap(); + + assert_eq!(parse_error(&response[0])["code"], "unknown_module"); + let event = capture + .events() + .into_iter() + .find(|event| { + event.target == "subc_core::control" + && event.fields.get("code") == Some(&"\"unknown_module\"".to_string()) + }) + .expect("route.open unknown-module refusal event"); + let logged = event.fields.get("module_id").expect("module_id field"); + assert!(!logged.bytes().any(|byte| byte < 0x20)); + assert_eq!(logged, r#""\u{1b}]52;c;AAAA\u{7}""#); + } + + #[tokio::test(flavor = "current_thread")] + async fn route_open_module_rejection_uses_daemon_counter_key() { + let registry = Arc::new(Registry::default()); + let forwarding = Arc::new(ForwardingTable::default()); + let handler = + ControlHandler::with_forwarding(Arc::clone(®istry), Arc::clone(&forwarding)); + let module_connection = ConnectionId::new(95); + let (module_ctx, mut module_rx) = route_ctx(module_connection); + handler + .handle_control_frame(&module_ctx, hello_frame("aft", PROTOCOL_VERSION, 395)) + .await + .unwrap(); + + let client_connection = ConnectionId::new(96); + let (client_ctx, _client_rx) = route_ctx(client_connection); + let capture = EventCapture::default(); + let _subscriber = + tracing::subscriber::set_default(tracing_subscriber::registry().with(capture.clone())); + let (route_task, bind) = relay_route_open( + &handler, + client_connection, + &client_ctx.egress, + &mut module_rx, + 396, + "aft", + "hostile-module-code", + ) + .await; + let hostile_code = "\u{1b}]52;c;AAAA\u{07}"; + let rejection = Frame::build( + FrameType::Error, + control_flags(), + 0, + 0, + bind.header.corr, + serde_json::to_vec(&ErrorBody::new(hostile_code, "module refused route.bind")).unwrap(), + ) + .unwrap(); + handler + .handle_control_frame(&module_ctx, rejection) + .await + .unwrap(); + + let response = route_task.await.unwrap(); + assert_eq!(parse_error(&response[0])["code"], hostile_code); + let counters = handler.counters().snapshot(); + assert_eq!( + counters["route_open_refused_by_code"], + json!({ "module_rejected": 1 }) + ); + assert!(counters["route_open_refused_by_code"] + .get(hostile_code) + .is_none()); + + let event = capture + .events() + .into_iter() + .find(|event| { + event.target == "subc_core::control" + && event.fields.get("code") == Some(&"\"module_rejected\"".to_string()) + }) + .expect("route.open module-rejection refusal event"); + let logged = event.fields.get("module_code").expect("module_code field"); + assert!(!logged.bytes().any(|byte| byte < 0x20)); + assert_eq!(logged, r#""\u{1b}]52;c;AAAA\u{7}""#); + } + #[tokio::test] async fn route_open_keeps_failed_unregistered_supervised_module_unavailable() { let registry = Arc::new(Registry::default()); diff --git a/crates/subc-core/src/observability.rs b/crates/subc-core/src/observability.rs index b057aa43..b7d42d24 100644 --- a/crates/subc-core/src/observability.rs +++ b/crates/subc-core/src/observability.rs @@ -12,6 +12,23 @@ use tracing::info; use crate::registry::ConnectionId; +const ROUTE_OPEN_REFUSAL_COUNTER_CODES: &[&str] = &[ + "module_warming", + "target_unavailable", + "module_removed", + "unknown_module", + "module_reloading", + "op_not_allowed", + "bad_consumer_identity", + "capability_forbidden", + "admission_facts_not_permitted", + "admission_facts_target_not_allowed", + "route_limit", + "forwarding_error", + "module_timeout", + "module_rejected", +]; + /// Shared count of authenticated socket connections accepted by the daemon. #[derive(Debug, Clone, Default)] pub struct ConnectedClients { @@ -68,6 +85,7 @@ pub struct DaemonCounters { // Per-module maps and the rate window are daemon-lifetime diagnostics only: // they deliberately reset on restart instead of becoming durable daemon state. module_frames_dropped_no_route_by_module: Arc>>, + route_open_refused_by_code: Arc>>, module_frames_dropped_no_route_window: Arc>, module_requests_dropped_stale_route: Arc, client_frames_dropped_stale_route: Arc, @@ -170,11 +188,16 @@ impl DaemonCounters { "module_frames_dropped_no_route_nonzero_minutes_last_10m".into(), drop_window.nonzero_minutes_last_10m(now).into(), ); - insert_nonempty_module_counts( + insert_nonempty_counts( &mut snapshot, "module_frames_dropped_no_route_by_module", &self.module_frames_dropped_no_route_by_module, ); + insert_nonempty_counts( + &mut snapshot, + "route_open_refused_by_code", + &self.route_open_refused_by_code, + ); snapshot.insert( "module_requests_dropped_stale_route".into(), self.module_requests_dropped_stale_route @@ -205,7 +228,7 @@ impl DaemonCounters { .load(Ordering::Relaxed) .into(), ); - insert_nonempty_module_counts( + insert_nonempty_counts( &mut snapshot, "goodbye_relay_module_dropped_by_module", &self.goodbye_relay_module_dropped_by_module, @@ -229,7 +252,7 @@ impl DaemonCounters { self.module_frames_dropped_no_route .fetch_add(1, Ordering::Relaxed); if let Some(module_id) = module_id { - increment_module_count(&self.module_frames_dropped_no_route_by_module, module_id); + increment_keyed_count(&self.module_frames_dropped_no_route_by_module, module_id); } self.module_frames_dropped_no_route_window .lock() @@ -237,6 +260,11 @@ impl DaemonCounters { .record(tokio::time::Instant::now()); } + pub(crate) fn increment_route_open_refused(&self, code: &'static str) { + debug_assert!(ROUTE_OPEN_REFUSAL_COUNTER_CODES.contains(&code)); + increment_keyed_count(&self.route_open_refused_by_code, code); + } + pub(crate) fn increment_module_requests_dropped_stale_route(&self) { self.module_requests_dropped_stale_route .fetch_add(1, Ordering::Relaxed); @@ -261,7 +289,7 @@ impl DaemonCounters { self.goodbye_relay_module_dropped .fetch_add(1, Ordering::Relaxed); if let Some(module_id) = module_id { - increment_module_count(&self.goodbye_relay_module_dropped_by_module, module_id); + increment_keyed_count(&self.goodbye_relay_module_dropped_by_module, module_id); } } @@ -276,21 +304,21 @@ impl DaemonCounters { } } -fn increment_module_count(counts: &Mutex>, module_id: &str) { +fn increment_keyed_count(counts: &Mutex>, key: &str) { *counts .lock() - .expect("module drop-count mutex poisoned") - .entry(module_id.to_string()) + .expect("keyed counter mutex poisoned") + .entry(key.to_string()) .or_default() += 1; } -fn insert_nonempty_module_counts( +fn insert_nonempty_counts( snapshot: &mut serde_json::Map, key: &str, counts: &Mutex>, ) { let counts: MutexGuard<'_, HashMap> = - counts.lock().expect("module drop-count mutex poisoned"); + counts.lock().expect("keyed counter mutex poisoned"); if !counts.is_empty() { snapshot.insert(key.to_string(), json!(&*counts)); }