Files
omarchycn/bin/omarchy-agent-usage-codex
4cd8a081cb Fix Codex usage collection on 0.149 (#7649)
* Fix Codex usage collector approval policy

* Capture codex argv with boundaries in the scanner test

The stub joined its arguments with "$*", so the assertion compared one
flattened string and could not tell five arguments from fewer containing
spaces. Passing "-s read-only" and "-a on-request" as single arguments --
which codex rejects as an unexpected argument -- passed the test. NUL
separation and an array comparison keep the boundaries the assertion is
about.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Co-authored-by: Codex XHigh <noreply@openai.com>

---------

Co-authored-by: Omabot <omabot@omarchy.org>
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
Co-authored-by: Codex XHigh <noreply@openai.com>
2026-08-25 16:07:17 +02:00

604 lines
21 KiB
Python
Executable File

#!/usr/bin/python3
# omarchy:summary=Print the Codex usage record as JSON
# omarchy:args=[--force] [--limits-only]
# omarchy:hidden=true
"""Collect Codex usage into one display-ready JSON record.
Local stats come from native Codex CLI session files, pi/omp sessions that
ran through openai-codex, and opencode sessions that ran on an OpenAI
provider; rate limits and the plan come from the Codex app-server RPC. The agents
panel only ever reads the JSON this prints.
"""
import argparse
import fcntl
import hashlib
import json
import os
import select
import shutil
import sqlite3
import subprocess
import sys
import tempfile
import time
from datetime import datetime, timedelta, timezone
from pathlib import Path
AGENT_ID = "codex"
AGENT_NAME = "Codex"
AUTH_HELP = "Run `codex login` to authenticate."
# A scan this recent is only reused to dedup concurrent collector runs (the
# update command backgrounds one per agent while the panel refreshes on its
# own); every periodic widget refresh lands a real rescan, however low
# refreshIntervalSec is set. --limits-only promises only fresh limits, so it
# may reuse a scan for up to 15 minutes.
SCAN_REUSE_SECONDS = 20
LIMITS_ONLY_REUSE_SECONDS = 900
def local_day(value):
if value is None:
return datetime.now().strftime("%Y-%m-%d")
if isinstance(value, (int, float)):
# pi message timestamps are milliseconds; Codex timestamps are usually seconds.
if value > 10_000_000_000:
value = value / 1000
return datetime.fromtimestamp(value).strftime("%Y-%m-%d")
text = str(value)
try:
if text.endswith("Z"):
dt = datetime.fromisoformat(text[:-1] + "+00:00")
else:
dt = datetime.fromisoformat(text)
if dt.tzinfo is not None:
dt = dt.astimezone()
return dt.strftime("%Y-%m-%d")
except Exception:
return datetime.now().strftime("%Y-%m-%d")
def number(value):
try:
return int(value or 0)
except Exception:
return 0
def model_name(raw):
value = str(raw or "codex")
return value if value else "codex"
def runtime_env():
home = str(Path.home())
path_parts = [
os.environ.get("PATH", ""),
f"{home}/.local/bin",
f"{home}/.npm-global/bin",
f"{home}/.local/share/mise/shims",
]
env = os.environ.copy()
env["PATH"] = os.pathsep.join(part for part in path_parts if part)
return env
ENV = runtime_env()
def find_command(name):
return shutil.which(name, path=ENV.get("PATH"))
now = datetime.now()
today = now.strftime("%Y-%m-%d")
recent_dates = [(now - timedelta(days=offset)).strftime("%Y-%m-%d") for offset in range(6, -1, -1)]
recent = {day: {"date": day, "messageCount": 0} for day in recent_dates}
today_tokens_by_model = {}
model_usage = {}
today_sessions = set()
active_days = set()
today_prompts = 0
today_total_tokens = 0
total_prompts = 0
total_sessions = set()
seen_pi_messages = set()
def add_usage(day, session_key, model, input_tokens, output_tokens, cache_read, cache_write):
global today_prompts, today_total_tokens, total_prompts
total = input_tokens + output_tokens + cache_read + cache_write
total_prompts += 1
total_sessions.add(session_key)
active_days.add(day)
bucket = model_usage.setdefault(model, {
"inputTokens": 0,
"outputTokens": 0,
"cacheReadInputTokens": 0,
"cacheCreationInputTokens": 0,
})
bucket["inputTokens"] += input_tokens
bucket["outputTokens"] += output_tokens
bucket["cacheReadInputTokens"] += cache_read
bucket["cacheCreationInputTokens"] += cache_write
if day in recent:
recent[day]["messageCount"] += total
if day == today:
today_prompts += 1
today_sessions.add(session_key)
today_total_tokens += total
today_tokens_by_model[model] = today_tokens_by_model.get(model, 0) + total
def scan_pi_sessions():
roots = [
Path.home() / ".pi" / "agent" / "sessions",
Path.home() / ".omp" / "agent" / "sessions",
]
rg = find_command("rg") or "rg"
for root in roots:
if not root.exists():
continue
try:
proc = subprocess.Popen(
[rg, "--json", "-e", r'"provider"\s*:\s*"openai-codex"', "-e", r'"api"\s*:\s*"openai-codex', str(root)],
stdout=subprocess.PIPE,
stderr=subprocess.DEVNULL,
text=True,
errors="replace",
env=ENV,
)
except FileNotFoundError:
return
assert proc.stdout is not None
for raw in proc.stdout:
try:
event = json.loads(raw)
if event.get("type") != "match":
continue
line = event.get("data", {}).get("lines", {}).get("text", "")
path = event.get("data", {}).get("path", {}).get("text", "pi-session")
entry = json.loads(line)
except Exception:
continue
if entry.get("type") != "message":
continue
message_key = path + ":" + str(entry.get("id") or "")
if message_key in seen_pi_messages:
continue
seen_pi_messages.add(message_key)
message = entry.get("message") or {}
if message.get("role") != "assistant":
continue
provider = str(message.get("provider") or "")
api = str(message.get("api") or "")
if provider != "openai-codex" and not api.startswith("openai-codex"):
continue
usage = message.get("usage") or {}
if not usage:
continue
total = number(usage.get("totalTokens"))
input_tokens = number(usage.get("input"))
output_tokens = number(usage.get("output"))
cache_read = number(usage.get("cacheRead"))
cache_write = number(usage.get("cacheWrite"))
if total and not (input_tokens or output_tokens or cache_read or cache_write):
input_tokens = total
if not (input_tokens or output_tokens or cache_read or cache_write):
continue
day = local_day(entry.get("timestamp") or message.get("timestamp"))
session_key = path
add_usage(day, session_key, model_name(message.get("model")), input_tokens, output_tokens, cache_read, cache_write)
try:
proc.wait(timeout=1)
except Exception:
proc.kill()
def scan_opencode_sessions():
# A subscription burned entirely through opencode leaves no native session
# files, but opencode records per-message provider, model, and token usage
# in its own database. Read-only: opencode may be writing right now.
# Returns whether the scan ran to completion: a scan cut short by a
# database error still contributes what it read, but must not be cached
# as if it were the whole story.
db = Path(os.environ.get("XDG_DATA_HOME") or (Path.home() / ".local" / "share")) / "opencode" / "opencode.db"
if not db.is_file():
return True
try:
conn = sqlite3.connect(db.resolve().as_uri() + "?mode=ro", uri=True, timeout=2)
except sqlite3.Error:
return False
try:
conn.execute("PRAGMA query_only = ON")
# OpenCode DBs grow huge (every historical message, JSON included), and
# Python-side json.loads of every row wasted ~600 MB of RSS on machines
# whose subscription never ran on OpenAI. The json_extract conditions are
# the authority for the rows that reach them: role == "assistant" and
# providerID == "openai", the same exact-match values the Python filter
# below checks. The guards in front are pure acceleration, not a perfect
# proxy for the Python filter:
# - The LIKE gates skip rows whose JSON cannot contain the two
# key/value pairs, avoiding the JSON parse of tens-of-MB blobs. They
# can differ from json.loads on duplicate keys (Python keeps the
# last, SQLite json_extract keeps the first) and ASCII-escaped
# values (\"ass\\u0069stant\" decodes for Python but not for LIKE),
# so a row the authority would accept can be gated out. Both cases
# are vanishingly rare in real opencode data.
# - json_valid guards the parse itself: json_extract() RAISES on
# malformed JSON instead of returning NULL, and one such row would
# otherwise abort the whole scan. SQLite does not promise that AND
# terms evaluate left to right, so the guard is a CASE around each
# json_extract rather than a separate AND term. Rows that are not
# well-formed JSON are skipped here; the per-row try/except below
# stays as the final safety net for rows that pass the SQL filter
# but fail json.loads.
for session_id, raw in conn.execute(
"SELECT session_id, data FROM message"
" WHERE data LIKE '%\"role\"%:%\"assistant\"%'"
" AND data LIKE '%\"providerID\"%:%\"openai\"%'"
" AND CASE WHEN json_valid(data) THEN json_extract(data, '$.role') END = 'assistant'"
" AND CASE WHEN json_valid(data) THEN json_extract(data, '$.providerID') END = 'openai'"
):
# One malformed row must not abort the scan, so every shape assumption
# lives inside the try.
try:
entry = json.loads(raw)
# Exact match: opencode provider ids are free-form, and a custom
# "openai-local" gateway is not this subscription.
if not isinstance(entry, dict) or entry.get("role") != "assistant":
continue
if str(entry.get("providerID") or "") != "openai":
continue
tokens = entry.get("tokens") or {}
cache = tokens.get("cache") or {}
input_tokens = number(tokens.get("input"))
# opencode keeps thinking tokens out of output; both are generated.
output_tokens = number(tokens.get("output")) + number(tokens.get("reasoning"))
cache_read = number(cache.get("read"))
cache_write = number(cache.get("write"))
if not (input_tokens or output_tokens or cache_read or cache_write):
continue
day = local_day((entry.get("time") or {}).get("created"))
model = model_name(str(entry.get("modelID") or "").rstrip("/").split("/")[-1])
except Exception:
continue
add_usage(day, "opencode:" + str(session_id), model, input_tokens, output_tokens, cache_read, cache_write)
except sqlite3.Error:
# Transient lock, schema migration, corruption: the numbers stop here,
# incomplete.
return False
finally:
conn.close()
return True
def scan_native_codex_sessions():
codex_home = Path(os.environ.get("CODEX_HOME") or (Path.home() / ".codex"))
roots = [codex_home / "sessions", codex_home / "archived_sessions"]
files = []
cutoff = time.time() - 30 * 24 * 60 * 60
for root in roots:
if not root.exists():
continue
for path in root.rglob("*.jsonl"):
try:
if path.stat().st_mtime >= cutoff:
files.append(path)
except OSError:
pass
for path in files:
current_model = "codex"
try:
with path.open(errors="replace") as handle:
for raw in handle:
try:
entry = json.loads(raw)
except Exception:
continue
if entry.get("type") == "turn_context":
payload = entry.get("payload") or {}
current_model = model_name(payload.get("model") or payload.get("model_slug") or current_model)
continue
payload = entry.get("payload") or entry
if entry.get("type") == "response_item" and isinstance(payload, dict):
payload = payload.get("payload") or payload
if not isinstance(payload, dict):
continue
if payload.get("type") != "token_count":
continue
info = payload.get("info") or {}
# total_token_usage is cumulative for the session. Adding every
# snapshot makes usage grow quadratically, so count the last turn.
usage = info.get("last_token_usage") or {}
cache_read = number(usage.get("cached_input_tokens"))
cache_write = number(usage.get("cache_write_input_tokens"))
# Cached tokens are included in input_tokens, and reasoning tokens
# are included in output_tokens. Keep the cache split without
# counting either category twice.
input_tokens = max(0, number(usage.get("input_tokens")) - cache_read - cache_write)
output_tokens = number(usage.get("output_tokens"))
if not (input_tokens or output_tokens or cache_read or cache_write):
continue
day = local_day(entry.get("timestamp") or path.stat().st_mtime)
add_usage(day, str(path), current_model, input_tokens, output_tokens, cache_read, cache_write)
except Exception:
continue
def cache_root():
root = Path(os.environ.get("XDG_CACHE_HOME") or (Path.home() / ".cache")) / "omarchy" / "agent-usage"
root.mkdir(parents=True, exist_ok=True)
return root
def scan_cache_paths():
codex_home = Path(os.environ.get("CODEX_HOME") or (Path.home() / ".codex"))
db = Path(os.environ.get("XDG_DATA_HOME") or (Path.home() / ".local" / "share")) / "opencode" / "opencode.db"
# The digest covers every data path the scan reads: the codex session
# roots, the opencode DB, and (via Path.home()) the pi/omp session roots.
digest = hashlib.sha1((str(Path.home()) + "\n" + str(codex_home) + "\n" + str(db)).encode("utf-8")).hexdigest()[:16]
root = cache_root()
return root / f"codex-scan-{digest}.json", root / f"codex-scan-{digest}.lock"
def read_fresh_json(path, max_age_seconds):
if max_age_seconds <= 0 or not path.exists():
return None
try:
# A negative age means the mtime is in the future: the clock moved
# backwards since the write, so the cache's freshness cannot be trusted.
age = time.time() - path.stat().st_mtime
if 0 <= age <= max_age_seconds:
return json.loads(path.read_text(encoding="utf-8"))
except Exception:
return None
return None
def write_json(path, payload):
# A temp name unique to this writer, not derived from the target: several
# collectors can run at once (the update command backgrounds one per agent,
# the panel refreshes on its own), and a shared temp path means the second
# replace finds the first one's file already moved away.
handle_fd, tmp_name = tempfile.mkstemp(dir=path.parent, prefix=path.name + ".", suffix=".tmp")
tmp = Path(tmp_name)
try:
with os.fdopen(handle_fd, "w", encoding="utf-8") as handle:
handle.write(json.dumps(payload, separators=(",", ":")) + "\n")
# mkstemp opens at 0600; nothing in the cache is sensitive, so open it
# up to the usual 0644.
tmp.chmod(0o644)
tmp.replace(path)
except BaseException:
tmp.unlink(missing_ok=True)
raise
# The cache payload is a versioned envelope around the local-stats dict, so a
# corrupted or foreign-shaped file is a cache miss (rescan + rewrite) instead
# of a crash or a garbage record.
def read_cached_stats(cache_file, max_age_seconds):
cached = read_fresh_json(cache_file, max_age_seconds)
if not isinstance(cached, dict) or cached.get("schemaVersion") != 1:
return None
# today* fields only mean "today" on the day they were scanned. A cache
# from another local date (midnight passed, or the clock moved) is a miss,
# not merely old, whatever its mtime says.
if cached.get("scanDate") != today:
return None
stats = cached.get("stats")
if not isinstance(stats, dict):
return None
if not all(key in stats for key in ("todayPrompts", "todayTotalTokens", "recentDays", "activeDates", "modelUsage")):
return None
return stats
def write_cached_stats(cache_file, stats):
try:
write_json(cache_file, {"schemaVersion": 1, "scanDate": today, "stats": stats})
except Exception as exc:
print(f"omarchy-agent-usage-codex: could not write usage cache ({exc})", file=sys.stderr)
def local_stats():
"""Snapshot the aggregated local usage into the record's stats dict."""
return {
"todayPrompts": today_prompts,
"todaySessions": len(today_sessions),
"todayTotalTokens": today_total_tokens,
"todayTokensByModel": today_tokens_by_model,
"recentDays": [recent[day] for day in recent_dates],
"totalPrompts": total_prompts,
"totalSessions": len(total_sessions),
# Days with any recorded usage, for the all-time "N days" summary. The
# dates travel too: merging snapshots from several machines needs their
# union, which a count alone cannot give.
"activeDays": len(active_days),
"activeDates": sorted(active_days),
"modelUsage": model_usage,
}
def run_local_scans():
scan_pi_sessions()
scan_native_codex_sessions()
complete = scan_opencode_sessions()
return local_stats(), complete
def cached_local_stats(max_age):
"""Local stats, with the cache as a pure optimization.
The cache must never take the collector down: any cache-layer failure
(unwritable cache root, lock errors, disk full) degrades to a direct scan
and a warning on stderr. The JSON record is the contract; the cache is not.
"""
try:
return _cached_local_stats(max_age)
except Exception as exc:
print(f"omarchy-agent-usage-codex: cache unavailable ({exc}); scanning directly", file=sys.stderr)
stats, _ = run_local_scans()
return stats
def _cached_local_stats(max_age):
cache_file, lock_file = scan_cache_paths()
cached = read_cached_stats(cache_file, max_age)
if cached is not None:
return cached
with lock_file.open("w") as lock:
fcntl.flock(lock, fcntl.LOCK_EX)
cached = read_cached_stats(cache_file, max_age)
if cached is not None:
return cached
stats, complete = run_local_scans()
# An interrupted scan still serves this run, but caching it would
# suppress the missing usage for every reader until the cache expires.
if complete:
write_cached_stats(cache_file, stats)
return stats
def rpc_request(proc, request_id, method, params=None, timeout=8):
payload = {"id": request_id, "method": method, "params": params or {}}
proc.stdin.write(json.dumps(payload) + "\n")
proc.stdin.flush()
deadline = time.time() + timeout
while time.time() < deadline:
ready, _, _ = select.select([proc.stdout], [], [], 0.25)
if not ready:
continue
line = proc.stdout.readline()
if not line:
break
try:
message = json.loads(line)
except Exception:
continue
if message.get("id") == request_id:
return message
raise TimeoutError(method)
def limit_window(window):
if not isinstance(window, dict):
return None
used = window.get("usedPercent")
if used is None:
return None
mins = number(window.get("windowDurationMins"))
if mins == 10080:
label = "Weekly (7-day)"
elif mins and mins % 60 == 0:
label = f"{mins // 60}h window"
elif mins:
label = f"{mins}m window"
else:
label = "Limit"
reset = window.get("resetsAt")
return {
"label": label,
"percent": float(used) / 100.0,
"resetsAt": datetime.fromtimestamp(number(reset), timezone.utc).isoformat() if reset else "",
}
def fetch_codex_rpc():
result = {"limits": [], "tierLabel": "", "usageStatusText": "", "authHelpText": AUTH_HELP}
codex = find_command("codex")
if not codex:
result["usageStatusText"] = "Codex unavailable"
result["authHelpText"] = "codex not found in PATH"
return result
try:
proc = subprocess.Popen(
[codex, "-s", "read-only", "-a", "on-request", "app-server"],
stdin=subprocess.PIPE,
stdout=subprocess.PIPE,
stderr=subprocess.DEVNULL,
text=True,
env=ENV,
)
except Exception as exc:
result["usageStatusText"] = "Codex unavailable"
result["authHelpText"] = str(exc)
return result
try:
rpc_request(proc, 1, "initialize", {"clientInfo": {"name": "omarchy-agent-usage", "version": "1"}}, timeout=8)
proc.stdin.write(json.dumps({"method": "initialized", "params": {}}) + "\n")
proc.stdin.flush()
account_msg = rpc_request(proc, 2, "account/read", timeout=4)
limits_msg = rpc_request(proc, 3, "account/rateLimits/read", timeout=4)
account = (account_msg.get("result") or {}).get("account") or {}
limits = (limits_msg.get("result") or {}).get("rateLimits") or {}
plan = limits.get("planType") or account.get("planType") or account.get("type") or ""
result["tierLabel"] = str(plan) if plan else ""
for window in (limits.get("primary"), limits.get("secondary")):
entry = limit_window(window)
if entry:
result["limits"].append(entry)
except Exception as exc:
result["usageStatusText"] = "Codex limits unavailable"
result["authHelpText"] = str(exc)
finally:
try:
proc.terminate()
proc.wait(timeout=1)
except Exception:
try:
proc.kill()
except Exception:
pass
return result
def main():
parser = argparse.ArgumentParser()
# --force rescans everything and rewrites the cache. --limits-only is kept
# for CLI compatibility with the panel's refreshLimits() call: only the
# limits probe must be fresh, so it may reuse a scan for far longer than a
# normal run, whose short window exists purely to dedup concurrent
# collector runs.
parser.add_argument("--force", action="store_true")
parser.add_argument("--limits-only", action="store_true")
args = parser.parse_args()
max_age = 0 if args.force else (LIMITS_ONLY_REUSE_SECONDS if args.limits_only else SCAN_REUSE_SECONDS)
stats = cached_local_stats(max_age)
rpc = fetch_codex_rpc()
record = {
"schemaVersion": 1,
"id": AGENT_ID,
"name": AGENT_NAME,
"updatedAt": datetime.now(timezone.utc).isoformat(),
"ready": True,
"hasLocalStats": True,
}
record.update(stats)
record.update(rpc)
print(json.dumps(record, separators=(",", ":")))
if __name__ == "__main__":
main()