Files
omarchy/bin/omarchy-agent-usage-codex
+14 75250d37ac Fix Codex limits, Claude counting, and agent usage reliability from community PRs (#14049)
* 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>
2026-10-02 22:03:07 -04:00

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()