use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use std::fs::OpenOptions;
use std::io::Write;
use std::path::Path;
pub const MARK_KIND: &str = "mark";
pub const PENDING_FILE: &str = "marks.pending";
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MarkRecord {
#[serde(default, skip_serializing_if = "String::is_empty")]
pub session_id: String,
pub kind: String,
pub wall_time: DateTime<Utc>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub elapsed_ms: Option<u64>,
pub note: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub data: Option<serde_json::Value>,
}
impl MarkRecord {
pub fn new(note: impl Into<String>) -> Self {
Self {
session_id: String::new(),
kind: MARK_KIND.to_string(),
wall_time: Utc::now(),
elapsed_ms: None,
note: note.into(),
data: None,
}
}
pub fn in_session(mut self, session_id: &str, started_at: DateTime<Utc>) -> Self {
self.session_id = session_id.to_string();
let ms = self
.wall_time
.signed_duration_since(started_at)
.num_milliseconds();
self.elapsed_ms = Some(ms.max(0) as u64);
self
}
pub fn with_data(mut self, data: serde_json::Value) -> Self {
self.data = Some(data);
self
}
}
pub fn append_line(path: &Path, line: &str) -> std::io::Result<()> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let mut file = OpenOptions::new().create(true).append(true).open(path)?;
let mut buf = String::with_capacity(line.len() + 1);
buf.push_str(line);
buf.push('\n');
file.write_all(buf.as_bytes())?;
file.flush()?;
Ok(())
}
pub fn append_to_timeline(session_dir: &Path, mark: &MarkRecord) -> std::io::Result<()> {
let line = serde_json::to_string(mark)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
append_line(&session_dir.join("timeline.jsonl"), &line)
}
pub fn notify_pending(session_dir: &Path, mark: &MarkRecord) -> std::io::Result<()> {
let line = serde_json::to_string(mark)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
append_line(&session_dir.join(PENDING_FILE), &line)
}
#[derive(Debug)]
pub struct MarkTail {
path: std::path::PathBuf,
offset: u64,
carry: String,
}
impl MarkTail {
pub fn new(path: impl Into<std::path::PathBuf>) -> Self {
let path = path.into();
let offset = std::fs::metadata(&path).map(|m| m.len()).unwrap_or(0);
Self {
path,
offset,
carry: String::new(),
}
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn poll(&mut self) -> Vec<String> {
let Ok(len) = std::fs::metadata(&self.path).map(|m| m.len()) else {
return Vec::new();
};
if len < self.offset {
self.offset = 0;
self.carry.clear();
}
if len == self.offset {
return Vec::new();
}
let Ok(bytes) = read_range(&self.path, self.offset, len) else {
return Vec::new();
};
self.offset = len;
self.carry.push_str(&String::from_utf8_lossy(&bytes));
let Some(at) = self.carry.rfind('\n') else {
return Vec::new();
};
let complete: String = self.carry.drain(..=at).collect();
complete.lines().map(|line| line.to_string()).collect()
}
}
fn read_range(path: &Path, from: u64, to: u64) -> std::io::Result<Vec<u8>> {
use std::io::{Read, Seek, SeekFrom};
let mut file = std::fs::File::open(path)?;
file.seek(SeekFrom::Start(from))?;
let mut buf = vec![0u8; (to - from) as usize];
file.read_exact(&mut buf)?;
Ok(buf)
}
pub fn slug(label: &str) -> String {
let mut out = String::with_capacity(label.len());
let mut last_dash = false;
for ch in label.chars() {
if ch.is_ascii_alphanumeric() || ch == '.' {
out.push(ch.to_ascii_lowercase());
last_dash = false;
} else if !last_dash {
out.push('-');
last_dash = true;
}
}
let trimmed = out.trim_matches('-');
if trimmed.len() <= 48 {
return trimmed.to_string();
}
let cut = &trimmed[..48];
match cut.rfind('-') {
Some(at) if at >= 24 => cut[..at].to_string(),
_ => cut.trim_end_matches('-').to_string(),
}
}
pub fn parse_label_line(line: &str) -> Option<MarkRecord> {
let trimmed = line.trim();
if trimmed.is_empty() {
return None;
}
let Ok(value) = serde_json::from_str::<serde_json::Value>(trimmed) else {
return Some(MarkRecord::new(trimmed));
};
let serde_json::Value::Object(ref map) = value else {
return Some(MarkRecord::new(trimmed));
};
let note = ["note", "label", "kind"]
.iter()
.find_map(|key| map.get(*key).and_then(|v| v.as_str()))
.unwrap_or(trimmed)
.to_string();
Some(MarkRecord::new(note).with_data(value))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn slug_keeps_words_and_collapses_the_rest() {
assert_eq!(slug("before-checkout"), "before-checkout");
assert_eq!(
slug("permission-denied @ orders"),
"permission-denied-orders"
);
assert_eq!(slug("Ready!"), "ready");
assert_eq!(slug("buyer.01.k3f9"), "buyer.01.k3f9");
assert_eq!(slug(" spaced out "), "spaced-out");
}
#[test]
fn slug_survives_a_label_made_only_of_punctuation() {
assert_eq!(slug("///"), "");
assert_eq!(slug(""), "");
}
#[test]
fn slug_truncates_on_a_word_boundary_when_it_can() {
let long = "before-we-place-the-order-and-then-check-the-receipt-page-loads";
let s = slug(long);
assert!(s.len() <= 48, "{s} is {} chars", s.len());
assert!(!s.ends_with('-'));
assert!(long.starts_with(&s));
}
#[test]
fn slug_truncates_hard_when_there_is_no_boundary() {
let long = "a".repeat(80);
assert_eq!(slug(&long).len(), 48);
}
#[test]
fn a_plain_line_is_the_label() {
let mark = parse_label_line("ready").unwrap();
assert_eq!(mark.note, "ready");
assert_eq!(mark.kind, MARK_KIND);
assert!(mark.data.is_none());
}
#[test]
fn a_json_line_keeps_its_payload_and_finds_a_label() {
let mark = parse_label_line(r#"{"kind":"route","route":"/orders","seq":3}"#).unwrap();
assert_eq!(mark.note, "route");
assert_eq!(mark.data.as_ref().unwrap()["route"], "/orders");
let labelled = parse_label_line(r#"{"note":"before-admin","seq":4}"#).unwrap();
assert_eq!(labelled.note, "before-admin");
}
#[test]
fn blank_lines_are_skipped_and_broken_json_is_a_label() {
assert!(parse_label_line(" ").is_none());
assert!(parse_label_line("").is_none());
let mark = parse_label_line(r#"{"kind":"#).unwrap();
assert_eq!(mark.note, r#"{"kind":"#);
}
#[test]
fn elapsed_is_measured_from_the_session_start() {
let started = Utc::now() - chrono::Duration::milliseconds(2500);
let mark = MarkRecord::new("x").in_session("s1", started);
assert_eq!(mark.session_id, "s1");
let ms = mark.elapsed_ms.unwrap();
assert!((2400..3000).contains(&ms), "elapsed_ms was {ms}");
}
#[test]
fn a_mark_made_before_the_session_started_clamps_to_zero() {
let started = Utc::now() + chrono::Duration::seconds(5);
let mark = MarkRecord::new("x").in_session("s1", started);
assert_eq!(mark.elapsed_ms, Some(0));
}
#[test]
fn append_line_writes_one_record_per_call() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("nested").join("out.jsonl");
append_line(&path, "one").unwrap();
append_line(&path, "two").unwrap();
let body = std::fs::read_to_string(&path).unwrap();
assert_eq!(body, "one\ntwo\n");
}
#[test]
fn timeline_and_pending_are_separate_files() {
let dir = tempfile::tempdir().unwrap();
let mark = MarkRecord::new("ready");
append_to_timeline(dir.path(), &mark).unwrap();
notify_pending(dir.path(), &mark).unwrap();
let timeline = std::fs::read_to_string(dir.path().join("timeline.jsonl")).unwrap();
let pending = std::fs::read_to_string(dir.path().join(PENDING_FILE)).unwrap();
assert!(timeline.contains(r#""kind":"mark""#));
assert_eq!(timeline.lines().count(), 1);
assert_eq!(pending.lines().count(), 1);
}
}