use crate::models::{
CodexQuotaSnapshot, CodexSessionRateLimits, CodexSessionWindow, QuotaSource, QuotaWindow,
};
use crate::utils::{is_codex_session_file, resolve_paths};
use anyhow::Result;
use serde_json::Value;
use std::fs::File;
use std::io::{BufRead, BufReader};
use std::path::{Path, PathBuf};
use walkdir::WalkDir;
const MAX_FILES: usize = 64;
pub fn latest_session_rate_limits() -> Result<Option<CodexQuotaSnapshot>> {
latest_session_rate_limits_in(&resolve_paths()?.codex_session_dir)
}
pub fn latest_session_rate_limits_in(
codex_session_dir: &Path,
) -> Result<Option<CodexQuotaSnapshot>> {
if !codex_session_dir.exists() {
return Ok(None);
}
let now = chrono::Local::now().timestamp();
for file in newest_codex_files(codex_session_dir, MAX_FILES) {
let values = read_jsonl_lenient(&file);
if let Some(snap) = extract_latest_rate_limits(&values, now) {
return Ok(Some(snap));
}
}
Ok(None)
}
fn read_jsonl_lenient(path: &Path) -> Vec<Value> {
let Ok(file) = File::open(path) else {
return Vec::new();
};
let mut values = Vec::new();
for line in BufReader::new(file).lines() {
let Ok(line) = line else {
continue;
};
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
if let Ok(value) = serde_json::from_str::<Value>(trimmed) {
values.push(value);
}
}
values
}
const MAX_DAY_DIRS: usize = 14;
fn newest_codex_files(dir: &Path, n: usize) -> Vec<PathBuf> {
let mut files: Vec<(std::time::SystemTime, PathBuf)> = Vec::new();
for day in recent_day_dirs(dir, MAX_DAY_DIRS) {
for entry in WalkDir::new(&day).into_iter().filter_map(|e| e.ok()) {
if !entry.file_type().is_file() {
continue;
}
let path = entry.path();
if !is_codex_session_file(path) {
continue;
}
let Ok(metadata) = entry.metadata() else {
continue;
};
if let Ok(modified) = metadata.modified() {
files.push((modified, path.to_path_buf()));
}
}
}
files.sort_by_key(|(modified, _)| std::cmp::Reverse(*modified));
files.truncate(n);
files.into_iter().map(|(_, p)| p).collect()
}
fn recent_day_dirs(sessions: &Path, limit: usize) -> Vec<PathBuf> {
let mut days = Vec::new();
for year in sorted_subdirs_desc(sessions) {
for month in sorted_subdirs_desc(&year) {
for day in sorted_subdirs_desc(&month) {
days.push(day);
if days.len() >= limit {
return days;
}
}
}
}
days
}
fn sorted_subdirs_desc(dir: &Path) -> Vec<PathBuf> {
let Ok(entries) = std::fs::read_dir(dir) else {
return Vec::new();
};
let mut subdirs: Vec<PathBuf> = entries
.filter_map(|e| e.ok())
.filter(|e| e.file_type().map(|t| t.is_dir()).unwrap_or(false))
.map(|e| e.path())
.collect();
subdirs.sort_by(|a, b| b.file_name().cmp(&a.file_name()));
subdirs
}
pub fn extract_latest_rate_limits(values: &[Value], now: i64) -> Option<CodexQuotaSnapshot> {
for v in values.iter().rev() {
let Some(payload) = v.get("payload") else {
continue;
};
let rl_val = payload
.get("rate_limits")
.or_else(|| payload.get("info").and_then(|i| i.get("rate_limits")));
let Some(rl_val) = rl_val else {
continue;
};
let Ok(rl) = serde_json::from_value::<CodexSessionRateLimits>(rl_val.clone()) else {
continue;
};
if rl.limit_id.as_deref().is_some_and(|id| id != "codex") {
continue;
}
let primary = rl
.primary
.as_ref()
.map(map_session_window)
.filter(|w| is_window_live(w, now));
let secondary = rl
.secondary
.as_ref()
.map(map_session_window)
.filter(|w| is_window_live(w, now));
if primary.is_none() && secondary.is_none() {
continue;
}
return Some(CodexQuotaSnapshot {
source: QuotaSource::SessionFallback,
fetched_at: now,
plan_type: rl.plan_type,
primary,
secondary,
credits_balance: None,
has_credits: None,
unlimited: None,
reset_credits_available: None,
reset_credit_expirations: None,
approx_messages: None,
spend_limit: None,
limit_reached: None,
needs_login: false,
});
}
None
}
fn is_window_live(w: &QuotaWindow, now: i64) -> bool {
w.resets_at_unix.is_some_and(|reset| reset > now)
}
fn map_session_window(w: &CodexSessionWindow) -> QuotaWindow {
QuotaWindow {
used_percent: w.used_percent.unwrap_or(0.0),
resets_at_unix: w.resets_at,
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn picks_newest_rate_limits() {
let values = vec![
json!({"payload":{"rate_limits":{"primary":{"used_percent":10.0,"window_minutes":300,"resets_at":111},"plan_type":"plus"}}}),
json!({"payload":{"type":"message"}}),
json!({"payload":{"rate_limits":{"primary":{"used_percent":42.0,"window_minutes":300,"resets_at":222},"secondary":{"used_percent":69.0,"window_minutes":10080,"resets_at":333},"plan_type":"plus"}}}),
];
let snap = extract_latest_rate_limits(&values, 0).unwrap();
assert_eq!(snap.source, QuotaSource::SessionFallback);
assert_eq!(snap.primary.as_ref().unwrap().used_percent, 42.0);
assert_eq!(snap.primary.as_ref().unwrap().resets_at_unix, Some(222));
assert_eq!(snap.secondary.as_ref().unwrap().used_percent, 69.0);
assert_eq!(snap.plan_type.as_deref(), Some("plus"));
}
#[test]
fn handles_info_nested_rate_limits() {
let values = vec![
json!({"payload":{"info":{"rate_limits":{"primary":{"used_percent":5.0,"resets_at":1}}}}}),
];
let snap = extract_latest_rate_limits(&values, 0).unwrap();
assert_eq!(snap.primary.unwrap().used_percent, 5.0);
}
#[test]
fn returns_none_without_rate_limits() {
let values = vec![json!({"payload":{"type":"message"}}), json!({"foo":1})];
assert!(extract_latest_rate_limits(&values, 0).is_none());
}
#[test]
fn skips_non_codex_limit_family() {
let values = vec![
json!({"payload":{"rate_limits":{"limit_id":"codex","primary":{"used_percent":12.0,"resets_at":1},"plan_type":"plus"}}}),
json!({"payload":{"rate_limits":{"limit_id":"codex_other","primary":{"used_percent":95.0,"resets_at":2}}}}),
];
let snap = extract_latest_rate_limits(&values, 0).unwrap();
assert_eq!(snap.primary.unwrap().used_percent, 12.0);
assert_eq!(snap.plan_type.as_deref(), Some("plus"));
}
#[test]
fn accepts_missing_limit_id() {
let values =
vec![json!({"payload":{"rate_limits":{"primary":{"used_percent":7.0,"resets_at":1}}}})];
let snap = extract_latest_rate_limits(&values, 0).unwrap();
assert_eq!(snap.primary.unwrap().used_percent, 7.0);
}
#[test]
fn rejects_fully_expired_snapshot() {
let values = vec![
json!({"payload":{"rate_limits":{"limit_id":"codex","primary":{"used_percent":90.0,"resets_at":500},"secondary":{"used_percent":80.0,"resets_at":400}}}}),
];
assert!(extract_latest_rate_limits(&values, 1000).is_none());
}
#[test]
fn drops_expired_window_keeps_live_one() {
let values = vec![
json!({"payload":{"rate_limits":{"limit_id":"codex","primary":{"used_percent":90.0,"resets_at":500},"secondary":{"used_percent":44.0,"resets_at":2000}}}}),
];
let snap = extract_latest_rate_limits(&values, 1000).unwrap();
assert!(snap.primary.is_none(), "expired 5h window dropped");
assert_eq!(snap.secondary.unwrap().used_percent, 44.0);
}
#[test]
fn drops_window_without_reset_time() {
let values = vec![
json!({"payload":{"rate_limits":{"limit_id":"codex","primary":{"used_percent":50.0}}}}),
];
assert!(extract_latest_rate_limits(&values, 1000).is_none());
}
#[test]
fn recent_day_dirs_is_bounded_and_newest_first() {
let dir = tempfile::tempdir().unwrap();
let sessions = dir.path();
for d in 1..=20u32 {
std::fs::create_dir_all(sessions.join(format!("2026/06/{d:02}"))).unwrap();
}
let recent = recent_day_dirs(sessions, 14);
assert_eq!(recent.len(), 14, "bounded to the limit, not all 20 days");
assert!(recent[0].ends_with("2026/06/20"), "newest day first");
assert!(recent[13].ends_with("2026/06/07"), "stops 14 days back");
}
#[test]
fn newest_files_within_date_dirs_sorted_and_capped() {
use std::io::Write;
use std::time::{Duration, SystemTime};
let dir = tempfile::tempdir().unwrap();
let sessions = dir.path();
let base = SystemTime::UNIX_EPOCH + Duration::from_secs(1_700_000_000);
let mk = |rel: &str, secs: u64| {
let path = sessions.join(rel);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
let mut file = std::fs::File::create(&path).unwrap();
file.write_all(b"{}\n").unwrap();
file.set_modified(base + Duration::from_secs(secs)).unwrap();
path
};
let old = mk("2026/06/26/rollout-old.jsonl", 0);
let new2 = mk("2026/06/27/rollout-b.jsonl", 10);
let new1 = mk("2026/06/27/rollout-a.jsonl", 20);
std::fs::write(sessions.join("2026/06/27/notes.txt"), "x").unwrap();
let newest = newest_codex_files(sessions, 2);
assert_eq!(
newest,
vec![new1, new2],
"newest mtime first, cap respected"
);
assert!(!newest.contains(&old), "older-day file dropped by the cap");
}
#[test]
fn lenient_reader_keeps_good_lines_despite_torn_tail() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("rollout-live.jsonl");
let body = concat!(
"{\"payload\":{\"rate_limits\":{\"primary\":{\"used_percent\":33.0,\"resets_at\":7}}}}\n",
"{\"payload\":{\"type\":\"message\"}}\n",
"{\"payload\":{\"rate_limits\":{\"primary\":{\"used_per",
);
std::fs::write(&path, body).unwrap();
let values = read_jsonl_lenient(&path);
assert_eq!(values.len(), 2, "torn tail dropped, complete lines kept");
let snap = extract_latest_rate_limits(&values, 0).unwrap();
assert_eq!(snap.primary.unwrap().used_percent, 33.0);
}
}