feat(review): live token counts for claude and codex reviewers · Entire

feat(review): live token counts for claude and codex reviewers

5860d3a→main·

peyton-alt·1w ago·7 files·+525 added/-23 removed

Salvages the live-token slices from closed #1370 onto current main. Reviewer token totals previously appeared only at run completion: claude's parser emitted one Tokens event at the terminal result envelope, and codex's exec --json stdout only carries usage on turn.completed (a review is usually a single turn).

Claude: assistant envelopes carry a usage snapshot taken at the start of each turn — input/cache counts are real but output_tokens is a 1-8 token 'initial decision' stub. Emit input-only Tokens{In, Out: 0} per usage-carrying assistant envelope so consumers see context growth during the run; the result envelope still delivers the final {In, Out} aggregrate. Emitting the misleading early output snapshot was rejected after capturing real claude stream-json.

Codex: tail the rollout transcript (located by thread.started's thread_id, same source codex's interactive UI reads) and emit cumulative Tokens per token_count event, deduped on movement. The parser stops the tailer and waits for it before closing the event channel. turn.completed now also emits Tokens per turn (multi-turn runs update at each boundary) with the previous post-loop emission kept as a defensive backstop.

Verified end-to-end via agent shims driving entire review: claude streaming envelopes with usage snapshots, and a codex shim writing rollout token_count lines picked up by the tailer, both to clean completion. Race-detector clean.

Co-Authored-By: Claude Fable 5 noreply@anthropic.com

Sessions

01KWYC8H544S7DV9G0YSEQNV9HView transcript

Changes

7

54 unmodified lines

// Emits Started first, Finished{Success:...} last (success follows result.is_error). // On a scanner error (torn stream), emits RunError then Finished{Success:false}. // // Tokens are emitted only at the terminal result envelope, not // incrementally — claude's per-assistant usage fields aren't cumulative // and summing them across messages would double-count. // Live-token semantics: Claude's assistant envelopes carry a usage snapshot // taken at the START of each turn — input_tokens/cache_* are populated but // output_tokens is essentially zero (a 1–8 token "initial decision" count // that does not update as text streams). The true per-turn output is only // surfaced on result (aggregate across all turns in the run) or on the // late message_delta event of --include-partial-messages mode. // // We emit Tokens{In: <sum>, Out: 0} on each assistant envelope so consumers // can show context size growing across multi-turn runs, and Tokens{In, Out} // on result with the final aggregate. Sending the misleading early // "Out: 1" snapshot was removed after capturing real claude stream-json.

// Package-private; called directly from this package's tests so they can // drive raw stdout fixtures through the parser without going through the

