Skip to content

Commit 9796fbd

Browse files
committed
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>
1 parent 8494413 commit 9796fbd

3 files changed

Lines changed: 120 additions & 6 deletions

File tree

‎loopx/control_plane/runtime/usage_statistics_codex.ts‎

Lines changed: 49 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,20 +1,64 @@
11
/** Read only timing envelopes from one explicitly bound session, in bounded chunks. */
22
import { open } from "node:fs/promises";
3+
import type { FileHandle } from "node:fs/promises";
34
import type { GoalObservation, Host } from "./usage_statistics_goal_contract.ts";
45
import { object } from "./usage_statistics_contract.ts";
56
export type CodexCursor = { offset: number; inode: string; since: number; seen: number; skipping?: boolean; open?: { id: string; start: number; confirmed: number } };
67
const BUDGET = 1024 * 1024;
8+
// The first line is not a short id record: Codex's recorder writes
9+
// `base_instructions` and the dynamic tool list into `session_meta`, and the
10+
// first lines observed in this project's Codex homes reach ~50 KiB. Frame the
11+
// line by LF under its own bounded budget instead of assuming one fixed 64 KiB
12+
// chunk holds it, which truncated the JSON and stranded the whole session.
13+
const HEADER_CHUNK = 65536;
14+
const HEADER_BUDGET = 2 * 1024 * 1024;
15+
type HeaderLine = { line: string } | { error: "incomplete" | "too_large" };
716
function timestamp(value: unknown): number { return typeof value === "string" ? Date.parse(value) : NaN; }
17+
/**
18+
* Read the opening record whole, framed by LF, up to a fixed budget.
19+
*
20+
* `incomplete` means the writer has not finished the record yet and the caller
21+
* should read it again later; `too_large` means the identity cannot be verified
22+
* at all, which is named rather than reported as a mismatch.
23+
*/
24+
async function readHeaderLine(file: FileHandle, size: number): Promise<HeaderLine> {
25+
const parts: Buffer[] = [];
26+
let offset = 0;
27+
while (offset < size && offset < HEADER_BUDGET) {
28+
const take = Math.min(HEADER_CHUNK, size - offset, HEADER_BUDGET - offset);
29+
const buffer = Buffer.alloc(take);
30+
const read = await file.read(buffer, 0, take, offset);
31+
if (read.bytesRead <= 0) break;
32+
const chunk = buffer.subarray(0, read.bytesRead);
33+
offset += read.bytesRead;
34+
const newline = chunk.indexOf(10);
35+
parts.push(newline < 0 ? chunk : chunk.subarray(0, newline));
36+
if (newline >= 0) return { line: Buffer.concat(parts).toString("utf8") };
37+
}
38+
return offset >= HEADER_BUDGET ? { error: "too_large" } : { error: "incomplete" };
39+
}
840
export async function readCodexTiming(path: string, thread: string, previous: CodexCursor | undefined, now: number, key: string, host: Host) {
941
const file = await open(path, "r");
1042
try {
1143
const stat = await file.stat();
12-
const header = Buffer.alloc(65536);
13-
const head = await file.read(header, 0, header.length, 0);
14-
const first = header.subarray(0, head.bytesRead).toString().split("\n")[0];
15-
const meta = JSON.parse(first);
16-
if (meta.type !== "session_meta" || (meta.payload?.id ?? meta.payload?.session_id) !== thread) throw new Error("usage_session_identity_mismatch");
1744
const inode = `${stat.dev}:${stat.ino}`;
45+
const header = await readHeaderLine(file, stat.size);
46+
if ("error" in header) {
47+
if (header.error === "incomplete") {
48+
// Nothing may be emitted before the identity is verified, and nothing is
49+
// lost: the record is simply not written yet, so read it again later.
50+
const cursor: CodexCursor = previous?.inode === inode && previous.offset <= stat.size
51+
? { ...structuredClone(previous), seen: now }
52+
: { offset: 0, inode, since: now, seen: now };
53+
return { cursor, observations: [] as GoalObservation[] };
54+
}
55+
throw new Error(`usage_session_header_too_large:${stat.size}`);
56+
}
57+
let meta: unknown;
58+
try { meta = JSON.parse(header.line); } catch { throw new Error("usage_session_header_unparsable"); }
59+
const record = object(meta) ? meta : {};
60+
const payload = object(record.payload) ? record.payload : {};
61+
if (record.type !== "session_meta" || (payload.id ?? payload.session_id) !== thread) throw new Error("usage_session_identity_mismatch");
1862
const reusable = previous?.inode === inode && previous.offset <= stat.size;
1963
const cursor: CodexCursor = reusable ? structuredClone(previous!) : { offset: Math.max(0, stat.size - BUDGET), inode, since: now, seen: now };
2064
cursor.seen = now;

‎tests/control_plane_ts/usage_statistics_sources.test.ts‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,47 @@ test("opt-out prevents opening bound transcript files; disable erases cycle curs
8585
await assert.rejects(readFile(path+".cycles"),/ENOENT/);
8686
});
8787

88+
test("a session_meta header larger than one chunk still binds and recovers timing", async t => {
89+
const path=join(await fixture(t),"long-header.jsonl");
90+
// Codex's recorder writes base_instructions and the tool list into this line;
91+
// 72 KiB already exceeds the old fixed 64 KiB read and is well inside the
92+
// framing budget, so it must bind exactly like a short header.
93+
const header=JSON.stringify({type:"session_meta",payload:{id:"thread",base_instructions:{text:"i".repeat(72000)},dynamic_tools:[{name:"shell"}]}});
94+
assert.ok(header.length>65536,"fixture header must exceed the fixed chunk this reader used");
95+
await writeFile(path,header+"\n"+event("task_started",at-1000,{turn_id:"first",started_at:new Date(at-1000).toISOString()}));
96+
let result=await readCodexTiming(path,"thread",undefined,at,key,"codex_app");
97+
assert.equal(result.observations.length,0);
98+
await assert.rejects(readCodexTiming(path,"other-thread",undefined,at,key,"codex_app"),/identity_mismatch/);
99+
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()}));
100+
result=await readCodexTiming(path,"thread",result.cursor,at+60001,key,"codex_app");
101+
assert.equal(result.observations.length,1);
102+
assert.equal(result.observations[0].end-result.observations[0].start,60000);
103+
assert.equal(result.observations[0].measurement,"codex_turn");
104+
});
105+
106+
test("an unfinished header waits for the rest of the line instead of failing the session", async t => {
107+
const path=join(await fixture(t),"partial-header.jsonl");
108+
const header=JSON.stringify({type:"session_meta",payload:{id:"thread",base_instructions:{text:"j".repeat(72000)}}});
109+
await writeFile(path,header.slice(0,70000));
110+
const waiting=await readCodexTiming(path,"thread",undefined,at,key,"codex_app");
111+
assert.equal(waiting.observations.length,0);
112+
assert.equal(waiting.cursor.offset,0,"an incomplete header must not advance the cursor past unread bytes");
113+
await writeFile(path,header+"\n"+event("task_started",at-1000,{turn_id:"first",started_at:new Date(at-1000).toISOString()}));
114+
const started=await readCodexTiming(path,"thread",waiting.cursor,at,key,"codex_app");
115+
await appendFile(path,event("task_complete",at+30000,{turn_id:"first",completed_at:new Date(at+30000).toISOString()}));
116+
const done=await readCodexTiming(path,"thread",started.cursor,at+30001,key,"codex_app");
117+
assert.equal(done.observations.length,1);
118+
// A start observed before the first scan is clamped to the cursor's `since`,
119+
// so the span reaches back only to the first read that saw this session.
120+
assert.equal(done.observations[0].end-done.observations[0].start,30000);
121+
});
122+
123+
test("a header above the framing budget is named instead of reported as a mismatch", async t => {
124+
const path=join(await fixture(t),"oversized-header.jsonl");
125+
await writeFile(path,JSON.stringify({type:"session_meta",payload:{id:"thread",base_instructions:{text:"k".repeat(2*1024*1024+16)}}})+"\n");
126+
await assert.rejects(readCodexTiming(path,"thread",undefined,at,key,"codex_app"),/usage_session_header_too_large/);
127+
});
128+
88129
test("oversized transcript content cannot strand later terminal timing", async t => {
89130
const path=join(await fixture(t),"large.jsonl");
90131
await writeFile(path,JSON.stringify({type:"session_meta",payload:{id:"thread"}})+"\n");

‎tests/test_usage_goal.py‎

Lines changed: 30 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -110,10 +110,19 @@ def test_metadata_cannot_redirect_timing_read_outside_selected_home(tmp_path, mo
110110
assert usage_goal._bound_codex_session(registry, "goal", "agent") is None
111111

112112

113-
def test_real_detached_cycle_and_bound_codex_event_reach_shared_ts_aggregator(tmp_path, monkeypatch):
113+
def _bound_cycle_through_detached_ts(tmp_path, monkeypatch, *, header_characters: int = 0):
114+
"""Drive the real quota observer -> detached TS chain over one bound session.
115+
116+
`header_characters` widens the session_meta line the way Codex's recorder
117+
does with base_instructions, so the same entry point covers a header that
118+
does not fit in one fixed read.
119+
"""
114120
from datetime import datetime, timezone
115121
from pathlib import Path
116122
registry, rollout = _bound_fixture(tmp_path, monkeypatch)
123+
if header_characters:
124+
rollout.write_text(json.dumps({"type": "session_meta", "payload": {
125+
"id": "thread-a", "base_instructions": {"text": "i" * header_characters}}}) + "\n")
117126
monkeypatch.setattr(usage_ping, "DEFAULT_RUNTIME_ROOT", tmp_path)
118127
for name in ("CI", "DO_NOT_TRACK", "LOOPX_USAGE_PING", "LOOPX_USAGE_POLICY"):
119128
monkeypatch.delenv(name, raising=False)
@@ -150,6 +159,11 @@ def publish(phase):
150159
goals = Path(str(usage_ping.state_path()) + ".goals")
151160
wait_for(lambda: len(json.loads(goals.read_text())["goals"]) == 2)
152161
preview = usage_ping.control("status")["goal_preview"]
162+
return preview, cycle_path, goals, publish
163+
164+
165+
def test_real_detached_cycle_and_bound_codex_event_reach_shared_ts_aggregator(tmp_path, monkeypatch):
166+
preview, cycle_path, goals, publish = _bound_cycle_through_detached_ts(tmp_path, monkeypatch)
153167
assert {row["measurement"] for row in preview["counters"]} == {"quota_cycle", "codex_turn"}
154168
assert all(row["host"] == "codex_app" for row in preview["counters"])
155169
assert "PRIVATE CONTENT" not in cycle_path.read_text() + goals.read_text()
@@ -158,3 +172,18 @@ def publish(phase):
158172
monkeypatch.setattr(usage_ping, "_detach", lambda *a, **k: pytest.fail("disabled observer spawned"))
159173
publish("start")
160174
assert not cycle_path.exists() and not goals.exists()
175+
176+
177+
def test_bound_session_with_a_long_metadata_header_still_reaches_ts(tmp_path, monkeypatch):
178+
"""Codex records base_instructions in session_meta; a long header must bind."""
179+
characters = 72_000
180+
header = json.dumps({"type": "session_meta", "payload": {
181+
"id": "thread-a", "base_instructions": {"text": "i" * characters}}}) + "\n"
182+
assert len(header) > 65536, "fixture header must exceed one fixed read"
183+
preview, _, _, publish = _bound_cycle_through_detached_ts(
184+
tmp_path, monkeypatch, header_characters=characters)
185+
assert {row["measurement"] for row in preview["counters"]} == {"quota_cycle", "codex_turn"}
186+
assert all(row["host"] == "codex_app" for row in preview["counters"])
187+
usage_ping.control("disable")
188+
monkeypatch.setattr(usage_ping, "_detach", lambda *a, **k: pytest.fail("disabled observer spawned"))
189+
publish("start")

0 commit comments

Comments
 (0)