use crate::context::Context;
use crate::usage::Snapshot;
use crate::{atomic, home};
use serde::{Deserialize, Serialize};
use std::io::Write;
use std::path::PathBuf;
const KEEP_FOR: i64 = 14 * 86_400;
const APART: i64 = 300;
const MOVED: f64 = 0.5;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Point {
pub at: i64,
pub windows: Vec<(String, f64)>,
#[serde(default)]
pub resets: Vec<(String, i64)>,
}
fn dir(ctx: &Context) -> PathBuf {
home::dir(ctx).join("readings")
}
fn path(ctx: &Context, account_uuid: &str) -> Option<PathBuf> {
let safe = !account_uuid.is_empty()
&& account_uuid
.bytes()
.all(|b| b.is_ascii_alphanumeric() || b == b'-' || b == b'_');
safe.then(|| dir(ctx).join(format!("{account_uuid}.ndjson")))
}
pub fn series(ctx: &Context, account_uuid: &str) -> Vec<Point> {
let Some(path) = path(ctx, account_uuid) else {
return Vec::new();
};
std::fs::read_to_string(path)
.unwrap_or_default()
.lines()
.filter_map(|line| serde_json::from_str::<Point>(line).ok())
.collect()
}
fn point_of(snapshot: &Snapshot, at: i64) -> Point {
Point {
at,
windows: snapshot
.windows
.iter()
.map(|w| (w.kind.clone(), w.percent))
.collect(),
resets: snapshot
.windows
.iter()
.filter_map(|w| w.resets_at.map(|r| (w.kind.clone(), r)))
.collect(),
}
}
fn worth_keeping(last: Option<&Point>, next: &Point) -> bool {
let Some(last) = last else {
return true;
};
if next.at - last.at >= APART {
return true;
}
next.windows.iter().any(|(kind, percent)| {
last.windows
.iter()
.find(|(k, _)| k == kind)
.is_none_or(|(_, was)| (percent - was).abs() >= MOVED)
})
}
pub fn record(ctx: &Context, account_uuid: &str, snapshot: &Snapshot) {
let Some(path) = path(ctx, account_uuid) else {
return;
};
let at = snapshot.observed_at.unwrap_or_else(|| ctx.now());
let next = point_of(snapshot, at);
let existing = series(ctx, account_uuid);
if !worth_keeping(existing.last(), &next) {
return;
}
if home::create_private(&dir(ctx)).is_err() {
return;
}
let Ok(line) = serde_json::to_string(&next) else {
return;
};
let stale = existing.iter().filter(|p| at - p.at > KEEP_FOR).count();
if stale > 0 {
let kept: Vec<String> = existing
.iter()
.filter(|p| at - p.at <= KEEP_FOR)
.filter_map(|p| serde_json::to_string(p).ok())
.chain(std::iter::once(line))
.collect();
let _ = atomic::write(
&path,
format!("{}\n", kept.join("\n")).as_bytes(),
atomic::Perms::Secret,
);
return;
}
if let Ok(mut file) = std::fs::OpenOptions::new()
.create(true)
.append(true)
.mode(0o600)
.open(&path)
{
let _ = writeln!(file, "{line}");
}
}
pub fn forget(ctx: &Context, account_uuid: &str) {
if let Some(path) = path(ctx, account_uuid) {
let _ = std::fs::remove_file(path);
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Runway {
Burning(i64),
Resting(i64),
Unknown,
}
impl Runway {
pub fn seconds(self) -> Option<i64> {
match self {
Runway::Burning(s) | Runway::Resting(s) => Some(s),
Runway::Unknown => None,
}
}
}
const ENOUGH_SPAN: i64 = 900;
const ENOUGH_POINTS: usize = 3;
pub fn runway(points: &[Point], now: i64) -> Runway {
let recent: Vec<&Point> = points.iter().filter(|p| now - p.at <= KEEP_FOR).collect();
if recent.len() < ENOUGH_POINTS {
return Runway::Unknown;
}
let (first, last) = (recent[0], recent[recent.len() - 1]);
let span = last.at - first.at;
if span < ENOUGH_SPAN {
return Runway::Unknown;
}
let mut shortest: Option<Runway> = None;
for (kind, percent) in &last.windows {
let was = first
.windows
.iter()
.find(|(k, _)| k == kind)
.map(|(_, p)| *p);
let resets_in = last
.resets
.iter()
.find(|(k, _)| k == kind)
.map(|(_, at)| at - now)
.filter(|left| *left > 0);
let filling = was.map_or(0.0, |was| percent - was) / span as f64;
let this = if filling > 0.0 && *percent < 100.0 {
let until_full = ((100.0 - percent) / filling) as i64;
match resets_in {
Some(reset) if reset < until_full => Runway::Resting(reset),
_ => Runway::Burning(until_full.max(0)),
}
} else {
match resets_in {
Some(reset) => Runway::Resting(reset),
None => continue,
}
};
if shortest
.and_then(Runway::seconds)
.is_none_or(|best| this.seconds().is_some_and(|s| s < best))
{
shortest = Some(this);
}
}
shortest.unwrap_or(Runway::Unknown)
}
pub fn runway_for(ctx: &Context, account_uuid: &str, now: i64) -> Runway {
runway(&series(ctx, account_uuid), now)
}
use std::os::unix::fs::OpenOptionsExt;
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_real_codex_account_has_a_history() {
let ctx = Context::new(std::path::PathBuf::from("/nowhere"))
.with_pitboard_home(std::path::PathBuf::from("/nowhere/.pitboard"));
assert!(
path(
&ctx,
"8c3f0f86-0a7c-4d52-9b0e-1f2a3b4c5d6e_user-AbC123dEf456"
)
.is_some()
);
}
use crate::time::{Clock, FixedClock};
use crate::usage::{Source, Window};
use std::sync::Arc;
const NOW: i64 = 1_760_000_000;
struct Scratch(PathBuf);
impl Drop for Scratch {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.0);
}
}
fn machine(name: &str) -> (Context, Arc<FixedClock>, Scratch) {
let root = std::env::temp_dir().join(format!(
"pitboard-history-{name}-{}-{:?}",
std::process::id(),
std::thread::current().id()
));
let _ = std::fs::remove_dir_all(&root);
let clock = Arc::new(FixedClock::at(NOW));
let ctx = Context::new(root.clone())
.with_pitboard_home(root.clone())
.with_clock(Arc::clone(&clock) as Arc<dyn Clock>);
home::ensure(&ctx).expect("a home");
(ctx, clock, Scratch(root))
}
fn reading(at: i64, five_hour: f64, resets_at: i64) -> Snapshot {
Snapshot {
observed_at: Some(at),
account_uuid: None,
source: Source::Live,
windows: vec![Window {
kind: "five_hour".into(),
scope: None,
percent: five_hour,
resets_at: Some(resets_at),
is_active: true,
severity: None,
length_seconds: None,
}],
}
}
fn points(samples: &[(i64, f64)], resets_at: i64) -> Vec<Point> {
samples
.iter()
.map(|(at, percent)| Point {
at: *at,
windows: vec![("five_hour".into(), *percent)],
resets: vec![("five_hour".into(), resets_at)],
})
.collect()
}
#[test]
fn a_reading_is_kept_and_read_back() {
let (ctx, _clock, _s) = machine("round-trip");
record(&ctx, "acc", &reading(NOW, 10.0, NOW + 3600));
record(&ctx, "acc", &reading(NOW + 600, 20.0, NOW + 3000));
let kept = series(&ctx, "acc");
assert_eq!(kept.len(), 2);
assert_eq!(kept[0].windows, vec![("five_hour".to_string(), 10.0)]);
assert_eq!(kept[1].at, NOW + 600);
}
#[test]
fn a_reading_that_says_nothing_new_is_not_kept() {
let (ctx, _clock, _s) = machine("quiet");
record(&ctx, "acc", &reading(NOW, 10.0, NOW + 3600));
record(&ctx, "acc", &reading(NOW + 10, 10.0, NOW + 3590));
record(&ctx, "acc", &reading(NOW + 20, 10.2, NOW + 3580));
assert_eq!(series(&ctx, "acc").len(), 1, "nothing moved");
record(&ctx, "acc", &reading(NOW + 30, 11.0, NOW + 3570));
assert_eq!(series(&ctx, "acc").len(), 2);
record(&ctx, "acc", &reading(NOW + 30 + APART, 11.0, NOW + 3000));
assert_eq!(series(&ctx, "acc").len(), 3);
}
#[test]
fn anything_older_than_a_fortnight_goes() {
let (ctx, _clock, _s) = machine("pruned");
for day in 0..20 {
record(
&ctx,
"acc",
&reading(NOW + day * 86_400, day as f64, NOW + day * 86_400 + 3600),
);
}
let kept = series(&ctx, "acc");
let newest = kept.last().expect("something").at;
assert!(
kept.iter().all(|p| newest - p.at <= KEEP_FOR),
"kept {} readings, oldest {}s back",
kept.len(),
newest - kept[0].at
);
assert!(kept.len() < 20);
}
#[test]
fn an_account_being_used_lasts_until_its_limit_fills() {
let series = points(
&[(NOW, 20.0), (NOW + 1800, 35.0), (NOW + 3600, 50.0)],
NOW + 86_400,
);
match runway(&series, NOW + 3600) {
Runway::Burning(seconds) => {
assert!(
(5_900..6_100).contains(&seconds),
"about a hundred minutes, got {seconds}"
);
}
other => panic!("got {other:?}"),
}
}
#[test]
fn an_account_nobody_is_using_lasts_until_its_window_resets() {
let series = points(
&[(NOW, 40.0), (NOW + 1800, 40.0), (NOW + 3600, 40.0)],
NOW + 7200,
);
assert_eq!(runway(&series, NOW + 3600), Runway::Resting(3600));
}
#[test]
fn a_reset_that_comes_first_is_what_the_account_lasts_until() {
let series = points(
&[(NOW, 20.0), (NOW + 1800, 25.0), (NOW + 3600, 30.0)],
NOW + 4200,
);
assert_eq!(
runway(&series, NOW + 3600),
Runway::Resting(600),
"it resets in ten minutes; it would take hours to fill"
);
}
#[test]
fn too_little_to_go_on_says_so_rather_than_guessing() {
assert_eq!(runway(&[], NOW), Runway::Unknown);
assert_eq!(
runway(
&points(&[(NOW, 10.0), (NOW + 600, 20.0)], NOW + 7200),
NOW + 600
),
Runway::Unknown,
"two readings are not a rate"
);
assert_eq!(
runway(
&points(
&[(NOW, 10.0), (NOW + 60, 15.0), (NOW + 120, 20.0)],
NOW + 7200
),
NOW + 120
),
Runway::Unknown,
"three readings two minutes apart are noise"
);
}
#[test]
fn a_window_that_reset_in_the_middle_does_not_produce_a_nonsense_rate() {
let series = points(
&[(NOW, 90.0), (NOW + 1800, 95.0), (NOW + 3600, 5.0)],
NOW + 7200,
);
assert_eq!(runway(&series, NOW + 3600), Runway::Resting(3600));
}
#[test]
fn what_is_known_about_an_account_goes_when_the_account_does() {
let (ctx, _clock, _s) = machine("forget");
record(&ctx, "acc", &reading(NOW, 10.0, NOW + 3600));
assert!(!series(&ctx, "acc").is_empty());
forget(&ctx, "acc");
assert!(series(&ctx, "acc").is_empty());
}
#[test]
fn a_name_that_is_not_an_identifier_is_refused_rather_than_escaped() {
let (ctx, _clock, _s) = machine("traversal");
record(&ctx, "../../etc/passwd", &reading(NOW, 10.0, NOW + 3600));
assert!(series(&ctx, "../../etc/passwd").is_empty());
assert!(path(&ctx, "").is_none());
}
}