use std::collections::HashMap;
use std::fs::{File, OpenOptions};
use std::io::{BufRead, BufReader, Write};
use std::path::{Path, PathBuf};
use std::sync::Mutex;
use serde_json::Value;
const SEGMENTS: u64 = 2;
pub const DEFAULT_MAX_LOG_BYTES: u64 = 32 * 1024 * 1024;
struct NodeFile {
file: File,
bytes: u64,
path: PathBuf,
warned: bool,
}
pub struct RecordLog {
dir: PathBuf,
max_bytes: u64,
nodes: Mutex<HashMap<String, NodeFile>>,
}
impl RecordLog {
pub fn new(dir: impl Into<PathBuf>, max_bytes: u64) -> Self {
RecordLog {
dir: dir.into(),
max_bytes: max_bytes.max(SEGMENTS),
nodes: Mutex::new(HashMap::new()),
}
}
fn segment_limit(&self) -> u64 {
(self.max_bytes / SEGMENTS).max(1)
}
pub fn append(&self, records: &[Value]) {
let mut nodes = self.nodes.lock().unwrap();
for rec in records {
let Some(path) = rec.get("path").and_then(Value::as_str) else {
continue;
};
let Some(rel) = safe_relative_path(path) else {
continue;
};
let line = rec.to_string();
self.append_one(&mut nodes, path, &rel, &line);
}
}
fn append_one(
&self,
nodes: &mut HashMap<String, NodeFile>,
path: &str,
rel: &Path,
line: &str,
) {
if !nodes.contains_key(path) {
let Some(nf) = self.open_node(rel) else {
return;
};
nodes.insert(path.to_string(), nf);
}
let limit = self.segment_limit();
let nf = nodes.get_mut(path).expect("just inserted");
let need = line.len() as u64 + 1;
if nf.bytes > 0 && nf.bytes + need > limit {
rotate(nf);
}
let res = nf
.file
.write_all(line.as_bytes())
.and_then(|()| nf.file.write_all(b"\n"));
match res {
Ok(()) => nf.bytes += need,
Err(e) => {
if !nf.warned {
nf.warned = true;
crate::msg!(
"warning: monitor record log: write failed for {} ({e}); \
this node's log stops here (training continues)",
nf.path.display(),
);
}
}
}
}
fn open_node(&self, rel: &Path) -> Option<NodeFile> {
let path = node_log_path(&self.dir, rel);
if let Some(parent) = path.parent()
&& let Err(e) = std::fs::create_dir_all(parent)
{
crate::msg!(
"warning: monitor record log: cannot create {} ({e}); \
record persistence disabled for this node",
parent.display(),
);
return None;
}
match OpenOptions::new().create(true).append(true).open(&path) {
Ok(file) => {
let bytes = file.metadata().map(|m| m.len()).unwrap_or(0);
Some(NodeFile {
file,
bytes,
path,
warned: false,
})
}
Err(e) => {
crate::msg!(
"warning: monitor record log: cannot open {} ({e}); \
record persistence disabled for this node",
path.display(),
);
None
}
}
}
pub fn tail(&self, path: &str, n: usize) -> Vec<Value> {
if n == 0 {
return Vec::new();
}
if let Ok(mut nodes) = self.nodes.lock()
&& let Some(nf) = nodes.get_mut(path)
{
let _ = nf.file.flush();
}
let Some(rel) = safe_relative_path(path) else {
return Vec::new();
};
let active = node_log_path(&self.dir, &rel);
let mut lines: std::collections::VecDeque<String> =
std::collections::VecDeque::with_capacity(n);
for seg in [rotated_path(&active), active] {
let Ok(file) = File::open(&seg) else { continue };
for line in BufReader::new(file).lines().map_while(Result::ok) {
if line.is_empty() {
continue;
}
if lines.len() == n {
lines.pop_front();
}
lines.push_back(line);
}
}
lines
.into_iter()
.filter_map(|l| serde_json::from_str(&l).ok())
.collect()
}
pub fn flush(&self) {
if let Ok(mut nodes) = self.nodes.lock() {
for nf in nodes.values_mut() {
let _ = nf.file.flush();
}
}
}
}
fn rotate(nf: &mut NodeFile) {
let _ = nf.file.flush();
let rotated = rotated_path(&nf.path);
if std::fs::rename(&nf.path, &rotated).is_err() {
return;
}
match OpenOptions::new().create(true).append(true).open(&nf.path) {
Ok(file) => {
nf.file = file;
nf.bytes = 0;
}
Err(e) => {
if !nf.warned {
nf.warned = true;
crate::msg!(
"warning: monitor record log: reopen after rotate failed for {} ({e})",
nf.path.display(),
);
}
}
}
}
fn rotated_path(active: &Path) -> PathBuf {
let mut s = active.as_os_str().to_os_string();
s.push(".1");
PathBuf::from(s)
}
fn node_log_path(dir: &Path, rel: &Path) -> PathBuf {
let mut s = dir.join(rel).into_os_string();
s.push(".log");
PathBuf::from(s)
}
fn safe_relative_path(path: &str) -> Option<PathBuf> {
if path.is_empty() {
return None;
}
let mut out = PathBuf::new();
for seg in path.split('/') {
if seg.is_empty() || seg == "." || seg == ".." {
return None;
}
let p = Path::new(seg);
if p.components().count() != 1
|| !matches!(
p.components().next(),
Some(std::path::Component::Normal(_))
)
{
return None;
}
out.push(seg);
}
Some(out)
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
struct TempDir(PathBuf);
impl TempDir {
fn new(name: &str) -> Self {
let p = std::env::temp_dir()
.join(format!("flodl-recordlog-{}-{name}", std::process::id()));
let _ = std::fs::remove_dir_all(&p);
std::fs::create_dir_all(&p).unwrap();
TempDir(p)
}
}
impl Drop for TempDir {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.0);
}
}
fn rec(path: &str, tick: u64) -> Value {
json!({ "v": 1, "kind": "node", "path": path, "tick": tick })
}
#[test]
fn path_is_the_filesystem_path() {
let d = TempDir::new("paths");
let log = RecordLog::new(&d.0, DEFAULT_MAX_LOG_BYTES);
log.append(&[
rec("root", 1),
rec("root/exa", 1),
rec("root/exa/rank0", 1),
]);
log.flush();
assert!(d.0.join("root.log").is_file());
assert!(d.0.join("root/exa.log").is_file());
assert!(d.0.join("root/exa/rank0.log").is_file());
}
#[test]
fn appends_one_json_per_line_per_node() {
let d = TempDir::new("append");
let log = RecordLog::new(&d.0, DEFAULT_MAX_LOG_BYTES);
log.append(&[rec("root", 1), rec("root/rank0", 1)]);
log.append(&[rec("root", 2), rec("root/rank0", 2)]);
log.flush();
let root = std::fs::read_to_string(d.0.join("root.log")).unwrap();
assert_eq!(root.lines().count(), 2);
for line in root.lines() {
let v: Value = serde_json::from_str(line).unwrap();
assert_eq!(v["path"], "root");
}
let r0 = std::fs::read_to_string(d.0.join("root/rank0.log")).unwrap();
assert_eq!(r0.lines().count(), 2);
}
#[test]
fn tail_returns_last_n_oldest_first() {
let d = TempDir::new("tail");
let log = RecordLog::new(&d.0, DEFAULT_MAX_LOG_BYTES);
for t in 1..=10 {
log.append(&[rec("root", t)]);
}
let got = log.tail("root", 3);
assert_eq!(got.len(), 3);
assert_eq!(got[0]["tick"], 8);
assert_eq!(got[2]["tick"], 10);
assert_eq!(log.tail("root", 100).len(), 10);
assert!(log.tail("root", 0).is_empty());
assert!(log.tail("root/nope", 5).is_empty());
}
#[test]
fn rotation_bounds_total_bytes_and_drops_oldest() {
let d = TempDir::new("rotate");
let log = RecordLog::new(&d.0, 400);
for t in 1..=200 {
log.append(&[rec("root", t)]);
}
log.flush();
let active = d.0.join("root.log");
let rotated = d.0.join("root.log.1");
let a = std::fs::metadata(&active).unwrap().len();
let r = std::fs::metadata(&rotated).map(|m| m.len()).unwrap_or(0);
assert!(a <= 200, "active segment {a} over the 200 B limit");
assert!(r <= 200, "rotated segment {r} over the 200 B limit");
assert!(a + r <= 400, "total {} over the cap", a + r);
let kept = log.tail("root", 1000);
assert!(!kept.is_empty());
let ticks: Vec<u64> = kept.iter().map(|v| v["tick"].as_u64().unwrap()).collect();
assert_eq!(*ticks.last().unwrap(), 200, "newest record retained");
assert!(ticks.len() < 200, "old records dropped (kept {})", ticks.len());
assert!(!ticks.contains(&1), "oldest record dropped");
for w in ticks.windows(2) {
assert_eq!(w[1], w[0] + 1, "ticks contiguous: {ticks:?}");
}
}
#[test]
fn reopen_resumes_byte_count_so_the_cap_still_holds() {
let d = TempDir::new("reopen");
{
let log = RecordLog::new(&d.0, 400);
for t in 1..=20 {
log.append(&[rec("root", t)]);
}
log.flush();
}
let log = RecordLog::new(&d.0, 400);
for t in 21..=200 {
log.append(&[rec("root", t)]);
}
log.flush();
let a = std::fs::metadata(d.0.join("root.log")).unwrap().len();
let r = std::fs::metadata(d.0.join("root.log.1"))
.map(|m| m.len())
.unwrap_or(0);
assert!(a + r <= 400, "total {} over the cap after reopen", a + r);
assert_eq!(
log.tail("root", 1).first().unwrap()["tick"],
200,
"newest record still retained across reopen",
);
}
#[test]
fn traversal_and_malformed_paths_are_refused() {
assert!(safe_relative_path("root/../../etc/passwd").is_none());
assert!(safe_relative_path("/etc/passwd").is_none());
assert!(safe_relative_path("..").is_none());
assert!(safe_relative_path("root//rank0").is_none());
assert!(safe_relative_path("").is_none());
assert_eq!(
safe_relative_path("root/exa/rank0"),
Some(PathBuf::from("root/exa/rank0")),
);
let d = TempDir::new("traversal");
let log = RecordLog::new(&d.0, DEFAULT_MAX_LOG_BYTES);
log.append(&[rec("../escape", 1), rec("root", 1)]);
log.flush();
assert!(d.0.join("root.log").is_file());
assert!(!d.0.parent().unwrap().join("escape.log").exists());
}
#[test]
fn records_without_a_path_are_skipped() {
let d = TempDir::new("nopath");
let log = RecordLog::new(&d.0, DEFAULT_MAX_LOG_BYTES);
log.append(&[json!({ "kind": "meta" }), rec("root", 1)]);
log.flush();
let root = std::fs::read_to_string(d.0.join("root.log")).unwrap();
assert_eq!(root.lines().count(), 1);
}
#[test]
fn unwritable_dir_never_panics() {
let d = TempDir::new("unwritable");
std::fs::write(d.0.join("root"), b"i am a file, not a dir").unwrap();
let log = RecordLog::new(&d.0, DEFAULT_MAX_LOG_BYTES);
log.append(&[rec("root/rank0", 1)]); log.flush();
assert!(log.tail("root/rank0", 5).is_empty());
}
#[test]
fn dotted_hostnames_do_not_collide() {
let d = TempDir::new("dottedhosts");
let log = RecordLog::new(&d.0, DEFAULT_MAX_LOG_BYTES);
log.append(&[rec("root/10.0.0.5", 1), rec("root/10.0.0.6", 2)]);
log.flush();
assert!(d.0.join("root/10.0.0.5.log").is_file());
assert!(d.0.join("root/10.0.0.6.log").is_file());
let five = log.tail("root/10.0.0.5", 10);
let six = log.tail("root/10.0.0.6", 10);
assert_eq!(five.len(), 1);
assert_eq!(six.len(), 1);
assert_eq!(five[0]["tick"], 1);
assert_eq!(six[0]["tick"], 2);
}
}