From 2c02203561d43ad87d4973282aa91eaf7287d0ba Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=AC=A7=E9=98=B3=E7=9D=BF?= Date: Thu, 4 Jun 2026 19:47:40 +0800 Subject: [PATCH] feat(skill): add codex usage uploader v2 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 新增独立的 codex-usage-uploader-v2 skill,用于上传脱敏后的 Codex 使用元数据,并通过 total_token_usage 快照去重避免 token 重复统计。 1. 新增 V2 skill 的 SKILL.md、OpenAI UI 元数据、ingest API 参考、上传脚本和回归测试。 2. 保留原 codex-usage-uploader skill 不变,README 增加 V2 安装路径和使用说明。 3. 已验证 V2 脚本测试通过,并通过 skill quick_validate 校验。 --- README.md | 22 + skills/codex-usage-uploader-v2/SKILL.md | 76 ++ .../agents/openai.yaml | 4 + .../references/ingest_api.md | 73 ++ .../scripts/codex_usage_uploader.py | 762 ++++++++++++++++++ .../scripts/test_codex_usage_uploader.py | 373 +++++++++ 6 files changed, 1310 insertions(+) create mode 100644 skills/codex-usage-uploader-v2/SKILL.md create mode 100644 skills/codex-usage-uploader-v2/agents/openai.yaml create mode 100644 skills/codex-usage-uploader-v2/references/ingest_api.md create mode 100644 skills/codex-usage-uploader-v2/scripts/codex_usage_uploader.py create mode 100644 skills/codex-usage-uploader-v2/scripts/test_codex_usage_uploader.py diff --git a/README.md b/README.md index 9320fe6..0bd2555 100644 --- a/README.md +++ b/README.md @@ -2,6 +2,11 @@ Codex skill for uploading sanitized local Codex Desktop or CLI usage metadata to an ingest API. It scans local Codex JSONL session logs incrementally, filters private content, sends metadata batches with Bearer authentication, and stores local progress state. +This repository contains two skill paths: + +- `skills/codex-usage-uploader`: original uploader. +- `skills/codex-usage-uploader-v2`: V2 uploader with token snapshot de-duplication. Use this when dashboards sum uploaded token fields and must avoid over-counting Codex cumulative `total_token_usage` or re-emitted `token_count` snapshots. + ## Install Install with Codex's built-in skill installer: @@ -14,6 +19,15 @@ python C:\Users\\.codex\skills\.system\skill-installer\scripts\install-ski Restart Codex after installation so the skill is discovered. +Install the V2 skill: + +```powershell +python C:\Users\\.codex\skills\.system\skill-installer\scripts\install-skill-from-github.py ` + --repo HardToFd/codex-usage-uploader-skill ` + --path skills/codex-usage-uploader-v2 ` + --branch develop +``` + ## Configure The uploader needs three values: @@ -30,6 +44,12 @@ Run a dry-run first: python C:\Users\\.codex\skills\codex-usage-uploader\scripts\codex_usage_uploader.py --dry-run ``` +For V2, use the V2 skill directory: + +```powershell +python C:\Users\\.codex\skills\codex-usage-uploader-v2\scripts\codex_usage_uploader.py --dry-run +``` + Then run the upload: ```powershell @@ -53,3 +73,5 @@ The uploader does not send raw user messages, assistant replies, reasoning conte See [skills/codex-usage-uploader/references/ingest_api.md](skills/codex-usage-uploader/references/ingest_api.md). +For V2 token snapshot de-duplication semantics, see [skills/codex-usage-uploader-v2/references/ingest_api.md](skills/codex-usage-uploader-v2/references/ingest_api.md). + diff --git a/skills/codex-usage-uploader-v2/SKILL.md b/skills/codex-usage-uploader-v2/SKILL.md new file mode 100644 index 0000000..8754ccc --- /dev/null +++ b/skills/codex-usage-uploader-v2/SKILL.md @@ -0,0 +1,76 @@ +--- +name: codex-usage-uploader-v2 +description: Upload sanitized Codex Desktop or CLI usage metadata to a configured ingest API using the V2 token snapshot de-duplication contract. Use when Codex needs to configure, run, backfill, troubleshoot, or automate local Codex session JSONL collection while avoiding duplicate token usage from cumulative `total_token_usage` and re-emitted `token_count` events. +--- + +# Codex Usage Uploader V2 + +## Overview + +Use this skill to collect local Codex session JSONL usage metadata and upload it to a dashboard ingest API. V2 keeps the original privacy boundary and adds a strict token snapshot de-duplication contract so dashboards can safely sum uploaded `token` fields without counting cumulative Codex snapshots twice. + +## Quick Start + +Run from the skill directory or pass the absolute script path: + +```bash +python scripts/codex_usage_uploader.py --endpoint https://collector.example.com/api/codex/usage-v2 --source-name alice --token +``` + +Environment variables are supported: + +```bash +set CODEX_USAGE_INGEST_URL=https://collector.example.com/api/codex/usage-v2 +set CODEX_USAGE_SOURCE_NAME=alice +set CODEX_USAGE_BEARER_TOKEN= +python scripts/codex_usage_uploader.py +``` + +Useful options: + +```bash +python scripts/codex_usage_uploader.py --dry-run --endpoint https://collector.example.com/api/codex/usage-v2 --source-name alice +python scripts/codex_usage_uploader.py --since 2026-05-01 --endpoint https://collector.example.com/api/codex/usage-v2 --source-name alice --token +python scripts/codex_usage_uploader.py --codex-home C:\Users\admin\.codex --timezone Asia/Shanghai --batch-size 500 --endpoint https://collector.example.com/api/codex/usage-v2 --source-name alice --token +``` + +If `python` is not on PATH, use the bundled Codex runtime if available. + +## Automation + +When the user asks to keep the dashboard updated, create a Codex cron automation that runs hourly and executes this skill's script with the configured environment variables. The default schedule is hourly in `Asia/Shanghai`. + +The automation prompt should be self-contained: + +```text +Use $codex-usage-uploader-v2 to run the Codex usage uploader with the configured CODEX_USAGE_INGEST_URL, CODEX_USAGE_SOURCE_NAME, and CODEX_USAGE_BEARER_TOKEN environment variables. Upload new local Codex usage metadata only with V2 token snapshot de-duplication. +``` + +## Privacy Boundary + +Do not upload raw message text, agent replies, reasoning content, shell stdout/stderr, full commands, tool output, or unified diffs. The script uploads counts, IDs, timestamps, status, durations, token fields, rate-limit fields, cwd/model/git context, command program summaries, and patch change-type statistics only. + +Accuracy tiers: + +- Exact: token fields, session ID, timestamp, event type, line number, rate limits, explicit duration/status fields. +- Context linked: cwd, model, effort, git, and turn context attached from the latest same-session context event. +- External: `source_name` is supplied by configuration. If it represents a person, configure a stable person or account label. + +Token counting contract: + +- Codex `token_count` log entries can include both `last_token_usage` and `total_token_usage`. +- `total_token_usage` is cumulative within the session and must not be summed across log entries. +- Codex can re-emit `token_count` when rate-limit state changes, sometimes with the same token usage snapshot. +- The uploader sends `token` from `last_token_usage` only for the first observed `total_token_usage` snapshot; repeated snapshots are skipped so dashboards can sum uploaded `token` fields as per-call increments. + +## API Reference + +Read `references/ingest_api.md` when implementing or validating the server-side ingest API. The server must treat `event_id` as an idempotency key and should return `accepted`, `duplicates`, and `errors` counts. + +## Troubleshooting + +- Missing endpoint: set `CODEX_USAGE_INGEST_URL` or pass `--endpoint`. +- Missing source: set `CODEX_USAGE_SOURCE_NAME` or pass `--source-name`. +- Missing token: set `CODEX_USAGE_BEARER_TOKEN` or pass `--token`; `--dry-run` does not require a token. +- Duplicate data: expected during file moves or re-scans; de-duplicate by `event_id`. Token usage snapshots are also filtered client-side when Codex re-emits the same cumulative usage with new rate-limit metadata. +- No events found: verify `$CODEX_HOME\sessions` or `$CODEX_HOME\archived_sessions` contains JSONL logs. diff --git a/skills/codex-usage-uploader-v2/agents/openai.yaml b/skills/codex-usage-uploader-v2/agents/openai.yaml new file mode 100644 index 0000000..1560527 --- /dev/null +++ b/skills/codex-usage-uploader-v2/agents/openai.yaml @@ -0,0 +1,4 @@ +interface: + display_name: "Codex Usage Uploader V2" + short_description: "Upload Codex usage with token de-duplication." + default_prompt: "Use $codex-usage-uploader-v2 to upload sanitized Codex usage metadata with token snapshot de-duplication." diff --git a/skills/codex-usage-uploader-v2/references/ingest_api.md b/skills/codex-usage-uploader-v2/references/ingest_api.md new file mode 100644 index 0000000..8c76588 --- /dev/null +++ b/skills/codex-usage-uploader-v2/references/ingest_api.md @@ -0,0 +1,73 @@ +# Codex Usage Ingest API V2 + +## Endpoint + +`POST ` + +Required headers: + +```http +Authorization: Bearer +Content-Type: application/json +``` + +## Request Body + +```json +{ + "source_name": "alice", + "machine_id": "stable-hash", + "codex_home": "C:\\Users\\alice\\.codex", + "collector_version": "2.0.0", + "collected_at": "2026-05-06T12:00:00+08:00", + "timezone": "Asia/Shanghai", + "batch_id": "uuid", + "events": [] +} +``` + +## Event Shape + +Every event contains: + +```json +{ + "event_id": "sha256-prefix", + "session_id": "uuid", + "timestamp": "2026-05-06T12:00:00.000Z", + "event_type": "token_count", + "line_no": 42, + "accuracy": "token_exact_context_linked", + "context": {} +} +``` + +Optional event sections: + +- `token`: de-duplicated per-call `last_token_usage` fields: `input_tokens`, `cached_input_tokens`, `output_tokens`, `reasoning_output_tokens`, `total_tokens`. +- `rate_limits`: sanitized Codex rate-limit metadata. +- `session`: sanitized session metadata. +- `task`: task lifecycle, duration, TTFT, abort, or error metadata. +- `tool`: tool name, namespace/server, call ID, duration, status, argument keys/counts, and output lengths only. +- `shell`: cwd, exit code, duration, status, command count, command types, and executable names only. +- `patch`: changed file count, change type counts, and file extension counts only. + +## Privacy Contract + +Clients must not send raw user messages, agent messages, reasoning content, shell stdout/stderr, full shell commands, tool output, API arguments with values, or unified diffs. + +## Response Body + +The server should return JSON: + +```json +{ + "accepted": 10, + "duplicates": 2, + "errors": [] +} +``` + +The server must de-duplicate by `event_id`. A repeated `event_id` from the same source should not increment usage totals twice. + +For token dashboards, sum only the uploaded `token` fields. Do not ingest or sum Codex log `total_token_usage`; it is cumulative within a session. Codex can also re-emit a `token_count` event when only rate-limit state changes, so clients filter repeated cumulative token snapshots before upload. diff --git a/skills/codex-usage-uploader-v2/scripts/codex_usage_uploader.py b/skills/codex-usage-uploader-v2/scripts/codex_usage_uploader.py new file mode 100644 index 0000000..73c2c01 --- /dev/null +++ b/skills/codex-usage-uploader-v2/scripts/codex_usage_uploader.py @@ -0,0 +1,762 @@ +#!/usr/bin/env python3 +import argparse +import copy +import hashlib +import json +import os +import pathlib +import re +import socket +import sys +import time +import uuid +import urllib.error +import urllib.request +from datetime import date, datetime, time as datetime_time, timedelta, timezone + +try: + from zoneinfo import ZoneInfo +except ModuleNotFoundError: + try: + from backports.zoneinfo import ZoneInfo + except ModuleNotFoundError: + def ZoneInfo(name): + fixed_offsets = { + "UTC": timezone.utc, + "Asia/Shanghai": timezone(timedelta(hours=8), name), + } + if name in fixed_offsets: + return fixed_offsets[name] + raise RuntimeError("Python 3.9+ or backports.zoneinfo is required for this timezone.") + + +VERSION = "2.0.0" +DEFAULT_BATCH_SIZE = 500 +SESSION_ID_RE = re.compile(r"([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})", re.I) +TOKEN_KEYS = ( + "input_tokens", + "cached_input_tokens", + "output_tokens", + "reasoning_output_tokens", + "total_tokens", +) +CONTENT_EVENT_TYPES = { + "user_message", + "message", + "agent_message", + "reasoning", + "compacted", + "context_compacted", +} + + +def parse_args(): + parser = argparse.ArgumentParser(description="Upload sanitized Codex usage metadata.") + parser.add_argument("--endpoint", default=os.environ.get("CODEX_USAGE_INGEST_URL")) + parser.add_argument("--source-name", default=os.environ.get("CODEX_USAGE_SOURCE_NAME")) + parser.add_argument("--token", default=os.environ.get("CODEX_USAGE_BEARER_TOKEN")) + parser.add_argument("--codex-home", default=os.environ.get("CODEX_HOME") or str(pathlib.Path.home() / ".codex")) + parser.add_argument("--timezone", default=os.environ.get("CODEX_USAGE_TIMEZONE") or "Asia/Shanghai") + parser.add_argument("--state-file", default=None) + parser.add_argument("--batch-size", type=int, default=DEFAULT_BATCH_SIZE) + parser.add_argument("--dry-run", action="store_true") + parser.add_argument("--since", default=None, help="Only emit events at or after this local date/time.") + parser.add_argument("--timeout", type=float, default=30.0) + parser.add_argument("--retries", type=int, default=3) + parser.add_argument("--retry-delay", type=float, default=1.0) + args = parser.parse_args() + + if not args.endpoint: + parser.error("--endpoint or CODEX_USAGE_INGEST_URL is required") + if not args.source_name: + parser.error("--source-name or CODEX_USAGE_SOURCE_NAME is required") + if not args.dry_run and not args.token: + parser.error("--token or CODEX_USAGE_BEARER_TOKEN is required unless --dry-run is used") + if args.batch_size <= 0: + parser.error("--batch-size must be greater than 0") + return args + + +def stable_hash(value, length=16): + return hashlib.sha256(str(value).encode("utf-8", errors="replace")).hexdigest()[:length] + + +def machine_id(codex_home): + raw = "|".join([ + socket.gethostname(), + os.environ.get("USERNAME") or os.environ.get("USER") or "", + str(pathlib.Path(codex_home).expanduser()), + ]) + return stable_hash(raw.lower(), 24) + + +def session_id_from_path(path): + match = SESSION_ID_RE.search(path.name) + return match.group(1).lower() if match else path.stem + + +def event_type(obj): + payload = obj.get("payload") if isinstance(obj, dict) else None + if isinstance(payload, dict) and payload.get("type"): + return payload.get("type") + return obj.get("type") if isinstance(obj, dict) else None + + +def iter_jsonl_files(codex_home): + roots = [codex_home / "sessions", codex_home / "archived_sessions"] + for root in roots: + if root.exists(): + for path in sorted(root.rglob("*.jsonl")): + yield path + + +def parse_datetime(value, tz): + if not value: + return None + if re.fullmatch(r"\d{4}-\d{2}-\d{2}", value): + day = date.fromisoformat(value) + return datetime.combine(day, datetime_time.min, tzinfo=tz) + normalized = value.replace("Z", "+00:00") + dt = datetime.fromisoformat(normalized) + if dt.tzinfo is None: + dt = dt.replace(tzinfo=tz) + return dt.astimezone(tz) + + +def clean(value): + if isinstance(value, dict): + cleaned = {key: clean(item) for key, item in value.items()} + return {key: item for key, item in cleaned.items() if item is not None and item != {} and item != []} + if isinstance(value, list): + return [clean(item) for item in value if item is not None] + return value + + +def duration_payload(value): + if not isinstance(value, dict): + return None + return clean({ + "secs": value.get("secs"), + "nanos": value.get("nanos"), + }) + + +def argument_keys(value): + parsed = None + if isinstance(value, dict): + parsed = value + elif isinstance(value, str): + try: + parsed = json.loads(value) + except json.JSONDecodeError: + return {"length": len(value)} + if isinstance(parsed, dict): + return {"keys": sorted(str(key) for key in parsed.keys()), "length": len(parsed)} + if isinstance(parsed, list): + return {"list_length": len(parsed)} + return None + + +def first_program(value): + if not value or not isinstance(value, str): + return None + token = value.strip().split()[0] if value.strip() else "" + if not token: + return None + token = token.strip("\"'") + if "\\" in token or "/" in token: + token = pathlib.PureWindowsPath(token).name + return token[:80] + + +def summarize_command(payload): + parsed = payload.get("parsed_cmd") + command_types = [] + programs = [] + if isinstance(parsed, list): + for item in parsed: + if not isinstance(item, dict): + continue + if item.get("type"): + command_types.append(str(item.get("type"))) + program = first_program(item.get("cmd")) + if program: + programs.append(program) + return clean({ + "cwd": payload.get("cwd"), + "exit_code": payload.get("exit_code"), + "duration": duration_payload(payload.get("duration")), + "status": payload.get("status"), + "command_count": len(parsed) if isinstance(parsed, list) else None, + "command_types": sorted(set(command_types)) if command_types else None, + "programs": sorted(set(programs))[:20] if programs else None, + }) + + +def summarize_patch(payload): + changes = payload.get("changes") + if not isinstance(changes, dict): + return None + change_types = {} + extensions = {} + for path_text, change in changes.items(): + if isinstance(change, dict): + change_type = change.get("type") or "unknown" + else: + change_type = "unknown" + change_types[change_type] = change_types.get(change_type, 0) + 1 + suffix = pathlib.PureWindowsPath(str(path_text)).suffix.lower() or "" + extensions[suffix] = extensions.get(suffix, 0) + 1 + return clean({ + "changed_files_count": len(changes), + "change_types": dict(sorted(change_types.items())), + "file_extensions": dict(sorted(extensions.items())), + "success": payload.get("success"), + "status": payload.get("status"), + }) + + +def summarize_result(result): + if not isinstance(result, dict): + return None + if "Ok" in result and isinstance(result["Ok"], dict): + ok = result["Ok"] + content = ok.get("content") + return clean({ + "ok": True, + "is_error": ok.get("isError"), + "content_items_count": len(content) if isinstance(content, list) else None, + "has_meta": isinstance(ok.get("_meta"), dict), + }) + if "Err" in result: + return {"ok": False} + return None + + +def summarize_tool(event_type_name, payload): + if event_type_name == "function_call": + return clean({ + "call_id": payload.get("call_id"), + "name": payload.get("name"), + "namespace": payload.get("namespace"), + "arguments": argument_keys(payload.get("arguments")), + }) + if event_type_name == "function_call_output": + output = payload.get("output") + return clean({ + "call_id": payload.get("call_id"), + "output_type": type(output).__name__, + "output_length": len(output) if isinstance(output, (str, list)) else None, + }) + if event_type_name == "custom_tool_call": + value = payload.get("input") + return clean({ + "call_id": payload.get("call_id"), + "name": payload.get("name"), + "status": payload.get("status"), + "input_length": len(value) if isinstance(value, str) else None, + }) + if event_type_name == "custom_tool_call_output": + output = payload.get("output") + return clean({ + "call_id": payload.get("call_id"), + "output_length": len(output) if isinstance(output, str) else None, + }) + if event_type_name == "mcp_tool_call_end": + invocation = payload.get("invocation") if isinstance(payload.get("invocation"), dict) else {} + return clean({ + "call_id": payload.get("call_id"), + "server": invocation.get("server"), + "name": invocation.get("tool"), + "arguments": argument_keys(invocation.get("arguments")), + "duration": duration_payload(payload.get("duration")), + "result": summarize_result(payload.get("result")), + }) + if event_type_name in {"dynamic_tool_call_request", "dynamic_tool_call_response"}: + return clean({ + "call_id": payload.get("call_id") or payload.get("callId"), + "turn_id": payload.get("turn_id") or payload.get("turnId"), + "name": payload.get("tool"), + "namespace": payload.get("namespace"), + "arguments": argument_keys(payload.get("arguments")), + "success": payload.get("success"), + "duration": duration_payload(payload.get("duration")), + "content_items_count": len(payload.get("content_items")) if isinstance(payload.get("content_items"), list) else None, + }) + if event_type_name in {"tool_search_call", "tool_search_output"}: + return clean({ + "call_id": payload.get("call_id"), + "status": payload.get("status"), + "execution": payload.get("execution"), + "arguments": argument_keys(payload.get("arguments")), + "tools_count": len(payload.get("tools")) if isinstance(payload.get("tools"), list) else None, + }) + if event_type_name in {"web_search_call", "web_search_end"}: + action = payload.get("action") if isinstance(payload.get("action"), dict) else {} + query = payload.get("query") or action.get("query") + queries = action.get("queries") + return clean({ + "call_id": payload.get("call_id"), + "status": payload.get("status"), + "action_type": action.get("type"), + "query_length": len(query) if isinstance(query, str) else None, + "queries_count": len(queries) if isinstance(queries, list) else None, + }) + if event_type_name == "view_image_tool_call": + path_value = payload.get("path") + return clean({ + "call_id": payload.get("call_id"), + "path_extension": pathlib.PureWindowsPath(path_value).suffix.lower() if isinstance(path_value, str) else None, + }) + return None + + +def update_context(context, event_type_name, payload): + next_context = copy.deepcopy(context or {}) + if event_type_name == "session_meta": + for key in ("cwd", "originator", "cli_version", "model_provider", "forked_from_id", "memory_mode"): + if payload.get(key) is not None: + next_context[key] = payload.get(key) + git = payload.get("git") + if isinstance(git, dict): + next_context["git"] = clean({ + "repository_url": git.get("repository_url"), + "commit_hash": git.get("commit_hash"), + "branch": git.get("branch"), + }) + elif event_type_name == "turn_context": + for key in ("turn_id", "cwd", "current_date", "timezone", "model", "effort"): + if payload.get(key) is not None: + next_context[key] = payload.get(key) + mode = payload.get("collaboration_mode") + settings = mode.get("settings") if isinstance(mode, dict) and isinstance(mode.get("settings"), dict) else {} + if settings.get("model"): + next_context["model"] = settings.get("model") + if settings.get("reasoning_effort"): + next_context["reasoning_effort"] = settings.get("reasoning_effort") + elif event_type_name == "task_started": + if payload.get("turn_id"): + next_context["turn_id"] = payload.get("turn_id") + if payload.get("model_context_window") is not None: + next_context["model_context_window"] = payload.get("model_context_window") + if payload.get("collaboration_mode_kind"): + next_context["collaboration_mode_kind"] = payload.get("collaboration_mode_kind") + elif event_type_name == "thread_name_updated": + thread_name = payload.get("thread_name") + next_context["thread"] = clean({ + "thread_id": payload.get("thread_id"), + "thread_name_hash": stable_hash(thread_name) if isinstance(thread_name, str) else None, + "thread_name_length": len(thread_name) if isinstance(thread_name, str) else None, + }) + + turn_id = payload.get("turn_id") or payload.get("turnId") + if turn_id: + next_context["turn_id"] = turn_id + return clean(next_context) + + +def public_context(context): + if not isinstance(context, dict): + return {} + allowed = { + "cwd", + "originator", + "cli_version", + "model_provider", + "forked_from_id", + "memory_mode", + "git", + "turn_id", + "current_date", + "timezone", + "model", + "effort", + "reasoning_effort", + "model_context_window", + "collaboration_mode_kind", + "thread", + } + return clean({key: copy.deepcopy(value) for key, value in context.items() if key in allowed}) + + +def event_accuracy(event_type_name): + if event_type_name == "token_count": + return "token_exact_context_linked" + if event_type_name in {"session_meta", "turn_context", "thread_name_updated"}: + return "context_exact" + if event_type_name in {"task_started", "task_complete", "turn_aborted", "error"}: + return "task_exact_context_linked" + if event_type_name in {"exec_command_end", "patch_apply_end"}: + return "event_exact_context_linked" + return "event_metadata_context_linked" + + +def build_event_identity(source_name, mid, session_id, line_no, timestamp, event_type_name, payload): + info = payload.get("info") if isinstance(payload.get("info"), dict) else {} + usage = info.get("last_token_usage") if isinstance(info.get("last_token_usage"), dict) else {} + identity = { + "source_name": source_name, + "machine_id": mid, + "session_id": session_id, + "line_no": line_no, + "timestamp": timestamp, + "event_type": event_type_name, + "call_id": payload.get("call_id") or payload.get("callId"), + "token_total": usage.get("total_tokens"), + "token_input": usage.get("input_tokens"), + "token_output": usage.get("output_tokens"), + } + encoded = json.dumps(identity, ensure_ascii=False, sort_keys=True, separators=(",", ":")) + return hashlib.sha256(encoded.encode("utf-8")).hexdigest() + + +def token_usage_fingerprint(usage): + if not isinstance(usage, dict): + return None + return tuple(usage.get(key, 0) or 0 for key in TOKEN_KEYS) + + +def token_total_fingerprint(payload): + if not isinstance(payload, dict) or payload.get("type") != "token_count": + return None + info = payload.get("info") if isinstance(payload.get("info"), dict) else {} + return token_usage_fingerprint(info.get("total_token_usage")) + + +def build_event(obj, path, line_no, context, source_name, mid): + payload = obj.get("payload") if isinstance(obj.get("payload"), dict) else {} + event_type_name = event_type(obj) + if not event_type_name or event_type_name in CONTENT_EVENT_TYPES: + return None + + timestamp = obj.get("timestamp") + sid = session_id_from_path(path) + event = { + "event_id": build_event_identity(source_name, mid, sid, line_no, timestamp, event_type_name, payload), + "session_id": sid, + "timestamp": timestamp, + "event_type": event_type_name, + "line_no": line_no, + "accuracy": event_accuracy(event_type_name), + "context": public_context(context), + } + + if event_type_name == "token_count": + info = payload.get("info") if isinstance(payload.get("info"), dict) else {} + usage = info.get("last_token_usage") if isinstance(info.get("last_token_usage"), dict) else {} + if not usage: + return None + event["token"] = {key: usage.get(key, 0) or 0 for key in TOKEN_KEYS} + if info.get("model_context_window") is not None: + event["context"]["model_context_window"] = info.get("model_context_window") + if isinstance(payload.get("rate_limits"), dict): + event["rate_limits"] = payload.get("rate_limits") + elif event_type_name == "session_meta": + event["session"] = clean({ + "id": payload.get("id"), + "originator": payload.get("originator"), + "cli_version": payload.get("cli_version"), + "model_provider": payload.get("model_provider"), + "forked_from_id": payload.get("forked_from_id"), + "memory_mode": payload.get("memory_mode"), + }) + elif event_type_name in {"turn_context", "thread_name_updated"}: + event["context_event"] = True + elif event_type_name == "task_started": + event["task"] = clean({ + "status": "started", + "turn_id": payload.get("turn_id"), + "started_at": payload.get("started_at"), + "model_context_window": payload.get("model_context_window"), + "collaboration_mode_kind": payload.get("collaboration_mode_kind"), + }) + elif event_type_name == "task_complete": + event["task"] = clean({ + "status": "complete", + "turn_id": payload.get("turn_id"), + "completed_at": payload.get("completed_at"), + "duration_ms": payload.get("duration_ms"), + "time_to_first_token_ms": payload.get("time_to_first_token_ms"), + "last_agent_message_length": len(payload.get("last_agent_message")) if isinstance(payload.get("last_agent_message"), str) else None, + }) + elif event_type_name == "turn_aborted": + event["task"] = clean({ + "status": "aborted", + "turn_id": payload.get("turn_id"), + "reason": payload.get("reason"), + "completed_at": payload.get("completed_at"), + "duration_ms": payload.get("duration_ms"), + }) + elif event_type_name == "error": + event["task"] = clean({ + "status": "error", + "message_length": len(payload.get("message")) if isinstance(payload.get("message"), str) else None, + "codex_error_info_length": len(payload.get("codex_error_info")) if isinstance(payload.get("codex_error_info"), str) else None, + }) + elif event_type_name == "exec_command_end": + event["shell"] = summarize_command(payload) + elif event_type_name == "patch_apply_end": + event["patch"] = summarize_patch(payload) + elif event_type_name == "item_completed": + item = payload.get("item") if isinstance(payload.get("item"), dict) else {} + event["task"] = clean({ + "status": "item_completed", + "thread_id": payload.get("thread_id"), + "turn_id": payload.get("turn_id"), + "item_type": item.get("type"), + "item_id": item.get("id"), + "item_text_length": len(item.get("text")) if isinstance(item.get("text"), str) else None, + }) + else: + tool = summarize_tool(event_type_name, payload) + if tool: + event["tool"] = tool + else: + return None + + return clean(event) + + +def load_state(path): + if not path.exists(): + return {"version": 1, "files": {}} + try: + with path.open("r", encoding="utf-8") as handle: + state = json.load(handle) + except (OSError, json.JSONDecodeError): + return {"version": 1, "files": {}} + if not isinstance(state, dict): + return {"version": 1, "files": {}} + state.setdefault("version", 1) + state.setdefault("files", {}) + return state + + +def save_state(path, state): + path.parent.mkdir(parents=True, exist_ok=True) + tmp = path.with_suffix(path.suffix + ".tmp") + with tmp.open("w", encoding="utf-8") as handle: + json.dump(state, handle, ensure_ascii=False, indent=2, sort_keys=True) + tmp.replace(path) + + +def reconstruct_last_token_total_fingerprint(path, offset): + if not offset: + return None + current_offset = 0 + last_fingerprint = None + try: + handle = path.open("rb") + except OSError: + return None + with handle: + for raw_line in handle: + next_offset = current_offset + len(raw_line) + if next_offset > offset: + break + current_offset = next_offset + try: + obj = json.loads(raw_line.decode("utf-8", errors="replace")) + except json.JSONDecodeError: + continue + payload = obj.get("payload") if isinstance(obj.get("payload"), dict) else {} + fingerprint = token_total_fingerprint(payload) + if fingerprint is not None: + last_fingerprint = fingerprint + return last_fingerprint + + +def read_file_events(path, record, tz, since_dt, source_name, mid): + record = record if isinstance(record, dict) else {} + stat = path.stat() + offset = int(record.get("offset") or 0) + line_no = int(record.get("line_no") or 0) + context = record.get("context") if isinstance(record.get("context"), dict) else {} + last_token_total_fingerprint = record.get("last_token_total_fingerprint") + if isinstance(last_token_total_fingerprint, list): + last_token_total_fingerprint = tuple(last_token_total_fingerprint) + else: + last_token_total_fingerprint = None + if stat.st_size < offset: + offset = 0 + line_no = 0 + context = {} + last_token_total_fingerprint = None + elif offset and last_token_total_fingerprint is None: + last_token_total_fingerprint = reconstruct_last_token_total_fingerprint(path, offset) + + events = [] + last_timestamp = record.get("last_timestamp") + current_offset = offset + with path.open("rb") as handle: + if offset: + handle.seek(offset) + for raw_line in handle: + current_offset += len(raw_line) + line_no += 1 + line = raw_line.decode("utf-8", errors="replace") + try: + obj = json.loads(line) + except json.JSONDecodeError: + continue + payload = obj.get("payload") if isinstance(obj.get("payload"), dict) else {} + type_name = event_type(obj) + if type_name: + context = update_context(context, type_name, payload) + timestamp = obj.get("timestamp") + if timestamp: + last_timestamp = timestamp + token_fingerprint = token_total_fingerprint(payload) + is_repeated_token_snapshot = False + if token_fingerprint is not None: + is_repeated_token_snapshot = token_fingerprint == last_token_total_fingerprint + last_token_total_fingerprint = token_fingerprint + if since_dt and timestamp: + try: + if parse_datetime(timestamp, tz) < since_dt: + continue + except ValueError: + continue + if is_repeated_token_snapshot: + continue + event = build_event(obj, path, line_no, context, source_name, mid) + if event: + events.append(event) + + next_record = { + "offset": current_offset, + "line_no": line_no, + "last_timestamp": last_timestamp, + "context": context, + "last_token_total_fingerprint": list(last_token_total_fingerprint) if last_token_total_fingerprint is not None else None, + } + return events, next_record + + +def collect_events(codex_home, state, tz, since_dt, source_name, mid): + codex_home = pathlib.Path(codex_home).expanduser() + next_state = copy.deepcopy(state or {"version": 1, "files": {}}) + next_state.setdefault("version", 1) + next_state.setdefault("files", {}) + all_events = [] + seen_ids = set() + files_scanned = 0 + for path in iter_jsonl_files(codex_home): + files_scanned += 1 + key = str(path) + file_events, next_record = read_file_events( + path, + next_state["files"].get(key, {}), + tz, + since_dt, + source_name, + mid, + ) + next_state["files"][key] = next_record + for event in file_events: + if event["event_id"] in seen_ids: + continue + seen_ids.add(event["event_id"]) + all_events.append(event) + stats = { + "files_scanned": files_scanned, + "events_collected": len(all_events), + } + return all_events, next_state, stats + + +def chunked(items, size): + for start in range(0, len(items), size): + yield items[start:start + size] + + +def build_batch(endpoint_args, codex_home, tz_name, mid, events): + return { + "source_name": endpoint_args.source_name, + "machine_id": mid, + "codex_home": str(codex_home), + "collector_version": VERSION, + "collected_at": datetime.now(ZoneInfo(tz_name)).isoformat(timespec="seconds"), + "timezone": tz_name, + "batch_id": str(uuid.uuid4()), + "events": events, + } + + +def post_json(endpoint, token, body, timeout): + data = json.dumps(body, ensure_ascii=False, separators=(",", ":")).encode("utf-8") + request = urllib.request.Request( + endpoint, + data=data, + headers={ + "Authorization": f"Bearer {token}", + "Content-Type": "application/json", + "User-Agent": f"codex-usage-uploader-v2/{VERSION}", + }, + method="POST", + ) + with urllib.request.urlopen(request, timeout=timeout) as response: + response_body = response.read().decode("utf-8", errors="replace") + if not response_body: + return {} + try: + return json.loads(response_body) + except json.JSONDecodeError: + return {"raw_response_length": len(response_body)} + + +def upload_events(args, codex_home, tz_name, mid, events): + responses = [] + for batch_events in chunked(events, args.batch_size): + body = build_batch(args, codex_home, tz_name, mid, batch_events) + last_error = None + for attempt in range(max(args.retries, 1)): + try: + responses.append(post_json(args.endpoint, args.token, body, args.timeout)) + last_error = None + break + except (urllib.error.URLError, TimeoutError, OSError) as exc: + last_error = exc + if attempt + 1 < max(args.retries, 1): + time.sleep(args.retry_delay) + if last_error: + raise RuntimeError(f"upload failed after {max(args.retries, 1)} attempt(s): {last_error}") + return responses + + +def main(): + args = parse_args() + tz = ZoneInfo(args.timezone) + codex_home = pathlib.Path(args.codex_home).expanduser() + state_file = pathlib.Path(args.state_file).expanduser() if args.state_file else codex_home / "codex_usage_uploader_state.json" + since_dt = parse_datetime(args.since, tz) if args.since else None + mid = machine_id(codex_home) + state = load_state(state_file) + events, next_state, stats = collect_events(codex_home, state, tz, since_dt, args.source_name, mid) + + if args.dry_run: + print(json.dumps({ + "dry_run": True, + "source_name": args.source_name, + "machine_id": mid, + "codex_home": str(codex_home), + "events_ready": len(events), + "stats": stats, + "sample_event": events[0] if events else None, + }, ensure_ascii=False, indent=2)) + return 0 + + responses = upload_events(args, codex_home, args.timezone, mid, events) if events else [] + save_state(state_file, next_state) + print(json.dumps({ + "uploaded": len(events), + "batches": len(responses), + "responses": responses, + "state_file": str(state_file), + "stats": stats, + }, ensure_ascii=False, indent=2)) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/skills/codex-usage-uploader-v2/scripts/test_codex_usage_uploader.py b/skills/codex-usage-uploader-v2/scripts/test_codex_usage_uploader.py new file mode 100644 index 0000000..b64f38a --- /dev/null +++ b/skills/codex-usage-uploader-v2/scripts/test_codex_usage_uploader.py @@ -0,0 +1,373 @@ +#!/usr/bin/env python3 +import importlib.util +import copy +import json +import subprocess +import sys +import tempfile +import threading +from http.server import BaseHTTPRequestHandler, HTTPServer +from pathlib import Path + + +SCRIPT = Path(__file__).with_name("codex_usage_uploader.py") +SESSION_ID = "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee" + + +def load_module(): + spec = importlib.util.spec_from_file_location("codex_usage_uploader", SCRIPT) + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +def event(timestamp, outer_type, payload): + return {"timestamp": timestamp, "type": outer_type, "payload": payload} + + +def write_session(codex_home, events): + session_dir = codex_home / "sessions" / "2026" / "05" / "06" + session_dir.mkdir(parents=True) + path = session_dir / f"rollout-2026-05-06T10-00-00-{SESSION_ID}.jsonl" + with path.open("w", encoding="utf-8") as handle: + for item in events: + handle.write(json.dumps(item, ensure_ascii=False) + "\n") + return path + + +def sample_events(): + return [ + event("2026-05-06T01:00:00.000Z", "session_meta", { + "id": SESSION_ID, + "cwd": "D:\\repo", + "originator": "codex_desktop", + "cli_version": "1.2.3", + "model_provider": "openai", + "git": { + "repository_url": "https://git.example.com/team/repo.git", + "commit_hash": "abc123", + "branch": "main", + }, + }), + event("2026-05-06T01:00:01.000Z", "turn_context", { + "turn_id": "turn-1", + "cwd": "D:\\repo", + "current_date": "2026-05-06", + "timezone": "Asia/Shanghai", + "model": "gpt-5.4", + "effort": "medium", + "collaboration_mode": { + "settings": { + "model": "gpt-5.4", + "reasoning_effort": "medium", + "developer_instructions": "DO_NOT_UPLOAD", + } + }, + }), + event("2026-05-06T01:00:02.000Z", "event_msg", { + "type": "token_count", + "info": { + "last_token_usage": { + "input_tokens": 100, + "cached_input_tokens": 30, + "output_tokens": 20, + "reasoning_output_tokens": 5, + "total_tokens": 120, + }, + "total_token_usage": { + "input_tokens": 9999, + "cached_input_tokens": 9999, + "output_tokens": 9999, + "reasoning_output_tokens": 9999, + "total_tokens": 9999, + }, + "model_context_window": 258400, + }, + "rate_limits": { + "limit_id": "primary", + "primary": {"used_percent": 12.5, "window_minutes": 300, "resets_at": 123}, + "secondary": None, + "credits": {"has_credits": True, "unlimited": False, "balance": None}, + "plan_type": "team", + }, + }), + event("2026-05-06T01:00:03.000Z", "event_msg", { + "type": "task_complete", + "turn_id": "turn-1", + "completed_at": 1000, + "duration_ms": 2500, + "time_to_first_token_ms": 300, + "last_agent_message": "AGENT_SECRET", + }), + event("2026-05-06T01:00:04.000Z", "event_msg", { + "type": "exec_command_end", + "call_id": "call-shell", + "turn_id": "turn-1", + "command": ["powershell", "-Command", "Write-Host STDOUT_SECRET"], + "parsed_cmd": [{"type": "shell", "cmd": "python secret_script.py --password STDOUT_SECRET"}], + "cwd": "D:\\repo", + "stdout": "STDOUT_SECRET", + "stderr": "STDERR_SECRET", + "aggregated_output": "STDOUT_SECRET", + "exit_code": 0, + "duration": {"secs": 1, "nanos": 2}, + "formatted_output": "STDOUT_SECRET", + "status": "completed", + }), + event("2026-05-06T01:00:05.000Z", "event_msg", { + "type": "patch_apply_end", + "call_id": "call-patch", + "turn_id": "turn-1", + "stdout": "PATCH_STDOUT_SECRET", + "stderr": "PATCH_STDERR_SECRET", + "success": True, + "status": "completed", + "changes": { + "D:\\repo\\main.py": { + "type": "update", + "unified_diff": "@@ DIFF_SECRET @@", + "move_path": None, + } + }, + }), + event("2026-05-06T01:00:06.000Z", "event_msg", { + "type": "function_call", + "name": "shell_command", + "namespace": "functions", + "call_id": "call-tool", + "arguments": "{\"command\":\"echo ARG_SECRET\",\"api_key\":\"ARG_SECRET\"}", + }), + event("2026-05-06T01:00:07.000Z", "event_msg", { + "type": "function_call_output", + "call_id": "call-tool", + "output": "TOOL_OUTPUT_SECRET", + }), + event("2026-05-06T01:00:08.000Z", "user_message", { + "type": "user_message", + "message": "USER_MESSAGE_SECRET", + "images": [], + "local_images": [], + "text_elements": [], + }), + ] + + +def collect_once(codex_home, state=None): + module = load_module() + state = state or {"version": 1, "files": {}} + return module.collect_events( + codex_home, + state, + module.ZoneInfo("Asia/Shanghai"), + None, + "alice", + "machine-1", + ) + + +def test_collects_token_and_metadata_without_private_content(): + with tempfile.TemporaryDirectory() as temp: + codex_home = Path(temp) + write_session(codex_home, sample_events()) + events, next_state, stats = collect_once(codex_home) + + assert stats["files_scanned"] == 1 + assert len(events) == 8 + token = next(item for item in events if item["event_type"] == "token_count") + assert token["token"]["total_tokens"] == 120 + assert token["token"]["input_tokens"] == 100 + assert "total_token_usage" not in json.dumps(token) + assert token["context"]["model"] == "gpt-5.4" + assert token["context"]["cwd"] == "D:\\repo" + assert token["rate_limits"]["primary"]["used_percent"] == 12.5 + + shell = next(item for item in events if item["event_type"] == "exec_command_end") + assert shell["shell"]["programs"] == ["python"] + assert shell["shell"]["exit_code"] == 0 + + patch = next(item for item in events if item["event_type"] == "patch_apply_end") + assert patch["patch"]["changed_files_count"] == 1 + assert patch["patch"]["change_types"] == {"update": 1} + + serialized = json.dumps(events, ensure_ascii=False) + for secret in [ + "AGENT_SECRET", + "STDOUT_SECRET", + "STDERR_SECRET", + "PATCH_STDOUT_SECRET", + "PATCH_STDERR_SECRET", + "DIFF_SECRET", + "ARG_SECRET", + "TOOL_OUTPUT_SECRET", + "USER_MESSAGE_SECRET", + "DO_NOT_UPLOAD", + ]: + assert secret not in serialized + + +def test_state_incremental_and_truncated_rescan_event_id_stability(): + with tempfile.TemporaryDirectory() as temp: + codex_home = Path(temp) + path = write_session(codex_home, sample_events()) + events, state, _ = collect_once(codex_home) + second_events, _, _ = collect_once(codex_home, state) + assert second_events == [] + + original_token_id = next(item["event_id"] for item in events if item["event_type"] == "token_count") + truncated = sample_events()[:3] + with path.open("w", encoding="utf-8") as handle: + for item in truncated: + handle.write(json.dumps(item, ensure_ascii=False) + "\n") + replayed_events, _, _ = collect_once(codex_home, state) + replayed_token_id = next(item["event_id"] for item in replayed_events if item["event_type"] == "token_count") + assert replayed_token_id == original_token_id + + +def test_repeated_token_snapshot_is_not_uploaded_twice(): + with tempfile.TemporaryDirectory() as temp: + codex_home = Path(temp) + events = sample_events()[:3] + duplicate = copy.deepcopy(events[2]) + duplicate["timestamp"] = "2026-05-06T01:00:02.500Z" + duplicate["payload"]["rate_limits"]["primary"]["used_percent"] = 13.0 + advanced = copy.deepcopy(events[2]) + advanced["timestamp"] = "2026-05-06T01:00:03.000Z" + advanced["payload"]["info"]["last_token_usage"] = { + "input_tokens": 40, + "cached_input_tokens": 10, + "output_tokens": 8, + "reasoning_output_tokens": 2, + "total_tokens": 48, + } + advanced["payload"]["info"]["total_token_usage"] = { + "input_tokens": 140, + "cached_input_tokens": 40, + "output_tokens": 28, + "reasoning_output_tokens": 7, + "total_tokens": 168, + } + write_session(codex_home, events + [duplicate, advanced]) + + collected, _, _ = collect_once(codex_home) + + token_events = [item for item in collected if item["event_type"] == "token_count"] + assert [item["token"]["total_tokens"] for item in token_events] == [120, 48] + + +def test_legacy_state_reconstructs_token_snapshot_before_offset(): + with tempfile.TemporaryDirectory() as temp: + codex_home = Path(temp) + path = write_session(codex_home, sample_events()[:3]) + _, state, _ = collect_once(codex_home) + legacy_state = copy.deepcopy(state) + for record in legacy_state["files"].values(): + record.pop("last_token_total_fingerprint", None) + + duplicate = copy.deepcopy(sample_events()[2]) + duplicate["timestamp"] = "2026-05-06T01:00:03.000Z" + with path.open("a", encoding="utf-8") as handle: + handle.write(json.dumps(duplicate, ensure_ascii=False) + "\n") + + collected, _, _ = collect_once(codex_home, legacy_state) + + assert [item for item in collected if item["event_type"] == "token_count"] == [] + + +def test_dry_run_cli_does_not_require_token(): + with tempfile.TemporaryDirectory() as temp: + codex_home = Path(temp) + write_session(codex_home, sample_events()[:3]) + result = subprocess.run( + [ + sys.executable, + str(SCRIPT), + "--codex-home", + str(codex_home), + "--endpoint", + "http://127.0.0.1:1/ingest", + "--source-name", + "alice", + "--dry-run", + ], + text=True, + capture_output=True, + check=True, + ) + data = json.loads(result.stdout) + assert data["dry_run"] is True + assert data["events_ready"] == 3 + + +def test_http_upload_uses_bearer_header_and_retries(): + captured = [] + + class Handler(BaseHTTPRequestHandler): + def do_POST(self): + length = int(self.headers["Content-Length"]) + body = self.rfile.read(length) + captured.append((self.headers.get("Authorization"), json.loads(body.decode("utf-8")))) + if len(captured) == 1: + self.send_response(500) + self.end_headers() + return + self.send_response(200) + self.send_header("Content-Type", "application/json") + self.end_headers() + self.wfile.write(b'{"accepted": 3, "duplicates": 0, "errors": []}') + + def log_message(self, format, *args): + return + + server = HTTPServer(("127.0.0.1", 0), Handler) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + with tempfile.TemporaryDirectory() as temp: + codex_home = Path(temp) / "codex" + state_file = Path(temp) / "state.json" + write_session(codex_home, sample_events()[:3]) + result = subprocess.run( + [ + sys.executable, + str(SCRIPT), + "--codex-home", + str(codex_home), + "--state-file", + str(state_file), + "--endpoint", + f"http://127.0.0.1:{server.server_port}/ingest", + "--source-name", + "alice", + "--token", + "TOKEN123", + "--retries", + "2", + "--retry-delay", + "0", + ], + text=True, + capture_output=True, + check=True, + ) + output = json.loads(result.stdout) + assert output["uploaded"] == 3 + assert state_file.exists() + finally: + server.shutdown() + thread.join(timeout=5) + + assert len(captured) == 2 + assert captured[0][0] == "Bearer TOKEN123" + assert captured[1][0] == "Bearer TOKEN123" + assert captured[1][1]["source_name"] == "alice" + assert len(captured[1][1]["events"]) == 3 + + +if __name__ == "__main__": + test_collects_token_and_metadata_without_private_content() + test_state_incremental_and_truncated_rescan_event_id_stability() + test_repeated_token_snapshot_is_not_uploaded_twice() + test_legacy_state_reconstructs_token_snapshot_before_offset() + test_dry_run_cli_does_not_require_token() + test_http_upload_uses_bearer_header_and_retries() + print("tests passed")