From 2b6e9b3af8b8f76795eba91ab5b520118aaae6ae Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Fri, 31 Jul 2026 21:39:05 -0600 Subject: [PATCH] fix(spl): make the supervised service observable at -v MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The verbose and debug flags were parsed and never read. A complete successful lifecycle — start, listen open, tunnel brokered, bytes both ways, posture flip, clean shutdown — produced zero bytes of output, while the Python it replaces logged thirty-three lifecycle messages. The supervisor launches this service as `journal spl -v`, so an operator diagnosing a live link had nothing to read. Wraps the callosum emitter rather than sprinkling log statements. The callosum vocabulary IS the lifecycle, so operator output and emitted events stay in lockstep by construction: a future event is observable the moment it is emitted, with nothing to remember. `health` is excluded at -v and included at -d, matching the Python, which did not log it either — it fires on every state change, every tunnel outcome and every thirty seconds, and would bury the transitions an operator is reading for. Verified live: the reconnect loop this immediately exposed turned out to be two daemons sharing one instance id in the test rig, which is precisely the class of problem the silence was hiding. --- core/crates/solstone-core-spl/src/callosum.rs | 130 ++++++++++++++++++ core/crates/solstone-core-spl/src/lib.rs | 2 +- .../solstone-core-spl/src/service_process.rs | 37 ++++- core/crates/solstone-core/src/main.rs | 10 +- 4 files changed, 167 insertions(+), 12 deletions(-) diff --git a/core/crates/solstone-core-spl/src/callosum.rs b/core/crates/solstone-core-spl/src/callosum.rs index d2e9a3db3..e728f6fa3 100644 --- a/core/crates/solstone-core-spl/src/callosum.rs +++ b/core/crates/solstone-core-spl/src/callosum.rs @@ -13,3 +13,133 @@ pub trait CallosumEmit: Send + Sync { /// Emits one named event with a JSON object payload. fn emit(&self, event: &'static str, payload: Value); } + +/// Operator verbosity for the supervised service. +/// +/// `journal spl` is launched by the supervisor as `journal spl -v`, and the +/// Python it replaces logged its whole lifecycle at that level. Silence is a +/// regression an operator only discovers while diagnosing a live link. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub enum Verbosity { + /// Only exceptional conditions reach stderr. + #[default] + Quiet, + /// Lifecycle transitions reach stderr. + Verbose, + /// Lifecycle transitions plus the periodic health snapshot. + Debug, +} + +impl Verbosity { + /// Resolves the level from the two CLI flags. `--debug` implies `--verbose`. + #[must_use] + pub fn from_flags(verbose: bool, debug: bool) -> Self { + if debug { + Self::Debug + } else if verbose { + Self::Verbose + } else { + Self::Quiet + } + } + + #[must_use] + fn logs(self, event: &str) -> bool { + match self { + Self::Quiet => false, + // `health` fires on every state change, every tunnel outcome and + // every 30s. Python did not log it either; at -v it would bury the + // transitions an operator is actually reading for. + Self::Verbose => event != "health", + Self::Debug => true, + } + } +} + +/// Mirrors every callosum event to stderr at the requested verbosity. +/// +/// Wrapping the emitter rather than sprinkling log statements keeps operator +/// output and the callosum vocabulary in lockstep by construction: a new +/// lifecycle event is observable the moment it is emitted, with nothing to +/// remember. +pub struct LoggingEmit { + inner: std::sync::Arc, + verbosity: Verbosity, +} + +impl LoggingEmit { + #[must_use] + pub fn new(inner: std::sync::Arc, verbosity: Verbosity) -> Self { + Self { inner, verbosity } + } +} + +impl CallosumEmit for LoggingEmit { + fn emit(&self, event: &'static str, payload: Value) { + if self.verbosity.logs(event) { + match payload.as_object() { + Some(fields) if !fields.is_empty() => { + let rendered = fields + .iter() + .map(|(key, value)| format!("{key}={value}")) + .collect::>() + .join(" "); + eprintln!("spl service: {event} {rendered}"); + } + _ => eprintln!("spl service: {event}"), + } + } + self.inner.emit(event, payload); + } +} + +#[cfg(test)] +mod logging_tests { + use super::*; + use std::sync::Mutex; + + #[derive(Default)] + struct Recorder(Mutex>); + + impl CallosumEmit for Recorder { + fn emit(&self, event: &'static str, _payload: Value) { + if let Ok(mut seen) = self.0.lock() { + seen.push(event); + } + } + } + + #[test] + fn debug_implies_verbose_and_quiet_is_the_default() { + assert_eq!(Verbosity::from_flags(false, false), Verbosity::Quiet); + assert_eq!(Verbosity::from_flags(true, false), Verbosity::Verbose); + assert_eq!(Verbosity::from_flags(false, true), Verbosity::Debug); + assert_eq!(Verbosity::from_flags(true, true), Verbosity::Debug); + assert_eq!(Verbosity::default(), Verbosity::Quiet); + } + + #[test] + fn verbose_reports_transitions_but_not_the_periodic_health_snapshot() { + assert!(Verbosity::Verbose.logs("connected")); + assert!(Verbosity::Verbose.logs("tunnel_pair")); + assert!(!Verbosity::Verbose.logs("health")); + assert!(Verbosity::Debug.logs("health")); + assert!(!Verbosity::Quiet.logs("connected")); + } + + #[test] + fn wrapping_never_swallows_an_event_at_any_verbosity() { + for verbosity in [Verbosity::Quiet, Verbosity::Verbose, Verbosity::Debug] { + let recorder = std::sync::Arc::new(Recorder::default()); + let emitter = LoggingEmit::new( + std::sync::Arc::clone(&recorder) as std::sync::Arc, + verbosity, + ); + emitter.emit("connecting", serde_json::json!({})); + emitter.emit("tunnel_pair", serde_json::json!({"tunnel_id": "t-1"})); + emitter.emit("health", serde_json::json!({"state": "connected"})); + let seen = recorder.0.lock().expect("recorder lock"); + assert_eq!(seen.as_slice(), ["connecting", "tunnel_pair", "health"]); + } + } +} diff --git a/core/crates/solstone-core-spl/src/lib.rs b/core/crates/solstone-core-spl/src/lib.rs index fe10b950e..eeb333783 100644 --- a/core/crates/solstone-core-spl/src/lib.rs +++ b/core/crates/solstone-core-spl/src/lib.rs @@ -28,7 +28,7 @@ mod ws_buffer; mod ws_sink; pub use admission::RelayAdmissionGate; -pub use callosum::CallosumEmit; +pub use callosum::{CallosumEmit, LoggingEmit, Verbosity}; pub use health::{ LINK_HEALTH_EVENT, OFFLINE_TUNNEL_REASONS, REASON_HOME_MISSING_MOBILE, REASON_LOCAL_PRIVATE_LISTENER_UNREACHABLE, REASON_RELAY_ADMISSION_SATURATED, diff --git a/core/crates/solstone-core-spl/src/service_process.rs b/core/crates/solstone-core-spl/src/service_process.rs index 8ef84f9a2..a2a843f15 100644 --- a/core/crates/solstone-core-spl/src/service_process.rs +++ b/core/crates/solstone-core-spl/src/service_process.rs @@ -32,7 +32,9 @@ use tokio::{ use crate::{ CallosumEmit, LinkServiceTokenRead, LinkStateRead, LoopbackConnect, LoopbackDialer, LoopbackStream, RelayClient, RelayClientConfig, RelayError, RelayServiceToken, ServiceDeps, - ServiceError, ServicePoll, ServiceToken, load_link_service_token, load_link_state, run_service, + ServiceError, ServicePoll, ServiceToken, + callosum::{LoggingEmit, Verbosity}, + load_link_service_token, load_link_state, run_service, }; const DEFAULT_RELAY_ENDPOINT: &str = "https://link.solstone.app"; @@ -76,20 +78,29 @@ impl NativeServiceError { /// /// Returns a class-only error if the runtime cannot start or the supervised /// service exits unexpectedly. -pub fn run_native_service(journal_root: PathBuf) -> Result<(), NativeServiceError> { +pub fn run_native_service( + journal_root: PathBuf, + verbosity: Verbosity, +) -> Result<(), NativeServiceError> { let runtime = tokio::runtime::Builder::new_multi_thread() .enable_all() .thread_name("spl-service") .build() .map_err(|_| NativeServiceError::Runtime)?; - runtime.block_on(run_native_service_async(journal_root)) + runtime.block_on(run_native_service_async(journal_root, verbosity)) } -async fn run_native_service_async(journal_root: PathBuf) -> Result<(), NativeServiceError> { +async fn run_native_service_async( + journal_root: PathBuf, + verbosity: Verbosity, +) -> Result<(), NativeServiceError> { let callosum = CallosumOutput::start(journal_root.join("health").join("callosum.sock")); + if verbosity != Verbosity::Quiet { + eprintln!("spl service: starting; watching link posture"); + } let (shutdown_send, shutdown_receive) = watch::channel(false); let signal_task = tokio::spawn(wait_for_shutdown_signal(shutdown_send)); - let mut deps = ProcessServiceDeps::new(journal_root, callosum, shutdown_receive); + let mut deps = ProcessServiceDeps::new(journal_root, callosum, verbosity, shutdown_receive); let result = run_service(&mut deps).await.map_err(classify_service_error); signal_task.abort(); @@ -129,6 +140,7 @@ async fn wait_for_shutdown_signal(shutdown: watch::Sender) { struct ProcessServiceDeps { journal_root: PathBuf, callosum: Arc, + verbosity: Verbosity, shutdown: watch::Receiver, } @@ -136,11 +148,13 @@ impl ProcessServiceDeps { fn new( journal_root: PathBuf, callosum: Arc, + verbosity: Verbosity, shutdown: watch::Receiver, ) -> Self { Self { journal_root, callosum, + verbosity, shutdown, } } @@ -230,7 +244,10 @@ impl ServiceDeps for ProcessServiceDeps { dispatch_read_deadline: DISPATCH_READ_DEADLINE, global_admission_ceiling: GLOBAL_ADMISSION_CEILING, }, - Arc::clone(&self.callosum) as Arc, + Arc::new(LoggingEmit::new( + Arc::clone(&self.callosum) as Arc, + self.verbosity, + )) as Arc, Arc::new(LocalLoopbackDialer), ); let running_client = client.clone(); @@ -1024,7 +1041,12 @@ mod tests { let (shutdown_send, shutdown_receive) = tokio::sync::watch::channel(false); drop(shutdown_send); let callosum = super::CallosumOutput::inactive(); - ProcessServiceDeps::new(root.to_path_buf(), callosum, shutdown_receive) + ProcessServiceDeps::new( + root.to_path_buf(), + callosum, + crate::callosum::Verbosity::Quiet, + shutdown_receive, + ) } #[tokio::test] @@ -1070,6 +1092,7 @@ mod tests { let mut deps = ProcessServiceDeps::new( journal.path().to_path_buf(), Arc::clone(&callosum), + crate::callosum::Verbosity::Quiet, shutdown_receive, ); diff --git a/core/crates/solstone-core/src/main.rs b/core/crates/solstone-core/src/main.rs index 12f9865ea..1709167e0 100644 --- a/core/crates/solstone-core/src/main.rs +++ b/core/crates/solstone-core/src/main.rs @@ -5,7 +5,8 @@ use std::process::ExitCode; use std::{env, ffi::OsStr, path::PathBuf}; use solstone_core_cli::{ - Command, IndexerOptions, JournalPathOptions, SplCommand, USAGE, evaluate_args, version_line, + Command, IndexerOptions, JournalPathOptions, ServiceOptions, SplCommand, USAGE, evaluate_args, + version_line, }; use solstone_core_indexer_store::db::reset_index; use solstone_core_indexer_store::scan::{ @@ -66,11 +67,11 @@ fn main() -> ExitCode { fn run_spl_process(command: SplCommand) -> ExitCode { match command { - SplCommand::Service(_) => run_spl_service(), + SplCommand::Service(options) => run_spl_service(options), } } -fn run_spl_service() -> ExitCode { +fn run_spl_service(options: ServiceOptions) -> ExitCode { let journal = match resolve_process_journal_path() { Ok(journal) => journal, Err(error) => { @@ -78,7 +79,8 @@ fn run_spl_service() -> ExitCode { return ExitCode::from(EXIT_TEMPFAIL); } }; - match solstone_core_spl::run_native_service(journal.path) { + let verbosity = solstone_core_spl::Verbosity::from_flags(options.verbose, options.debug); + match solstone_core_spl::run_native_service(journal.path, verbosity) { Ok(()) => ExitCode::SUCCESS, Err(error) => { eprintln!("spl service failed: {}", error.class()); -- 2.51.2