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::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 offload: Option<mira_core::offload::Target>,
pub queue: usize,
pub shards: 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),
offload: None,
queue: 128,
shards: 1,
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: Arc<[mpsc::Sender<Job<R>>]>,
pub(crate) turn: Arc<AtomicU64>,
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(),
turn: self.turn.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,
shard_open_since: [AtomicU64; MAX_SHARDS],
shard_stalled_since: [AtomicU64; MAX_SHARDS],
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),
shard_open_since: [const { AtomicU64::new(0) }; MAX_SHARDS],
shard_stalled_since: [const { AtomicU64::new(0) }; MAX_SHARDS],
warned: AtomicU64::new(0),
refuse_warned: AtomicU64::new(0),
}
}
fn oldest(slots: &[AtomicU64; MAX_SHARDS]) -> u64 {
slots
.iter()
.map(|t| t.load(Relaxed))
.filter(|&t| t != 0)
.min()
.unwrap_or(0)
}
fn set_open_since(&self, shard: usize, at: u64) {
self.shard_open_since[shard].store(at, Relaxed);
self.open_since
.store(Self::oldest(&self.shard_open_since), Relaxed);
}
fn mark_stalled(&self, shard: usize) {
let _ = self.shard_stalled_since[shard].compare_exchange(
0,
now_secs().max(1),
Relaxed,
Relaxed,
);
self.stalled_since
.store(Self::oldest(&self.shard_stalled_since), Relaxed);
}
#[cfg(test)]
pub fn forget_open(&self) {
for t in &self.shard_open_since {
t.store(0, Relaxed);
}
self.open_since.store(0, Relaxed);
}
fn clear_stalled(&self, shard: usize) {
self.shard_stalled_since[shard].store(0, Relaxed);
self.stalled_since
.store(Self::oldest(&self.shard_stalled_since), 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> Ingest<R> {
fn reserve(&self) -> Option<mpsc::Permit<'_, Job<R>>> {
self.tx.iter().find_map(|tx| tx.try_reserve().ok())
}
}
impl<R: prost::Message> Ingest<R> {
pub async fn submit(&self, req: R) -> Result<(), Rejected> {
let _t_submit = mira_core::diag::Scope::new(&mira_core::diag::SUBMIT_TOTAL);
let t_admit = std::time::Instant::now();
let (ack, wait) = oneshot::channel();
let permit = match self.reserve() {
Some(p) => p,
None => {
let i = self.turn.fetch_add(1, Relaxed) as usize % self.tx.len();
match tokio::time::timeout(ADMIT_WAIT, self.tx[i].reserve()).await {
Ok(Ok(p)) => p,
Ok(Err(_)) => return Err(Rejected::Closed),
Err(_) => {
self.rejects.record_shed();
return Err(Rejected::Busy);
}
}
}
};
mira_core::diag::SUBMIT_ADMIT.record(t_admit.elapsed().as_nanos() as u64);
if let Some(wal) = &self.wal {
let t_enc = std::time::Instant::now();
let body = req.encode_to_vec();
mira_core::diag::WAL_ENCODE.record(t_enc.elapsed().as_nanos() as u64);
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,
});
let _t_ack = mira_core::diag::Scope::new(&mira_core::diag::SUBMIT_WAIT_ACK);
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> {
if let Some(w) = &self.wal {
w.reframed(self.signal, seq);
}
let req = match R::decode(body) {
Ok(r) => r,
Err(e) => {
if let Some(w) = &self.wal {
w.published(self.signal, &[seq]);
}
return Err(Rejected::Failed(e.to_string()));
}
};
let job = Job {
req,
ack: oneshot::channel().0,
wal_seq: Some(seq),
};
match self.reserve() {
Some(p) => {
p.send(job);
Ok(())
}
None => self.tx[0].blocking_send(job).map_err(|_| Rejected::Closed),
}
}
}
pub const SIGNALS: [&str; 3] = ["logs", "traces", "metrics"];
pub const MAX_SHARDS: usize = 16;
pub fn shard_count(configured: usize, cores: usize) -> usize {
match configured {
0 => (cores / 2).clamp(1, MAX_SHARDS),
n => n.min(MAX_SHARDS),
}
}
pub struct Flushers(tokio::task::JoinSet<()>);
impl Drop for Flushers {
fn drop(&mut self) {
self.0.detach_all();
}
}
impl Flushers {
#[cfg(test)]
pub fn abort(&mut self) {
self.0.abort_all();
}
#[cfg(test)]
pub fn wedged() -> Self {
let mut set = tokio::task::JoinSet::new();
set.spawn(std::future::pending());
Self(set)
}
}
impl std::future::Future for Flushers {
type Output = Result<(), tokio::task::JoinError>;
fn poll(
mut self: std::pin::Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Self::Output> {
use std::task::Poll;
loop {
match self.0.poll_join_next(cx) {
Poll::Ready(Some(Ok(()))) => {}
Poll::Ready(Some(Err(e))) => return Poll::Ready(Err(e)),
Poll::Ready(None) => return Poll::Ready(Ok(())),
Poll::Pending => return Poll::Pending,
}
}
}
}
pub fn spawn<B: SignalBuilder>(cfg: &Arc<Config>) -> (Ingest<B::Request>, OpenSlot, Flushers) {
let shards = cfg.shards.clamp(1, MAX_SHARDS);
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 resume = match block::scan(&cfg.data_dir, B::SIGNAL) {
Ok(blocks) => Some(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");
None
}
};
let rejects = rejects_for(B::SIGNAL);
let depth = cfg.queue.div_ceil(shards).max(1);
let mut txs = Vec::with_capacity(shards);
let mut slots = Vec::with_capacity(shards);
let mut tasks = tokio::task::JoinSet::new();
for shard in 0..shards {
let (tx, rx) = mpsc::channel(depth);
let (ask, asks) = mpsc::channel(8);
let slot = Shard {
cur: Arc::default(),
ask,
};
txs.push(tx);
if let Some(resume) = resume {
tasks.spawn(flusher::<B>(
rx,
asks,
cfg.clone(),
slot.clone(),
shard,
shards,
resume,
));
}
slots.push(slot);
}
let ingest = Ingest {
tx: txs.into(),
turn: Arc::default(),
rejects,
wal: cfg.wal.clone(),
signal: wal::Signal::named(B::SIGNAL).expect("every signal has a log discriminant"),
};
(
ingest,
OpenSlot {
shards: slots.into(),
},
Flushers(tasks),
)
}
#[derive(Clone)]
struct Shard {
cur: Arc<Mutex<Option<Arc<Open>>>>,
ask: mpsc::Sender<oneshot::Sender<Option<Arc<Open>>>>,
}
#[derive(Clone, Default)]
pub struct OpenSlot {
shards: Arc<[Shard]>,
}
impl OpenSlot {
pub async fn fresh(&self) -> Vec<Arc<Open>> {
let mut out = Vec::with_capacity(self.shards.len());
let mut waiting = Vec::with_capacity(self.shards.len());
for s in self.shards.iter() {
let (tx, rx) = oneshot::channel();
match s.ask.try_send(tx) {
Ok(()) => waiting.push((s, rx)),
Err(_) => out.extend(s.get()),
}
}
for (s, rx) in waiting {
out.extend(rx.await.unwrap_or_else(|_| s.get()));
}
out
}
}
impl Shard {
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, wal.node())?
.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: Shard,
shard: usize,
shards: usize,
resume: u64,
) {
let rejects = rejects_for(B::SIGNAL);
let mut seq = resume.saturating_add(shard as u64);
let signal = wal::Signal::named(B::SIGNAL).expect("every signal has a log discriminant");
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_seqs: Vec<u64> = Vec::new();
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_seqs.push(seq);
}
match builder.append_request(&job.req) {
Ok(0) => {
let _ = job.ack.send(Ok(()));
}
Ok(_) => {
if waiters.is_empty() {
rejects.set_open_since(shard, now_secs());
}
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.set_open_since(shard, 0);
}
continue;
}
rejects.set_open_since(shard, 0);
let sealed = match builder.finish() {
Ok(s) => s,
Err(e) => {
let msg = e.to_string();
rejects.mark_stalled(shard);
for w in waiters.drain(..) {
let _ = w.send(Err(Rejected::Unavailable(msg.clone())));
}
wal_seqs.clear();
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 += shards as u64;
let block_seqs = std::mem::take(&mut wal_seqs);
let block_wal_hi = match &cfg.wal {
Some(w) => w.watermark_for(signal, &block_seqs),
None => 0,
};
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))) => {
if let Some(w) = &cfg.wal {
w.published(signal, &block_seqs);
}
rejects.published.fetch_add(1, Relaxed);
rejects.rows.fetch_add(rows as u64, Relaxed);
rejects.bytes.fetch_add(bytes, Relaxed);
rejects.clear_stalled(shard);
tracing::info!(signal = B::SIGNAL, rows, bytes, seq = this_seq, path = %path.display(), "block published");
Ok(())
}
Ok(Err(e)) => {
rejects.mark_stalled(shard);
tracing::error!(signal = B::SIGNAL, seq = this_seq, error = %e, "block not published");
Err(e.to_string())
}
Err(e) => {
rejects.mark_stalled(shard);
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: &Shard,
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 offload = cfg.offload.clone();
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 = match &offload {
Some(t) => block::expire_with(&dir, s, cutoff, &|b| t.push(s, b).map(|_| ())),
None => 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)) => report_sweep(results, &reclaimed),
Err(e) => tracing::warn!(error = %e, "retention task panicked"),
}
}
}
fn report_sweep(
results: [(
&str,
mira_core::error::Result<usize>,
mira_core::error::Result<usize>,
); 3],
reclaimed: &mira_core::error::Result<Vec<PathBuf>>,
) {
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"),
}
}
}
#[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 seqs(dir: &std::path::Path) -> Vec<u64> {
let mut v: Vec<u64> = block::scan(dir, "logs")
.unwrap()
.iter()
.map(|b| b.seq)
.collect();
v.sort_unstable();
v
}
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, node).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, node).unwrap(), [2, 0, 0]);
let again = Wal::replay(
&dir,
node,
block::wal_watermarks(&dir, node).unwrap(),
|_, _, _| unreachable!("a frame a block already claims must never be replayed again"),
)
.unwrap();
assert_eq!((again.replayed, again.skipped), (0, 2));
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();
let segments = |dir: &std::path::Path| {
std::fs::read_dir(dir.join(".wal"))
.unwrap()
.filter(|e| {
e.as_ref()
.unwrap()
.path()
.extension()
.is_some_and(|x| x == "wal")
})
.count()
};
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!(segments(&dir), 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!(
segments(&dir),
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: [tx].into(),
turn: Arc::default(),
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: [tx].into(),
turn: Arc::default(),
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>(&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()
}));
until(|| blocks(&dir) == 0).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 retention_with_offload_puts_the_block_in_the_store_before_unlinking_it() {
let (c, dir) = cfg("retention-offload");
let store = dir.join("cold");
let (tx, _open, h) = spawn::<LogsBuilder>(&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 target = mira_core::offload::Target::parse(&format!("file://{}", store.display()))
.expect("file:// target");
spawn_retention(Arc::new(Config {
data_dir: dir.clone(),
retention: Duration::ZERO,
offload: Some(target.clone()),
..Default::default()
}));
for _ in 0..200 {
if blocks(&dir) == 0 {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert_eq!(blocks(&dir), 0, "the local copy still goes");
assert_eq!(
target.list("logs").unwrap().len(),
1,
"and the store has it, listed by name with no index written"
);
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>(&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)
}
#[test]
fn a_sweep_reports_what_it_did_and_a_quiet_one_says_nothing() {
let bad = || mira_core::error::Error::BadMagic {
path: PathBuf::from("/data/logs/0000"),
};
listening(|| {
report_sweep(
[
("logs", Ok(0), Ok(0)),
("traces", Ok(0), Ok(0)),
("metrics", Ok(0), Ok(0)),
],
&Ok(Vec::new()),
);
report_sweep(
[
("logs", Ok(2), Ok(0)),
("traces", Ok(0), Ok(1)),
("metrics", Err(bad()), Err(bad())),
],
&Err(bad()),
);
});
}
async fn until(mut done: impl FnMut() -> bool) -> bool {
for _ in 0..6_000 {
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>>(&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, node).unwrap(),
[0, 0, 0],
"a block that was never published claims no sequence"
);
assert!(
open.fresh().await.is_empty(),
"and advertises no rows the read path could no longer produce"
);
let replayed = Wal::replay(
&dir,
node,
block::wal_watermarks(&dir, node).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_empty(),
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>(&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.shard_stalled_since[0].store(1, Relaxed);
r.mark_stalled(0);
assert_eq!(r.stalled_since.load(Relaxed), 1);
let fresh = Rejects::new("logs");
assert_eq!(fresh.stalled_since.load(Relaxed), 0);
fresh.mark_stalled(0);
assert_ne!(fresh.stalled_since.load(Relaxed), 0);
assert_eq!(stalled(), None);
}
#[test]
fn a_signals_clocks_report_the_worst_shard_not_the_latest_one() {
let r = Rejects::new("logs");
r.set_open_since(0, 100);
r.set_open_since(1, 500);
assert_eq!(r.open_since.load(Relaxed), 100);
r.set_open_since(0, 0);
assert_eq!(r.open_since.load(Relaxed), 500);
r.set_open_since(1, 0);
assert_eq!(
r.open_since.load(Relaxed),
0,
"nothing open is zero, not min"
);
r.shard_stalled_since[1].store(900, Relaxed);
r.mark_stalled(0);
let both = r.stalled_since.load(Relaxed);
assert!(both > 0 && both <= 900, "the older of the two, got {both}");
r.clear_stalled(0);
assert_eq!(r.stalled_since.load(Relaxed), 900);
r.clear_stalled(1);
assert_eq!(r.stalled_since.load(Relaxed), 0);
}
#[test]
fn shards_are_counted_from_cores_and_clamped_at_both_ends() {
assert_eq!(shard_count(0, 1), 1);
assert_eq!(shard_count(0, 2), 1);
assert_eq!(shard_count(0, 12), 6);
assert_eq!(shard_count(0, 128), MAX_SHARDS);
assert_eq!(shard_count(1, 128), 1);
assert_eq!(shard_count(4, 2), 4);
assert_eq!(shard_count(999, 2), MAX_SHARDS);
}
fn sharded(name: &str, shards: usize) -> (Arc<Config>, PathBuf) {
let (c, dir) = cfg(name);
let node = block::node_id(name);
let mut c = Arc::try_unwrap(c).ok().expect("freshly built");
c.node = node;
c.shards = shards;
c.wal = Some(Arc::new(Wal::open(&dir, node).unwrap()));
c.max_block_age = Duration::from_secs(30);
(Arc::new(c), dir)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn shards_partition_the_sequence_space() {
let (c, dir) = sharded("shardseq", 4);
let (tx, _open, h) = spawn::<LogsBuilder>(&c);
let mut sent = Vec::new();
for _ in 0..8 {
let tx = tx.clone();
sent.push(tokio::spawn(async move { tx.submit(wide(40_000)).await }));
}
for s in sent {
s.await
.unwrap()
.unwrap_or_else(|_| panic!("nothing may be shed: the wait is 5s"));
}
drop(tx);
h.await.unwrap();
let seqs = seqs(&dir);
assert_eq!(seqs.len(), 8, "one block per export, none lost");
let mut uniq = seqs.clone();
uniq.dedup();
assert_eq!(uniq, seqs, "two shards reused a sequence: {seqs:?}");
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_trickle_stays_on_one_shard_and_strides_its_sequences() {
let (c, dir) = sharded("shardtrickle", 4);
let (tx, _open, h) = spawn::<LogsBuilder>(&c);
for _ in 0..4 {
tx.submit(crate::e2e::logs_export("checkout", 2_000, 4))
.await
.unwrap_or_else(|_| panic!("acknowledged"));
}
for _ in 0..2 {
tx.submit(wide(40_000))
.await
.unwrap_or_else(|_| panic!("acknowledged"));
}
drop(tx);
h.await.unwrap();
assert_eq!(
seqs(&dir),
vec![0, 4],
"one shard's blocks, striding by the shard count — four shards must \
not mean four files for a load one shard can take"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn spawn_bounds_a_shard_count_that_never_went_through_shard_count() {
for (configured, stride) in [(0, 1), (usize::MAX, MAX_SHARDS)] {
let (c, dir) = sharded(&format!("shardclamp{configured}"), configured);
let (tx, _open, h) = spawn::<LogsBuilder>(&c);
for _ in 0..2 {
tx.submit(wide(40_000))
.await
.unwrap_or_else(|_| panic!("acknowledged"));
}
drop(tx);
h.await.unwrap();
assert_eq!(
seqs(&dir),
vec![0, stride as u64],
"{configured} shards must be clamped to {stride}"
);
let _ = std::fs::remove_dir_all(&dir);
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_crash_after_an_out_of_order_seal_replays_every_frame_no_block_holds() {
let (mut c, dir) = sharded("shardcrash", 2);
Arc::get_mut(&mut c).unwrap().queue = 2;
let node = c.node;
let wal = Arc::clone(c.wal.as_ref().unwrap());
let (tx, _open, mut h) = spawn::<LogsBuilder>(&c);
tx.submit(crate::e2e::logs_export("checkout", 2_000, 4))
.await
.unwrap_or_else(|_| panic!("the log took it"));
let held = tx.tx[0].clone().reserve_owned().await.unwrap();
for _ in 0..2 {
tx.submit(wide(40_000))
.await
.unwrap_or_else(|_| panic!("the log took it"));
}
assert!(
until(|| blocks(&dir) == 1).await,
"shard 1's first block never landed"
);
h.abort();
drop(held);
let published = block::scan(&dir, "logs").unwrap();
assert_eq!(published.len(), 1);
assert_eq!(
published[0].wal_hi, 0,
"a block holding frame 1 may not claim past frame 0, which shard 0 \
still has — this is the assertion `max(seq) + 1` fails"
);
assert_eq!(
block::wal_watermarks(&dir, node).unwrap(),
[0, 0, 0],
"and the reduction over the directory says the same"
);
let mut got = Vec::new();
let replayed = Wal::replay(
&dir,
node,
block::wal_watermarks(&dir, node).unwrap(),
|_, seq, body| {
got.push((seq, body.to_vec()));
Ok(())
},
)
.unwrap();
assert_eq!(
got.iter().map(|(s, _)| *s).collect::<Vec<_>>(),
vec![0, 1, 2],
"every frame, none skipped"
);
assert_eq!((replayed.replayed, replayed.skipped), (3, 0));
let (tx, _open, h) = spawn::<LogsBuilder>(&c);
for (seq, body) in got {
tx.replay(&body, seq)
.unwrap_or_else(|_| panic!("the flusher took frame {seq}"));
}
drop(tx);
h.await.unwrap();
assert_eq!(
wal.watermark_for(wal::Signal::Logs, &[]),
3,
"every frame is published, so nothing is pending"
);
let again = Wal::replay(
&dir,
node,
block::wal_watermarks(&dir, node).unwrap(),
|_, seq, _| unreachable!("frame {seq} is in a block already"),
)
.unwrap();
assert_eq!((again.replayed, again.skipped), (0, 3));
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn every_shard_answers_the_read_path() {
let (mut c, dir) = sharded("shardfresh", 4);
Arc::get_mut(&mut c).unwrap().queue = 4;
let (tx, open, h) = spawn::<LogsBuilder>(&c);
let mut sent = Vec::new();
for i in 0..8 {
let tx = tx.clone();
sent.push(tokio::spawn(async move {
tx.submit(crate::e2e::logs_export("checkout", 2_000 + i * 10, 4))
.await
}));
}
for s in sent {
s.await.unwrap().unwrap_or_else(|_| panic!("acknowledged"));
}
let rows: usize = open.fresh().await.iter().map(|o| o.sealed.num_rows).sum();
assert_eq!(
rows, 32,
"eight exports of four records, all of them findable"
);
drop(tx);
h.await.unwrap();
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_full_shard_spills_into_the_next_one() {
let (tx0, _rx0) = mpsc::channel::<Job<ExportLogsServiceRequest>>(1);
let (tx1, mut rx1) = mpsc::channel::<Job<ExportLogsServiceRequest>>(1);
let _held = tx0.clone().reserve_owned().await.unwrap();
let ingest = Ingest {
tx: [tx0, tx1].into(),
turn: Arc::default(),
rejects: Box::leak(Box::new(Rejects::new("logs"))),
wal: None,
signal: wal::Signal::Logs,
};
let sent = tokio::spawn({
let i = ingest.clone();
async move { i.submit(ExportLogsServiceRequest::default()).await }
});
let job = rx1
.recv()
.await
.expect("shard 1 gets what shard 0 cannot take");
let _ = job.ack.send(Ok(()));
sent.await
.unwrap()
.unwrap_or_else(|_| panic!("admitted, not shed"));
assert_eq!(
ingest.rejects.shed.load(Relaxed),
0,
"spilling to a free shard is not shedding"
);
}
}