Skip to content

Commit b575203

Browse files
committed
fix(cli): match forked Codex replay against the parent usage stream
A forked rollout opens by replaying its parent's token history. The same-second heuristic only held when the whole burst landed in one second, so long replays and nested forks counted the parent's usage again — cached input especially, inflating totals several-fold. Codex copies the token counts verbatim but rewrites the timestamps, so match the leading usage against the parent's own stream instead, bounded at the fork instant so usage the parent recorded afterwards cannot mask the child's own events. The parent is located by session id, which rollout filenames embed, so the index costs a directory listing and is built only once a fork is actually seen. Its stream is read unfiltered: a nested fork replays the history its parent had itself copied, and the prefix only lines up against the whole thing. The same-second heuristic stays as the fallback for when the parent log is unavailable or the copied history was compacted. Events the UUIDv7 creation anchor already dropped still consume their slot in the prefix, or the rollout's first real turn would be compared against the head of it and could be deleted as replay. Claude-Session: https://claude.ai/code/session_019u3RKiNJY1sRknveHJy9iJ
1 parent 137fe2a commit b575203

2 files changed

Lines changed: 635 additions & 48 deletions

File tree

‎packages/cli/src/adapters/codex.ts‎

Lines changed: 292 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -72,17 +72,54 @@ async function parseCodexSessionFile(
7272
// (e.g. when only rate_limits metadata changes). Dedupe to avoid double-counting tokens
7373
// and inflating estimated cost.
7474
let lastTokenUsageKey: string | undefined
75-
// Forked subagent rollouts (session_meta.source.subagent.thread_spawn) begin by
76-
// REPLAYING the parent session's entire token history, all re-stamped at the
77-
// subagent's creation second. The lastTokenUsageKey dedup above only catches
78-
// *consecutive identical* records within a file — it cannot catch this replay,
79-
// whose cumulative counts differ line to line, so the parent's usage (cached
80-
// input especially) was being counted once per subagent file and inflating
81-
// totals several-fold. Detect the replay block up front and skip its leading
82-
// run of token_count events. Mirrors ccusage detect_subagent_replay_second
83-
// (adapter/codex/parser.rs).
75+
// Forked rollouts (subagent thread_spawn, branch/goal/resume forks) begin by
76+
// REPLAYING the parent session's entire token history. The lastTokenUsageKey
77+
// dedup above only catches *consecutive identical* records within a file — it
78+
// cannot catch this replay, whose cumulative counts differ line to line, so the
79+
// parent's usage (cached input especially) was being counted once per forked
80+
// file and inflating totals several-fold.
81+
//
82+
// Codex copies the token counts verbatim but rewrites the timestamps, so match
83+
// the leading usage against the parent's own stream and drop what lines up. The
84+
// older same-second heuristic only held when the whole replay burst landed in
85+
// one second — it missed long replays and nested forks. It stays as the fallback
86+
// for when the parent log is unavailable or the copied history was compacted and
87+
// no longer starts at the parent's first event. Mirrors ccusage CodexReplayPlan
88+
// + detect_replay_second (adapter/codex/replay.rs, parser.rs).
8489
const replaySecond = detectSubagentReplaySecond(text, lines)
85-
let skipReplay = replaySecond !== undefined
90+
const replayPrefix = await codexReplayPrefix(filePath, lines)
91+
let replay: { kind: 'matching', index: number } | { kind: 'second' } | { kind: 'done' }
92+
= replayPrefix === undefined
93+
? (replaySecond === undefined ? { kind: 'done' } : { kind: 'second' })
94+
: { kind: 'matching', index: 0 }
95+
96+
// Whether a usage event is the rollout's own rather than replayed history.
97+
// Each arm either returns or advances the state toward 'done', so the loop only
98+
// re-runs to apply the event to the state it switched to.
99+
const admitReplayedUsage = (usageKey: string, eventTs: string): boolean => {
100+
for (;;) {
101+
if (replay.kind === 'matching') {
102+
if ((replayPrefix ?? [])[replay.index] === usageKey) {
103+
replay = { kind: 'matching', index: replay.index + 1 }
104+
return false
105+
}
106+
// Nothing matched, so the parent stream cannot anchor this replay: the log
107+
// is unavailable, or Codex rewrote the copied history. Fall back to the
108+
// rewritten-second burst — but only when not a single event lined up,
109+
// since a mid-prefix break means the child's own usage has started.
110+
replay = replay.index === 0 && replaySecond !== undefined ? { kind: 'second' } : { kind: 'done' }
111+
continue
112+
}
113+
if (replay.kind === 'second') {
114+
if (eventTs.slice(0, 19) === replaySecond) {
115+
return false
116+
}
117+
replay = { kind: 'done' }
118+
continue
119+
}
120+
return true
121+
}
122+
}
86123
// Branch/goal/resume forks copy the parent rollout's lines VERBATIM into the
87124
// new file, keeping the original timestamps — nothing is re-stamped, so the
88125
// same-second heuristic above can't catch them. But the file's own events can
@@ -230,9 +267,18 @@ async function parseCodexSessionFile(
230267
// measured against the copied prior cumulative, not zero.
231268
if (creationMs !== undefined && Date.parse(ts) < creationMs) {
232269
if (payloadType === 'token_count') {
233-
const copiedTotal = objectField(objectField(payload, 'info'), 'total_token_usage')
234-
if (Object.keys(copiedTotal).length > 0) {
235-
previousTotals = readCodexCumulative(copiedTotal)
270+
const step = codexTokenCountUsage(payload, previousTotals)
271+
previousTotals = step.nextTotals
272+
// Keep the parent-stream match aligned with what the anchor already
273+
// dropped. Both layers target the same copied history, so a copied event
274+
// the anchor removed must still consume its slot in the parent prefix —
275+
// otherwise the rollout's first real event gets compared against the
276+
// prefix head and can be mistaken for replay. Only advance on a match:
277+
// this line is discarded either way, so it must not push the state
278+
// machine off `matching` for the events that follow.
279+
if (step.usage && replay.kind === 'matching'
280+
&& (replayPrefix ?? [])[replay.index] === codexUsageKey(step.usage)) {
281+
replay = { kind: 'matching', index: replay.index + 1 }
236282
}
237283
}
238284
continue
@@ -308,46 +354,18 @@ async function parseCodexSessionFile(
308354
if (tierFromInfo) {
309355
serviceTier = tierFromInfo.toLowerCase()
310356
}
311-
// Per-turn usage: prefer last_token_usage; otherwise derive the delta from
312-
// the cumulative total_token_usage minus the running baseline, so token_count
313-
// events that carry only a cumulative total are counted instead of dropped.
314-
const totalUsage = objectField(info, 'total_token_usage')
315-
const hasTotal = Object.keys(totalUsage).length > 0
316-
// Codex re-emits a last_token_usage snapshot when only metadata around it
317-
// changed (rate limits, service tier). The cumulative is the authority: if
318-
// total_token_usage did not advance, no new tokens were spent, so the
319-
// snapshot is a repeat of one already counted. The lastTokenUsageKey dedup
320-
// below only catches *consecutive* repeats; this also catches re-emissions
321-
// separated by other token_count events. Mirrors ccusage's
322-
// cumulative_advanced filter (adapter/codex/parser.rs).
323-
const cumulativeAdvanced = !hasTotal
324-
|| previousTotals === undefined
325-
|| !sameCodexCumulative(previousTotals, readCodexCumulative(totalUsage))
326-
const usage = (cumulativeAdvanced ? tokenUsageFromPayload(payload) : undefined)
327-
?? (hasTotal ? codexUsageDelta(totalUsage, previousTotals, info) : undefined)
357+
const step = codexTokenCountUsage(payload, previousTotals)
328358
// Advance the baseline from every total_token_usage we see — including
329359
// replayed events we skip — so the first real event's delta is measured
330360
// against the right prior cumulative.
331-
if (hasTotal) {
332-
previousTotals = readCodexCumulative(totalUsage)
333-
}
361+
previousTotals = step.nextTotals
362+
const usage = step.usage
334363
if (usage) {
335-
// Drop the leading run of replayed parent-history token_count events in a
336-
// forked subagent rollout (all stamped at the replay second). The first
337-
// event at a later second is the subagent's own usage and ends the skip.
338-
if (skipReplay) {
339-
if (ts.slice(0, 19) === replaySecond) {
340-
break
341-
}
342-
skipReplay = false
364+
const usageKey = codexUsageKey(usage)
365+
// Drop the parent history a forked rollout replayed on open.
366+
if (!admitReplayedUsage(usageKey, ts)) {
367+
break
343368
}
344-
const usageKey = [
345-
usage.tokensInput,
346-
usage.tokensCachedInput,
347-
usage.tokensOutput,
348-
usage.tokensReasoningOutput,
349-
usage.tokensTotal,
350-
].join(':')
351369
if (usageKey === lastTokenUsageKey) {
352370
break
353371
}
@@ -765,6 +783,185 @@ function rolloutCreationMs(filePath: string): number | undefined {
765783
return Number.isFinite(ms) && ms > 0 ? ms : undefined
766784
}
767785

786+
// ── forked-session replay: match the copied prefix against the parent stream ──
787+
788+
interface CodexStreamEvent {
789+
key: string
790+
tsMs: number | undefined
791+
}
792+
793+
// The `session_meta` line names the session this rollout forked from: branch and
794+
// resume forks set `forked_from_id`, subagent spawns nest the id under
795+
// `source.subagent.thread_spawn.parent_thread_id`. The line's own timestamp is the
796+
// fork instant. Mirrors ccusage read_codex_session_metadata (adapter/codex/replay.rs).
797+
function codexForkInfo(lines: string[]): { parentId: string | undefined, forkedAtMs: number | undefined } {
798+
const raw = lines.length > 0 ? parseJsonLine(lines[0]) : undefined
799+
if (!raw || stringField(raw, 'type') !== 'session_meta') {
800+
return { parentId: undefined, forkedAtMs: undefined }
801+
}
802+
const payload = objectField(raw, 'payload')
803+
const threadSpawn = objectField(objectField(objectField(payload, 'source'), 'subagent'), 'thread_spawn')
804+
const parentId = stringField(payload, 'forked_from_id') || stringField(threadSpawn, 'parent_thread_id')
805+
const forkedAt = timestampFrom(raw.timestamp)
806+
return {
807+
parentId: parentId || undefined,
808+
forkedAtMs: forkedAt ? Date.parse(forkedAt) : undefined,
809+
}
810+
}
811+
812+
// Rollout filenames embed the session id: `rollout-<local time>-<session uuid>.jsonl`,
813+
// where the uuid equals `session_meta.payload.id` (verified against real rollouts).
814+
// That lets a parent be located from its id without opening candidate files.
815+
const ROLLOUT_SESSION_ID_RE = /^rollout-\d{4}-\d{2}-\d{2}T\d{2}-\d{2}-\d{2}-([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})\.jsonl$/
816+
817+
function codexSessionIdFromFileName(filePath: string): string | undefined {
818+
return path.basename(filePath).match(ROLLOUT_SESSION_ID_RE)?.[1]
819+
}
820+
821+
// Codex lays sessions out as `<CODEX_HOME>/sessions/YYYY/MM/DD/` and rotates old
822+
// ones into `archived_sessions/`, so a parent can sit under either root and under
823+
// any date. Outside that layout (relocated logs, tests) search the child's own
824+
// directory. Active sessions come first so they win over an archived copy.
825+
function codexRolloutSearchRoots(childPath: string): string[] {
826+
const dir = path.dirname(childPath)
827+
const parts = dir.split(path.sep)
828+
const index = Math.max(parts.lastIndexOf('sessions'), parts.lastIndexOf('archived_sessions'))
829+
if (index <= 0) {
830+
return [dir]
831+
}
832+
const codexRoot = parts.slice(0, index).join(path.sep)
833+
return [path.join(codexRoot, 'sessions'), path.join(codexRoot, 'archived_sessions')]
834+
}
835+
836+
const codexSessionIndexCache = new Map<string, Promise<Map<string, string>>>()
837+
838+
async function buildCodexSessionIndex(roots: string[]): Promise<Map<string, string>> {
839+
const index = new Map<string, string>()
840+
for (const root of roots) {
841+
let files: string[]
842+
try {
843+
files = await listJsonlFiles(root)
844+
}
845+
catch {
846+
continue
847+
}
848+
for (const file of files) {
849+
const id = codexSessionIdFromFileName(file)
850+
if (id && !index.has(id)) {
851+
index.set(id, file)
852+
}
853+
}
854+
}
855+
return index
856+
}
857+
858+
async function codexRolloutPathForSession(childPath: string, sessionId: string): Promise<string | undefined> {
859+
const roots = codexRolloutSearchRoots(childPath)
860+
const cacheKey = roots.join('\0')
861+
let index = codexSessionIndexCache.get(cacheKey)
862+
if (!index) {
863+
index = buildCodexSessionIndex(roots)
864+
codexSessionIndexCache.set(cacheKey, index)
865+
}
866+
const sessionPaths = await index
867+
return sessionPaths.get(sessionId)
868+
}
869+
870+
interface CodexUsageStreamCacheEntry { mtimeMs: number, size: number, stream: CodexStreamEvent[] }
871+
872+
const codexUsageStreamCache = new Map<string, CodexUsageStreamCacheEntry>()
873+
874+
/**
875+
* The usage a session recorded, in order, with NO replay filtering applied.
876+
*
877+
* A nested fork replays its parent's whole stream — including the history that
878+
* parent had itself copied from the grandparent — so the prefix only lines up
879+
* against the *unfiltered* stream. Mirrors ccusage read_usage_events, which visits
880+
* the parent with `replayed_prefix: None` (adapter/codex/replay.rs).
881+
*
882+
* Only session-format `token_count` lines are read: a rollout must carry a
883+
* session_meta to be named as someone's parent, so headless `codex exec` files
884+
* never appear here.
885+
*/
886+
async function codexUsageStream(filePath: string): Promise<CodexStreamEvent[]> {
887+
let info
888+
try {
889+
info = await stat(filePath)
890+
}
891+
catch {
892+
return []
893+
}
894+
const cached = codexUsageStreamCache.get(filePath)
895+
if (cached && cached.mtimeMs === info.mtimeMs && cached.size === info.size) {
896+
return cached.stream
897+
}
898+
899+
let text: string
900+
try {
901+
text = await readFile(filePath, 'utf8')
902+
}
903+
catch {
904+
return []
905+
}
906+
907+
const stream: CodexStreamEvent[] = []
908+
let previousTotals: CodexCumulative | undefined
909+
for (const line of text.split('\n')) {
910+
if (!line) {
911+
continue
912+
}
913+
const raw = parseJsonLine(line)
914+
if (!raw || stringField(raw, 'type') !== 'event_msg') {
915+
continue
916+
}
917+
const payload = objectField(raw, 'payload')
918+
if (stringField(payload, 'type') !== 'token_count') {
919+
continue
920+
}
921+
const step = codexTokenCountUsage(payload, previousTotals)
922+
previousTotals = step.nextTotals
923+
if (!step.usage) {
924+
continue
925+
}
926+
const ts = timestampFrom(raw.timestamp) || timestampFrom(payload.timestamp)
927+
stream.push({ key: codexUsageKey(step.usage), tsMs: ts ? Date.parse(ts) : undefined })
928+
}
929+
930+
codexUsageStreamCache.set(filePath, { mtimeMs: info.mtimeMs, size: info.size, stream })
931+
return stream
932+
}
933+
934+
/**
935+
* Usage keys a forked rollout replayed from its parent, in order.
936+
*
937+
* Returns undefined for rollouts that are not forks, and an empty array for forks
938+
* whose parent log is unavailable — which tells the caller to fall back to the
939+
* rewritten-second heuristic. Mirrors CodexReplayPlan::replay_prefix.
940+
*/
941+
async function codexReplayPrefix(filePath: string, lines: string[]): Promise<string[] | undefined> {
942+
const { parentId, forkedAtMs } = codexForkInfo(lines)
943+
if (!parentId) {
944+
return undefined
945+
}
946+
// A session listing itself as its own parent would match its whole stream and
947+
// drop every event it recorded.
948+
if (codexSessionIdFromFileName(filePath) === parentId) {
949+
return []
950+
}
951+
const parentPath = await codexRolloutPathForSession(filePath, parentId)
952+
if (!parentPath || path.resolve(parentPath) === path.resolve(filePath)) {
953+
return []
954+
}
955+
const stream = await codexUsageStream(parentPath)
956+
// Usage the parent recorded after the fork was never replayed, so it must not
957+
// mask the child's own events.
958+
const after = forkedAtMs === undefined
959+
? -1
960+
: stream.findIndex(event => event.tsMs !== undefined && event.tsMs > forkedAtMs)
961+
const replayed = after === -1 ? stream : stream.slice(0, after)
962+
return replayed.map(event => event.key)
963+
}
964+
768965
// Codex spawns a subagent into its own rollout file that opens by replaying the
769966
// parent session's token history, re-stamped at the subagent's creation second.
770967
// Return that second when this file is such a replay so the parser can drop the
@@ -841,6 +1038,53 @@ function readCodexCumulative(usage: Record<string, unknown>): CodexCumulative {
8411038
}
8421039
}
8431040

1041+
type CodexUsageMetrics = NonNullable<ReturnType<typeof tokenUsageFromPayload>>
1042+
1043+
// Identity of a usage event: the token counts themselves. Used both to dedupe
1044+
// consecutive repeats and to line a forked rollout's replayed history up against
1045+
// its parent's stream, where Codex copies the counts but rewrites everything else.
1046+
function codexUsageKey(usage: CodexUsageMetrics): string {
1047+
return [
1048+
usage.tokensInput,
1049+
usage.tokensCachedInput,
1050+
usage.tokensOutput,
1051+
usage.tokensReasoningOutput,
1052+
usage.tokensTotal,
1053+
].join(':')
1054+
}
1055+
1056+
/**
1057+
* Per-turn usage carried by one `token_count` payload, plus the cumulative
1058+
* baseline to carry into the next one.
1059+
*
1060+
* Prefers `last_token_usage`; otherwise derives the delta from the cumulative
1061+
* `total_token_usage` minus the running baseline, so token_count events that
1062+
* carry only a cumulative total are counted instead of dropped.
1063+
*
1064+
* Codex re-emits a last_token_usage snapshot when only metadata around it changed
1065+
* (rate limits, service tier). The cumulative is the authority: if
1066+
* total_token_usage did not advance, no new tokens were spent, so the snapshot is
1067+
* a repeat of one already counted. Mirrors ccusage's cumulative_advanced filter
1068+
* (adapter/codex/parser.rs).
1069+
*/
1070+
function codexTokenCountUsage(
1071+
payload: Record<string, unknown>,
1072+
previousTotals: CodexCumulative | undefined,
1073+
): { usage: CodexUsageMetrics | undefined, nextTotals: CodexCumulative | undefined } {
1074+
const info = objectField(payload, 'info')
1075+
const totalUsage = objectField(info, 'total_token_usage')
1076+
const hasTotal = Object.keys(totalUsage).length > 0
1077+
const cumulativeAdvanced = !hasTotal
1078+
|| previousTotals === undefined
1079+
|| !sameCodexCumulative(previousTotals, readCodexCumulative(totalUsage))
1080+
const usage = (cumulativeAdvanced ? tokenUsageFromPayload(payload) : undefined)
1081+
?? (hasTotal ? codexUsageDelta(totalUsage, previousTotals, info) : undefined)
1082+
return {
1083+
usage,
1084+
nextTotals: hasTotal ? readCodexCumulative(totalUsage) : previousTotals,
1085+
}
1086+
}
1087+
8441088
function sameCodexCumulative(a: CodexCumulative, b: CodexCumulative): boolean {
8451089
return a.input === b.input
8461090
&& a.cached === b.cached

0 commit comments

Comments
 (0)