+14








75250d37ac
* Read Codex app-server replies from the raw fd (#13703) * Resolve Codex through mise which instead of running the lazy launcher (#13109) * Skip the Codex app-server probe when there are no credentials (#13106) Adapted: credentials are checked in the home being probed rather than in the CODEX_HOME environment variable, since each registered account is probed in its own home, so a signed-out secondary account isn't hidden behind the primary's login. A home without credentials reports "Waiting for auth" like any other signed-out home. The credentials store setting is read with tomllib, so a single-quoted value counts too. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Show the Codex CLI's own error when its app-server dies (#8977) Detect an app-server that exits or stops answering, and report the end of its stderr instead of a bare RPC method name. Rebased onto the raw-fd reply reader; the switch from "-a on-request" to "-a never" is left out, keeping the current approval flags. * Count pi sessions when HOME is a git checkout (#13209) * Count only OpenAI-backed native sessions as Codex usage (#12032) * Deduplicate Pi usage across forked sessions (#8602) * Skip unchanged native Codex token snapshots (#10531) * Count omp and pi profile sessions in the agent usage collectors (#9546) `omp --profile=<name>` (and pi's equivalent) relocates the whole agent tree under <base>/profiles/<name>/. The Claude and Codex collectors only ever scanned <base>/agent/sessions, so a subscription driven entirely through a profile was invisible to the agents panel: no tokens by day, no tokens by model, no prompt or session counts. Discover the profile roots alongside the default one. Sessions are keyed by file path, so a profile adds sessions instead of double-counting the default root, and a missing or unreadable profiles directory leaves the existing behavior untouched. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014zFbJcDEEpV5BAmsH6kAB3 * Skip unrelated Codex session lines before JSON parsing (#12803) Adapted: session_meta lines also pass the pre-filter, since the provider filter from #12032 reads them to skip rollouts served by a non-OpenAI provider. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Read only the Codex session files that changed since the last scan (#12595) Native Codex rollouts keep per-file totals between runs, replayed while a file's mtime and size are unchanged. Rebased onto the session_meta provider filter, snapshot dedup, and line pre-filter, which now live in the per-file reader. pi and omp sessions are left out of the per-file cache: a forked pi session repeats its parent's messages, so they are deduplicated across the whole tree on every scan. * Count streamed Claude messages by their highest-output usage line (#10606) Claude Code writes a streamed assistant response as several transcript lines that share one message id, one per content block. Each line carries a usage object. The first line's output_tokens is a placeholder, often 1, and the last line has the real count. Input and cache fields usually match across the lines. The scanner dedupes by message id and keeps the first line it sees, so it under-counts output tokens. On a machine with 2,577 transcripts it reported 39.0M output tokens against 60.1M used, a 35% shortfall. Input and both cache fields differed by under 0.01%. Keep the line with the highest output count, with the last one scanned winning a tie. The whole line is kept because a response can fall back to another model mid-stream. Those lines are separate snapshots with different cache figures and a different model, and taking a maximum per field across them over-counts cache tokens and credits the wrong model. The zero-usage check now runs before dedup, so a zero-usage first line no longer claims a message id and hides a later line with real usage. Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> Co-authored-by: GPT-6 Astra <noreply@openai.com> * Index Claude transcripts so the agents refresh reads only what was appended (#8313) omarchy-agent-usage-claude re-parsed every line of every transcript under ~/.claude/projects on each refresh: no mtime cutoff, no memory of the last pass. The agents widget is on by default and ticks every 15 minutes, so the cost grew for the life of the machine. After one month here that was 803 files, 640 MB, 127k lines and 57k JSON parses per tick, about 1 core-second, pushed through the page cache every quarter hour forever. Keep a per-file index next to the scan cache: the unique usage records already parsed out of each transcript and the byte offset they end at. A file whose size and mtime match is not opened; a file that grew is read from the stored offset; a file that shrank or was rewritten is read from the start. --force drops the index and rescans from scratch. The summary is built from the indexed records in the same directory order the walk always used. That matters: when a resumed session carries earlier messages, the same message id appears in two files with different usage, and the first file visited wins. 91 ids differed on this machine; sorting the walk moved one model's output total by 25k tokens. Output is now byte-identical to the previous scan on a frozen copy of the corpus, cold, warm, and after an append. Warm refresh: 1.0 s -> 0.10 s of CPU, of which the scan itself is 70 ms; the index for this corpus is 2.9 MB. Adapted: - Rebased onto #10606: the highest-output rule for streamed messages now lives where the index parses records, and decides between files too. - The index records the timezone it was written in, and a change rereads every transcript, since its records hold local days. - A file only counts as appended to when its inode and the hash of what was already read still match, so a transcript replaced by a larger one, or rewritten in place, is read from the start. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Label a Claude Team seat by its subscription, not its rate-limit tier (#11109) The collector built the plan label from the OAuth rateLimitTier first, so a Team premium seat, which runs on default_claude_max_5x, showed in the agents panel as "Max 5x". Lead with subscriptionType and keep the multiplier as its qualifier: Max still reads "Max 5x", a Team seat reads "Team 5x". Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * Label the Claude plan from the profile the CLI refreshes (#7225) Adapted: the profile is found the same way current_account_id() finds it, now shared as profile_path(): ~/.claude.json for the default home, the home's own .claude.json otherwise. The original fell back to ~/.claude.json for any home without CLAUDE_CONFIG_DIR set, so a secondary account read the primary's tier. The profile's tier also keeps the subscription in the label, so a Team seat stays "Team" (#11109). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Call a lapsed Claude access token paused, not signed out (#8093) * Refresh Claude usage after the clock moves backwards (#9956) * Bound unreadable Claude transcript warnings (#12414) * Count Claude usage from opencode v2 sessions (#13894) * Reload agent usage records when an inotify watch fails to rearm (#10067) * Reload agent usage records after each update run instead of on a timer Rather than #10067's two-minute timer per record, reload every record when the omarchy-agent-usage-update process exits, the moment its files can have been replaced. A reload that finds a file unchanged keeps its record, so the panel isn't stirred up by identical data. The grep test now runs the QML functions. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Show the agent status when the trouble line has no help text (#8497) * Clear stale agent login guidance after a successful probe (#8892) * Clear the Grok login hint after a successful probe #8892 cleared the default login hint after a successful probe in the Claude and Codex collectors; Grok's collector had the same stale hint. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Read Fireworks credentials from pi's auth.json (#7455) The Fireworks collector skipped pi, Omarchy's default agent, when walking its credential ladder, so a machine signed in to Fireworks only through pi (/login fireworks) never showed the tab. Insert the key pi stores in $PI_CODING_AGENT_DIR/auth.json (default ~/.pi/agent) between the firectl auth.ini and the opencode fallback. pi keys can be literals, $ENV_VAR/${ENV_VAR} references, or !command shell lookups. The collector resolves the first two; command lookups stay pi-only and are skipped rather than sent to the API verbatim. * Call a lapsed Grok access token paused, not signed out Grok's access token lives six hours and Grok mints a new one from its refresh token whenever it starts, so a lapsed one is routine. Reporting it as an expired sign-in made the panel offer Sign-in required several times a day, sending people through grok login for nothing. With a refresh token present it now reads as paused, like Claude's. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Keep showing Grok's last limits while it sits idle While Grok hasn't run, nothing on the machine has spent its allowance, so with a refresh token on hand the last numbers still stand: they show as current rather than dimmed under a status line. A weekly window that reset in the meantime starts over at 0%, a whole number of weeks on. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Ask for a Grok sign-in once its refresh token is past 30 days A refresh token older than Grok's 30-day sign-in can't renew anything, so the panel offers Sign-in required again instead of showing the last limits as current. With nothing cached yet it says to start Grok, rather than showing an empty section without a word. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Check both ends of what the Claude index read before resuming a transcript A transcript rewritten in place could grow and change only after its first kilobytes, and the index took it for an append. It now compares the last kilobytes before the resume point too. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Simplify the agent usage collectors - Codex: pass the forced-scan choice down instead of a module global, make the per-file reader's cache arguments required, shrink the cache record check, and drop guards for shapes that can't occur: an empty launcher path, realpath raising, mise itself being a lazy launcher, multi-line `mise which` output, and probing without a temp file for stderr. - Claude: decide an append by the digest of both ends of what was read alone; the inode and mtime checks it made redundant are gone. - Snapshot: the device id falls back to the hostname, which always exists. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Share fixture setup in the agent usage scanner tests Every fixture home lives under one scratch directory with a single cleanup trap, instead of a trap rewritten with a longer list for each new home, and the Codex test builds its signed-in homes with one helper. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Treat a replaced Claude transcript as new even when its ends match Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Probe Codex without its error text when there's no temporary space Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: tossbaws <17258053+tossbaws@users.noreply.github.com> Co-authored-by: surim0n <suritech@gmail.com> Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com> Co-authored-by: anonwurcod <anonwurcod@proton.me> Co-authored-by: Kevin Rajan <7121943+kvnloo@users.noreply.github.com> Co-authored-by: Nate Ashby <nate.ashby11@gmail.com> Co-authored-by: Aris Gysel <aris.gysel@me.com> Co-authored-by: Brams <76213579+Brams-s@users.noreply.github.com> Co-authored-by: This_Is_NPC <gabrielfollone27@gmail.com> Co-authored-by: sanjyay <102979855+sanjyay@users.noreply.github.com> Co-authored-by: PapistProtocol <12738904+PapistProtocol@users.noreply.github.com> Co-authored-by: steez <stevedimakos97@gmail.com> Co-authored-by: GPT-6 Astra <noreply@openai.com> Co-authored-by: Ryan Yogan <ryanyogan@gmail.com> Co-authored-by: Oli Denton <41393837+omdenton@users.noreply.github.com> Co-authored-by: Igor Kramar <i@ikramar.ru> Co-authored-by: Martin Eidensten <martin@meibe.se> Co-authored-by: Romain Perron <rdj.perron@gmail.com> Co-authored-by: Omarchy Contributor <contributor@users.noreply.github.com> Co-authored-by: manuaudio <manu@arimaka.com> Co-authored-by: Tyler South <tsouth2@gmail.com> Co-authored-by: whathek <Hek846@users.noreply.github.com> Co-authored-by: Ty Richards <me@tyrichards.com>
1078 lines
40 KiB
Python
Executable File
1078 lines
40 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
|
|
import tomllib
|
|
from datetime import datetime, timedelta, timezone
|
|
from pathlib import Path
|
|
|
|
AGENT_ID = "codex"
|
|
AGENT_NAME = "Codex"
|
|
AUTH_HELP = "Run `codex login` to authenticate."
|
|
|
|
# session_meta.model_provider for the built-in OpenAI backend, which is the
|
|
# only one this subscription pays for. Any other id points somewhere else:
|
|
# `--oss`'s ollama or lmstudio, or a custom provider in config.toml.
|
|
NATIVE_PROVIDER = "openai"
|
|
|
|
# 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
|
|
|
|
# Per-file totals of native Codex sessions survive between runs, so a rescan
|
|
# only reads the session files that changed. Bumping the schema retires every
|
|
# record written by an older layout.
|
|
FILE_CACHE_SCHEMA = 1
|
|
|
|
|
|
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"))
|
|
|
|
|
|
def mise_shims_dir():
|
|
data_home = ENV.get("MISE_DATA_DIR") or os.path.join(
|
|
ENV.get("XDG_DATA_HOME") or os.path.join(str(Path.home()), ".local", "share"), "mise")
|
|
return os.path.join(data_home, "shims")
|
|
|
|
|
|
def is_lazy_launcher(path):
|
|
# omarchy-mise-install writes a regular file whose first step is `mise use
|
|
# -g`, and mise's shims exec `mise x`; running either installs the tool it
|
|
# wraps. A shim is a symlink to the mise binary itself, so the symlink
|
|
# check comes after the shim checks: a symlink elsewhere is the user's own
|
|
# binary.
|
|
if os.path.dirname(os.path.abspath(path)) == mise_shims_dir():
|
|
return True
|
|
if os.path.islink(path):
|
|
mise = find_command("mise")
|
|
return bool(mise) and os.path.realpath(path) == os.path.realpath(mise)
|
|
try:
|
|
with open(path, "rb") as handle:
|
|
head = handle.read(4096)
|
|
return b"mise use" in head or b"mise x " in head
|
|
except OSError:
|
|
return False
|
|
|
|
|
|
def find_codex_binary():
|
|
for directory in ENV.get("PATH", "").split(os.pathsep):
|
|
candidate = os.path.join(directory, "codex")
|
|
if os.path.isfile(candidate) and os.access(candidate, os.X_OK) and not is_lazy_launcher(candidate):
|
|
return candidate
|
|
# `mise which` only prints a binary that is already installed, so it is the
|
|
# safe way to resolve a codex managed by mise without running the launcher.
|
|
# `mise which` can resolve a configured `latest` over the network; the usage
|
|
# probe stays offline.
|
|
mise = find_command("mise")
|
|
if mise:
|
|
try:
|
|
out = subprocess.run([mise, "which", "codex"], capture_output=True, text=True, env=dict(ENV, MISE_OFFLINE="1"), timeout=10)
|
|
resolved = out.stdout.strip()
|
|
if out.returncode == 0 and os.path.isfile(resolved) and not is_lazy_launcher(resolved):
|
|
return resolved
|
|
except Exception:
|
|
pass
|
|
return None
|
|
|
|
|
|
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()
|
|
|
|
|
|
# One session file's totals, the unit the file cache stores: day -> model ->
|
|
# [input, output, cacheRead, cacheWrite, prompts]. Day-keyed and with no
|
|
# "today" in it, so a record stays true however long it is kept.
|
|
def new_file_record(stat_result):
|
|
return {"mtime": stat_result.st_mtime, "size": stat_result.st_size, "days": {}}
|
|
|
|
|
|
def record_usage(record, day, model, input_tokens, output_tokens, cache_read, cache_write):
|
|
totals = record["days"].setdefault(day, {}).setdefault(model, [0, 0, 0, 0, 0])
|
|
totals[0] += input_tokens
|
|
totals[1] += output_tokens
|
|
totals[2] += cache_read
|
|
totals[3] += cache_write
|
|
totals[4] += 1
|
|
|
|
|
|
def merge_file_record(session_key, record):
|
|
"""Fold one file's totals, freshly read or replayed from cache, into the run."""
|
|
global today_prompts, today_total_tokens, total_prompts
|
|
for day, models in record["days"].items():
|
|
for model, totals in models.items():
|
|
input_tokens, output_tokens, cache_read, cache_write, prompts = totals
|
|
total = input_tokens + output_tokens + cache_read + cache_write
|
|
total_prompts += prompts
|
|
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 += prompts
|
|
today_sessions.add(session_key)
|
|
today_total_tokens += total
|
|
today_tokens_by_model[model] = today_tokens_by_model.get(model, 0) + total
|
|
|
|
|
|
def add_usage(day, session_key, model, input_tokens, output_tokens, cache_read, cache_write):
|
|
record = {"days": {}}
|
|
record_usage(record, day, model, input_tokens, output_tokens, cache_read, cache_write)
|
|
merge_file_record(session_key, record)
|
|
|
|
|
|
def pi_parent_session(path):
|
|
try:
|
|
with open(path, encoding="utf-8", errors="replace") as session_file:
|
|
header = json.loads(session_file.readline())
|
|
except (OSError, json.JSONDecodeError):
|
|
return None
|
|
|
|
if not isinstance(header, dict):
|
|
return None
|
|
parent = header.get("parentSession")
|
|
if not isinstance(parent, str) or not parent:
|
|
return None
|
|
return os.path.abspath(os.path.expanduser(parent))
|
|
|
|
|
|
# `omp --profile=<name>` (and pi's equivalent) relocates the whole agent tree
|
|
# under <base>/profiles/<name>/, so a subscription driven entirely through a
|
|
# profile leaves the default root empty and its usage uncounted. Each root is
|
|
# scanned separately and sessions are keyed by path, so a profile adds
|
|
# sessions rather than double-counting the default one.
|
|
def pi_session_roots():
|
|
roots = []
|
|
for base in (Path.home() / ".pi", Path.home() / ".omp"):
|
|
roots.append(base / "agent" / "sessions")
|
|
profiles = base / "profiles"
|
|
try:
|
|
# Sorted so the scan order does not depend on directory order.
|
|
roots.extend(sorted(child / "agent" / "sessions" for child in profiles.iterdir() if child.is_dir()))
|
|
except OSError:
|
|
# No profiles directory, or it is unreadable: the default root stands.
|
|
pass
|
|
return roots
|
|
|
|
|
|
def scan_pi_sessions():
|
|
roots = pi_session_roots()
|
|
rg = find_command("rg") or "rg"
|
|
session_keys = {}
|
|
today_session_keys = {}
|
|
for root in roots:
|
|
if not root.exists():
|
|
continue
|
|
try:
|
|
proc = subprocess.Popen(
|
|
# --no-ignore: these are data files, not a source tree. When $HOME is
|
|
# itself a git checkout (a common dotfiles setup), the parent repo's
|
|
# ignore rules would otherwise silently exclude every session file and
|
|
# the scan would count zero usage.
|
|
[rg, "--json", "--no-ignore", "-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", "")
|
|
raw_path = event.get("data", {}).get("path", {}).get("text", "pi-session")
|
|
path = os.path.abspath(os.path.expanduser(raw_path))
|
|
entry = json.loads(line)
|
|
except Exception:
|
|
continue
|
|
|
|
if entry.get("type") != "message":
|
|
continue
|
|
message_id = str(entry.get("id") or "")
|
|
message_timestamp = str(entry.get("timestamp") or "")
|
|
# Forks retain both fields; IDs alone can collide across sessions.
|
|
if message_id and message_timestamp:
|
|
message_key = ("entry", message_id, message_timestamp)
|
|
else:
|
|
message_key = ("path", path, message_id)
|
|
duplicate = message_key in seen_pi_messages
|
|
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_keys.setdefault(path, set()).add(message_key)
|
|
if day == today:
|
|
today_session_keys.setdefault(path, set()).add(message_key)
|
|
if duplicate:
|
|
continue
|
|
add_usage(day, path, model_name(message.get("model")), input_tokens, output_tokens, cache_read, cache_write)
|
|
|
|
try:
|
|
proc.wait(timeout=1)
|
|
except Exception:
|
|
proc.kill()
|
|
|
|
# Resolve session ownership after scanning so rg's file order cannot affect it.
|
|
parents = {path: pi_parent_session(path) for path in session_keys}
|
|
total_sessions.difference_update(session_keys)
|
|
today_sessions.difference_update(session_keys)
|
|
for path, keys in session_keys.items():
|
|
inherited = set()
|
|
visited = {path}
|
|
parent = parents.get(path)
|
|
while parent and parent not in visited:
|
|
visited.add(parent)
|
|
inherited.update(session_keys.get(parent, ()))
|
|
parent = parents.get(parent)
|
|
|
|
unique_keys = keys - inherited
|
|
if unique_keys:
|
|
total_sessions.add(path)
|
|
if unique_keys & today_session_keys.get(path, set()):
|
|
today_sessions.add(path)
|
|
|
|
|
|
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(cached_files, scanned_files):
|
|
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:
|
|
stat_result = path.stat()
|
|
except OSError:
|
|
continue
|
|
if stat_result.st_mtime >= cutoff:
|
|
files.append((path, stat_result))
|
|
|
|
for path, stat_result in files:
|
|
key = str(path)
|
|
record = reusable_record(cached_files, key, stat_result)
|
|
if record is None:
|
|
record = read_native_codex_session(path, stat_result)
|
|
if record["mtime"] is not None:
|
|
scanned_files[key] = record
|
|
merge_file_record(key, record)
|
|
|
|
|
|
def read_native_codex_session(path, stat_result):
|
|
record = new_file_record(stat_result)
|
|
current_model = "codex"
|
|
seen_meta = False
|
|
previous_total_usage = None
|
|
try:
|
|
with path.open(errors="replace") as handle:
|
|
for raw in handle:
|
|
# Cheap pre-filter before JSON parsing keeps files with unrelated
|
|
# lines inexpensive. session_meta still has to get through: it says
|
|
# which provider served the rollout.
|
|
if '"token_count"' not in raw and '"turn_context"' not in raw and '"session_meta"' not in raw:
|
|
continue
|
|
try:
|
|
entry = json.loads(raw)
|
|
except Exception:
|
|
continue
|
|
if entry.get("type") == "session_meta":
|
|
# A forked rollout copies its parent's session_meta after its own,
|
|
# so only the first one, as Codex reads it, says who served this.
|
|
if seen_meta:
|
|
continue
|
|
seen_meta = True
|
|
# Codex CLI happily fronts any OpenAI-compatible backend (`--oss`,
|
|
# or a custom model_provider in config.toml). Those turns bill the
|
|
# local box or a third party, never this subscription, so drop the
|
|
# whole rollout the way the pi and opencode scans drop foreign
|
|
# providers.
|
|
meta = entry.get("payload") or {}
|
|
provider = str(meta.get("model_provider") or "")
|
|
# Rollouts predating the field carry no provider at all; count
|
|
# those rather than silently lose the history.
|
|
if provider and provider != NATIVE_PROVIDER:
|
|
break
|
|
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
|
|
# Quota updates can repeat the previous request's last_token_usage.
|
|
# Only an unchanged cumulative snapshot proves this is a repeat:
|
|
# separate requests may have identical last-token counts, and older
|
|
# records may not carry cumulative counters at all. Keep this state
|
|
# per rollout and accept changed counters, including counter resets.
|
|
total_usage = info.get("total_token_usage")
|
|
if isinstance(total_usage, dict) and total_usage:
|
|
if total_usage == previous_total_usage:
|
|
continue
|
|
previous_total_usage = total_usage
|
|
day = local_day(entry.get("timestamp") or stat_result.st_mtime)
|
|
record_usage(record, day, current_model, input_tokens, output_tokens, cache_read, cache_write)
|
|
except Exception:
|
|
# Keep what was read before the failure, but caching a partial record would
|
|
# hide the rest until the file changes, so hand back one the cache refuses.
|
|
record["mtime"] = None
|
|
return record
|
|
|
|
|
|
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_digest():
|
|
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.
|
|
return hashlib.sha1((str(Path.home()) + "\n" + str(codex_home) + "\n" + str(db)).encode("utf-8")).hexdigest()[:16]
|
|
|
|
|
|
def scan_cache_paths():
|
|
digest = scan_cache_digest()
|
|
root = cache_root()
|
|
return root / f"codex-scan-{digest}.json", root / f"codex-scan-{digest}.lock"
|
|
|
|
|
|
# Records hold local days, so they are only true in the timezone that wrote them.
|
|
def file_cache_zone():
|
|
return [os.environ.get("TZ"), list(time.tzname), time.timezone, time.altzone]
|
|
|
|
|
|
def file_cache_path():
|
|
return cache_root() / f"codex-files-{scan_cache_digest()}.json"
|
|
|
|
|
|
def valid_file_record(value):
|
|
try:
|
|
counts = [totals for models in value["days"].values() for totals in models.values()]
|
|
return (
|
|
type(value["mtime"]) in (int, float) and type(value["size"]) is int
|
|
and all(len(totals) == 5 and all(type(n) is int for n in totals) for totals in counts)
|
|
)
|
|
except Exception:
|
|
return False
|
|
|
|
|
|
def read_file_cache():
|
|
"""Per-file totals of native sessions from earlier scans, keyed by path.
|
|
|
|
Session files are append-only, so a file whose mtime and size both match the
|
|
record cannot have grown a turn since: its totals can be replayed instead of
|
|
read. That is what keeps a refresh off the gigabytes of history it already
|
|
counted. Unlike the aggregate scan cache these records carry no "today" and
|
|
never expire on their own -- only the file they describe or a change of
|
|
timezone can invalidate one, so midnight and the clock moving leave them true.
|
|
|
|
pi and omp sessions stay out of it: a forked pi session repeats its parent's
|
|
messages, so they are only deduplicated across the whole tree at once.
|
|
"""
|
|
try:
|
|
payload = json.loads(file_cache_path().read_text(encoding="utf-8"))
|
|
except Exception:
|
|
return {}
|
|
if not isinstance(payload, dict) or payload.get("schemaVersion") != FILE_CACHE_SCHEMA:
|
|
return {}
|
|
if payload.get("zone") != file_cache_zone():
|
|
return {}
|
|
files = payload.get("files")
|
|
if not isinstance(files, dict):
|
|
return {}
|
|
return {path: record for path, record in files.items() if valid_file_record(record)}
|
|
|
|
|
|
def write_file_cache(files):
|
|
# Only the files this scan saw are written back, so deleted sessions and
|
|
# sessions that aged out of the 30-day window leave the cache instead of
|
|
# growing it forever.
|
|
try:
|
|
write_json(file_cache_path(), {"schemaVersion": FILE_CACHE_SCHEMA, "zone": file_cache_zone(), "files": files})
|
|
except Exception as exc:
|
|
print(f"omarchy-agent-usage-codex: could not write file cache ({exc})", file=sys.stderr)
|
|
|
|
|
|
def reusable_record(cached_files, key, stat_result):
|
|
record = cached_files.get(key)
|
|
if record is None:
|
|
return None
|
|
if record["mtime"] != stat_result.st_mtime or record["size"] != stat_result.st_size:
|
|
return None
|
|
return record
|
|
|
|
|
|
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.
|
|
# Version 2 invalidates totals counted before native notification deduplication.
|
|
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") != 2:
|
|
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": 2, "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,
|
|
}
|
|
|
|
|
|
# A forced scan (no reuse at all) reads every session file again.
|
|
def run_local_scans(max_age):
|
|
cached_files = read_file_cache() if max_age > 0 else {}
|
|
scanned_files = {}
|
|
scan_pi_sessions()
|
|
scan_native_codex_sessions(cached_files, scanned_files)
|
|
complete = scan_opencode_sessions()
|
|
write_file_cache(scanned_files)
|
|
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(max_age)
|
|
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(max_age)
|
|
# 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_send(proc, payload, operation):
|
|
"""Write one JSON-RPC frame. Raise immediately if the app-server is gone."""
|
|
try:
|
|
proc.stdin.write(json.dumps(payload) + "\n")
|
|
proc.stdin.flush()
|
|
except OSError as exc:
|
|
try:
|
|
proc.stdin.close()
|
|
except OSError:
|
|
pass
|
|
raise RuntimeError(f"Codex app-server exited before {operation}") from exc
|
|
|
|
|
|
def rpc_request(proc, request_id, method, params=None, timeout=8):
|
|
payload = {"id": request_id, "method": method, "params": params or {}}
|
|
rpc_send(proc, payload, method)
|
|
# Read the raw fd and split lines ourselves. The app-server can emit
|
|
# notifications in the same write as a reply; a buffered readline() would
|
|
# pull them all into Python's buffer where select() cannot see them, and
|
|
# the reply would sit there until the deadline.
|
|
fd = proc.stdout.fileno()
|
|
pending = getattr(proc, "_rpc_pending", b"")
|
|
deadline = time.monotonic() + timeout
|
|
try:
|
|
while True:
|
|
while b"\n" in pending:
|
|
line, pending = pending.split(b"\n", 1)
|
|
try:
|
|
message = json.loads(line)
|
|
except Exception:
|
|
continue
|
|
if isinstance(message, dict) and message.get("id") == request_id:
|
|
return message
|
|
remaining = deadline - time.monotonic()
|
|
if remaining <= 0:
|
|
break
|
|
ready, _, _ = select.select([fd], [], [], min(0.25, remaining))
|
|
if not ready:
|
|
if proc.poll() is not None:
|
|
raise RuntimeError(f"Codex app-server exited before {method}")
|
|
continue
|
|
chunk = os.read(fd, 65536)
|
|
if not chunk:
|
|
raise RuntimeError(f"Codex app-server exited before {method}")
|
|
pending += chunk
|
|
finally:
|
|
proc._rpc_pending = pending
|
|
# Process still up but silent: keep the method name searchable, but name the stall.
|
|
raise TimeoutError(f"Codex app-server did not answer {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 codex_rpc_help(proc, stderr_file, exc):
|
|
"""Explain an app-server RPC failure: a live server's stall, a dead one's own
|
|
stderr, or the login hint when it died saying nothing."""
|
|
try:
|
|
proc.wait(timeout=1)
|
|
except Exception:
|
|
pass
|
|
|
|
text = str(exc).strip()
|
|
if proc.poll() is None:
|
|
# Still up: stall or unparseable payload. Not an auth problem, and any
|
|
# stderr so far is its logging, not why it stopped.
|
|
return text[:300] or AUTH_HELP
|
|
|
|
try:
|
|
stderr_file.seek(0)
|
|
detail = stderr_file.read()
|
|
except Exception:
|
|
detail = ""
|
|
lines = [line.strip() for line in detail.splitlines() if line.strip()]
|
|
if lines:
|
|
# The fatal error comes last, after whatever the CLI logged on the way up.
|
|
return f"codex app-server exited: {' '.join(lines)[-275:]}"
|
|
# Exited cleanly with empty stderr — typical "not logged in" shape.
|
|
return AUTH_HELP
|
|
|
|
|
|
def reset_credits(result):
|
|
granted = (result.get("rateLimitResetCredits") or {}).get("credits") or []
|
|
available = [c for c in granted if isinstance(c, dict) and c.get("status") == "available"]
|
|
if not available:
|
|
return None
|
|
expiries = [number(c.get("expiresAt")) for c in available if number(c.get("expiresAt")) > 0]
|
|
return {
|
|
"available": len(available),
|
|
"nextExpiresAt": datetime.fromtimestamp(min(expiries), timezone.utc).isoformat() if expiries else "",
|
|
}
|
|
|
|
|
|
# Keyring and auto stores keep credentials outside auth.json, so only the file
|
|
# store can be known to be empty without asking Codex.
|
|
def has_codex_credentials(codex_home, env):
|
|
if env.get("CODEX_ACCESS_TOKEN"):
|
|
return True
|
|
try:
|
|
with (codex_home / "config.toml").open("rb") as handle:
|
|
config = tomllib.load(handle)
|
|
except FileNotFoundError:
|
|
config = {}
|
|
except Exception:
|
|
# A config this can't read may still name another store; let Codex say.
|
|
return True
|
|
if str(config.get("cli_auth_credentials_store") or "file") != "file":
|
|
return True
|
|
return (codex_home / "auth.json").is_file()
|
|
|
|
|
|
def fetch_codex_rpc(home=None):
|
|
result = {"limits": [], "tierLabel": "", "usageStatusText": "", "authHelpText": AUTH_HELP}
|
|
env = dict(ENV, CODEX_HOME=str(home)) if home else ENV
|
|
codex = find_codex_binary()
|
|
if not codex:
|
|
result["usageStatusText"] = "Codex unavailable"
|
|
result["authHelpText"] = "codex not found in PATH"
|
|
return result
|
|
|
|
# No auth.json under the file store and no CODEX_ACCESS_TOKEN means
|
|
# account/read can only fail — and starting the app-server is not free: it
|
|
# syncs the plugin list, a git fetch per refresh that also leaves
|
|
# .tmp/git-* folders behind. The home is the one being probed, so a
|
|
# signed-out secondary account reads as signed out even while ~/.codex
|
|
# holds a login.
|
|
if not has_codex_credentials(Path(env.get("CODEX_HOME") or (Path.home() / ".codex")), env):
|
|
result["usageStatusText"] = "Waiting for auth"
|
|
result["authHelpText"] = AUTH_HELP
|
|
return result
|
|
|
|
# stderr goes to a regular file, not a PIPE: an unread pipe can fill and
|
|
# deadlock a chatty app-server. It is only read once the process has exited.
|
|
# Without temporary space the probe still runs, just without its error text.
|
|
try:
|
|
stderr_file = tempfile.TemporaryFile("w+", errors="replace")
|
|
except OSError:
|
|
stderr_file = None
|
|
try:
|
|
proc = subprocess.Popen(
|
|
[codex, "-s", "read-only", "-a", "on-request", "app-server"],
|
|
stdin=subprocess.PIPE,
|
|
stdout=subprocess.PIPE,
|
|
stderr=stderr_file or subprocess.DEVNULL,
|
|
text=True,
|
|
env=env,
|
|
)
|
|
except Exception as exc:
|
|
if stderr_file:
|
|
stderr_file.close()
|
|
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)
|
|
rpc_send(proc, {"method": "initialized", "params": {}}, "initialized")
|
|
limits_msg = rpc_request(proc, 2, "account/rateLimits/read", timeout=8)
|
|
|
|
# A home nobody is signed in to answers with an error rather than limits.
|
|
# Say so the same way the Claude collector does, so the panel offers to
|
|
# sign that account in again.
|
|
error = limits_msg.get("error")
|
|
if isinstance(error, dict):
|
|
message = str(error.get("message") or "")
|
|
if "auth" in message.lower():
|
|
result["usageStatusText"] = "Waiting for auth"
|
|
result["authHelpText"] = AUTH_HELP
|
|
else:
|
|
result["usageStatusText"] = "Codex limits unavailable"
|
|
result["authHelpText"] = message
|
|
return result
|
|
limits = (limits_msg.get("result") or {}).get("rateLimits") or {}
|
|
result["authHelpText"] = ""
|
|
plan = limits.get("planType") or ""
|
|
|
|
# The limits name the plan themselves. account/read is only a fallback
|
|
# for when they don't, and never a reason to lose them: Codex 0.158's
|
|
# app-server can leave it unanswered for good.
|
|
if not plan:
|
|
try:
|
|
account_msg = rpc_request(proc, 3, "account/read", timeout=2)
|
|
account = (account_msg.get("result") or {}).get("account") or {}
|
|
plan = account.get("planType") or account.get("type") or ""
|
|
except (TimeoutError, RuntimeError):
|
|
pass
|
|
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)
|
|
|
|
# OpenAI hands out free "full reset" credits that wipe the rate limits on
|
|
# demand and lapse after a month. Worth knowing about before a limit bites.
|
|
credits = reset_credits(limits_msg.get("result") or {})
|
|
if credits:
|
|
result["resetCredits"] = credits
|
|
except Exception as exc:
|
|
result["usageStatusText"] = "Codex limits unavailable"
|
|
result["authHelpText"] = codex_rpc_help(proc, stderr_file, exc)
|
|
finally:
|
|
try:
|
|
proc.terminate()
|
|
proc.wait(timeout=1)
|
|
except Exception:
|
|
try:
|
|
proc.kill()
|
|
except Exception:
|
|
pass
|
|
if stderr_file:
|
|
stderr_file.close()
|
|
return result
|
|
|
|
|
|
# The accounts `omarchy agent account` registered, read straight from its
|
|
# registry. Only a registry holding a second account matters: with one, the
|
|
# record is exactly what it always was.
|
|
def registered_accounts():
|
|
state = Path(os.environ.get("XDG_STATE_HOME") or (Path.home() / ".local" / "state"))
|
|
try:
|
|
registry = json.loads((state / "omarchy" / "agents" / "accounts" / "codex.json").read_text(encoding="utf-8"))
|
|
except Exception:
|
|
return []
|
|
accounts = [a for a in registry.get("accounts") or [] if isinstance(a, dict) and a.get("id")]
|
|
if len(accounts) < 2:
|
|
return []
|
|
active_id = registry.get("active") or accounts[0]["id"]
|
|
if not any(a["id"] == active_id for a in accounts):
|
|
active_id = accounts[0]["id"]
|
|
for account in accounts:
|
|
account["active"] = account["id"] == active_id
|
|
account["switch"] = {
|
|
"mode": "auto" if registry.get("switch") == "auto" else "manual",
|
|
"threshold": registry.get("threshold") or 95,
|
|
}
|
|
return accounts
|
|
|
|
|
|
# Each account's app-server runs in that account's home, so it answers for that
|
|
# login and refreshes that login's tokens there.
|
|
def account_limits(account):
|
|
# The primary is ~/.codex by definition, not whatever CODEX_HOME this run
|
|
# happened to inherit — that would report another account's limits as Main's.
|
|
rpc = fetch_codex_rpc(account.get("home") or (Path.home() / ".codex"))
|
|
return {
|
|
"id": account["id"],
|
|
"label": str(account.get("label") or account["id"]),
|
|
"email": str(account.get("email") or ""),
|
|
"plan": rpc["tierLabel"] or str(account.get("plan") or ""),
|
|
"active": account["active"],
|
|
"primary": bool(account.get("primary")),
|
|
"limits": rpc["limits"],
|
|
"stale": not rpc["limits"],
|
|
"usageStatusText": rpc["usageStatusText"],
|
|
"authHelpText": rpc["authHelpText"],
|
|
"resetCredits": rpc.get("resetCredits"),
|
|
}
|
|
|
|
|
|
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)
|
|
registered = registered_accounts()
|
|
accounts = [account_limits(account) for account in registered]
|
|
current = next((a for a in accounts if a["active"]), None)
|
|
if current:
|
|
# The top-level fields keep describing the account new sessions use.
|
|
rpc = {
|
|
"limits": current["limits"],
|
|
"tierLabel": current["plan"],
|
|
"usageStatusText": current["usageStatusText"],
|
|
"authHelpText": current["authHelpText"],
|
|
}
|
|
if current.get("resetCredits"):
|
|
rpc["resetCredits"] = current["resetCredits"]
|
|
else:
|
|
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)
|
|
if accounts:
|
|
record["accountSwitch"] = registered[0]["switch"]
|
|
record["accounts"] = accounts
|
|
print(json.dumps(record, separators=(",", ":")))
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|