use dial9_trace_format::encoder::{Encoder, RawEncoder};
use crate::clock::clock_pair;
use crate::collector::Batch;
use crate::format::{ClockSyncEvent, SegmentMetadataEvent};
use crate::fs::{ActiveHandle, Fs, RemoveReason};
use crate::primitives::fs;
use crate::rate_limit::rate_limited;
use crate::sealed::SegmentRef;
use std::collections::VecDeque;
use std::io::BufWriter;
use std::marker::PhantomData;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::{Duration, Instant};
use metrique_timesource::time_source;
mod mode_sealed {
pub trait Sealed {}
}
pub trait BufferMode: mode_sealed::Sealed + Send + 'static {
const IS_DISK: bool;
}
#[derive(Debug)]
#[non_exhaustive]
pub struct Disk;
#[derive(Debug)]
#[non_exhaustive]
pub struct Memory;
impl mode_sealed::Sealed for Disk {}
impl mode_sealed::Sealed for Memory {}
impl BufferMode for Disk {
const IS_DISK: bool = true;
}
impl BufferMode for Memory {
const IS_DISK: bool = false;
}
pub type DiskBuffer = SegmentWriter<Disk>;
pub type MemoryBuffer = SegmentWriter<Memory>;
const DIAL9_VERSION_KEY: &str = "dial9.dial9-tokio-telemetry.version";
const DIAL9_VERSION_VALUE: &str = env!("CARGO_PKG_VERSION");
const PROCESS_AVAILABLE_PARALLELISM_KEY: &str = "process.available_parallelism";
#[derive(Clone)]
struct SegmentMetadata {
entries: Vec<(String, String)>,
}
impl Default for SegmentMetadata {
fn default() -> Self {
let mut entries = vec![(
DIAL9_VERSION_KEY.to_string(),
DIAL9_VERSION_VALUE.to_string(),
)];
match std::thread::available_parallelism() {
Ok(parallelism) => entries.push((
PROCESS_AVAILABLE_PARALLELISM_KEY.to_string(),
parallelism.get().to_string(),
)),
Err(e) => rate_limited!(Duration::from_secs(60), {
tracing::warn!("failed to read process available parallelism: {e}");
}),
}
Self { entries }
}
}
impl SegmentMetadata {
fn new(user_entries: Vec<(String, String)>) -> Self {
let mut s = Self::default();
s.merge(user_entries.into_iter());
s
}
fn merge(&mut self, entries: impl Iterator<Item = (String, String)>) -> bool {
let mut merged: Vec<(String, String)> = entries.collect();
for (k, v) in &self.entries {
if !merged.iter().any(|(mk, _)| mk == k) {
merged.push((k.clone(), v.clone()));
}
}
if merged == self.entries {
return false;
}
self.entries = merged;
true
}
}
const DEFAULT_ROTATION_PERIOD: Duration = Duration::from_secs(60);
const DEFAULT_DRAIN_INTERVAL: Duration = Duration::from_secs(30);
const SEGMENT_STEM: &str = "trace";
const BYTES_PER_MIB: u64 = 1024 * 1024;
const MAX_FILE_SIZE_CAP: u64 = 100 * BYTES_PER_MIB;
fn derive_max_file_size(max_total_size: u64) -> u64 {
(max_total_size / 4).min(MAX_FILE_SIZE_CAP)
}
pub struct SegmentWriter<Mode: BufferMode = Disk> {
dir: PathBuf,
stem: String,
max_file_size: u64,
max_total_size: u64,
rotation_period: Duration,
next_rotation_time: Option<Instant>,
closed_files: VecDeque<(SegmentRef, u64)>,
active_path: PathBuf,
state: WriterState,
next_index: u32,
segment_metadata: SegmentMetadata,
dropped_events: usize,
has_real_events: bool,
drain_interval: Duration,
next_drain_time: Instant,
fs: Arc<Fs>,
boot_id: Option<String>,
_namespace_lock: Option<std::fs::File>,
_mode: PhantomData<Mode>,
}
impl<M: BufferMode> std::fmt::Debug for SegmentWriter<M> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SegmentWriter")
.field("dir", &self.dir)
.field("stem", &self.stem)
.field("max_file_size", &self.max_file_size)
.field("max_total_size", &self.max_total_size)
.finish_non_exhaustive()
}
}
#[allow(clippy::large_enum_variant)]
enum WriterState {
Active {
writer: RawEncoder<BufWriter<ActiveHandle>>,
need_metadata: bool,
},
Finished,
}
#[bon::bon]
impl SegmentWriter<Disk> {
#[builder(builder_type = DiskBufferBuilder, finish_fn = build)]
pub fn builder(
base_path: impl Into<PathBuf>,
max_file_size: Option<u64>,
max_total_size: u64,
rotation_period: Option<Duration>,
segment_metadata: Option<Vec<(String, String)>>,
) -> std::io::Result<Self> {
Self::create(
base_path,
max_file_size.unwrap_or_else(|| derive_max_file_size(max_total_size)),
max_total_size,
rotation_period.unwrap_or(DEFAULT_ROTATION_PERIOD),
segment_metadata
.map(SegmentMetadata::new)
.unwrap_or_default(),
)
}
fn create(
base_path: impl Into<PathBuf>,
max_file_size: u64,
max_total_size: u64,
rotation_period: Duration,
segment_metadata: SegmentMetadata,
) -> std::io::Result<Self> {
if rotation_period == Duration::from_secs(0) {
return Err(std::io::Error::other("Rotation period must not be zero"));
}
let dir = base_path.into();
if !dir.as_os_str().is_empty() {
fs::create_dir_all(&dir)?;
}
let stem = SEGMENT_STEM.to_string();
let fs = Fs::new_disk(&dir, stem.as_str());
let discovered = fs.discover_existing()?;
let first_index = discovered.next_active_index;
let next_index = first_index
.checked_add(1)
.ok_or_else(|| std::io::Error::other("trace segment index overflow"))?;
let first_path = Self::active_path(&dir, &stem, first_index);
let handle = fs.create_segment(&first_path)?;
let state = Self::prepare_segment(BufWriter::new(handle))?;
let now = time_source().instant().as_std();
let drain_interval = rotation_period.min(DEFAULT_DRAIN_INTERVAL);
let mut writer = Self {
dir,
stem,
max_file_size,
max_total_size,
rotation_period,
next_rotation_time: Self::next_rotation_from(now, rotation_period),
closed_files: discovered.closed_files,
active_path: first_path,
state,
next_index,
segment_metadata,
dropped_events: 0,
has_real_events: false,
drain_interval,
next_drain_time: now + drain_interval,
fs,
boot_id: None,
_namespace_lock: None,
_mode: PhantomData,
};
writer.evict_oldest()?;
Ok(writer)
}
pub fn set_namespace(&mut self, boot_id: String, lock: std::fs::File) {
self.boot_id = Some(boot_id);
self._namespace_lock = Some(lock);
}
pub fn single_file(path: impl Into<PathBuf>) -> std::io::Result<Self> {
let path = path.into();
let dir = path
.parent()
.filter(|p| !p.as_os_str().is_empty())
.unwrap_or(Path::new("."))
.to_path_buf();
let stem = path
.file_stem()
.and_then(|s| s.to_str())
.unwrap_or(SEGMENT_STEM)
.to_string();
let fs = Fs::new_disk(&dir, stem.as_str());
let active_path = Self::active_path(&dir, &stem, 0);
let handle = fs.create_segment(&active_path)?;
let state = Self::prepare_segment(BufWriter::new(handle))?;
let now = time_source().instant().as_std();
Ok(Self {
dir,
stem,
max_file_size: u64::MAX,
max_total_size: u64::MAX,
rotation_period: Duration::MAX,
next_rotation_time: None,
closed_files: VecDeque::new(),
active_path,
state,
next_index: 1,
segment_metadata: SegmentMetadata::default(),
dropped_events: 0,
has_real_events: false,
drain_interval: DEFAULT_DRAIN_INTERVAL,
next_drain_time: now + DEFAULT_DRAIN_INTERVAL,
fs,
boot_id: None,
_namespace_lock: None,
_mode: PhantomData,
})
}
}
fn pick_segment_size(max_total_size: u64) -> u64 {
const MIN_SLOTS: u64 = 8;
(max_total_size / MIN_SLOTS).max(1)
}
#[bon::bon]
impl SegmentWriter<Memory> {
pub fn new(max_total_size: u64) -> std::io::Result<Self> {
Self::create_in_memory(
max_total_size,
pick_segment_size(max_total_size),
DEFAULT_ROTATION_PERIOD,
SegmentMetadata::default(),
)
}
#[builder(builder_type = MemoryBufferBuilder, finish_fn = build)]
pub fn builder(
max_total_size: u64,
max_segment_size: Option<u64>,
rotation_period: Option<Duration>,
segment_metadata: Option<Vec<(String, String)>>,
) -> std::io::Result<Self> {
let seg_size = max_segment_size.unwrap_or_else(|| pick_segment_size(max_total_size));
Self::create_in_memory(
max_total_size,
seg_size,
rotation_period.unwrap_or(DEFAULT_ROTATION_PERIOD),
segment_metadata
.map(SegmentMetadata::new)
.unwrap_or_default(),
)
}
fn create_in_memory(
max_total_size: u64,
max_segment_size: u64,
rotation_period: Duration,
segment_metadata: SegmentMetadata,
) -> std::io::Result<Self> {
if max_total_size == 0 {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"max_total_size must be > 0",
));
}
if max_segment_size == 0 {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"max_segment_size must be > 0",
));
}
if rotation_period == Duration::from_secs(0) {
return Err(std::io::Error::other("Rotation period must not be zero"));
}
let min_total = (crate::fs::PIPELINE_RESERVE_SEGMENTS + 1).saturating_mul(max_segment_size);
if max_total_size < min_total {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!(
"max_total_size ({max_total_size}) must be >= {min_total} \
({} × max_segment_size: 1 active + 1 in-flight + 1 ring slot)",
crate::fs::PIPELINE_RESERVE_SEGMENTS + 1
),
));
}
let fs = Fs::new_in_memory(max_total_size, max_segment_size)?;
let dir = PathBuf::from("mem");
let stem = SEGMENT_STEM.to_string();
let active_path = Self::active_path(&dir, &stem, 0);
let handle = fs.create_segment(&active_path)?;
let state = Self::prepare_segment(BufWriter::new(handle))?;
let now = time_source().instant().as_std();
let drain_interval = rotation_period.min(DEFAULT_DRAIN_INTERVAL);
Ok(Self {
dir,
stem,
max_file_size: max_segment_size,
max_total_size,
rotation_period,
next_rotation_time: Self::next_rotation_from(now, rotation_period),
closed_files: VecDeque::new(),
active_path,
state,
next_index: 1,
segment_metadata,
dropped_events: 0,
has_real_events: false,
drain_interval,
next_drain_time: now + drain_interval,
fs,
boot_id: None,
_namespace_lock: None,
_mode: PhantomData,
})
}
}
impl<M: BufferMode> SegmentWriter<M> {
pub fn boot_id(&self) -> Option<&str> {
self.boot_id.as_deref()
}
pub fn trace_dir(&self) -> &Path {
&self.dir
}
pub fn trace_stem(&self) -> &str {
&self.stem
}
pub fn current_active_path(&self) -> &Path {
&self.active_path
}
fn prepare_segment(writer: BufWriter<ActiveHandle>) -> std::io::Result<WriterState> {
let mut encoder = Encoder::new_to(writer)?;
let (mono, real) = clock_pair();
encoder.write(&ClockSyncEvent {
timestamp_ns: mono,
realtime_ns: real,
})?;
Ok(WriterState::Active {
writer: encoder.into_raw_encoder(),
need_metadata: true,
})
}
fn write_metadata_if_needed(&mut self) -> std::io::Result<()> {
match &mut self.state {
WriterState::Active {
writer,
need_metadata,
} => {
if *need_metadata {
Self::write_segment_metadata(writer, &self.segment_metadata.entries)?;
}
*need_metadata = false;
Ok(())
}
WriterState::Finished => Ok(()),
}
}
fn write_segment_metadata(
writer: &mut RawEncoder<BufWriter<ActiveHandle>>,
entries: &[(String, String)],
) -> std::io::Result<()> {
let mut enc = Encoder::new();
let entries = entries.to_vec();
let (mono, real) = clock_pair();
enc.write(&SegmentMetadataEvent {
timestamp_ns: mono,
entries,
})?;
enc.write(&ClockSyncEvent {
timestamp_ns: mono,
realtime_ns: real,
})?;
writer.write_raw(&enc.finish())?;
Ok(())
}
fn active_path(dir: &Path, stem: &str, index: u32) -> PathBuf {
dir.join(format!("{stem}.{index}.bin.active"))
}
fn next_rotation_from(now: Instant, period: Duration) -> Option<Instant> {
(period != Duration::MAX).then(|| now + period)
}
fn rotate(&mut self) -> std::io::Result<()> {
if matches!(self.state, WriterState::Finished) {
return Ok(());
}
let now = time_source().instant().as_std();
self.next_rotation_time = Self::next_rotation_from(now, self.rotation_period);
self.next_drain_time = now + self.drain_interval;
let WriterState::Active {
writer: mut raw, ..
} = std::mem::replace(&mut self.state, WriterState::Finished)
else {
return Ok(());
};
let _ = raw.flush();
let closed_size = raw.bytes_written();
let current_index = self.next_index - 1;
let bw: BufWriter<ActiveHandle> = raw.into_inner();
let handle: ActiveHandle = bw
.into_inner()
.unwrap_or_else(|e| e.into_inner().into_parts().0);
match self.fs.seal(handle, &self.active_path, current_index) {
Ok(seg_ref) => {
if M::IS_DISK {
self.closed_files.push_back((seg_ref, closed_size));
}
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
rate_limited!(Duration::from_secs(60), {
tracing::warn!(
"active trace file {} disappeared before sealing; \
abandoning segment and starting a fresh one",
self.active_path.display()
);
});
}
Err(e) => {
return Err(e);
}
}
let new_path = Self::active_path(&self.dir, &self.stem, self.next_index);
self.next_index += 1;
let handle: ActiveHandle = self.fs.create_segment(&new_path)?;
self.state = match Self::prepare_segment(BufWriter::new(handle)) {
Ok(s) => s,
Err(e) => {
let _ = self.fs.remove_active(&new_path);
return Err(e);
}
};
self.active_path = new_path;
self.has_real_events = false;
tracing::debug!(
segment_index = self.next_index - 1,
"rotated to new trace segment"
);
self.evict_oldest()?;
Ok(())
}
fn total_size(&self) -> u64 {
if !M::IS_DISK {
return 0;
}
let closed: u64 = self.closed_files.iter().map(|(_, s)| s).sum();
let active = match &self.state {
WriterState::Active { writer, .. } => writer.bytes_written(),
WriterState::Finished => 0,
};
closed + active
}
fn evict_oldest(&mut self) -> std::io::Result<()> {
if !M::IS_DISK {
return Ok(());
}
while self.total_size() > self.max_total_size && !self.closed_files.is_empty() {
if let Some((seg_ref, _size)) = self.closed_files.pop_front() {
self.fs.remove_sealed(&seg_ref, RemoveReason::Eviction);
}
}
if self.total_size() > self.max_total_size {
self.state = WriterState::Finished;
}
Ok(())
}
fn maybe_rotate(&mut self) -> std::io::Result<()> {
let WriterState::Active { writer: raw, .. } = &self.state else {
return Ok(());
};
if raw.bytes_written() > self.max_file_size {
self.rotate()?;
}
Ok(())
}
}
impl<M: BufferMode> SegmentWriter<M> {
#[cfg(feature = "pipeline")]
pub(crate) fn fs_handle(&self) -> Option<Arc<Fs>> {
Some(Arc::clone(&self.fs))
}
pub fn flush(&mut self) -> std::io::Result<()> {
if let WriterState::Active { writer: raw, .. } = &mut self.state {
raw.flush()?;
}
Ok(())
}
#[cfg(test)]
pub(crate) fn segment_metadata(&self) -> &[(String, String)] {
&self.segment_metadata.entries
}
pub fn update_segment_metadata(&mut self, entries: impl IntoIterator<Item = (String, String)>) {
if self.segment_metadata.merge(entries.into_iter()) {
match &mut self.state {
WriterState::Active { need_metadata, .. } => *need_metadata = true,
WriterState::Finished => {}
}
}
}
pub(crate) fn write_current_segment_metadata(&mut self) -> std::io::Result<()> {
self.write_metadata_if_needed()
}
pub(crate) fn should_drain(&self) -> bool {
self.has_real_events && time_source().instant().as_std() >= self.next_drain_time
}
pub(crate) fn drained(&mut self) -> std::io::Result<bool> {
if !self.has_real_events {
return Ok(false);
}
let now = time_source().instant().as_std();
if self
.next_rotation_time
.is_some_and(|deadline| now >= deadline)
{
self.rotate()?;
return Ok(true);
}
self.next_drain_time = now + self.drain_interval;
Ok(false)
}
pub fn finalize(&mut self) -> std::io::Result<()> {
if matches!(self.state, WriterState::Finished) {
rate_limited!(Duration::from_secs(60), {
tracing::warn!("writer is already closed.");
});
self.fs.mark_writer_done();
return Ok(());
}
let _ = self.flush();
let WriterState::Active { writer: raw, .. } =
std::mem::replace(&mut self.state, WriterState::Finished)
else {
self.fs.mark_writer_done();
return Ok(());
};
let bytes_written = raw.bytes_written();
let bw: BufWriter<ActiveHandle> = raw.into_inner();
let handle: ActiveHandle = bw
.into_inner()
.unwrap_or_else(|e| e.into_inner().into_parts().0);
let current_index = self.next_index - 1;
if self.has_real_events {
match self.fs.seal(handle, &self.active_path, current_index) {
Ok(seg_ref) => {
if M::IS_DISK {
self.closed_files.push_back((seg_ref, bytes_written));
}
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
rate_limited!(Duration::from_secs(60), {
tracing::warn!(
"active trace file {} disappeared before finalize; \
dropping segment",
self.active_path.display()
);
});
}
Err(e) => {
self.fs.mark_writer_done();
return Err(e);
}
}
} else {
tracing::debug!(
"removing empty final segment {}",
self.active_path.display()
);
if let Err(e) = self.fs.remove_active(&self.active_path)
&& e.kind() != std::io::ErrorKind::NotFound
{
self.fs.mark_writer_done();
return Err(e);
}
}
if let Err(e) = self.evict_oldest() {
self.fs.mark_writer_done();
return Err(e);
}
self.fs.mark_writer_done();
Ok(())
}
crate::test_util_pub! {
fn write_encoded_batch(&mut self, batch: &Batch) -> std::io::Result<()> {
self.write_metadata_if_needed()?;
let WriterState::Active { writer: raw, .. } = &mut self.state else {
self.dropped_events += batch.event_count() as usize;
return Ok(());
};
if batch.event_count() > 0 {
raw.write_raw(batch.encoded_bytes())?;
self.has_real_events = true;
self.maybe_rotate()?;
}
Ok(())
}
}
}
impl<M: BufferMode> Drop for SegmentWriter<M> {
fn drop(&mut self) {
if self.dropped_events > 0 {
rate_limited!(Duration::from_secs(60), {
tracing::info!(
target: "dial9_telemetry",
dropped_events = self.dropped_events,
"SegmentWriter dropped events after finalization"
);
});
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use dial9_trace_format::TraceEvent;
use std::collections::HashMap;
use std::io::Read;
use tempfile::TempDir;
#[derive(TraceEvent)]
#[traceevent(wire_slot)]
struct TestEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
value: u64,
}
#[derive(Debug)]
enum Decoded {
ClockSync {
timestamp_ns: u64,
realtime_ns: u64,
},
SegmentMetadata {
timestamp_ns: u64,
entries: HashMap<String, String>,
},
Data {
timestamp_ns: u64,
},
}
fn decode_all(data: &[u8]) -> Vec<Decoded> {
use dial9_trace_format::decoder::{DecodedFrameRef, Decoder};
use dial9_trace_format::types::FieldValueRef;
let mut dec = Decoder::new(data).expect("valid trace header");
let mut out = Vec::new();
while let Some(frame) = dec.next_frame_ref().expect("decode frame") {
let DecodedFrameRef::Event {
type_id,
timestamp_ns,
values,
} = frame
else {
continue;
};
let ts = timestamp_ns;
let name = dec.registry().get(type_id).map(|s| s.name());
match name {
Some("ClockSyncEvent") => {
let realtime_ns = match values.first() {
Some(FieldValueRef::Varint(v)) => *v,
other => panic!("ClockSyncEvent realtime_ns: {other:?}"),
};
out.push(Decoded::ClockSync {
timestamp_ns: ts,
realtime_ns,
});
}
Some("SegmentMetadataEvent") => {
let entries = match values.first() {
Some(FieldValueRef::StringMap(m)) => m
.iter()
.map(|(k, v)| (k.to_string(), v.to_string()))
.collect(),
other => panic!("SegmentMetadataEvent entries: {other:?}"),
};
out.push(Decoded::SegmentMetadata {
timestamp_ns: ts,
entries,
});
}
_ => out.push(Decoded::Data { timestamp_ns: ts }),
}
}
out
}
fn test_batch() -> Batch {
let mut enc = Encoder::new_to(Vec::new()).unwrap();
enc.write(&TestEvent {
timestamp_ns: 1000,
value: 0,
})
.unwrap();
Batch::new(enc.into_inner(), 1)
}
fn rotating_file(base: &std::path::Path, i: u32) -> String {
format!("{}.{}.bin", base.display(), i)
}
fn read_trace_events(path: &str) -> Vec<Decoded> {
let data = std::fs::read(path).unwrap();
decode_all(&data)
.into_iter()
.filter(|e| matches!(e, Decoded::Data { .. }))
.collect()
}
fn total_disk_usage(dir: &std::path::Path) -> u64 {
std::fs::read_dir(dir)
.unwrap()
.filter_map(|e| e.ok())
.filter(|e| {
let p = e.path();
p.extension()
.is_some_and(|ext| ext == "bin" || ext == "active")
})
.map(|e| e.metadata().unwrap().len())
.sum()
}
fn single_event_file_size() -> u64 {
let dir = TempDir::new().unwrap();
let path = dir.path().join("probe.bin");
let mut w = DiskBuffer::single_file(&path).unwrap();
w.write_encoded_batch(&test_batch()).unwrap();
w.flush().unwrap();
std::fs::metadata(w.current_active_path()).unwrap().len()
}
#[test]
fn test_writer_creation() {
let dir = TempDir::new().unwrap();
let path = dir.path().join("test_trace_v2.bin");
let writer = DiskBuffer::single_file(&path);
assert!(writer.is_ok());
}
#[test]
fn derive_max_file_size_caps_large_budgets_at_100_mib() {
assert_eq!(
derive_max_file_size(1024 * BYTES_PER_MIB),
100 * BYTES_PER_MIB
);
}
#[test]
fn derive_max_file_size_uses_quarter_of_small_budgets() {
assert_eq!(derive_max_file_size(64 * BYTES_PER_MIB), 16 * BYTES_PER_MIB);
}
#[test]
fn builder_defaults_max_file_size_from_total_size() {
let dir = TempDir::new().unwrap();
let total = 64 * BYTES_PER_MIB;
let writer = DiskBuffer::builder()
.base_path(dir.path())
.max_total_size(total)
.build()
.expect("builder should succeed without max_file_size");
assert_eq!(writer.max_file_size, derive_max_file_size(total));
}
#[test]
fn builder_honors_explicit_max_file_size() {
let dir = TempDir::new().unwrap();
let writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(7 * BYTES_PER_MIB)
.max_total_size(64 * BYTES_PER_MIB)
.build()
.expect("builder should succeed");
assert_eq!(writer.max_file_size, 7 * BYTES_PER_MIB);
}
#[test]
fn test_write_event() {
let dir = TempDir::new().unwrap();
let path = dir.path().join("test_event_v2.bin");
let mut writer = DiskBuffer::single_file(&path).unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
let metadata = std::fs::metadata(writer.current_active_path()).unwrap();
assert!(
metadata.len() > 0,
"file should not be empty after writing an event"
);
}
#[test]
fn test_write_batch_sizes() {
let dir = TempDir::new().unwrap();
let path = dir.path().join("test_batch_v2.bin");
let mut writer = DiskBuffer::single_file(&path).unwrap();
let one_event_size = single_event_file_size();
for _ in 0..2 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.flush().unwrap();
let metadata = std::fs::metadata(writer.current_active_path()).unwrap();
assert!(metadata.len() > one_event_size);
}
#[test]
fn test_binary_format_header() {
let dir = TempDir::new().unwrap();
let path = dir.path().join("test_format_v2.bin");
let writer = DiskBuffer::single_file(&path).unwrap();
let active = writer.current_active_path().to_owned();
drop(writer);
let mut file = std::fs::File::open(&active).unwrap();
let mut magic = [0u8; 4];
file.read_exact(&mut magic).unwrap();
assert_eq!(&magic, b"TRC\0");
}
#[test]
fn test_rotating_writer_creation() {
let dir = TempDir::new().unwrap();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(1024)
.max_total_size(4096)
.build()
.unwrap();
writer.finalize().unwrap();
assert!(
!dir.path().join("trace.0.bin").exists(),
"empty segment should not be sealed"
);
assert!(
!dir.path().join("trace.0.bin.active").exists(),
"active file should be removed"
);
}
#[test]
fn test_rotating_writer_rotation() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.build()
.unwrap();
for _ in 0..3 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.finalize().unwrap();
let total: usize = (0..10)
.map(|i| {
let f = rotating_file(&base, i);
if std::path::Path::new(&f).exists() {
read_trace_events(&f).len()
} else {
0
}
})
.sum();
assert_eq!(total, 3);
}
#[test]
fn test_rotating_writer_eviction() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let one_event = single_event_file_size();
let max_file_size = one_event;
let max_total_size = max_file_size * 3;
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(max_file_size)
.max_total_size(max_total_size)
.build()
.unwrap();
for _ in 0..10 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.finalize().unwrap();
assert!(total_disk_usage(dir.path()) <= max_total_size);
assert!(!std::path::Path::new(&rotating_file(&base, 0)).exists());
}
#[test]
fn test_rotating_writer_stops_when_over_budget() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let one_event = single_event_file_size();
let max_file_size = one_event;
let max_total_size = one_event + 5;
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(max_file_size)
.max_total_size(max_total_size)
.build()
.unwrap();
for _ in 0..100 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.finalize().unwrap();
let total: usize = (0..100)
.map(|i| {
let f = rotating_file(&base, i);
if std::path::Path::new(&f).exists() {
read_trace_events(&f).len()
} else {
0
}
})
.sum();
assert!(
total < 100,
"should have stopped writing, got {total} events"
);
}
#[test]
fn test_writer_stops_on_tiny_overshoot_after_eviction() {
let dir = TempDir::new().unwrap();
let max_file_size = 200;
let num_files = 100u64;
let max_total_size = max_file_size * num_files;
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(max_file_size)
.max_total_size(max_total_size)
.build()
.unwrap();
for i in 0..5000 {
writer.write_encoded_batch(&test_batch()).unwrap();
if matches!(writer.state, WriterState::Finished) {
panic!(
"Writer stopped at batch {i}! total_size={}, max_total_size={}, \
closed_files={}. \
write_encoded_batch should try eviction before stopping.",
writer.total_size(),
max_total_size,
writer.closed_files.len()
);
}
}
}
#[test]
fn test_rotating_writer_file_naming() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.build()
.unwrap();
for _ in 0..5 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.finalize().unwrap();
assert!(
std::path::Path::new(&rotating_file(&base, 0)).exists(),
"File 0 should exist"
);
let total: usize = (0..10)
.map(|i| {
let f = rotating_file(&base, i);
if std::path::Path::new(&f).exists() {
read_trace_events(&f).len()
} else {
0
}
})
.sum();
assert_eq!(total, 5);
}
#[test]
fn test_write_batch_across_rotation_boundary() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.build()
.unwrap();
for _ in 0..3 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.finalize().unwrap();
let total: usize = (0..10)
.map(|i| {
let f = rotating_file(&base, i);
if std::path::Path::new(&f).exists() {
read_trace_events(&f).len()
} else {
0
}
})
.sum();
assert_eq!(total, 3);
}
#[test]
fn test_rotated_files_have_valid_headers() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.build()
.unwrap();
for _ in 0..3 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.finalize().unwrap();
let total: usize = (0..10)
.map(|i| {
let f = rotating_file(&base, i);
if std::path::Path::new(&f).exists() {
read_trace_events(&f).len() } else {
0
}
})
.sum();
assert_eq!(total, 3);
}
#[test]
fn test_flush_after_stop() {
let dir = TempDir::new().unwrap();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(10_000)
.max_total_size(50)
.build()
.unwrap();
for _ in 0..5 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
assert!(writer.flush().is_ok());
assert!(writer.flush().is_ok());
}
#[test]
fn test_mixed_event_sizes() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.build()
.unwrap();
for _ in 0..3 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.finalize().unwrap();
let mut total = 0;
for i in 0..10 {
let f = rotating_file(&base, i);
if std::path::Path::new(&f).exists() {
total += read_trace_events(&f).len();
}
}
assert_eq!(total, 3);
}
#[test]
fn test_event_exactly_on_max_file_size_boundary() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.build()
.unwrap();
for _ in 0..2 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.finalize().unwrap();
let total: usize = (0..10)
.map(|i| {
let f = rotating_file(&base, i);
if std::path::Path::new(&f).exists() {
read_trace_events(&f).len()
} else {
0
}
})
.sum();
assert_eq!(total, 2);
}
#[test]
fn test_active_suffix_while_writing() {
let dir = TempDir::new().unwrap();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(1024)
.max_total_size(100000)
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
let active = dir.path().join("trace.0.bin.active");
assert!(active.exists(), "active file should exist while writing");
let sealed = dir.path().join("trace.0.bin");
assert!(!sealed.exists(), "sealed file should not exist yet");
}
#[test]
fn test_rotation_seals_previous_file() {
let dir = TempDir::new().unwrap();
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
assert!(
dir.path().join("trace.0.bin").exists(),
"rotated file should be sealed"
);
assert!(
!dir.path().join("trace.0.bin.active").exists(),
"rotated file should not be active"
);
assert!(
dir.path().join("trace.1.bin.active").exists(),
"current file should be active"
);
assert!(
!dir.path().join("trace.1.bin").exists(),
"current file should not be sealed"
);
}
#[test]
fn test_finalize_renames_current_file() {
let dir = TempDir::new().unwrap();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(1024)
.max_total_size(100000)
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.finalize().unwrap();
assert!(
dir.path().join("trace.0.bin").exists(),
"file should be sealed after finalize()"
);
assert!(
!dir.path().join("trace.0.bin.active").exists(),
"active file should be gone after finalize()"
);
}
#[test]
fn test_finalize_removes_empty_segment_after_rotation() {
let dir = TempDir::new().unwrap();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(1)
.max_total_size(100_000)
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
assert!(dir.path().join("trace.0.bin").exists());
assert!(dir.path().join("trace.1.bin.active").exists());
writer.finalize().unwrap();
assert!(
!dir.path().join("trace.1.bin").exists(),
"empty segment should not be sealed"
);
assert!(
!dir.path().join("trace.1.bin.active").exists(),
"empty active file should be removed"
);
assert!(dir.path().join("trace.0.bin").exists());
}
#[test]
fn test_single_file_no_active_suffix() {
let dir = TempDir::new().unwrap();
let path = dir.path().join("test.bin");
let mut writer = DiskBuffer::single_file(&path).unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.finalize().unwrap();
assert!(dir.path().join("test.0.bin").exists());
assert!(!dir.path().join("test.0.bin.active").exists());
}
#[test]
#[cfg(feature = "pipeline")]
fn test_single_file_sealed_segment_discoverable_by_worker() {
use crate::sealed::find_sealed_segments;
let dir = TempDir::new().unwrap();
let path = dir.path().join("trace.bin");
let mut writer = DiskBuffer::single_file(&path).unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.finalize().unwrap();
let segments = find_sealed_segments(dir.path(), "trace").unwrap();
assert_eq!(
segments.len(),
1,
"worker should find exactly one sealed segment"
);
assert_eq!(segments[0].path, dir.path().join("trace.0.bin"));
}
#[test]
fn test_segment_metadata_roundtrip() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(100_000)
.max_total_size(100_000)
.segment_metadata(vec![
("service".into(), "checkout-api".into()),
("host".into(), "i-0abc123".into()),
])
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.finalize().unwrap();
let all_events = decode_all(&std::fs::read(format!("{}.0.bin", base.display())).unwrap());
let metadata: Vec<_> = all_events
.iter()
.filter_map(|e| match e {
Decoded::SegmentMetadata { entries, .. } => Some(entries.clone()),
_ => None,
})
.collect();
assert_eq!(metadata.len(), 1);
assert!(
metadata[0].get("service").map(String::as_str) == Some("checkout-api"),
"missing service entry: {:?}",
metadata[0]
);
assert!(
metadata[0].get("host").map(String::as_str) == Some("i-0abc123"),
"missing host entry: {:?}",
metadata[0]
);
assert_eq!(
metadata[0].get(DIAL9_VERSION_KEY).map(String::as_str),
Some(DIAL9_VERSION_VALUE),
"missing built-in dial9.dial9-tokio-telemetry.version: {:?}",
metadata[0]
);
}
#[test]
fn test_segment_metadata_written_in_every_rotated_file() {
let dir = TempDir::new().unwrap();
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.segment_metadata(vec![("k".into(), "v".into())])
.build()
.unwrap();
for _ in 0..5 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.flush().unwrap();
writer.finalize().unwrap();
let mut files: Vec<_> = std::fs::read_dir(dir.path())
.unwrap()
.filter_map(|e| e.ok())
.map(|e| e.path())
.filter(|p| p.extension().is_some_and(|ext| ext == "bin"))
.collect();
files.sort();
assert!(files.len() >= 2, "expected at least 2 files from rotation");
for file in &files {
let all_events = decode_all(&std::fs::read(file).unwrap());
let has_metadata = all_events.iter().any(|e| match e {
Decoded::SegmentMetadata { entries, .. } => {
entries.get("k").map(String::as_str) == Some("v")
}
_ => false,
});
assert!(has_metadata, "{}: expected SegmentMetadata", file.display());
}
}
#[test]
fn test_dynamic_metadata_merged_on_rotation() {
let dir = TempDir::new().unwrap();
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.segment_metadata(vec![("service".into(), "myapp".into())])
.build()
.unwrap();
let mut merged = writer.segment_metadata().to_vec();
merged.push(("runtime.main".into(), "0,1,2,3".into()));
writer.update_segment_metadata(merged);
for _ in 0..4 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.flush().unwrap();
writer.finalize().unwrap();
let mut files: Vec<_> = std::fs::read_dir(dir.path())
.unwrap()
.filter_map(|e| e.ok())
.map(|e| e.path())
.filter(|p| p.extension().is_some_and(|ext| ext == "bin"))
.collect();
files.sort();
assert!(files.len() >= 2, "expected at least 2 files from rotation");
for file in &files[1..] {
let all_events = decode_all(&std::fs::read(file).unwrap());
let meta: Vec<_> = all_events
.iter()
.filter_map(|e| match e {
Decoded::SegmentMetadata { entries, .. } => Some(entries.clone()),
_ => None,
})
.collect();
assert_eq!(
meta.len(),
1,
"{}: expected 1 metadata event",
file.display()
);
assert!(
meta[0].get("service").map(String::as_str) == Some("myapp"),
"{}: missing static metadata",
file.display()
);
assert!(
meta[0].get("runtime.main").map(String::as_str) == Some("0,1,2,3"),
"{}: missing dynamic runtime worker metadata",
file.display()
);
}
}
#[test]
fn test_segment_metadata_empty_entries() {
let dir = TempDir::new().unwrap();
let path = dir.path().join("trace.bin");
let mut writer = DiskBuffer::single_file(&path).unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
let all_events = decode_all(&std::fs::read(writer.current_active_path()).unwrap());
let data_count = all_events
.iter()
.filter(|e| matches!(e, Decoded::Data { .. }))
.count();
assert_eq!(data_count, 1);
let metadata: Vec<_> = all_events
.iter()
.filter_map(|e| match e {
Decoded::SegmentMetadata { entries, .. } => Some(entries),
_ => None,
})
.collect();
assert_eq!(metadata.len(), 1);
assert_eq!(
metadata[0].get(DIAL9_VERSION_KEY).map(String::as_str),
Some(DIAL9_VERSION_VALUE)
);
}
#[test]
fn test_eviction_removes_gz_variant() {
let dir = TempDir::new().unwrap();
let one_event = single_event_file_size();
let max_file_size = one_event;
let max_total_size = max_file_size * 100;
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(max_file_size)
.max_total_size(max_total_size)
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
let seg0 = dir.path().join("trace.0.bin");
let seg0_gz = dir.path().join("trace.0.bin.gz");
assert!(seg0.exists(), "trace.0.bin should exist after rotation");
std::fs::rename(&seg0, &seg0_gz).unwrap();
writer.max_total_size = max_file_size;
for _ in 0..3 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.finalize().unwrap();
assert!(!seg0_gz.exists(), "trace.0.bin.gz should have been evicted");
}
#[test]
fn test_eviction_keeps_most_recent_segment_when_over_budget() {
let dir = TempDir::new().unwrap();
let one_event = single_event_file_size();
let max_file_size = u64::MAX;
let max_total_size = one_event / 2;
assert!(
max_total_size < one_event,
"test setup: budget must be smaller than a single segment"
);
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(max_file_size)
.max_total_size(max_total_size)
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
assert!(
writer.total_size() > max_total_size,
"single segment ({}) should exceed budget ({max_total_size})",
writer.total_size()
);
writer.evict_oldest().unwrap();
assert!(
matches!(writer.state, WriterState::Finished),
"writer should stop once even the most-recent segment exceeds budget"
);
assert!(
std::path::Path::new(&writer.current_active_path()).exists(),
"the most-recent segment must not be evicted"
);
assert!(
total_disk_usage(dir.path()) > max_total_size,
"retained segment is expected to push on-disk usage over the budget"
);
}
#[tokio::test(start_paused = true)]
async fn test_time_rotation_triggers_on_expired_boundary() {
use metrique_timesource::{TimeSource, tokio::set_time_source_for_current_runtime};
let _guard = set_time_source_for_current_runtime(TimeSource::tokio(std::time::UNIX_EPOCH));
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(u64::MAX)
.max_total_size(100_000)
.rotation_period(Duration::from_secs(60))
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
let initial_index = writer.next_index;
tokio::time::advance(Duration::from_secs(61)).await;
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.drained().unwrap();
assert!(
writer.next_index > initial_index,
"expected time-based rotation to trigger"
);
writer.finalize().unwrap();
let total: usize = (0..10)
.map(|i| {
let f = rotating_file(&base, i);
if std::path::Path::new(&f).exists() {
read_trace_events(&f).len()
} else {
0
}
})
.sum();
assert_eq!(total, 2);
}
#[tokio::test(start_paused = true)]
async fn test_first_rotation_uses_monotonic_period_not_wallclock_alignment() {
use metrique_timesource::{TimeSource, tokio::set_time_source_for_current_runtime};
let start_wall = std::time::UNIX_EPOCH + Duration::from_secs(22);
let _guard = set_time_source_for_current_runtime(TimeSource::tokio(start_wall));
let dir = TempDir::new().unwrap();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(u64::MAX)
.max_total_size(100_000)
.rotation_period(Duration::from_secs(60))
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
let initial_index = writer.next_index;
tokio::time::advance(Duration::from_secs(50)).await;
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.drained().unwrap();
assert_eq!(
writer.next_index, initial_index,
"rotation must not fire before one full rotation_period of monotonic time has elapsed",
);
tokio::time::advance(Duration::from_secs(11)).await;
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.drained().unwrap();
assert!(
writer.next_index > initial_index,
"rotation should fire once a full rotation_period of monotonic time has elapsed",
);
writer.finalize().unwrap();
}
#[tokio::test(start_paused = true)]
async fn test_time_rotation_skips_when_no_real_events() {
use metrique_timesource::{TimeSource, tokio::set_time_source_for_current_runtime};
let _guard = set_time_source_for_current_runtime(TimeSource::tokio(std::time::UNIX_EPOCH));
let dir = TempDir::new().unwrap();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(u64::MAX)
.max_total_size(100_000)
.rotation_period(Duration::from_secs(60))
.build()
.unwrap();
tokio::time::advance(Duration::from_secs(120)).await;
let empty_batch = Batch::new(vec![], 0);
writer.write_encoded_batch(&empty_batch).unwrap();
assert_eq!(
writer.next_index, 1,
"should not rotate when no real events exist"
);
writer.finalize().unwrap();
}
#[test]
fn test_size_rotation_still_works_with_time_disabled() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.rotation_period(std::time::Duration::MAX)
.build()
.unwrap();
for _ in 0..3 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.finalize().unwrap();
let total: usize = (0..10)
.map(|i| {
let f = rotating_file(&base, i);
if std::path::Path::new(&f).exists() {
read_trace_events(&f).len()
} else {
0
}
})
.sum();
assert_eq!(total, 3);
}
#[tokio::test(start_paused = true)]
async fn test_time_rotation_respects_eviction_budget() {
use metrique_timesource::{TimeSource, tokio::set_time_source_for_current_runtime};
let _guard = set_time_source_for_current_runtime(TimeSource::tokio(std::time::UNIX_EPOCH));
let dir = TempDir::new().unwrap();
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(u64::MAX)
.max_total_size(one_event * 3)
.rotation_period(Duration::from_secs(60))
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
for _ in 0..5 {
tokio::time::advance(Duration::from_secs(61)).await;
writer.write_encoded_batch(&test_batch()).unwrap();
writer.drained().unwrap();
}
writer.finalize().unwrap();
assert!(
total_disk_usage(dir.path()) <= one_event * 3,
"disk usage should stay within budget"
);
}
#[test]
fn test_builder_rotation_period_default() {
let dir = TempDir::new().unwrap();
let writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(1024)
.max_total_size(100_000)
.build()
.unwrap();
assert_eq!(writer.rotation_period, DEFAULT_ROTATION_PERIOD);
}
#[test]
fn test_new_uses_default_rotation_period() {
let dir = TempDir::new().unwrap();
let writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(1024)
.max_total_size(100_000)
.build()
.unwrap();
assert_eq!(writer.rotation_period, DEFAULT_ROTATION_PERIOD);
}
#[tokio::test(start_paused = true)]
async fn test_finalize_after_time_rotation() {
use metrique_timesource::{TimeSource, tokio::set_time_source_for_current_runtime};
let _guard = set_time_source_for_current_runtime(TimeSource::tokio(std::time::UNIX_EPOCH));
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(u64::MAX)
.max_total_size(100_000)
.rotation_period(Duration::from_secs(60))
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
tokio::time::advance(Duration::from_secs(61)).await;
writer.write_encoded_batch(&test_batch()).unwrap();
writer.drained().unwrap();
writer.finalize().unwrap();
let total: usize = (0..10)
.map(|i| {
let f = rotating_file(&base, i);
if std::path::Path::new(&f).exists() {
read_trace_events(&f).len()
} else {
0
}
})
.sum();
assert_eq!(total, 2);
}
#[tokio::test(start_paused = true)]
async fn test_stale_boundary_does_not_rotate_first_event() {
use metrique_timesource::{TimeSource, tokio::set_time_source_for_current_runtime};
let _guard = set_time_source_for_current_runtime(TimeSource::tokio(std::time::UNIX_EPOCH));
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(u64::MAX)
.max_total_size(100_000)
.rotation_period(Duration::from_secs(60))
.build()
.unwrap();
tokio::time::advance(Duration::from_secs(300)).await;
writer.write_encoded_batch(&test_batch()).unwrap();
assert_eq!(
writer.next_index, 1,
"first event after idle gap should not trigger immediate rotation"
);
writer.write_encoded_batch(&test_batch()).unwrap();
assert_eq!(
writer.next_index, 1,
"second event should still be in the same segment"
);
writer.finalize().unwrap();
let events = read_trace_events(&rotating_file(&base, 0));
assert_eq!(events.len(), 2, "both events should be in segment 0");
}
#[test]
fn test_clock_sync_precedes_first_data_event() {
use crate::sealed::LEGACY_EPOCH_NS_FLOOR;
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(100_000)
.max_total_size(100_000)
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.finalize().unwrap();
let data = std::fs::read(rotating_file(&base, 0)).unwrap();
let all = decode_all(&data);
let first_data_idx = all
.iter()
.position(|e| matches!(e, Decoded::Data { .. }))
.expect("expected at least one data event");
let first_clock_sync_idx = all
.iter()
.position(|e| matches!(e, Decoded::ClockSync { .. }))
.expect("expected a ClockSyncEvent in the file");
assert!(first_clock_sync_idx < first_data_idx);
match &all[first_clock_sync_idx] {
Decoded::ClockSync { realtime_ns, .. } => {
assert!(*realtime_ns >= LEGACY_EPOCH_NS_FLOOR);
}
_ => unreachable!(),
}
}
#[test]
fn test_segment_metadata_timestamp_is_monotonic_scale() {
use crate::sealed::LEGACY_EPOCH_NS_FLOOR;
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(100_000)
.max_total_size(100_000)
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.finalize().unwrap();
let data = std::fs::read(rotating_file(&base, 0)).unwrap();
let all = decode_all(&data);
let seg_ts = all
.iter()
.find_map(|e| match e {
Decoded::SegmentMetadata { timestamp_ns, .. } => Some(*timestamp_ns),
_ => None,
})
.expect("SegmentMetadata");
assert!(
seg_ts < LEGACY_EPOCH_NS_FLOOR,
"SegmentMetadata.timestamp_nanos ({seg_ts}) should be monotonic-scale"
);
}
#[test]
fn test_clock_sync_written_in_every_rotated_file() {
let dir = TempDir::new().unwrap();
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.build()
.unwrap();
for _ in 0..5 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.flush().unwrap();
writer.finalize().unwrap();
let mut files: Vec<_> = std::fs::read_dir(dir.path())
.unwrap()
.filter_map(|e| e.ok())
.map(|e| e.path())
.filter(|p| p.extension().is_some_and(|ext| ext == "bin"))
.collect();
files.sort();
assert!(files.len() >= 2, "expected at least 2 files from rotation");
for file in &files {
let all = decode_all(&std::fs::read(file).unwrap());
let has_clock_sync = all.iter().any(|e| matches!(e, Decoded::ClockSync { .. }));
assert!(
has_clock_sync,
"{}: expected ClockSyncEvent",
file.display()
);
}
}
#[test]
fn test_legacy_trace_without_clock_sync_still_decodes() {
let mut enc = Encoder::new_to(Vec::new()).unwrap();
enc.write(&SegmentMetadataEvent {
timestamp_ns: 1,
entries: vec![("k".into(), "v".into())],
})
.unwrap();
enc.write(&TestEvent {
timestamp_ns: 1000,
value: 0,
})
.unwrap();
let buf = enc.into_inner();
let all = decode_all(&buf);
assert!(
all.iter().any(|e| matches!(e, Decoded::Data { .. })),
"expected data event to decode"
);
assert!(
!all.iter().any(|e| matches!(e, Decoded::ClockSync { .. })),
"legacy trace must not contain ClockSync"
);
}
#[test]
fn test_clock_sync_offset_recovers_wall_clock_for_recent_event() {
use std::time::{SystemTime, UNIX_EPOCH};
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(100_000)
.max_total_size(100_000)
.build()
.unwrap();
let park_ts = crate::clock::clock_monotonic_ns();
let mut enc = Encoder::new_to(Vec::new()).unwrap();
enc.write(&TestEvent {
timestamp_ns: park_ts,
value: 0,
})
.unwrap();
writer
.write_encoded_batch(&Batch::new(enc.into_inner(), 1))
.unwrap();
writer.flush().unwrap();
writer.finalize().unwrap();
let all = decode_all(&std::fs::read(rotating_file(&base, 0)).unwrap());
let (sync_mono, sync_real) = all
.iter()
.find_map(|e| match e {
Decoded::ClockSync {
timestamp_ns,
realtime_ns,
} => Some((*timestamp_ns, *realtime_ns)),
_ => None,
})
.expect("ClockSync");
let park_from_file = all
.iter()
.find_map(|e| match e {
Decoded::Data { timestamp_ns } => Some(*timestamp_ns),
_ => None,
})
.expect("data event");
let offset = sync_real as i128 - sync_mono as i128;
let reconstructed_wall_ns = park_from_file as i128 + offset;
let now_ns = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_nanos() as i128;
let diff = (reconstructed_wall_ns - now_ns).abs();
assert!(
diff < 5_000_000_000,
"reconstructed wall clock {reconstructed_wall_ns} diverges from now {now_ns} by {diff}ns"
);
}
#[test]
fn test_update_segment_metadata_appears_in_trace() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(100_000)
.max_total_size(100_000)
.build()
.unwrap();
writer.update_segment_metadata(vec![
("bucket".into(), "my-bucket".into()),
("service_name".into(), "my-svc".into()),
]);
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.finalize().unwrap();
let all = decode_all(&std::fs::read(rotating_file(&base, 0)).unwrap());
let metadata: Vec<_> = all
.iter()
.filter_map(|e| match e {
Decoded::SegmentMetadata { entries, .. } => Some(entries.clone()),
_ => None,
})
.collect();
assert!(!metadata.is_empty(), "expected SegmentMetadata event");
assert!(
metadata.last().unwrap().get("bucket").map(String::as_str) == Some("my-bucket"),
"S3 metadata should be in segment"
);
assert!(
metadata
.last()
.unwrap()
.get("service_name")
.map(String::as_str)
== Some("my-svc"),
"S3 metadata should be in segment"
);
}
#[test]
fn test_merge_preserves_s3_metadata_across_runtime_updates() {
let dir = TempDir::new().unwrap();
let one_event = single_event_file_size();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.build()
.unwrap();
writer.update_segment_metadata(vec![
("bucket".into(), "my-bucket".into()),
("service_name".into(), "my-svc".into()),
]);
writer.update_segment_metadata(vec![("runtime.main".into(), "0,1".into())]);
for _ in 0..4 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.flush().unwrap();
writer.finalize().unwrap();
let mut files: Vec<_> = std::fs::read_dir(dir.path())
.unwrap()
.filter_map(|e| e.ok())
.map(|e| e.path())
.filter(|p| p.extension().is_some_and(|ext| ext == "bin"))
.collect();
files.sort();
assert!(files.len() >= 2, "expected rotation");
for file in &files[1..] {
let all = decode_all(&std::fs::read(file).unwrap());
let meta: Vec<_> = all
.iter()
.filter_map(|e| match e {
Decoded::SegmentMetadata { entries, .. } => Some(entries.clone()),
_ => None,
})
.collect();
let last = meta.last().expect("expected SegmentMetadata");
assert!(
last.get("bucket").map(String::as_str) == Some("my-bucket"),
"{}: S3 metadata lost after merge",
file.display()
);
assert!(
last.get("runtime.main").map(String::as_str) == Some("0,1"),
"{}: runtime metadata missing",
file.display()
);
}
}
#[test]
fn test_update_segment_metadata_no_op_when_unchanged() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(100_000)
.max_total_size(100_000)
.build()
.unwrap();
let entries = vec![("k".into(), "v".into())];
writer.update_segment_metadata(entries.clone());
writer.write_encoded_batch(&test_batch()).unwrap();
writer.update_segment_metadata(entries.clone());
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.finalize().unwrap();
let all = decode_all(&std::fs::read(rotating_file(&base, 0)).unwrap());
let metadata_count = all
.iter()
.filter(|e| matches!(e, Decoded::SegmentMetadata { .. }))
.count();
assert_eq!(
metadata_count, 1,
"identical update_segment_metadata should not trigger another write"
);
}
#[test]
fn test_dial9_version_in_segment_metadata() {
let dir = TempDir::new().unwrap();
let path = dir.path().join("trace.bin");
let mut writer = DiskBuffer::single_file(&path).unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.finalize().unwrap();
let sealed = dir.path().join("trace.0.bin");
let all = decode_all(&std::fs::read(&sealed).unwrap());
let version_value = all.iter().find_map(|e| match e {
Decoded::SegmentMetadata { entries, .. } => entries.get(DIAL9_VERSION_KEY).cloned(),
_ => None,
});
assert_eq!(
version_value.as_deref(),
Some(env!("CARGO_PKG_VERSION")),
"expected dial9.dial9-tokio-telemetry.version entry matching CARGO_PKG_VERSION"
);
}
#[test]
fn test_available_parallelism_in_segment_metadata() {
let expected = std::thread::available_parallelism().map(|n| n.get().to_string());
let dir = TempDir::new().unwrap();
let path = dir.path().join("trace.bin");
let mut writer = DiskBuffer::single_file(&path).unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.finalize().unwrap();
let sealed = dir.path().join("trace.0.bin");
let all = decode_all(&std::fs::read(&sealed).unwrap());
let value = all.iter().find_map(|e| match e {
Decoded::SegmentMetadata { entries, .. } => {
entries.get(PROCESS_AVAILABLE_PARALLELISM_KEY).cloned()
}
_ => None,
});
match expected {
Ok(expected) => assert_eq!(
value.as_deref(),
Some(expected.as_str()),
"expected process.available_parallelism entry matching std::thread::available_parallelism()"
),
Err(_) => assert!(
value.is_none(),
"process.available_parallelism should be omitted when available_parallelism() fails"
),
}
}
#[test]
fn test_dial9_version_user_override_wins() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(100_000)
.max_total_size(100_000)
.segment_metadata(vec![(DIAL9_VERSION_KEY.into(), "builder-override".into())])
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.rotate().unwrap();
writer.update_segment_metadata(vec![(DIAL9_VERSION_KEY.into(), "runtime-override".into())]);
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
writer.finalize().unwrap();
let read_version = |idx: u32| -> String {
let all = decode_all(&std::fs::read(rotating_file(&base, idx)).unwrap());
all.iter()
.find_map(|e| match e {
Decoded::SegmentMetadata { entries, .. } => {
entries.get(DIAL9_VERSION_KEY).cloned()
}
_ => None,
})
.expect("expected dial9.dial9-tokio-telemetry.version entry")
};
assert_eq!(read_version(0), "builder-override");
assert_eq!(read_version(1), "runtime-override");
}
#[tokio::test(start_paused = true)]
async fn test_drained_recovers_when_active_file_deleted() {
use metrique_timesource::{TimeSource, tokio::set_time_source_for_current_runtime};
let _guard = set_time_source_for_current_runtime(TimeSource::tokio(std::time::UNIX_EPOCH));
let dir = TempDir::new().unwrap();
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(u64::MAX)
.max_total_size(100_000)
.rotation_period(Duration::from_secs(60))
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
let active_path = writer.current_active_path().to_owned();
assert!(active_path.exists());
std::fs::remove_file(&active_path).unwrap();
tokio::time::advance(Duration::from_secs(61)).await;
assert!(writer.should_drain(), "should_drain should fire");
writer
.drained()
.expect("drained() must recover from missing .active file");
assert!(
!writer.should_drain(),
"should_drain must return false after recovery (otherwise flush loop spins)"
);
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
assert!(
writer.current_active_path().exists(),
"writer must have a fresh active file after recovery"
);
writer.finalize().unwrap();
}
#[tokio::test(start_paused = true)]
async fn test_drained_recovers_when_parent_dir_deleted() {
use metrique_timesource::{TimeSource, tokio::set_time_source_for_current_runtime};
let _guard = set_time_source_for_current_runtime(TimeSource::tokio(std::time::UNIX_EPOCH));
let dir = TempDir::new().unwrap();
let trace_dir = dir.path().join("traces");
std::fs::create_dir_all(&trace_dir).unwrap();
let mut writer = DiskBuffer::builder()
.base_path(&trace_dir)
.max_file_size(u64::MAX)
.max_total_size(100_000)
.rotation_period(Duration::from_secs(60))
.build()
.unwrap();
writer.write_encoded_batch(&test_batch()).unwrap();
writer.flush().unwrap();
std::fs::remove_dir_all(&trace_dir).unwrap();
assert!(!writer.current_active_path().exists());
tokio::time::advance(Duration::from_secs(61)).await;
assert!(writer.should_drain());
let _ = writer.drained();
assert!(
!writer.should_drain(),
"should_drain must return false after a failed rotation \
(otherwise the flush loop spins on every 5ms tick)"
);
tokio::time::advance(Duration::from_millis(5)).await;
let _ = writer.drained();
assert!(!writer.should_drain());
}
#[test]
fn test_restart_seeds_closed_files_and_evicts() {
let dir = TempDir::new().unwrap();
let base = dir.path().join("trace");
let one_event = single_event_file_size();
{
let mut w = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(100_000)
.build()
.unwrap();
for _ in 0..4 {
w.write_encoded_batch(&test_batch()).unwrap();
}
w.finalize().unwrap();
}
let bin_count_before = (0..20)
.filter(|i| std::path::Path::new(&rotating_file(&base, *i)).exists())
.count();
assert!(
bin_count_before >= 2,
"lifetime 1 should leave multiple sealed segments"
);
let new_budget = one_event + 1; let writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(new_budget)
.build()
.unwrap();
assert!(
total_disk_usage(dir.path()) <= new_budget,
"disk usage exceeds shrunk budget after restart: {}",
total_disk_usage(dir.path())
);
let next_active_path = writer.current_active_path();
assert!(next_active_path.exists());
assert!(
next_active_path
.to_str()
.is_some_and(|s| s.ends_with(".bin.active"))
);
}
#[test]
fn test_restart_discards_stale_active_files() {
let dir = TempDir::new().unwrap();
let orphan = dir.path().join("trace.99.bin.active");
std::fs::write(&orphan, b"orphaned").unwrap();
let _w = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(1024)
.max_total_size(100_000)
.build()
.unwrap();
assert!(
!orphan.exists(),
"stale .active should be discarded on construction"
);
}
#[test]
fn test_restart_counts_gz_siblings_toward_budget() {
let dir = TempDir::new().unwrap();
let bin = dir.path().join("trace.0.bin");
let gz = dir.path().join("trace.0.bin.gz");
std::fs::write(&bin, vec![0u8; 4096]).unwrap();
std::fs::write(&gz, vec![0u8; 1024]).unwrap();
let _w = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(100_000)
.max_total_size(100)
.build()
.unwrap();
assert!(!bin.exists(), ".bin should be evicted under restart budget");
assert!(!gz.exists(), ".bin.gz must be evicted with its .bin family");
}
#[test]
fn test_finalize_evicts_to_budget() {
let dir = TempDir::new().unwrap();
let one_event = single_event_file_size();
let max_total_size = one_event * 2;
let mut writer = DiskBuffer::builder()
.base_path(dir.path())
.max_file_size(one_event)
.max_total_size(max_total_size)
.build()
.unwrap();
for _ in 0..10 {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.finalize().unwrap();
assert!(
total_disk_usage(dir.path()) <= max_total_size,
"finalize must leave disk usage within budget"
);
}
#[test]
fn in_memory_builder_wires_custom_options() {
let writer = MemoryBuffer::builder()
.max_total_size(8 * 1024 * 1024)
.max_segment_size(64 * 1024)
.rotation_period(Duration::from_secs(30))
.segment_metadata(vec![("svc".into(), "test".into())])
.build()
.unwrap();
assert_eq!(writer.max_file_size, 64 * 1024);
assert_eq!(writer.rotation_period, Duration::from_secs(30));
assert!(
writer
.segment_metadata
.entries
.iter()
.any(|(k, v)| k == "svc" && v == "test")
);
}
#[test]
fn in_memory_rejects_zero_total_size() {
let err = MemoryBuffer::new(0).unwrap_err();
assert_eq!(err.kind(), std::io::ErrorKind::InvalidInput);
}
#[test]
fn in_memory_builder_enforces_3x_segment_min_total_size() {
let seg: u64 = 2048;
let err = MemoryBuffer::builder()
.max_total_size(3 * seg - 1)
.max_segment_size(seg)
.build()
.unwrap_err();
assert_eq!(err.kind(), std::io::ErrorKind::InvalidInput);
MemoryBuffer::builder()
.max_total_size(3 * seg)
.max_segment_size(seg)
.build()
.expect("3× segment must be accepted");
}
#[cfg(feature = "pipeline")]
mod mem_e2e_tests {
use super::*;
use crate::pipeline::{ProcessError, SegmentData, SegmentProcessor};
use crate::worker::WorkerLoop;
use std::future::Future;
use std::pin::Pin;
use std::sync::{Arc, Mutex};
use std::time::Duration;
struct CapturingProcessor {
segments: Arc<Mutex<Vec<Vec<u8>>>>,
}
impl CapturingProcessor {
fn new() -> (Self, Arc<Mutex<Vec<Vec<u8>>>>) {
let segments = Arc::new(Mutex::new(Vec::new()));
(
Self {
segments: segments.clone(),
},
segments,
)
}
}
impl SegmentProcessor for CapturingProcessor {
fn name(&self) -> &'static str {
"Capture"
}
fn process(
&mut self,
data: SegmentData,
) -> Pin<Box<dyn Future<Output = Result<SegmentData, ProcessError>> + Send + '_>>
{
self.segments
.lock()
.unwrap()
.push(data.payload().clone().into_vec());
Box::pin(async move { Ok(data) })
}
}
async fn run_mem_e2e(mut writer: MemoryBuffer, events: usize) -> Vec<Vec<u8>> {
let fs = writer.fs_handle().expect("memory writer exposes its Fs");
for _ in 0..events {
writer.write_encoded_batch(&test_batch()).unwrap();
}
writer.finalize().unwrap();
let (capture, captured) = CapturingProcessor::new();
let stop = tokio_util::sync::CancellationToken::new();
let mut worker = WorkerLoop::new(
fs,
Duration::from_millis(5),
vec![Box::new(capture)],
stop,
metrique_writer::sink::DevNullSink::boxed(),
None,
)
.await
.expect("initialize worker");
worker.run().await;
let segments = captured.lock().unwrap();
segments.clone()
}
fn count_payload_events(segments: &[Vec<u8>]) -> usize {
segments
.iter()
.flat_map(|s| decode_all(s))
.filter(|e| matches!(e, Decoded::Data { .. }))
.count()
}
#[tokio::test]
async fn mem_writer_e2e_delivers_all_events() {
const EVENTS: usize = 25;
let segments = run_mem_e2e(MemoryBuffer::new(1 << 20).unwrap(), EVENTS).await;
assert!(!segments.is_empty(), "worker captured no segments");
assert_eq!(
count_payload_events(&segments),
EVENTS,
"every written event must reach the processor"
);
}
#[tokio::test]
async fn mem_writer_e2e_delivers_all_events_across_rotations() {
const EVENTS: usize = 60;
let writer = MemoryBuffer::builder()
.max_total_size(16 * 1024 * 1024)
.max_segment_size(256)
.build()
.unwrap();
let segments = run_mem_e2e(writer, EVENTS).await;
assert!(
segments.len() >= 2,
"tiny segments must force rotation, got {} segment(s)",
segments.len()
);
assert_eq!(
count_payload_events(&segments),
EVENTS,
"every event across all rotated segments must reach the processor"
);
}
}
}