From bf41cec80255566f381e5e8ac8fa6a1088dee6a7 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 01:13:57 +0800 Subject: [PATCH 1/4] feat(usage): observe common quota cycles and bound Codex turn timing Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../usage-statistics-settings.tsx | 2 +- .../migrations/0003-goal-duration-sources.sql | 8 + apps/usage-collector/schema.sql | 6 + apps/usage-collector/src/basic-usage.ts | 24 +-- apps/usage-collector/src/collector.js | 1 + apps/usage-collector/test/collector.test.mjs | 23 ++- loopx/chat_runtime.py | 2 +- loopx/cli_commands/quota.py | 18 +++ .../control_plane/runtime/usage_statistics.ts | 17 ++- .../runtime/usage_statistics_cli.ts | 11 +- .../runtime/usage_statistics_codex.ts | 67 +++++++++ .../runtime/usage_statistics_cycles.ts | 68 +++++++++ .../runtime/usage_statistics_goal_contract.ts | 29 ++-- .../runtime/usage_statistics_goals.ts | 68 +++++---- loopx/control_plane/turn_driver/executor.py | 2 +- loopx/usage_goal.py | 94 +++++++++++- loopx/usage_ping.py | 6 +- .../test_quota_settlement_cli.py | 48 +++++- .../usage_statistics_goals.test.ts | 14 +- .../usage_statistics_sources.test.ts | 141 ++++++++++++++++++ tests/test_usage_goal.py | 101 +++++++++++++ tsconfig.control-plane.json | 1 + 22 files changed, 674 insertions(+), 77 deletions(-) create mode 100644 apps/usage-collector/migrations/0003-goal-duration-sources.sql create mode 100644 loopx/control_plane/runtime/usage_statistics_codex.ts create mode 100644 loopx/control_plane/runtime/usage_statistics_cycles.ts create mode 100644 tests/control_plane_ts/usage_statistics_sources.test.ts diff --git a/apps/presentation/dashboard/src/features/personal-workspace/usage-statistics-settings.tsx b/apps/presentation/dashboard/src/features/personal-workspace/usage-statistics-settings.tsx index dde07c2004..9f5a6fbad7 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/usage-statistics-settings.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/usage-statistics-settings.tsx @@ -26,7 +26,7 @@ export function UsageStatisticsSettings() { ? "用于决定平台支持和改进命令体验。每天向 LoopX 的 Cloudflare 收集服务发送随机安装标识、版本、系统、CPU 架构、Python 版本和安装渠道;固定的 CLI 功能、结果、耗时区间和错误类别在本机按天汇总后另行发送,不带安装标识。" : "Helps prioritize platform support and CLI improvements. A daily heartbeat sends a random installation ID, version, OS, CPU architecture, Python version and install channel to the LoopX Cloudflare collector. Fixed CLI feature, result, duration and error counts are aggregated locally by day and sent separately without the ID."}

{zh ? "不采集提示词、代码、路径、命令参数、Goal 内容或原始错误。当前功能计数仅覆盖 CLI;命令成功不等于 Goal 完成。" : "No prompts, code, paths, arguments, Goal contents or raw errors. Feature counts currently cover CLI only; command success is not Goal completion."}

-

{zh ? "另按天汇总受管 Turn 与普通 Goal 对话的执行跨度和 Host 调用累计时长区间(并行重叠只计一次),包括未结束的 Goal,不带 Goal 或安装标识。从本机开始观测后计量,可能漏计;不代表 Goal 完成、CPU 用时或计费时长。原生 /goal 和外部附加任务暂不覆盖。" : "Also aggregates observed Goal execution-span and Host-call duration buckets by day, including unfinished Goals; overlaps count once. No Goal or installation ID. Measurement begins locally and may undercount; it is not completion evidence, CPU time or billing. Covers managed Turns and regular owner Goal chat; native /goal and externally attached tasks are excluded."}

+

{zh ? "按天分别汇总所有 Host 的 quota→spend 推进周期、已绑定 Codex 任务的本地轮次时间、受管 Turn 与普通 Goal 对话的 Host 调用时间。上传固定 Host 类别及跨度/时长区间,不上传会话内容、Goal 或安装标识。三种口径重叠,不能相加;可能漏计,不代表完成、CPU 用时或计费。" : "Daily, separate span/duration buckets for all Hosts using quota→spend, local timing events from bound Codex tasks, and direct Host calls in managed Turns and regular owner Goal chat. Sends fixed Host categories, never session contents, Goal or installation IDs. The three overlapping populations cannot be added; partial observations are not completion, CPU time or billing."}

