use std::collections::HashMap;
use std::path::{Path, PathBuf};
use tokio::io::{AsyncReadExt, AsyncSeekExt};
use super::config::LogLevel;
use super::offsets::OffsetStore;
use super::parse::{parse_line, LogEvent};
use crate::error::Result;
const LOG_FILENAME_PREFIX: &str = "ant-node.";
const LOG_FILENAME_SUFFIX: &str = ".log";
const MAX_CHUNK_BYTES: usize = 1024 * 1024;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TailedEvent {
pub node_id: u32,
pub file_name: String,
pub byte_offset: u64,
pub event: LogEvent,
}
impl TailedEvent {
#[must_use]
pub fn document_id(&self, installation_id: &str) -> String {
format!(
"{}-{}-{}-{}",
installation_id, self.node_id, self.file_name, self.byte_offset
)
}
}
pub struct LogTailer {
node_id: u32,
log_dir: PathBuf,
primed: bool,
last_seen_len: HashMap<String, u64>,
}
impl LogTailer {
#[must_use]
pub fn new(node_id: u32, log_dir: PathBuf) -> Self {
Self {
node_id,
log_dir,
primed: false,
last_seen_len: HashMap::new(),
}
}
#[must_use]
pub fn node_id(&self) -> u32 {
self.node_id
}
#[must_use]
pub fn log_dir(&self) -> &Path {
&self.log_dir
}
pub fn mark_primed(&mut self) {
self.primed = true;
}
pub async fn poll(
&mut self,
offsets: &mut OffsetStore,
min_level: LogLevel,
) -> Result<PollOutcome> {
let files = self.discover_files().await?;
let mut outcome = PollOutcome::default();
for (index, path) in files.iter().enumerate() {
let is_newest = index + 1 == files.len();
match self.poll_file(path, offsets, min_level, is_newest).await {
Ok(mut file_outcome) => {
outcome.events.append(&mut file_outcome.events);
outcome.dropped_by_level += file_outcome.dropped_by_level;
}
Err(error) => {
tracing::debug!(
"log forwarding: skipping {} this poll: {error}",
path.display()
);
}
}
}
let live: Vec<String> = files.iter().map(|p| p.display().to_string()).collect();
self.last_seen_len.retain(|key, _| live.contains(key));
self.primed = true;
outcome.live_files = live;
Ok(outcome)
}
async fn discover_files(&self) -> Result<Vec<PathBuf>> {
let mut entries = match tokio::fs::read_dir(&self.log_dir).await {
Ok(entries) => entries,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(error) => return Err(error.into()),
};
let mut files = Vec::new();
while let Some(entry) = entries.next_entry().await? {
let name = entry.file_name().to_string_lossy().to_string();
if name.starts_with(LOG_FILENAME_PREFIX) && name.ends_with(LOG_FILENAME_SUFFIX) {
files.push(entry.path());
}
}
files.sort();
Ok(files)
}
async fn poll_file(
&mut self,
path: &Path,
offsets: &mut OffsetStore,
min_level: LogLevel,
is_newest: bool,
) -> Result<PollOutcome> {
let key = path.display().to_string();
let file_name = path
.file_name()
.map(|n| n.to_string_lossy().to_string())
.unwrap_or_else(|| key.clone());
let mut file = tokio::fs::File::open(path).await?;
let len = file.metadata().await?.len();
let mut start = match offsets.get(&key) {
Some(offset) => offset,
None if !self.primed || !is_newest => len,
None => 0,
};
if len < start {
tracing::debug!("log forwarding: {key} shrank; restarting from the beginning");
start = 0;
}
let was_growing = self.last_seen_len.insert(key.clone(), len) != Some(len);
if len == start {
offsets.set(&key, start);
return Ok(PollOutcome::default());
}
let to_read = usize::try_from(len - start)
.unwrap_or(MAX_CHUNK_BYTES)
.min(MAX_CHUNK_BYTES);
file.seek(std::io::SeekFrom::Start(start)).await?;
let mut buffer = vec![0u8; to_read];
let read = file.read_exact(&mut buffer).await.map(|_| to_read)?;
buffer.truncate(read);
let reached_eof = start + read as u64 >= len;
let usable = match buffer.iter().rposition(|b| *b == b'\n') {
Some(index) => index + 1,
None if read == MAX_CHUNK_BYTES => read,
None => {
offsets.set(&key, start);
return Ok(PollOutcome::default());
}
};
let text = String::from_utf8_lossy(&buffer[..usable]).to_string();
let outcome = self.collect_events(
&text,
start,
&file_name,
min_level,
reached_eof && !was_growing,
);
offsets.set(&key, outcome.next_offset);
Ok(outcome.into_poll_outcome())
}
fn collect_events(
&self,
text: &str,
chunk_start: u64,
file_name: &str,
min_level: LogLevel,
release_final_event: bool,
) -> CollectOutcome {
let mut outcome = CollectOutcome {
next_offset: chunk_start,
..CollectOutcome::default()
};
let mut pending: Option<TailedEvent> = None;
let mut cursor = chunk_start;
for line in text.split_inclusive('\n') {
let line_start = cursor;
cursor += line.len() as u64;
let content = line.trim_end_matches(['\n', '\r']);
match parse_line(content) {
Some(event) => {
if let Some(previous) = pending.take() {
outcome.push(previous, min_level);
}
outcome.next_offset = line_start;
pending = Some(TailedEvent {
node_id: self.node_id,
file_name: file_name.to_string(),
byte_offset: line_start,
event,
});
}
None => match pending.as_mut() {
Some(held) => held.event.push_continuation(content),
None => outcome.next_offset = cursor,
},
}
}
match pending {
Some(event) if release_final_event => {
outcome.push(event, min_level);
outcome.next_offset = cursor;
}
Some(_) => {}
None => outcome.next_offset = cursor,
}
outcome
}
}
#[derive(Debug, Default)]
pub struct PollOutcome {
pub events: Vec<TailedEvent>,
pub dropped_by_level: u64,
pub live_files: Vec<String>,
}
#[derive(Debug, Default)]
struct CollectOutcome {
events: Vec<TailedEvent>,
dropped_by_level: u64,
next_offset: u64,
}
impl CollectOutcome {
fn push(&mut self, event: TailedEvent, min_level: LogLevel) {
if event.event.level < min_level {
self.dropped_by_level += 1;
} else {
self.events.push(event);
}
}
fn into_poll_outcome(self) -> PollOutcome {
PollOutcome {
events: self.events,
dropped_by_level: self.dropped_by_level,
live_files: Vec::new(),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Write;
fn line(level: &str, message: &str) -> String {
format!("2026-08-19T20:50:00.123456Z {level} ant_node::node: {message}\n")
}
struct Fixture {
_dir: tempfile::TempDir,
log_dir: PathBuf,
offsets_path: PathBuf,
}
impl Fixture {
fn new() -> Self {
let dir = tempfile::tempdir().unwrap();
let log_dir = dir.path().join("logs");
std::fs::create_dir_all(&log_dir).unwrap();
let offsets_path = dir.path().join("offsets.json");
Self {
_dir: dir,
log_dir,
offsets_path,
}
}
fn append(&self, file_name: &str, contents: &str) {
let mut file = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(self.log_dir.join(file_name))
.unwrap();
file.write_all(contents.as_bytes()).unwrap();
}
fn offsets(&self) -> OffsetStore {
OffsetStore::load(&self.offsets_path)
}
fn tailer(&self) -> LogTailer {
LogTailer::new(7, self.log_dir.clone())
}
}
async fn drain(
tailer: &mut LogTailer,
offsets: &mut OffsetStore,
min_level: LogLevel,
) -> PollOutcome {
let mut merged = tailer.poll(offsets, min_level).await.unwrap();
let mut second = tailer.poll(offsets, min_level).await.unwrap();
merged.events.append(&mut second.events);
merged.dropped_by_level += second.dropped_by_level;
merged
}
#[tokio::test]
async fn the_first_poll_joins_existing_files_at_their_end() {
let fixture = Fixture::new();
fixture.append("ant-node.2026-08-19.log", &line("INFO", "historic"));
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
let outcome = tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
assert!(outcome.events.is_empty(), "history must not be shipped");
fixture.append("ant-node.2026-08-19.log", &line("INFO", "fresh"));
let outcome = drain(&mut tailer, &mut offsets, LogLevel::Info).await;
assert_eq!(outcome.events.len(), 1);
assert_eq!(outcome.events[0].event.message, "fresh");
}
#[tokio::test]
async fn a_restart_resumes_from_the_persisted_offset_without_duplicating() {
let fixture = Fixture::new();
fixture.append("ant-node.2026-08-19.log", &line("INFO", "first"));
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
fixture.append("ant-node.2026-08-19.log", &line("INFO", "second"));
let before = drain(&mut tailer, &mut offsets, LogLevel::Info).await;
assert_eq!(before.events.len(), 1);
offsets.save().unwrap();
let mut restarted = fixture.tailer();
restarted.mark_primed();
let mut reloaded = fixture.offsets();
let outcome = drain(&mut restarted, &mut reloaded, LogLevel::Info).await;
assert!(outcome.events.is_empty(), "nothing new, nothing re-sent");
fixture.append("ant-node.2026-08-19.log", &line("INFO", "third"));
let outcome = drain(&mut restarted, &mut reloaded, LogLevel::Info).await;
assert_eq!(outcome.events.len(), 1);
assert_eq!(outcome.events[0].event.message, "third");
}
#[tokio::test]
async fn lines_written_while_the_daemon_was_down_are_not_lost() {
let fixture = Fixture::new();
fixture.append("ant-node.2026-08-19.log", &line("INFO", "before"));
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
offsets.save().unwrap();
fixture.append("ant-node.2026-08-19.log", &line("INFO", "during downtime"));
let mut restarted = fixture.tailer();
restarted.mark_primed();
let mut reloaded = fixture.offsets();
let outcome = drain(&mut restarted, &mut reloaded, LogLevel::Info).await;
assert_eq!(outcome.events.len(), 1);
assert_eq!(outcome.events[0].event.message, "during downtime");
}
#[tokio::test]
async fn a_node_started_after_enabling_is_captured_from_its_first_line() {
let fixture = Fixture::new();
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
let outcome = tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
assert!(outcome.events.is_empty());
fixture.append(
"ant-node.2026-08-19.log",
&format!(
"{}{}",
line("INFO", "starting version=0.17.2 commit=abc1234"),
line("INFO", "listening for connections"),
),
);
let outcome = drain(&mut tailer, &mut offsets, LogLevel::Info).await;
let messages: Vec<&str> = outcome
.events
.iter()
.map(|e| e.event.message.as_str())
.collect();
assert_eq!(
messages,
vec![
"starting version=0.17.2 commit=abc1234",
"listening for connections"
],
"the node's first line must not be skipped"
);
assert_eq!(outcome.events[0].byte_offset, 0);
assert_eq!(outcome.events[0].event.version.as_deref(), Some("0.17.2"));
assert_eq!(outcome.events[0].event.commit.as_deref(), Some("abc1234"));
}
#[tokio::test]
async fn enabling_after_the_node_started_skips_its_startup_line() {
let fixture = Fixture::new();
fixture.append(
"ant-node.2026-08-19.log",
&line("INFO", "starting version=0.17.2 commit=abc1234"),
);
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
let outcome = drain(&mut tailer, &mut offsets, LogLevel::Info).await;
assert!(
outcome.events.is_empty(),
"pre-consent lines are not uploaded"
);
fixture.append("ant-node.2026-08-19.log", &line("INFO", "later activity"));
let outcome = drain(&mut tailer, &mut offsets, LogLevel::Info).await;
let messages: Vec<&str> = outcome
.events
.iter()
.map(|e| e.event.message.as_str())
.collect();
assert_eq!(messages, vec!["later activity"]);
}
#[tokio::test]
async fn a_new_days_file_is_read_from_the_beginning() {
let fixture = Fixture::new();
fixture.append("ant-node.2026-08-19.log", &line("INFO", "yesterday"));
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
fixture.append("ant-node.2026-08-20.log", &line("INFO", "today"));
let outcome = drain(&mut tailer, &mut offsets, LogLevel::Info).await;
assert_eq!(outcome.events.len(), 1);
assert_eq!(outcome.events[0].event.message, "today");
assert_eq!(outcome.events[0].file_name, "ant-node.2026-08-20.log");
}
#[tokio::test]
async fn a_truncated_file_restarts_from_the_beginning() {
let fixture = Fixture::new();
fixture.append("ant-node.2026-08-19.log", &line("INFO", "original content"));
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
std::fs::write(
fixture.log_dir.join("ant-node.2026-08-19.log"),
line("INFO", "new"),
)
.unwrap();
let outcome = drain(&mut tailer, &mut offsets, LogLevel::Info).await;
assert_eq!(outcome.events.len(), 1);
assert_eq!(outcome.events[0].event.message, "new");
}
#[tokio::test]
async fn a_partially_written_line_is_held_until_it_is_complete() {
let fixture = Fixture::new();
fixture.append("ant-node.2026-08-19.log", &line("INFO", "complete"));
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
fixture.append(
"ant-node.2026-08-19.log",
"2026-08-19T20:50:01.000000Z INFO ant_node::node: half a li",
);
let outcome = tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
assert!(outcome.events.is_empty(), "a partial line is not an event");
fixture.append("ant-node.2026-08-19.log", "ne here\n");
tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
let outcome = tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
assert_eq!(outcome.events.len(), 1);
assert_eq!(outcome.events[0].event.message, "half a line here");
}
#[tokio::test]
async fn a_panic_and_its_backtrace_stay_one_event() {
let fixture = Fixture::new();
fixture.append("ant-node.2026-08-19.log", &line("INFO", "before the panic"));
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
fixture.append(
"ant-node.2026-08-19.log",
&format!(
"{}thread 'main' panicked\n at src/node.rs:42\n",
line("ERROR", "it broke")
),
);
tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
let outcome = tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
assert_eq!(outcome.events.len(), 1);
assert_eq!(
outcome.events[0].event.message,
"it broke\nthread 'main' panicked\n at src/node.rs:42"
);
}
#[tokio::test]
async fn events_below_the_minimum_level_are_dropped_and_counted() {
let fixture = Fixture::new();
fixture.append("ant-node.2026-08-19.log", &line("INFO", "seed"));
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
fixture.append(
"ant-node.2026-08-19.log",
&format!(
"{}{}{}",
line("DEBUG", "chatter"),
line("TRACE", "more chatter"),
line("WARN", "worth keeping")
),
);
let outcome = drain(&mut tailer, &mut offsets, LogLevel::Info).await;
let messages: Vec<&str> = outcome
.events
.iter()
.map(|e| e.event.message.as_str())
.collect();
assert_eq!(messages, vec!["worth keeping"]);
assert_eq!(outcome.dropped_by_level, 2);
}
#[tokio::test]
async fn a_missing_log_directory_is_not_an_error() {
let dir = tempfile::tempdir().unwrap();
let mut tailer = LogTailer::new(1, dir.path().join("never-created"));
let mut offsets = OffsetStore::load(&dir.path().join("offsets.json"));
let outcome = tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
assert!(outcome.events.is_empty());
}
#[tokio::test]
async fn unrelated_files_in_the_log_directory_are_ignored() {
let fixture = Fixture::new();
fixture.append("ant-node.2026-08-19.log", &line("INFO", "seed"));
fixture.append("notes.txt", "not a log file\n");
fixture.append("ant-node.2026-08-19.log.gz", "compressed\n");
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
assert_eq!(offsets.len(), 1, "only the rolling log file is tracked");
}
#[tokio::test]
async fn a_retention_deleted_file_drops_out_of_the_reported_live_set() {
let fixture = Fixture::new();
fixture.append("ant-node.2026-08-18.log", &line("INFO", "old"));
fixture.append("ant-node.2026-08-19.log", &line("INFO", "current"));
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
let outcome = tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
assert_eq!(outcome.live_files.len(), 2);
assert_eq!(offsets.len(), 2);
std::fs::remove_file(fixture.log_dir.join("ant-node.2026-08-18.log")).unwrap();
let outcome = tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
assert_eq!(
outcome.live_files.len(),
1,
"the deleted file is no longer live"
);
assert!(outcome.live_files[0].ends_with("ant-node.2026-08-19.log"));
}
#[tokio::test]
async fn polling_does_not_evict_another_nodes_offsets() {
let fixture = Fixture::new();
fixture.append("ant-node.2026-08-19.log", &line("INFO", "mine"));
let mut tailer = fixture.tailer();
let mut offsets = fixture.offsets();
offsets.set("/some/other/node/logs/ant-node.2026-08-19.log", 4242);
tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
assert_eq!(
offsets.get("/some/other/node/logs/ant-node.2026-08-19.log"),
Some(4242),
"another node's position must survive this tailer's poll"
);
}
const INSTALL_A: &str = "0123456789abcdef";
const INSTALL_B: &str = "fedcba9876543210";
#[test]
fn the_document_id_is_stable_and_position_derived() {
let event = TailedEvent {
node_id: 7,
file_name: "ant-node.2026-08-19.log".to_string(),
byte_offset: 104_857,
event: parse_line(&line("INFO", "hello")).unwrap(),
};
assert_eq!(
event.document_id(INSTALL_A),
"0123456789abcdef-7-ant-node.2026-08-19.log-104857"
);
assert_eq!(
event.document_id(INSTALL_A),
event.clone().document_id(INSTALL_A)
);
assert!(
event.document_id(INSTALL_A).len() <= 512,
"Elasticsearch caps _id at 512 bytes"
);
}
#[test]
fn identical_positions_on_two_installations_do_not_collide() {
let same_event = || TailedEvent {
node_id: 1,
file_name: "ant-node.2026-08-19.log".to_string(),
byte_offset: 0,
event: parse_line(&line("INFO", "starting version=0.17.2")).unwrap(),
};
assert_ne!(
same_event().document_id(INSTALL_A),
same_event().document_id(INSTALL_B),
"the same local position on two machines must not share a document id"
);
}
#[test]
fn document_ids_differ_across_nodes_files_and_positions() {
let base = TailedEvent {
node_id: 7,
file_name: "ant-node.2026-08-19.log".to_string(),
byte_offset: 100,
event: parse_line(&line("INFO", "hello")).unwrap(),
};
let other_node = TailedEvent {
node_id: 8,
..base.clone()
};
let other_file = TailedEvent {
file_name: "ant-node.2026-08-20.log".to_string(),
..base.clone()
};
let other_offset = TailedEvent {
byte_offset: 200,
..base.clone()
};
let ids = [
base.document_id(INSTALL_A),
other_node.document_id(INSTALL_A),
other_file.document_id(INSTALL_A),
other_offset.document_id(INSTALL_A),
];
let unique: std::collections::HashSet<&String> = ids.iter().collect();
assert_eq!(unique.len(), 4);
}
#[tokio::test]
async fn a_file_larger_than_one_chunk_is_read_to_the_end() {
let fixture = Fixture::new();
let mut tailer = fixture.tailer();
tailer.mark_primed();
let mut offsets = fixture.offsets();
let one = line(
"INFO",
"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx",
);
let per_chunk = (MAX_CHUNK_BYTES / one.len()) + 1;
let total = per_chunk * 3;
let mut blob = String::with_capacity(total * one.len());
for i in 0..total {
blob.push_str(&format!(
"2026-08-19T20:50:00.123456Z INFO ant_node::node: event {i}\n"
));
}
let file_len = blob.len() as u64;
fixture.append("ant-node.2026-08-19.log", &blob);
let mut seen = 0usize;
for _ in 0..50 {
let outcome = tailer.poll(&mut offsets, LogLevel::Info).await.unwrap();
seen += outcome.events.len();
}
let key = fixture
.log_dir
.join("ant-node.2026-08-19.log")
.display()
.to_string();
let final_offset = offsets.get(&key).unwrap_or(0);
eprintln!(
"file_len={file_len} final_offset={final_offset} events_seen={seen} expected={total} MAX_CHUNK_BYTES={MAX_CHUNK_BYTES}"
);
assert_eq!(
seen, total,
"every event must be shipped; stalled at offset {final_offset} of {file_len}"
);
}
#[tokio::test]
async fn older_unseen_dailies_are_not_backfilled() {
let fixture = Fixture::new();
for day in ["16", "17", "18"] {
fixture.append(
&format!("ant-node.2026-08-{day}.log"),
&line("INFO", &format!("retained history from the {day}th")),
);
}
fixture.append("ant-node.2026-08-19.log", &line("INFO", "today"));
let mut tailer = fixture.tailer();
tailer.mark_primed();
let mut offsets = fixture.offsets();
let outcome = drain(&mut tailer, &mut offsets, LogLevel::Info).await;
let messages: Vec<&str> = outcome
.events
.iter()
.map(|event| event.event.message.as_str())
.collect();
assert_eq!(
messages,
vec!["today"],
"only the newest daily is backfilled; the retained window must be left alone"
);
}
#[tokio::test]
async fn the_newest_daily_is_still_backfilled_from_its_start() {
let fixture = Fixture::new();
fixture.append("ant-node.2026-08-18.log", &line("INFO", "yesterday"));
let mut tailer = fixture.tailer();
tailer.mark_primed();
let mut offsets = fixture.offsets();
drain(&mut tailer, &mut offsets, LogLevel::Info).await;
fixture.append("ant-node.2026-08-19.log", &line("INFO", "today"));
let outcome = drain(&mut tailer, &mut offsets, LogLevel::Info).await;
assert_eq!(outcome.events.len(), 1);
assert_eq!(outcome.events[0].event.message, "today");
}
}