type claudeMessage struct { Content []claudeBlock json:"content" // Usage on assistant envelopes is the per-turn-START snapshot — input // counts are populated but output_tokens reflects only the model's // initial decision, not the streamed text. Final aggregate usage // arrives on the result envelope. Reuses messageUsage (declared in // types.go) to stay aligned with the transcript-parser usage shape. Usage messageUsage json:"usage" }

type claudeBlock struct {

}


Mcmd/entire/cli/agent/claudecode/reviewer.go+31/-3

// TestParseClaudeOutput_EmitsInputOnlyDuringRun captures the live-token // contract for Claude: the assistant envelopes' usage block carries a // per-turn-start snapshot where output_tokens is essentially zero (1–8 // tokens of initial decision), so we surface input only on those envelopes // and let the terminal result envelope carry the true {In, Out} total.

func TestParseClaudeOutput_EmitsInputOnlyDuringRun(t *testing.T) { t.Parallel() f, err := os.Open("testdata/stream_with_deltas.jsonl") if err != nil { t.Fatal(err) } defer f.Close()

var events []reviewtypes.Event for ev := range parseClaudeOutput(f) { events = append(events, ev) }

var tokens []reviewtypes.Tokens sawFinished := false for _, e := range events { switch ev := e.(type) { case reviewtypes.Tokens: if sawFinished { t.Errorf("Tokens event arrived AFTER Finished — wrong ordering") } tokens = append(tokens, ev) case reviewtypes.Finished: sawFinished = true } } if len(tokens) < 2 { t.Fatalf("expected >=2 Tokens events (assistant snapshots + final), got %d", len(tokens)) }

// All Tokens BEFORE the final one are assistant-envelope snapshots: // input is populated; Out must be exactly 0 to avoid showing the // misleading per-turn-start initial-decision count in the TUI. for i := range len(tokens) - 1 { tk := tokens[i] if tk.Out != 0 { t.Errorf("tokens[%d].Out = %d, want 0 (assistant envelope snapshot — output not yet known)", i, tk.Out) } if tk.In == 0 { t.Errorf("tokens[%d].In = 0, want >0 (assistant envelope carries input)", i) } }

// The terminal Tokens event is from result and carries the true // aggregate output count (2511 in the fixture). last := tokens[len(tokens)-1] if last.Out == 0 { t.Error("final Tokens.Out = 0, want non-zero (result envelope final tally)") } if last.In == 0 { t.Error("final Tokens.In = 0, want non-zero") } } // collectEvents drains an event channel into a slice. func collectEvents(ch <-chan reviewtypes.Event) []reviewtypes.Event { var events []reviewtypes.Event }


Mcmd/entire/cli/agent/claudecode/reviewer_test.go+74/-3

1 2 3 4 5 6 7

{"type":"system","subtype":"init","cwd":"/redacted/worktree","session_id":"a905e63f-aaaa-aaaa-aaaa-aaaaaaaaaaaa","model":"claude-haiku-4-5","permissionMode":"plan","output_style":"default","apiKeySource":"none","uuid":"redacted-uuid-1"} {"type":"assistant","message":{"model":"claude-haiku-4-5-20251001","id":"msg_turn1","type":"message","role":"assistant","content":[{"type":"thinking","thinking":"Analyzing the request..."}],"stop_reason":null,"usage":{"input_tokens":10,"cache_creation_input_tokens":56267,"cache_read_input_tokens":0,"output_tokens":6,"service_tier":"standard"}},"session_id":"a905e63f-aaaa-aaaa-aaaa-aaaaaaaaaaaa","uuid":"redacted-uuid-2"} {"type":"assistant","message":{"model":"claude-haiku-4-5-20251001","id":"msg_turn1","type":"message","role":"assistant","content":[{"type":"text","text":"I'll outline a plan first."}],"stop_reason":null,"usage":{"input_tokens":10,"cache_creation_input_tokens":56267,"cache_read_input_tokens":0,"output_tokens":6,"service_tier":"standard"}},"session_id":"a905e63f-aaaa-aaaa-aaaa-aaaaaaaaaaaa","uuid":"redacted-uuid-3"} {"type":"assistant","message":{"model":"claude-haiku-4-5-20251001","id":"msg_turn1","type":"message","role":"assistant","content":[{"type":"tool_use","id":"toolu_01","name":"Write","input":{"file_path":"plan.md","content":"plan body"}}],"stop_reason":null,"usage":{"input_tokens":10,"cache_creation_input_tokens":56267,"cache_read_input_tokens":0,"output_tokens":6,"service_tier":"standard"}},"session_id":"a905e63f-aaaa-aaaa-aaaa-aaaaaaaaaaaa","uuid":"redacted-uuid-4"} {"type":"assistant","message":{"model":"claude-haiku-4-5-20251001","id":"msg_turn2","type":"message","role":"assistant","content":[{"type":"text","text":"Plan created, ready to proceed."}],"stop_reason":null,"usage":{"input_tokens":5,"cache_creation_input_tokens":10066,"cache_read_input_tokens":46555,"output_tokens":1,"service_tier":"standard"}},"session_id":"a905e63f-aaaa-aaaa-aaaa-aaaaaaaaaaaa","uuid":"redacted-uuid-5"} {"type":"assistant","message":{"model":"claude-haiku-4-5-20251001","id":"msg_turn3","type":"message","role":"assistant","content":[{"type":"text","text":"Found 3 issues."}],"stop_reason":null,"usage":{"input_tokens":6,"cache_creation_input_tokens":107,"cache_read_input_tokens":56621,"output_tokens":2,"service_tier":"standard"}},"session_id":"a905e63f-aaaa-aaaa-aaaa-aaaaaaaaaaaa","uuid":"redacted-uuid-6"} {"type":"result","subtype":"success","is_error":false,"duration_ms":29272,"num_turns":3,"result":"Found 3 issues.","stop_reason":"end_turn","session_id":"a905e63f-aaaa-aaaa-aaaa-aaaaaaaaaaaa","total_cost_usd":0.105,"usage":{"input_tokens":21,"cache_creation_input_tokens":66440,"cache_read_input_tokens":103176,"output_tokens":2511,"service_tier":"standard"},"uuid":"redacted-uuid-7"}


Acmd/entire/cli/agent/claudecode/testdata/stream_with_deltas.jsonl+7

1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134

package codex

import ( "bytes" "context" "encoding/json" "log/slog" "os" "time"

"github.com/entireio/cli/cmd/entire/cli/logging" reviewtypes "github.com/entireio/cli/cmd/entire/cli/review/types" )

// Polling/tailing cadence for the rollout token tailer. const ( rolloutPollInterval = 300 * time.Millisecond rolloutPollAttempts = 100 // ~30s for codex to create the rollout file rolloutTailInterval = 400 * time.Millisecond rolloutReadChunk = 8192 )

// tailRolloutTokens resolves the codex rollout transcript for threadID and // tails it, emitting a cumulative reviewtypes.Tokens event for every // token_count codex writes (~once per model turn). codex's exec --json // stdout only carries usage on turn.completed envelopes, and a review is // usually a single turn — so without this, consumers see no token movement // until the run ends. The rollout file is the same source codex's // interactive UI reads for its live token counter. // // token_count.total_token_usage is a running total, so each emission is an // absolute count — matching consumers' overwrite-not-sum semantics. // Duplicate totals are suppressed so we only emit on real movement. // // Returns when stop is closed (the stdout stream ended) or the rollout file // never appears. The caller must wait for this to return before closing the // event channel (see parseCodexOutputBuf), so a send here can never race a // channel close. func tailRolloutTokens(threadID string, out chan<- reviewtypes.Event, stop <-chan struct{}) { ctx := context.Background() sessionDir, err := (&CodexAgent{}).GetSessionDir("") if err != nil { logging.Debug(ctx, "codex token tail: session dir unresolved", slog.String("error", err.Error())) return } path := waitForRollout(sessionDir, threadID, stop) if path == "" { return } f, err := os.Open(path) //nolint:gosec // path is a glob match under codex's session dir, not user input if err != nil { logging.Debug(ctx, "codex token tail: open rollout failed", slog.String("error", err.Error())) return } defer f.Close()

// Tail via os.File.Read rather than bufio.Reader: bufio is sticky on EOF // and would never observe lines codex appends after we first catch up. var pending []byte chunk := make([]byte, rolloutReadChunk) lastIn, lastOut := -1, -1 ticker := time.NewTicker(rolloutTailInterval) defer ticker.Stop() for { for { n, readErr := f.Read(chunk) if n > 0 { pending = append(pending, chunk[:n]...) for { idx := bytes.IndexByte(pending, '\n') if idx < 0 { break } line := pending[:idx] pending = pending[idx+1:] in, outTok, ok := parseRolloutTokenCount(line) if !ok || (in == lastIn && outTok == lastOut) { continue } lastIn, lastOut = in, outTok select { case out <- reviewtypes.Tokens{In: in, Out: outTok}: case <-stop: return } } } if readErr != nil { break // EOF or error — wait for the file to grow, then retry } } select { case <-stop: return case <-ticker.C: } } }

// waitForRollout polls for the rollout file matching threadID, returning its // path or "" if stop fires or the attempts are exhausted. func waitForRollout(sessionDir, threadID string, stop <-chan struct{}) string { for range rolloutPollAttempts { if path := findRolloutBySessionID(sessionDir, threadID); path != "" { return path } select { case <-stop: return "" case <-time.After(rolloutPollInterval): } } return "" }

// parseRolloutTokenCount extracts cumulative input/output token totals from one // rollout JSONL line. ok is false for any line that isn't a token_count event // carrying total_token_usage. Reuses the rolloutLine/eventMsgPayload/ // tokenCountInfo shapes from transcript.go so the two readers can't drift. func parseRolloutTokenCount(data []byte) (in, out int, ok bool) { var line rolloutLine if json.Unmarshal(data, &line) != nil || line.Type != "event_msg" { return 0, 0, false } var evt eventMsgPayload if json.Unmarshal(line.Payload, &evt) != nil || evt.Type != "token_count" || len(evt.Info) == 0 { return 0, 0, false } var info tokenCountInfo if json.Unmarshal(evt.Info, &info) != nil || info.TotalTokenUsage == nil { return 0, 0, false } return info.TotalTokenUsage.InputTokens, info.TotalTokenUsage.OutputTokens, true }


Acmd/entire/cli/agent/codex/review_tokens.go+134

1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 8 unmodified lines


Mcmd/entire/cli/agent/codex/review_tokens_test.go+161