diff --git a/crates/session-store/src/worker_metadata.rs b/crates/session-store/src/worker_metadata.rs index 3549dbb9..52bca7ae 100644 --- a/crates/session-store/src/worker_metadata.rs +++ b/crates/session-store/src/worker_metadata.rs @@ -14,8 +14,24 @@ use crate::{SegmentId, SessionId}; use serde::{Deserialize, Serialize}; +use std::collections::HashMap; use std::fs; -use std::path::PathBuf; +use std::path::{Path, PathBuf}; +use std::sync::{Arc, Mutex, OnceLock, Weak}; + +fn metadata_lock(path: &Path) -> Arc> { + static LOCKS: OnceLock>>>> = OnceLock::new(); + let mut locks = LOCKS + .get_or_init(|| Mutex::new(HashMap::new())) + .lock() + .expect("metadata lock registry poisoned"); + if let Some(lock) = locks.get(path).and_then(Weak::upgrade) { + return lock; + } + let lock = Arc::new(Mutex::new(())); + locks.insert(path.to_path_buf(), Arc::downgrade(&lock)); + lock +} /// Errors from Worker metadata persistence. #[derive(Debug, thiserror::Error)] @@ -348,6 +364,7 @@ pub trait WorkerMetadataStore: Send + Sync { pub struct WorkerAggregateStore { root: PathBuf, worker_name: String, + update_lock: Arc>, } impl WorkerAggregateStore { @@ -359,7 +376,11 @@ impl WorkerAggregateStore { let worker_name = worker_name.into(); validate_worker_name(&worker_name)?; fs::create_dir_all(&root)?; - Ok(Self { root, worker_name }) + Ok(Self { + update_lock: metadata_lock(&root), + root, + worker_name, + }) } fn validate_name(&self, worker_name: &str) -> Result<(), WorkerStoreError> { @@ -426,6 +447,47 @@ impl WorkerMetadataStore for WorkerAggregateStore { Ok(Some(metadata)) } + fn update_by_name( + &self, + worker_name: &str, + update: F, + ) -> Result + where + F: FnOnce(&mut WorkerMetadata), + { + let _guard = self + .update_lock + .lock() + .expect("metadata update lock poisoned"); + let mut metadata = self + .read_by_name(worker_name)? + .unwrap_or_else(|| WorkerMetadata::new(worker_name, None)); + update(&mut metadata); + self.write(&metadata)?; + Ok(metadata) + } + + fn compare_and_swap_active( + &self, + worker_name: &str, + expected: &WorkerActiveSegmentRef, + replacement: WorkerActiveSegmentRef, + ) -> Result { + let _guard = self + .update_lock + .lock() + .expect("metadata update lock poisoned"); + let Some(mut metadata) = self.read_by_name(worker_name)? else { + return Ok(false); + }; + if metadata.active.as_ref() != Some(expected) { + return Ok(false); + } + metadata.active = Some(replacement); + self.write(&metadata)?; + Ok(true) + } + fn list_names(&self) -> Result, WorkerStoreError> { Ok(if self.metadata_path().is_file() { vec![self.worker_name.clone()] @@ -452,6 +514,7 @@ impl WorkerMetadataStore for WorkerAggregateStore { #[derive(Clone)] pub struct FsWorkerStore { root: PathBuf, + update_lock: Arc>, } impl FsWorkerStore { @@ -459,7 +522,10 @@ impl FsWorkerStore { pub fn new(root: impl Into) -> Result { let root = root.into(); fs::create_dir_all(&root)?; - Ok(Self { root }) + Ok(Self { + update_lock: metadata_lock(&root), + root, + }) } fn worker_dir(&self, worker_name: &str) -> Result { @@ -475,12 +541,32 @@ impl FsWorkerStore { impl WorkerMetadataStore for FsWorkerStore { fn write(&self, metadata: &WorkerMetadata) -> Result<(), WorkerStoreError> { let path = self.metadata_path(&metadata.worker_name)?; - if let Some(parent) = path.parent() { - fs::create_dir_all(parent)?; + let mut content = serde_json::to_vec_pretty(metadata)?; + content.push(b'\n'); + let parent = path.parent().expect("metadata path has parent"); + fs::create_dir_all(parent)?; + let temp = parent.join(format!( + ".metadata.json.tmp-{}-{}", + std::process::id(), + uuid::Uuid::now_v7() + )); + let result = (|| -> Result<(), WorkerStoreError> { + use std::io::Write; + let mut file = std::fs::OpenOptions::new() + .write(true) + .create_new(true) + .open(&temp)?; + file.write_all(&content)?; + file.sync_all()?; + drop(file); + fs::rename(&temp, &path)?; + std::fs::File::open(parent)?.sync_all()?; + Ok(()) + })(); + if result.is_err() { + let _ = fs::remove_file(temp); } - let content = serde_json::to_vec_pretty(metadata)?; - fs::write(path, content)?; - Ok(()) + result } fn read_by_name(&self, worker_name: &str) -> Result, WorkerStoreError> { @@ -493,6 +579,47 @@ impl WorkerMetadataStore for FsWorkerStore { Ok(Some(serde_json::from_str(&content)?)) } + fn update_by_name( + &self, + worker_name: &str, + update: F, + ) -> Result + where + F: FnOnce(&mut WorkerMetadata), + { + let _guard = self + .update_lock + .lock() + .expect("metadata update lock poisoned"); + let mut metadata = self + .read_by_name(worker_name)? + .unwrap_or_else(|| WorkerMetadata::new(worker_name, None)); + update(&mut metadata); + self.write(&metadata)?; + Ok(metadata) + } + + fn compare_and_swap_active( + &self, + worker_name: &str, + expected: &WorkerActiveSegmentRef, + replacement: WorkerActiveSegmentRef, + ) -> Result { + let _guard = self + .update_lock + .lock() + .expect("metadata update lock poisoned"); + let Some(mut metadata) = self.read_by_name(worker_name)? else { + return Ok(false); + }; + if metadata.active.as_ref() != Some(expected) { + return Ok(false); + } + metadata.active = Some(replacement); + self.write(&metadata)?; + Ok(true) + } + fn list_names(&self) -> Result, WorkerStoreError> { let mut names = Vec::new(); if !self.root.exists() { @@ -902,4 +1029,38 @@ mod tests { assert_eq!(restored.reclaimed_children.len(), 1); assert_eq!(restored.reclaimed_children[0].scope_delegated, vec![scope]); } + + #[test] + fn active_segment_cas_allows_exactly_one_concurrent_winner() { + let temp = tempfile::tempdir().unwrap(); + let store = FsWorkerStore::new(temp.path()).unwrap(); + let session_id = crate::new_session_id(); + let old = WorkerActiveSegmentRef::active_segment(session_id, crate::new_segment_id()); + store + .write(&WorkerMetadata::new("agent", Some(old.clone()))) + .unwrap(); + let barrier = Arc::new(std::sync::Barrier::new(3)); + let handles = [crate::new_segment_id(), crate::new_segment_id()].map(|segment_id| { + let store = store.clone(); + let old = old.clone(); + let barrier = barrier.clone(); + std::thread::spawn(move || { + barrier.wait(); + store + .compare_and_swap_active( + "agent", + &old, + WorkerActiveSegmentRef::active_segment(session_id, segment_id), + ) + .unwrap() + }) + }); + barrier.wait(); + let winners = handles + .into_iter() + .map(|handle| handle.join().unwrap()) + .filter(|won| *won) + .count(); + assert_eq!(winners, 1); + } } diff --git a/crates/tui/src/app.rs b/crates/tui/src/app.rs index 938a205f..86722b63 100644 --- a/crates/tui/src/app.rs +++ b/crates/tui/src/app.rs @@ -279,7 +279,7 @@ pub struct App { run_error_messages: Vec, /// Current compaction identity/revision used to fence snapshot/live updates. active_compaction: Option<(String, u64)>, - compaction_progress: Option, + pub compaction_progress: Option, /// Presentation-only Internal Worker projections keyed by session identity. /// They are rendered in separate selectable views and never mixed into `blocks`. pub internal_workers: Vec, diff --git a/crates/tui/src/ui.rs b/crates/tui/src/ui.rs index 58b790f1..b23209df 100644 --- a/crates/tui/src/ui.rs +++ b/crates/tui/src/ui.rs @@ -151,7 +151,7 @@ fn run_status_line(app: &App, now: Instant) -> Line<'static> { format!("{} reqs", app.run_requests) }; - Line::from(vec![ + let mut spans = vec![ Span::styled( RUN_SPINNER_FRAMES[spinner_index], Style::default() @@ -159,6 +159,15 @@ fn run_status_line(app: &App, now: Instant) -> Line<'static> { .add_modifier(Modifier::BOLD), ), Span::raw(" "), + ]; + if let Some(progress) = &app.compaction_progress { + spans.push(Span::styled( + format!("Compacting ยท {:?}", progress.phase).to_lowercase(), + Style::default().fg(Color::Cyan), + )); + spans.push(Span::styled(" | ", Style::default().fg(Color::DarkGray))); + } + spans.extend([ Span::styled( fmt_run_elapsed(elapsed.as_secs()), Style::default().fg(Color::Gray), @@ -177,7 +186,8 @@ fn run_status_line(app: &App, now: Instant) -> Line<'static> { fmt_tokens(app.run_output_tokens), Style::default().fg(Color::Yellow), ), - ]) + ]); + Line::from(spans) } fn fmt_run_elapsed(secs: u64) -> String { diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index b46f922d..dc713560 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -5610,7 +5610,18 @@ impl Worker { session_id: old_loc.session_id, segment_id: new_segment_id, }; - self.compare_and_swap_worker_metadata_segment(old_loc, new_location)?; + // Move the live writer lease under the same exclusive compaction owner + // before publishing the durable pointer. If the CAS loses, restore the + // derived lease and leave the staged Segment unreachable. + if self.scope_allocation.is_some() { + worker_allocation::update_segment(&self.manifest.worker.name, new_segment_id)?; + } + if let Err(error) = self.compare_and_swap_worker_metadata_segment(old_loc, new_location) { + if self.scope_allocation.is_some() { + worker_allocation::update_segment(&self.manifest.worker.name, old_loc.segment_id)?; + } + return Err(error); + } // All live mutations after the durable commit are infallible and happen // before the replacement SegmentStart is broadcast. This keeps the @@ -5632,24 +5643,6 @@ impl Worker { .expect("usage_history poisoned") .clear(); - // workers.json is live writer-bookkeeping derived from the durable - // metadata pointer above. A stale value still retains the worker-name - // lease and cannot authorize another writer; after a process exit it is - // reclaimed through the normal stale-allocation path. Do not report a - // committed compaction as failed solely because this derived projection - // could not be refreshed. - if self.scope_allocation.is_some() - && let Err(error) = - worker_allocation::update_segment(&self.manifest.worker.name, new_segment_id) - { - warn!( - worker = %self.manifest.worker.name, - segment_id = %new_segment_id, - error = %error, - "compaction committed but live writer allocation projection could not be refreshed" - ); - } - // Broadcast only after every live authority points at the committed // replacement. Runtime-owned extensions stay visible and restorable. self.sink