{state ? <> diff --git a/apps/usage-collector/migrations/0003-goal-duration-sources.sql b/apps/usage-collector/migrations/0003-goal-duration-sources.sql new file mode 100644 index 0000000000..37d88bbee7 --- /dev/null +++ b/apps/usage-collector/migrations/0003-goal-duration-sources.sql @@ -0,0 +1,8 @@ +CREATE TABLE IF NOT EXISTS goal_duration_counts ( + day TEXT NOT NULL, measurement TEXT NOT NULL, host TEXT NOT NULL, + span TEXT NOT NULL, duration TEXT NOT NULL, count INTEGER NOT NULL, + PRIMARY KEY (day, measurement, host, span, duration) +); +INSERT INTO goal_duration_counts (day, measurement, host, span, duration, count) +SELECT day, 'host_call', 'unknown', span, execution, count FROM goal_usage_counts +WHERE true ON CONFLICT (day, measurement, host, span, duration) DO NOTHING; diff --git a/apps/usage-collector/schema.sql b/apps/usage-collector/schema.sql index 5999f44bcb..e5bb349ee1 100644 --- a/apps/usage-collector/schema.sql +++ b/apps/usage-collector/schema.sql @@ -33,3 +33,9 @@ CREATE TABLE IF NOT EXISTS goal_usage_counts ( day TEXT NOT NULL, span TEXT NOT NULL, execution TEXT NOT NULL, count INTEGER NOT NULL, PRIMARY KEY (day, span, execution) ); + +CREATE TABLE IF NOT EXISTS goal_duration_counts ( + day TEXT NOT NULL, measurement TEXT NOT NULL, host TEXT NOT NULL, + span TEXT NOT NULL, duration TEXT NOT NULL, count INTEGER NOT NULL, + PRIMARY KEY (day, measurement, host, span, duration) +); diff --git a/apps/usage-collector/src/basic-usage.ts b/apps/usage-collector/src/basic-usage.ts index 9fc079c7f1..6e65069b59 100644 --- a/apps/usage-collector/src/basic-usage.ts +++ b/apps/usage-collector/src/basic-usage.ts @@ -25,19 +25,25 @@ export async function aggregateStats(db: Database, since: string) { } export { validGoalAggregate } from "../../../loopx/control_plane/runtime/usage_statistics_goal_contract.ts"; +import { MEASUREMENTS } from "../../../loopx/control_plane/runtime/usage_statistics_goal_contract.ts"; import type { GoalAggregate } from "../../../loopx/control_plane/runtime/usage_statistics_goal_contract.ts"; export async function recordGoals(db: Database, value: GoalAggregate, day: string) { await db.batch(value.counters.map(row => db.prepare( - "INSERT INTO goal_usage_counts (day, span, execution, count) VALUES (?1, ?2, ?3, ?4) " + - "ON CONFLICT (day, span, execution) DO UPDATE SET count = count + excluded.count", - ).bind(day, row.span, row.execution, row.count))); + "INSERT INTO goal_duration_counts (day, measurement, host, span, duration, count) VALUES (?1, ?2, ?3, ?4, ?5, ?6) " + + "ON CONFLICT (day, measurement, host, span, duration) DO UPDATE SET count = count + excluded.count", + ).bind(day, row.measurement, row.host, row.span, row.duration, row.count))); } export async function goalStats(db: Database, since: string) { - const totals: Record> = {}; - for (const column of ["span", "execution"]) { - const rows = await db.prepare(`SELECT ${column} AS key, SUM(count) AS n FROM goal_usage_counts WHERE day >= ?1 GROUP BY ${column}`).bind(since).all(); - totals[column] = {}; - for (const row of rows.results) if (Number(row.n) >= 5) totals[column][String(row.key)] = Number(row.n); + // Each measurement is an independent population. Never add the three clocks. + const measurements: Record>> = {}; + for (const measurement of MEASUREMENTS) { + const totals: Record> = {}; + for (const column of ["span", "duration", "host"]) { + const rows = await db.prepare(`SELECT ${column} AS key, SUM(count) AS n FROM goal_duration_counts WHERE day >= ?1 AND measurement = ?2 GROUP BY ${column}`).bind(since, measurement).all(); + totals[column] = {}; + for (const row of rows.results) if (Number(row.n) >= 5) totals[column][String(row.key)] = Number(row.n); + } + measurements[measurement] = totals; } - return { schema: "loopx_goal_usage_stats_v1", definition: "Lossy observed Goal-day snapshots received in the last 30 UTC days, including unfinished Goals. Span is first-to-last observed execution; execution is union of Host-call intervals, not CPU time. Lower bounds since measurement began; managed Turns and regular owner Goal chat only. Not unique Goals, completion evidence or billing. Cells below 5 omitted.", totals }; + return { schema: "loopx_goal_usage_stats_v1", definition: "Lossy observed Goal-day samples in the last 30 receipt days, grouped by measurement. quota_cycle is admitted quota-to-successful-spend elapsed time for any Host, including pauses; codex_turn uses bound session timing; host_call uses direct invocation checkpoints. These populations overlap: do not add counts or durations. Span is first-to-last observed interval. Not unique Goals, completion, CPU time or billing. Cells below 5 omitted.", measurements }; } diff --git a/apps/usage-collector/src/collector.js b/apps/usage-collector/src/collector.js index a880d51e03..80fa4778e5 100644 --- a/apps/usage-collector/src/collector.js +++ b/apps/usage-collector/src/collector.js @@ -135,6 +135,7 @@ export async function purge(db, day) { const cutoff = shiftDays(day, -RETENTION_DAYS); await db.batch([ db.prepare("DELETE FROM pings WHERE day < ?1").bind(cutoff), + db.prepare("DELETE FROM goal_duration_counts WHERE day < ?1").bind(shiftDays(day, -30)), db.prepare("DELETE FROM goal_usage_counts WHERE day < ?1").bind(shiftDays(day, -30)), db.prepare("DELETE FROM usage_counts WHERE day < ?1").bind(shiftDays(day, -30)), db.prepare("DELETE FROM installs WHERE install_id NOT IN (SELECT DISTINCT install_id FROM pings)"), diff --git a/apps/usage-collector/test/collector.test.mjs b/apps/usage-collector/test/collector.test.mjs index b07debf541..615173d586 100644 --- a/apps/usage-collector/test/collector.test.mjs +++ b/apps/usage-collector/test/collector.test.mjs @@ -168,18 +168,18 @@ test("Goal duration ingestion uses real SQL, rejects identity, suppresses small const send = value => new Request("https://collector.example/v1/goals", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify(value), }); - const value = { schema: "loopx_goal_usage_aggregate_v1", counters: [{ span: "lt_30d", execution: "lt_1d", count: 3 }] }; + const value = { schema: "loopx_goal_usage_aggregate_v1", counters: [{ measurement: "host_call", host: "unknown", span: "lt_30d", duration: "lt_1d", count: 3 }] }; for (const field of ["goal_id", "install_id", "timestamp", "prompt"]) { assert.equal((await handle(send({ ...value, [field]: "private" }), db)).status, 400); } assert.equal((await handle(send(value), db, at("2026-09-01"))).status, 204); const stats = () => handle(new Request("https://collector.example/v1/goal-stats"), db, at("2026-09-02")).then(r => r.json()); - assert.deepEqual((await stats()).totals, { span: {}, execution: {} }); + assert.deepEqual((await stats()).measurements.host_call, { span: {}, duration: {}, host: {} }); await handle(send(value), db, at("2026-09-01")); - assert.deepEqual((await stats()).totals, { span: { lt_30d: 6 }, execution: { lt_1d: 6 } }); - assert.deepEqual(Object.keys(db.raw.get("SELECT * FROM goal_usage_counts")).sort(), ["count", "day", "execution", "span"]); + assert.deepEqual((await stats()).measurements.host_call, { span: { lt_30d: 6 }, duration: { lt_1d: 6 }, host: { unknown: 6 } }); + assert.deepEqual(Object.keys(db.raw.get("SELECT * FROM goal_duration_counts")).sort(), ["count", "day", "duration", "host", "measurement", "span"]); await purge(db, "2026-10-02"); - assert.equal(db.raw.get("SELECT COUNT(*) n FROM goal_usage_counts").n, 0); + assert.equal(db.raw.get("SELECT COUNT(*) n FROM goal_duration_counts").n, 0); }); test("Goal migration is additive and preserves existing aggregate counters", () => { @@ -191,3 +191,16 @@ test("Goal migration is additive and preserves existing aggregate counters", () assert.equal(db.prepare("SELECT count(*) n FROM goal_usage_counts").get().n, 0); db.close(); }); + +test("measurement migration preserves historical counts and is safe to repeat", () => { + const db = new DatabaseSync(":memory:"); + db.exec(readFileSync(new URL("../migrations/0002-goal-usage.sql", import.meta.url), "utf8")); + db.exec("INSERT INTO goal_usage_counts VALUES ('2026-09-01','lt_7d','lt_6h',8)"); + const migration = readFileSync(new URL("../migrations/0003-goal-duration-sources.sql", import.meta.url), "utf8"); + db.exec(migration); db.exec(migration); + assert.deepEqual({...db.prepare("SELECT * FROM goal_duration_counts").get()}, { + day:"2026-09-01", measurement:"host_call", host:"unknown", span:"lt_7d", duration:"lt_6h", count:8, + }); + assert.equal(db.prepare("SELECT count FROM goal_usage_counts").get().count,8); + db.close(); +}); diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py index 0188bdf4ad..59befe2fb7 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -1279,7 +1279,7 @@ def event_sink(kind: str, payload: dict[str, Any]) -> None: # Manager/native/external conversations have different execution # and waiting boundaries. Do not invent Goal timing for them. observed_goal = str(session.get("goal_id") or "") if scope["kind"] == "owner_goal" else "" - with observe_goal_execution(self.store.root.parent, observed_goal): + with observe_goal_execution(self.store.root.parent, observed_goal, host=str(session.get("agent_id") or "unknown")): if attachments: if not isinstance(adapter, CodexAppServerAdapter): raise ValueError("image attachments currently require the Codex Agent endpoint") diff --git a/loopx/cli_commands/quota.py b/loopx/cli_commands/quota.py index 578d797c26..746bbd8ead 100644 --- a/loopx/cli_commands/quota.py +++ b/loopx/cli_commands/quota.py @@ -330,6 +330,8 @@ def handle_quota_command( print_payload: PrintPayload, append_cli_rollout_event: RolloutEventAppender, ) -> int: + import time + usage_quota_started = time.time_ns() // 1_000_000 heartbeat_turn_id: str | None = None heartbeat_receipt_existing: dict[str, object] | None = None heartbeat_receipt_existing_status = "replayed" @@ -755,6 +757,22 @@ def handle_quota_command( goal_id=args.goal_id, agent_id=args.agent_id, ) + if context is not None and payload.get("ok"): + phase = ( + "start" if args.quota_command == "should-run" and payload.get("should_run") is True + else "spend" if args.quota_command == "spend-slot" and bool(args.execute) + and payload.get("appended") else None + ) + if phase: + from ..usage_goal import observe_quota_cycle + observe_quota_cycle( + registry_path=registry_path, runtime_root=context.runtime_root, + goal_id=args.goal_id, agent_id=args.agent_id, + turn_id=_effective_spend_turn_instance_id(payload, heartbeat_turn_id=heartbeat_turn_id), + phase=phase, at=usage_quota_started if phase == "start" else time.time_ns() // 1_000_000, + host=str(getattr(args, "host_surface", None) or getattr(args, "runtime_profile", None) + or ("codex_app" if getattr(args, "codex_app", False) else "unknown")), + ) if bool(getattr(args, "turn_envelope", False)): payload = _render_turn_envelope_payload( payload, diff --git a/loopx/control_plane/runtime/usage_statistics.ts b/loopx/control_plane/runtime/usage_statistics.ts index 79cc679fef..1c66d9d005 100644 --- a/loopx/control_plane/runtime/usage_statistics.ts +++ b/loopx/control_plane/runtime/usage_statistics.ts @@ -11,9 +11,12 @@ import { recordGoalUsage, goalPreview } from "./usage_statistics_goals.ts"; import { validGoalAggregate } from "./usage_statistics_goal_contract.ts"; import type { GoalAggregate, GoalObservation } from "./usage_statistics_goal_contract.ts"; +import { cycleObservations } from "./usage_statistics_cycles.ts"; +import type { CycleObservation } from "./usage_statistics_cycles.ts"; + export const STATE_SCHEMA = "loopx_usage_ping_state_v1"; export const DEFAULT_ENDPOINT = "https://loopx-usage-collector.huangrt01.workers.dev/v1/ping"; -export const NOTICE_VERSION = 2; +export const NOTICE_VERSION = 3; export type Env = Record; export type Context = { env: Env; version: string; python: string; channel: string; now?: Date }; type Notice = { version: number; endpoint: string; policy: string }; @@ -91,7 +94,7 @@ export async function inspect(path: string, ctx: Context) { aggregate_preview: state.consent === "disabled" || !state.counters?.length ? null : { schema: AGGREGATE_SCHEMA, counters: state.counters }, goal_preview: state.consent === "disabled" ? null : await goalPreview(path + ".goals", state.generation).catch(() => null), aggregate_day: state.day ?? null, - disclosure: "LoopX basic usage statistics are on by default after this notice. Daily heartbeats send a random installation ID, version, OS, CPU architecture, Python version and install channel to the configured LoopX collector (Cloudflare). Fixed CLI feature/result/duration/error counts are sent separately without an ID. Observed Goal execution-span and Host-call duration buckets are separately aggregated without Goal or installation IDs. Measurements cover managed Turns and regular owner Goal chat, from observation onward; they are lower bounds, not completion or billing evidence. No prompts, code, paths, arguments, Goal contents or raw errors. Disable all with loopx usage-ping disable or LOOPX_USAGE_PING=0; inspect with loopx usage-ping status. Consent-required distributions wait for explicit enable. Recipient: " + (endpoint(ctx.env) || "not configured") }; + disclosure: "LoopX basic usage statistics are on by default after this notice. Daily heartbeats send a random installation ID, version, OS, CPU architecture, Python version and install channel to the configured LoopX collector (Cloudflare). Fixed CLI feature/result/duration/error counts are sent separately without an ID. Goal span/duration buckets and fixed Host labels are aggregated without Goal or installation IDs. Common quota-to-spend cycles cover every Host using the quota CLI; bound Codex tasks add local timing-event reads; managed Turns and regular owner Goal chat add direct Host-call timing. These overlapping measurements are separate, partial and not completion or billing evidence. Raw session content is never uploaded. No prompts, code, paths, arguments, Goal contents or raw errors. Disable all with loopx usage-ping disable or LOOPX_USAGE_PING=0; inspect with loopx usage-ping status. Consent-required distributions wait for explicit enable. Recipient: " + (endpoint(ctx.env) || "not configured") }; } export async function configure(path: string, ctx: Context, action: "enable" | "disable" | "acknowledge", expectedNotice?: unknown) { await withFileMutationLock(path, async () => { @@ -100,6 +103,7 @@ export async function configure(path: string, ctx: Context, action: "enable" | " if (action === "disable") { await save(path, state); await rm(path + ".goals", { force: true }); + await rm(path + ".cycles", { force: true }); return; } if (action === "acknowledge" && JSON.stringify(expectedNotice) !== JSON.stringify(notice(ctx))) throw new Error("usage_notice_changed"); @@ -108,6 +112,7 @@ export async function configure(path: string, ctx: Context, action: "enable" | " if (state.notice && !sameNotice(state, ctx)) { state.counters = []; await rm(path + ".goals", { force: true }); + await rm(path + ".cycles", { force: true }); state.generation = randomUUID(); if (state.notice.endpoint !== endpoint(ctx.env)) state.install_id = randomUUID(); } @@ -125,7 +130,7 @@ const post: Post = async (url, payload) => (await fetch(url, { })).status; /** Called in a detached process with one allowlisted observation, never raw argv/output. */ -export async function observe(path: string, ctx: Context, generation: string, counter: Counter | null, send: Post = post, goal?: GoalObservation) { +export async function observe(path: string, ctx: Context, generation: string, counter: Counter | null, send: Post = post, goal?: GoalObservation, cycle?: CycleObservation) { if (counter !== null && (!validCounter(counter) || counter.count !== 1)) return { sent: false, reason: "invalid_observation" }; let heartbeat: Ping | null = null; let aggregate: Aggregate | null = null; @@ -136,7 +141,11 @@ export async function observe(path: string, ctx: Context, generation: string, co const blocked = blockedBy(state, ctx); if (blocked || !generation || generation !== state.generation) return false; if (state.day && state.day > today) return false; - try { goals = await recordGoalUsage(path + ".goals", generation, (ctx.now ?? new Date()).getTime(), goal); } + try { + const now = (ctx.now ?? new Date()).getTime(); + const intervals = cycle ? await cycleObservations(path + ".cycles", generation, now, cycle) : []; + goals = await recordGoalUsage(path + ".goals", generation, now, [...(goal ? [goal] : []), ...intervals]); + } catch { /* A damaged optional measurement cannot block other diagnostics. */ } // Flush only a completed day's local aggregate. No event times or per-install key leave this boundary. if (state.day && state.day < today && state.counters?.length) { diff --git a/loopx/control_plane/runtime/usage_statistics_cli.ts b/loopx/control_plane/runtime/usage_statistics_cli.ts index 58c56b0c8f..feff1703be 100644 --- a/loopx/control_plane/runtime/usage_statistics_cli.ts +++ b/loopx/control_plane/runtime/usage_statistics_cli.ts @@ -3,7 +3,10 @@ import { durationBucket, FEATURES, object } from "./usage_statistics_contract.ts import type { Context } from "./usage_statistics.ts"; import type { Counter } from "./usage_statistics_contract.ts"; -import { validGoalObservation } from "./usage_statistics_goal_contract.ts"; +import type { GoalObservation } from "./usage_statistics_goal_contract.ts"; +import { validGoalObservation, hostCategory } from "./usage_statistics_goal_contract.ts"; + +import { validCycle } from "./usage_statistics_cycles.ts"; try { let input = ""; @@ -20,8 +23,10 @@ try { result = await configure(request.path, ctx, request.action, request.notice); } else if (request.action === "start") { result = await observe(request.path, ctx, String(request.generation), null); - } else if (request.action === "goal" && validGoalObservation(request.observation, Date.now())) { - result = await observe(request.path, ctx, String(request.generation), null, undefined, request.observation); + } else if (request.action === "cycle" && validCycle(request.observation, Date.now())) { + result = await observe(request.path, ctx, String(request.generation), null, undefined, undefined, request.observation); + } else if (request.action === "goal" && object(request.observation) && validGoalObservation({ ...request.observation, host: hostCategory(request.observation.host) }, Date.now())) { + result = await observe(request.path, ctx, String(request.generation), null, undefined, { ...request.observation, host: hostCategory(request.observation.host) } as GoalObservation); } else if (request.action === "observe" && typeof request.feature === "string" && typeof request.elapsed_ms === "number" && Number.isFinite(request.elapsed_ms) && request.elapsed_ms >= 0) { result = await observe(request.path, ctx, String(request.generation), { diff --git a/loopx/control_plane/runtime/usage_statistics_codex.ts b/loopx/control_plane/runtime/usage_statistics_codex.ts new file mode 100644 index 0000000000..093b8eeea1 --- /dev/null +++ b/loopx/control_plane/runtime/usage_statistics_codex.ts @@ -0,0 +1,67 @@ +/** Read only timing envelopes from one explicitly bound session, in bounded chunks. */ +import { open } from "node:fs/promises"; +import type { GoalObservation, Host } from "./usage_statistics_goal_contract.ts"; +import { object } from "./usage_statistics_contract.ts"; +export type CodexCursor = { offset: number; inode: string; since: number; seen: number; skipping?: boolean; open?: { id: string; start: number; confirmed: number } }; +const BUDGET = 1024 * 1024; +function timestamp(value: unknown): number { return typeof value === "string" ? Date.parse(value) : NaN; } +export async function readCodexTiming(path: string, thread: string, previous: CodexCursor | undefined, now: number, key: string, host: Host) { + const file = await open(path, "r"); + try { + const stat = await file.stat(); + const header = Buffer.alloc(65536); + const head = await file.read(header, 0, header.length, 0); + const first = header.subarray(0, head.bytesRead).toString().split("\n")[0]; + const meta = JSON.parse(first); + if (meta.type !== "session_meta" || (meta.payload?.id ?? meta.payload?.session_id) !== thread) throw new Error("usage_session_identity_mismatch"); + const inode = `${stat.dev}:${stat.ino}`; + const reusable = previous?.inode === inode && previous.offset <= stat.size; + const cursor: CodexCursor = reusable ? structuredClone(previous!) : { offset: Math.max(0, stat.size - BUDGET), inode, since: now, seen: now }; + cursor.seen = now; + const buffer = Buffer.alloc(Math.min(BUDGET, stat.size - cursor.offset)); + const read = await file.read(buffer, 0, buffer.length, cursor.offset); + const bytes = buffer.subarray(0, read.bytesRead); + const last = bytes.lastIndexOf(10); + if (last < 0) { + // Large content records must not permanently strand later timing events. + if (read.bytesRead === BUDGET) { cursor.offset += read.bytesRead; cursor.skipping = true; cursor.open = undefined; } + return { cursor, observations: [] as GoalObservation[] }; + } + let body = bytes.subarray(0, last + 1).toString("utf8"); + if (cursor.skipping || !reusable && cursor.offset > 0) body = body.slice(body.indexOf("\n") + 1); + delete cursor.skipping; + cursor.offset += last + 1; + const observations: GoalObservation[] = []; + const emit = (start: number, end: number) => { + start = Math.max(start, cursor.since); + if (Number.isSafeInteger(start) && Number.isSafeInteger(end) && end > start && end <= now + 1000 && end - start <= 7 * 86400000) + observations.push({ key, start, end, measurement: "codex_turn", host }); + }; + for (const line of body.split("\n")) { + if (!line) continue; + let event: unknown; + try { event = JSON.parse(line); } catch { cursor.open = undefined; continue; } + if (!object(event) || event.type !== "event_msg" || !object(event.payload)) continue; + const payload = event.payload; + const id = typeof payload.turn_id === "string" && payload.turn_id.length <= 128 ? payload.turn_id : ""; + if (payload.type === "task_started" && id) { + const start = timestamp(payload.started_at ?? event.timestamp); + if (Number.isFinite(start)) cursor.open = { id, start, confirmed: Math.max(start, cursor.since) }; + } else if (["task_complete", "task_completed", "turn_aborted"].includes(String(payload.type)) && id) { + const start = timestamp(payload.started_at); + const end = timestamp(payload.completed_at ?? event.timestamp); + // New Codex records carry their own exact start; older records require + // the matching open Turn. Never pair an unrelated terminal by proximity. + const matched = cursor.open?.id === id ? cursor.open : undefined; + if (Number.isFinite(start)) emit(start, end); + else if (matched) emit(matched.confirmed, end); + if (matched) cursor.open = undefined; + } else if (cursor.open && payload.type === "token_count") { + const confirmed = timestamp(event.timestamp); + emit(cursor.open.confirmed, confirmed); + if (Number.isFinite(confirmed)) cursor.open.confirmed = Math.max(cursor.open.confirmed, confirmed); + } + } + return { cursor, observations }; + } finally { await file.close(); } +} diff --git a/loopx/control_plane/runtime/usage_statistics_cycles.ts b/loopx/control_plane/runtime/usage_statistics_cycles.ts new file mode 100644 index 0000000000..24e3f51c17 --- /dev/null +++ b/loopx/control_plane/runtime/usage_statistics_cycles.ts @@ -0,0 +1,68 @@ +/** Diagnostic cycle pairing shared by every Host; never admission or settlement authority. */ +import { readFile, chmod } from "node:fs/promises"; +import { atomicWriteJson } from "../effect_runtime_io.ts"; +import type { JsonObject } from "../effect_program.ts"; +import { object } from "./usage_statistics_contract.ts"; +import { hostCategory } from "./usage_statistics_goal_contract.ts"; +import type { GoalObservation } from "./usage_statistics_goal_contract.ts"; +import { readCodexTiming } from "./usage_statistics_codex.ts"; +import type { CodexCursor } from "./usage_statistics_codex.ts"; +export type CycleObservation = { key: string; lane: string; turn: string | null; phase: "start" | "spend"; at: number; host: string; codex?: { path: string; id: string } }; +type Cycle = { id: string; start?: number; end?: number; touched: number; exact: boolean; host: string; floor?: number }; +type State = { generation: string; cycles: Cycle[]; cursors: Record }; +const WEEK = 7 * 86400000; +export function validCycle(value: unknown, now: number): value is CycleObservation { + return object(value) && [value.key, value.lane].every(v => typeof v === "string" && /^[a-f0-9]{64}$/.test(v)) + && (value.turn === null || typeof value.turn === "string" && /^[a-f0-9]{64}$/.test(value.turn)) + && ["start", "spend"].includes(String(value.phase)) && Number.isSafeInteger(value.at) + && Number(value.at) <= now + 1000 && Number(value.at) >= now - 86400000 && typeof value.host === "string"; +} +export async function cycleObservations(path: string, generation: string, now: number, input: CycleObservation): Promise { + let state: State = { generation, cycles: [], cursors: {} }; + try { + const text = await readFile(path, "utf8"); + if (text.length > 256000) throw new Error("usage_cycle_limit"); + const prior = JSON.parse(text); + if (prior.generation === generation && Array.isArray(prior.cycles) && prior.cycles.length <= 128 && object(prior.cursors) && Object.keys(prior.cursors).length <= 64) state = prior; + } catch (error) { if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; } + state.cycles = state.cycles.filter(c => now - c.touched <= WEEK); + const id = `${input.key}:${input.lane}:${input.turn ?? "lane"}`; + let cycle = state.cycles.find(c => c.id === id); + if (!cycle && input.phase === "start" && state.cycles.length >= 128) { + // Completed diagnostics cannot crowd out new work forever. An execution + // replay never supplies a fresh spend marker at the CLI boundary. + const closed = state.cycles.filter(c => c.end !== undefined).sort((a,b)=>a.touched-b.touched)[0]; + if (closed) state.cycles = state.cycles.filter(c => c !== closed); + } + if (!cycle && state.cycles.length < 128) { + cycle = { id, touched: input.at, exact: input.turn !== null, host: hostCategory(input.host) }; state.cycles.push(cycle); + } + const observations: GoalObservation[] = []; + const host = hostCategory(input.host); + if (cycle) { + if (!cycle.exact && input.phase === "start" && cycle.end !== undefined && input.at > cycle.end) { cycle.floor = cycle.end; delete cycle.start; delete cycle.end; cycle.host = host; } + if (input.at > (cycle.floor ?? 0)) { + if (input.phase === "start") cycle.start = Math.min(cycle.start ?? input.at, input.at); + else cycle.end = Math.min(cycle.end ?? input.at, input.at); + cycle.touched = Math.max(cycle.touched, input.at); + } + if (cycle.start !== undefined && cycle.end !== undefined && cycle.end >= cycle.start && cycle.end - cycle.start <= WEEK) + observations.push({ key: input.key, start: cycle.start, end: cycle.end, measurement: "quota_cycle", host: hostCategory(cycle.host) }); + } + if (input.codex && typeof input.codex.path === "string" && typeof input.codex.id === "string") { + // Keys and paths stay machine-local. No transcript strings enter state or wire data. + const cursorKey = `${input.key}:${input.lane}:${input.codex.id}`; + state.cursors = Object.fromEntries(Object.entries(state.cursors).filter(([, c])=>now-c.seen<=WEEK)); + if (!state.cursors[cursorKey] && Object.keys(state.cursors).length >= 64) { + const oldest = Object.entries(state.cursors).sort((a,b)=>a[1].seen-b[1].seen)[0][0]; + delete state.cursors[oldest]; + } + try { + const timing = await readCodexTiming(input.codex.path, input.codex.id, state.cursors[cursorKey], now, input.key, host); + state.cursors[cursorKey] = timing.cursor; + observations.push(...timing.observations); + } catch { /* Finer Host timing is optional; the universal cycle still works. */ } + } + await atomicWriteJson(path, state as unknown as JsonObject); await chmod(path, 0o600); + return observations; +} diff --git a/loopx/control_plane/runtime/usage_statistics_goal_contract.ts b/loopx/control_plane/runtime/usage_statistics_goal_contract.ts index 8b05a65811..09b535dd8d 100644 --- a/loopx/control_plane/runtime/usage_statistics_goal_contract.ts +++ b/loopx/control_plane/runtime/usage_statistics_goal_contract.ts @@ -3,9 +3,18 @@ import { object } from "./usage_statistics_contract.ts"; export const GOAL_SCHEMA = "loopx_goal_usage_aggregate_v1"; export const GOAL_DURATIONS = ["lt_1m", "lt_10m", "lt_1h", "lt_6h", "lt_1d", "lt_7d", "lt_30d", "gte_30d"] as const; export type GoalDuration = typeof GOAL_DURATIONS[number]; -export type GoalCount = { span: GoalDuration; execution: GoalDuration; count: number }; +export const MEASUREMENTS = ["host_call", "codex_turn", "quota_cycle"] as const; +export const HOSTS = ["codex_app", "codex_cli", "claude_code", "dsh", "opencode", "trae", "other", "unknown"] as const; +export type Measurement = typeof MEASUREMENTS[number]; +export type Host = typeof HOSTS[number]; +export function hostCategory(value: unknown): Host { + if ((HOSTS as readonly unknown[]).includes(value)) return value as Host; + const aliases: Record = { "codex-app": "codex_app", "codex-app-ssh": "codex_app", codex_app: "codex_app", codex_app_heartbeat: "codex_app", codex_app_ssh_goal: "codex_app", "codex-cli": "codex_cli", "codex-cli-tui": "codex_cli", codex_cli: "codex_cli", codex: "codex_cli", "codex-ide-plugin": "codex_cli", "deepseek-harness-native": "dsh", "claude-code": "claude_code", claude_code: "claude_code", dsh: "dsh", opencode: "opencode", trae_app: "trae", "generic-cli": "other", generic_cli: "other" }; + return typeof value === "string" && Object.hasOwn(aliases, value) ? aliases[value] : "unknown"; +} +export type GoalCount = { measurement: Measurement; host: Host; span: GoalDuration; duration: GoalDuration; count: number }; export type GoalAggregate = { schema: typeof GOAL_SCHEMA; counters: GoalCount[] }; -export type GoalObservation = { key: string; start: number; end: number }; +export type GoalObservation = { key: string; start: number; end: number; measurement: Measurement; host: Host }; const DAY = 86400000; export function goalDuration(ms: number): GoalDuration { const limits = [60000, 600000, 3600000, 21600000, DAY, 7 * DAY, 30 * DAY]; @@ -13,22 +22,24 @@ export function goalDuration(ms: number): GoalDuration { } export function validGoalAggregate(value: unknown): value is GoalAggregate { if (!object(value) || Object.keys(value).sort().join() !== "counters,schema" || value.schema !== GOAL_SCHEMA - || !Array.isArray(value.counters) || !value.counters.length || value.counters.length > 64) return false; + || !Array.isArray(value.counters) || !value.counters.length || value.counters.length > 128) return false; const keys = new Set(); return value.counters.every(row => { - if (!object(row) || Object.keys(row).sort().join() !== "count,execution,span" - || !(GOAL_DURATIONS as readonly unknown[]).includes(row.span) || !(GOAL_DURATIONS as readonly unknown[]).includes(row.execution) + if (!object(row) || Object.keys(row).sort().join() !== "count,duration,host,measurement,span" + || !(MEASUREMENTS as readonly unknown[]).includes(row.measurement) || !(HOSTS as readonly unknown[]).includes(row.host) + || !(GOAL_DURATIONS as readonly unknown[]).includes(row.span) || !(GOAL_DURATIONS as readonly unknown[]).includes(row.duration) || !Number.isInteger(row.count) || Number(row.count) < 1 || Number(row.count) > 128) return false; - const key = `${row.span}:${row.execution}`; + const key = `${row.measurement}:${row.host}:${row.span}:${row.duration}`; if (keys.has(key)) return false; keys.add(key); return true; }); } export function validGoalObservation(value: unknown, now: number): value is GoalObservation { - return object(value) && Object.keys(value).sort().join() === "end,key,start" + return object(value) && Object.keys(value).sort().join() === "end,host,key,measurement,start" + && (MEASUREMENTS as readonly unknown[]).includes(value.measurement) && (HOSTS as readonly unknown[]).includes(value.host) && typeof value.key === "string" && /^[a-f0-9]{64}$/.test(value.key) && Number.isSafeInteger(value.start) && Number.isSafeInteger(value.end) && Number(value.start) > 0 && Number(value.start) <= Number(value.end) - && Number(value.end) <= now + 1000 && Number(value.end) >= now - DAY - && Number(value.end) - Number(value.start) <= 120000; + && Number(value.end) <= now + 1000 && Number(value.end) >= now - 7 * DAY + && Number(value.end) - Number(value.start) <= (value.measurement === "host_call" ? 120000 : 7 * DAY); } diff --git a/loopx/control_plane/runtime/usage_statistics_goals.ts b/loopx/control_plane/runtime/usage_statistics_goals.ts index d70f82c02f..34f0d916ce 100644 --- a/loopx/control_plane/runtime/usage_statistics_goals.ts +++ b/loopx/control_plane/runtime/usage_statistics_goals.ts @@ -1,13 +1,13 @@ -/** Local, lossy execution observations; never a Goal lifecycle or billing authority. */ +/** Local, lossy duration observations; never a Goal lifecycle or billing authority. */ import { readFile, chmod } from "node:fs/promises"; import { atomicWriteJson } from "../effect_runtime_io.ts"; import type { JsonObject } from "../effect_program.ts"; import { object } from "./usage_statistics_contract.ts"; -import { GOAL_SCHEMA, goalDuration, validGoalObservation } from "./usage_statistics_goal_contract.ts"; -import type { GoalAggregate, GoalObservation, GoalCount } from "./usage_statistics_goal_contract.ts"; +import { GOAL_SCHEMA, HOSTS, MEASUREMENTS, goalDuration, validGoalObservation } from "./usage_statistics_goal_contract.ts"; +import type { GoalAggregate, GoalObservation, GoalCount, Measurement, Host } from "./usage_statistics_goal_contract.ts"; type Interval = [number, number]; -type MeasuredGoal = { key: string; first: number; last: number; intervals: Interval[]; total: number; watermark: number; day: string; reported?: string }; +type MeasuredGoal = { key: string; measurement: Measurement; host: Host; first: number; last: number; intervals: Interval[]; total: number; watermark: number; day: string; reported?: string }; type GoalState = { generation: string; goals: MeasuredGoal[] }; const DAY = 86400000; const MAX_GOALS = 64; @@ -29,6 +29,7 @@ async function load(path: string, generation: string): Promise { const value = JSON.parse(text) as GoalState; if (value.generation !== generation) return { generation, goals: [] }; if (!Array.isArray(value.goals) || value.goals.length > MAX_GOALS || value.goals.some(g => !object(g) + || !MEASUREMENTS.includes(g.measurement) || !HOSTS.includes(g.host) || typeof g.key !== "string" || !/^[a-f0-9]{64}$/.test(g.key) || ![g.first, g.last, g.total, g.watermark].every(n => Number.isSafeInteger(n) && n >= 0) || g.first > g.last || typeof g.day !== "string" || !/^\d{4}-\d{2}-\d{2}$/.test(g.day) @@ -45,9 +46,9 @@ function snapshot(goals: MeasuredGoal[]): GoalAggregate | null { const rows = new Map(); for (const goal of goals) { const span = goalDuration(goal.last - goal.first); - const execution = goalDuration(goal.total + goal.intervals.reduce((sum, [a, b]) => sum + b - a, 0)); - const key = `${span}:${execution}`; - const row = rows.get(key) ?? { span, execution, count: 0 }; + const duration = goalDuration(goal.total + goal.intervals.reduce((sum, [a, b]) => sum + b - a, 0)); + const key = `${goal.measurement}:${goal.host}:${span}:${duration}`; + const row = rows.get(key) ?? { measurement: goal.measurement, host: goal.host, span, duration, count: 0 }; row.count++; rows.set(key, row); } return rows.size ? { schema: GOAL_SCHEMA, counters: [...rows.values()] } : null; @@ -56,37 +57,44 @@ export async function goalPreview(path: string, generation: string): Promise g.reported !== g.day)); } /** Caller holds the common consent lock. Claims before sending: no retry identifiers. */ -export async function recordGoalUsage(path: string, generation: string, now: number, observation?: GoalObservation): Promise { +export async function recordGoalUsage(path: string, generation: string, now: number, observation?: GoalObservation | GoalObservation[]): Promise { const state = await load(path, generation); const today = new Date(now).toISOString().slice(0, 10); - const ready = state.goals.filter(g => g.day < today && g.reported !== g.day && now - Date.parse(g.day) < 8 * DAY); - const payload = snapshot(ready); - for (const goal of ready) goal.reported = goal.day; // A quiet Goal is not reported again each day. Keep observed cumulative time // for up to 90 days of inactivity, then restart measurement if it returns. state.goals = state.goals.filter(g => now - g.last < 90 * DAY); - if (observation && validGoalObservation(observation, now)) { - const observedDay = new Date(observation.end).toISOString().slice(0, 10); - let goal = state.goals.find(g => g.key === observation.key); - if (!goal && state.goals.length < MAX_GOALS) { - goal = { key: observation.key, first: observation.start, last: observation.end, intervals: [], total: 0, watermark: 0, day: observedDay }; - state.goals.push(goal); - } - if (goal && observation.start >= goal.watermark && goal.day <= observedDay && goal.reported !== observedDay) { - // Compact only intervals older than the accepted late-observation window. - const boundary = now - DAY - 120000; - const retired = goal.intervals.filter(([, b]) => b < boundary); - goal.total += retired.reduce((sum, [a, b]) => sum + b - a, 0); - goal.intervals = goal.intervals.filter(([, b]) => b >= boundary); - goal.watermark = Math.max(goal.watermark, boundary); - const intervals = union(goal.intervals, [observation.start, observation.end]); - if (intervals.length <= MAX_INTERVALS) { - goal.intervals = intervals; - goal.first = Math.min(goal.first, observation.start); goal.last = Math.max(goal.last, observation.end); - goal.day = observedDay; + const observations = (Array.isArray(observation) ? observation : observation ? [observation] : []).filter(item => validGoalObservation(item, now)); + function apply(items: GoalObservation[]) { + for (const observation of items) { + const observedDay = new Date(observation.end).toISOString().slice(0, 10); + let goal = state.goals.find(g => g.key === observation.key && g.measurement === observation.measurement && g.host === observation.host); + if (!goal && state.goals.length < MAX_GOALS) { + goal = { key: observation.key, measurement: observation.measurement, host: observation.host, first: observation.start, last: observation.end, intervals: [], total: 0, watermark: 0, day: observedDay }; + state.goals.push(goal); + } + if (goal && observation.start >= goal.watermark && goal.day <= observedDay && goal.reported !== observedDay) { + // Compact only intervals older than the accepted late-observation window. + const boundary = now - 14 * DAY; + const retired = goal.intervals.filter(([, b]) => b < boundary); + goal.total += retired.reduce((sum, [a, b]) => sum + b - a, 0); + goal.intervals = goal.intervals.filter(([, b]) => b >= boundary); + goal.watermark = Math.max(goal.watermark, boundary); + const intervals = union(goal.intervals, [observation.start, observation.end]); + if (intervals.length <= MAX_INTERVALS) { + goal.intervals = intervals; + goal.first = Math.min(goal.first, observation.start); goal.last = Math.max(goal.last, observation.end); + goal.day = observedDay; + } } } } + // Apply late terminal evidence before claiming that day's snapshot. Current + // day's intervals are applied afterwards so they cannot erase a closed day. + apply(observations.filter(item => new Date(item.end).toISOString().slice(0, 10) < today)); + const ready = state.goals.filter(g => g.day < today && g.reported !== g.day && now - Date.parse(g.day) < 8 * DAY); + const payload = snapshot(ready); + for (const goal of ready) goal.reported = goal.day; + apply(observations.filter(item => new Date(item.end).toISOString().slice(0, 10) === today)); await atomicWriteJson(path, state as unknown as JsonObject); await chmod(path, 0o600); return payload; } diff --git a/loopx/control_plane/turn_driver/executor.py b/loopx/control_plane/turn_driver/executor.py index bb0fbd1c64..c39ce99f3c 100644 --- a/loopx/control_plane/turn_driver/executor.py +++ b/loopx/control_plane/turn_driver/executor.py @@ -772,7 +772,7 @@ def _host_result_stage( if confirm_start is not None: confirm_start() from ...usage_goal import observe_goal_execution - with observe_goal_execution(usage_runtime_root or project, usage_goal_id): + with observe_goal_execution(usage_runtime_root or project, usage_goal_id, host=str((plan.get("host") or {}).get("kind") or "unknown")): host_observation = ( _run_host_runner(request, runner=host_runner) if host_runner is not None diff --git a/loopx/usage_goal.py b/loopx/usage_goal.py index ec67210e54..7fe40f31b4 100644 --- a/loopx/usage_goal.py +++ b/loopx/usage_goal.py @@ -5,7 +5,7 @@ """ from __future__ import annotations -from contextlib import contextmanager +from contextlib import closing, contextmanager import hashlib import json import os @@ -18,7 +18,7 @@ @contextmanager -def observe_goal_execution(runtime_root: Path, goal_id: str) -> Iterator[None]: +def observe_goal_execution(runtime_root: Path, goal_id: str, *, host: str = "unknown") -> Iterator[None]: stop = threading.Event() publish = None try: @@ -50,7 +50,7 @@ def checkpoint() -> None: previous = elapsed usage_ping._detach(usage_ping._request( "goal", path, generation=generation, - observation={"key": key, "start": start, "end": wall + elapsed}, + observation={"key": key, "start": start, "end": wall + elapsed, "measurement": "host_call", "host": host}, )) finally: lock.release() @@ -74,3 +74,91 @@ def periodically() -> None: # harmless: the TS interval union deduplicates it. if publish is not None: publish() + + +def observe_quota_cycle(*, registry_path: Path, runtime_root: Path, goal_id: str, + agent_id: str | None, turn_id: str | None, phase: str, + at: int, host: str) -> None: + """Detach binding discovery and session metadata lookup from quota latency.""" + try: + import sys + state = json.loads(usage_ping.state_path().read_text()) + if state.get("consent") == "disabled" or not state.get("generation"): + return + if os.environ.get("LOOPX_USAGE_PING") == "0" or os.environ.get("DO_NOT_TRACK") == "1" or os.environ.get("CI") == "true": + return + usage_ping._detach({"registry": str(registry_path), "runtime": str(runtime_root), + "goal": goal_id, "agent": agent_id, "turn": turn_id, + "phase": phase, "at": at, "host": host, + "path": str(usage_ping.state_path()), "generation": state["generation"]}, + command=[sys.executable, "-m", "loopx.usage_goal"]) + except Exception: + pass + + +def _bound_codex_session(registry_path: Path, goal_id: str, agent_id: str | None): + """Use accepted exact binding and only the selected Codex home; never infer by cwd.""" + import sqlite3 + from .control_plane.projects.registry_codec import load_registry + from .registry import find_registry_goal + from .thread_agent_binding import collect_accepted_bindings, resolve_thread_agent_binding + goal = find_registry_goal(load_registry(registry_path), goal_id) + if not goal or not agent_id: + return None + bindings = [b for b in collect_accepted_bindings([goal]) + if b["agent_id"] == agent_id and b["host_surface"] in + {"codex-app", "codex-app-ssh", "codex-cli-tui", "codex-ide-plugin"}] + current = os.environ.get("CODEX_THREAD_ID") + if current: + bindings = [b for b in bindings if b["thread_id"] == current] + if len(bindings) != 1: + return None + binding = bindings[0] + if resolve_thread_agent_binding(goal, host_surface=binding["host_surface"], thread_id=binding["thread_id"])["status"] != "bound": + return None + home = Path(os.environ.get("CODEX_HOME") or "~/.codex").expanduser().resolve() + # SQLite is a read-only Host metadata adapter, not LoopX state authority. + for database in sorted(home.glob("state_*.sqlite"), reverse=True)[:4]: + try: + with closing(sqlite3.connect(database.as_uri() + "?mode=ro", uri=True, timeout=0.1)) as connection: + row = connection.execute("SELECT rollout_path FROM threads WHERE id = ?", (binding["thread_id"],)).fetchone() + if row: + path = Path(row[0]).resolve() + if any(path.is_relative_to(home / directory) for directory in ("sessions", "archived_sessions")): + return {"path": str(path), "id": binding["thread_id"]}, binding["host_surface"] + except (OSError, ValueError, sqlite3.Error): + continue + return None + + +def _dispatch_cycle(request) -> None: + path = Path(request["path"]) + state = json.loads(path.read_text()) + if state.get("consent") == "disabled" or state.get("generation") != request["generation"]: + return + generation = request["generation"] + + def digest(*parts): + return hashlib.sha256(json.dumps([generation, *parts]).encode()).hexdigest() + + observation = {"key": digest(str(Path(request["runtime"]).resolve()), request["goal"]), + "lane": digest(request["agent"] or "unscoped"), + "turn": digest(request["turn"]) if request["turn"] else None, + "phase": request["phase"], "at": request["at"], "host": request["host"]} + try: + binding = _bound_codex_session(Path(request["registry"]), request["goal"], request["agent"]) + if binding: + observation["codex"], observation["host"] = binding + except Exception: + pass # Common cycle collection must survive unavailable Host metadata. + usage_ping._detach(usage_ping._request("cycle", path, generation=generation, observation=observation)) + + +if __name__ == "__main__": + import sys + try: + raw = sys.stdin.buffer.read(8193) + if len(raw) <= 8192: + _dispatch_cycle(json.loads(raw)) + except Exception: + pass # No private paths or transcript errors on CLI output. diff --git a/loopx/usage_ping.py b/loopx/usage_ping.py index d21f6c2a48..402dae5d38 100644 --- a/loopx/usage_ping.py +++ b/loopx/usage_ping.py @@ -71,7 +71,7 @@ def begin(command: str) -> tuple[str, float] | None: state = json.loads(path.read_text()) if path.exists() else {} if state.get("consent") == "disabled": return None - if (state.get("notice") or {}).get("version") != 2: + if (state.get("notice") or {}).get("version") != 3: # Unattended machines remain silent until the owner sees the notice # or explicitly enables from CLI/settings. JSON stdout stays clean. if not sys.stderr.isatty(): @@ -112,14 +112,14 @@ def finish(ticket: tuple[str, float] | None, command: str, code: int, error: Bas pass # Telemetry cannot replace the command's result. -def _detach(request: dict[str, Any]) -> None: +def _detach(request: dict[str, Any], *, command: list[str] | None = None) -> None: kwargs: dict[str, Any] = {"stdin": subprocess.PIPE, "stdout": subprocess.DEVNULL, "stderr": subprocess.DEVNULL, "close_fds": True} if os.name == "nt": kwargs["creationflags"] = getattr(subprocess, "DETACHED_PROCESS", 0) | getattr(subprocess, "CREATE_NEW_PROCESS_GROUP", 0) else: kwargs["start_new_session"] = True - child = subprocess.Popen(_command(), **kwargs) + child = subprocess.Popen(command or _command(), **kwargs) assert child.stdin is not None child.stdin.write(json.dumps(request).encode()) child.stdin.close() diff --git a/tests/control_plane/test_quota_settlement_cli.py b/tests/control_plane/test_quota_settlement_cli.py index cb9e642aeb..491f103638 100644 --- a/tests/control_plane/test_quota_settlement_cli.py +++ b/tests/control_plane/test_quota_settlement_cli.py @@ -639,8 +639,22 @@ def _strip_heartbeat_workspace_causality(runtime: Path) -> None: def test_gitless_goal_refresh_and_quota_spend_settle_end_to_end( - tmp_path: Path, + tmp_path: Path, monkeypatch, ) -> None: + # Production quota CLI -> detached Python discovery -> TS cycle owner. + # The isolated home also proves telemetry never reads the operator's sessions. + from loopx import usage_ping + import time + home = tmp_path / "isolated-home" + machine = home / ".codex" / "loopx" + machine.mkdir(parents=True) + monkeypatch.setenv("HOME", str(home)) + monkeypatch.setenv("CODEX_HOME", str(home / ".codex")) + monkeypatch.setenv("LOOPX_USAGE_PING_ENDPOINT", "http://127.0.0.1:1/v1/ping") + for key in ("CI", "DO_NOT_TRACK", "LOOPX_USAGE_PING", "LOOPX_USAGE_POLICY"): + monkeypatch.delenv(key, raising=False) + usage_path = machine / "usage-ping.json" + usage_ping.control("enable", usage_path) project, runtime, registry_path = _write_fixture( tmp_path, required_capability="filesystem_write", @@ -675,6 +689,18 @@ def test_gitless_goal_refresh_and_quota_spend_settle_end_to_end( == "required" ) + # Preview/failed spend before validated delivery must not finish measurement. + cycles_path = Path(str(usage_path) + ".cycles") + deadline = time.monotonic() + 8 + while not cycles_path.exists() and time.monotonic() < deadline: + time.sleep(0.03) + assert "start" in json.loads(cycles_path.read_text())["cycles"][0] + for execute in (False, True): + _run_cli(registry_path, runtime, "quota", "spend-slot", "--goal-id", GOAL_ID, + "--slots", "1", "--source", "heartbeat", *binding, + *(["--execute"] if execute else []), "--scan-path", str(project), cwd=project) + assert "end" not in json.loads(cycles_path.read_text())["cycles"][0] + refresh_rc, refresh = _run_cli( registry_path, runtime, @@ -726,6 +752,26 @@ def test_gitless_goal_refresh_and_quota_spend_settle_end_to_end( assert spend["delivery_workspace_validated"] is True assert spend["delivery_workspace"]["workspace_identity"] == f"loopx:{GOAL_ID}" assert _spend_run_count(runtime) == 1 + cycles_path = Path(str(usage_path) + ".cycles") + deadline = time.monotonic() + 8 + while time.monotonic() < deadline: + if cycles_path.exists(): + cycles = json.loads(cycles_path.read_text())["cycles"] + if cycles and cycles[0].get("end"): + break + time.sleep(0.03) + else: + raise AssertionError("public quota/spend CLI did not complete a telemetry cycle") + assert len(cycles) == 1 and cycles[0]["exact"] is True + assert cycles[0]["start"] < cycles[0]["end"] + before_replay = cycles_path.read_bytes() + replay_rc, replay = _run_cli(registry_path, runtime, "quota", "spend-slot", "--goal-id", GOAL_ID, + "--slots", "1", "--source", "heartbeat", *binding, + "--execute", "--scan-path", str(project), cwd=project) + assert replay_rc == 0 and replay["idempotent_replay"] is True + assert cycles_path.read_bytes() == before_replay + assert usage_ping.control("status", usage_path)["goal_preview"]["counters"][0]["measurement"] == "quota_cycle" + usage_ping.control("disable", usage_path) def test_codex_app_refresh_stages_validated_memory_and_spend_finalizes_hook( diff --git a/tests/control_plane_ts/usage_statistics_goals.test.ts b/tests/control_plane_ts/usage_statistics_goals.test.ts index 410ab6275f..e853f4809f 100644 --- a/tests/control_plane_ts/usage_statistics_goals.test.ts +++ b/tests/control_plane_ts/usage_statistics_goals.test.ts @@ -16,7 +16,7 @@ async function fixture(t: test.TestContext) { return join(root, "usage.json"); } const ctx = (now: number): Context => ({ env: { LOOPX_USAGE_PING_ENDPOINT: "http://127.0.0.1:1/v1/ping" }, version: "1.0.0", python: "3.13", channel: "source", now: new Date(now) }); -const checkpoint = (start: number, end: number) => ({ key, start, end }); +const checkpoint = (start: number, end: number) => ({ key, start, end, measurement: "host_call" as const, host: "unknown" as const }); test("union counts concurrent/nested intervals once and preserves idle gaps under replay and reordering", () => { const observations: [number, number][] = [[0, 10], [4, 6], [8, 12], [20, 25], [0, 10]]; @@ -32,7 +32,7 @@ test("unfinished Goals report once per observed day; concurrent retries add only await recordGoalUsage(path, "generation", base + 90000, checkpoint(base, base + 60000)); // Retry after a ten-hour pause: span includes the pause, execution does not. await recordGoalUsage(path, "generation", base + 10 * 3600000 + 30000, checkpoint(base + 10 * 3600000, base + 10 * 3600000 + 30000)); - const expected = { schema: GOAL_SCHEMA, counters: [{ span: "lt_1d", execution: "lt_10m", count: 1 }] }; + const expected = { schema: GOAL_SCHEMA, counters: [{ measurement: "host_call", host: "unknown", span: "lt_1d", duration: "lt_10m", count: 1 }] }; assert.deepEqual(await goalPreview(path, "generation"), expected); assert.deepEqual(await recordGoalUsage(path, "generation", base + day), expected); assert.equal(await recordGoalUsage(path, "generation", base + day), null); @@ -43,13 +43,13 @@ test("unfinished Goals report once per observed day; concurrent retries add only test("restart retains measurements, compaction retains totals and rejects old replay", async t => { const path = await fixture(t); await recordGoalUsage(path, "g", base + 60000, checkpoint(base, base + 60000)); - await recordGoalUsage(path, "g", base + 3 * day, checkpoint(base + 3 * day - 60000, base + 3 * day)); + await recordGoalUsage(path, "g", base + 15 * day, checkpoint(base + 15 * day - 60000, base + 15 * day)); const state = JSON.parse(await readFile(path, "utf8")); assert.equal(state.goals[0].total, 60000); assert.equal(state.goals[0].intervals.length, 1); - await recordGoalUsage(path, "g", base + 3 * day, checkpoint(base, base + 60000)); + await recordGoalUsage(path, "g", base + 15 * day, checkpoint(base, base + 60000)); assert.equal(JSON.parse(await readFile(path, "utf8")).goals[0].total, 60000); - assert.deepEqual((await goalPreview(path, "g"))?.counters, [{ span: "lt_7d", execution: "lt_10m", count: 1 }]); + assert.deepEqual((await goalPreview(path, "g"))?.counters, [{ measurement: "host_call", host: "unknown", span: "lt_30d", duration: "lt_10m", count: 1 }]); // Consent generation reset has no continuity with the previous measurement. assert.equal(await goalPreview(path, "new-generation"), null); }); @@ -58,7 +58,7 @@ test("no execution after last confirmed prefix is inferred when host disappears" const path = await fixture(t); await recordGoalUsage(path, "g", base + 10000, checkpoint(base, base + 10000)); const result = await recordGoalUsage(path, "g", base + 5 * day); - assert.deepEqual(result?.counters, [{ span: "lt_1m", execution: "lt_1m", count: 1 }]); + assert.deepEqual(result?.counters, [{ measurement: "host_call", host: "unknown", span: "lt_1m", duration: "lt_1m", count: 1 }]); assert.equal(await recordGoalUsage(path, "g", base + 6 * day), null); assert.equal(validGoalObservation(checkpoint(base, base + 10 * day), base + 10 * day), false); assert.equal(validGoalObservation(checkpoint(base + 10, base), base), false); @@ -97,7 +97,7 @@ test("new scope needs notice again without undoing explicit disable", async t => }); test("closed contract rejects identifiers, timestamps and extra fields; long durations are not clipped at a minute", () => { - const valid = { schema: GOAL_SCHEMA, counters: [{ span: "gte_30d", execution: "lt_7d", count: 5 }] }; + const valid = { schema: GOAL_SCHEMA, counters: [{ measurement: "host_call", host: "unknown", span: "gte_30d", duration: "lt_7d", count: 5 }] }; assert.ok(validGoalAggregate(valid)); for (const extra of [{ goal_id: "private" }, { install_id: key }, { timestamp: base }]) { assert.equal(validGoalAggregate({ ...valid, ...extra }), false); diff --git a/tests/control_plane_ts/usage_statistics_sources.test.ts b/tests/control_plane_ts/usage_statistics_sources.test.ts new file mode 100644 index 0000000000..b06e6e3efe --- /dev/null +++ b/tests/control_plane_ts/usage_statistics_sources.test.ts @@ -0,0 +1,141 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { mkdtemp, writeFile, appendFile, readFile, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { cycleObservations } from "../../loopx/control_plane/runtime/usage_statistics_cycles.ts"; +import type { CycleObservation } from "../../loopx/control_plane/runtime/usage_statistics_cycles.ts"; +import { readCodexTiming } from "../../loopx/control_plane/runtime/usage_statistics_codex.ts"; +import { recordGoalUsage, goalPreview } from "../../loopx/control_plane/runtime/usage_statistics_goals.ts"; +import { configure, inspect, observe } from "../../loopx/control_plane/runtime/usage_statistics.ts"; +const at = Date.parse("2026-09-01T10:00:00Z"); +const key = "a".repeat(64), lane = "b".repeat(64), turn = "c".repeat(64); +const cycle = (phase: "start" | "spend", time: number, extra = {}): CycleObservation => ({ key, lane, turn, phase, at: time, host: "codex-app", ...extra }); +async function fixture(t: test.TestContext) { + const root = await mkdtemp(join(tmpdir(), "loopx-source-usage-")); + t.after(()=>rm(root,{recursive:true,force:true})); return root; +} +const event = (type: string, time: number, fields = {}) => JSON.stringify({ type: "event_msg", timestamp: new Date(time).toISOString(), payload: { type, ...fields } }) + "\n"; + +test("universal cycles preserve first quota and first successful spend across all Host labels and reordered transport", async t => { + const root = await fixture(t); + for (const host of ["codex-app", "claude-code", "dsh", "opencode", "generic-cli"]) { + const path = join(root, host); + assert.deepEqual(await cycleObservations(path, "g", at, cycle("start", at, {host})), []); + await cycleObservations(path, "g", at+100, cycle("start", at+100, {host})); + const result = await cycleObservations(path, "g", at+1000, cycle("spend", at+1000, {host})); + assert.equal(result[0].start, at); assert.equal(result[0].end, at+1000); assert.equal(result[0].measurement, "quota_cycle"); + const replay = await cycleObservations(path, "g", at+2000, cycle("spend", at+2000, {host})); + assert.equal(replay[0].end, at+1000); + } + const path = join(root,"reordered"); + await cycleObservations(path,"g",at+1000,cycle("spend",at+1000)); + const delayed = await cycleObservations(path,"g",at+1100,cycle("start",at)); + assert.equal(delayed[0].end-delayed[0].start,1000); + assert.deepEqual(await cycleObservations(join(root,"missing"),"g",at,cycle("spend",at)),[]); +}); + +test("missing spend never invents elapsed work; exact cycles cannot borrow another lane or turn", async t => { + const path=join(await fixture(t),"cycles"); + await cycleObservations(path,"g",at,cycle("start",at)); + assert.deepEqual(await cycleObservations(path,"g",at+1000,cycle("spend",at+1000,{lane:"d".repeat(64)})),[]); + assert.deepEqual(await cycleObservations(path,"g",at+1000,cycle("spend",at+1000,{turn:"e".repeat(64)})),[]); + assert.deepEqual(await cycleObservations(path,"g",at+8*86400000,cycle("spend",at+8*86400000)),[]); +}); + +test("bound Codex timing is incremental, content-free, identity checked and handles interrupted/partial records", async t => { + const path=join(await fixture(t),"rollout.jsonl"); + await writeFile(path,JSON.stringify({type:"session_meta",payload:{id:"thread"}})+"\n"+event("task_started",at-1000,{turn_id:"first",started_at:new Date(at-1000).toISOString()})); + let result=await readCodexTiming(path,"thread",undefined,at,key,"codex_app"); + assert.equal(result.observations.length,0); + await appendFile(path,event("task_complete",at+60000,{turn_id:"first",started_at:new Date(at-1000).toISOString(),completed_at:new Date(at+60000).toISOString(),last_agent_message:"DO NOT PERSIST PRIVATE CONTENT"})); + result=await readCodexTiming(path,"thread",result.cursor,at+60001,key,"codex_app"); + assert.equal(result.observations[0].start,at); assert.equal(result.observations[0].end,at+60000); + assert.ok(!JSON.stringify(result).includes("PRIVATE CONTENT")); + assert.equal((await readCodexTiming(path,"thread",result.cursor,at+60001,key,"codex_app")).observations.length,0); + await assert.rejects(readCodexTiming(path,"wrong",undefined,at,key,"codex_app"),/identity_mismatch/); + const partial=event("turn_aborted",at+120000,{turn_id:"second",started_at:new Date(at+70000).toISOString()}); + await appendFile(path,partial.slice(0,-2)); + const incomplete=await readCodexTiming(path,"thread",result.cursor,at+120000,key,"codex_app"); + assert.equal(incomplete.observations.length,0); + await appendFile(path,partial.slice(-2)); + const terminal=await readCodexTiming(path,"thread",incomplete.cursor,at+120001,key,"codex_app"); + assert.equal(terminal.observations[0].end-terminal.observations[0].start,50000); +}); + +test("three clocks remain separate and same-source overlaps are deduplicated", async t => { + const path=join(await fixture(t),"goals"); + const intervals = ["quota_cycle","codex_turn","host_call"].map(measurement=>({key,start:at,end:at+60000,measurement,host:"codex_app"} as const)); + await recordGoalUsage(path,"g",at+60001,intervals as Parameters[3]); + await recordGoalUsage(path,"g",at+60001,intervals as Parameters[3]); + const rows=(await goalPreview(path,"g"))!.counters; + assert.equal(rows.length,3); assert.ok(rows.every(r=>r.count===1 && r.duration==="lt_10m")); +}); + +test("opt-out prevents opening bound transcript files; disable erases cycle cursors", async t => { + const path=join(await fixture(t),"usage.json"); + const ctx={env:{LOOPX_USAGE_PING_ENDPOINT:"http://127.0.0.1:1/v1/ping"},version:"1.0.0",python:"3.13",channel:"source",now:new Date(at)}; + await configure(path,ctx,"enable"); const generation=JSON.parse(await readFile(path,"utf8")).generation; + const input=cycle("start",at,{codex:{path:"/nonexistent/private.jsonl",id:"thread"}}); + await observe(path,{...ctx,env:{...ctx.env,DO_NOT_TRACK:"1"}},generation,null,async()=>{throw new Error("must not send");},undefined,input); + await assert.rejects(readFile(path+".cycles"),/ENOENT/); + await observe(path,ctx,generation,null,async()=>204,undefined,input); + assert.ok(await readFile(path+".cycles")); + await configure(path,ctx,"disable"); + await assert.rejects(readFile(path+".cycles"),/ENOENT/); +}); + +test("oversized transcript content cannot strand later terminal timing", async t => { + const path=join(await fixture(t),"large.jsonl"); + await writeFile(path,JSON.stringify({type:"session_meta",payload:{id:"thread"}})+"\n"); + let result=await readCodexTiming(path,"thread",undefined,at,key,"codex_app"); + await appendFile(path,JSON.stringify({type:"response_item",payload:{text:"x".repeat(2*1024*1024)}})+"\n"+event("task_complete",at+60000,{turn_id:"turn",started_at:new Date(at).toISOString()})); + const found=[]; + for(let i=0;i<4;i++) { + result=await readCodexTiming(path,"thread",result.cursor,at+60001,key,"codex_app"); + found.push(...result.observations); + } + assert.equal(found.length,1); assert.equal(found[0].end-found[0].start,60000); +}); + +test("unbound sequential cycles restart but exact replay keeps its original Host population", async t => { + const path=join(await fixture(t),"cycles"); + await cycleObservations(path,"g",at,cycle("start",at,{turn:null})); + await cycleObservations(path,"g",at+1000,cycle("spend",at+1000,{turn:null})); + await cycleObservations(path,"g",at+2000,cycle("start",at+2000,{turn:null})); + await cycleObservations(path,"g",at+2500,cycle("spend",at+1000,{turn:null})); + await cycleObservations(path,"g",at+2500,cycle("start",at,{turn:null})); + const next=await cycleObservations(path,"g",at+3000,cycle("spend",at+3000,{turn:null})); + assert.equal(next[0].start,at+2000); + await cycleObservations(path,"g",at,cycle("start",at)); + const replay=await cycleObservations(path,"g",at+4000,cycle("spend",at+4000,{host:"unknown"})); + assert.equal(replay[0].host,"codex_app"); +}); + + +test("scope expansion renews disclosure and fences old observations without undoing disable", async t => { + const path=join(await fixture(t),"usage.json"); + const ctx={env:{LOOPX_USAGE_PING_ENDPOINT:"http://127.0.0.1:1/v1/ping"},version:"1.0.0",python:"3.13",channel:"source",now:new Date(at)}; + await configure(path,ctx,"enable"); + const prior=JSON.parse(await readFile(path,"utf8")); prior.notice.version=2; + await writeFile(path,JSON.stringify(prior)); + assert.equal((await inspect(path,ctx)).blocked_by,"notice_required"); + await observe(path,ctx,prior.generation,null,async()=>{throw new Error("unexpected send");},undefined,cycle("start",at)); + await assert.rejects(readFile(path+".cycles"),/ENOENT/); + const enabled=await configure(path,ctx,"enable"); assert.equal(enabled.notice.version,3); + const current=JSON.parse(await readFile(path,"utf8")); assert.notEqual(current.generation,prior.generation); + await configure(path,ctx,"disable"); + assert.equal((await inspect(path,ctx)).consent,"disabled"); +}); + + +test("completed diagnostic cycles cannot exhaust capacity and suppress future work", async t => { + const path=join(await fixture(t),"cycles"); + for(let n=1;n<=129;n++) { + const identity=n.toString(16).padStart(64,"0"); + await cycleObservations(path,"g",at+n*1000,cycle("start",at+n*1000,{turn:identity})); + const result=await cycleObservations(path,"g",at+n*1000+500,cycle("spend",at+n*1000+500,{turn:identity})); + assert.equal(result[0].end-result[0].start,500); + } + assert.equal(JSON.parse(await readFile(path,"utf8")).cycles.length,128); +}); diff --git a/tests/test_usage_goal.py b/tests/test_usage_goal.py index bd33e638f1..73b5d2b2be 100644 --- a/tests/test_usage_goal.py +++ b/tests/test_usage_goal.py @@ -57,3 +57,104 @@ def set(self): assert all(row["observation"]["start"] <= row["observation"]["end"] for row in observed) assert "private-goal" not in json.dumps(observed) assert str(tmp_path) not in json.dumps([row["observation"] for row in observed]) + + +def _bound_fixture(tmp_path, monkeypatch): + import sqlite3 + home = tmp_path / "codex-home" + sessions = home / "sessions" + sessions.mkdir(parents=True) + rollout = sessions / "bound.jsonl" + rollout.write_text(json.dumps({"type": "session_meta", "payload": {"id": "thread-a"}}) + "\n") + database = home / "state_5.sqlite" + with sqlite3.connect(database) as conn: + conn.execute("CREATE TABLE threads(id TEXT, rollout_path TEXT)") + conn.execute("INSERT INTO threads VALUES (?, ?)", ("thread-a", str(rollout))) + registry = tmp_path / "registry.json" + registry.write_text(json.dumps({"schema_version": "0.1", "goals": [{"id": "goal", "coordination": { + "registered_agents": ["agent"], "thread_agent_bindings": [ + {"thread_id": "thread-a", "host_surface": "codex-app", "agent_id": "agent"}, + ], + }}]})) + monkeypatch.setenv("CODEX_HOME", str(home)) + monkeypatch.delenv("CODEX_THREAD_ID", raising=False) + return registry, rollout + + +def test_binding_discovery_uses_exact_agent_thread_and_selected_home_only(tmp_path, monkeypatch): + registry, rollout = _bound_fixture(tmp_path, monkeypatch) + assert usage_goal._bound_codex_session(registry, "goal", "agent") == ( + {"path": str(rollout), "id": "thread-a"}, "codex-app") + assert usage_goal._bound_codex_session(registry, "goal", "other") is None + monkeypatch.setenv("CODEX_THREAD_ID", "unbound-thread") + assert usage_goal._bound_codex_session(registry, "goal", "agent") is None + monkeypatch.delenv("CODEX_THREAD_ID") + doc = json.loads(registry.read_text()) + doc["goals"][0]["coordination"]["thread_agent_bindings"].append( + {"thread_id": "thread-b", "host_surface": "codex-app", "agent_id": "agent"}) + registry.write_text(json.dumps(doc)) + assert usage_goal._bound_codex_session(registry, "goal", "agent") is None + monkeypatch.setenv("CODEX_THREAD_ID", "thread-a") + assert usage_goal._bound_codex_session(registry, "goal", "agent") is not None + monkeypatch.setenv("CODEX_HOME", str(tmp_path / "other-home")) + assert usage_goal._bound_codex_session(registry, "goal", "agent") is None + + +def test_metadata_cannot_redirect_timing_read_outside_selected_home(tmp_path, monkeypatch): + import sqlite3 + registry, rollout = _bound_fixture(tmp_path, monkeypatch) + outside = tmp_path / "outside.jsonl" + outside.write_text(rollout.read_text()) + with sqlite3.connect(rollout.parent.parent / "state_5.sqlite") as conn: + conn.execute("UPDATE threads SET rollout_path = ?", (str(outside),)) + assert usage_goal._bound_codex_session(registry, "goal", "agent") is None + + +def test_real_detached_cycle_and_bound_codex_event_reach_shared_ts_aggregator(tmp_path, monkeypatch): + from datetime import datetime, timezone + from pathlib import Path + registry, rollout = _bound_fixture(tmp_path, monkeypatch) + monkeypatch.setattr(usage_ping, "DEFAULT_RUNTIME_ROOT", tmp_path) + for name in ("CI", "DO_NOT_TRACK", "LOOPX_USAGE_PING", "LOOPX_USAGE_POLICY"): + monkeypatch.delenv(name, raising=False) + monkeypatch.setenv("LOOPX_USAGE_PING_ENDPOINT", "http://127.0.0.1:1/v1/ping") + usage_ping.control("enable") + cycle_path = Path(str(usage_ping.state_path()) + ".cycles") + + def wait_for(predicate): + deadline = time.monotonic() + 8 + while time.monotonic() < deadline: + try: + if predicate(): + return + except (OSError, ValueError): + pass + time.sleep(0.03) + pytest.fail("detached timing did not reach TS owner") + + now = time.time_ns() // 1_000_000 + def publish(phase): + usage_goal.observe_quota_cycle(registry_path=registry, runtime_root=tmp_path / "runtime", goal_id="goal", + agent_id="agent", turn_id="turn-a", phase=phase, + at=time.time_ns() // 1_000_000, host="unknown") + publish("start") + wait_for(lambda: bool(json.loads(cycle_path.read_text())["cursors"])) + end = time.time_ns() // 1_000_000 + with rollout.open("a") as stream: + stream.write(json.dumps({"type": "event_msg", "timestamp": datetime.now(timezone.utc).isoformat(), "payload": { + "type": "task_complete", "turn_id": "turn-a", "started_at": datetime.fromtimestamp(now / 1000, timezone.utc).isoformat(), + "completed_at": datetime.fromtimestamp(end / 1000, timezone.utc).isoformat(), + "last_agent_message": "PRIVATE CONTENT MUST NOT LEAVE THE SESSION", + }}) + "\n") + publish("spend") + goals = Path(str(usage_ping.state_path()) + ".goals") + wait_for(lambda: len(json.loads(goals.read_text())["goals"]) == 2) + preview = usage_ping.control("status")["goal_preview"] + assert {row["measurement"] for row in preview["counters"]} == {"quota_cycle", "codex_turn"} + assert all(row["host"] == "codex_app" for row in preview["counters"]) + assert "PRIVATE CONTENT" not in cycle_path.read_text() + goals.read_text() + assert "thread-a" not in json.dumps(preview) + usage_ping.control("disable") + monkeypatch.setattr(usage_ping, "_detach", lambda *a, **k: pytest.fail("disabled observer spawned")) + publish("start") + assert not cycle_path.exists() and not goals.exists() diff --git a/tsconfig.control-plane.json b/tsconfig.control-plane.json index ad12babe1f..37cd8ab98d 100644 --- a/tsconfig.control-plane.json +++ b/tsconfig.control-plane.json @@ -16,6 +16,7 @@ "loopx/control_plane/runtime/usage_statistics*.ts", "tests/control_plane_ts/usage_statistics.test.ts", "tests/control_plane_ts/usage_statistics_goals.test.ts", + "tests/control_plane_ts/usage_statistics_sources.test.ts", "loopx/control_plane/effect_program.ts", "loopx/control_plane/effect_runtime_errors.ts", "loopx/control_plane/effect_runtime_handlers.ts", From d2626e36d8e87fedd600c9c6708089ce788e823c Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 01:14:17 +0800 Subject: [PATCH 2/4] docs(usage): explain independent duration populations and rollout Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- apps/usage-collector/README.md | 11 ++- docs/reference/usage-ping.md | 114 +++++++++++++++++------------ docs/reference/usage-ping.zh-CN.md | 82 ++++++++++++--------- 3 files changed, 124 insertions(+), 83 deletions(-) diff --git a/apps/usage-collector/README.md b/apps/usage-collector/README.md index fb26213807..dda2a01897 100644 --- a/apps/usage-collector/README.md +++ b/apps/usage-collector/README.md @@ -8,7 +8,7 @@ The TypeScript client/collector allowlist lives in |---|---| | `POST /v1/ping` | Daily random-ID heartbeat with version/OS/CPU/Python/channel; ≤1 KiB | | `POST /v1/aggregate` | Fixed CLI counts, no installation ID or join key; ≤16 KiB | -| `POST /v1/goals` | Observed Goal-day span/execution buckets, no identity; ≤16 KiB | +| `POST /v1/goals` | Independent Goal/measurement/Host-day span/duration buckets, no identity; ≤16 KiB | | `GET /v1/goal-stats` | Independent 30-day duration histograms; cells below 5 omitted | | `GET /v0/stats` | Deduplicated active/new installations, including retained v0 clients; version/OS/CPU/channel breakdown | | `GET /v1/aggregate-stats` | Independent 30-day feature/result/duration/error totals; cells below 5 omitted | @@ -49,8 +49,13 @@ npx wrangler d1 migrations apply loopx-usage --remote npx wrangler deploy ``` -Existing v1 installations also apply `0002-goal-usage.sql` before deploying -the Goal-duration Worker. Back up first; this only adds `goal_usage_counts`. +Existing v1 installations also apply `0002-goal-usage.sql` and then +`0003-goal-duration-sources.sql` before deploying the Goal-duration Worker. +Back up first. The latter adds `goal_duration_counts` and copies prior counts +as `host_call`/`unknown`, preserving the old table for rollback. +`quota_cycle`, `codex_turn` and `host_call` are overlapping populations and +must never be summed. The unreleased Goal payload requires measurement/Host +labels and uses `duration` instead of `execution`. A Worker rollback can leave that additive table intact. Qualify `/v1/ping`, `/v1/aggregate`, `/v1/goals`, all stats endpoints, and invalid-field/size diff --git a/docs/reference/usage-ping.md b/docs/reference/usage-ping.md index 158e62b98d..e688a202e5 100644 --- a/docs/reference/usage-ping.md +++ b/docs/reference/usage-ping.md @@ -145,54 +145,76 @@ collector, and LoopX continues to work without telemetry. ## Observed Goal duration -The Goal channel adds two duration histograms to basic statistics. It helps -answer whether observed work continues for hours or days, and how much Host -execution those spans contain. It does not identify a person or a Goal. - -- **Span:** first to most recent observed Host execution, including intervening - pauses. It stops growing while no execution is observed. -- **Execution:** union of observed Host-call intervals for one Goal on one - machine. Concurrent or nested calls overlap only once; retry execution counts, - settlement-only replay does not. Network/tool/approval waits inside a Host call - are included; this is neither CPU time nor billing time. -- **Sampling:** one cumulative snapshot per locally observed Goal-day, flushed - after that UTC day closes when another observation or normal usage occurs. - Unfinished Goals are included. A continuously executing Host checkpoints every - minute and can flush without another CLI command. Quiet Goals are not counted - again every day. These counts are **Goal-day observations, not unique Goals**; - the collector cannot join a Goal across days or machines. - -Coverage is managed `turn run-once` Host execution and regular owner Goal chat. -Native `/goal`, externally attached agent sessions, manager and external-audience -conversations are excluded because they do not share these timing boundaries. -File, SQLite and PostgreSQL use the same observer; no provider state is queried -or changed. Measurement starts when first observed after notice acknowledgment, -not at historical Goal creation. Disabling, changing recipient, or clearing local -state restarts measurement. There is no historical backfill or completion claim. - -All durations are lower bounds on observed work: confirmed prefixes survive a -crash; missing final checkpoints, lock contention, network failure, suspended -hosts and collection limits can lose observations. Never extrapolate a crashed -Host as still executing. Local storage holds at most 64 Goals and 512 disjoint -recent intervals per Goal; old intervals compact into totals, and 90-day inactive -Goals expire. Delayed observations older than one day are discarded; pending -snapshots older than seven days are discarded. Do not use this channel for -liveness detection, quotas, acceptance, or accounting. - -The only outgoing Goal payload is: +Goal timing answers whether observed work continues across hours or days. It +reports three **independent populations**; never add their durations or counts: + +| Measurement | Boundary and coverage | Interpretation | +| --- | --- | --- | +| `quota_cycle` | Every Host using the shared quota CLI: first allowed `should-run` to successful executed `spend-slot`, including Codex App | Coarse progression time, including intervening pauses and waits | +| `codex_turn` | Exact accepted Codex task binding, timing events in that selected local Codex home | Finer execution intervals, including tool and approval waits | +| `host_call` | Managed `turn run-once` and regular owner Goal chat around actual Host invocation | Directly instrumented call intervals, including network/tool waits | + +Repeated allowed quota reads preserve the first start. Denial, spend preview, +failed settlement and receipt repair do not finish a cycle. Successful settlement +replays emit no new timing marker and cannot extend its end. Exact Turn identity separates concurrent cycles; +without one, only a single inferred cycle per Goal/agent lane is measured. This +fallback cannot distinguish concurrent unbound cycles. Missing spend produces no +finished interval. These are diagnostic observations, never quota authority. + +Codex discovery uses accepted Goal/agent/task bindings and the selected +`CODEX_HOME` read-only metadata database. It does not search other homes or infer +ownership from cwd. Timing extraction runs in detached processes triggered by +quota observations. It reads at most 1 MiB of new JSONL per observation (the tail +on first discovery), retains only timing cursors and open Turn identity locally, +and never uploads transcript content. A terminal written after spend is picked +up on a later observation; this is not a global session watcher. Missing or +ambiguous bindings and unavailable files leave the common quota cycle working. +No historical backfill or extrapolation of crashed sessions occurs. + +Within each Goal/measurement/Host series, **span** is first to most recent +observed activity, including pauses; **duration** is the union of observed +intervals. Parallel or nested overlap counts once within that series. Both stop +growing without new evidence. They are not Goal age, CPU time, completion or +billing evidence. Fixed Host labels are `codex_app`, `codex_cli`, `claude_code`, +`dsh`, `opencode`, `trae`, `other`, `unknown`; custom names never go on the wire. + +One cumulative snapshot per observed series/UTC day is claimed after that day +closes, when another observation or normal usage occurs. Unfinished Goals count; +quiet Goals are not counted daily. Managed Host calls checkpoint every minute. +Counts are **Goal/measurement/Host-day observations**, not unique Goals or users. +The collector cannot join a Goal across days or machines. + +The observer is independent of File/SQLite/PostgreSQL state ownership and never +changes authority state. Collection starts after notice acknowledgment. Disable, +recipient changes or clearing local state restart measurement. Crashes, missing +checkpoints, contention, offline collectors and limits can lose observations; +these are partial measurements, not a complete execution accounting system. +Local limits are 64 series, 512 recent disjoint intervals per series, 128 cycles, +64 Codex cursors. Closed cycles and least-recently-read cursors yield capacity to +new work; inactive cursors expire after seven days. Intervals older than 14 days compact into totals; series expire +after 90 inactive days. Observations may arrive seven days late; coarse/fine +intervals longer than seven days are discarded. Direct Host checkpoints longer +than two minutes are discarded as unproven scheduling suspension. Unsent daily +snapshots expire after seven days. No retries require an outgoing identity. + +The outgoing Goal payload is strictly allowlisted: ```json -{"schema":"loopx_goal_usage_aggregate_v1","counters":[{"span":"lt_7d","execution":"lt_6h","count":1}]} +{"schema":"loopx_goal_usage_aggregate_v1","counters":[{"measurement":"quota_cycle","host":"codex_app","span":"lt_7d","duration":"lt_6h","count":1}]} ``` Both durations use `lt_1m`, `lt_10m`, `lt_1h`, `lt_6h`, `lt_1d`, `lt_7d`, -`lt_30d`, `gte_30d`. No Goal ID, installation ID, source path, name, event time -or free text is sent. `/v1/goals` accepts the strict payload; `/v1/goal-stats` -returns independent marginal histograms for the last 30 receipt days and omits -cells below five. The existing usage settings switch, environment opt-outs and -consent policy control all three channels. Settings and `loopx usage-ping status` -show `goal_preview`; this is a current local snapshot, not a delivery receipt. -The expanded scope requires notice version 2; previous explicit disable persists. - -Deploy collector migration `0002-goal-usage.sql` and its Worker before shipping -the client. This additive table leaves existing heartbeats and CLI counts intact. +`lt_30d`, `gte_30d`. No Goal ID, installation ID, path, name, event time or free +text is sent. `/v1/goals` ingests counts; `/v1/goal-stats` returns separate +measurement histograms over 30 receipt days, omitting cells below five. + +The existing settings switch, environment opt-outs and consent policy control +all channels and local timing reads. Settings and `loopx usage-ping status` +show `goal_preview`, a local snapshot rather than a delivery receipt. Expanded +scope requires notice version 3; an existing explicit disable persists. + +Before shipping the client, back up D1, apply `0002-goal-usage.sql` and +`0003-goal-duration-sources.sql`, then deploy the Worker. The latter migrates +previous Goal counts into `host_call`/`unknown` without deleting the old table; +existing heartbeat and CLI counts remain intact. The unreleased Goal v1 payload +now requires measurement/Host labels and `duration` in place of `execution`. diff --git a/docs/reference/usage-ping.zh-CN.md b/docs/reference/usage-ping.zh-CN.md index bcd04d32ba..34464cb573 100644 --- a/docs/reference/usage-ping.zh-CN.md +++ b/docs/reference/usage-ping.zh-CN.md @@ -118,43 +118,57 @@ observability。但网络服务商仍处理连接信息,分开数据包不保 ## Goal 执行时长观测 -Goal 统计用于判断实际观测到的工作能否跨小时、跨天持续,以及其中有多少 -Host 执行时间。它不会识别某个用户或 Goal。 - -- **执行跨度(span)**:首次到最近一次观测执行的时间,包含中间的暂停; - 没有新执行时不会继续增长,不是从历史创建日期推算的 Goal 年龄。 -- **累计执行(execution)**:同机同 Goal 的 Host 调用时间区间求并集, - 并行和嵌套重叠只算一次,真实重试计入,仅补结算的重放不计入。 - 调用中的网络、工具和审批等待也包含在内,不是 CPU 用时或计费时长。 -- **采样**:每个本机有执行观测的 Goal 每 UTC 日最多产生一份累计快照, - 次日发生观测或普通使用时发送。持续运行的 Host 每分钟提交检查点, - 无须等待下一条 CLI 命令;未完成的 Goal 也会计入。长期静默的 Goal - 不会每天重复计数。统计单位是 **Goal-day 观测,不是去重 Goal 数**。 - -覆盖受管 `turn run-once` 与普通用户 Goal 对话。原生 `/goal`、外部附加的 -Agent、管理者和外部受众对话暂不覆盖,避免将不同执行边界混作同一种统计。 -File、SQLite、PostgreSQL 共用这层旁观机制,不读取或改写 provider 状态。 -从告知确认后第一次实际观测开始计量;关闭、变更收集地址、清空本机状态会 -重新开始,不回填历史,也不把停止或命令成功解释为 Goal 完成。 - -这些时长是已观测工作的下界:崩溃只保留已确认的检查点,缺少尾部检查点、 -锁竞争、网络失败、主机挂起和容量限制都可能漏计;不会把崩溃后的时间继续 -算作执行。本机最多保留 64 个 Goal、每个 Goal 512 个近期不相交区间,旧区间 -压缩为累计数,90 天无执行的 Goal 到期删除。超过一天的迟到观测、超过七天的 -待发快照丢弃。不能用此统计判定存活、验收、配额或计费。 +用三种**独立口径**判断工作是否跨小时、跨天持续。它们会重叠,不能相加: + +| measurement | 范围 | 含义 | +| --- | --- | --- | +| `quota_cycle` | 使用通用 quota CLI 的所有 Host,包括 Codex App | 首次允许推进的 `should-run` 到成功执行 `spend-slot`,包含暂停和等待 | +| `codex_turn` | 已接受的 Codex 任务绑定,当前 Codex home 内的真实时间事件 | 更细的执行区间,也包含工具和审批等待 | +| `host_call` | 受管 `turn run-once`、普通用户 Goal 对话的实际 Host 调用 | 直接打点的调用区间,包含网络与工具等待 | + +重复读 quota 不重置起点;拒绝推进、spend 预览、失败结算和回执修复不会结束周期。 +成功结算重放不生成新的时间标记,也不延长终点。有精确 Turn 身份时按 Turn 区分;没有时每个 Goal/agent +通道只推断一个周期,无法区分同通道并发的无身份周期。缺少 spend 不虚构结束。 +它是旁观统计,不影响准入、结算或业务状态。 + +Codex 发现使用既有 Goal/agent/task 绑定和所选 `CODEX_HOME` 的只读元数据, +不跨 home 搜索,不按工作目录猜测归属。quota 观测启动脱离前台的进程,每次最多 +读取 1 MiB 新 JSONL(首次读取尾部),本地仅保留时间游标和未结束轮次身份。 +原始会话内容不上传。spend 后才写入的结束事件需要等下次观测,并非全局实时监听。 +绑定缺失、歧义或文件不可用不会阻断通用周期统计;不回填历史,不外推崩溃后的时间。 + +每个 Goal/口径/Host 分别计算 **span**(首次至最近观测活动,包含中间暂停)和 +**duration**(已观测区间的并集)。同口径同 Host 内并行重叠只计一次;没有新证据 +就不增长。它们不是 Goal 年龄、CPU 用时、完成证据或计费时长。Host 只允许固定枚举 +`codex_app`、`codex_cli`、`claude_code`、`dsh`、`opencode`、`trae`、`other`、 +`unknown`,不发送自定义名称。 + +每个有活动的统计序列每 UTC 日最多一份累计快照,次日发生观测或普通使用时发送。 +未完成 Goal 也计入;静默 Goal 不每天重复计数。直接 Host 打点每分钟提交检查点。 +单位是 **Goal/口径/Host-day 观测**,不是去重 Goal 或用户数,接收端不能跨天、 +跨机器关联同一个 Goal。 + +这层统计独立于 File/SQLite/PostgreSQL,不改写权威状态。告知确认后才开始观测; +关闭、更换收集地址或清空本机状态会重新开始。崩溃、缺少检查点、锁竞争、网络失败 +和容量限制均可能漏计。最多保留 64 条序列、每条 512 个近期不相交区间、128 个周期、 +64 个 Codex 游标。已结束周期及最久未读的游标会让位给新工作;七天未读的游标过期。14 天前的区间压缩为累计数,90 天无活动的序列过期。 +接纳七天内迟到观测;超过七天的周期/轮次区间、超过两分钟的直接 Host 检查点丢弃。 +待发快照七天过期,不使用出站身份做重试去重。 出站样例: ```json -{"schema":"loopx_goal_usage_aggregate_v1","counters":[{"span":"lt_7d","execution":"lt_6h","count":1}]} +{"schema":"loopx_goal_usage_aggregate_v1","counters":[{"measurement":"quota_cycle","host":"codex_app","span":"lt_7d","duration":"lt_6h","count":1}]} ``` -两个时长均分为小于 1 分钟、10 分钟、1 小时、6 小时、1 天、7 天、30 天及 -不少于 30 天八档。不发送 Goal ID、安装 ID、名称、路径、事件时间或自由文本。 -`/v1/goals` 接收严格白名单数据;`/v1/goal-stats` 提供最近 30 个接收日的独立 -分布,少于 5 的单元不公开。设置页与 `loopx usage-ping status` 的 `goal_preview` -可查看本机当前快照;它不是发送回执。 - -沿用同一统计开关、环境变量和同意策略。扩大统计范围需要重新展示第 2 版告知, -不会重置原有的关闭选择。发布客户端前先备份 D1,应用 `0002-goal-usage.sql` 并 -部署 Worker;新增表不会删除既有心跳或 CLI 计数。 +两种时长均分为小于 1 分钟、10 分钟、1 小时、6 小时、1 天、7 天、30 天及不少于 +30 天八档。不发送 Goal ID、安装 ID、名称、路径、事件时间或自由文本。 +`/v1/goals` 接收白名单数据;`/v1/goal-stats` 按口径分别给出最近 30 个接收日的 +分布,少于 5 的单元不公开。设置与 `loopx usage-ping status` 的 `goal_preview` +是本机当前快照,不是发送回执。 + +沿用统一开关、环境变量和同意策略,关闭也停止本地时间读取。扩大范围需要第 3 版 +告知,保留原有关闭选择。发布客户端前先备份 D1,依次应用 `0002-goal-usage.sql`、 +`0003-goal-duration-sources.sql` 并部署 Worker。后者将旧计数转入 `host_call` / +`unknown`,保留旧表以便回滚,不影响心跳和 CLI 计数。尚未发布的 Goal v1 协议 +现在要求口径和 Host 字段,并以 `duration` 替代 `execution`。 From 8494413f03824596d200b72140ee2fe8bc183c95 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 01:22:33 +0800 Subject: [PATCH 3/4] refactor(quota): isolate observation and response projection from dispatch Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/quota.py | 66 +++++++++++++++++-------------------- loopx/usage_goal.py | 22 ++++++++++++- tests/test_turn_envelope.py | 5 +-- 3 files changed, 54 insertions(+), 39 deletions(-) diff --git a/loopx/cli_commands/quota.py b/loopx/cli_commands/quota.py index 746bbd8ead..2ba1351da7 100644 --- a/loopx/cli_commands/quota.py +++ b/loopx/cli_commands/quota.py @@ -1,9 +1,12 @@ from __future__ import annotations import argparse +import time from collections.abc import Callable, Mapping from pathlib import Path +from ..usage_goal import observe_quota_result + from ..capabilities.explore.composition_frontier import ( project_live_explore_composition_frontier, ) @@ -301,17 +304,33 @@ def _attach_turn_start_hook_dispatch( -def _render_turn_envelope_payload( +def _project_quota_cli_payload( payload: dict[str, object], + args: argparse.Namespace, + detail_sections: frozenset[str], scheduler_context: object, ) -> dict[str, object]: - """Render the Turn envelope, degrading to the typed payload on rejection. + """Project already-decided facts; preserve typed failures on envelope rejection. The envelope is an additive hot-path view over a decided payload. A typed validation/failure payload has no interaction contract to project, so a renderer rejection keeps the typed diagnostic itself (with the skip reason) instead of masking it with a crash (issue #3687). """ + if not bool(getattr(args, "turn_envelope", False)): + if args.quota_command == "should-run": + return compact_quota_should_run_cli_payload( + payload, + include_todo_summary_detail="agent-todos" in detail_sections, + include_user_todo_summary_detail="user-todos" in detail_sections, + include_goal_boundary_detail="goal-boundary" in detail_sections, + include_vision_detail="vision" in detail_sections, + ) + if args.quota_command == "monitor-poll": + return compact_quota_monitor_poll_cli_payload( + payload, include_decision_detail="decisions" in detail_sections, + ) + return payload try: return build_turn_envelope( payload, @@ -330,7 +349,6 @@ def handle_quota_command( print_payload: PrintPayload, append_cli_rollout_event: RolloutEventAppender, ) -> int: - import time usage_quota_started = time.time_ns() // 1_000_000 heartbeat_turn_id: str | None = None heartbeat_receipt_existing: dict[str, object] | None = None @@ -757,40 +775,16 @@ def handle_quota_command( goal_id=args.goal_id, agent_id=args.agent_id, ) - if context is not None and payload.get("ok"): - phase = ( - "start" if args.quota_command == "should-run" and payload.get("should_run") is True - else "spend" if args.quota_command == "spend-slot" and bool(args.execute) - and payload.get("appended") else None - ) - if phase: - from ..usage_goal import observe_quota_cycle - observe_quota_cycle( - registry_path=registry_path, runtime_root=context.runtime_root, - goal_id=args.goal_id, agent_id=args.agent_id, - turn_id=_effective_spend_turn_instance_id(payload, heartbeat_turn_id=heartbeat_turn_id), - phase=phase, at=usage_quota_started if phase == "start" else time.time_ns() // 1_000_000, - host=str(getattr(args, "host_surface", None) or getattr(args, "runtime_profile", None) - or ("codex_app" if getattr(args, "codex_app", False) else "unknown")), - ) - if bool(getattr(args, "turn_envelope", False)): - payload = _render_turn_envelope_payload( - payload, - context.scheduler_context if context is not None else None, - ) - elif args.quota_command == "should-run": - payload = compact_quota_should_run_cli_payload( - payload, - include_todo_summary_detail="agent-todos" in detail_sections, - include_user_todo_summary_detail="user-todos" in detail_sections, - include_goal_boundary_detail="goal-boundary" in detail_sections, - include_vision_detail="vision" in detail_sections, - ) - elif args.quota_command == "monitor-poll": - payload = compact_quota_monitor_poll_cli_payload( - payload, - include_decision_detail="decisions" in detail_sections, + if context is not None: + observe_quota_result( + args, payload, registry_path=registry_path, runtime_root=context.runtime_root, + turn_id=_effective_spend_turn_instance_id(payload, heartbeat_turn_id=heartbeat_turn_id), + started_at=usage_quota_started, ) + payload = _project_quota_cli_payload( + payload, args, detail_sections, + context.scheduler_context if context is not None else None, + ) if args.quota_command == "should-run" and context is not None: attach_host_poll_receipt( context.status_payload, diff --git a/loopx/usage_goal.py b/loopx/usage_goal.py index 7fe40f31b4..cc680576db 100644 --- a/loopx/usage_goal.py +++ b/loopx/usage_goal.py @@ -5,6 +5,7 @@ """ from __future__ import annotations +import argparse from contextlib import closing, contextmanager import hashlib import json @@ -12,7 +13,7 @@ from pathlib import Path import threading import time -from collections.abc import Iterator +from collections.abc import Iterator, Mapping from . import usage_ping @@ -76,6 +77,25 @@ def periodically() -> None: publish() +def observe_quota_result(args: argparse.Namespace, payload: Mapping[str, object], *, + registry_path: Path, runtime_root: Path, + turn_id: str | None, started_at: int) -> None: + """Translate already-decided CLI facts; never decide admission or settlement.""" + if not payload.get("ok"): + return + if args.quota_command == "should-run" and payload.get("should_run") is True: + phase, at = "start", started_at + elif args.quota_command == "spend-slot" and args.execute and payload.get("appended"): + phase, at = "spend", time.time_ns() // 1_000_000 + else: + return # Preview, failure and replay are not fresh execution evidence. + host = str(getattr(args, "host_surface", None) or getattr(args, "runtime_profile", None) + or ("codex_app" if getattr(args, "codex_app", False) else "unknown")) + observe_quota_cycle(registry_path=registry_path, runtime_root=runtime_root, + goal_id=args.goal_id, agent_id=args.agent_id, turn_id=turn_id, + phase=phase, at=at, host=host) + + def observe_quota_cycle(*, registry_path: Path, runtime_root: Path, goal_id: str, agent_id: str | None, turn_id: str | None, phase: str, at: int, host: str) -> None: diff --git a/tests/test_turn_envelope.py b/tests/test_turn_envelope.py index 220df5c9f0..0a30f37f0f 100644 --- a/tests/test_turn_envelope.py +++ b/tests/test_turn_envelope.py @@ -1047,7 +1047,8 @@ def test_protocol_packet_compatibility_does_not_bypass_host_signature_check( @pytest.mark.parametrize("has_packet", [False, True]) def test_envelope_fallback_preserves_typed_failure_with_or_without_packet(has_packet: bool) -> None: - from loopx.cli_commands.quota import _render_turn_envelope_payload + from loopx.cli_commands.quota import _project_quota_cli_payload + from argparse import Namespace failure: dict[str, Any] = { "ok": False, "decision": "skip", "should_run": False, @@ -1058,7 +1059,7 @@ def test_envelope_fallback_preserves_typed_failure_with_or_without_packet(has_pa "schema_version": "protocol_action_packet_v0", "summary": HISTORICAL_V0_SUMMARY, } before = deepcopy(failure) - rendered = _render_turn_envelope_payload(failure, None) + rendered = _project_quota_cli_payload(failure, Namespace(turn_envelope=True), frozenset(), None) assert "interaction_contract must be an object" in rendered.pop("turn_envelope_skipped") assert rendered == before assert failure == before From 9796fbdb326d22a477b66b520265fcd056019d2a Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 02:54:36 +0800 Subject: [PATCH 4/4] fix(usage): frame the bound Codex header on LF instead of one 64 KiB read The Codex reader took a fixed 65,536-byte chunk, split it at the first newline and parsed the result. A `session_meta` line is not a short id record -- Codex's recorder writes `base_instructions` and the dynamic tool list into it, and observed first lines already reach ~50 KiB -- so a legitimate header above that chunk truncated the JSON, threw before any cursor was built, and left the session with `quota_cycle` but never `codex_turn`. Retrying re-read the same truncated head, so the loss was permanent. Read the opening record by LF framing under its own bounded budget (`HEADER_BUDGET`, 2 MiB) in `HEADER_CHUNK` reads, parse only what the identity check needs, and give the two non-bindable outcomes a name: an unfinished header returns without advancing the cursor so the next scan reads it whole, and a header past the budget fails as `usage_session_header_too_large` instead of as an identity mismatch. Validation: the sources test adds a 72 KiB header case (binds, keeps the identity mismatch rejection and recovers timing), an unfinished-header case that recovers on the next scan, and an over-budget case that names the reason; `tests/test_usage_goal.py` drives the same long header through the real quota observer -> detached TS chain and asserts both measurements arrive. Reverting to the fixed chunk fails all four. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../runtime/usage_statistics_codex.ts | 54 +++++++++++++++++-- .../usage_statistics_sources.test.ts | 41 ++++++++++++++ tests/test_usage_goal.py | 31 ++++++++++- 3 files changed, 120 insertions(+), 6 deletions(-) diff --git a/loopx/control_plane/runtime/usage_statistics_codex.ts b/loopx/control_plane/runtime/usage_statistics_codex.ts index 093b8eeea1..274a2668a7 100644 --- a/loopx/control_plane/runtime/usage_statistics_codex.ts +++ b/loopx/control_plane/runtime/usage_statistics_codex.ts @@ -1,20 +1,64 @@ /** Read only timing envelopes from one explicitly bound session, in bounded chunks. */ import { open } from "node:fs/promises"; +import type { FileHandle } from "node:fs/promises"; import type { GoalObservation, Host } from "./usage_statistics_goal_contract.ts"; import { object } from "./usage_statistics_contract.ts"; export type CodexCursor = { offset: number; inode: string; since: number; seen: number; skipping?: boolean; open?: { id: string; start: number; confirmed: number } }; const BUDGET = 1024 * 1024; +// The first line is not a short id record: Codex's recorder writes +// `base_instructions` and the dynamic tool list into `session_meta`, and the +// first lines observed in this project's Codex homes reach ~50 KiB. Frame the +// line by LF under its own bounded budget instead of assuming one fixed 64 KiB +// chunk holds it, which truncated the JSON and stranded the whole session. +const HEADER_CHUNK = 65536; +const HEADER_BUDGET = 2 * 1024 * 1024; +type HeaderLine = { line: string } | { error: "incomplete" | "too_large" }; function timestamp(value: unknown): number { return typeof value === "string" ? Date.parse(value) : NaN; } +/** + * Read the opening record whole, framed by LF, up to a fixed budget. + * + * `incomplete` means the writer has not finished the record yet and the caller + * should read it again later; `too_large` means the identity cannot be verified + * at all, which is named rather than reported as a mismatch. + */ +async function readHeaderLine(file: FileHandle, size: number): Promise { + const parts: Buffer[] = []; + let offset = 0; + while (offset < size && offset < HEADER_BUDGET) { + const take = Math.min(HEADER_CHUNK, size - offset, HEADER_BUDGET - offset); + const buffer = Buffer.alloc(take); + const read = await file.read(buffer, 0, take, offset); + if (read.bytesRead <= 0) break; + const chunk = buffer.subarray(0, read.bytesRead); + offset += read.bytesRead; + const newline = chunk.indexOf(10); + parts.push(newline < 0 ? chunk : chunk.subarray(0, newline)); + if (newline >= 0) return { line: Buffer.concat(parts).toString("utf8") }; + } + return offset >= HEADER_BUDGET ? { error: "too_large" } : { error: "incomplete" }; +} export async function readCodexTiming(path: string, thread: string, previous: CodexCursor | undefined, now: number, key: string, host: Host) { const file = await open(path, "r"); try { const stat = await file.stat(); - const header = Buffer.alloc(65536); - const head = await file.read(header, 0, header.length, 0); - const first = header.subarray(0, head.bytesRead).toString().split("\n")[0]; - const meta = JSON.parse(first); - if (meta.type !== "session_meta" || (meta.payload?.id ?? meta.payload?.session_id) !== thread) throw new Error("usage_session_identity_mismatch"); const inode = `${stat.dev}:${stat.ino}`; + const header = await readHeaderLine(file, stat.size); + if ("error" in header) { + if (header.error === "incomplete") { + // Nothing may be emitted before the identity is verified, and nothing is + // lost: the record is simply not written yet, so read it again later. + const cursor: CodexCursor = previous?.inode === inode && previous.offset <= stat.size + ? { ...structuredClone(previous), seen: now } + : { offset: 0, inode, since: now, seen: now }; + return { cursor, observations: [] as GoalObservation[] }; + } + throw new Error(`usage_session_header_too_large:${stat.size}`); + } + let meta: unknown; + try { meta = JSON.parse(header.line); } catch { throw new Error("usage_session_header_unparsable"); } + const record = object(meta) ? meta : {}; + const payload = object(record.payload) ? record.payload : {}; + if (record.type !== "session_meta" || (payload.id ?? payload.session_id) !== thread) throw new Error("usage_session_identity_mismatch"); const reusable = previous?.inode === inode && previous.offset <= stat.size; const cursor: CodexCursor = reusable ? structuredClone(previous!) : { offset: Math.max(0, stat.size - BUDGET), inode, since: now, seen: now }; cursor.seen = now; diff --git a/tests/control_plane_ts/usage_statistics_sources.test.ts b/tests/control_plane_ts/usage_statistics_sources.test.ts index b06e6e3efe..db847705bc 100644 --- a/tests/control_plane_ts/usage_statistics_sources.test.ts +++ b/tests/control_plane_ts/usage_statistics_sources.test.ts @@ -85,6 +85,47 @@ test("opt-out prevents opening bound transcript files; disable erases cycle curs await assert.rejects(readFile(path+".cycles"),/ENOENT/); }); +test("a session_meta header larger than one chunk still binds and recovers timing", async t => { + const path=join(await fixture(t),"long-header.jsonl"); + // Codex's recorder writes base_instructions and the tool list into this line; + // 72 KiB already exceeds the old fixed 64 KiB read and is well inside the + // framing budget, so it must bind exactly like a short header. + const header=JSON.stringify({type:"session_meta",payload:{id:"thread",base_instructions:{text:"i".repeat(72000)},dynamic_tools:[{name:"shell"}]}}); + assert.ok(header.length>65536,"fixture header must exceed the fixed chunk this reader used"); + await writeFile(path,header+"\n"+event("task_started",at-1000,{turn_id:"first",started_at:new Date(at-1000).toISOString()})); + let result=await readCodexTiming(path,"thread",undefined,at,key,"codex_app"); + assert.equal(result.observations.length,0); + await assert.rejects(readCodexTiming(path,"other-thread",undefined,at,key,"codex_app"),/identity_mismatch/); + await appendFile(path,event("task_complete",at+60000,{turn_id:"first",started_at:new Date(at-1000).toISOString(),completed_at:new Date(at+60000).toISOString()})); + result=await readCodexTiming(path,"thread",result.cursor,at+60001,key,"codex_app"); + assert.equal(result.observations.length,1); + assert.equal(result.observations[0].end-result.observations[0].start,60000); + assert.equal(result.observations[0].measurement,"codex_turn"); +}); + +test("an unfinished header waits for the rest of the line instead of failing the session", async t => { + const path=join(await fixture(t),"partial-header.jsonl"); + const header=JSON.stringify({type:"session_meta",payload:{id:"thread",base_instructions:{text:"j".repeat(72000)}}}); + await writeFile(path,header.slice(0,70000)); + const waiting=await readCodexTiming(path,"thread",undefined,at,key,"codex_app"); + assert.equal(waiting.observations.length,0); + assert.equal(waiting.cursor.offset,0,"an incomplete header must not advance the cursor past unread bytes"); + await writeFile(path,header+"\n"+event("task_started",at-1000,{turn_id:"first",started_at:new Date(at-1000).toISOString()})); + const started=await readCodexTiming(path,"thread",waiting.cursor,at,key,"codex_app"); + await appendFile(path,event("task_complete",at+30000,{turn_id:"first",completed_at:new Date(at+30000).toISOString()})); + const done=await readCodexTiming(path,"thread",started.cursor,at+30001,key,"codex_app"); + assert.equal(done.observations.length,1); + // A start observed before the first scan is clamped to the cursor's `since`, + // so the span reaches back only to the first read that saw this session. + assert.equal(done.observations[0].end-done.observations[0].start,30000); +}); + +test("a header above the framing budget is named instead of reported as a mismatch", async t => { + const path=join(await fixture(t),"oversized-header.jsonl"); + await writeFile(path,JSON.stringify({type:"session_meta",payload:{id:"thread",base_instructions:{text:"k".repeat(2*1024*1024+16)}}})+"\n"); + await assert.rejects(readCodexTiming(path,"thread",undefined,at,key,"codex_app"),/usage_session_header_too_large/); +}); + test("oversized transcript content cannot strand later terminal timing", async t => { const path=join(await fixture(t),"large.jsonl"); await writeFile(path,JSON.stringify({type:"session_meta",payload:{id:"thread"}})+"\n"); diff --git a/tests/test_usage_goal.py b/tests/test_usage_goal.py index 73b5d2b2be..110d5ad371 100644 --- a/tests/test_usage_goal.py +++ b/tests/test_usage_goal.py @@ -110,10 +110,19 @@ def test_metadata_cannot_redirect_timing_read_outside_selected_home(tmp_path, mo assert usage_goal._bound_codex_session(registry, "goal", "agent") is None -def test_real_detached_cycle_and_bound_codex_event_reach_shared_ts_aggregator(tmp_path, monkeypatch): +def _bound_cycle_through_detached_ts(tmp_path, monkeypatch, *, header_characters: int = 0): + """Drive the real quota observer -> detached TS chain over one bound session. + + `header_characters` widens the session_meta line the way Codex's recorder + does with base_instructions, so the same entry point covers a header that + does not fit in one fixed read. + """ from datetime import datetime, timezone from pathlib import Path registry, rollout = _bound_fixture(tmp_path, monkeypatch) + if header_characters: + rollout.write_text(json.dumps({"type": "session_meta", "payload": { + "id": "thread-a", "base_instructions": {"text": "i" * header_characters}}}) + "\n") monkeypatch.setattr(usage_ping, "DEFAULT_RUNTIME_ROOT", tmp_path) for name in ("CI", "DO_NOT_TRACK", "LOOPX_USAGE_PING", "LOOPX_USAGE_POLICY"): monkeypatch.delenv(name, raising=False) @@ -150,6 +159,11 @@ def publish(phase): goals = Path(str(usage_ping.state_path()) + ".goals") wait_for(lambda: len(json.loads(goals.read_text())["goals"]) == 2) preview = usage_ping.control("status")["goal_preview"] + return preview, cycle_path, goals, publish + + +def test_real_detached_cycle_and_bound_codex_event_reach_shared_ts_aggregator(tmp_path, monkeypatch): + preview, cycle_path, goals, publish = _bound_cycle_through_detached_ts(tmp_path, monkeypatch) assert {row["measurement"] for row in preview["counters"]} == {"quota_cycle", "codex_turn"} assert all(row["host"] == "codex_app" for row in preview["counters"]) assert "PRIVATE CONTENT" not in cycle_path.read_text() + goals.read_text() @@ -158,3 +172,18 @@ def publish(phase): monkeypatch.setattr(usage_ping, "_detach", lambda *a, **k: pytest.fail("disabled observer spawned")) publish("start") assert not cycle_path.exists() and not goals.exists() + + +def test_bound_session_with_a_long_metadata_header_still_reaches_ts(tmp_path, monkeypatch): + """Codex records base_instructions in session_meta; a long header must bind.""" + characters = 72_000 + header = json.dumps({"type": "session_meta", "payload": { + "id": "thread-a", "base_instructions": {"text": "i" * characters}}}) + "\n" + assert len(header) > 65536, "fixture header must exceed one fixed read" + preview, _, _, publish = _bound_cycle_through_detached_ts( + tmp_path, monkeypatch, header_characters=characters) + assert {row["measurement"] for row in preview["counters"]} == {"quota_cycle", "codex_turn"} + assert all(row["host"] == "codex_app" for row in preview["counters"]) + usage_ping.control("disable") + monkeypatch.setattr(usage_ping, "_detach", lambda *a, **k: pytest.fail("disabled observer spawned")) + publish("start")