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());