use crate::paths;
use crate::util::now_rfc3339;
use serde_json::{json, Value};
use std::collections::HashMap;
use std::fs::OpenOptions;
use std::io::Write;
use std::path::PathBuf;
fn path(worker: &str) -> PathBuf {
paths::worker_journal_file(worker)
}
pub fn append(worker: &str, run_id: &str, kind: &str, detail: Value) {
let _ = append_required(worker, run_id, kind, detail);
}
pub fn append_required(
worker: &str,
run_id: &str,
kind: &str,
detail: Value,
) -> Result<(), String> {
let p = path(worker);
if let Some(dir) = p.parent() {
std::fs::create_dir_all(dir).map_err(|error| error.to_string())?;
}
let ev = json!({ "ts": now_rfc3339(), "worker": worker, "run_id": run_id, "kind": kind, "detail": detail });
let mut f = OpenOptions::new()
.create(true)
.append(true)
.open(&p)
.map_err(|error| error.to_string())?;
writeln!(f, "{ev}").map_err(|error| error.to_string())?;
f.sync_data().map_err(|error| error.to_string())
}
pub fn for_run(worker: &str, run_id: &str) -> Vec<Value> {
std::fs::read_to_string(path(worker))
.unwrap_or_default()
.lines()
.filter_map(|l| serde_json::from_str::<Value>(l).ok())
.filter(|e| e.get("run_id").and_then(|v| v.as_str()) == Some(run_id))
.collect()
}
pub fn tail(worker: &str, n: usize) -> Vec<Value> {
let text = std::fs::read_to_string(path(worker)).unwrap_or_default();
let mut evs: Vec<Value> = text
.lines()
.filter_map(|l| serde_json::from_str(l).ok())
.collect();
let len = evs.len();
if len > n {
evs.split_off(len - n)
} else {
evs
}
}
pub fn run_workers() -> HashMap<String, String> {
fn files(dir: &std::path::Path, out: &mut Vec<PathBuf>) {
for entry in std::fs::read_dir(dir).into_iter().flatten().flatten() {
let path = entry.path();
if path.is_dir() {
files(&path, out);
} else if path.file_name().and_then(|n| n.to_str()) == Some("events.jsonl") {
out.push(path);
}
}
}
let mut paths = Vec::new();
files(&crate::paths::workers_data_dir(), &mut paths);
let mut out = HashMap::new();
for path in paths {
let text = std::fs::read_to_string(path).unwrap_or_default();
for event in text
.lines()
.filter_map(|line| serde_json::from_str::<Value>(line).ok())
{
let run_id = event.get("run_id").and_then(Value::as_str).unwrap_or("");
let worker = event.get("worker").and_then(Value::as_str).unwrap_or("");
if !run_id.is_empty() && !worker.is_empty() {
out.insert(
run_id.to_string(),
worker.strip_prefix("org/").unwrap_or(worker).to_string(),
);
}
}
}
out
}