diff --git a/Cargo.lock b/Cargo.lock index 74073a90..6dfd0248 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -254,7 +254,7 @@ dependencies = [ [[package]] name = "ck-bus" -version = "0.1.11" +version = "0.1.13" dependencies = [ "async-nats", "async-trait", @@ -1807,7 +1807,7 @@ dependencies = [ [[package]] name = "subc-client-rs" -version = "0.18.5" +version = "0.19.0" dependencies = [ "async-trait", "serde", @@ -1822,7 +1822,7 @@ dependencies = [ [[package]] name = "subc-control" -version = "0.19.0" +version = "0.20.0" dependencies = [ "serde", "serde_json", @@ -1831,7 +1831,7 @@ dependencies = [ [[package]] name = "subc-core" -version = "0.20.22" +version = "0.20.23" dependencies = [ "base64", "cortexkit-log", @@ -1857,7 +1857,7 @@ dependencies = [ [[package]] name = "subc-daemon" -version = "0.21.3" +version = "0.21.4" dependencies = [ "cortexkit-log", "cortexkit-paths", diff --git a/crates/ck-bus/Cargo.toml b/crates/ck-bus/Cargo.toml index 60d95aec..c2d665c1 100644 --- a/crates/ck-bus/Cargo.toml +++ b/crates/ck-bus/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "ck-bus" -version = "0.1.11" +version = "0.1.13" edition = "2021" publish = false description = "Supervised owner of the CortexKit NATS message plane." @@ -28,7 +28,7 @@ data-encoding = "=2.11.1" nkeys = "=0.4.5" serde_json = "1" sha2 = "0.10" -subc-client-rs = { path = "../subc-client-rs", version = "0.18" } +subc-client-rs = { path = "../subc-client-rs", version = "0.19" } subc-protocol = { path = "../subc-protocol", version = "0.25.0" } tokio = { version = "1", features = ["macros", "process", "rt-multi-thread", "sync", "time"] } @@ -36,7 +36,7 @@ tokio = { version = "1", features = ["macros", "process", "rt-multi-thread", "sy futures-util = "0.3" tempfile = "3" tokio = { version = "1", features = ["io-util", "net"] } -subc-control = { path = "../subc-control", version = "0.19" } +subc-control = { path = "../subc-control", version = "0.20" } subc-daemon = { path = "../subc-daemon", version = "0.21.0", features = ["test-support"] } subc-jsonc = { path = "../subc-jsonc", version = "0.1.0" } subc-transport = { path = "../subc-transport", version = "0.7.0" } diff --git a/crates/subc-client-rs/Cargo.toml b/crates/subc-client-rs/Cargo.toml index b9304514..e84331ab 100644 --- a/crates/subc-client-rs/Cargo.toml +++ b/crates/subc-client-rs/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "subc-client-rs" -version = "0.18.5" +version = "0.19.0" edition = "2021" publish = true description = "Shared serve + consume client for Rust subc modules." @@ -11,7 +11,7 @@ repository = "https://github.com/cortexkit/subconscious" async-trait = "0.1" serde = { version = "1", features = ["derive"] } serde_json = "1" -subc-control = { path = "../subc-control", version = "0.19" } +subc-control = { path = "../subc-control", version = "0.20" } subc-protocol = { path = "../subc-protocol", version = "0.25.0" } subc-transport = { path = "../subc-transport", version = "0.7" } tokio = { version = "1", features = ["io-util", "macros", "net", "rt", "sync", "time"] } diff --git a/crates/subc-control/Cargo.toml b/crates/subc-control/Cargo.toml index 278da801..075a893e 100644 --- a/crates/subc-control/Cargo.toml +++ b/crates/subc-control/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "subc-control" -version = "0.19.0" +version = "0.20.0" edition = "2021" publish = true description = "Client-facing subc control-plane wire shapes." diff --git a/crates/subc-control/src/lib.rs b/crates/subc-control/src/lib.rs index daedb9cd..7f3dd531 100644 --- a/crates/subc-control/src/lib.rs +++ b/crates/subc-control/src/lib.rs @@ -609,6 +609,39 @@ pub struct SupervisorObservedProcess { pub running_image: RunningImageAgreement, } +/// Independent comparisons of configured path and running image at list time. +/// An absent verdict means the daemon predates this field, not agreement. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct PendingReloadVerdict { + pub path: ReloadPathAgreement, + pub image: RunningImageAgreement, +} + +/// Whether the running process was spawned from the currently configured program. +#[derive(Debug, Clone, PartialEq)] +pub enum ReloadPathAgreement { + Match, + Mismatch { + configured: PathBuf, + spawned_from: PathBuf, + }, + Unavailable { + reason: ReloadPathUnavailableReason, + }, + Unknown { + tag: String, + body: OrderedJsonObject, + }, +} + +open_string_enum! { + /// Why configured and spawned paths cannot be compared. + ReloadPathUnavailableReason { + NotRunning => "not_running", + SpawnedPathUnavailable => "spawned_path_unavailable", + } +} + /// Daemon provenance paired with its runtime process observation. #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct SupervisorDaemonProvenance { @@ -779,6 +812,19 @@ enum RunningImageAgreementWire { }, } +#[derive(Debug, Serialize, Deserialize)] +#[serde(tag = "status", rename_all = "snake_case")] +enum ReloadPathAgreementWire { + Match, + Mismatch { + configured: PathBuf, + spawned_from: PathBuf, + }, + Unavailable { + reason: ReloadPathUnavailableReason, + }, +} + #[derive(Debug, Serialize, Deserialize)] #[serde(tag = "method", rename_all = "snake_case")] enum RunningImageEvidenceWire { @@ -1185,6 +1231,61 @@ impl Serialize for RunningImageAgreement { } } +impl Serialize for ReloadPathAgreement { + fn serialize(&self, serializer: S) -> Result + where + S: Serializer, + { + match self { + Self::Match => ReloadPathAgreementWire::Match.serialize(serializer), + Self::Mismatch { + configured, + spawned_from, + } => ReloadPathAgreementWire::Mismatch { + configured: configured.clone(), + spawned_from: spawned_from.clone(), + } + .serialize(serializer), + Self::Unavailable { reason } => ReloadPathAgreementWire::Unavailable { + reason: reason.clone(), + } + .serialize(serializer), + Self::Unknown { body, .. } => body.serialize(serializer), + } + } +} + +impl<'de> Deserialize<'de> for ReloadPathAgreement { + fn deserialize(deserializer: D) -> Result + where + D: Deserializer<'de>, + { + let (tag, body) = read_tagged(deserializer, "status")?; + match tag.as_str() { + "match" => Ok(Self::Match), + "mismatch" => { + match serde_json::from_value(body.into_value()).map_err(D::Error::custom)? { + ReloadPathAgreementWire::Mismatch { + configured, + spawned_from, + } => Ok(Self::Mismatch { + configured, + spawned_from, + }), + _ => unreachable!(), + } + } + "unavailable" => match serde_json::from_value(body.into_value()) + .map_err(D::Error::custom)? + { + ReloadPathAgreementWire::Unavailable { reason } => Ok(Self::Unavailable { reason }), + _ => unreachable!(), + }, + _ => Ok(Self::Unknown { tag, body }), + } + } +} + impl<'de> Deserialize<'de> for RunningImageAgreement { fn deserialize(deserializer: D) -> Result where @@ -1708,7 +1809,7 @@ pub enum ModuleProtocol { None, } -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct SupervisorEntry { pub module_id: String, pub state: String, @@ -1729,6 +1830,10 @@ pub struct SupervisorEntry { #[serde(default)] pub protocol: ModuleProtocol, pub health: SupervisorHealthStatus, + /// Computed from the stored launch spec and observed process at list time; + /// None means an older daemon did not report this comparison. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub pending_reload: Option, /// When the daemon last collected this module's health, as unix /// milliseconds. Absent means NEVER PROBED (a module inside its first probe /// window, whose `health` is therefore `Unknown` rather than good), not diff --git a/crates/subc-control/tests/golden/supervisor_entry_reload_verdict_states.json b/crates/subc-control/tests/golden/supervisor_entry_reload_verdict_states.json new file mode 100644 index 00000000..f875bdfb --- /dev/null +++ b/crates/subc-control/tests/golden/supervisor_entry_reload_verdict_states.json @@ -0,0 +1,102 @@ +[ + { + "enabled": true, + "health": "degraded", + "last_exit_ms": 1700000000123, + "last_probe_ms": 1700000000000, + "live": true, + "max_restarts": 3, + "module_id": "aft-tools", + "pending_reload": { + "image": { + "evidence": { + "digest": "old", + "method": "linux_proc_sha256" + }, + "status": "match" + }, + "path": { + "configured": "/bin/new", + "spawned_from": "/bin/old", + "status": "mismatch" + } + }, + "protocol": "subc", + "restart_count": 2, + "state": "running" + }, + { + "enabled": true, + "health": "degraded", + "last_exit_ms": 1700000000123, + "last_probe_ms": 1700000000000, + "live": true, + "max_restarts": 3, + "module_id": "aft-tools", + "pending_reload": { + "image": { + "disk": { + "digest": "new", + "method": "linux_proc_sha256" + }, + "running": { + "digest": "old", + "method": "linux_proc_sha256" + }, + "status": "mismatch" + }, + "path": { + "status": "match" + } + }, + "protocol": "subc", + "restart_count": 2, + "state": "running" + }, + { + "enabled": true, + "health": "degraded", + "last_exit_ms": 1700000000123, + "last_probe_ms": 1700000000000, + "live": true, + "max_restarts": 3, + "module_id": "aft-tools", + "pending_reload": { + "image": { + "reason": "not_running", + "status": "unavailable" + }, + "path": { + "reason": "not_running", + "status": "unavailable" + } + }, + "protocol": "subc", + "restart_count": 2, + "state": "running" + }, + { + "enabled": true, + "health": "degraded", + "last_exit_ms": 1700000000123, + "last_probe_ms": 1700000000000, + "live": true, + "max_restarts": 3, + "module_id": "aft-tools", + "pending_reload": { + "image": { + "evidence": { + "digest": "same", + "method": "linux_proc_sha256" + }, + "status": "match" + }, + "path": { + "status": "match" + } + }, + "protocol": "subc", + "restart_count": 2, + "state": "running" + } +] diff --git a/crates/subc-control/tests/golden_json.rs b/crates/subc-control/tests/golden_json.rs index e5ecb80b..41d97654 100644 --- a/crates/subc-control/tests/golden_json.rs +++ b/crates/subc-control/tests/golden_json.rs @@ -5,8 +5,9 @@ use serde_json::Value; use subc_control::{ CatalogEntry, ClientControlPush, ClientControlRequest, ClientControlResponse, ConsumerIdentity, DaemonBuildProvenance, DaemonObservedProcess, LiveSpawn, ModuleDeclaredProvenance, - ModuleProtocol, NotReadyReason, PollKind, RouteCloseReason, RunningImageAgreement, - RunningImageEvidence, SpawnCursor, SpawnEvent, SpawnEventKind, SpawnSnapshot, + ModuleProtocol, NotReadyReason, PendingReloadVerdict, PollKind, ReloadPathAgreement, + ReloadPathUnavailableReason, RouteCloseReason, RunningImageAgreement, RunningImageEvidence, + RunningImageUnavailableReason, SpawnCursor, SpawnEvent, SpawnEventKind, SpawnSnapshot, StderrCaptureState, StderrTail, StderrTailEntry, SupervisorDaemonProvenance, SupervisorEntry, SupervisorHealthEntry, SupervisorHealthStatus, SupervisorModuleProvenance, SupervisorObservedProcess, SupervisorRescanResult, SupervisorRoute, SupervisorRouteConsumer, @@ -74,6 +75,51 @@ fn control_wire_shapes_match_golden_json_and_round_trip() { }, ); assert_golden("supervisor_entry", &supervisor_entry()); + let evidence = |digest: &str| RunningImageEvidence::LinuxProcSha256 { + digest: digest.to_string(), + }; + let verdicts = [ + PendingReloadVerdict { + path: ReloadPathAgreement::Mismatch { + configured: "/bin/new".into(), + spawned_from: "/bin/old".into(), + }, + image: RunningImageAgreement::Match { + evidence: evidence("old"), + }, + }, + PendingReloadVerdict { + path: ReloadPathAgreement::Match, + image: RunningImageAgreement::Mismatch { + running: evidence("old"), + disk: evidence("new"), + }, + }, + PendingReloadVerdict { + path: ReloadPathAgreement::Unavailable { + reason: ReloadPathUnavailableReason::NotRunning, + }, + image: RunningImageAgreement::Unavailable { + reason: RunningImageUnavailableReason::NotRunning, + }, + }, + PendingReloadVerdict { + path: ReloadPathAgreement::Match, + image: RunningImageAgreement::Match { + evidence: evidence("same"), + }, + }, + ]; + assert_golden( + "supervisor_entry_reload_verdict_states", + &verdicts + .into_iter() + .map(|verdict| SupervisorEntry { + pending_reload: Some(verdict), + ..supervisor_entry() + }) + .collect::>(), + ); assert_golden( "supervisor_spawn_event", &SpawnEvent { @@ -891,6 +937,7 @@ fn supervisor_entry() -> SupervisorEntry { live: true, protocol: ModuleProtocol::Subc, health: SupervisorHealthStatus::Degraded, + pending_reload: None, last_probe_ms: Some(1_700_000_000_000), last_exit_code: None, last_exit_signal: None, @@ -913,6 +960,31 @@ fn supervisor_entry() -> SupervisorEntry { } } +#[test] +fn supervisor_entry_without_reload_verdict_is_unknown_not_clear() { + let entry: SupervisorEntry = serde_json::from_str( + r#"{"module_id":"legacy","state":"running","enabled":true,"live":true,"health":"unknown"}"#, + ) + .expect("old daemon supervisor entry decodes"); + assert_eq!(entry.pending_reload, None); +} + +#[test] +fn reload_path_unknown_variant_and_reason_preserve_forward_wire() { + let future = serde_json::json!({ + "path": {"status": "future_path", "detail": {"x": 1}}, + "image": {"status": "unavailable", "reason": "future_probe"} + }); + let decoded: PendingReloadVerdict = serde_json::from_value(future.clone()).unwrap(); + assert!( + matches!(decoded.path, ReloadPathAgreement::Unknown { ref tag, .. } if tag == "future_path") + ); + assert!( + matches!(decoded.image, RunningImageAgreement::Unavailable { reason: RunningImageUnavailableReason::Unknown(ref reason) } if reason == "future_probe") + ); + assert_eq!(serde_json::to_value(decoded).unwrap(), future); +} + /// The windowed budget on the wire: the same count and cap, plus the span they /// are counted over. Pinned as its own case because the two shapes mean /// different things -- "2 of 3 crashes ever" vs "2 of 3 in the last ten diff --git a/crates/subc-core/Cargo.toml b/crates/subc-core/Cargo.toml index 1d436abb..a16dbbf8 100644 --- a/crates/subc-core/Cargo.toml +++ b/crates/subc-core/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "subc-core" -version = "0.20.22" +version = "0.20.23" edition = "2021" license = "MIT" publish = false diff --git a/crates/subc-core/src/bin/ck.rs b/crates/subc-core/src/bin/ck.rs index 88e04845..2a06f2de 100644 --- a/crates/subc-core/src/bin/ck.rs +++ b/crates/subc-core/src/bin/ck.rs @@ -5608,15 +5608,19 @@ fn print_module_table(modules: &[Value], verbose: bool) { .iter() .map(|module| { vec![ - display_field(module, "module_id"), - display_field(module, "state"), + terminal_safe_string(&display_field(module, "module_id")), + terminal_safe_string(&display_field(module, "state")), enabled_word(module.get("enabled").and_then(Value::as_bool)), live_word(module), - human_health_status(&display_field(module, "health")), + terminal_safe_string(&human_health_status(&display_field(module, "health"))), + reload_list_marker(module).to_string(), ] }) .collect::>(); - print_table(&["module", "state", "enabled", "live", "health"], rows); + print_table( + &["module", "state", "enabled", "live", "health", "reload"], + rows, + ); return; } @@ -5624,13 +5628,105 @@ fn print_module_table(modules: &[Value], verbose: bool) { .iter() .map(|module| { vec![ - display_field(module, "module_id"), - module_status_text(module), - human_health_status(&display_field(module, "health")), + terminal_safe_string(&display_field(module, "module_id")), + terminal_safe_string(&module_status_text(module)), + terminal_safe_string(&human_health_status(&display_field(module, "health"))), + reload_list_marker(module).to_string(), ] }) .collect::>(); - print_table(&["module", "status", "health"], rows); + print_table(&["module", "status", "health", "reload"], rows); +} + +fn reload_list_marker(module: &Value) -> &'static str { + let Some(verdict) = module.get("pending_reload") else { + return "unknown"; + }; + let path = verdict.pointer("/path/status").and_then(Value::as_str); + let image = verdict.pointer("/image/status").and_then(Value::as_str); + let image_reason = verdict.pointer("/image/reason").and_then(Value::as_str); + match (path, image, image_reason) { + (Some("mismatch"), Some("mismatch"), _) => "pending (path+image)", + (Some("mismatch"), _, _) => "pending (path)", + (Some("match"), Some("mismatch"), _) => "pending (image)", + (Some("match"), Some("match"), _) => "nothing", + (Some("match"), Some("unavailable"), Some("unsupported_platform")) => "nothing (image n/a)", + _ => "unknown", + } +} + +fn print_reload_verdict(module: &Value) { + let Some(verdict) = module.get("pending_reload") else { + println!(" pending reload: unknown (daemon did not report a verdict)"); + return; + }; + let path = verdict.get("path"); + let path_status = path + .and_then(|value| value.get("status")) + .and_then(Value::as_str); + match path_status { + Some("mismatch") => println!( + " configured program: pending reload (configured {}, running from {})", + path.and_then(|value| value.get("configured")) + .and_then(Value::as_str) + .map(terminal_safe_string) + .unwrap_or_else(|| "unknown".into()), + path.and_then(|value| value.get("spawned_from")) + .and_then(Value::as_str) + .map(terminal_safe_string) + .unwrap_or_else(|| "unknown".into()), + ), + Some("match") => println!(" configured program: matches running process"), + Some("unavailable") => println!( + " configured program: unknown ({})", + path.and_then(|value| value.get("reason")) + .and_then(Value::as_str) + .map(terminal_safe_string) + .unwrap_or_else(|| "reason unavailable".into()), + ), + _ => println!( + " configured program: unknown ({})", + terminal_safe_string(path_status.unwrap_or("missing status")) + ), + } + println!( + " running image: {}", + image_reload_sentence(verdict.get("image")) + ); +} + +fn image_reload_sentence(image: Option<&Value>) -> String { + let image_status = image + .and_then(|value| value.get("status")) + .and_then(Value::as_str); + match image_status { + Some("mismatch") => "pending reload (file at spawned path changed since spawn)".to_string(), + Some("match") => "matches file at spawned path".to_string(), + Some("unavailable") + if image + .and_then(|value| value.get("reason")) + .and_then(Value::as_str) + == Some("unsupported_platform") => + { + "not checked on this platform (unsupported_platform)".to_string() + } + Some("unavailable") => format!( + "unknown ({})", + image + .and_then(|value| value.get("reason")) + .and_then(Value::as_str) + .map(terminal_safe_string) + .unwrap_or_else(|| "reason unavailable".into()), + ), + _ => format!( + "unknown ({})", + terminal_safe_string(image_status.unwrap_or("missing status")) + ), + } +} + +fn preview_rescan_changed_label() -> &'static str { + "would change (pending reload)" } fn print_rescan_table(result: &Value) { @@ -5650,6 +5746,7 @@ fn print_rescan_table(result: &Value) { .map(|ids| { ids.iter() .filter_map(Value::as_str) + .map(terminal_safe_string) .collect::>() .join(", ") }) @@ -5663,7 +5760,13 @@ fn print_rescan_table(result: &Value) { let warnings = result .get("capability_warnings") .and_then(Value::as_array) - .map(|items| items.iter().filter_map(Value::as_str).collect::>()) + .map(|items| { + items + .iter() + .filter_map(Value::as_str) + .map(terminal_safe_string) + .collect::>() + }) .unwrap_or_default(); if added.is_none() @@ -5688,14 +5791,7 @@ fn print_rescan_table(result: &Value) { ); } if let Some(ids) = changed { - println!( - "{}: {ids}", - if preview { - "would restart" - } else { - "restarted" - } - ); + println!("{}: {ids}", preview_rescan_changed_label()); } if let Some(ids) = enabled { println!( @@ -5724,6 +5820,7 @@ fn print_applied_rescan_table(result: &Value) { .map(|ids| { ids.iter() .filter_map(Value::as_str) + .map(terminal_safe_string) .collect::>() .join(", ") }) @@ -5748,9 +5845,19 @@ fn print_applied_rescan_table(result: &Value) { ], ]; print_table(&["change", "modules / count"], rows); + if let Some(changed) = result + .get("changed_pending_reload") + .and_then(Value::as_array) + { + if !changed.is_empty() { + println!( + "restart deferred; apply with `ck module restart ` for each changed module" + ); + } + } if let Some(warnings) = result.get("capability_warnings").and_then(Value::as_array) { for warning in warnings.iter().filter_map(Value::as_str) { - println!("warning: {warning}"); + println!("warning: {}", terminal_safe_string(warning)); } } let restart_required = module_ids("restart_required"); @@ -5767,15 +5874,15 @@ fn print_status_table( observed: Option<&Value>, verbose: bool, ) { - let module_id = display_field(module, "module_id"); + let module_id = terminal_safe_string(&display_field(module, "module_id")); let health_status = health .map(|entry| display_field(entry, "status")) .filter(|value| value != "-") .unwrap_or_else(|| display_field(module, "health")); println!( "{module_id} — {}, {}", - module_status_text(module), - health_sentence_status(&health_status) + terminal_safe_string(&module_status_text(module)), + terminal_safe_string(&health_sentence_status(&health_status)) ); let pid = observed @@ -5806,13 +5913,14 @@ fn print_status_table( let binary = observed .and_then(|value| value.get("spawned_from")) .and_then(Value::as_str) - .map(|path| display_home_path(Path::new(path))) + .map(|path| terminal_safe_string(&display_home_path(Path::new(path)))) .unwrap_or_else(|| "none".to_string()); let image = observed .and_then(|value| value.get("running_image")) .map(running_image_clause) .unwrap_or_else(|| "running image status unknown".to_string()); println!(" binary: {binary} ({image})"); + print_reload_verdict(module); if health_status != "ok" { if let Some(detail) = health.and_then(health_operator_detail) { println!(" health: {detail}"); @@ -7579,6 +7687,106 @@ mod tests { ); } + #[test] + fn rescan_preview_changed_label_pairs_with_deferred_apply_label() { + let preview = preview_rescan_changed_label(); + assert_eq!(preview, "would change (pending reload)"); + assert!(!preview.contains("restart")); + let applied = "changed-pending-reload"; + assert!(applied.contains("pending-reload")); + } + + #[test] + fn reload_list_marker_distinguishes_path_image_unknown_and_clear() { + let path = serde_json::json!({"pending_reload": {"path": {"status":"mismatch", "configured":"/new", "spawned_from":"/old"}, "image": {"status":"unavailable", "reason":"hash_failed"}}}); + assert_eq!(reload_list_marker(&path), "pending (path)"); + let image = serde_json::json!({"pending_reload": {"path": {"status":"match"}, "image": {"status":"mismatch"}}}); + assert_eq!(reload_list_marker(&image), "pending (image)"); + let unknown = serde_json::json!({"pending_reload": {"path": {"status":"unavailable", "reason":"not_running"}, "image": {"status":"unavailable", "reason":"not_running"}}}); + assert_eq!(reload_list_marker(&unknown), "unknown"); + assert_eq!(reload_list_marker(&serde_json::json!({})), "unknown"); + let clear = serde_json::json!({"pending_reload": {"path": {"status":"match"}, "image": {"status":"match"}}}); + assert_eq!(reload_list_marker(&clear), "nothing"); + } + + #[test] + fn reload_list_marker_keeps_every_unavailable_image_unknown() { + for reason in [ + "not_running", + "running_executable_unreadable", + "spawned_path_unreadable", + "hash_failed", + "process_identity_unconfirmed", + "unsupported_platform_next", + "future_probe_reason", + ] { + let module = serde_json::json!({"pending_reload": { + "path": {"status": "match"}, + "image": {"status": "unavailable", "reason": reason}, + }}); + assert_eq!(reload_list_marker(&module), "unknown", "{reason}"); + let sentence = image_reload_sentence(module.pointer("/pending_reload/image")); + assert_eq!(sentence, format!("unknown ({reason})")); + } + } + + #[test] + fn reload_list_marker_keeps_future_image_unavailability_unknown() { + let module = serde_json::json!({"pending_reload": { + "path": {"status": "match"}, + "image": {"status": "unavailable", "reason": "future_probe_reason"}, + }}); + assert_eq!(reload_list_marker(&module), "unknown"); + assert_eq!( + image_reload_sentence(module.pointer("/pending_reload/image")), + "unknown (future_probe_reason)" + ); + } + + #[test] + fn reload_list_marker_keeps_pending_path_when_image_unavailable() { + let module = serde_json::json!({"pending_reload": { + "path": {"status": "mismatch", "configured": "/new", "spawned_from": "/old"}, + "image": {"status": "unavailable", "reason": "unsupported_platform"}, + }}); + assert_eq!(reload_list_marker(&module), "pending (path)"); + assert_eq!( + image_reload_sentence(module.pointer("/pending_reload/image")), + "not checked on this platform (unsupported_platform)" + ); + } + + #[test] + fn reload_list_marker_reports_image_not_supported_when_path_matches() { + let module = serde_json::json!({"pending_reload": { + "path": {"status": "match"}, + "image": {"status": "unavailable", "reason": "unsupported_platform"}, + }}); + assert_eq!(reload_list_marker(&module), "nothing (image n/a)"); + assert_eq!( + image_reload_sentence(module.pointer("/pending_reload/image")), + "not checked on this platform (unsupported_platform)" + ); + } + + #[test] + fn reload_list_marker_keeps_unknown_path_when_image_not_supported() { + let module = serde_json::json!({"pending_reload": { + "path": {"status": "unavailable", "reason": "not_running"}, + "image": {"status": "unavailable", "reason": "unsupported_platform"}, + }}); + assert_eq!(reload_list_marker(&module), "unknown"); + assert_eq!( + image_reload_sentence(module.pointer("/pending_reload/image")), + "not checked on this platform (unsupported_platform)" + ); + let image_mismatch = serde_json::json!({"pending_reload": { + "path": {"status": "unavailable", "reason": "spawned_path_unavailable"}, + "image": {"status": "mismatch"}, + }}); + assert_eq!(reload_list_marker(&image_mismatch), "unknown"); + } + /// A module keeps its pre-adoption r1 history beside its r2 segments, so the /// merged view reads one directory holding both grammars. The r1 lines must /// be dated and interleave by time: left undated, they sort after every r2 diff --git a/crates/subc-core/tests/ck_cli.rs b/crates/subc-core/tests/ck_cli.rs index ad1e6068..16dc704a 100644 --- a/crates/subc-core/tests/ck_cli.rs +++ b/crates/subc-core/tests/ck_cli.rs @@ -1948,10 +1948,12 @@ async fn module_list_json_uses_subc_override_and_shows_stub() { let text_output = ck_with_subc(&server.connection_file_path, ["module", "list"]); assert_exit(&text_output, 0); let text_stdout = text(&text_output.stdout); - assert_eq!( - text_stdout, - "module status health \nck-list-stub running unknown\n" - ); + let expected_list = if cfg!(any(target_os = "linux", target_os = "macos")) { + "module status health reload \nck-list-stub running unknown nothing\n" + } else { + "module status health reload \nck-list-stub running unknown nothing (image n/a)\n" + }; + assert_eq!(text_stdout, expected_list); assert!( !json_stdout.contains("next:"), "JSON output must not gain human footer: {json_stdout}" @@ -1983,16 +1985,67 @@ async fn module_list_renders_status_words_not_wire_booleans() { let output = ck_with_subc(&server.connection_file_path, ["module", "list"]); assert_exit(&output, 0); - assert_eq!( - text(&output.stdout), - "module status health \ninsula running degraded\n" - ); + let expected_list = if cfg!(any(target_os = "linux", target_os = "macos")) { + "module status health reload \ninsula running degraded nothing\n" + } else { + "module status health reload \ninsula running degraded nothing (image n/a)\n" + }; + assert_eq!(text(&output.stdout), expected_list); assert!(!text(&output.stdout).contains("true")); assert!(!text(&output.stdout).contains("false")); module.stop().await.unwrap(); } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn module_list_keeps_configured_path_mismatch_pending_on_every_platform() { + let server = TestServer::start().await; + let supervisor = supervisor(&server); + let module_id = "path-mismatch"; + let module = spawn_stub(&server, &supervisor, module_id).await; + let replacement_home = TempDir::new("ck-list-replacement"); + let replacement = replacement_home.path().join( + Path::new(env!("CARGO_BIN_EXE_fake-aft-stub")) + .file_name() + .unwrap(), + ); + fs::copy(env!("CARGO_BIN_EXE_fake-aft-stub"), &replacement).unwrap(); + let mut configured = stub_spec(module_id); + configured.program = replacement.clone(); + module.update_spec_for_test(configured).await.unwrap(); + + let json = assert_json_success(ck_with_subc( + &server.connection_file_path, + ["module", "list", "--json"], + )); + let entry = &json["modules"][0]; + assert_eq!(entry["module_id"], module_id); + assert_eq!(entry["pending_reload"]["path"]["status"], "mismatch"); + assert_eq!( + entry["pending_reload"]["path"]["configured"], + json!(replacement) + ); + if !cfg!(any(target_os = "linux", target_os = "macos")) { + assert_eq!(entry["pending_reload"]["image"]["status"], "unavailable"); + assert_eq!( + entry["pending_reload"]["image"]["reason"], + "unsupported_platform" + ); + } + + let output = ck_with_subc(&server.connection_file_path, ["module", "list"]); + assert_exit(&output, 0); + let stdout = text(&output.stdout); + assert!( + stdout.lines().any(|line| { + line.starts_with(module_id) && line.trim_end().ends_with("pending (path)") + }), + "{stdout}" + ); + + module.stop().await.unwrap(); +} + /// What an operator is told about a module that speaks no subc wire. /// /// `live` stays a boolean on the wire, but for this module it answers a weaker @@ -2127,12 +2180,17 @@ async fn module_status_renders_key_value_block_byte_for_byte() { "start age renders as an age: {rendered_age:?} vs {started:?}" ); assert_eq!(before, format!("aft — running, degraded\n pid {pid}")); + let image_verdict = if cfg!(any(target_os = "linux", target_os = "macos")) { + "matches file at spawned path" + } else { + "not checked on this platform (unsupported_platform)" + }; // The budget renders with the window it is counted over (`in 10m`), because // the count alone reads as a lifetime total and stopped being one. assert_eq!( rest, format!( - "0 of 1 in 10m · drain 25 ms · restart backoff 10 ms to 30s\n last exit: none\n drain gauges: 0 drains with undeclared gauge\n binary: {binary} ({image})\nmetrics: run `ck health aft`\n" + "0 of 1 in 10m · drain 25 ms · restart backoff 10 ms to 30s\n last exit: none\n drain gauges: 0 drains with undeclared gauge\n binary: {binary} ({image})\n configured program: matches running process\n running image: {image_verdict}\nmetrics: run `ck health aft`\n" ) ); @@ -3205,6 +3263,7 @@ fn scripted_supervisor_entry(module_id: &str, drain_timeout_ms: Option) -> live: true, protocol: ModuleProtocol::Subc, health: SupervisorHealthStatus::Ok, + pending_reload: None, last_probe_ms: None, last_exit_code: None, last_exit_signal: None, diff --git a/crates/subc-core/tests/daemon_config.rs b/crates/subc-core/tests/daemon_config.rs index 494babe1..6f6b06fc 100644 --- a/crates/subc-core/tests/daemon_config.rs +++ b/crates/subc-core/tests/daemon_config.rs @@ -8,8 +8,8 @@ use std::{ use serde_json::{json, Value}; use subc_control::{ - ClientControlRequest, ClientControlResponse, ConsumerIdentity, SupervisorEntry, - SupervisorRescanResult, + ClientControlRequest, ClientControlResponse, ConsumerIdentity, ReloadPathAgreement, + SupervisorEntry, SupervisorRescanResult, }; use subc_daemon::{ bootstrap::{run_with_config, run_with_daemon_config_path, BootstrapConfig}, @@ -554,14 +554,109 @@ async fn rescan_changed_spec_is_pending_until_reload_uses_it() { let old_route = open_route(&mut old_client, module_id, 501).await; let before_pid = call_tool(&mut old_client, old_route, 502, "_test.pid").await; - let changed = stub_module( + let mut changed = stub_module( module_id, true, [("FAKE_AFT_TOOLCALL_RESULT", "after-reload")], ); + let replacement = daemon.temp_dir.join("replacement-fake-aft-stub"); + fs::copy(env!("CARGO_BIN_EXE_fake-aft-stub"), &replacement).unwrap(); + changed["program"] = json!(replacement.to_string_lossy()); fs::write(&daemon.config_path, config_doc([changed])).unwrap(); - let result = supervisor_rescan(&daemon.connection_file_path, 503).await; - assert_eq!(result.changed_pending_reload, [module_id]); + let preview = ck_under_test_command() + .args(["module", "rescan", "--dry-run", "--subc"]) + .arg(&daemon.connection_file_path) + .output() + .unwrap(); + assert!( + preview.status.success(), + "{}", + String::from_utf8_lossy(&preview.stderr) + ); + let preview_text = String::from_utf8(preview.stdout).unwrap(); + assert!( + preview_text.contains("would change (pending reload)"), + "{preview_text}" + ); + assert!(!preview_text.contains("would restart"), "{preview_text}"); + let applied = ck_under_test_command() + .args(["module", "rescan", "--subc"]) + .arg(&daemon.connection_file_path) + .output() + .unwrap(); + assert!( + applied.status.success(), + "{}", + String::from_utf8_lossy(&applied.stderr) + ); + let applied_text = String::from_utf8(applied.stdout).unwrap(); + assert!( + applied_text.contains("changed-pending-reload"), + "{applied_text}" + ); + assert!(applied_text.contains(module_id), "{applied_text}"); + assert!( + applied_text.contains("ck module restart "), + "{applied_text}" + ); + let pending = wait_for_supervisor_entry( + &daemon.connection_file_path, + module_id, + |entry| { + matches!( + entry.pending_reload.as_ref().map(|verdict| &verdict.path), + Some(ReloadPathAgreement::Mismatch { .. }) + ) + }, + STATE_TIMEOUT, + ) + .await; + assert!( + matches!(pending.pending_reload.unwrap().path, ReloadPathAgreement::Mismatch { configured, .. } if configured == replacement) + ); + let again = supervisor_rescan(&daemon.connection_file_path, 510).await; + assert!(again.changed_pending_reload.is_empty()); + let still_pending = supervisor_modules(&daemon.connection_file_path, 511).await; + assert!(matches!( + still_pending + .iter() + .find(|entry| entry.module_id == module_id) + .and_then(|entry| entry.pending_reload.as_ref()) + .map(|verdict| &verdict.path), + Some(ReloadPathAgreement::Mismatch { .. }) + )); + let list = ck_under_test_command() + .args(["module", "list", "--subc"]) + .arg(&daemon.connection_file_path) + .output() + .unwrap(); + assert!( + list.status.success(), + "{}", + String::from_utf8_lossy(&list.stderr) + ); + assert!(String::from_utf8(list.stdout) + .unwrap() + .contains("pending (path)")); + let status = ck_under_test_command() + .args(["module", "status", module_id, "--subc"]) + .arg(&daemon.connection_file_path) + .output() + .unwrap(); + assert!( + status.status.success(), + "{}", + String::from_utf8_lossy(&status.stderr) + ); + let status_text = String::from_utf8(status.stdout).unwrap(); + assert!( + status_text.contains("configured program: pending reload"), + "{status_text}" + ); + assert!( + status_text.contains(&replacement.display().to_string()), + "{status_text}" + ); assert_eq!( call_tool(&mut old_client, old_route, 504, "_test.pid").await, before_pid @@ -584,6 +679,18 @@ async fn rescan_changed_spec_is_pending_until_reload_uses_it() { response, ClientControlResponse::SupervisorAck { applied: true, .. } )); + wait_for_supervisor_entry( + &daemon.connection_file_path, + module_id, + |entry| { + matches!( + entry.pending_reload.as_ref().map(|verdict| &verdict.path), + Some(ReloadPathAgreement::Match) + ) + }, + STATE_TIMEOUT, + ) + .await; let mut new_client = wait_for_client(&daemon.connection_file_path, START_TIMEOUT).await; let new_route = open_route(&mut new_client, module_id, 507).await; diff --git a/crates/subc-daemon/Cargo.toml b/crates/subc-daemon/Cargo.toml index b36cbda0..33475b0d 100644 --- a/crates/subc-daemon/Cargo.toml +++ b/crates/subc-daemon/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "subc-daemon" -version = "0.21.3" +version = "0.21.4" edition = "2021" publish = true description = "Embeddable subc daemon: bootstrap, module supervision, and opaque-byte splice routing." @@ -22,7 +22,7 @@ rlimit = "0.11" serde = { version = "1", features = ["derive"] } serde_json = "1" sha2 = "0.10" -subc-control = { path = "../subc-control", version = "0.19" } +subc-control = { path = "../subc-control", version = "0.20" } subc-jsonc = { path = "../subc-jsonc", version = "0.1.0" } subc-protocol = { path = "../subc-protocol", version = "0.25.0" } subc-transport = { path = "../subc-transport", version = "0.7.0" } diff --git a/crates/subc-daemon/src/control.rs b/crates/subc-daemon/src/control.rs index 8e1e3abc..1394773e 100644 --- a/crates/subc-daemon/src/control.rs +++ b/crates/subc-daemon/src/control.rs @@ -1,7 +1,7 @@ use std::{ collections::{BTreeMap, BTreeSet, HashMap, HashSet}, fmt, - path::PathBuf, + path::{Path, PathBuf}, sync::{Arc, Mutex}, time::{Duration, Instant as StdInstant}, }; @@ -10,9 +10,10 @@ use serde::{Deserialize, Serialize}; use subc_control::{ ops, CapabilityRequirementStatus, CatalogEntry, ClientControlPush, ClientControlRequest, ClientControlResponse, ConsumerIdentity, DaemonBuildProvenance, DaemonObservedProcess, - ModuleDeclaredProvenance, ModuleProtocol, NotReadyReason, PollKind, RouteCloseReason, - SpawnCursor, StderrCaptureState, StderrTail, StderrTailEntry, SupervisorDaemonProvenance, - SupervisorEntry, SupervisorHealthEntry, SupervisorModuleProvenance, SupervisorObservedProcess, + ModuleDeclaredProvenance, ModuleProtocol, NotReadyReason, PendingReloadVerdict, PollKind, + ReloadPathAgreement, ReloadPathUnavailableReason, RouteCloseReason, SpawnCursor, + StderrCaptureState, StderrTail, StderrTailEntry, SupervisorDaemonProvenance, SupervisorEntry, + SupervisorHealthEntry, SupervisorModuleProvenance, SupervisorObservedProcess, SupervisorRescanResult, SupervisorRoute, SupervisorRouteConsumer, SupervisorRouteModule, }; use subc_protocol::{ @@ -151,6 +152,33 @@ pub const DEFAULT_ROUTE_BIND_BREAKER_COOLDOWN: Duration = Duration::from_secs(20 const DEFAULT_HEALTH_PROBE_TIMEOUT: Duration = Duration::from_secs(5); const SLOW_CONTROL_DISPATCH_THRESHOLD: Duration = Duration::from_secs(1); +fn reload_verdict( + configured: &Path, + spawned_from: Option<&Path>, + image: subc_control::RunningImageAgreement, +) -> PendingReloadVerdict { + let path = match spawned_from { + Some(spawned_from) if configured == spawned_from => ReloadPathAgreement::Match, + Some(spawned_from) => ReloadPathAgreement::Mismatch { + configured: configured.to_path_buf(), + spawned_from: spawned_from.to_path_buf(), + }, + None => ReloadPathAgreement::Unavailable { + reason: if matches!( + image, + subc_control::RunningImageAgreement::Unavailable { + reason: subc_control::RunningImageUnavailableReason::NotRunning + } + ) { + ReloadPathUnavailableReason::NotRunning + } else { + ReloadPathUnavailableReason::SpawnedPathUnavailable + }, + }, + }; + PendingReloadVerdict { path, image } +} + #[derive(Clone)] struct DaemonProvenanceFacts { build: DaemonBuildProvenance, @@ -2017,7 +2045,7 @@ impl ControlHandler { route_epoch, kind, } => self.handle_route_poll(ctx, frame, route_channel, route_epoch, kind), - ClientControlRequest::SupervisorList {} => self.handle_supervisor_list(frame), + ClientControlRequest::SupervisorList {} => self.handle_supervisor_list(frame).await, ClientControlRequest::SupervisorSpawnSnapshot {} => { self.handle_supervisor_spawn_snapshot(frame) } @@ -3198,46 +3226,57 @@ impl ControlHandler { } } - fn handle_supervisor_list(&self, frame: Frame) -> Result, RouterError> { + async fn handle_supervisor_list(&self, frame: Frame) -> Result, RouterError> { let generation = self .registry .generation() .map_err(|err| RouterError::backend(0, frame.header.corr, err.to_string()))?; - let modules = self - .supervisor - .list() - .into_iter() - .map(|module| { - let status = module.status_for_control("list").map_err(|err| { - RouterError::backend( - 0, - frame.header.corr, - format!("failed to read supervisor status: {err}"), - ) - })?; - Ok(SupervisorEntry { - module_id: status.module_id, - state: status.state.to_string(), - enabled: status.enabled, - live: status.live, - protocol: status.protocol, - health: status.health.status, - last_probe_ms: status.health.last_probe_ms, - last_exit_code: status.last_exit.as_ref().and_then(|e| e.code), - last_exit_signal: status.last_exit.as_ref().and_then(|e| e.signal), - last_exit_ms: status.last_exit.as_ref().map(|e| e.at_ms), - last_exit_kind: status.last_exit.as_ref().map(|e| e.kind.into()), - restart_count: Some(status.restart_count), - max_restarts: Some(status.max_restarts), - lifetime_restarts: Some(status.lifetime_restarts), - spawn_generation: Some(status.spawn_generation), - restart_window_secs: Some(status.restart_window.as_secs()), - drain_timeout_ms: Some(status.drain_timeout.as_millis() as u64), - restart_backoff_ms: Some(status.restart_backoff.as_millis() as u64), - restart_max_backoff_ms: Some(status.restart_max_backoff.as_millis() as u64), - }) - }) - .collect::, RouterError>>()?; + let mut modules = Vec::new(); + for module in self.supervisor.list() { + let status = module.status_for_control("list").map_err(|err| { + RouterError::backend( + 0, + frame.header.corr, + format!("failed to read supervisor status: {err}"), + ) + })?; + let (configured, _) = module.configuration().map_err(|err| { + RouterError::backend( + 0, + frame.header.corr, + format!("failed to read module configuration: {err}"), + ) + })?; + // Status and configuration snapshots release their locks before the image probe awaits. + let image = module.running_image_agreement().await; + let pending_reload = Some(reload_verdict( + &configured.program, + status.spawned_from.as_deref(), + image, + )); + modules.push(SupervisorEntry { + module_id: status.module_id, + state: status.state.to_string(), + enabled: status.enabled, + live: status.live, + protocol: status.protocol, + health: status.health.status, + pending_reload, + last_probe_ms: status.health.last_probe_ms, + last_exit_code: status.last_exit.as_ref().and_then(|e| e.code), + last_exit_signal: status.last_exit.as_ref().and_then(|e| e.signal), + last_exit_ms: status.last_exit.as_ref().map(|e| e.at_ms), + last_exit_kind: status.last_exit.as_ref().map(|e| e.kind.into()), + restart_count: Some(status.restart_count), + max_restarts: Some(status.max_restarts), + lifetime_restarts: Some(status.lifetime_restarts), + spawn_generation: Some(status.spawn_generation), + restart_window_secs: Some(status.restart_window.as_secs()), + drain_timeout_ms: Some(status.drain_timeout.as_millis() as u64), + restart_backoff_ms: Some(status.restart_backoff.as_millis() as u64), + restart_max_backoff_ms: Some(status.restart_max_backoff.as_millis() as u64), + }); + } let response = ClientControlResponse::SupervisorList { generation, modules, @@ -9336,6 +9375,94 @@ mod tests { assert_eq!(handler.provenance_probe_override, Some(expected)); } + #[test] + fn reload_verdict_detects_configured_program_different_from_spawned_path() { + let verdict = reload_verdict( + std::path::Path::new("/bin/new"), + Some(std::path::Path::new("/bin/old")), + subc_control::RunningImageAgreement::Unavailable { + reason: subc_control::RunningImageUnavailableReason::HashFailed, + }, + ); + assert!(matches!( + verdict.path, + subc_control::ReloadPathAgreement::Mismatch { configured, spawned_from } + if configured == std::path::Path::new("/bin/new") + && spawned_from == std::path::Path::new("/bin/old") + )); + } + + #[test] + fn reload_verdict_detects_replaced_image_at_same_path() { + let image = subc_control::RunningImageAgreement::Mismatch { + running: subc_control::RunningImageEvidence::LinuxProcSha256 { + digest: "old".into(), + }, + disk: subc_control::RunningImageEvidence::LinuxProcSha256 { + digest: "new".into(), + }, + }; + let verdict = reload_verdict( + std::path::Path::new("/bin/same"), + Some(std::path::Path::new("/bin/same")), + image.clone(), + ); + assert_eq!(verdict.path, subc_control::ReloadPathAgreement::Match); + assert_eq!(verdict.image, image); + } + + #[test] + fn reload_verdict_preserves_stopped_and_unavailable_reasons() { + let image = subc_control::RunningImageAgreement::Unavailable { + reason: subc_control::RunningImageUnavailableReason::NotRunning, + }; + let verdict = reload_verdict(std::path::Path::new("/bin/same"), None, image.clone()); + assert_eq!( + verdict.path, + subc_control::ReloadPathAgreement::Unavailable { + reason: subc_control::ReloadPathUnavailableReason::NotRunning, + } + ); + assert_eq!(verdict.image, image); + + let unconfirmed = subc_control::RunningImageAgreement::Unavailable { + reason: subc_control::RunningImageUnavailableReason::ProcessIdentityUnconfirmed, + }; + let verdict = reload_verdict( + std::path::Path::new("/bin/same"), + Some(std::path::Path::new("/bin/same")), + unconfirmed.clone(), + ); + assert_eq!(verdict.path, subc_control::ReloadPathAgreement::Match); + assert_eq!(verdict.image, unconfirmed); + } + + #[test] + fn reload_verdict_preserves_each_image_unavailability_reason() { + use subc_control::RunningImageUnavailableReason as Reason; + + for reason in [ + Reason::NotRunning, + Reason::UnsupportedPlatform, + Reason::RunningExecutableUnreadable, + Reason::SpawnedPathUnreadable, + Reason::HashFailed, + Reason::ProcessIdentityUnconfirmed, + Reason::Unknown("future_probe_reason".to_string()), + ] { + let image = subc_control::RunningImageAgreement::Unavailable { + reason: reason.clone(), + }; + let verdict = reload_verdict( + std::path::Path::new("/bin/same"), + Some(std::path::Path::new("/bin/same")), + image.clone(), + ); + assert_eq!(verdict.path, subc_control::ReloadPathAgreement::Match); + assert_eq!(verdict.image, image, "{reason:?}"); + } + } + #[tokio::test] async fn malformed_control_bodies_return_invalid_control_body() { let handler = ControlHandler::default();