use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering::Relaxed};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use mira_core::SignalBuilder;
use mira_core::block;
use mira_core::signal::Open;
use mira_core::wal::{self, Wal};
use tokio::sync::{mpsc, oneshot};
use tokio::task::JoinHandle;
use tokio::time::{Instant, sleep_until};
pub struct Config {
pub data_dir: PathBuf,
pub node: u32,
pub target_block_bytes: usize,
pub max_block_age: Duration,
pub retention: Duration,
pub queue: usize,
pub wal: Option<Arc<Wal>>,
}
impl Default for Config {
fn default() -> Self {
Self {
data_dir: PathBuf::from("./data"),
node: block::node_id("mira"),
target_block_bytes: 32 << 20,
max_block_age: Duration::from_secs(2),
retention: Duration::from_secs(7 * 24 * 3600),
queue: 128,
wal: None,
}
}
}
pub(crate) struct Job<R> {
req: R,
ack: oneshot::Sender<Result<(), Rejected>>,
wal_seq: Option<u64>,
}
pub struct Ingest<R> {
pub(crate) tx: mpsc::Sender<Job<R>>,
pub(crate) rejects: &'static Rejects,
pub(crate) wal: Option<Arc<Wal>>,
pub(crate) signal: wal::Signal,
}
impl<R> Clone for Ingest<R> {
fn clone(&self) -> Self {
Self {
tx: self.tx.clone(),
rejects: self.rejects,
wal: self.wal.clone(),
signal: self.signal,
}
}
}
pub enum Rejected {
Busy,
Closed,
Unavailable(String),
Failed(String),
}
fn now_secs() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs()
}
fn once_a_second(gate: &AtomicU64) -> bool {
let now = now_secs();
gate.swap(now, Relaxed) != now
}
pub const UNREADY_AFTER: Duration = Duration::from_secs(120);
const ADMIT_WAIT: Duration = Duration::from_secs(5);
pub struct Rejects {
pub signal: &'static str,
pub shed: AtomicU64,
pub failed: AtomicU64,
pub refused: AtomicU64,
pub published: AtomicU64,
pub rows: AtomicU64,
pub bytes: AtomicU64,
pub open_since: AtomicU64,
pub stalled_since: AtomicU64,
warned: AtomicU64,
refuse_warned: AtomicU64,
}
impl Rejects {
const fn new(signal: &'static str) -> Self {
Self {
signal,
shed: AtomicU64::new(0),
failed: AtomicU64::new(0),
refused: AtomicU64::new(0),
published: AtomicU64::new(0),
rows: AtomicU64::new(0),
bytes: AtomicU64::new(0),
open_since: AtomicU64::new(0),
stalled_since: AtomicU64::new(0),
warned: AtomicU64::new(0),
refuse_warned: AtomicU64::new(0),
}
}
fn mark_stalled(&self) {
let _ = self
.stalled_since
.compare_exchange(0, now_secs().max(1), Relaxed, Relaxed);
}
fn record_shed(&self) {
self.shed.fetch_add(1, Relaxed);
if once_a_second(&self.warned) {
tracing::warn!(
signal = self.signal,
"ingest queue full; shedding exports (senders are told to retry)"
);
}
}
}
pub static REJECTS: [Rejects; SIGNALS.len()] = [
Rejects::new(SIGNALS[0]),
Rejects::new(SIGNALS[1]),
Rejects::new(SIGNALS[2]),
];
fn rejects_for(signal: &str) -> &'static Rejects {
REJECTS
.iter()
.find(|r| r.signal == signal)
.expect("every signal that has a builder has a counter slot")
}
fn stall_of(r: &Rejects, now: u64) -> Option<u64> {
match r.stalled_since.load(Relaxed) {
0 => None,
since => {
let secs = now.saturating_sub(since);
(secs >= UNREADY_AFTER.as_secs()).then_some(secs)
}
}
}
pub fn stalled() -> Option<(&'static str, u64)> {
let now = now_secs();
REJECTS
.iter()
.find_map(|r| stall_of(r, now).map(|secs| (r.signal, secs)))
}
impl<R: prost::Message> Ingest<R> {
pub async fn submit(&self, req: R) -> Result<(), Rejected> {
let (ack, wait) = oneshot::channel();
let permit = match self.tx.try_reserve() {
Ok(p) => p,
Err(mpsc::error::TrySendError::Full(())) => {
match tokio::time::timeout(ADMIT_WAIT, self.tx.reserve()).await {
Ok(Ok(p)) => p,
Ok(Err(_)) => return Err(Rejected::Closed),
Err(_) => {
self.rejects.record_shed();
return Err(Rejected::Busy);
}
}
}
Err(mpsc::error::TrySendError::Closed(())) => return Err(Rejected::Closed),
};
if let Some(wal) = &self.wal {
let body = req.encode_to_vec();
return match wal.append_then(self.signal, &body, move |seq| {
permit.send(Job {
req,
ack,
wal_seq: Some(seq),
});
}) {
Ok(_) => Ok(()),
Err(e) => {
self.rejects.failed.fetch_add(1, Relaxed);
Err(match e {
mira_core::Error::WalFrameTooLarge { .. } => {
Rejected::Failed(e.to_string())
}
_ => Rejected::Unavailable(e.to_string()),
})
}
};
}
permit.send(Job {
req,
ack,
wal_seq: None,
});
match wait.await {
Ok(Ok(())) => Ok(()),
Ok(Err(r)) => {
self.rejects.failed.fetch_add(1, Relaxed);
Err(r)
}
Err(_) => Err(Rejected::Closed),
}
}
}
impl<R: prost::Message + Default> Ingest<R> {
pub fn replay(&self, body: &[u8], seq: u64) -> Result<(), Rejected> {
let req = R::decode(body).map_err(|e| Rejected::Failed(e.to_string()))?;
self.tx
.blocking_send(Job {
req,
ack: oneshot::channel().0,
wal_seq: Some(seq),
})
.map_err(|_| Rejected::Closed)
}
}
pub const SIGNALS: [&str; 3] = ["logs", "traces", "metrics"];
pub fn spawn<B: SignalBuilder>(cfg: Arc<Config>) -> (Ingest<B::Request>, OpenSlot, JoinHandle<()>) {
let (tx, rx) = mpsc::channel(cfg.queue);
let rejects = rejects_for(B::SIGNAL);
let ingest = Ingest {
tx,
rejects,
wal: cfg.wal.clone(),
signal: wal::Signal::named(B::SIGNAL).expect("every signal has a log discriminant"),
};
let (ask, asks) = mpsc::channel(8);
let slot = OpenSlot {
cur: Arc::default(),
ask,
};
(
ingest,
slot.clone(),
tokio::spawn(flusher::<B>(rx, asks, cfg, slot)),
)
}
#[derive(Clone)]
pub struct OpenSlot {
cur: Arc<Mutex<Option<Arc<Open>>>>,
ask: mpsc::Sender<oneshot::Sender<Option<Arc<Open>>>>,
}
impl Default for OpenSlot {
fn default() -> Self {
Self {
cur: Arc::default(),
ask: mpsc::channel(1).0,
}
}
}
impl OpenSlot {
pub async fn fresh(&self) -> Option<Arc<Open>> {
let (tx, rx) = oneshot::channel();
match self.ask.try_send(tx) {
Ok(()) => rx.await.unwrap_or_else(|_| self.get()),
Err(_) => self.get(),
}
}
pub fn get(&self) -> Option<Arc<Open>> {
self.lock().clone()
}
fn put(&self, v: Option<Arc<Open>>) {
*self.lock() = v;
}
fn lock(&self) -> std::sync::MutexGuard<'_, Option<Arc<Open>>> {
self.cur.lock().unwrap_or_else(|e| e.into_inner())
}
}
pub type OpenSlots = [OpenSlot; SIGNALS.len()];
pub fn spawn_retention(cfg: Arc<Config>) {
tokio::spawn(retention(cfg.clone()));
if cfg.wal.is_some() {
tokio::spawn(wal_maintenance(cfg));
}
}
pub const WAL_SYNC_PERIOD: Duration = Duration::from_millis(250);
async fn wal_maintenance(cfg: Arc<Config>) {
let Some(wal) = cfg.wal.clone() else { return };
let mut tick = tokio::time::interval(WAL_SYNC_PERIOD);
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
let mut ticks: u64 = 0;
loop {
tick.tick().await;
ticks += 1;
wal_sweep(wal.clone(), cfg.data_dir.clone(), ticks % 240 == 0).await;
}
}
async fn wal_sweep(wal: Arc<Wal>, dir: PathBuf, truncating: bool) {
let done = tokio::task::spawn_blocking(move || {
wal.sync()?;
if !truncating {
return Ok(0);
}
let covered = block::wal_watermarks(&dir)?.into_iter().min().unwrap_or(0);
wal.truncate(covered)
})
.await;
match done {
Ok(Ok(0)) => {}
Ok(Ok(n)) => tracing::info!(segments = n, "write-ahead log segments removed"),
Ok(Err(e)) => tracing::warn!(error = %e, "write-ahead log maintenance failed"),
Err(e) => tracing::warn!(error = %e, "write-ahead log maintenance panicked"),
}
}
async fn flusher<B: SignalBuilder>(
mut rx: mpsc::Receiver<Job<B::Request>>,
mut asks: mpsc::Receiver<oneshot::Sender<Option<Arc<Open>>>>,
cfg: Arc<Config>,
open_slot: OpenSlot,
) {
let rejects = rejects_for(B::SIGNAL);
let mut seq = match block::scan(&cfg.data_dir, B::SIGNAL) {
Ok(blocks) => blocks.iter().map(|b| b.seq).max().map_or(0, |s| s + 1),
Err(e) => {
tracing::error!(signal = B::SIGNAL, error = %e, "cannot scan data directory");
return;
}
};
match block::sweep_staging(&cfg.data_dir, B::SIGNAL, cfg.node) {
Ok(0) => {}
Ok(n) => tracing::info!(signal = B::SIGNAL, count = n, "swept stale staging dirs"),
Err(e) => tracing::warn!(signal = B::SIGNAL, error = %e, "cannot sweep staging dirs"),
}
let mut builder = B::default();
let mut waiters: Vec<oneshot::Sender<Result<(), Rejected>>> = Vec::new();
let mut batch = Vec::with_capacity(64);
let mut carry: Vec<Job<B::Request>> = Vec::new();
let mut deadline = Instant::now() + cfg.max_block_age;
let mut open = true;
let mut wal_hi: u64 = 0;
let mut asked: Vec<oneshot::Sender<Option<Arc<Open>>>> = Vec::new();
let mut snapped: Option<usize> = None;
while open || !carry.is_empty() {
answer::<B>(
&builder,
&mut asked,
&mut snapped,
!carry.is_empty() || !rx.is_empty(),
&open_slot,
cfg.node,
seq,
);
let mut aged = false;
if carry.is_empty() {
tokio::select! {
n = rx.recv_many(&mut batch, 64) => {
if n == 0 {
open = false;
}
}
who = asks.recv() => {
if let Some(who) = who {
asked.push(who);
}
}
_ = sleep_until(deadline) => aged = true,
}
}
while let Ok(who) = asks.try_recv() {
asked.push(who);
}
let mut jobs = std::mem::take(&mut carry);
jobs.append(&mut batch);
let mut dict_full = false;
for job in jobs {
let empty = builder.is_empty();
if !empty && (dict_full || !builder.has_headroom_for(&job.req)) {
dict_full = true;
carry.push(job);
continue;
}
if waiters.is_empty() {
deadline = Instant::now() + cfg.max_block_age;
}
if let Some(seq) = job.wal_seq {
wal_hi = wal_hi.max(seq + 1);
}
match builder.append_request(&job.req) {
Ok(0) => {
let _ = job.ack.send(Ok(()));
}
Ok(_) => {
if waiters.is_empty() {
rejects.open_since.store(now_secs(), Relaxed);
}
waiters.push(job.ack);
}
Err(e) => {
if empty {
let _ = builder.finish();
}
let refused = rejects.refused.fetch_add(1, Relaxed) + 1;
if once_a_second(&rejects.refuse_warned) {
tracing::error!(
signal = B::SIGNAL,
error = %e,
refused,
"export permanently refused; its records are gone. The sender \
is told not to retry, so nothing will bring them back — the \
request does not fit an empty block, which means splitting it \
at the sender is the only fix"
);
}
let _ = job.ack.send(Err(Rejected::Failed(e.to_string())));
}
}
}
let full = dict_full || builder.approx_bytes() >= cfg.target_block_bytes;
if builder.is_empty() || !(full || (aged && !waiters.is_empty()) || !open) {
if waiters.is_empty() {
deadline = Instant::now() + cfg.max_block_age;
rejects.open_since.store(0, Relaxed);
}
continue;
}
rejects.open_since.store(0, Relaxed);
let sealed = match builder.finish() {
Ok(s) => s,
Err(e) => {
let msg = e.to_string();
rejects.mark_stalled();
for w in waiters.drain(..) {
let _ = w.send(Err(Rejected::Unavailable(msg.clone())));
}
wal_hi = 0;
open_slot.put(None);
tracing::error!(signal = B::SIGNAL, error = %msg, "block discarded");
continue;
}
};
let dir = cfg.data_dir.clone();
let node = cfg.node;
let this_seq = seq;
seq += 1;
let block_wal_hi = std::mem::take(&mut wal_hi);
let rows = sealed.num_rows;
let result = tokio::task::spawn_blocking(move || {
block::publish(&dir, B::SIGNAL, node, this_seq, block_wal_hi, &sealed)
.map(|b| (dir_bytes(&b.dir), b.dir))
})
.await;
open_slot.put(None);
snapped = None;
let outcome = match result {
Ok(Ok((bytes, path))) => {
rejects.published.fetch_add(1, Relaxed);
rejects.rows.fetch_add(rows as u64, Relaxed);
rejects.bytes.fetch_add(bytes, Relaxed);
rejects.stalled_since.store(0, Relaxed);
tracing::info!(signal = B::SIGNAL, rows, bytes, seq = this_seq, path = %path.display(), "block published");
Ok(())
}
Ok(Err(e)) => {
rejects.mark_stalled();
tracing::error!(signal = B::SIGNAL, seq = this_seq, error = %e, "block not published");
Err(e.to_string())
}
Err(e) => {
rejects.mark_stalled();
Err(format!("flush task panicked: {e}"))
}
};
for w in waiters.drain(..) {
let _ = w.send(outcome.clone().map_err(Rejected::Unavailable));
}
deadline = Instant::now() + cfg.max_block_age;
}
}
fn answer<B: SignalBuilder>(
builder: &B,
asked: &mut Vec<oneshot::Sender<Option<Arc<Open>>>>,
snapped: &mut Option<usize>,
pending: bool,
slot: &OpenSlot,
node: u32,
seq: u64,
) {
if asked.is_empty() || pending {
return;
}
let bytes = builder.approx_bytes();
if builder.is_empty() {
slot.put(None);
*snapped = None;
} else if *snapped != Some(bytes) {
*snapped = Some(bytes);
match builder.snapshot() {
Ok(sealed) => slot.put(Some(Arc::new(Open { node, seq, sealed }))),
Err(e) => {
slot.put(None);
tracing::debug!(signal = B::SIGNAL, error = %e, "open block not snapshotted");
}
}
}
let cur = slot.get();
for who in asked.drain(..) {
let _ = who.send(cur.clone());
}
}
fn dir_bytes(dir: &Path) -> u64 {
std::fs::read_dir(dir)
.into_iter()
.flatten()
.flatten()
.filter_map(|e| e.metadata().ok())
.map(|m| m.len())
.sum()
}
const MIN_FREE: f64 = 0.10;
fn reclaim(dir: &Path, min_free: f64) -> mira_core::error::Result<Vec<PathBuf>> {
let mut dropped = Vec::new();
let mut free = block::free_fraction(dir)?;
if free >= min_free {
return Ok(dropped);
}
let mut blocks = Vec::new();
for s in SIGNALS {
blocks.extend(block::scan(dir, s)?);
}
blocks.sort_by_key(|b| (b.max_ts, b.seq));
for b in blocks {
if free >= min_free {
break;
}
match std::fs::remove_dir_all(&b.dir) {
Ok(()) => {
tracing::warn!(
block = %b.dir.display(),
free = format!("{free:.3}"),
"volume is nearly full; dropped a block that had not reached its retention"
);
dropped.push(b.dir);
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => tracing::warn!(
block = %b.dir.display(),
error = %e,
"cannot drop block to reclaim space; skipping it",
),
}
free = block::free_fraction(dir)?;
}
if !dropped.is_empty() {
tracing::warn!(
blocks = dropped.len(),
free = format!("{free:.3}"),
"dropped blocks ahead of their retention to keep the volume writable; \
retention is longer than this disk can hold at the current ingest rate"
);
}
Ok(dropped)
}
async fn retention(cfg: Arc<Config>) {
let mut tick = tokio::time::interval(Duration::from_secs(60));
loop {
tick.tick().await;
let dir = cfg.data_dir.clone();
let ttl = cfg.retention;
let node = cfg.node;
let swept = tokio::task::spawn_blocking(move || {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos() as i64;
let cutoff = now.saturating_sub(i64::try_from(ttl.as_nanos()).unwrap_or(i64::MAX));
let results = SIGNALS.map(|s| {
let dropped = block::expire(&dir, s, cutoff);
let cold = block::compact(&dir, s, node, now - block::COLD_AFTER_NS);
(s, dropped, cold)
});
(results, reclaim(&dir, MIN_FREE))
})
.await;
match swept {
Ok((results, reclaimed)) => {
if let Err(e) = reclaimed {
tracing::warn!(error = %e, "cannot read free space; retention is TTL-only this sweep");
}
for (signal, dropped, cold) in results {
match dropped {
Ok(0) => {}
Ok(n) => tracing::info!(signal, blocks = n, "retention dropped blocks"),
Err(e) => tracing::warn!(signal, error = %e, "retention failed"),
}
match cold {
Ok(0) => {}
Ok(n) => tracing::info!(signal, blocks = n, "compacted blocks to zstd"),
Err(e) => tracing::warn!(signal, error = %e, "compaction failed"),
}
}
}
Err(e) => tracing::warn!(error = %e, "retention task panicked"),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use mira_core::logs::LogsBuilder;
use mira_proto::collector::logs::v1::ExportLogsServiceRequest;
use mira_proto::common::v1::{AnyValue, KeyValue, any_value};
use mira_proto::logs::v1::{LogRecord, ResourceLogs, ScopeLogs};
fn cfg(name: &str) -> (Arc<Config>, PathBuf) {
let dir = std::env::temp_dir().join(format!("mira-pipe-{name}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
(
Arc::new(Config {
data_dir: dir.clone(),
max_block_age: Duration::from_millis(50),
..Default::default()
}),
dir,
)
}
fn blocks(dir: &std::path::Path) -> usize {
block::scan(dir, "logs").map_or(0, |b| b.len())
}
fn wide(n: usize) -> ExportLogsServiceRequest {
ExportLogsServiceRequest {
resource_logs: vec![ResourceLogs {
scope_logs: vec![ScopeLogs {
log_records: vec![LogRecord {
time_unix_nano: 1_000,
attributes: (0..n)
.map(|i| KeyValue {
key: format!("k{i}"),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue("v".into())),
}),
})
.collect(),
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
}
}
#[tokio::test]
async fn a_request_that_does_not_fit_lands_in_the_next_block() {
let (c, dir) = cfg("carry");
let (tx, _open, h) = spawn::<LogsBuilder>(c);
let (a, b) = tokio::join!(tx.submit(wide(40_000)), tx.submit(wide(40_000)));
assert!(a.is_ok() && b.is_ok(), "both callers must be acknowledged");
drop(tx);
h.await.unwrap();
assert_eq!(blocks(&dir), 2, "the deferred request got its own block");
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn an_impossible_request_fails_only_itself() {
let (c, dir) = cfg("toowide");
let (tx, _open, h) = spawn::<LogsBuilder>(c);
let before = tx.rejects.failed.load(Relaxed);
let destroyed = tx.rejects.refused.load(Relaxed);
let refusal = tx.submit(wide(70_000)).await;
assert!(
matches!(&refusal, Err(Rejected::Failed(e))
if e.contains("65535") || e.contains("dictionary")),
"70k distinct keys cannot fit a u16 dictionary, and the refusal has \
to be the permanent one that names why"
);
assert!(tx.rejects.failed.load(Relaxed) > before);
assert!(tx.rejects.refused.load(Relaxed) > destroyed);
tx.submit(crate::e2e::logs_export("checkout", 2_000, 4))
.await
.unwrap_or_else(|_| panic!("the pipeline is still open for business"));
drop(tx);
h.await.unwrap();
assert_eq!(blocks(&dir), 1, "only the good export was published");
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn an_empty_export_is_acknowledged_without_a_block() {
let (c, dir) = cfg("empty");
let (tx, _open, h) = spawn::<LogsBuilder>(c);
tx.submit(ExportLogsServiceRequest::default())
.await
.unwrap_or_else(|_| panic!("an empty export is not an error"));
drop(tx);
h.await.unwrap();
assert_eq!(blocks(&dir), 0, "nothing to seal, so nothing was sealed");
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_logged_export_is_acknowledged_before_its_block_is_sealed() {
let dir = std::env::temp_dir().join(format!("mira-pipe-wal-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let node = block::node_id("waltest");
let wal = Arc::new(Wal::open(&dir, node).unwrap());
let c = Arc::new(Config {
data_dir: dir.clone(),
node,
max_block_age: Duration::from_secs(1),
wal: Some(Arc::clone(&wal)),
..Default::default()
});
let (tx, _open, h) = spawn::<LogsBuilder>(c);
let started = std::time::Instant::now();
tx.submit(crate::e2e::logs_export("checkout", 2_000, 4))
.await
.unwrap_or_else(|_| panic!("the log accepted it"));
let acked = started.elapsed();
assert!(
acked < Duration::from_millis(100),
"acknowledged in {acked:?}, which is the block age, not the log"
);
assert_eq!(blocks(&dir), 0, "the ack did not wait for a block");
assert_eq!(wal.next_seq(), 1, "the export is a frame");
drop(tx);
h.await.unwrap();
let published = block::scan(&dir, "logs").unwrap();
assert_eq!(published.len(), 1);
assert_eq!(published[0].wal_hi, 1);
assert_eq!(block::wal_watermarks(&dir).unwrap(), [1, 0, 0]);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_replayed_frame_is_claimed_by_the_block_that_finally_stores_it() {
let dir = std::env::temp_dir().join(format!("mira-pipe-replay-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let node = block::node_id("replaytest");
let wal = Wal::open(&dir, node).unwrap();
let body = {
use prost::Message as _;
crate::e2e::logs_export("checkout", 2_000, 4).encode_to_vec()
};
wal.append(wal::Signal::Logs, &body).unwrap();
wal.append(wal::Signal::Logs, &body).unwrap();
drop(wal);
let wal = Arc::new(Wal::open(&dir, node).unwrap());
let c = Arc::new(Config {
data_dir: dir.clone(),
node,
max_block_age: Duration::from_millis(50),
wal: Some(Arc::clone(&wal)),
..Default::default()
});
let (tx, _open, h) = spawn::<LogsBuilder>(c);
let replayed = {
let tx = tx.clone();
let dir = dir.clone();
tokio::task::spawn_blocking(move || {
Wal::replay(&dir, node, [0, 0, 0], |_, seq, body| {
assert!(tx.replay(body, seq).is_ok(), "the flusher took it");
Ok(())
})
.unwrap()
})
.await
.unwrap()
};
assert_eq!(replayed.replayed, 2);
drop(tx);
h.await.unwrap();
assert_eq!(block::wal_watermarks(&dir).unwrap(), [2, 0, 0]);
let mut handed_back = 0;
let again = Wal::replay(
&dir,
node,
block::wal_watermarks(&dir).unwrap(),
|_, _, _| {
handed_back += 1;
Ok(())
},
)
.unwrap();
assert_eq!(
(again.replayed, again.skipped, handed_back),
(0, 2, 0),
"a frame a block already claims must never be replayed again"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_block_that_cannot_be_published_is_a_retryable_answer() {
let (c, dir) = cfg("unpublishable");
let (tx, _open, h) = spawn::<LogsBuilder>(c);
std::fs::write(dir.join(".tmp"), b"not a directory").unwrap();
let answer = tx.submit(wide(1)).await;
assert!(
matches!(&answer, Err(Rejected::Unavailable(why)) if !why.is_empty()),
"a publish that failed has to be answered retryably, and with a reason: \
`Failed` would have the exporter drop the batch, and an empty string \
leaves the operator reading the sender's log for a disk fault"
);
drop(tx);
h.await.unwrap();
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn an_export_too_large_for_a_frame_is_refused_permanently() {
let dir = std::env::temp_dir().join(format!("mira-pipe-huge-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let node = block::node_id("hugetest");
let c = Arc::new(Config {
data_dir: dir.clone(),
node,
wal: Some(Arc::new(Wal::open(&dir, node).unwrap())),
..Default::default()
});
let (tx, _open, h) = spawn::<LogsBuilder>(c);
let mut req = wide(1);
req.resource_logs[0].scope_logs[0].log_records[0].body = Some(AnyValue {
value: Some(any_value::Value::StringValue("x".repeat(64 << 20))),
});
let answer = tx.submit(req).await;
assert!(
matches!(&answer, Err(Rejected::Failed(why)) if why.contains("frame")),
"an export that can never be framed has to be refused permanently and \
named as a framing limit; retrying it burns the link on the same bytes"
);
drop(tx);
h.await.unwrap();
assert_eq!(blocks(&dir), 0, "nothing was framed, so nothing was stored");
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_wal_sweep_syncs_every_tick_and_only_truncates_on_the_slow_one() {
let dir = std::env::temp_dir().join(format!("mira-pipe-sweep-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let node = block::node_id("sweeptest");
let wal = Arc::new(Wal::open(&dir, node).unwrap());
wal.append(wal::Signal::Logs, b"a frame").unwrap();
wal_sweep(Arc::clone(&wal), dir.clone(), false).await;
wal_sweep(Arc::clone(&wal), dir.clone(), true).await;
assert_eq!(wal.next_seq(), 1, "a sweep renumbers nothing");
assert_eq!(
std::fs::read_dir(dir.join(".wal")).unwrap().count(),
1,
"the open segment is never dropped"
);
let bad = dir.join("unreadable");
std::fs::create_dir_all(&bad).unwrap();
std::fs::write(bad.join("logs"), b"not a directory").unwrap();
wal_sweep(Arc::clone(&wal), bad, true).await;
wal.append(wal::Signal::Logs, b"another").unwrap();
assert_eq!(wal.next_seq(), 2);
let stale = dir.join(".wal").join(format!("{node:08x}-{:020}.wal", 9));
std::fs::File::create(&stale).unwrap();
wal_sweep(Arc::clone(&wal), dir.clone(), false).await;
assert!(stale.exists(), "a sync is not a truncation");
wal_sweep(Arc::clone(&wal), dir.clone(), true).await;
assert!(!stale.exists(), "the slow tick removed the dead segment");
assert_eq!(
std::fs::read_dir(dir.join(".wal")).unwrap().count(),
1,
"and left the open one, which is still holding two unclaimed frames"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test(start_paused = true)]
async fn a_full_queue_sheds_and_a_closed_one_says_so() {
let (tx, rx) = mpsc::channel::<Job<ExportLogsServiceRequest>>(1);
let rejects = &REJECTS[0];
let ingest = Ingest {
tx,
rejects,
wal: None,
signal: wal::Signal::Logs,
};
let req = || ExportLogsServiceRequest::default();
let before = rejects.shed.load(Relaxed);
let pending = tokio::spawn({
let i = ingest.clone();
async move { i.submit(req()).await }
});
while rx.capacity() > 0 {
tokio::task::yield_now().await;
}
assert!(matches!(ingest.submit(req()).await, Err(Rejected::Busy)));
assert_eq!(rejects.shed.load(Relaxed), before + 1);
drop(rx);
assert!(matches!(pending.await.unwrap(), Err(Rejected::Closed)));
assert!(matches!(ingest.submit(req()).await, Err(Rejected::Closed)));
}
#[tokio::test(start_paused = true)]
async fn a_queue_that_drains_inside_the_wait_admits_instead_of_shedding() {
let (tx, mut rx) = mpsc::channel::<Job<ExportLogsServiceRequest>>(1);
let rejects: &'static Rejects = Box::leak(Box::new(Rejects::new("logs")));
let ingest = Ingest {
tx,
rejects,
wal: None,
signal: wal::Signal::Logs,
};
let req = || ExportLogsServiceRequest::default();
let first = tokio::spawn({
let i = ingest.clone();
async move { i.submit(req()).await }
});
while rx.capacity() > 0 {
tokio::task::yield_now().await;
}
let waiter = tokio::spawn({
let i = ingest.clone();
async move { i.submit(req()).await }
});
tokio::time::sleep(Duration::from_secs(4)).await;
let job = rx.recv().await.expect("the filler's job");
let _ = job.ack.send(Ok(()));
assert!(matches!(first.await.unwrap(), Ok(())));
let job = rx.recv().await.expect("the waiter's job");
let _ = job.ack.send(Ok(()));
assert!(matches!(waiter.await.unwrap(), Ok(())));
assert_eq!(rejects.shed.load(Relaxed), 0, "nothing was shed");
}
#[tokio::test]
async fn an_unreadable_data_directory_stops_the_flusher_at_startup() {
let (c, dir) = cfg("unscannable");
std::fs::write(dir.join("logs"), b"not a directory").unwrap();
let (tx, _open, h) = spawn::<LogsBuilder>(c);
h.await.unwrap();
assert!(matches!(
tx.submit(ExportLogsServiceRequest::default()).await,
Err(Rejected::Closed)
));
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn retention_drops_expired_blocks_on_its_first_pass() {
let (c, dir) = cfg("retention");
let (tx, _open, h) = spawn::<LogsBuilder>(Arc::clone(&c));
tx.submit(crate::e2e::logs_export("checkout", 1_000, 4))
.await
.unwrap_or_else(|_| panic!("export"));
drop(tx);
h.await.unwrap();
assert_eq!(blocks(&dir), 1);
spawn_retention(Arc::new(Config {
data_dir: dir.clone(),
retention: Duration::ZERO,
..Default::default()
}));
for _ in 0..200 {
if blocks(&dir) == 0 {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert_eq!(blocks(&dir), 0, "a block older than its TTL is unlinked");
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_full_volume_drops_the_oldest_blocks_before_their_ttl() {
let (c, dir) = cfg("space");
let (tx, _open, h) = spawn::<LogsBuilder>(Arc::clone(&c));
for ts in [3_000_000, 1_000_000, 2_000_000] {
tx.submit(crate::e2e::logs_export("checkout", ts, 4))
.await
.unwrap_or_else(|_| panic!("export"));
}
drop(tx);
h.await.unwrap();
assert_eq!(blocks(&dir), 3);
let logs = rejects_for("logs");
assert!(logs.published.load(Relaxed) >= 3, "blocks are counted");
assert!(logs.rows.load(Relaxed) >= 12, "rows are counted");
assert!(logs.bytes.load(Relaxed) > 0, "bytes on disk are counted");
assert!(reclaim(&dir, 0.0).unwrap().is_empty());
let mut want = block::scan(&dir, "logs").unwrap();
want.sort_by_key(|b| b.max_ts);
let want: Vec<PathBuf> = want.into_iter().map(|b| b.dir).collect();
assert_eq!(reclaim(&dir, 2.0).unwrap(), want);
assert_eq!(blocks(&dir), 0);
let _ = std::fs::remove_dir_all(&dir);
}
fn listening<T>(f: impl FnOnce() -> T) -> T {
let sub = tracing_subscriber::fmt()
.with_max_level(tracing::Level::TRACE)
.with_test_writer()
.finish();
tracing::subscriber::with_default(sub, f)
}
async fn until(mut done: impl FnMut() -> bool) -> bool {
for _ in 0..200 {
if done() {
return true;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
done()
}
fn fake_block(dir: &Path, signal: &str, max_ts: i64, seq: u64) -> PathBuf {
let partition = dir.join(signal).join("p=1970-01-01-00");
std::fs::create_dir_all(&partition).unwrap();
let block = partition.join(format!(
"{:020}-{max_ts:020}-{:08x}-{seq:012}-{:020}",
0, 7, 0
));
std::fs::create_dir_all(&block).unwrap();
block
}
#[test]
fn every_shed_export_is_counted_and_at_most_one_a_second_is_logged() {
let gate = AtomicU64::new(0);
assert!(once_a_second(&gate), "the first caller in a second speaks");
assert!(!once_a_second(&gate), "and everyone behind it is silent");
let r = Rejects::new("logs");
listening(|| {
for _ in 0..3 {
r.record_shed();
}
});
assert_eq!(r.shed.load(Relaxed), 3, "every shed export is counted");
assert_ne!(
r.warned.load(Relaxed),
0,
"the gate is armed, so the next thousand this second are silent"
);
}
#[tokio::test]
async fn a_log_failure_that_is_not_the_senders_fault_is_answered_retryably() {
let dir = std::env::temp_dir().join(format!("mira-pipe-walgone-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let node = block::node_id("walgone");
let wal = Arc::new(Wal::open(&dir, node).unwrap());
wal.append(wal::Signal::Logs, &vec![0u8; 64 << 20]).unwrap();
let c = Arc::new(Config {
data_dir: dir.clone(),
node,
wal: Some(Arc::clone(&wal)),
..Default::default()
});
let (tx, _open, h) = spawn::<LogsBuilder>(c);
std::fs::remove_dir_all(dir.join(".wal")).unwrap();
let answer = tx.submit(wide(1)).await;
assert!(
matches!(&answer, Err(Rejected::Unavailable(why)) if !why.is_empty()),
"a log that cannot write must be retryable and say why"
);
drop(tx);
h.await.unwrap();
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn wal_maintenance_without_a_log_has_nothing_to_do() {
let (c, dir) = cfg("nowal");
assert!(
c.wal.is_none(),
"the shipped default for this test's config"
);
tokio::time::timeout(Duration::from_millis(250), wal_maintenance(c))
.await
.expect("a node with no log has no maintenance loop to run");
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_staging_directory_a_crash_left_behind_is_swept_at_boot() {
let (c, dir) = cfg("staging");
let node = c.node;
let stale = dir
.join(".tmp")
.join(format!("logs-{node:08x}-000000000007"));
std::fs::create_dir_all(&stale).unwrap();
std::fs::write(stale.join("logs.arrow"), b"half a block").unwrap();
let theirs = dir
.join(".tmp")
.join(format!("logs-{:08x}-000000000007", 0));
std::fs::create_dir_all(&theirs).unwrap();
let (tx, _open, h) = spawn::<LogsBuilder>(c);
tx.submit(ExportLogsServiceRequest::default())
.await
.unwrap_or_else(|_| panic!("the flusher booted"));
drop(tx);
h.await.unwrap();
assert!(!stale.exists(), "the leaked staging directory is gone");
assert!(
theirs.exists(),
"another replica's is not this node's to take"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[derive(Default)]
struct Brittle<const SEAL: bool>(LogsBuilder);
impl<const SEAL: bool> SignalBuilder for Brittle<SEAL> {
type Request = ExportLogsServiceRequest;
const SIGNAL: &'static str = "logs";
fn has_headroom_for(&self, req: &Self::Request) -> bool {
self.0.has_headroom_for(req)
}
fn append_request(&mut self, req: &Self::Request) -> mira_core::error::Result<usize> {
self.0.append_request(req)
}
fn approx_bytes(&self) -> usize {
self.0.approx_bytes()
}
fn is_empty(&self) -> bool {
self.0.is_empty()
}
fn finish(&mut self) -> mira_core::error::Result<mira_core::signal::Sealed> {
if SEAL {
let _ = self.0.finish();
return Err(mira_core::Error::DictionaryFull("attr_key"));
}
self.0.finish()
}
fn snapshot(&self) -> mira_core::error::Result<mira_core::signal::Sealed> {
if SEAL {
return self.0.snapshot();
}
Err(mira_core::Error::DictionaryFull("attr_key"))
}
}
#[tokio::test]
async fn a_block_that_cannot_be_sealed_nacks_retryably_and_leaves_its_frames_in_the_log() {
let (c, dir) = cfg("brittle-seal");
let (tx, _open, h) = spawn::<Brittle<true>>(Arc::clone(&c));
let answer = tx.submit(wide(1)).await;
assert!(
matches!(&answer, Err(Rejected::Unavailable(why)) if why.contains("dictionary")),
"a seal that failed is the block's fault, not this caller's"
);
drop(tx);
h.await.unwrap();
assert_eq!(blocks(&dir), 0, "nothing was published");
let _ = std::fs::remove_dir_all(&dir);
let dir = std::env::temp_dir().join(format!("mira-pipe-brittle-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let node = block::node_id("brittletest");
let wal = Arc::new(Wal::open(&dir, node).unwrap());
let c = Arc::new(Config {
data_dir: dir.clone(),
node,
max_block_age: Duration::from_millis(50),
wal: Some(Arc::clone(&wal)),
..Default::default()
});
let (tx, open, h) = spawn::<Brittle<true>>(c);
tx.submit(wide(1))
.await
.unwrap_or_else(|_| panic!("the log took it, whatever the block does later"));
drop(tx);
h.await.unwrap();
assert_eq!(
block::wal_watermarks(&dir).unwrap(),
[0, 0, 0],
"a block that was never published claims no sequence"
);
assert!(
open.get().is_none(),
"and advertises no rows the read path could no longer produce"
);
let replayed = Wal::replay(
&dir,
node,
block::wal_watermarks(&dir).unwrap(),
|_, _, _| Ok(()),
)
.unwrap();
assert_eq!(
(replayed.replayed, replayed.skipped),
(1, 0),
"the acknowledged export survived the block that could not hold it"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn an_open_block_that_cannot_be_snapshotted_shows_nothing_rather_than_stale_rows() {
for (breaks, want) in [(true, false), (false, true)] {
let dir =
std::env::temp_dir().join(format!("mira-pipe-snap{breaks}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let node = block::node_id("snaptest");
let c = Arc::new(Config {
data_dir: dir.clone(),
node,
max_block_age: Duration::from_secs(30),
wal: Some(Arc::new(Wal::open(&dir, node).unwrap())),
..Default::default()
});
let (tx, open, h) = if breaks {
spawn::<Brittle<false>>(c)
} else {
spawn::<Brittle<true>>(c)
};
tx.submit(crate::e2e::logs_export("checkout", 2_000, 4))
.await
.unwrap_or_else(|_| panic!("acknowledged by the log"));
assert_eq!(
open.fresh().await.is_some(),
want,
"breaks={breaks}: a snapshot that failed must clear the slot"
);
drop(tx);
h.await.unwrap();
let _ = std::fs::remove_dir_all(&dir);
}
}
#[test]
fn reclaim_stops_as_soon_as_the_volume_is_back_over_the_floor() {
let dir = std::env::temp_dir().join(format!("mira-pipe-ballast-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let oldest = fake_block(&dir, "logs", 1_000, 0);
let newer = fake_block(&dir, "logs", 2_000, 1);
let newest = fake_block(&dir, "traces", 3_000, 2);
let empty = block::free_fraction(&dir).unwrap();
{
use std::io::Write;
for i in 0..128 {
let mut f = std::fs::File::create(oldest.join(format!("{i}.arrow"))).unwrap();
f.write_all(&vec![0u8; 1 << 20]).unwrap();
f.sync_all().unwrap();
}
}
let full = block::free_fraction(&dir).unwrap();
assert!(full < empty, "128 MiB moved the needle: {full} vs {empty}");
let floor = (full + empty) / 2.0;
let dropped = listening(|| reclaim(&dir, floor).unwrap());
assert_eq!(
dropped,
vec![oldest],
"the oldest block, and then it stopped"
);
assert!(
newer.exists() && newest.exists(),
"nothing else was touched"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_block_that_cannot_be_dropped_does_not_stop_the_ones_behind_it() {
let dir = std::env::temp_dir().join(format!("mira-pipe-undrop-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let partition = dir.join("logs").join("p=1970-01-01-00");
std::fs::create_dir_all(&partition).unwrap();
let impostor = partition.join(format!("{:020}-{:020}-{:08x}-{:012}-{:020}", 0, 1, 7, 0, 0));
std::fs::write(&impostor, b"not a block").unwrap();
let dropped = listening(|| reclaim(&dir, 2.0).unwrap());
assert!(
dropped.is_empty() && impostor.exists(),
"nothing was dropped, and the sweep still returned"
);
let real = fake_block(&dir, "traces", 5_000, 3);
std::os::unix::fs::symlink(dir.join("traces"), dir.join("metrics")).unwrap();
let dropped = listening(|| reclaim(&dir, 2.0).unwrap());
assert_eq!(
dropped,
vec![real.clone()],
"the block is reported once, and the second sighting is not an error"
);
assert!(!real.exists());
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_sweep_that_fails_never_takes_the_retention_loop_with_it() {
let dir = std::env::temp_dir().join(format!("mira-pipe-sweepfail-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("logs"), b"not a directory").unwrap();
let doomed = fake_block(&dir, "traces", 1_000, 0);
let sweeping = tokio::spawn(retention(Arc::new(Config {
data_dir: dir.clone(),
retention: Duration::ZERO,
..Default::default()
})));
assert!(
until(|| !doomed.exists()).await,
"the signal that could be swept was swept, whatever the broken one did"
);
assert!(!sweeping.is_finished(), "and the loop kept its next tick");
sweeping.abort();
let sweeping = tokio::spawn(retention(Arc::new(Config {
data_dir: dir.join("never-created"),
retention: Duration::ZERO,
..Default::default()
})));
tokio::time::sleep(Duration::from_millis(200)).await;
assert!(
!sweeping.is_finished(),
"a volume it cannot even measure is not a reason to stop measuring it"
);
sweeping.abort();
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn an_absurd_retention_keeps_every_block_and_still_compacts_the_cold_ones() {
let (c, dir) = cfg("forever");
let (tx, _open, h) = spawn::<LogsBuilder>(Arc::clone(&c));
tx.submit(crate::e2e::logs_export("checkout", 1_000, 4))
.await
.unwrap_or_else(|_| panic!("export"));
drop(tx);
h.await.unwrap();
assert_eq!(blocks(&dir), 1);
let cfg = Arc::new(Config {
data_dir: dir.clone(),
retention: Duration::from_secs(200_000 * 86_400),
..Default::default()
});
let sweeping = tokio::spawn(retention(cfg));
let cold = block::scan(&dir, "logs").unwrap()[0].dir.join("cold");
assert!(
until(|| cold.exists()).await,
"the block was compacted rather than deleted"
);
assert!(
!sweeping.is_finished(),
"one sweep, and the loop is still there"
);
sweeping.abort();
assert_eq!(
blocks(&dir),
1,
"a retention longer than i64 nanoseconds keeps everything"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_stall_is_only_reportable_once_it_has_outlasted_the_recovery() {
let r = Rejects::new("logs");
let now = 1_700_000_000;
assert_eq!(stall_of(&r, now), None, "a healthy signal is never unready");
r.stalled_since.store(now, Relaxed);
let after = UNREADY_AFTER.as_secs();
assert_eq!(
stall_of(&r, now),
None,
"one failed publish is not an outage"
);
assert_eq!(stall_of(&r, now + after - 1), None);
assert_eq!(stall_of(&r, now + after), Some(after));
r.stalled_since.store(1, Relaxed);
r.mark_stalled();
assert_eq!(r.stalled_since.load(Relaxed), 1);
let fresh = Rejects::new("logs");
assert_eq!(fresh.stalled_since.load(Relaxed), 0);
fresh.mark_stalled();
assert_ne!(fresh.stalled_since.load(Relaxed), 0);
assert_eq!(stalled(), None);
}
}