#[cfg(feature = "pipeline")]
use std::collections::HashSet;
use std::collections::{BTreeMap, HashMap};
use std::io;
use std::path::{Path, PathBuf};
use std::time::Duration;
use crate::primitives::fs;
use crate::primitives::sync::Mutex;
use crate::primitives::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use crate::rate_limit::rate_limited;
#[cfg(feature = "pipeline")]
use crate::sealed::find_sealed_segments;
use crate::sealed::{SealedSegment, SegmentArtifact, SegmentRef, parse_segment_artifact};
use super::{ActiveHandle, DiscoveredArtifacts, RemoveReason};
#[cfg(feature = "pipeline")]
use super::{TakenFiles, TakenSegment};
pub(crate) struct DiskFs {
dir: PathBuf,
stem: String,
claimed: Mutex<HashMap<u32, u64>>,
dropped: AtomicU64,
writer_done: AtomicBool,
}
impl DiskFs {
pub(crate) fn new(dir: impl Into<PathBuf>, stem: impl Into<String>) -> Self {
Self {
dir: dir.into(),
stem: stem.into(),
claimed: Mutex::new(HashMap::new()),
dropped: AtomicU64::new(0),
writer_done: AtomicBool::new(false),
}
}
pub(super) fn create_segment(&self, path: &Path) -> io::Result<ActiveHandle> {
match fs::File::create(path) {
Ok(f) => Ok(ActiveHandle::Disk(f)),
Err(e) if e.kind() == io::ErrorKind::NotFound => {
if let Some(parent) = path.parent()
&& !parent.as_os_str().is_empty()
{
fs::create_dir_all(parent)?;
}
fs::File::create(path).map(ActiveHandle::Disk)
}
Err(e) => Err(e),
}
}
pub(super) fn seal(
&self,
active_handle: ActiveHandle,
active_path: &Path,
index: u32,
) -> io::Result<SegmentRef> {
drop(active_handle);
let sealed_path = strip_active_suffix(active_path);
match fs::rename(active_path, &sealed_path) {
Ok(()) => Ok(SegmentRef::Disk(SealedSegment {
path: sealed_path,
index,
})),
Err(e) => Err(e),
}
}
pub(super) fn remove_sealed(&self, seg: &SegmentRef, reason: RemoveReason) {
if let Some(path) = seg.disk_path() {
remove_segment_family(path);
}
self.claimed.lock().unwrap().remove(&seg.index());
if matches!(reason, RemoveReason::Eviction) {
self.dropped.fetch_add(1, Ordering::Relaxed);
}
}
pub(super) fn remove_active(&self, path: &Path) -> io::Result<()> {
match fs::remove_file(path) {
Ok(()) => Ok(()),
Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(()),
Err(e) => {
rate_limited!(Duration::from_secs(60), {
tracing::warn!(
target: "dial9_worker",
error = %e,
path = %path.display(),
"failed to remove active segment (best-effort)"
);
});
Ok(())
}
}
}
#[cfg(feature = "pipeline")]
pub(super) fn release_claim(&self, index: u32) {
self.claimed.lock().unwrap().remove(&index);
}
#[cfg(feature = "pipeline")]
pub(super) fn writer_done(&self) -> bool {
self.writer_done.load(Ordering::Acquire)
}
pub(super) fn mark_writer_done(&self) {
self.writer_done.store(true, Ordering::Release);
}
#[cfg(feature = "pipeline")]
pub(super) fn take_files(&self) -> TakenFiles {
let on_disk = match find_sealed_segments(&self.dir, &self.stem) {
Ok(s) => s,
Err(e) => {
rate_limited!(Duration::from_secs(60), {
tracing::warn!(
target: "dial9_worker",
error = %e,
"failed to scan for sealed segments"
);
});
return empty_taken_files(self.dropped.swap(0, Ordering::AcqRel));
}
};
let on_disk_indices: HashSet<u32> = on_disk.iter().map(|s| s.index).collect();
let already_claimed: HashSet<u32> = {
let claimed = self.claimed.lock().unwrap();
claimed.keys().copied().collect()
};
let mut new_claims: Vec<(u32, u64)> = Vec::new();
let mut new_segments: Vec<TakenSegment> = Vec::new();
for seg in &on_disk {
if already_claimed.contains(&seg.index) {
continue;
}
let size = match fs::metadata(&seg.path) {
Ok(m) => m.len(),
Err(e) => {
rate_limited!(Duration::from_secs(60), {
tracing::warn!(
target: "dial9_worker",
error = %e,
path = %seg.path.display(),
"failed to stat sealed segment; recording size 0 \
(in_flight_bytes will undercount this segment)"
);
});
0
}
};
new_claims.push((seg.index, size));
new_segments.push(TakenSegment::disk(seg.clone()));
}
let (in_flight_segments, in_flight_bytes) = {
let mut claimed = self.claimed.lock().unwrap();
claimed.retain(|idx, _| on_disk_indices.contains(idx));
for (idx, size) in new_claims {
claimed.insert(idx, size);
}
(claimed.len() as u64, claimed.values().sum::<u64>())
};
TakenFiles {
segments: new_segments,
queued_segments: None,
queued_bytes: None,
in_flight_segments,
in_flight_bytes,
in_flight_bytes_peak: None,
segments_dropped: self.dropped.swap(0, Ordering::AcqRel),
}
}
}
impl DiskFs {
pub(super) fn discover_existing(&self) -> io::Result<DiscoveredArtifacts> {
let mut retained_sizes: BTreeMap<u32, u64> = BTreeMap::new();
if !self.dir.exists() {
return Ok(DiscoveredArtifacts::default());
}
for entry in fs::read_dir(&self.dir)? {
let entry = entry?;
let path = entry.path();
let metadata = match entry.metadata() {
Ok(m) => m,
Err(e) if e.kind() == io::ErrorKind::NotFound => continue,
Err(e) => return Err(e),
};
if !metadata.is_file() {
continue;
}
let Some(file_name) = path.file_name().and_then(|n| n.to_str()) else {
continue;
};
match parse_segment_artifact(file_name, &self.stem) {
Some(SegmentArtifact::Retained { index }) => {
*retained_sizes.entry(index).or_default() += metadata.len();
}
Some(SegmentArtifact::Active) => {
tracing::warn!(
target: "dial9_worker",
path = %path.display(),
"discarding stale active trace segment from a previous writer"
);
match fs::remove_file(&path) {
Ok(()) => {}
Err(e) if e.kind() == io::ErrorKind::NotFound => {}
Err(e) => return Err(e),
}
}
None => {}
}
}
let next_active_index = match retained_sizes.last_key_value() {
Some((&idx, _)) => idx
.checked_add(1)
.ok_or_else(|| io::Error::other("trace segment index overflow"))?,
None => 0,
};
let closed_files = retained_sizes
.into_iter()
.map(|(index, size)| {
let path = self.dir.join(format!("{}.{}.bin", self.stem, index));
(SegmentRef::Disk(SealedSegment { path, index }), size)
})
.collect();
Ok(DiscoveredArtifacts {
closed_files,
next_active_index,
})
}
}
fn remove_segment_family(path: &Path) {
let Some(file_name) = path.file_name().and_then(|n| n.to_str()) else {
return;
};
let Some(parent) = path.parent() else {
return;
};
let entries = match fs::read_dir(parent) {
Ok(e) => e,
Err(e) if e.kind() == io::ErrorKind::NotFound => return,
Err(e) => {
rate_limited!(Duration::from_secs(60), {
tracing::warn!(
target: "dial9_worker",
error = %e,
parent = %parent.display(),
"failed to scan parent for trace family eviction"
);
});
return;
}
};
for entry in entries.flatten() {
let name = entry.file_name();
let Some(name_str) = name.to_str() else {
continue;
};
let is_family = name_str == file_name
|| name_str
.strip_prefix(file_name)
.is_some_and(|s| s.starts_with('.'));
if !is_family {
continue;
}
match fs::remove_file(&entry.path()) {
Ok(()) => {}
Err(e) if e.kind() == io::ErrorKind::NotFound => {}
Err(e) => {
rate_limited!(Duration::from_secs(60), {
tracing::warn!(
target: "dial9_worker",
error = %e,
path = %entry.path().display(),
"failed to remove trace artifact"
);
});
}
}
}
}
fn strip_active_suffix(path: &Path) -> PathBuf {
let s = path.to_str().unwrap_or_default();
if let Some(without) = s.strip_suffix(".active") {
PathBuf::from(without)
} else {
path.to_path_buf()
}
}
#[cfg(feature = "pipeline")]
fn empty_taken_files(segments_dropped: u64) -> TakenFiles {
TakenFiles {
segments: vec![],
queued_segments: None,
queued_bytes: None,
in_flight_segments: 0,
in_flight_bytes: 0,
in_flight_bytes_peak: None,
segments_dropped,
}
}
#[cfg(all(test, feature = "pipeline"))]
mod tests {
use super::*;
use crate::fs::Fs;
use assert2::check;
#[test]
fn disk_fs_claim_dedup() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("trace.0.bin"), b"seg0").unwrap();
std::fs::write(dir.path().join("trace.1.bin"), b"seg1").unwrap();
let fs = Fs::Disk(DiskFs::new(dir.path(), "trace"));
let t1 = fs.take_files();
check!(t1.segments.len() == 2);
let t2 = fs.take_files();
check!(t2.segments.is_empty());
}
#[test]
fn disk_fs_scan_prunes_claim_when_file_deleted() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("trace.0.bin");
std::fs::write(&path, b"seg0").unwrap();
let fs = Fs::Disk(DiskFs::new(dir.path(), "trace"));
let t1 = fs.take_files();
check!(t1.segments.len() == 1);
check!(t1.in_flight_segments == 1);
std::fs::remove_file(&path).unwrap();
let t2 = fs.take_files();
check!(
t2.segments.is_empty(),
"vanished file must not be re-dispatched"
);
check!(t2.in_flight_segments == 0, "stale claim must be pruned");
check!(t2.in_flight_bytes == 0);
}
#[test]
fn disk_fs_release_claim_redispatches() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("trace.0.bin"), b"seg0").unwrap();
let disk = DiskFs::new(dir.path(), "trace");
let t1 = disk.take_files();
check!(t1.segments.len() == 1);
let seg = &t1.segments[0].seg_ref;
disk.release_claim(seg.index());
let t2 = disk.take_files();
check!(
t2.segments.len() == 1,
"released claim should be re-dispensed"
);
}
#[test]
fn disk_fs_eviction_bumps_dropped() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("trace.0.bin"), b"data").unwrap();
let fs = Fs::Disk(DiskFs::new(dir.path(), "trace"));
let t = fs.take_files();
check!(t.segments.len() == 1);
let seg = t.segments.into_iter().next().unwrap().seg_ref;
check!(t.segments_dropped == 0);
fs.remove_sealed(&seg, RemoveReason::Eviction);
let t2 = fs.take_files();
check!(t2.segments_dropped == 1);
let t3 = fs.take_files();
check!(t3.segments_dropped == 0);
}
#[test]
fn disk_fs_terminal_does_not_bump_dropped() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("trace.0.bin"), b"data").unwrap();
let fs = Fs::Disk(DiskFs::new(dir.path(), "trace"));
let t = fs.take_files();
let seg = t.segments.into_iter().next().unwrap().seg_ref;
fs.remove_sealed(&seg, RemoveReason::Terminal);
let t2 = fs.take_files();
check!(t2.segments_dropped == 0);
}
#[test]
fn discover_existing_empty_dir() {
let dir = tempfile::tempdir().unwrap();
let disk = DiskFs::new(dir.path(), "trace");
let d = disk.discover_existing().unwrap();
check!(d.next_active_index == 0);
check!(d.closed_files.is_empty());
}
#[test]
fn discover_existing_sums_artifact_family_per_index() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("trace.0.bin"), vec![0u8; 100]).unwrap();
std::fs::write(dir.path().join("trace.0.bin.gz"), vec![0u8; 30]).unwrap();
std::fs::write(dir.path().join("trace.2.bin"), vec![0u8; 50]).unwrap();
let disk = DiskFs::new(dir.path(), "trace");
let d = disk.discover_existing().unwrap();
check!(d.next_active_index == 3, "max(0,2)+1 = 3");
let by_index: std::collections::HashMap<u32, u64> = d
.closed_files
.iter()
.map(|(seg, size)| (seg.index(), *size))
.collect();
check!(by_index.get(&0) == Some(&130), ".bin + .bin.gz summed");
check!(by_index.get(&2) == Some(&50));
}
#[test]
fn discover_existing_discards_stale_active() {
let dir = tempfile::tempdir().unwrap();
let stale = dir.path().join("trace.7.bin.active");
std::fs::write(&stale, b"orphan").unwrap();
let disk = DiskFs::new(dir.path(), "trace");
let _ = disk.discover_existing().unwrap();
check!(!stale.exists(), "stale .active must be discarded");
}
#[test]
fn discover_existing_ignores_unrelated_files() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("other.0.bin"), b"x").unwrap();
std::fs::write(dir.path().join("README"), b"x").unwrap();
std::fs::write(dir.path().join("trace.0.bin"), b"x").unwrap();
let disk = DiskFs::new(dir.path(), "trace");
let d = disk.discover_existing().unwrap();
check!(d.closed_files.len() == 1);
check!(d.next_active_index == 1);
}
#[test]
fn remove_segment_family_removes_bin_and_gz_siblings() {
let dir = tempfile::tempdir().unwrap();
let bin = dir.path().join("trace.3.bin");
let gz = dir.path().join("trace.3.bin.gz");
let unrelated = dir.path().join("trace.4.bin");
std::fs::write(&bin, b"x").unwrap();
std::fs::write(&gz, b"x").unwrap();
std::fs::write(&unrelated, b"x").unwrap();
remove_segment_family(&bin);
check!(!bin.exists());
check!(!gz.exists());
check!(unrelated.exists(), "sibling with different index untouched");
}
#[test]
fn strip_active_suffix_removes_suffix() {
let p = Path::new("/tmp/trace.0.bin.active");
check!(strip_active_suffix(p) == PathBuf::from("/tmp/trace.0.bin"));
}
#[test]
fn strip_active_suffix_no_suffix() {
let p = Path::new("/tmp/trace.0.bin");
check!(strip_active_suffix(p) == PathBuf::from("/tmp/trace.0.bin"));
}
}
#[cfg(all(test, shuttle))]
mod shuttle_tests {
use super::*;
use crate::primitives::sync::Arc;
use crate::primitives::sync::atomic::AtomicUsize;
const COUNT: u32 = 3;
const SCANS: usize = 3;
fn seal_one(disk: &DiskFs, dir: &Path, stem: &str, index: u32) {
let active_path = dir.join(format!("{stem}.{index}.bin.active"));
let handle = disk.create_segment(&active_path).unwrap();
disk.seal(handle, &active_path, index).unwrap();
}
crate::shuttle_test! {
num_iters = 1_000, depth = 3, should_panic,
expect_panic = "every sealed segment must be claimed exactly once",
replay = "91022bc1cfb7e1e792c7bc6e802449922481242992a424499224490000";
fn shuttle_claim_dedup() {
let dir = tempfile::tempdir().unwrap();
let stem = "trace";
let disk = Arc::new(DiskFs::new(dir.path(), stem));
for i in 0..COUNT {
seal_one(&disk, dir.path(), stem, i);
}
let (tx, rx) = crate::primitives::sync::mpsc::sync_channel::<SegmentRef>(COUNT as usize);
let dispatched = Arc::new(AtomicUsize::new(0));
let claimer = {
let disk = disk.clone();
let dispatched = dispatched.clone();
crate::primitives::thread::spawn(move || {
for _ in 0..SCANS {
let taken = disk.take_files();
for seg in taken.segments {
dispatched.fetch_add(1, Ordering::Relaxed);
let _ = tx.send(seg.seg_ref);
}
shuttle::thread::yield_now();
}
})
};
let remover = crate::primitives::thread::spawn(move || {
let mut removed = 0usize;
while removed < COUNT as usize {
let seg_ref = rx.recv().unwrap();
disk.remove_sealed(&seg_ref, RemoveReason::Terminal);
removed += 1;
}
removed
});
claimer.join().unwrap();
let removed_total = remover.join().unwrap();
assert_eq!(
dispatched.load(Ordering::Relaxed),
COUNT as usize,
"every sealed segment must be claimed exactly once"
);
assert_eq!(removed_total, COUNT as usize);
}
}
}