#!/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=` (and pi's equivalent) relocates the whole agent tree # under /profiles//, 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()