fix: serialize compaction pointer commits

This commit is contained in:
2026-09-16 04:41:27 +09:00
parent df22526a0d
commit c0fe20e8a2
4 changed files with 194 additions and 30 deletions
+168 -7
View File
@@ -14,8 +14,24 @@
use crate::{SegmentId, SessionId}; use crate::{SegmentId, SessionId};
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::fs; 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<Mutex<()>> {
static LOCKS: OnceLock<Mutex<HashMap<PathBuf, Weak<Mutex<()>>>>> = 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. /// Errors from Worker metadata persistence.
#[derive(Debug, thiserror::Error)] #[derive(Debug, thiserror::Error)]
@@ -348,6 +364,7 @@ pub trait WorkerMetadataStore: Send + Sync {
pub struct WorkerAggregateStore { pub struct WorkerAggregateStore {
root: PathBuf, root: PathBuf,
worker_name: String, worker_name: String,
update_lock: Arc<Mutex<()>>,
} }
impl WorkerAggregateStore { impl WorkerAggregateStore {
@@ -359,7 +376,11 @@ impl WorkerAggregateStore {
let worker_name = worker_name.into(); let worker_name = worker_name.into();
validate_worker_name(&worker_name)?; validate_worker_name(&worker_name)?;
fs::create_dir_all(&root)?; 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> { fn validate_name(&self, worker_name: &str) -> Result<(), WorkerStoreError> {
@@ -426,6 +447,47 @@ impl WorkerMetadataStore for WorkerAggregateStore {
Ok(Some(metadata)) Ok(Some(metadata))
} }
fn update_by_name<F>(
&self,
worker_name: &str,
update: F,
) -> Result<WorkerMetadata, WorkerStoreError>
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<bool, WorkerStoreError> {
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<Vec<String>, WorkerStoreError> { fn list_names(&self) -> Result<Vec<String>, WorkerStoreError> {
Ok(if self.metadata_path().is_file() { Ok(if self.metadata_path().is_file() {
vec![self.worker_name.clone()] vec![self.worker_name.clone()]
@@ -452,6 +514,7 @@ impl WorkerMetadataStore for WorkerAggregateStore {
#[derive(Clone)] #[derive(Clone)]
pub struct FsWorkerStore { pub struct FsWorkerStore {
root: PathBuf, root: PathBuf,
update_lock: Arc<Mutex<()>>,
} }
impl FsWorkerStore { impl FsWorkerStore {
@@ -459,7 +522,10 @@ impl FsWorkerStore {
pub fn new(root: impl Into<PathBuf>) -> Result<Self, WorkerStoreError> { pub fn new(root: impl Into<PathBuf>) -> Result<Self, WorkerStoreError> {
let root = root.into(); let root = root.into();
fs::create_dir_all(&root)?; fs::create_dir_all(&root)?;
Ok(Self { root }) Ok(Self {
update_lock: metadata_lock(&root),
root,
})
} }
fn worker_dir(&self, worker_name: &str) -> Result<PathBuf, WorkerStoreError> { fn worker_dir(&self, worker_name: &str) -> Result<PathBuf, WorkerStoreError> {
@@ -475,12 +541,32 @@ impl FsWorkerStore {
impl WorkerMetadataStore for FsWorkerStore { impl WorkerMetadataStore for FsWorkerStore {
fn write(&self, metadata: &WorkerMetadata) -> Result<(), WorkerStoreError> { fn write(&self, metadata: &WorkerMetadata) -> Result<(), WorkerStoreError> {
let path = self.metadata_path(&metadata.worker_name)?; let path = self.metadata_path(&metadata.worker_name)?;
if let Some(parent) = path.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)?; fs::create_dir_all(parent)?;
} let temp = parent.join(format!(
let content = serde_json::to_vec_pretty(metadata)?; ".metadata.json.tmp-{}-{}",
fs::write(path, content)?; 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(()) Ok(())
})();
if result.is_err() {
let _ = fs::remove_file(temp);
}
result
} }
fn read_by_name(&self, worker_name: &str) -> Result<Option<WorkerMetadata>, WorkerStoreError> { fn read_by_name(&self, worker_name: &str) -> Result<Option<WorkerMetadata>, WorkerStoreError> {
@@ -493,6 +579,47 @@ impl WorkerMetadataStore for FsWorkerStore {
Ok(Some(serde_json::from_str(&content)?)) Ok(Some(serde_json::from_str(&content)?))
} }
fn update_by_name<F>(
&self,
worker_name: &str,
update: F,
) -> Result<WorkerMetadata, WorkerStoreError>
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<bool, WorkerStoreError> {
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<Vec<String>, WorkerStoreError> { fn list_names(&self) -> Result<Vec<String>, WorkerStoreError> {
let mut names = Vec::new(); let mut names = Vec::new();
if !self.root.exists() { if !self.root.exists() {
@@ -902,4 +1029,38 @@ mod tests {
assert_eq!(restored.reclaimed_children.len(), 1); assert_eq!(restored.reclaimed_children.len(), 1);
assert_eq!(restored.reclaimed_children[0].scope_delegated, vec![scope]); 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);
}
} }
+1 -1
View File
@@ -279,7 +279,7 @@ pub struct App {
run_error_messages: Vec<String>, run_error_messages: Vec<String>,
/// Current compaction identity/revision used to fence snapshot/live updates. /// Current compaction identity/revision used to fence snapshot/live updates.
active_compaction: Option<(String, u64)>, active_compaction: Option<(String, u64)>,
compaction_progress: Option<protocol::InFlightCompaction>, pub compaction_progress: Option<protocol::InFlightCompaction>,
/// Presentation-only Internal Worker projections keyed by session identity. /// Presentation-only Internal Worker projections keyed by session identity.
/// They are rendered in separate selectable views and never mixed into `blocks`. /// They are rendered in separate selectable views and never mixed into `blocks`.
pub internal_workers: Vec<InternalWorkerView>, pub internal_workers: Vec<InternalWorkerView>,
+12 -2
View File
@@ -151,7 +151,7 @@ fn run_status_line(app: &App, now: Instant) -> Line<'static> {
format!("{} reqs", app.run_requests) format!("{} reqs", app.run_requests)
}; };
Line::from(vec![ let mut spans = vec![
Span::styled( Span::styled(
RUN_SPINNER_FRAMES[spinner_index], RUN_SPINNER_FRAMES[spinner_index],
Style::default() Style::default()
@@ -159,6 +159,15 @@ fn run_status_line(app: &App, now: Instant) -> Line<'static> {
.add_modifier(Modifier::BOLD), .add_modifier(Modifier::BOLD),
), ),
Span::raw(" "), 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( Span::styled(
fmt_run_elapsed(elapsed.as_secs()), fmt_run_elapsed(elapsed.as_secs()),
Style::default().fg(Color::Gray), 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), fmt_tokens(app.run_output_tokens),
Style::default().fg(Color::Yellow), Style::default().fg(Color::Yellow),
), ),
]) ]);
Line::from(spans)
} }
fn fmt_run_elapsed(secs: u64) -> String { fn fmt_run_elapsed(secs: u64) -> String {
+12 -19
View File
@@ -5610,7 +5610,18 @@ impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
session_id: old_loc.session_id, session_id: old_loc.session_id,
segment_id: new_segment_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 // All live mutations after the durable commit are infallible and happen
// before the replacement SegmentStart is broadcast. This keeps the // before the replacement SegmentStart is broadcast. This keeps the
@@ -5632,24 +5643,6 @@ impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
.expect("usage_history poisoned") .expect("usage_history poisoned")
.clear(); .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 // Broadcast only after every live authority points at the committed
// replacement. Runtime-owned extensions stay visible and restorable. // replacement. Runtime-owned extensions stay visible and restorable.
self.sink self.sink