From 5305a3b1b8af962dcbf2bddcf673bdfb426b7037 Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Thu, 30 Jul 2026 10:32:20 -0600 Subject: [PATCH] recover a stale observer key by re-registering, and say so Add a serialized, cooldown-bounded single-attempt re-register reachable from listing, segment upload, and event relay 401s over the shared Arc. Persist replacement identity before swapping Inner, and publish a recovery generation so try_probe skips the open breaker's remaining cooldown exactly once. Fix AC 17 rather than recording a residual: save_config now preserves key and stream from disk by default, while narrow save_identity is the only identity writer. A stale whole-struct write from solstone settings therefore cannot clobber a recovered key, and future non-identity writers are safe by omission. Deliberately change registration_403_latches_revoked to registration_403_does_not_latch_revoked because local_request_only is a register-route guard, not evidence of revocation; ingest, listing, and event 403s still latch. Preserve ensure_registered_skips_when_key_present with only its constructor adaptation. Keep revoked fully short-circuiting so a 403-latched observer attempts no recovery on any path. Co-Authored-By: Claude Opus 5 (1M context) --- crates/solstone-linux/src/cli.rs | 34 +- crates/solstone-linux/src/config.rs | 123 ++- .../src/observer_contract_tests.rs | 8 +- crates/solstone-linux/src/run.rs | 21 +- crates/solstone-linux/src/sync.rs | 345 ++++++-- crates/solstone-linux/src/test_support.rs | 48 +- crates/solstone-linux/src/upload.rs | 765 ++++++++++++++++-- 7 files changed, 1194 insertions(+), 150 deletions(-) diff --git a/crates/solstone-linux/src/cli.rs b/crates/solstone-linux/src/cli.rs index df10d22..3f9493a 100644 --- a/crates/solstone-linux/src/cli.rs +++ b/crates/solstone-linux/src/cli.rs @@ -5,11 +5,11 @@ use crate::{ capture_stats::{ compute_quarantine_stats, compute_status_capture_stats, format_quarantine_line, }, - config::{Config, ConfigPaths, DEFAULT_SERVER_URL, load_config, save_config}, + config::{Config, ConfigPaths, DEFAULT_SERVER_URL, load_config, save_config, save_identity}, session_env::{self, Output, Runner}, streams::stream_name, sync_health::{derive_health, load_facts}, - upload::UploadClient, + upload::{UploadClient, key_prefix}, }; use clap::{Parser, Subcommand}; use std::{ @@ -258,7 +258,13 @@ impl Registrar for RealRegistrar { .build() .map_err(|error| error.to_string())?; let _guard = runtime.enter(); - let client = UploadClient::new(config, host, "linux", env!("CARGO_PKG_VERSION")); + let client = UploadClient::new( + config, + host, + "linux", + env!("CARGO_PKG_VERSION"), + std::sync::Arc::new(crate::run::SystemClock::new()), + ); Ok(runtime.block_on(client.ensure_registered(config))) } } @@ -319,7 +325,13 @@ fn cmd_setup( } if let Some(token) = token { config.key = token; - if let Err(error) = save_config(&config) { + 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)) + { let _ = write_line(errors, format!("Error saving config: {error}")); return 1; } @@ -385,10 +397,6 @@ fn setup_footer(output: &mut dyn Write, config: &Config) -> io::Result<()> { ) } -fn key_prefix(key: &str) -> String { - key.chars().take(8).collect() -} - trait PromptIo { fn read_line(&mut self, prompt: &str) -> io::Result; fn write_line(&mut self, line: &str) -> io::Result<()>; @@ -792,7 +800,15 @@ mod tests { if self.result { config.key = "newkey00".into(); config.stream = "locked-stream".into(); - save_config(config).unwrap(); + save_identity( + &ConfigPaths { + base_dir: Some(config.base_dir.clone()), + config_dir: Some(config.config_dir.clone()), + }, + &config.key, + &config.stream, + ) + .unwrap(); } Ok(self.result) } diff --git a/crates/solstone-linux/src/config.rs b/crates/solstone-linux/src/config.rs index 320853b..c468a4b 100644 --- a/crates/solstone-linux/src/config.rs +++ b/crates/solstone-linux/src/config.rs @@ -3,11 +3,25 @@ use serde::{Deserialize, Serialize}; use serde_json::{Map, Value}; -use std::{env, fs, io, os::unix::fs::PermissionsExt, path::PathBuf}; +use std::{ + env, fs, io, + os::unix::fs::PermissionsExt, + path::PathBuf, + sync::{ + Mutex, MutexGuard, OnceLock, + atomic::{AtomicU64, Ordering}, + }, + thread, + time::{Duration, Instant}, +}; 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); +static CONFIG_WRITE_LOCK: OnceLock> = OnceLock::new(); +static CONFIG_TEMP_SEQUENCE: AtomicU64 = AtomicU64::new(0); #[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] pub struct Config { @@ -98,7 +112,7 @@ pub struct LoadedConfig { pub warnings: Vec, } -#[derive(Default)] +#[derive(Clone, Default)] pub struct ConfigPaths { pub base_dir: Option, pub config_dir: Option, @@ -263,10 +277,33 @@ pub fn load_config(paths: ConfigPaths) -> LoadedConfig { LoadedConfig { config, warnings } } -pub fn save_config(config: &Config) -> io::Result<()> { +fn acquire_config_write_lock() -> io::Result> { + let lock = CONFIG_WRITE_LOCK.get_or_init(|| Mutex::new(())); + let deadline = Instant::now() + CONFIG_WRITE_LOCK_TIMEOUT; + loop { + match lock.try_lock() { + Ok(guard) => return Ok(guard), + Err(std::sync::TryLockError::WouldBlock) if Instant::now() < deadline => { + thread::sleep(CONFIG_WRITE_LOCK_POLL); + } + Err(std::sync::TryLockError::WouldBlock) => { + return Err(io::Error::new( + io::ErrorKind::TimedOut, + "timed out waiting to write config", + )); + } + Err(std::sync::TryLockError::Poisoned(_)) => { + return Err(io::Error::other("config write lock is poisoned")); + } + } + } +} + +fn write_config(config: &Config) -> io::Result<()> { config.ensure_dirs()?; let path = config.config_path(); - let temporary = path.with_extension(format!("{}.tmp", std::process::id())); + let sequence = CONFIG_TEMP_SEQUENCE.fetch_add(1, Ordering::Relaxed); + let temporary = path.with_extension(format!("{}.{}.tmp", std::process::id(), sequence)); let mut text = serde_json::to_string_pretty(config).map_err(io::Error::other)?; text.push('\n'); fs::write(&temporary, text)?; @@ -274,6 +311,33 @@ pub fn save_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. + 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; + } + write_config(&merged) +} + +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) +} + fn migrate(config: &Config) -> io::Result<()> { let old_dir = config.base_dir.join("config"); if config.config_dir == old_dir || config.config_path().exists() { @@ -702,6 +766,57 @@ mod tests { .exists() ); } + + // AC 11/17: identity-only recovery and a stale whole-config writer preserve both change sets. + #[test] + fn stale_settings_snapshot_preserves_recovered_identity() { + let t = tempfile::tempdir().unwrap(); + let initial = Config { + base_dir: t.path().into(), + config_dir: t.path().join("cfg"), + server_url: "https://journal".into(), + cache_retention_days: 7, + ..Config::default() + }; + save_config(&initial).unwrap(); + let config_paths = paths(t.path()); + save_identity(&config_paths, "STALE-KEY", "desktop-old").unwrap(); + let mut stale_settings = load_config(config_paths.clone()).config; + + save_identity(&config_paths, "NEW-KEY", "desktop-new").unwrap(); + stale_settings.cache_retention_days = 30; + save_config(&stale_settings).unwrap(); + + let saved = load_config(config_paths).config; + assert_eq!(saved.key, "NEW-KEY"); + assert_eq!(saved.stream, "desktop-new"); + assert_eq!(saved.cache_retention_days, 30); + } + + // AC 12: recovery writes identity only and cannot revert newer non-identity disk state. + #[test] + fn save_identity_preserves_newer_server_url() { + let t = tempfile::tempdir().unwrap(); + let config_paths = paths(t.path()); + let mut config = Config { + base_dir: t.path().into(), + config_dir: t.path().join("cfg"), + server_url: "https://old".into(), + ..Config::default() + }; + save_config(&config).unwrap(); + save_identity(&config_paths, "STALE-KEY", "desktop").unwrap(); + config = load_config(config_paths.clone()).config; + config.server_url = "https://new".into(); + save_config(&config).unwrap(); + + save_identity(&config_paths, "NEW-KEY", "desktop-new").unwrap(); + let saved = load_config(config_paths).config; + assert_eq!(saved.server_url, "https://new"); + assert_eq!(saved.key, "NEW-KEY"); + assert_eq!(saved.stream, "desktop-new"); + } + // AC: migration failure is returned as one warning and leaves legacy data intact. #[test] fn migration_warning() { diff --git a/crates/solstone-linux/src/observer_contract_tests.rs b/crates/solstone-linux/src/observer_contract_tests.rs index 985628c..980f3d5 100644 --- a/crates/solstone-linux/src/observer_contract_tests.rs +++ b/crates/solstone-linux/src/observer_contract_tests.rs @@ -428,7 +428,13 @@ fn config(server: &MockServer, temp: &TempDir) -> Config { } fn client(config: &Config) -> UploadClient { - UploadClient::new(config, "archon", "linux", "1.4.0") + UploadClient::new( + config, + "archon", + "linux", + "1.4.0", + std::sync::Arc::new(crate::test_support::MutableClock::new(0.0, 0.0)), + ) } async fn assert_upload_contract( diff --git a/crates/solstone-linux/src/run.rs b/crates/solstone-linux/src/run.rs index df1ff99..8840961 100644 --- a/crates/solstone-linux/src/run.rs +++ b/crates/solstone-linux/src/run.rs @@ -250,13 +250,14 @@ fn run_capture( let video = construct_video(Arc::clone(¬ifier), || VideoBackend::new(&config)) .map_err(ObserverError::VideoStart)?; let (audio, mute) = PulseAudioCapture::spawn().map_err(ObserverError::Io)?; + let clock = SystemClock::new(); let upload = Arc::new(UploadClient::new( &config, host.clone(), "linux", env!("CARGO_PKG_VERSION"), + Arc::new(clock.clone()), )); - let clock = SystemClock::new(); let sync = SyncService::start(config.clone(), Arc::clone(&upload), Arc::new(clock.clone())); let sync_trigger = sync.trigger_handle(); let sync_sampler = sync.sampler_handle(); @@ -543,7 +544,7 @@ pub(crate) struct SystemClock { started: Instant, } impl SystemClock { - fn new() -> Self { + pub(crate) fn new() -> Self { Self { wall: SystemTime::now() .duration_since(UNIX_EPOCH) @@ -813,7 +814,13 @@ mod tests { config_dir: t.path().join("config"), ..Config::default() }; - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(SystemClock::new()), + )); let count = Arc::new(AtomicUsize::new(0)); let mut sink = UploadEventSink { client: Arc::clone(&client), @@ -893,7 +900,13 @@ mod tests { config_dir: t.path().join("config"), ..Config::default() }; - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(SystemClock::new()), + )); let extra_owner = Arc::clone(&client); let result = stop_upload_sender(client, Duration::from_millis(1)).await; assert!(result.is_err()); diff --git a/crates/solstone-linux/src/sync.rs b/crates/solstone-linux/src/sync.rs index 407c64b..778431e 100644 --- a/crates/solstone-linux/src/sync.rs +++ b/crates/solstone-linux/src/sync.rs @@ -240,6 +240,7 @@ struct SyncWorker { last_contact_flush: f64, registration_refused: bool, draining_shutdown: bool, + last_recovery_generation: u64, #[cfg(test)] fail_next_pass: bool, } @@ -254,6 +255,7 @@ impl SyncWorker { recent_error_count: Arc, ) -> Self { let synced_days = load_synced_days(&config.state_dir()); + let last_recovery_generation = client.recovery_generation(); Self { config, client, @@ -275,6 +277,7 @@ impl SyncWorker { last_contact_flush: 0.0, registration_refused: false, draining_shutdown: false, + last_recovery_generation, #[cfg(test)] fail_next_pass: false, } @@ -340,8 +343,13 @@ impl SyncWorker { self.clear_progress(); return false; } + let generation = self.client.recovery_generation(); + let skip_cooldown = generation != self.last_recovery_generation; + if skip_cooldown { + self.last_recovery_generation = generation; + } let elapsed = self.clock.monotonic_seconds() - self.circuit_open_since; - if elapsed < self.circuit_cooldown { + if !skip_cooldown && elapsed < self.circuit_cooldown { self.set_progress( format!("{:.0}s until probe", self.circuit_cooldown - elapsed), false, @@ -352,6 +360,8 @@ impl SyncWorker { // Named deviation: Python's datetime.now() is uninjected; day derivation uses the injected wall clock. let today = timestamp_parts(self.clock.wall_seconds()).0; let result = self.client.get_server_segments(&today).await; + // Consume a recovery that completed inside this probe so it cannot skip a future breaker. + self.last_recovery_generation = self.client.recovery_generation(); if result.error_type.is_none() { self.record_contact(true); self.circuit_open = false; @@ -992,7 +1002,7 @@ fn remove_if_empty(path: &Path) { mod tests { use super::*; use crate::{ - test_support::{MockServer, wait_for_requests}, + test_support::{MockServer, MutableClock, wait_for_requests}, upload::ListingFile, }; use serde_json::{Value, json}; @@ -1012,35 +1022,6 @@ mod tests { } } - struct MutableClock { - wall: std::sync::atomic::AtomicU64, - mono: std::sync::atomic::AtomicU64, - } - - impl MutableClock { - fn new(wall: f64, mono: f64) -> Self { - Self { - wall: std::sync::atomic::AtomicU64::new(wall.to_bits()), - mono: std::sync::atomic::AtomicU64::new(mono.to_bits()), - } - } - fn set_wall(&self, value: f64) { - self.wall.store(value.to_bits(), Ordering::Release); - } - fn set_mono(&self, value: f64) { - self.mono.store(value.to_bits(), Ordering::Release); - } - } - - impl Clock for MutableClock { - fn wall_seconds(&self) -> f64 { - f64::from_bits(self.wall.load(Ordering::Acquire)) - } - fn monotonic_seconds(&self) -> f64 { - f64::from_bits(self.mono.load(Ordering::Acquire)) - } - } - struct FixedClock { wall: f64, mono: f64, @@ -1126,14 +1107,21 @@ mod tests { config_dir: temp.path().join("config"), ..Config::default() }; - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let clock = Arc::new(FixedClock { + wall: 1_800_000_000.0, + mono: 100.0, + }); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + clock.clone(), + )); let worker = SyncWorker::new( config, client, - Arc::new(FixedClock { - wall: 1_800_000_000.0, - mono: 100.0, - }), + clock, SyncControl { notify: Arc::new(Notify::new()), pending_trigger: Arc::new(AtomicBool::new(false)), @@ -1777,7 +1765,16 @@ mod tests { ..Config::default() }; let server = MockServer::new(vec![]).await; - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let clock: Arc = Arc::new(FixedClock { wall: 1_800_000_000.0, mono: 0.0, @@ -1845,7 +1842,16 @@ mod tests { ..Config::default() }; save_synced_days(&config.state_dir(), &HashSet::from(["20260101".to_owned()])).unwrap(); - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let service = SyncService::start( config, client, @@ -2063,6 +2069,137 @@ mod tests { assert!(upload_hits(&server) >= 1); } + // AC 4/16: listing alone repairs a distinct key and skips the still-active breaker wait once. + #[tokio::test] + async fn recovery_generation_skips_breaker_wait_once_with_empty_event_queue() { + let temp = tempfile::tempdir().unwrap(); + let (server, mut worker) = test_worker( + &temp, + vec![ + (401, json!({})), + (200, json!({"key":"NEW-KEY","name":"desktop-new"})), + (200, json!({"items":[],"total":0})), + ], + -1, + ) + .await; + let clock = Arc::new(MutableClock::new(1_800_000_000.0, 100.0)); + worker.clock = clock.clone(); + + worker.sync_pass(true).await; + assert!(worker.circuit_open); + assert_eq!(worker.client.recovery_generation(), 1); + assert!(clock.monotonic_seconds() < worker.circuit_open_since + worker.circuit_cooldown); + let before = server.request_count("/app/observer/ingest/segments/"); + assert!(worker.try_probe().await); + assert_eq!( + server.request_count("/app/observer/ingest/segments/"), + before + 1 + ); + assert_eq!(clock.monotonic_seconds(), 100.0); + } + + // AC 14: positive-control the sync-path record, then prove it contains prefixes but no full key. + #[tokio::test] + async fn sync_recovery_log_has_name_and_prefixes_without_full_keys() { + const STALE: &str = "STALE-KEY-FULL"; + const NEW: &str = "NEW-KEY-FULL"; + let temp = tempfile::tempdir().unwrap(); + let server = MockServer::new(vec![ + (401, json!({})), + (200, json!({"key":NEW,"name":"desktop-new"})), + ]) + .await; + let config = Config { + server_url: server.url.clone(), + key: STALE.into(), + stream: "desktop".into(), + sync_retry_delays: vec![0], + base_dir: temp.path().to_path_buf(), + config_dir: temp.path().join("config"), + ..Config::default() + }; + let clock = Arc::new(FixedClock { + wall: 1_800_000_000.0, + mono: 100.0, + }); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + clock.clone(), + )); + let mut worker = SyncWorker::new( + config, + client, + clock, + SyncControl { + notify: Arc::new(Notify::new()), + pending_trigger: Arc::new(AtomicBool::new(false)), + running: Arc::new(AtomicBool::new(true)), + }, + Arc::new(Mutex::new(SyncFacts::default())), + Arc::new(AtomicU8::new(0)), + ); + 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 { worker.sync_pass(true).await }.with_subscriber(subscriber)) + .await + .unwrap(); + let captured = String::from_utf8(output.lock().unwrap().clone()).unwrap(); + assert!( + captured.contains("Journal identity repair completed"), + "{captured}" + ); + assert!(captured.contains("desktop-new"), "{captured}"); + assert!(captured.contains("STALE-KE"), "{captured}"); + assert!(captured.contains("NEW-KEY-"), "{captured}"); + assert!(!captured.contains(STALE), "{captured}"); + assert!(!captured.contains(NEW), "{captured}"); + } + + // AC 15: one or 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"}))]; + responses.extend((0..19).map(|_| (401, json!({})))); + let temp = tempfile::tempdir().unwrap(); + let (server, worker) = test_worker(&temp, responses, -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 { + for _ in 0..20 { + 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: AC 2+5 intentionally break // tests/test_sync.py::test_auth_opens_immediately parity: both statuses open immediately, // but only a revoked 403 is permanent. @@ -2394,7 +2531,16 @@ mod tests { config_dir: temp.path().join("config"), ..Config::default() }; - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let service = SyncService::start( config, Arc::clone(&client), @@ -2425,7 +2571,16 @@ mod tests { }, ) .unwrap(); - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let service = SyncService::start( config.clone(), client, @@ -2874,7 +3029,16 @@ mod tests { config_dir: temp.path().join("config"), ..Config::default() }; - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let service = SyncService::start( config, client, @@ -2924,7 +3088,16 @@ mod tests { }; save_facts(&config.state_dir(), &facts).unwrap(); assert_eq!(load_facts(&config.state_dir()).pending_confirmed, Some(-5)); - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let service = SyncService::start( config, client, @@ -2949,7 +3122,16 @@ mod tests { config_dir: temp.path().join("config"), ..Config::default() }; - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let service = SyncService::start( config, client, @@ -2976,7 +3158,16 @@ mod tests { config_dir: temp.path().join("config"), ..Config::default() }; - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let service = SyncService::start( config, client, @@ -3013,7 +3204,16 @@ mod tests { ..Config::default() }; save_synced_days(&config.state_dir(), &HashSet::from(["20260101".to_owned()])).unwrap(); - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let clock = Arc::new(MutableClock::new(1_800_000_000.0, 0.0)); let service = SyncService::start(config, client, clock.clone()); service.trigger(); @@ -3064,7 +3264,16 @@ mod tests { config_dir: temp.path().join("config"), ..Config::default() }; - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let service = SyncService::start( config, client, @@ -3103,7 +3312,16 @@ mod tests { config_dir: temp.path().join("config"), ..Config::default() }; - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let service = SyncService::start( config.clone(), client, @@ -3189,7 +3407,16 @@ mod tests { ..Config::default() }; save_synced_days(&config.state_dir(), &HashSet::from(["20260101".to_owned()])).unwrap(); - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let service = SyncService::start( config, client, @@ -3224,7 +3451,16 @@ mod tests { let (server, mut worker) = test_worker(&temp, vec![], -1).await; worker.config.key.clear(); let config = worker.config.clone(); - worker.client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + worker.client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let notify = Arc::clone(&worker.notify); let running = Arc::clone(&worker.running); let facts = worker.facts.lock().unwrap().clone(); @@ -3280,7 +3516,16 @@ mod tests { }, ) .unwrap(); - let client = Arc::new(UploadClient::new(&config, "host", "linux", "test")); + let client = Arc::new(UploadClient::new( + &config, + "host", + "linux", + "test", + Arc::new(FixedClock { + wall: 0.0, + mono: 0.0, + }), + )); let service = SyncService::start( config.clone(), client, diff --git a/crates/solstone-linux/src/test_support.rs b/crates/solstone-linux/src/test_support.rs index f2ad520..45a1c20 100644 --- a/crates/solstone-linux/src/test_support.rs +++ b/crates/solstone-linux/src/test_support.rs @@ -1,6 +1,7 @@ // SPDX-License-Identifier: AGPL-3.0-only // Copyright (c) 2026 sol pbc +use crate::observer::Clock; use http_body_util::{BodyExt, Full}; use hyper::{ Request, Response, @@ -14,7 +15,10 @@ use std::{ collections::VecDeque, convert::Infallible, pin::Pin, - sync::{Arc, Mutex}, + sync::{ + Arc, Mutex, + atomic::{AtomicU64, Ordering}, + }, task::{Context, Poll}, }; use tokio::{ @@ -25,6 +29,38 @@ use tokio::{ type BoxError = Box; +pub(crate) struct MutableClock { + wall: AtomicU64, + mono: AtomicU64, +} + +impl MutableClock { + pub(crate) fn new(wall: f64, mono: f64) -> Self { + Self { + wall: AtomicU64::new(wall.to_bits()), + mono: AtomicU64::new(mono.to_bits()), + } + } + + pub(crate) fn set_wall(&self, value: f64) { + self.wall.store(value.to_bits(), Ordering::Release); + } + + pub(crate) fn set_mono(&self, value: f64) { + self.mono.store(value.to_bits(), Ordering::Release); + } +} + +impl Clock for MutableClock { + fn wall_seconds(&self) -> f64 { + f64::from_bits(self.wall.load(Ordering::Acquire)) + } + + fn monotonic_seconds(&self) -> f64 { + f64::from_bits(self.mono.load(Ordering::Acquire)) + } +} + struct ReceiverBody(mpsc::Receiver>); impl Body for ReceiverBody { @@ -64,6 +100,7 @@ pub(crate) enum Action { OwnedRaw(u16, String), Disconnect, Stream(u16, mpsc::Receiver>), + GatedResponse(u16, Value, Arc), } impl MockServer { @@ -146,6 +183,15 @@ impl MockServer { Action::Stream(status, receiver) => { (status, ReceiverBody(receiver).boxed()) } + Action::GatedResponse(status, value, gate) => { + gate.notified().await; + ( + status, + Full::new(Bytes::from(value.to_string())) + .map_err(|never| match never {}) + .boxed(), + ) + } }; Ok::<_, std::io::Error>( Response::builder() diff --git a/crates/solstone-linux/src/upload.rs b/crates/solstone-linux/src/upload.rs index 9ef5afb..df4ace2 100644 --- a/crates/solstone-linux/src/upload.rs +++ b/crates/solstone-linux/src/upload.rs @@ -2,8 +2,9 @@ // Copyright (c) 2026 sol pbc use crate::{ - config::{Config, save_config}, + config::{Config, ConfigPaths, save_identity}, event_sender::{EventSender, SILENT_QUEUE_MAX}, + observer::Clock, sync_health::ErrorType, }; use reqwest::{Client, StatusCode, multipart}; @@ -13,7 +14,7 @@ use std::{ path::{Path, PathBuf}, sync::{ Arc, Mutex, - atomic::{AtomicBool, Ordering}, + atomic::{AtomicBool, AtomicU64, Ordering}, }, time::Duration, }; @@ -25,6 +26,10 @@ const STREAM_TYPE: &str = "desktop"; const OBSERVER_PROTOCOL_VERSION_HEADER: &str = "X-Solstone-Protocol-Version"; const DEFAULT_RETRY_DELAYS: [i64; 4] = [5, 30, 120, 300]; const MAX_IMMEDIATE_ATTEMPTS: usize = 2; +const RECOVERY_COOLDOWN: Duration = Duration::from_secs(300); +const TELLING_WINDOW: Duration = Duration::from_secs(300); +const TELLING_BURST_LIMIT: usize = 12; +pub(crate) const KEY_PREFIX_CHARS: usize = 8; #[derive(Clone, Debug, PartialEq, Eq)] pub struct UploadResult { @@ -91,6 +96,27 @@ pub(crate) struct Inner { version: String, retry_delays: Vec, immediate_attempts: usize, + paths: ConfigPaths, + clock: Arc, + recovery_lock: tokio::sync::Mutex<()>, + last_recovery_attempt: Mutex>, + recovery_generation: AtomicU64, + telling: Mutex, + dropped_events: AtomicU64, +} + +#[derive(Default)] +struct TellingState { + window_start: Option, + records: usize, + rejection_warned: bool, + suppressed_records: usize, +} + +enum RegisterAttempt { + Registered { key: String, name: String }, + GuardRefused { reason_code: Option }, + Failed, } pub struct UploadClient { @@ -110,8 +136,9 @@ impl UploadClient { hostname: impl Into, platform: impl Into, version: impl Into, + clock: Arc, ) -> Self { - Self::with_silent_capacity(config, hostname, platform, version, SILENT_QUEUE_MAX) + Self::with_silent_capacity(config, hostname, platform, version, clock, SILENT_QUEUE_MAX) } fn with_silent_capacity( @@ -119,6 +146,7 @@ impl UploadClient { hostname: impl Into, platform: impl Into, version: impl Into, + clock: Arc, silent_capacity: usize, ) -> Self { let retry_delays = if config.sync_retry_delays.is_empty() { @@ -140,6 +168,16 @@ impl UploadClient { immediate_attempts: config .sync_max_retries .clamp(1, MAX_IMMEDIATE_ATTEMPTS as i64) as usize, + paths: ConfigPaths { + base_dir: Some(config.base_dir.clone()), + config_dir: Some(config.config_dir.clone()), + }, + clock, + recovery_lock: tokio::sync::Mutex::new(()), + last_recovery_attempt: Mutex::new(None), + recovery_generation: AtomicU64::new(0), + telling: Mutex::new(TellingState::default()), + dropped_events: AtomicU64::new(0), }); let event_sender = EventSender::with_capacity(Arc::clone(&inner), silent_capacity); Self { @@ -156,6 +194,10 @@ impl UploadClient { !self.inner.key.lock().unwrap().is_empty() } + pub(crate) fn recovery_generation(&self) -> u64 { + self.inner.recovery_generation.load(Ordering::Acquire) + } + pub fn request_stop(&self) { self.inner.cancellation.cancel(); } @@ -172,77 +214,37 @@ impl UploadClient { if self.inner.url.is_empty() { return false; } - let stream = self.inner.stream.lock().unwrap().clone(); - let mut descriptor = json!({ - "platform": self.inner.platform, - "hostname": self.inner.hostname, - "stream_type": STREAM_TYPE, - "version": self.inner.version, - }); - if !stream.is_empty() { - descriptor["label"] = Value::String(stream); + let _registration = self.inner.recovery_lock.lock().await; + if self.is_registered() { + return true; } let attempts = 3.min(self.inner.retry_delays.len()); - let url = format!("{}/app/observer/register", self.inner.url); for attempt in 0..attempts { - let response = self - .inner - .client - .post(&url) - .json(&descriptor) - .timeout(EVENT_TIMEOUT) - .send() - .await; - match response { - Ok(response) if response.status() == StatusCode::OK => { - let body = match response.json::().await { - Ok(body) => body, - Err(error) => { - tracing::warn!( - attempt = attempt + 1, - %error, - "Registration attempt returned malformed JSON" - ); - if attempt + 1 < attempts { - tokio::time::sleep(retry_delay(&self.inner.retry_delays, attempt)) - .await; - } - continue; - } - }; - let (Some(key), Some(name)) = ( - body.get("key").and_then(Value::as_str), - body.get("name").and_then(Value::as_str), - ) else { - // Named deviation: Python raises KeyError for a malformed - // success body; Rust reports registration failure without panicking. - return false; - }; - config.key = key.to_owned(); - config.stream = name.to_owned(); - if let Err(error) = save_config(config) { - // Named deviation: Python propagates the persistence error; - // Rust leaves the client unregistered so a later call can retry. - tracing::error!(%error, "Failed to persist observer registration"); + match self.inner.register_once().await { + RegisterAttempt::Registered { key, name } => { + if let Err(error) = save_identity(&self.inner.paths, &key, &name) { + // Named deviation: Python propagates the persistence error; Rust keeps all + // three identity stores unchanged so a later call can retry. + tracing::error!(%error, "Failed to persist registration"); return false; } - *self.inner.key.lock().unwrap() = key.to_owned(); - *self.inner.stream.lock().unwrap() = name.to_owned(); - tracing::info!(name, "Registered observer"); + *self.inner.key.lock().unwrap() = key.clone(); + *self.inner.stream.lock().unwrap() = name.clone(); + config.key = key; + config.stream = name.clone(); + tracing::info!(name, "Registered"); return true; } - Ok(response) if response.status() == StatusCode::FORBIDDEN => { - self.inner.revoked.store(true, Ordering::Release); - tracing::error!("Registration rejected (403)"); + RegisterAttempt::GuardRefused { reason_code } => { + tracing::error!( + status = 403, + reason_code = reason_code.as_deref(), + "Journal refused local registration" + ); return false; } - Ok(response) => tracing::warn!( - attempt = attempt + 1, - status = %response.status(), - "Registration attempt failed" - ), - Err(error) => { - tracing::warn!(attempt = attempt + 1, %error, "Registration attempt failed") + RegisterAttempt::Failed => { + tracing::warn!(attempt = attempt + 1, "Registration attempt failed") } } if attempt + 1 < attempts { @@ -353,6 +355,8 @@ impl UploadClient { last_status = Some(status.as_u16()); if status == StatusCode::FORBIDDEN { self.inner.revoked.store(true, Ordering::Release); + } else if status == StatusCode::UNAUTHORIZED { + self.inner.recover_after_401("upload", false).await; } if error_type != ErrorType::Transient { tracing::error!(%status, ?error_type, "Upload rejected"); @@ -404,6 +408,8 @@ impl UploadClient { let error_type = Self::classify_error(Some(status.as_u16()), false); if status == StatusCode::FORBIDDEN { self.inner.revoked.store(true, Ordering::Release); + } else if status == StatusCode::UNAUTHORIZED { + self.inner.recover_after_401("listing", false).await; } tracing::warn!(%status, ?error_type, "Segments query failed"); return query_failure(error_type, Some(status.as_u16())); @@ -449,6 +455,190 @@ impl UploadClient { } impl Inner { + async fn register_once(&self) -> RegisterAttempt { + let stream = self.stream.lock().unwrap().clone(); + let mut descriptor = json!({ + "platform": self.platform, + "hostname": self.hostname, + "stream_type": STREAM_TYPE, + "version": self.version, + }); + if !stream.is_empty() { + descriptor["label"] = Value::String(stream); + } + let request = self + .client + .post(format!("{}/app/observer/register", self.url)) + .json(&descriptor) + .timeout(EVENT_TIMEOUT) + .send(); + let response = tokio::select! { + response = request => response, + () = self.cancellation.cancelled() => return RegisterAttempt::Failed, + }; + let Ok(response) = response else { + return RegisterAttempt::Failed; + }; + if response.status() == StatusCode::FORBIDDEN { + let reason_code = response.json::().await.ok().and_then(|body| { + body.get("reason_code") + .and_then(Value::as_str) + .map(str::to_owned) + }); + return RegisterAttempt::GuardRefused { reason_code }; + } + if response.status() != StatusCode::OK { + return RegisterAttempt::Failed; + } + let Ok(body) = response.json::().await else { + return RegisterAttempt::Failed; + }; + let (Some(key), Some(name)) = ( + body.get("key").and_then(Value::as_str), + body.get("name").and_then(Value::as_str), + ) else { + return RegisterAttempt::Failed; + }; + RegisterAttempt::Registered { + key: key.to_owned(), + name: name.to_owned(), + } + } + + fn recovery_on_cooldown(&self, now: f64) -> bool { + self.last_recovery_attempt + .lock() + .unwrap() + .is_some_and(|last| now - last < RECOVERY_COOLDOWN.as_secs_f64()) + } + + fn prepare_telling(&self, now: f64) -> std::sync::MutexGuard<'_, TellingState> { + let mut telling = self.telling.lock().unwrap(); + if telling + .window_start + .is_none_or(|start| now - start >= TELLING_WINDOW.as_secs_f64()) + { + *telling = TellingState { + window_start: Some(now), + ..TellingState::default() + }; + } + telling + } + + fn tell_first_rejection(&self, route: &str, key: &str, now: f64) { + let mut telling = self.prepare_telling(now); + if telling.rejection_warned { + return; + } + telling.rejection_warned = true; + if telling.records >= TELLING_BURST_LIMIT { + telling.suppressed_records += 1; + return; + } + telling.records += 1; + drop(telling); + tracing::warn!( + reason = "unauthorized", + route, + status = 401, + key_prefix = key_prefix(key), + recovery_generation = self.recovery_generation.load(Ordering::Acquire), + "Journal rejected the current key; attempting identity repair" + ); + } + + fn tell_outcome(&self, outcome: &str, name: &str, old_key: &str, new_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; + drop(telling); + let dropped_events_lower_bound = self.dropped_events.swap(0, Ordering::AcqRel); + tracing::warn!( + outcome, + name, + old_key_prefix = key_prefix(old_key), + new_key_prefix = key_prefix(new_key), + recovery_generation = self.recovery_generation.load(Ordering::Acquire), + dropped_events_lower_bound, + dropped_events_note = + "lower bound; queued, superseded, or otherwise rejected events are not included", + "Journal identity repair completed" + ); + } + + 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; + drop(telling); + tracing::error!( + status = 403, + reason_code, + key_prefix = key_prefix(key), + recovery_generation = self.recovery_generation.load(Ordering::Acquire), + "Journal refused local identity repair" + ); + } + + async fn recover_after_401(&self, route: &str, rejected_event: bool) { + // A 403 latch is an authorization boundary: no trigger may attempt recovery once set. + if self.revoked.load(Ordering::Acquire) { + return; + } + if rejected_event { + self.dropped_events.fetch_add(1, Ordering::AcqRel); + } + let now = self.clock.monotonic_seconds(); + let current_key = self.key.lock().unwrap().clone(); + self.tell_first_rejection(route, ¤t_key, now); + if self.recovery_on_cooldown(now) { + return; + } + + let _recovery = self.recovery_lock.lock().await; + // Recheck both authorization and cooldown after waiting for a concurrent trigger. + if self.revoked.load(Ordering::Acquire) { + return; + } + let now = self.clock.monotonic_seconds(); + if self.recovery_on_cooldown(now) { + return; + } + *self.last_recovery_attempt.lock().unwrap() = Some(now); + // Mitigation 1 compares against the key rejected by this attempt, never a prior response. + let rejected_key = self.key.lock().unwrap().clone(); + let attempt = self.register_once().await; + if self.revoked.load(Ordering::Acquire) { + return; + } + match attempt { + RegisterAttempt::Registered { key, name } if key == rejected_key => { + self.tell_outcome("same_key_returned", &name, &rejected_key, &key, now); + } + RegisterAttempt::Registered { key, name } => { + if let Err(error) = save_identity(&self.paths, &key, &name) { + tracing::error!(%error, "Failed to persist repaired journal identity"); + return; + } + *self.key.lock().unwrap() = key.clone(); + *self.stream.lock().unwrap() = name.clone(); + self.recovery_generation.fetch_add(1, Ordering::AcqRel); + self.tell_outcome("recovered", &name, &rejected_key, &key, now); + } + RegisterAttempt::GuardRefused { reason_code } => { + self.tell_guard_refusal(reason_code.as_deref(), &rejected_key, now); + } + RegisterAttempt::Failed => {} + } + } + pub(crate) async fn relay_event( &self, tract: &str, @@ -478,6 +668,8 @@ impl Inner { Ok(response) => { if response.status() == StatusCode::FORBIDDEN { self.revoked.store(true, Ordering::Release); + } else if response.status() == StatusCode::UNAUTHORIZED { + self.recover_after_401("event", true).await; } false } @@ -486,6 +678,10 @@ impl Inner { } } +pub(crate) fn key_prefix(key: &str) -> String { + key.chars().take(KEY_PREFIX_CHARS).collect() +} + fn retry_delay(delays: &[i64], attempt: usize) -> Duration { Duration::from_secs(delays[attempt.min(delays.len() - 1)].max(0) as u64) } @@ -566,10 +762,25 @@ mod tests { use super::*; use crate::{ config::{ConfigPaths, load_config}, - test_support::{Action, MockServer, wait_for_requests}, + test_support::{Action, MockServer, MutableClock, wait_for_requests}, }; use tempfile::TempDir; use tokio::net::TcpListener; + use tracing::instrument::WithSubscriber; + + #[derive(Clone)] + struct Buffer(Arc>>); + + impl std::io::Write for Buffer { + fn write(&mut self, bytes: &[u8]) -> std::io::Result { + self.0.lock().unwrap().extend_from_slice(bytes); + Ok(bytes.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } fn config(server: &MockServer, temp: &TempDir) -> Config { Config { @@ -582,7 +793,13 @@ mod tests { } fn client(config: &Config) -> UploadClient { - UploadClient::new(config, "host-a", "linux", "0.1.0") + UploadClient::new( + config, + "host-a", + "linux", + "0.1.0", + Arc::new(MutableClock::new(0.0, 0.0)), + ) } fn write_file(temp: &TempDir, name: &str, body: &[u8]) -> PathBuf { @@ -620,6 +837,7 @@ mod tests { } // tests/test_upload.py::test_ensure_registered_skips_when_key_present + // AC 18: constructor clock plumbing deliberately preserves the existing-key short circuit. #[tokio::test] async fn ensure_registered_skips_when_key_present() { let server = MockServer::new(vec![]).await; @@ -629,6 +847,309 @@ mod tests { assert!(server.requests().is_empty()); } + // AC 1/2: a stale-key 401 registers, persists a distinct replacement, and the next relay uses it. + #[tokio::test] + async fn stale_key_401_issues_register_request() { + const STALE: &str = "STALE-KEY-111"; + const NEW: &str = "NEW-KEY-222"; + let server = MockServer::new(vec![ + (401, json!({})), + (200, json!({"key":NEW,"name":"desktop-new"})), + (200, json!({})), + ]) + .await; + let temp = TempDir::new().unwrap(); + let config = Config { + key: STALE.into(), + stream: "desktop".into(), + ..config(&server, &temp) + }; + save_identity( + &ConfigPaths { + base_dir: Some(config.base_dir.clone()), + config_dir: Some(config.config_dir.clone()), + }, + STALE, + "desktop", + ) + .unwrap(); + let client = client(&config); + + assert!(!client.relay_event("observe", "status", Map::new()).await); + assert_eq!(server.request_count("/app/observer/register"), 1); + assert_eq!(client.inner.key.lock().unwrap().as_str(), NEW); + assert_eq!(load_config(client.inner.paths.clone()).config.key, NEW); + assert!(client.relay_event("observe", "status", Map::new()).await); + let requests = server.requests(); + assert_eq!( + requests[2] + .headers + .get("authorization") + .unwrap() + .to_str() + .unwrap(), + format!("Bearer {NEW}") + ); + } + + // AC 3/7: an idempotent same-key response is not recovery and still arms cooldown. + #[tokio::test] + async fn same_key_registration_is_not_recovery_and_arms_cooldown() { + let server = MockServer::new(vec![ + (401, json!({})), + (200, json!({"key":"STALE-KEY","name":"desktop"})), + (401, json!({})), + ]) + .await; + let temp = TempDir::new().unwrap(); + let config = Config { + key: "STALE-KEY".into(), + stream: "desktop".into(), + ..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 { + assert!(!client.relay_event("observe", "status", Map::new()).await); + assert_eq!(client.recovery_generation(), 0); + assert!(!client.relay_event("observe", "status", Map::new()).await); + } + .with_subscriber(subscriber) + .await; + assert_eq!(server.request_count("/app/observer/register"), 1); + let captured = String::from_utf8(output.lock().unwrap().clone()).unwrap(); + assert!( + captured.contains("outcome=\"same_key_returned\""), + "{captured}" + ); + assert!(!captured.contains("outcome=\"recovered\""), "{captured}"); + } + + // AC 6: the real sync-query/event race shares one in-flight recovery attempt. + #[tokio::test] + async fn sync_and_event_401_share_one_inflight_recovery() { + let gate = Arc::new(tokio::sync::Notify::new()); + let server = MockServer::new_actions(vec![ + Action::GatedResponse(401, json!({}), Arc::clone(&gate)), + Action::GatedResponse(401, json!({}), Arc::clone(&gate)), + Action::Response(200, json!({"key":"NEW-KEY","name":"desktop-new"})), + ]) + .await; + let temp = TempDir::new().unwrap(); + let config = Config { + key: "STALE-KEY".into(), + stream: "desktop".into(), + ..config(&server, &temp) + }; + let client = Arc::new(client(&config)); + let listing = { + let client = Arc::clone(&client); + tokio::spawn(async move { client.get_server_segments("20260101").await }) + }; + let event = { + let client = Arc::clone(&client); + tokio::spawn(async move { client.relay_event("observe", "status", Map::new()).await }) + }; + wait_for_requests(&server, 2).await; + tokio::task::yield_now().await; + gate.notify_waiters(); + let (listing, event) = tokio::join!(listing, event); + let listing = listing.unwrap(); + let event = event.unwrap(); + assert!(listing.segments.is_none()); + assert!(!event); + assert_eq!(server.request_count("/app/observer/register"), 1); + } + + // AC 7: a burst is bounded by cooldown and the next 300-second window starts fresh. + #[tokio::test] + async fn recovery_cooldown_bounds_register_requests_and_resets() { + let mut responses = vec![ + (401, json!({})), + (200, json!({"key":"STALE-KEY","name":"desktop"})), + ]; + responses.extend((0..19).map(|_| (401, json!({})))); + responses.extend([ + (401, json!({})), + (200, json!({"key":"STALE-KEY","name":"desktop"})), + ]); + let server = MockServer::new(responses).await; + let temp = TempDir::new().unwrap(); + let config = Config { + key: "STALE-KEY".into(), + stream: "desktop".into(), + ..config(&server, &temp) + }; + let clock = Arc::new(MutableClock::new(0.0, 0.0)); + let client = UploadClient::new(&config, "host", "linux", "test", clock.clone()); + for _ in 0..20 { + assert!(!client.relay_event("observe", "status", Map::new()).await); + } + assert_eq!(server.request_count("/app/observer/register"), 1); + clock.set_mono(300.0); + assert!(!client.relay_event("observe", "status", Map::new()).await); + assert_eq!(server.request_count("/app/observer/register"), 2); + } + + // AC 9: every failed register shape preserves disk, runtime, caller config, and stale bearer use. + #[tokio::test] + async fn failed_registration_preserves_all_stores_and_stale_bearer() { + for failed in [ + Action::Disconnect, + Action::Response(500, json!({})), + Action::Response(200, json!({"name":"desktop"})), + ] { + let server = MockServer::new_actions(vec![ + Action::Response(401, json!({})), + failed, + Action::Response(200, json!({})), + ]) + .await; + let temp = TempDir::new().unwrap(); + let config = Config { + key: "STALE-KEY".into(), + 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 before = config.clone(); + let client = client(&config); + + assert!(!client.relay_event("observe", "status", Map::new()).await); + assert_eq!(config, before); + assert_eq!(client.inner.key.lock().unwrap().as_str(), "STALE-KEY"); + assert_eq!( + load_config(client.inner.paths.clone()).config.key, + "STALE-KEY" + ); + assert!(client.relay_event("observe", "status", Map::new()).await); + let requests = server.requests(); + assert_eq!( + requests.last().unwrap().headers["authorization"], + "Bearer STALE-KEY" + ); + } + } + + // AC 8: cancellation around a successful register exposes only complete stale/new identities. + #[tokio::test] + async fn cancelled_recovery_exposes_no_third_disk_identity() { + const STALE: &str = "STALE-KEY"; + const NEW: &str = "NEW-KEY"; + for cancel in [true, false] { + let gate = Arc::new(tokio::sync::Notify::new()); + let server = MockServer::new_actions(vec![ + Action::Response(401, json!({})), + Action::GatedResponse( + 200, + json!({"key":NEW,"name":"desktop-new"}), + Arc::clone(&gate), + ), + ]) + .await; + let temp = TempDir::new().unwrap(); + let config = Config { + key: STALE.into(), + stream: "desktop-old".into(), + ..config(&server, &temp) + }; + save_identity( + &ConfigPaths { + base_dir: Some(config.base_dir.clone()), + config_dir: Some(config.config_dir.clone()), + }, + STALE, + "desktop-old", + ) + .unwrap(); + let client = Arc::new(client(&config)); + let relay = { + let client = Arc::clone(&client); + tokio::spawn( + async move { client.relay_event("observe", "status", Map::new()).await }, + ) + }; + wait_for_requests(&server, 2).await; + if cancel { + client.request_stop(); + } + gate.notify_one(); + assert!(!relay.await.unwrap()); + let saved = load_config(client.inner.paths.clone()).config; + let pair = (saved.key.as_str(), saved.stream.as_str()); + assert!( + pair == (STALE, "desktop-old") || pair == (NEW, "desktop-new"), + "unexpected reachable identity {pair:?}" + ); + assert_eq!( + pair, + if cancel { + (STALE, "desktop-old") + } else { + (NEW, "desktop-new") + } + ); + } + } + + // AC 10: already-green regression pin; failed registration keeps keyless relay offline. + #[tokio::test] + async fn failed_registration_keeps_empty_key_and_relay_offline() { + let server = MockServer::new(vec![(500, json!({}))]).await; + let temp = TempDir::new().unwrap(); + let mut config = config(&server, &temp); + config.key.clear(); + config.sync_retry_delays = vec![0]; + let client = client(&config); + assert!(!client.ensure_registered(&mut config).await); + let before = server.requests().len(); + assert!(!client.relay_event("observe", "status", Map::new()).await); + assert_eq!(server.requests().len(), before); + assert!(!client.is_registered()); + } + + // AC 15: zero rejections and 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; + 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 { + for _ in 0..20 { + assert!(client.relay_event("observe", "status", Map::new()).await); + } + } + .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 @@ -1190,15 +1711,9 @@ mod tests { async fn assert_403_latches(path: &str) { let server = MockServer::new(vec![(403, json!({}))]).await; let temp = TempDir::new().unwrap(); - let mut config = config(&server, &temp); + let config = config(&server, &temp); let client = client(&config); match path { - "registration" => { - config.key.clear(); - let client = self::client(&config); - assert!(!client.ensure_registered(&mut config).await); - assert!(client.is_revoked()); - } "upload" => { let media = write_file(&temp, "a.flac", b"a"); assert!(!client.upload_segment("d", "s", &[media]).await.success); @@ -1216,10 +1731,55 @@ mod tests { } } - // AC: registration 403 latches revoked + // AC 18: register-route guard refusal deliberately no longer latches revocation. #[tokio::test] - async fn registration_403_latches_revoked() { - assert_403_latches("registration").await; + async fn registration_403_does_not_latch_revoked() { + let server = + MockServer::new(vec![(403, json!({"reason_code":"local_request_only"}))]).await; + let temp = TempDir::new().unwrap(); + let mut config = config(&server, &temp); + config.key.clear(); + let client = client(&config); + assert!(!client.ensure_registered(&mut config).await); + assert!(!client.is_revoked()); + assert!(config.key.is_empty()); + } + + // AC 5: register guard refusal preserves the stale identity; ingest 403 still revokes it. + #[tokio::test] + async fn register_guard_refusal_does_not_revoke_but_ingest_403_does() { + let server = MockServer::new(vec![ + (401, json!({})), + (403, json!({"reason_code":"local_request_only"})), + (403, json!({})), + ]) + .await; + let temp = TempDir::new().unwrap(); + let config = Config { + key: "STALE-KEY".into(), + stream: "desktop".into(), + ..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 { + assert!(!client.relay_event("observe", "status", Map::new()).await); + assert!(!client.is_revoked()); + assert_eq!(client.inner.key.lock().unwrap().as_str(), "STALE-KEY"); + assert!(!client.relay_event("observe", "status", Map::new()).await); + assert!(client.is_revoked()); + } + .with_subscriber(subscriber) + .await; + let captured = String::from_utf8(output.lock().unwrap().clone()).unwrap(); + assert!(captured.contains("Journal refused local identity repair")); + assert!(captured.contains("local_request_only")); } // AC: upload 403 latches revoked #[tokio::test] @@ -1265,8 +1825,14 @@ mod tests { async fn submit_status_is_nonblocking_while_relay_is_blocked() { let (server, gate) = MockServer::gated().await; let temp = TempDir::new().unwrap(); - let mut client = - UploadClient::with_silent_capacity(&config(&server, &temp), "host", "linux", "v", 1); + let mut client = UploadClient::with_silent_capacity( + &config(&server, &temp), + "host", + "linux", + "v", + Arc::new(MutableClock::new(0.0, 0.0)), + 1, + ); client.enqueue_status(Map::from_iter([("seq".into(), json!(1))])); wait_for_requests(&server, 1).await; let start = tokio::time::Instant::now(); @@ -1302,8 +1868,14 @@ mod tests { async fn silent_overflow_rejects_incoming_and_stop_is_bounded() { let (server, gate) = MockServer::gated().await; let temp = TempDir::new().unwrap(); - let mut client = - UploadClient::with_silent_capacity(&config(&server, &temp), "host", "linux", "v", 1); + let mut client = UploadClient::with_silent_capacity( + &config(&server, &temp), + "host", + "linux", + "v", + Arc::new(MutableClock::new(0.0, 0.0)), + 1, + ); assert!(client.enqueue_stream_silent(Map::from_iter([("node_id".into(), json!(1))]))); assert!(!client.enqueue_stream_silent(Map::from_iter([("node_id".into(), json!(2))]))); wait_for_requests(&server, 1).await; @@ -1315,4 +1887,35 @@ mod tests { let body: Value = serde_json::from_slice(&server.requests()[0].body).unwrap(); assert_eq!(body["node_id"], 1); } + + // AC 13: event recovery performs one register attempt while the bounded queue remains usable. + #[tokio::test] + async fn event_recovery_is_single_attempt_and_queue_stays_bounded() { + let gate = Arc::new(tokio::sync::Notify::new()); + let server = MockServer::new_actions(vec![ + Action::Response(401, json!({})), + Action::GatedResponse(500, json!({}), Arc::clone(&gate)), + ]) + .await; + let temp = TempDir::new().unwrap(); + let config = Config { + key: "STALE-KEY".into(), + stream: "desktop".into(), + ..config(&server, &temp) + }; + let mut client = client(&config); + client.enqueue_status(Map::new()); + wait_for_requests(&server, 2).await; + for sequence in 0..20 { + assert!( + client + .enqueue_stream_silent(Map::from_iter([("sequence".into(), json!(sequence),)])) + ); + } + let started = tokio::time::Instant::now(); + gate.notify_one(); + assert_eq!(client.stop(Duration::from_secs(1)).await, 0); + assert!(started.elapsed() < Duration::from_secs(1)); + assert_eq!(server.request_count("/app/observer/register"), 1); + } } -- 2.51.2