diff --git a/crates/solstone-linux/src/cli.rs b/crates/solstone-linux/src/cli.rs index 3f9493a..3b86563 100644 --- a/crates/solstone-linux/src/cli.rs +++ b/crates/solstone-linux/src/cli.rs @@ -5,7 +5,10 @@ use crate::{ capture_stats::{ compute_quarantine_stats, compute_status_capture_stats, format_quarantine_line, }, - config::{Config, ConfigPaths, DEFAULT_SERVER_URL, load_config, save_config, save_identity}, + config::{ + Config, ConfigPaths, DEFAULT_SERVER_URL, load_config, save_config, + save_config_with_identity, + }, session_env::{self, Output, Runner}, streams::stream_name, sync_health::{derive_health, load_facts}, @@ -325,13 +328,9 @@ fn cmd_setup( } if let Some(token) = token { config.key = token; - let identity_paths = ConfigPaths { - base_dir: Some(config.base_dir.clone()), - config_dir: Some(config.config_dir.clone()), - }; - if let Err(error) = save_config(&config) - .and_then(|()| save_identity(&identity_paths, &config.key, &config.stream)) - { + // The one-line API swap is unavoidable after save_config became identity-preserving: + // setup still performs its original single whole-config write when given a token. + if let Err(error) = save_config_with_identity(&config) { let _ = write_line(errors, format!("Error saving config: {error}")); return 1; } @@ -630,6 +629,7 @@ fn cmd_run(interval: Option) -> i32 { #[cfg(test)] mod tests { use super::*; + use crate::config::save_identity; use clap::CommandFactory; use std::{cell::Cell, path::Path}; diff --git a/crates/solstone-linux/src/config.rs b/crates/solstone-linux/src/config.rs index c468a4b..399e178 100644 --- a/crates/solstone-linux/src/config.rs +++ b/crates/solstone-linux/src/config.rs @@ -19,7 +19,7 @@ pub const DEFAULT_SERVER_URL: &str = "http://localhost:5015"; pub const DEFAULT_SYNC_STALE_THRESHOLD: i64 = 600; const DEFAULT_RETRY_DELAYS: [i64; 4] = [5, 30, 120, 300]; const CONFIG_WRITE_LOCK_TIMEOUT: Duration = Duration::from_millis(100); -const CONFIG_WRITE_LOCK_POLL: Duration = Duration::from_millis(1); +const CONFIG_WRITE_LOCK_POLL: Duration = Duration::from_micros(100); static CONFIG_WRITE_LOCK: OnceLock> = OnceLock::new(); static CONFIG_TEMP_SEQUENCE: AtomicU64 = AtomicU64::new(0); @@ -278,6 +278,9 @@ pub fn load_config(paths: ConfigPaths) -> LoadedConfig { } fn acquire_config_write_lock() -> io::Result> { + // The guard protects filesystem syscalls only (read, write, chmod, and rename), never an + // await or interactive prompt. Normal contention is therefore sub-millisecond; the bounded + // deadline is a backstop for a stalled writer rather than an expected wait. let lock = CONFIG_WRITE_LOCK.get_or_init(|| Mutex::new(())); let deadline = Instant::now() + CONFIG_WRITE_LOCK_TIMEOUT; loop { @@ -311,31 +314,64 @@ fn write_config(config: &Config) -> io::Result<()> { fs::rename(temporary, path) } -pub fn save_config(config: &Config) -> io::Result<()> { - // This bounded critical section contains filesystem syscalls only. It never spans a prompt, - // an await, or any caller-visible guard, so an abandoned interactive settings session cannot - // prevent the running app from repairing its identity. +enum IdentityWrite<'a> { + PreserveDisk, + Provided { key: &'a str, stream: &'a str }, +} + +fn save_config_inner( + paths: &ConfigPaths, + source: Option<&Config>, + identity: IdentityWrite<'_>, +) -> io::Result<()> { let _guard = acquire_config_write_lock()?; - let mut merged = config.clone(); - if config.config_path().exists() { - let disk = load_config(ConfigPaths { - base_dir: Some(config.base_dir.clone()), - config_dir: Some(config.config_dir.clone()), - }) - .config; - merged.key = disk.key; - merged.stream = disk.stream; + let mut merged = source + .cloned() + .unwrap_or_else(|| load_config(paths.clone()).config); + match identity { + IdentityWrite::PreserveDisk if merged.config_path().exists() => { + let disk = load_config(paths.clone()).config; + merged.key = disk.key; + merged.stream = disk.stream; + } + // On the first write there is no disk identity to preserve, so save_config necessarily + // writes the caller's identity. Existing configs remain safe-by-default. + IdentityWrite::PreserveDisk => {} + IdentityWrite::Provided { key, stream } => { + merged.key = key.to_owned(); + merged.stream = stream.to_owned(); + } } write_config(&merged) } +pub fn save_config(config: &Config) -> io::Result<()> { + save_config_inner( + &ConfigPaths { + base_dir: Some(config.base_dir.clone()), + config_dir: Some(config.config_dir.clone()), + }, + Some(config), + IdentityWrite::PreserveDisk, + ) +} + +pub fn save_config_with_identity(config: &Config) -> io::Result<()> { + save_config_inner( + &ConfigPaths { + base_dir: Some(config.base_dir.clone()), + config_dir: Some(config.config_dir.clone()), + }, + Some(config), + IdentityWrite::Provided { + key: &config.key, + stream: &config.stream, + }, + ) +} + pub fn save_identity(paths: &ConfigPaths, key: &str, stream: &str) -> io::Result<()> { - // See save_config: lock acquisition is bounded and the guard covers filesystem work only. - let _guard = acquire_config_write_lock()?; - let mut config = load_config(paths.clone()).config; - config.key = key.to_owned(); - config.stream = stream.to_owned(); - write_config(&config) + save_config_inner(paths, None, IdentityWrite::Provided { key, stream }) } fn migrate(config: &Config) -> io::Result<()> { diff --git a/crates/solstone-linux/src/sync.rs b/crates/solstone-linux/src/sync.rs index 778431e..47b6d83 100644 --- a/crates/solstone-linux/src/sync.rs +++ b/crates/solstone-linux/src/sync.rs @@ -2041,9 +2041,10 @@ mod tests { ); } - // AC 1: listing_401_recovers_and_uploads_segment proves a rejected listing is recoverable. + // Legacy breaker criterion: a listing 401 opens a recoverable breaker and a later probe + // resumes segment upload; this does not exercise stale-key repair. #[tokio::test] - async fn listing_401_recovers_and_uploads_segment() { + async fn listing_401_opens_breaker_and_later_probe_resumes_upload() { let temp = tempfile::tempdir().unwrap(); create_segment(&temp, "120000_300", b"screen"); let (server, mut worker) = test_worker( @@ -2164,7 +2165,7 @@ mod tests { assert!(!captured.contains(NEW), "{captured}"); } - // AC 15: one or twenty sync-path rejections emit exactly one first-rejection warning. + // AC 15: twenty sync-path rejections emit exactly one first-rejection warning. #[tokio::test] async fn first_rejection_warning_is_once_per_window() { let mut responses = vec![(401, json!({})), (200, json!({"key":"K","name":"desktop"}))]; @@ -2200,7 +2201,44 @@ mod tests { ); } - // Named deviation: AC 2+5 intentionally break + // AC 15: one sync-path rejection emits exactly one first-rejection warning. + #[tokio::test] + async fn one_rejection_emits_one_first_rejection_warning() { + let temp = tempfile::tempdir().unwrap(); + let (server, worker) = test_worker( + &temp, + vec![(401, json!({})), (200, json!({"key":"K","name":"desktop"}))], + -1, + ) + .await; + let client = Arc::clone(&worker.client); + let output = Arc::new(Mutex::new(Vec::new())); + let writer = Buffer(Arc::clone(&output)); + let subscriber = tracing_subscriber::fmt() + .without_time() + .with_ansi(false) + .with_writer(move || writer.clone()) + .finish(); + tokio::spawn( + async move { + let _ = client.get_server_segments("20260101").await; + } + .with_subscriber(subscriber), + ) + .await + .unwrap(); + assert_eq!(server.request_count("/app/observer/register"), 1); + let captured = String::from_utf8(output.lock().unwrap().clone()).unwrap(); + assert_eq!( + captured + .matches("Journal rejected the current key; attempting identity repair") + .count(), + 1, + "{captured}" + ); + } + + // Named deviation from the legacy breaker AC 2+5: // tests/test_sync.py::test_auth_opens_immediately parity: both statuses open immediately, // but only a revoked 403 is permanent. #[tokio::test] @@ -2220,9 +2258,10 @@ mod tests { } } - // AC 3: upload_401_is_recoverable_and_retries proves the rejected POST is retried. + // Legacy breaker criterion: an upload 401 opens a recoverable breaker and a later probe + // permits the segment POST to be retried; this does not exercise successful stale-key repair. #[tokio::test] - async fn upload_401_is_recoverable_and_retries() { + async fn upload_401_opens_recoverable_breaker_and_later_retries() { let temp = tempfile::tempdir().unwrap(); create_segment(&temp, "120000_300", b"screen"); let (server, mut worker) = test_worker( diff --git a/crates/solstone-linux/src/upload.rs b/crates/solstone-linux/src/upload.rs index df4ace2..023c8b4 100644 --- a/crates/solstone-linux/src/upload.rs +++ b/crates/solstone-linux/src/upload.rs @@ -110,7 +110,6 @@ struct TellingState { window_start: Option, records: usize, rejection_warned: bool, - suppressed_records: usize, } enum RegisterAttempt { @@ -236,10 +235,14 @@ impl UploadClient { return true; } RegisterAttempt::GuardRefused { reason_code } => { + // One-shot CLI setup deliberately does not consume the daemon telling window. tracing::error!( status = 403, reason_code = reason_code.as_deref(), - "Journal refused local registration" + key_prefix = key_prefix(&self.inner.key.lock().unwrap()), + recovery_generation = + self.inner.recovery_generation.load(Ordering::Acquire), + "Journal refused local identity repair" ); return false; } @@ -533,7 +536,6 @@ impl Inner { } telling.rejection_warned = true; if telling.records >= TELLING_BURST_LIMIT { - telling.suppressed_records += 1; return; } telling.records += 1; @@ -549,14 +551,13 @@ impl Inner { } fn tell_outcome(&self, outcome: &str, name: &str, old_key: &str, new_key: &str, now: f64) { + let dropped_events_lower_bound = self.dropped_events.swap(0, Ordering::AcqRel); let mut telling = self.prepare_telling(now); if telling.records >= TELLING_BURST_LIMIT { - telling.suppressed_records += 1; return; } telling.records += 1; drop(telling); - let dropped_events_lower_bound = self.dropped_events.swap(0, Ordering::AcqRel); tracing::warn!( outcome, name, @@ -573,7 +574,6 @@ impl Inner { fn tell_guard_refusal(&self, reason_code: Option<&str>, key: &str, now: f64) { let mut telling = self.prepare_telling(now); if telling.records >= TELLING_BURST_LIMIT { - telling.suppressed_records += 1; return; } telling.records += 1; @@ -1045,9 +1045,12 @@ mod tests { } } - // AC 8: cancellation around a successful register exposes only complete stale/new identities. + // AC 8: persistence after a successful register is synchronous and await-free, so + // cancellation cannot be scheduled between the response and persistence. Cancellation before + // the response leaves the exact stale pair; uninterrupted completion writes the exact new + // pair. Those are therefore the only two identities reachable on disk. #[tokio::test] - async fn cancelled_recovery_exposes_no_third_disk_identity() { + async fn cancellation_boundary_exposes_only_stale_or_new_disk_identity() { const STALE: &str = "STALE-KEY"; const NEW: &str = "NEW-KEY"; for cancel in [true, false] { @@ -1122,7 +1125,7 @@ mod tests { assert!(!client.is_registered()); } - // AC 15: zero rejections and twenty clean relays emit no rejection warning. + // AC 15: twenty clean relays emit no rejection warning. #[tokio::test] async fn clean_relays_emit_no_rejection_warning() { let server = MockServer::new((0..20).map(|_| (200, json!({}))).collect()).await; @@ -1150,6 +1153,28 @@ mod tests { ); } + // AC 15: a directly-awaited path with no relay and no rejection is silent. + #[tokio::test] + async fn zero_rejections_emit_no_rejection_warning() { + let server = MockServer::new(vec![]).await; + let temp = TempDir::new().unwrap(); + let config = config(&server, &temp); + let _client = client(&config); + let output = Arc::new(Mutex::new(Vec::new())); + let writer = Buffer(Arc::clone(&output)); + let subscriber = tracing_subscriber::fmt() + .without_time() + .with_ansi(false) + .with_writer(move || writer.clone()) + .finish(); + async {}.with_subscriber(subscriber).await; + let captured = String::from_utf8(output.lock().unwrap().clone()).unwrap(); + assert!( + !captured.contains("Journal rejected the current key"), + "{captured}" + ); + } + // tests/test_upload.py::test_upload_segment_uses_bearer_and_keyless_route // tests/test_upload.py::test_upload_segment_declares_content_types // AC: multipart received bytes preserve fields, filenames, content types, and file content @@ -1760,6 +1785,15 @@ mod tests { stream: "desktop".into(), ..config(&server, &temp) }; + save_identity( + &ConfigPaths { + base_dir: Some(config.base_dir.clone()), + config_dir: Some(config.config_dir.clone()), + }, + &config.key, + &config.stream, + ) + .unwrap(); let client = client(&config); let output = Arc::new(Mutex::new(Vec::new())); let writer = Buffer(Arc::clone(&output)); @@ -1772,6 +1806,9 @@ mod tests { assert!(!client.relay_event("observe", "status", Map::new()).await); assert!(!client.is_revoked()); assert_eq!(client.inner.key.lock().unwrap().as_str(), "STALE-KEY"); + let saved = load_config(client.inner.paths.clone()).config; + assert_eq!(saved.key, "STALE-KEY"); + assert_eq!(saved.stream, "desktop"); assert!(!client.relay_event("observe", "status", Map::new()).await); assert!(client.is_revoked()); } @@ -1888,7 +1925,9 @@ mod tests { assert_eq!(body["node_id"], 1); } - // AC 13: event recovery performs one register attempt while the bounded queue remains usable. + // AC 13: event recovery performs exactly one register request while the bounded queue remains + // usable. There is no retry loop on this path, so a stalled worker is bounded by the single + // register request's EVENT_TIMEOUT without making this a 30-second real-time test. #[tokio::test] async fn event_recovery_is_single_attempt_and_queue_stays_bounded() { let gate = Arc::new(tokio::sync::Notify::new()); @@ -1906,6 +1945,7 @@ mod tests { let mut client = client(&config); client.enqueue_status(Map::new()); wait_for_requests(&server, 2).await; + assert_eq!(server.request_count("/app/observer/register"), 1); for sequence in 0..20 { assert!( client