From 58b644537ba026602015d0b2147d2b6496ed03ea Mon Sep 17 00:00:00 2001 From: Claude Date: Mon, 5 Oct 2026 11:24:04 +0000 Subject: [PATCH 1/3] feat(model): add AI usage backfill API client and helpers - SendHTTPRequestJSON now returns a typed HTTPStatusError for non-2xx responses so callers can tell retryable errors from client errors. The error text is unchanged. - Wire types and senders for the server's /api/v1/cc/backfill endpoints. Backfill events use pointer fields so a failed tool call is sent as "success": false instead of being dropped by omitempty. - Shared helpers for the transcript parsers: stable "bf1:" event ids from natural keys, a JSONL reader without a line length cap, file discovery, session selection, and request packing that keeps whole sessions together and marks a session completed on its last batch. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LMCwtubXzBYhQTnX44vhzF --- model/aicode_backfill_common.go | 299 +++++++++++++++++++++++++++ model/aicode_backfill_common_test.go | 126 +++++++++++ model/aicode_backfill_types.go | 108 ++++++++++ model/aicode_otel_claude_settings.go | 15 +- model/api.base.go | 16 +- model/api_aicode_backfill.go | 71 +++++++ model/api_aicode_backfill_test.go | 119 +++++++++++ 7 files changed, 747 insertions(+), 7 deletions(-) create mode 100644 model/aicode_backfill_common.go create mode 100644 model/aicode_backfill_common_test.go create mode 100644 model/aicode_backfill_types.go create mode 100644 model/api_aicode_backfill.go create mode 100644 model/api_aicode_backfill_test.go diff --git a/model/aicode_backfill_common.go b/model/aicode_backfill_common.go new file mode 100644 index 0000000..eaaf4b0 --- /dev/null +++ b/model/aicode_backfill_common.go @@ -0,0 +1,299 @@ +package model + +import ( + "bufio" + "bytes" + "crypto/sha256" + "encoding/hex" + "errors" + "io" + "io/fs" + "os" + "path/filepath" + "runtime" + "sort" + "strings" + "time" +) + +// backfillSettleWindow keeps sessions that were active very recently out of +// a backfill: their transcripts may still grow, and since uploads are +// idempotent per event id, a partial usage line would win forever. +const backfillSettleWindow = 30 * time.Minute + +// BackfillOptions controls which local sessions a backfill picks up. +type BackfillOptions struct { + Since time.Time // sessions starting before this are skipped (zero: no bound) + Until time.Time // sessions starting at or after this are skipped (zero: no bound) + NoPrompts bool // upload prompt lengths but not the prompt text + Now time.Time +} + +// BackfillSession is one Claude Code or Codex session rebuilt from local +// transcripts, with its events sorted by time. +type BackfillSession struct { + ID string + ClientType string + Start time.Time + End time.Time + Events []AICodeBackfillEvent +} + +// BackfillParseResult is what a transcript parser found on disk. +type BackfillParseResult struct { + Sessions []*BackfillSession + Files int + BadLines int +} + +// backfillHash derives both the event id and the in-memory dedup key from a +// natural key of the event (message id, tool call id, ...). The session id is +// deliberately not part of it, so copies of the same item in resumed or +// forked transcripts collapse into one event. +func backfillHash(clientType, naturalKey string) [sha256.Size]byte { + return sha256.Sum256([]byte(clientType + "\x00" + naturalKey)) +} + +func backfillEventID(clientType, naturalKey string) string { + sum := backfillHash(clientType, naturalKey) + return "bf1:" + hex.EncodeToString(sum[:20]) +} + +type backfillKey [16]byte + +func newBackfillKey(clientType, naturalKey string) backfillKey { + sum := backfillHash(clientType, naturalKey) + var k backfillKey + copy(k[:], sum[:16]) + return k +} + +func newBackfillEvent(clientType, eventType, naturalKey, sessionID string, ts time.Time) AICodeBackfillEvent { + return AICodeBackfillEvent{ + EventID: backfillEventID(clientType, naturalKey), + EventType: eventType, + ClientType: clientType, + Timestamp: ts.Unix(), + SessionID: sessionID, + } +} + +// backfillMachine describes this machine the same way the live OTEL resource +// attributes do (see claudeSettingsResourceAttributes). +type backfillMachine struct { + OSType string + HostArch string + UserName string + MachineName string +} + +func currentBackfillMachine() backfillMachine { + userName, machineName := aiCodeResourceIdentity() + return backfillMachine{ + OSType: runtime.GOOS, + HostArch: runtime.GOARCH, + UserName: userName, + MachineName: machineName, + } +} + +func (m backfillMachine) apply(e *AICodeBackfillEvent) { + e.OSType = m.OSType + e.HostArch = m.HostArch + e.UserName = m.UserName + e.MachineName = m.MachineName + e.TeamID = "shelltime" +} + +// readJSONLLines calls fn for every non-blank line of a JSONL file. Lines +// can be several megabytes (Claude Code stores tool results inline), so no +// line length limit is applied. fn owns the slice it receives. +func readJSONLLines(path string, fn func(line []byte)) error { + f, err := os.Open(path) + if err != nil { + return err + } + defer f.Close() + + r := bufio.NewReaderSize(f, 1<<20) + for { + line, err := r.ReadBytes('\n') + line = bytes.TrimRight(line, "\r\n") + if len(bytes.TrimSpace(line)) > 0 { + fn(line) + } + if errors.Is(err, io.EOF) { + return nil + } + if err != nil { + return err + } + } +} + +// collectJSONL lists every *.jsonl file under the given roots, sorted. +// Files last modified more than a day before since cannot hold sessions in +// range and are skipped. Missing roots are ignored. +func collectJSONL(roots []string, since time.Time) ([]string, error) { + seenRoots := map[string]bool{} + seenFiles := map[string]bool{} + var files []string + + for _, root := range roots { + resolved, err := filepath.EvalSymlinks(root) + if err != nil { + if errors.Is(err, fs.ErrNotExist) { + continue + } + return nil, err + } + if seenRoots[resolved] { + continue + } + seenRoots[resolved] = true + + err = filepath.WalkDir(resolved, func(path string, d fs.DirEntry, walkErr error) error { + if walkErr != nil { + // Unreadable directory: skip it rather than abort the backfill. + if d != nil && d.IsDir() { + return filepath.SkipDir + } + return nil + } + if !d.Type().IsRegular() || !strings.HasSuffix(d.Name(), ".jsonl") { + return nil + } + if !since.IsZero() { + info, err := d.Info() + if err == nil && info.ModTime().Before(since.Add(-24*time.Hour)) { + return nil + } + } + if !seenFiles[path] { + seenFiles[path] = true + files = append(files, path) + } + return nil + }) + if err != nil { + return nil, err + } + } + + sort.Strings(files) + return files, nil +} + +// splitPathList splits a comma-separated list of paths, as used by the +// CLAUDE_CONFIG_DIR and CODEX_HOME environment variables. +func splitPathList(value string) []string { + var out []string + for _, p := range strings.Split(value, ",") { + if p = strings.TrimSpace(p); p != "" { + out = append(out, p) + } + } + return out +} + +// SelectBackfillSessions keeps the sessions that start within the requested +// range and are not still active. It returns them and how many active +// sessions were held back. +func SelectBackfillSessions(sessions []*BackfillSession, opts BackfillOptions) ([]*BackfillSession, int) { + now := opts.Now + if now.IsZero() { + now = time.Now() + } + settled := now.Add(-backfillSettleWindow) + + var kept []*BackfillSession + active := 0 + for _, s := range sessions { + if len(s.Events) == 0 { + continue + } + if !opts.Since.IsZero() && s.Start.Before(opts.Since) { + continue + } + if !opts.Until.IsZero() && !s.Start.Before(opts.Until) { + continue + } + if s.End.After(settled) { + active++ + continue + } + kept = append(kept, s) + } + return kept, active +} + +// PackBackfillBatches groups sessions into upload requests of at most +// maxEvents events and maxCompleted completed sessions. Whole sessions are +// kept together; only a session larger than maxEvents is split, and a +// session is listed in completedSessionIds only on the request carrying its +// last event, so the server builds its summary once all events are stored. +func PackBackfillBatches(clientType string, sessions []*BackfillSession, maxEvents, maxCompleted int) []AICodeBackfillRequest { + ordered := make([]*BackfillSession, len(sessions)) + copy(ordered, sessions) + sort.SliceStable(ordered, func(i, j int) bool { + if !ordered[i].Start.Equal(ordered[j].Start) { + return ordered[i].Start.Before(ordered[j].Start) + } + return ordered[i].ID < ordered[j].ID + }) + + var batches []AICodeBackfillRequest + cur := AICodeBackfillRequest{ClientType: clientType} + flush := func() { + if len(cur.Events) > 0 || len(cur.CompletedSessionIDs) > 0 { + batches = append(batches, cur) + } + cur = AICodeBackfillRequest{ClientType: clientType} + } + + for _, s := range ordered { + events := s.Events + if len(cur.CompletedSessionIDs) >= maxCompleted || + (len(cur.Events) > 0 && len(cur.Events)+len(events) > maxEvents) { + flush() + } + for len(events) > 0 { + room := maxEvents - len(cur.Events) + if room <= 0 { + flush() + room = maxEvents + } + n := min(room, len(events)) + cur.Events = append(cur.Events, events[:n]...) + events = events[n:] + if len(events) > 0 { + flush() + } + } + cur.CompletedSessionIDs = append(cur.CompletedSessionIDs, s.ID) + } + flush() + return batches +} + +// finalizeBackfillSession sorts a session's events and sets its time range. +func finalizeBackfillSession(s *BackfillSession) { + sort.SliceStable(s.Events, func(i, j int) bool { + return s.Events[i].Timestamp < s.Events[j].Timestamp + }) + if len(s.Events) == 0 { + return + } + first := time.Unix(s.Events[0].Timestamp, 0) + last := time.Unix(s.Events[len(s.Events)-1].Timestamp, 0) + if s.Start.IsZero() || first.Before(s.Start) { + s.Start = first + } + if s.End.IsZero() || last.After(s.End) { + s.End = last + } +} + +func intRef(v int) *int { return &v } + +func boolRef(v bool) *bool { return &v } diff --git a/model/aicode_backfill_common_test.go b/model/aicode_backfill_common_test.go new file mode 100644 index 0000000..0422449 --- /dev/null +++ b/model/aicode_backfill_common_test.go @@ -0,0 +1,126 @@ +package model + +import ( + "fmt" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func testBackfillSession(id string, start time.Time, events int) *BackfillSession { + s := &BackfillSession{ID: id, ClientType: AICodeClientClaudeCode} + for i := 0; i < events; i++ { + ts := start.Add(time.Duration(i) * time.Second) + s.Events = append(s.Events, newBackfillEvent(AICodeClientClaudeCode, AICodeEventApiRequest, fmt.Sprintf("%s-%d", id, i), id, ts)) + } + finalizeBackfillSession(s) + return s +} + +func TestBackfillEventIDIsStable(t *testing.T) { + // Event ids must never change between CLI versions, or re-running a + // backfill would upload everything twice. + assert.Equal(t, "bf1:efb74f1d4dd4e799fe39035fce73f1d5da2c3a5c", backfillEventID(AICodeClientClaudeCode, "msg_1:req_1")) + assert.Regexp(t, `^bf1:[0-9a-f]{40}$`, backfillEventID(AICodeClientCodex, "call_1")) + assert.NotEqual(t, backfillEventID(AICodeClientClaudeCode, "x"), backfillEventID(AICodeClientCodex, "x")) + assert.Equal(t, newBackfillKey(AICodeClientCodex, "x"), newBackfillKey(AICodeClientCodex, "x")) +} + +func TestPackBackfillBatches(t *testing.T) { + base := time.Date(2026, 9, 1, 10, 0, 0, 0, time.UTC) + sessions := []*BackfillSession{ + testBackfillSession("c", base.Add(2*time.Hour), 3), + testBackfillSession("a", base, 4), + testBackfillSession("big", base.Add(time.Hour), 12), + } + + batches := PackBackfillBatches(AICodeClientClaudeCode, sessions, 5, 50) + + // a (4) fits alone; big (12) is split 5+5+2; c (3) joins big's last chunk. + require.Len(t, batches, 4) + assert.Len(t, batches[0].Events, 4) + assert.Equal(t, []string{"a"}, batches[0].CompletedSessionIDs) + assert.Len(t, batches[1].Events, 5) + assert.Empty(t, batches[1].CompletedSessionIDs) + assert.Len(t, batches[2].Events, 5) + assert.Empty(t, batches[2].CompletedSessionIDs) + assert.Len(t, batches[3].Events, 5) + assert.Equal(t, []string{"big", "c"}, batches[3].CompletedSessionIDs) + for _, b := range batches { + assert.Equal(t, AICodeClientClaudeCode, b.ClientType) + assert.LessOrEqual(t, len(b.Events), 5) + } +} + +func TestPackBackfillBatchesRespectsCompletedCap(t *testing.T) { + base := time.Date(2026, 9, 1, 10, 0, 0, 0, time.UTC) + var sessions []*BackfillSession + for i := 0; i < 5; i++ { + sessions = append(sessions, testBackfillSession(fmt.Sprintf("s%d", i), base.Add(time.Duration(i)*time.Minute), 1)) + } + + batches := PackBackfillBatches(AICodeClientCodex, sessions, 500, 2) + + require.Len(t, batches, 3) + assert.Equal(t, []string{"s0", "s1"}, batches[0].CompletedSessionIDs) + assert.Equal(t, []string{"s2", "s3"}, batches[1].CompletedSessionIDs) + assert.Equal(t, []string{"s4"}, batches[2].CompletedSessionIDs) +} + +func TestSelectBackfillSessions(t *testing.T) { + now := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + old := testBackfillSession("old", now.AddDate(0, 0, -20), 2) + mid := testBackfillSession("mid", now.AddDate(0, 0, -5), 2) + active := testBackfillSession("active", now.Add(-10*time.Minute), 2) + empty := &BackfillSession{ID: "empty"} + + kept, activeCount := SelectBackfillSessions([]*BackfillSession{old, mid, active, empty}, BackfillOptions{Now: now}) + assert.Equal(t, []*BackfillSession{old, mid}, kept) + assert.Equal(t, 1, activeCount) + + kept, _ = SelectBackfillSessions([]*BackfillSession{old, mid}, BackfillOptions{Now: now, Since: now.AddDate(0, 0, -10)}) + assert.Equal(t, []*BackfillSession{mid}, kept) + + kept, _ = SelectBackfillSessions([]*BackfillSession{old, mid}, BackfillOptions{Now: now, Until: now.AddDate(0, 0, -10)}) + assert.Equal(t, []*BackfillSession{old}, kept) +} + +func TestReadJSONLLinesHandlesHugeLines(t *testing.T) { + path := filepath.Join(t.TempDir(), "big.jsonl") + huge := `{"x":"` + strings.Repeat("a", 10<<20) + `"}` + require.NoError(t, os.WriteFile(path, []byte("{\"a\":1}\r\n\n"+huge+"\n{\"b\":2}"), 0o644)) + + var lens []int + require.NoError(t, readJSONLLines(path, func(line []byte) { lens = append(lens, len(line)) })) + + assert.Equal(t, []int{7, len(huge), 7}, lens) +} + +func TestCollectJSONL(t *testing.T) { + root := t.TempDir() + require.NoError(t, os.MkdirAll(filepath.Join(root, "proj", "sub"), 0o755)) + recent := filepath.Join(root, "proj", "recent.jsonl") + stale := filepath.Join(root, "proj", "sub", "stale.jsonl") + require.NoError(t, os.WriteFile(recent, []byte("{}"), 0o644)) + require.NoError(t, os.WriteFile(stale, []byte("{}"), 0o644)) + require.NoError(t, os.WriteFile(filepath.Join(root, "proj", "notes.txt"), []byte("x"), 0o644)) + staleTime := time.Now().AddDate(0, 0, -30) + require.NoError(t, os.Chtimes(stale, staleTime, staleTime)) + + link := filepath.Join(t.TempDir(), "link") + require.NoError(t, os.Symlink(root, link)) + + all, err := collectJSONL([]string{root, link, filepath.Join(root, "missing")}, time.Time{}) + require.NoError(t, err) + assert.Len(t, all, 2) + + inRange, err := collectJSONL([]string{root}, time.Now().AddDate(0, 0, -7)) + require.NoError(t, err) + assert.Len(t, inRange, 1) + assert.Equal(t, "recent.jsonl", filepath.Base(inRange[0])) +} diff --git a/model/aicode_backfill_types.go b/model/aicode_backfill_types.go new file mode 100644 index 0000000..5ee007e --- /dev/null +++ b/model/aicode_backfill_types.go @@ -0,0 +1,108 @@ +package model + +import "time" + +// Client types as stored by the server. The live OTEL processor sends +// "claude-code", which the server maps to claude_code; backfill sends the +// canonical values directly. +const ( + AICodeClientClaudeCode = "claude_code" + AICodeClientCodex = "codex" +) + +// Backfill session statuses reported by the server. +const ( + AICodeBackfillStatusBackfilled = "backfilled" + AICodeBackfillStatusLive = "live" + AICodeBackfillStatusArchived = "archived" +) + +// Server limits for one backfill request. +const ( + AICodeBackfillMaxEvents = 500 + AICodeBackfillMaxCompleted = 50 + AICodeBackfillMaxSessionIDs = 1000 +) + +// AICodeBackfillEvent is one historical event sent to POST /api/v1/cc/backfill. +// It mirrors the server's CCEventData. Unlike AICodeOtelEvent, flags, counts +// and costs are pointers so that false and zero values survive omitempty +// (a failed tool call must be sent as "success": false). +type AICodeBackfillEvent struct { + EventID string `json:"eventId"` + EventType string `json:"eventType"` + ClientType string `json:"clientType"` + Timestamp int64 `json:"timestamp"` // unix seconds + + Model string `json:"model,omitempty"` + + Prompt string `json:"prompt,omitempty"` + PromptLength *int `json:"promptLength,omitempty"` + + ToolName string `json:"toolName,omitempty"` + Success *bool `json:"success,omitempty"` + DurationMs *int `json:"durationMs,omitempty"` + Error string `json:"error,omitempty"` + ToolParameters map[string]any `json:"toolParameters,omitempty"` + ToolArguments map[string]any `json:"toolArguments,omitempty"` + CallID string `json:"callId,omitempty"` + + CostUSD *float64 `json:"costUsd,omitempty"` + InputTokens *int `json:"inputTokens,omitempty"` + OutputTokens *int `json:"outputTokens,omitempty"` + CacheReadTokens *int `json:"cacheReadTokens,omitempty"` + CacheCreationTokens *int `json:"cacheCreationTokens,omitempty"` + ReasoningTokens *int `json:"reasoningTokens,omitempty"` + + EventKind string `json:"eventKind,omitempty"` + Provider string `json:"provider,omitempty"` + ApprovalPolicy string `json:"approvalPolicy,omitempty"` + SandboxPolicy string `json:"sandboxPolicy,omitempty"` + ReasoningEffort string `json:"reasoningEffort,omitempty"` + + SessionID string `json:"sessionId"` + ConversationID string `json:"conversationId,omitempty"` + AppVersion string `json:"appVersion,omitempty"` + OSType string `json:"osType,omitempty"` + HostArch string `json:"hostArch,omitempty"` + Pwd string `json:"pwd,omitempty"` + UserName string `json:"userName,omitempty"` + MachineName string `json:"machineName,omitempty"` + TeamID string `json:"teamId,omitempty"` +} + +// AICodeBackfillRequest is the body of POST /api/v1/cc/backfill. +type AICodeBackfillRequest struct { + ClientType string `json:"clientType"` + Events []AICodeBackfillEvent `json:"events"` + CompletedSessionIDs []string `json:"completedSessionIds,omitempty"` + AISummary bool `json:"aiSummary,omitempty"` +} + +type AICodeBackfillSessionStatus struct { + SessionID string `json:"sessionId"` + Status string `json:"status"` + EventCount int `json:"eventCount,omitempty"` +} + +type AICodeBackfillResponse struct { + Success bool `json:"success"` + Accepted int `json:"accepted"` + Summarized int `json:"summarized"` + SkippedSessions []AICodeBackfillSessionStatus `json:"skippedSessions"` +} + +type AICodeBackfillSessionsRequest struct { + ClientType string `json:"clientType"` + SessionIDs []string `json:"sessionIds"` +} + +type AICodeBackfillSessionsResponse struct { + Sessions []AICodeBackfillSessionStatus `json:"sessions"` +} + +type AICodeBackfillCompleteRequest struct { + ClientType string `json:"clientType"` + From time.Time `json:"from"` + To time.Time `json:"to"` +} diff --git a/model/aicode_otel_claude_settings.go b/model/aicode_otel_claude_settings.go index 50369d6..f2c73a4 100644 --- a/model/aicode_otel_claude_settings.go +++ b/model/aicode_otel_claude_settings.go @@ -56,12 +56,19 @@ func claudeSettingsOtelEnvVars() []claudeSettingsEnvVar { } func claudeSettingsResourceAttributes() string { - username := os.Getenv("USER") + username, hostname := aiCodeResourceIdentity() + return fmt.Sprintf("user.name=%s,machine.name=%s,team.id=shelltime", username, hostname) +} + +// aiCodeResourceIdentity returns the user and machine names reported as the +// user.name and machine.name OTEL resource attributes. +func aiCodeResourceIdentity() (userName, machineName string) { + userName = os.Getenv("USER") if u, err := user.Current(); err == nil && u.Username != "" { - username = u.Username + userName = u.Username } - hostname, _ := os.Hostname() - return fmt.Sprintf("user.name=%s,machine.name=%s,team.id=shelltime", username, hostname) + machineName, _ = os.Hostname() + return userName, machineName } func (s *ClaudeSettingsAICodeOtelEnvService) Install() error { diff --git a/model/api.base.go b/model/api.base.go index c7a5208..923d0ea 100644 --- a/model/api.base.go +++ b/model/api.base.go @@ -4,7 +4,6 @@ import ( "bytes" "context" "encoding/json" - "errors" "fmt" "io" "log/slog" @@ -24,6 +23,17 @@ type HTTPRequestOptions[T any, R any] struct { Timeout time.Duration // Optional, defaults to 10 seconds } +// HTTPStatusError is returned by SendHTTPRequestJSON for non-2xx responses, +// so callers can tell retryable server errors from client errors. +type HTTPStatusError struct { + StatusCode int + Message string +} + +func (e *HTTPStatusError) Error() string { + return e.Message +} + // SendHTTPRequestJSON is a generic HTTP request function that sends JSON data and unmarshals the response func SendHTTPRequestJSON[T any, R any](opts HTTPRequestOptions[T, R]) error { ctx, span := modelTracer.Start(opts.Context, "http.send.json") @@ -83,10 +93,10 @@ func SendHTTPRequestJSON[T any, R any](opts HTTPRequestOptions[T, R]) error { err = json.Unmarshal(buf, &msg) if err != nil { slog.Error("Failed to parse error response", slog.Any("err", err)) - return fmt.Errorf("HTTP error: %d", resp.StatusCode) + return &HTTPStatusError{StatusCode: resp.StatusCode, Message: fmt.Sprintf("HTTP error: %d", resp.StatusCode)} } slog.Error("Error response", slog.String("message", msg.ErrorMessage)) - return errors.New(msg.ErrorMessage) + return &HTTPStatusError{StatusCode: resp.StatusCode, Message: msg.ErrorMessage} } // Only try to unmarshal if we have a response struct diff --git a/model/api_aicode_backfill.go b/model/api_aicode_backfill.go new file mode 100644 index 0000000..9cbeea7 --- /dev/null +++ b/model/api_aicode_backfill.go @@ -0,0 +1,71 @@ +package model + +import ( + "context" + "net/http" + "time" +) + +const aiCodeBackfillTimeout = 60 * time.Second + +// FetchAICodeBackfillStatus asks which of the given sessions the server +// already has. Sessions it knows nothing about are not in the result. +// POST /api/v1/cc/backfill/sessions +func FetchAICodeBackfillStatus(ctx context.Context, endpoint Endpoint, clientType string, sessionIDs []string) ([]AICodeBackfillSessionStatus, error) { + ctx, span := modelTracer.Start(ctx, "aicode_backfill.status") + defer span.End() + + var resp AICodeBackfillSessionsResponse + err := SendHTTPRequestJSON(HTTPRequestOptions[AICodeBackfillSessionsRequest, AICodeBackfillSessionsResponse]{ + Context: ctx, + Endpoint: endpoint, + Method: http.MethodPost, + Path: "/api/v1/cc/backfill/sessions", + Payload: AICodeBackfillSessionsRequest{ClientType: clientType, SessionIDs: sessionIDs}, + Response: &resp, + Timeout: aiCodeBackfillTimeout, + }) + if err != nil { + return nil, err + } + return resp.Sessions, nil +} + +// SendAICodeBackfill uploads one batch of historical events. +// POST /api/v1/cc/backfill +func SendAICodeBackfill(ctx context.Context, endpoint Endpoint, req AICodeBackfillRequest) (*AICodeBackfillResponse, error) { + ctx, span := modelTracer.Start(ctx, "aicode_backfill.send") + defer span.End() + + var resp AICodeBackfillResponse + err := SendHTTPRequestJSON(HTTPRequestOptions[AICodeBackfillRequest, AICodeBackfillResponse]{ + Context: ctx, + Endpoint: endpoint, + Method: http.MethodPost, + Path: "/api/v1/cc/backfill", + Payload: req, + Response: &resp, + Timeout: aiCodeBackfillTimeout, + }) + if err != nil { + return nil, err + } + return &resp, nil +} + +// CompleteAICodeBackfill tells the server a backfill finished so it can +// refresh activity data and caches for the backfilled range. +// POST /api/v1/cc/backfill/complete +func CompleteAICodeBackfill(ctx context.Context, endpoint Endpoint, req AICodeBackfillCompleteRequest) error { + ctx, span := modelTracer.Start(ctx, "aicode_backfill.complete") + defer span.End() + + return SendHTTPRequestJSON(HTTPRequestOptions[AICodeBackfillCompleteRequest, struct{}]{ + Context: ctx, + Endpoint: endpoint, + Method: http.MethodPost, + Path: "/api/v1/cc/backfill/complete", + Payload: req, + Timeout: aiCodeBackfillTimeout, + }) +} diff --git a/model/api_aicode_backfill_test.go b/model/api_aicode_backfill_test.go new file mode 100644 index 0000000..d7eb652 --- /dev/null +++ b/model/api_aicode_backfill_test.go @@ -0,0 +1,119 @@ +package model + +import ( + "context" + "errors" + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestFetchAICodeBackfillStatus(t *testing.T) { + var gotPath, gotAuth string + var req AICodeBackfillSessionsRequest + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotPath = r.URL.Path + gotAuth = r.Header.Get("Authorization") + readJSONBody(t, r, &req) + _, _ = w.Write([]byte(`{"sessions":[{"sessionId":"s1","status":"backfilled","eventCount":7}]}`)) + })) + defer server.Close() + + statuses, err := FetchAICodeBackfillStatus(context.Background(), Endpoint{Token: "tok", APIEndpoint: server.URL}, AICodeClientCodex, []string{"s1", "s2"}) + + require.NoError(t, err) + assert.Equal(t, "/api/v1/cc/backfill/sessions", gotPath) + assert.Equal(t, "CLI tok", gotAuth) + assert.Equal(t, AICodeClientCodex, req.ClientType) + assert.Equal(t, []string{"s1", "s2"}, req.SessionIDs) + require.Len(t, statuses, 1) + assert.Equal(t, AICodeBackfillSessionStatus{SessionID: "s1", Status: AICodeBackfillStatusBackfilled, EventCount: 7}, statuses[0]) +} + +func TestSendAICodeBackfill(t *testing.T) { + var gotPath string + var body map[string]any + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotPath = r.URL.Path + readJSONBody(t, r, &body) + _, _ = w.Write([]byte(`{"success":true,"accepted":1,"summarized":1,"skippedSessions":[{"sessionId":"s2","status":"live"}]}`)) + })) + defer server.Close() + + failed := false + resp, err := SendAICodeBackfill(context.Background(), Endpoint{APIEndpoint: server.URL}, AICodeBackfillRequest{ + ClientType: AICodeClientClaudeCode, + Events: []AICodeBackfillEvent{{ + EventID: "bf1:abc", + EventType: AICodeEventToolResult, + SessionID: "s1", + Success: &failed, + }}, + CompletedSessionIDs: []string{"s1"}, + }) + + require.NoError(t, err) + assert.Equal(t, "/api/v1/cc/backfill", gotPath) + assert.Equal(t, 1, resp.Accepted) + assert.Equal(t, []AICodeBackfillSessionStatus{{SessionID: "s2", Status: AICodeBackfillStatusLive}}, resp.SkippedSessions) + + events := body["events"].([]any) + event := events[0].(map[string]any) + // A failed tool call must reach the server as false, not be dropped. + assert.Equal(t, false, event["success"]) + assert.NotContains(t, event, "inputTokens") + assert.Equal(t, []any{"s1"}, body["completedSessionIds"]) + assert.NotContains(t, body, "aiSummary") +} + +func TestCompleteAICodeBackfill(t *testing.T) { + var gotPath string + var req AICodeBackfillCompleteRequest + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotPath = r.URL.Path + readJSONBody(t, r, &req) + _, _ = w.Write([]byte(`{"success":true}`)) + })) + defer server.Close() + + from := time.Date(2026, 9, 1, 0, 0, 0, 0, time.UTC) + to := time.Date(2026, 9, 30, 0, 0, 0, 0, time.UTC) + err := CompleteAICodeBackfill(context.Background(), Endpoint{APIEndpoint: server.URL}, AICodeBackfillCompleteRequest{ClientType: AICodeClientClaudeCode, From: from, To: to}) + + require.NoError(t, err) + assert.Equal(t, "/api/v1/cc/backfill/complete", gotPath) + assert.True(t, from.Equal(req.From)) + assert.True(t, to.Equal(req.To)) +} + +func TestSendHTTPRequestJSON_ReturnsHTTPStatusError(t *testing.T) { + tests := []struct { + name string + status int + body string + wantMsg string + }{ + {name: "server error message", status: http.StatusServiceUnavailable, body: `{"code":503,"error":"pricing unavailable"}`, wantMsg: "pricing unavailable"}, + {name: "unparseable body", status: http.StatusNotFound, body: `404 page not found`, wantMsg: "HTTP error: 404"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(tt.status) + _, _ = w.Write([]byte(tt.body)) + })) + defer server.Close() + + err := CompleteAICodeBackfill(context.Background(), Endpoint{APIEndpoint: server.URL}, AICodeBackfillCompleteRequest{}) + + var statusErr *HTTPStatusError + require.True(t, errors.As(err, &statusErr)) + assert.Equal(t, tt.status, statusErr.StatusCode) + assert.Equal(t, tt.wantMsg, err.Error()) + }) + } +} From 81a67ba4e4175001ada696397ffed5f0870a64a7 Mon Sep 17 00:00:00 2001 From: Claude Date: Mon, 5 Oct 2026 11:31:35 +0000 Subject: [PATCH 2/3] feat(model): parse Claude Code and Codex transcripts for backfill Rebuild historical sessions from local transcripts as backfill events matching what the live OTEL pipeline records. Claude Code (~/.claude/projects, ~/.config/claude/projects or CLAUDE_CONFIG_DIR): one api_request per message id + request id, keeping the most complete usage of a streamed response; human prompts only (no meta, subagent, compaction, local-command or notification lines); tool results paired with their tool calls, keeping only file paths of the parameters. Items copied into resumed sessions are credited once, to the session that started first. Codex (~/.codex/sessions, archived_sessions or CODEX_HOME): conversation_starts, user prompts, response.completed token counts from last_token_usage or the cumulative delta (skipping null and repeated counts), and tool results with exit code and duration. Items copied into forked rollouts are dropped. Batches are also bounded by encoded size, and long prompts and tool arguments are capped, so a batch stays under the server's body limit. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LMCwtubXzBYhQTnX44vhzF --- model/aicode_backfill_claude.go | 531 ++++++++++++++++++++++++ model/aicode_backfill_claude_test.go | 275 ++++++++++++ model/aicode_backfill_codex.go | 598 +++++++++++++++++++++++++++ model/aicode_backfill_codex_test.go | 232 +++++++++++ model/aicode_backfill_common.go | 72 +++- model/aicode_backfill_common_test.go | 21 +- model/aicode_backfill_types.go | 4 +- 7 files changed, 1712 insertions(+), 21 deletions(-) create mode 100644 model/aicode_backfill_claude.go create mode 100644 model/aicode_backfill_claude_test.go create mode 100644 model/aicode_backfill_codex.go create mode 100644 model/aicode_backfill_codex_test.go diff --git a/model/aicode_backfill_claude.go b/model/aicode_backfill_claude.go new file mode 100644 index 0000000..12df285 --- /dev/null +++ b/model/aicode_backfill_claude.go @@ -0,0 +1,531 @@ +package model + +import ( + "encoding/json" + "os" + "path/filepath" + "regexp" + "strings" + "time" + "unicode/utf8" +) + +// claudeSyntheticModel marks assistant lines Claude Code writes itself +// (errors, "no response requested") rather than receiving from the API. +const claudeSyntheticModel = "" + +// ClaudeProjectRoots returns the directories holding Claude Code transcripts: +// the comma-separated CLAUDE_CONFIG_DIR list if set, otherwise both +// $XDG_CONFIG_HOME/claude (default ~/.config/claude) and ~/.claude. +func ClaudeProjectRoots() []string { + var bases []string + if env := os.Getenv("CLAUDE_CONFIG_DIR"); strings.TrimSpace(env) != "" { + bases = splitPathList(env) + } else { + home, _ := os.UserHomeDir() + xdg := os.Getenv("XDG_CONFIG_HOME") + if xdg == "" { + xdg = filepath.Join(home, ".config") + } + bases = []string{filepath.Join(xdg, "claude"), filepath.Join(home, ".claude")} + } + + roots := make([]string, 0, len(bases)) + for _, base := range bases { + roots = append(roots, filepath.Join(base, "projects")) + } + return roots +} + +type claudeTranscriptLine struct { + Type string `json:"type"` + UUID string `json:"uuid"` + SessionID string `json:"sessionId"` + Timestamp string `json:"timestamp"` + Cwd string `json:"cwd"` + Version string `json:"version"` + IsSidechain bool `json:"isSidechain"` + IsMeta bool `json:"isMeta"` + IsCompactSummary bool `json:"isCompactSummary"` + IsAPIErrorMessage bool `json:"isApiErrorMessage"` + RequestID string `json:"requestId"` + CostUSD *float64 `json:"costUSD"` + Origin *claudeOrigin `json:"origin"` + Message *claudeMessage `json:"message"` + ToolUseResult json.RawMessage `json:"toolUseResult"` +} + +type claudeOrigin struct { + Kind string `json:"kind"` +} + +type claudeMessage struct { + ID string `json:"id"` + Model string `json:"model"` + Content json.RawMessage `json:"content"` + Usage *claudeUsage `json:"usage"` +} + +type claudeUsage struct { + InputTokens int `json:"input_tokens"` + OutputTokens int `json:"output_tokens"` + CacheCreationInputTokens int `json:"cache_creation_input_tokens"` + CacheReadInputTokens int `json:"cache_read_input_tokens"` +} + +func (u claudeUsage) total() int { + return u.InputTokens + u.OutputTokens + u.CacheCreationInputTokens + u.CacheReadInputTokens +} + +type claudeContentBlock struct { + Type string `json:"type"` + Text string `json:"text"` + ID string `json:"id"` + Name string `json:"name"` + Input json.RawMessage `json:"input"` + ToolUseID string `json:"tool_use_id"` + IsError bool `json:"is_error"` +} + +type claudeAPIRequest struct { + naturalKey string + sessions []string + ts time.Time + model string + usage claudeUsage + costUSD *float64 + isError bool + errorText string +} + +type claudeToolUse struct { + name string + params map[string]any +} + +type claudeToolResult struct { + sessions []string + ts time.Time + isError bool + durationMs *int +} + +type claudePrompt struct { + sessions []string + ts time.Time + text string +} + +type claudeSessionMeta struct { + start time.Time + cwd string + version string + versionTs time.Time +} + +// claudeParser rebuilds sessions from Claude Code transcripts. Items are +// deduplicated globally: one API response is written as several lines that +// repeat the same usage, and resumed sessions copy earlier history into new +// files. A duplicated item is credited to the earliest session it appears in. +type claudeParser struct { + apiReqs map[backfillKey]*claudeAPIRequest + apiOrder []backfillKey + toolUses map[string]claudeToolUse + toolResults map[string]*claudeToolResult + resultOrder []string + prompts map[string]*claudePrompt + promptOrder []string + sessions map[string]*claudeSessionMeta + badLines int +} + +// ParseClaudeTranscripts reads every Claude Code transcript under roots and +// rebuilds the sessions found in them as backfill events. +func ParseClaudeTranscripts(roots []string, opts BackfillOptions) (*BackfillParseResult, error) { + files, err := collectJSONL(roots, opts.Since) + if err != nil { + return nil, err + } + + p := &claudeParser{ + apiReqs: map[backfillKey]*claudeAPIRequest{}, + toolUses: map[string]claudeToolUse{}, + toolResults: map[string]*claudeToolResult{}, + prompts: map[string]*claudePrompt{}, + sessions: map[string]*claudeSessionMeta{}, + } + for _, file := range files { + if err := readJSONLLines(file, p.handleLine); err != nil { + return nil, err + } + } + + return &BackfillParseResult{ + Sessions: p.build(opts), + Files: len(files), + BadLines: p.badLines, + }, nil +} + +func (p *claudeParser) handleLine(line []byte) { + var l claudeTranscriptLine + if err := json.Unmarshal(line, &l); err != nil { + p.badLines++ + return + } + if l.SessionID == "" || l.Timestamp == "" { + return + } + ts, err := time.Parse(time.RFC3339Nano, l.Timestamp) + if err != nil { + p.badLines++ + return + } + + meta := p.sessions[l.SessionID] + if meta == nil { + meta = &claudeSessionMeta{start: ts} + p.sessions[l.SessionID] = meta + } + if ts.Before(meta.start) { + meta.start = ts + } + if meta.cwd == "" && l.Cwd != "" { + meta.cwd = l.Cwd + } + if l.Version != "" && !ts.Before(meta.versionTs) { + meta.version = l.Version + meta.versionTs = ts + } + + if l.Message == nil { + return + } + switch l.Type { + case "assistant": + p.handleAssistant(&l, ts) + case "user": + p.handleUser(&l, ts) + } +} + +func (p *claudeParser) handleAssistant(l *claudeTranscriptLine, ts time.Time) { + blocks, _ := decodeClaudeContent(l.Message.Content) + for _, b := range blocks { + if b.Type == "tool_use" && b.ID != "" { + if _, ok := p.toolUses[b.ID]; !ok { + p.toolUses[b.ID] = claudeToolUse{name: b.Name, params: claudeToolParams(b.Input)} + } + } + } + + usage := l.Message.Usage + if usage == nil { + return + } + if l.Message.Model == claudeSyntheticModel && !l.IsAPIErrorMessage { + return + } + + naturalKey := claudeRequestKey(l) + key := newBackfillKey(AICodeClientClaudeCode, naturalKey) + if existing, ok := p.apiReqs[key]; ok { + existing.sessions = appendUnique(existing.sessions, l.SessionID) + // Early lines of a streamed response can carry partial usage; keep + // the most complete one. + if usage.total() > existing.usage.total() { + existing.usage = *usage + if l.CostUSD != nil { + existing.costUSD = l.CostUSD + } + } + return + } + + req := &claudeAPIRequest{ + naturalKey: naturalKey, + sessions: []string{l.SessionID}, + ts: ts, + model: l.Message.Model, + usage: *usage, + costUSD: l.CostUSD, + isError: l.IsAPIErrorMessage, + } + if req.isError { + req.errorText = capBackfillText(claudeBlocksText(blocks), backfillMaxErrorBytes) + if req.model == claudeSyntheticModel { + req.model = "" + } + } + p.apiReqs[key] = req + p.apiOrder = append(p.apiOrder, key) +} + +func (p *claudeParser) handleUser(l *claudeTranscriptLine, ts time.Time) { + blocks, text := decodeClaudeContent(l.Message.Content) + + hasToolResult := false + for _, b := range blocks { + if b.Type != "tool_result" || b.ToolUseID == "" { + continue + } + hasToolResult = true + if existing, ok := p.toolResults[b.ToolUseID]; ok { + existing.sessions = appendUnique(existing.sessions, l.SessionID) + continue + } + p.toolResults[b.ToolUseID] = &claudeToolResult{ + sessions: []string{l.SessionID}, + ts: ts, + isError: b.IsError, + durationMs: claudeToolDuration(l.ToolUseResult), + } + p.resultOrder = append(p.resultOrder, b.ToolUseID) + } + if hasToolResult || l.UUID == "" { + return + } + + // Only what a person typed counts as a prompt: not meta lines, subagent + // prompts, compaction summaries or injected notifications. + if l.IsMeta || l.IsSidechain || l.IsCompactSummary { + return + } + if l.Origin != nil && l.Origin.Kind != "human" { + return + } + if blocks != nil { + text = claudeBlocksText(blocks) + } + prompt, ok := normalizeClaudePrompt(text) + if !ok { + return + } + + if existing, ok := p.prompts[l.UUID]; ok { + existing.sessions = appendUnique(existing.sessions, l.SessionID) + return + } + p.prompts[l.UUID] = &claudePrompt{sessions: []string{l.SessionID}, ts: ts, text: prompt} + p.promptOrder = append(p.promptOrder, l.UUID) +} + +// earliestSession picks the session that started first among those an item +// appeared in, so copies in resumed sessions are credited to the original. +func (p *claudeParser) earliestSession(candidates []string) string { + best := candidates[0] + for _, c := range candidates[1:] { + if p.sessions[c].start.Before(p.sessions[best].start) { + best = c + } + } + return best +} + +func (p *claudeParser) build(opts BackfillOptions) []*BackfillSession { + out := map[string]*BackfillSession{} + session := func(id string) *BackfillSession { + s := out[id] + if s == nil { + s = &BackfillSession{ID: id, ClientType: AICodeClientClaudeCode} + out[id] = s + } + return s + } + + for _, key := range p.apiOrder { + r := p.apiReqs[key] + sid := p.earliestSession(r.sessions) + eventType := AICodeEventApiRequest + if r.isError { + eventType = AICodeEventApiError + } + e := newBackfillEvent(AICodeClientClaudeCode, eventType, "api:"+r.naturalKey, sid, r.ts) + e.Model = r.model + if r.isError { + e.Error = r.errorText + } else { + e.InputTokens = intRef(r.usage.InputTokens) + e.OutputTokens = intRef(r.usage.OutputTokens) + e.CacheReadTokens = intRef(r.usage.CacheReadInputTokens) + e.CacheCreationTokens = intRef(r.usage.CacheCreationInputTokens) + e.CostUSD = r.costUSD + } + session(sid).Events = append(session(sid).Events, e) + } + + for _, uuid := range p.promptOrder { + pr := p.prompts[uuid] + sid := p.earliestSession(pr.sessions) + e := newBackfillEvent(AICodeClientClaudeCode, AICodeEventUserPrompt, "prompt:"+uuid, sid, pr.ts) + e.PromptLength = intRef(utf8.RuneCountInString(pr.text)) + if !opts.NoPrompts { + e.Prompt = capBackfillText(pr.text, backfillMaxPromptBytes) + } + session(sid).Events = append(session(sid).Events, e) + } + + for _, id := range p.resultOrder { + res := p.toolResults[id] + use, ok := p.toolUses[id] + if !ok { + continue + } + sid := p.earliestSession(res.sessions) + e := newBackfillEvent(AICodeClientClaudeCode, AICodeEventToolResult, "tool:"+id, sid, res.ts) + e.ToolName = use.name + e.Success = boolRef(!res.isError) + e.DurationMs = res.durationMs + e.ToolParameters = use.params + session(sid).Events = append(session(sid).Events, e) + } + + machine := currentBackfillMachine() + sessions := make([]*BackfillSession, 0, len(out)) + for id, s := range out { + meta := p.sessions[id] + for i := range s.Events { + machine.apply(&s.Events[i]) + s.Events[i].Pwd = meta.cwd + s.Events[i].AppVersion = meta.version + } + finalizeBackfillSession(s) + sessions = append(sessions, s) + } + return sessions +} + +// claudeRequestKey identifies one API response. Claude Code writes a line +// per content block, all sharing message.id and requestId. +func claudeRequestKey(l *claudeTranscriptLine) string { + switch { + case l.Message.ID != "" && l.RequestID != "": + return l.Message.ID + ":" + l.RequestID + case l.Message.ID != "": + return l.Message.ID + default: + return "uuid:" + l.UUID + } +} + +// decodeClaudeContent returns the content blocks of a message, or its text +// when the content is a plain string. +func decodeClaudeContent(raw json.RawMessage) ([]claudeContentBlock, string) { + if len(raw) == 0 { + return nil, "" + } + if raw[0] == '"' { + var text string + if err := json.Unmarshal(raw, &text); err == nil { + return nil, text + } + return nil, "" + } + var blocks []claudeContentBlock + if err := json.Unmarshal(raw, &blocks); err != nil { + return nil, "" + } + if blocks == nil { + blocks = []claudeContentBlock{} + } + return blocks, "" +} + +func claudeBlocksText(blocks []claudeContentBlock) string { + var parts []string + for _, b := range blocks { + if b.Type == "text" && strings.TrimSpace(b.Text) != "" { + parts = append(parts, b.Text) + } + } + return strings.Join(parts, "\n") +} + +var ( + claudeCommandNamePattern = regexp.MustCompile(`([^<]*)`) + claudeCommandArgsPattern = regexp.MustCompile(`(?s)(.*?)`) +) + +// normalizeClaudePrompt turns the text of a user line into the prompt the +// person typed. Output of local commands and interruption markers are not +// prompts; slash commands are reduced to "/name args". +func normalizeClaudePrompt(text string) (string, bool) { + text = strings.TrimSpace(text) + if text == "" { + return "", false + } + for _, prefix := range []string{ + "", + "", + "", + "[Request interrupted", + } { + if strings.HasPrefix(text, prefix) { + return "", false + } + } + if m := claudeCommandNamePattern.FindStringSubmatch(text); m != nil && strings.HasPrefix(text, "ok", nil), + claudeUser("sess-a", "u4", claudeTS(0, 15), []any{text("[Request interrupted by user]")}, nil), + claudeUser("sess-a", "u5", claudeTS(0, 16), "agent finished", obj{"origin": obj{"kind": "task-notification"}}), + obj{ + "type": "assistant", "uuid": "a4", "sessionId": "sess-a", "timestamp": claudeTS(0, 17), + "message": obj{"id": "msg_syn", "model": claudeSyntheticModel, "content": []any{text("No response requested.")}, "usage": usage(0, 0, 0, 0)}, + }, + obj{ + "type": "assistant", "uuid": "a5", "sessionId": "sess-a", "timestamp": claudeTS(0, 18), "isApiErrorMessage": true, + "message": obj{"id": "msg_err", "model": claudeSyntheticModel, "content": []any{text("API Error: 529 overloaded")}, "usage": usage(0, 0, 0, 0)}, + }, + claudeUser("sess-a", "u6", claudeTS(1, 0), + "/review\nreview\nPR 12", + obj{"origin": obj{"kind": "human"}}), + `{"type":"assistant","truncated`, + claudeAssistant("sess-a", "a6", claudeTS(1, 5), "msg_3", "req_3", + []any{obj{"type": "tool_use", "id": "toolu_2", "name": "Bash", "input": obj{"command": "export TOKEN=abc"}}}, usage(5, 5, 0, 2000), nil), + claudeUser("sess-a", "u-r2", claudeTS(1, 9), + []any{obj{"type": "tool_result", "tool_use_id": "toolu_2", "content": "ok"}}, + obj{"toolUseResult": obj{"stdout": "ok", "durationMs": 1234}}), + ) + + // A subagent transcript: lines carry the parent session and isSidechain. + writeJSONL(t, filepath.Join(projects, "-tmp-proj", "sess-a", "subagents", "agent-x.jsonl"), + claudeUser("sess-a", "sub-u1", claudeTS(1, 20), "Explore the repo", obj{"isSidechain": true}), + claudeAssistant("sess-a", "sub-a1", claudeTS(1, 25), "msg_5", "req_5", []any{text("Found it")}, usage(20, 30, 0, 0), obj{"isSidechain": true}), + ) + + // A resumed session that copied earlier history before continuing. + writeJSONL(t, filepath.Join(projects, "-tmp-proj", "sess-b.jsonl"), + claudeAssistant("sess-b", "a3", claudeTS(0, 12), "msg_2", "req_2", []any{text("Done.")}, usage(3, 7, 0, 1100), nil), + claudeUser("sess-b", "u7", claudeTS(30, 0), "Now add tests", obj{"origin": obj{"kind": "human"}}), + claudeAssistant("sess-b", "b1", claudeTS(30, 5), "msg_4", "req_4", []any{text("Added.")}, usage(1, 2, 3, 4), nil), + ) + + // A second config dir with an older transcript format (no origin, string content). + root2 := t.TempDir() + writeJSONL(t, filepath.Join(root2, "projects", "-other", "sess-c.jsonl"), + obj{"type": "user", "uuid": "c1", "sessionId": "sess-c", "timestamp": claudeTS(40, 0), "cwd": "/tmp/other", "version": "1.0.0", + "message": obj{"role": "user", "content": "Old style prompt"}}, + obj{"type": "assistant", "uuid": "c2", "sessionId": "sess-c", "timestamp": claudeTS(40, 5), "costUSD": 0.5, + "message": obj{"id": "msg_c", "model": "claude-3-5-sonnet-20241022", "content": []any{text("Hi")}, "usage": obj{"input_tokens": 7, "output_tokens": 8}}}, + ) + return root, root2 +} + +func TestParseClaudeTranscripts(t *testing.T) { + root, root2 := writeClaudeFixtures(t) + t.Setenv("CLAUDE_CONFIG_DIR", root+", "+root2) + + res, err := ParseClaudeTranscripts(ClaudeProjectRoots(), BackfillOptions{}) + require.NoError(t, err) + + assert.Equal(t, 4, res.Files) + assert.Equal(t, 1, res.BadLines) + require.Len(t, res.Sessions, 3) + + a := sessionByID(t, res, "sess-a") + prompts := eventsOfType(a, AICodeEventUserPrompt) + require.Len(t, prompts, 2) + assert.Equal(t, "Fix the bug", prompts[0].Prompt) + assert.Equal(t, 11, *prompts[0].PromptLength) + assert.Equal(t, "/review PR 12", prompts[1].Prompt) + + requests := eventsOfType(a, AICodeEventApiRequest) + require.Len(t, requests, 4, "msg_1 once, msg_2 credited to the original session, msg_3, subagent msg_5") + assert.Equal(t, 50, *requests[0].OutputTokens, "the most complete usage line wins") + assert.Equal(t, 1000, *requests[0].CacheReadTokens) + assert.Equal(t, 100, *requests[0].CacheCreationTokens) + assert.Nil(t, requests[0].CostUSD) + assert.InDelta(t, 0.01, *requests[1].CostUSD, 1e-9) + assert.Equal(t, 20, *requests[3].InputTokens, "subagent usage belongs to the parent session") + + apiErrors := eventsOfType(a, AICodeEventApiError) + require.Len(t, apiErrors, 1) + assert.Equal(t, "API Error: 529 overloaded", apiErrors[0].Error) + assert.Empty(t, apiErrors[0].Model) + assert.Nil(t, apiErrors[0].InputTokens) + + tools := eventsOfType(a, AICodeEventToolResult) + require.Len(t, tools, 2) + assert.Equal(t, "Edit", tools[0].ToolName) + assert.Equal(t, false, *tools[0].Success) + assert.Nil(t, tools[0].DurationMs) + assert.Equal(t, map[string]any{"file_path": "/tmp/proj/a.go"}, tools[0].ToolParameters) + assert.Equal(t, "Bash", tools[1].ToolName) + assert.Equal(t, true, *tools[1].Success) + assert.Equal(t, 1234, *tools[1].DurationMs) + assert.Nil(t, tools[1].ToolParameters, "commands are never uploaded") + + for _, e := range a.Events { + assert.Equal(t, AICodeClientClaudeCode, e.ClientType) + assert.Equal(t, "sess-a", e.SessionID) + assert.Equal(t, "/tmp/proj", e.Pwd) + assert.Equal(t, "2.1.0", e.AppVersion) + assert.Equal(t, "shelltime", e.TeamID) + assert.Regexp(t, `^bf1:[0-9a-f]{40}$`, e.EventID) + } + assert.Equal(t, time.Date(2026, 9, 10, 10, 0, 0, 0, time.UTC), a.Start.UTC()) + + b := sessionByID(t, res, "sess-b") + assert.Len(t, eventsOfType(b, AICodeEventApiRequest), 1, "copied history is not counted twice") + assert.Len(t, eventsOfType(b, AICodeEventUserPrompt), 1) + + c := sessionByID(t, res, "sess-c") + require.Len(t, c.Events, 2) + assert.Equal(t, "Old style prompt", eventsOfType(c, AICodeEventUserPrompt)[0].Prompt) + cReq := eventsOfType(c, AICodeEventApiRequest)[0] + assert.InDelta(t, 0.5, *cReq.CostUSD, 1e-9) + assert.Equal(t, 0, *cReq.CacheReadTokens) + assert.Equal(t, "/tmp/other", cReq.Pwd) + assert.Equal(t, "1.0.0", cReq.AppVersion) +} + +func TestParseClaudeTranscriptsIDsDoNotDependOnSession(t *testing.T) { + root, root2 := writeClaudeFixtures(t) + res, err := ParseClaudeTranscripts([]string{filepath.Join(root, "projects"), filepath.Join(root2, "projects")}, BackfillOptions{}) + require.NoError(t, err) + + ids := map[string]bool{} + for _, s := range res.Sessions { + for _, e := range s.Events { + assert.False(t, ids[e.EventID], "duplicate event id %s", e.EventID) + ids[e.EventID] = true + } + } + assert.True(t, ids[backfillEventID(AICodeClientClaudeCode, "api:msg_2:req_2")]) + assert.True(t, ids[backfillEventID(AICodeClientClaudeCode, "prompt:u1")]) + assert.True(t, ids[backfillEventID(AICodeClientClaudeCode, "tool:toolu_1")]) +} + +func TestParseClaudeTranscriptsNoPrompts(t *testing.T) { + root, _ := writeClaudeFixtures(t) + res, err := ParseClaudeTranscripts([]string{filepath.Join(root, "projects")}, BackfillOptions{NoPrompts: true}) + require.NoError(t, err) + + prompts := eventsOfType(sessionByID(t, res, "sess-a"), AICodeEventUserPrompt) + require.Len(t, prompts, 2) + assert.Empty(t, prompts[0].Prompt) + assert.Equal(t, 11, *prompts[0].PromptLength) +} + +func TestParseClaudeTranscriptsCapsLongPrompts(t *testing.T) { + root := t.TempDir() + long := strings.Repeat("é", backfillMaxPromptBytes) + writeJSONL(t, filepath.Join(root, "p", "s.jsonl"), + claudeUser("s", "u1", claudeTS(0, 0), long, obj{"origin": obj{"kind": "human"}})) + + res, err := ParseClaudeTranscripts([]string{root}, BackfillOptions{}) + require.NoError(t, err) + + e := res.Sessions[0].Events[0] + assert.LessOrEqual(t, len(e.Prompt), backfillMaxPromptBytes) + assert.True(t, strings.HasPrefix(long, e.Prompt)) + assert.Equal(t, backfillMaxPromptBytes, *e.PromptLength, "the length reports the full prompt") +} + +func TestClaudeProjectRootsDefaults(t *testing.T) { + home := t.TempDir() + t.Setenv("HOME", home) + t.Setenv("CLAUDE_CONFIG_DIR", "") + t.Setenv("XDG_CONFIG_HOME", "") + + assert.Equal(t, []string{ + filepath.Join(home, ".config", "claude", "projects"), + filepath.Join(home, ".claude", "projects"), + }, ClaudeProjectRoots()) +} diff --git a/model/aicode_backfill_codex.go b/model/aicode_backfill_codex.go new file mode 100644 index 0000000..34083ff --- /dev/null +++ b/model/aicode_backfill_codex.go @@ -0,0 +1,598 @@ +package model + +import ( + "bufio" + "encoding/json" + "fmt" + "os" + "path/filepath" + "regexp" + "sort" + "strconv" + "strings" + "time" + "unicode/utf8" +) + +// codexFallbackModel is used for token counts logged before any turn +// context named the model, matching ccusage's choice for old rollouts. +const codexFallbackModel = "gpt-5" + +// codexResponseCompleted is the event kind the live Codex OTEL pipeline +// attaches to the sse_event carrying a response's final token counts; the +// server only sums tokens of Codex events with this kind. +const codexResponseCompleted = "response.completed" + +// CodexSessionRoots returns the directories holding Codex rollouts: the +// sessions and archived_sessions folders of each CODEX_HOME entry (comma +// separated), or of ~/.codex. +func CodexSessionRoots() []string { + homes := splitPathList(os.Getenv("CODEX_HOME")) + if len(homes) == 0 { + home, _ := os.UserHomeDir() + homes = []string{filepath.Join(home, ".codex")} + } + roots := make([]string, 0, len(homes)*2) + for _, h := range homes { + roots = append(roots, filepath.Join(h, "sessions"), filepath.Join(h, "archived_sessions")) + } + return roots +} + +type codexRolloutLine struct { + Timestamp string `json:"timestamp"` + Type string `json:"type"` + Payload json.RawMessage `json:"payload"` +} + +type codexSessionMeta struct { + ID string `json:"id"` + Cwd string `json:"cwd"` + CLIVersion string `json:"cli_version"` + ModelProvider string `json:"model_provider"` +} + +type codexTurnContext struct { + Cwd string `json:"cwd"` + Model string `json:"model"` + Effort string `json:"effort"` + ApprovalPolicy string `json:"approval_policy"` + SandboxPolicy json.RawMessage `json:"sandbox_policy"` +} + +type codexUsage struct { + InputTokens int `json:"input_tokens"` + CachedInputTokens int `json:"cached_input_tokens"` + OutputTokens int `json:"output_tokens"` + ReasoningOutputTokens int `json:"reasoning_output_tokens"` + TotalTokens int `json:"total_tokens"` +} + +func (u codexUsage) isZero() bool { + return u.InputTokens == 0 && u.CachedInputTokens == 0 && u.OutputTokens == 0 && u.ReasoningOutputTokens == 0 +} + +// minus returns u - prev with every field clamped at zero. +func (u codexUsage) minus(prev codexUsage) codexUsage { + return codexUsage{ + InputTokens: max(u.InputTokens-prev.InputTokens, 0), + CachedInputTokens: max(u.CachedInputTokens-prev.CachedInputTokens, 0), + OutputTokens: max(u.OutputTokens-prev.OutputTokens, 0), + ReasoningOutputTokens: max(u.ReasoningOutputTokens-prev.ReasoningOutputTokens, 0), + TotalTokens: max(u.TotalTokens-prev.TotalTokens, 0), + } +} + +type codexEventMsg struct { + Type string `json:"type"` + Message string `json:"message"` + Kind string `json:"kind"` + Info *codexTokenInfo `json:"info"` + Model string `json:"model"` +} + +type codexTokenInfo struct { + TotalTokenUsage *codexUsage `json:"total_token_usage"` + LastTokenUsage *codexUsage `json:"last_token_usage"` + Model string `json:"model"` +} + +type codexResponseItem struct { + Type string `json:"type"` + Role string `json:"role"` + Name string `json:"name"` + Arguments string `json:"arguments"` + Input string `json:"input"` + Action json.RawMessage `json:"action"` + CallID string `json:"call_id"` + Output json.RawMessage `json:"output"` + Content []struct { + Type string `json:"type"` + Text string `json:"text"` + } `json:"content"` +} + +type codexPendingCall struct { + name string + args map[string]any + ts time.Time +} + +// codexFileState is the per-rollout state needed to turn its lines into +// events: the session, the latest turn context, the previous cumulative +// token usage and tool calls waiting for their output. +type codexFileState struct { + path string + sessionID string + cwd string + cliVersion string + provider string + metaTs time.Time + model string + effort string + approval string + sandbox string + turnSeen bool + firstTurn codexTurnContext + prevTotal *codexUsage + pending map[string]codexPendingCall + events []AICodeBackfillEvent + userEvents int + userFallback []codexPromptCandidate +} + +type codexPromptCandidate struct { + text string + rawTs string + ts time.Time +} + +type codexParser struct { + opts BackfillOptions + seen map[backfillKey]struct{} + sessions map[string]*BackfillSession + meta map[string]*codexFileState + badLines int +} + +var codexRolloutUUIDPattern = regexp.MustCompile(`([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})\.jsonl$`) + +// ParseCodexRollouts reads every Codex rollout under roots and rebuilds the +// sessions found in them as backfill events. +func ParseCodexRollouts(roots []string, opts BackfillOptions) (*BackfillParseResult, error) { + files, err := collectJSONL(roots, opts.Since) + if err != nil { + return nil, err + } + // Forked rollouts copy their parent's history. Reading files in the + // order they started credits copied items to the original session. + starts := make(map[string]string, len(files)) + for _, f := range files { + starts[f] = codexFirstTimestamp(f) + } + sort.SliceStable(files, func(i, j int) bool { return starts[files[i]] < starts[files[j]] }) + + p := &codexParser{ + opts: opts, + seen: map[backfillKey]struct{}{}, + sessions: map[string]*BackfillSession{}, + meta: map[string]*codexFileState{}, + } + for _, f := range files { + if err := p.parseFile(f); err != nil { + return nil, err + } + } + + machine := currentBackfillMachine() + sessions := make([]*BackfillSession, 0, len(p.sessions)) + for id, s := range p.sessions { + meta := p.meta[id] + for i := range s.Events { + machine.apply(&s.Events[i]) + s.Events[i].ConversationID = id + s.Events[i].Pwd = meta.cwd + s.Events[i].AppVersion = meta.cliVersion + } + finalizeBackfillSession(s) + if len(s.Events) > 0 { + sessions = append(sessions, s) + } + } + return &BackfillParseResult{Sessions: sessions, Files: len(files), BadLines: p.badLines}, nil +} + +// codexFirstTimestamp returns the timestamp of a rollout's first line, or +// "" when it can't be read. +func codexFirstTimestamp(path string) string { + f, err := os.Open(path) + if err != nil { + return "" + } + defer f.Close() + line, _ := bufio.NewReader(f).ReadBytes('\n') + var l codexRolloutLine + if json.Unmarshal(line, &l) != nil { + return "" + } + return l.Timestamp +} + +func (p *codexParser) parseFile(path string) error { + st := &codexFileState{path: path, pending: map[string]codexPendingCall{}} + if m := codexRolloutUUIDPattern.FindStringSubmatch(filepath.Base(path)); m != nil { + st.sessionID = m[1] + } + + err := readJSONLLines(path, func(line []byte) { + var l codexRolloutLine + if err := json.Unmarshal(line, &l); err != nil { + p.badLines++ + return + } + if l.Timestamp == "" || len(l.Payload) == 0 { + return + } + ts, err := time.Parse(time.RFC3339Nano, l.Timestamp) + if err != nil { + p.badLines++ + return + } + switch l.Type { + case "session_meta": + p.handleSessionMeta(st, l.Payload, ts) + case "turn_context": + p.handleTurnContext(st, l.Payload) + case "event_msg": + p.handleEventMsg(st, l.Payload, l.Timestamp, ts) + case "response_item": + p.handleResponseItem(st, l.Payload, l.Timestamp, ts) + } + }) + if err != nil { + return err + } + if st.sessionID == "" { + return nil + } + + // Old rollouts don't log user_message events; fall back to the user + // messages of the conversation itself. + if st.userEvents == 0 { + for _, c := range st.userFallback { + p.addPrompt(st, c.text, c.rawTs, c.ts) + } + } + if !st.metaTs.IsZero() { + p.addConversationStart(st) + } + + if _, ok := p.meta[st.sessionID]; !ok { + p.meta[st.sessionID] = st + } + if len(st.events) == 0 { + return nil + } + s := p.sessions[st.sessionID] + if s == nil { + s = &BackfillSession{ID: st.sessionID, ClientType: AICodeClientCodex} + p.sessions[st.sessionID] = s + } + s.Events = append(s.Events, st.events...) + return nil +} + +// keep reports whether an item is new, recording it as seen. Copies of the +// same item in forked rollouts are dropped. +func (p *codexParser) keep(naturalKey string) bool { + k := newBackfillKey(AICodeClientCodex, naturalKey) + if _, ok := p.seen[k]; ok { + return false + } + p.seen[k] = struct{}{} + return true +} + +func (p *codexParser) handleSessionMeta(st *codexFileState, raw json.RawMessage, ts time.Time) { + if !st.metaTs.IsZero() { + // Forked rollouts can repeat their parent's metadata; the first + // entry describes this file. + return + } + var meta codexSessionMeta + if err := json.Unmarshal(raw, &meta); err != nil { + p.badLines++ + return + } + if meta.ID != "" { + st.sessionID = meta.ID + } + st.metaTs = ts + st.cwd = meta.Cwd + st.cliVersion = meta.CLIVersion + st.provider = meta.ModelProvider +} + +func (p *codexParser) handleTurnContext(st *codexFileState, raw json.RawMessage) { + var tc codexTurnContext + if err := json.Unmarshal(raw, &tc); err != nil { + p.badLines++ + return + } + if tc.Model != "" { + st.model = tc.Model + } + st.effort = tc.Effort + st.approval = tc.ApprovalPolicy + st.sandbox = codexSandboxMode(tc.SandboxPolicy) + if st.cwd == "" { + st.cwd = tc.Cwd + } + if !st.turnSeen { + st.turnSeen = true + st.firstTurn = tc + } +} + +func (p *codexParser) handleEventMsg(st *codexFileState, raw json.RawMessage, rawTs string, ts time.Time) { + var msg codexEventMsg + if err := json.Unmarshal(raw, &msg); err != nil { + p.badLines++ + return + } + switch msg.Type { + case "user_message": + st.userEvents++ + if msg.Kind != "" && msg.Kind != "plain" { + return + } + p.addPrompt(st, msg.Message, rawTs, ts) + case "token_count": + p.handleTokenCount(st, &msg, rawTs, ts) + } +} + +func (p *codexParser) addPrompt(st *codexFileState, text, rawTs string, ts time.Time) { + text = strings.TrimSpace(text) + if text == "" || codexIsContextMessage(text) { + return + } + naturalKey := "prompt:" + rawTs + "|" + text + if !p.keep(naturalKey) { + return + } + e := newBackfillEvent(AICodeClientCodex, AICodeEventUserPrompt, naturalKey, st.sessionID, ts) + e.PromptLength = intRef(utf8.RuneCountInString(text)) + if !p.opts.NoPrompts { + e.Prompt = capBackfillText(text, backfillMaxPromptBytes) + } + st.events = append(st.events, e) +} + +func (p *codexParser) handleTokenCount(st *codexFileState, msg *codexEventMsg, rawTs string, ts time.Time) { + info := msg.Info + if info == nil { + return + } + total := info.TotalTokenUsage + if total != nil && st.prevTotal != nil && *total == *st.prevTotal { + // Codex re-emits the last count (e.g. alongside rate limit updates). + return + } + + var usage codexUsage + switch { + case info.LastTokenUsage != nil: + usage = *info.LastTokenUsage + case total != nil && st.prevTotal != nil: + usage = total.minus(*st.prevTotal) + case total != nil: + usage = *total + } + if total != nil { + prev := *total + st.prevTotal = &prev + } + if usage.isZero() { + return + } + + model := st.model + if model == "" { + model = info.Model + } + if model == "" { + model = msg.Model + } + if model == "" { + model = codexFallbackModel + } + + // The model is left out of the key: a forked rollout may copy a count + // before the turn context that names its model. + naturalKey := fmt.Sprintf("tokens:%s|%d|%d|%d|%d|%d", rawTs, + usage.InputTokens, usage.CachedInputTokens, usage.OutputTokens, usage.ReasoningOutputTokens, usage.TotalTokens) + if !p.keep(naturalKey) { + return + } + + e := newBackfillEvent(AICodeClientCodex, AICodeEventSSEEvent, naturalKey, st.sessionID, ts) + e.EventKind = codexResponseCompleted + e.Model = model + // Same semantics as the live pipeline: input includes cached tokens. + e.InputTokens = intRef(usage.InputTokens) + e.CacheReadTokens = intRef(usage.CachedInputTokens) + e.OutputTokens = intRef(usage.OutputTokens) + e.ReasoningTokens = intRef(usage.ReasoningOutputTokens) + st.events = append(st.events, e) +} + +func (p *codexParser) handleResponseItem(st *codexFileState, raw json.RawMessage, rawTs string, ts time.Time) { + var item codexResponseItem + if err := json.Unmarshal(raw, &item); err != nil { + p.badLines++ + return + } + switch item.Type { + case "message": + if item.Role != "user" { + return + } + var parts []string + for _, c := range item.Content { + if c.Type == "input_text" && strings.TrimSpace(c.Text) != "" { + parts = append(parts, c.Text) + } + } + st.userFallback = append(st.userFallback, codexPromptCandidate{text: strings.Join(parts, "\n"), rawTs: rawTs, ts: ts}) + case "function_call": + st.pending[item.CallID] = codexPendingCall{name: item.Name, args: codexToolArguments(item.Arguments), ts: ts} + case "custom_tool_call": + st.pending[item.CallID] = codexPendingCall{name: item.Name, args: codexCapArgs(map[string]any{"input": item.Input}), ts: ts} + case "local_shell_call": + var action map[string]any + _ = json.Unmarshal(item.Action, &action) + st.pending[item.CallID] = codexPendingCall{name: "local_shell", args: codexCapArgs(action), ts: ts} + case "function_call_output", "custom_tool_call_output": + call, ok := st.pending[item.CallID] + if !ok || item.CallID == "" { + return + } + delete(st.pending, item.CallID) + naturalKey := "tool:" + item.CallID + if !p.keep(naturalKey) { + return + } + e := newBackfillEvent(AICodeClientCodex, AICodeEventToolResult, naturalKey, st.sessionID, ts) + e.ToolName = call.name + e.CallID = item.CallID + e.ToolArguments = call.args + e.Success, e.DurationMs = codexToolOutcome(item.Output) + st.events = append(st.events, e) + } +} + +func (p *codexParser) addConversationStart(st *codexFileState) { + naturalKey := "conv:" + st.sessionID + if !p.keep(naturalKey) { + return + } + e := newBackfillEvent(AICodeClientCodex, AICodeEventConversationStarts, naturalKey, st.sessionID, st.metaTs) + e.Provider = st.provider + if st.turnSeen { + e.Model = st.firstTurn.Model + e.ReasoningEffort = st.firstTurn.Effort + e.ApprovalPolicy = st.firstTurn.ApprovalPolicy + e.SandboxPolicy = codexSandboxMode(st.firstTurn.SandboxPolicy) + } + st.events = append(st.events, e) +} + +// codexIsContextMessage reports whether a user message is context Codex +// injects (instructions, environment) rather than something typed. +func codexIsContextMessage(text string) bool { + for _, prefix := range []string{"", "", "", "# AGENTS.md instructions"} { + if strings.HasPrefix(text, prefix) { + return true + } + } + return false +} + +// codexSandboxMode reads the sandbox policy, which is a plain string in some +// Codex versions and an object with a mode in others. +func codexSandboxMode(raw json.RawMessage) string { + if len(raw) == 0 { + return "" + } + var s string + if json.Unmarshal(raw, &s) == nil { + return s + } + var obj struct { + Mode string `json:"mode"` + Type string `json:"type"` + } + if json.Unmarshal(raw, &obj) == nil { + if obj.Mode != "" { + return obj.Mode + } + return obj.Type + } + return "" +} + +func codexToolArguments(arguments string) map[string]any { + if arguments == "" { + return nil + } + var args map[string]any + if err := json.Unmarshal([]byte(arguments), &args); err != nil { + return nil + } + return codexCapArgs(args) +} + +// codexCapArgs drops tool arguments too large to upload, such as whole +// patches, keeping batches small. +func codexCapArgs(args map[string]any) map[string]any { + if len(args) == 0 { + return nil + } + buf, err := json.Marshal(args) + if err != nil || len(buf) > backfillMaxToolArgsBytes { + return nil + } + return args +} + +var ( + codexExitCodePattern = regexp.MustCompile(`(?m)^Exit code: (-?\d+)`) + codexWallTimePattern = regexp.MustCompile(`(?m)^Wall time: ([0-9.]+) seconds`) +) + +// codexToolOutcome reads success and duration from a tool output. Shell +// tools report them either as JSON metadata or as "Exit code:" / "Wall +// time:" lines; other tools report neither, and both stay nil. +func codexToolOutcome(raw json.RawMessage) (*bool, *int) { + if len(raw) == 0 { + return nil, nil + } + var text string + if json.Unmarshal(raw, &text) != nil { + // Some versions store the output as an object. + text = string(raw) + } + + var structured struct { + Metadata *struct { + ExitCode *int `json:"exit_code"` + DurationSeconds *float64 `json:"duration_seconds"` + } `json:"metadata"` + } + if json.Unmarshal([]byte(text), &structured) == nil && structured.Metadata != nil { + var success *bool + var duration *int + if structured.Metadata.ExitCode != nil { + success = boolRef(*structured.Metadata.ExitCode == 0) + } + if structured.Metadata.DurationSeconds != nil { + duration = intRef(int(*structured.Metadata.DurationSeconds * 1000)) + } + return success, duration + } + + var success *bool + var duration *int + if m := codexExitCodePattern.FindStringSubmatch(text); m != nil { + if code, err := strconv.Atoi(m[1]); err == nil { + success = boolRef(code == 0) + } + } + if m := codexWallTimePattern.FindStringSubmatch(text); m != nil { + if secs, err := strconv.ParseFloat(m[1], 64); err == nil { + duration = intRef(int(secs * 1000)) + } + } + return success, duration +} diff --git a/model/aicode_backfill_codex_test.go b/model/aicode_backfill_codex_test.go new file mode 100644 index 0000000..7ed00d8 --- /dev/null +++ b/model/aicode_backfill_codex_test.go @@ -0,0 +1,232 @@ +package model + +import ( + "encoding/json" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +const ( + codexSessionA = "11111111-1111-4111-8111-111111111111" + codexSessionB = "22222222-2222-4222-8222-222222222222" + codexSessionC = "33333333-3333-4333-8333-333333333333" +) + +func codexTS(hour, min, sec int) string { + return time.Date(2026, 9, 10, hour, min, sec, 0, time.UTC).Format("2006-01-02T15:04:05.000Z") +} + +func codexLine(ts, typ string, payload any) obj { + return obj{"timestamp": ts, "type": typ, "payload": payload} +} + +func codexTokens(ts string, total, last obj) obj { + info := obj{"model_context_window": 272000} + if total != nil { + info["total_token_usage"] = total + } + if last != nil { + info["last_token_usage"] = last + } + return codexLine(ts, "event_msg", obj{"type": "token_count", "info": info}) +} + +func codexUsageObj(in, cached, out, reasoning, total int) obj { + return obj{"input_tokens": in, "cached_input_tokens": cached, "output_tokens": out, "reasoning_output_tokens": reasoning, "total_tokens": total} +} + +func codexUserItem(ts, text string) obj { + return codexLine(ts, "response_item", obj{"type": "message", "role": "user", "content": []any{obj{"type": "input_text", "text": text}}}) +} + +func codexOutput(ts, callID, output string) obj { + return codexLine(ts, "response_item", obj{"type": "function_call_output", "call_id": callID, "output": output}) +} + +func mustJSON(t *testing.T, v any) string { + buf, err := json.Marshal(v) + require.NoError(t, err) + return string(buf) +} + +func writeCodexFixtures(t *testing.T) string { + home := t.TempDir() + day := filepath.Join(home, "sessions", "2026", "09", "10") + first := codexUsageObj(1000, 800, 100, 40, 1100) + shellArgs := mustJSON(t, obj{"command": []any{"bash", "-lc", "go test ./..."}, "workdir": "/tmp/codex"}) + + writeJSONL(t, filepath.Join(day, "rollout-2026-09-10T10-00-00-"+codexSessionA+".jsonl"), + codexLine(codexTS(10, 0, 0), "session_meta", obj{"id": codexSessionA, "cwd": "/tmp/codex", "cli_version": "0.40.0", "model_provider": "openai", "originator": "codex_cli_rs"}), + codexUserItem(codexTS(10, 0, 1), "\n /tmp/codex\n"), + codexUserItem(codexTS(10, 0, 1), "fix the tests"), + codexLine(codexTS(10, 0, 1), "event_msg", obj{"type": "user_message", "message": "fix the tests", "images": []any{}}), + codexLine(codexTS(10, 0, 2), "turn_context", obj{"cwd": "/tmp/codex", "model": "gpt-5-codex", "effort": "medium", "approval_policy": "on-request", "sandbox_policy": obj{"mode": "workspace-write"}}), + codexLine(codexTS(10, 0, 3), "event_msg", obj{"type": "token_count", "info": nil}), + codexLine(codexTS(10, 0, 4), "response_item", obj{"type": "function_call", "name": "shell", "arguments": shellArgs, "call_id": "call_1"}), + codexTokens(codexTS(10, 0, 5), first, first), + // Re-emitted count with unchanged totals. + codexTokens(codexTS(10, 0, 6), first, first), + codexOutput(codexTS(10, 0, 7), "call_1", mustJSON(t, obj{"output": "FAIL", "metadata": obj{"exit_code": 1, "duration_seconds": 2.5}})), + codexLine(codexTS(10, 0, 8), "response_item", obj{"type": "custom_tool_call", "name": "apply_patch", "call_id": "call_2", "input": "*** Begin Patch\n" + strings.Repeat("+line\n", 1000)}), + codexLine(codexTS(10, 0, 9), "response_item", obj{"type": "custom_tool_call_output", "call_id": "call_2", "output": "Exit code: 0\nWall time: 0.3 seconds\nOutput:\nSuccess"}), + // Only the cumulative total: the turn's usage is the difference. + codexTokens(codexTS(10, 0, 10), codexUsageObj(2500, 1800, 300, 90, 2800), nil), + codexLine(codexTS(10, 0, 11), "response_item", obj{"type": "function_call", "name": "docs__search", "arguments": `{"q":"retry"}`, "call_id": "call_3"}), + codexOutput(codexTS(10, 0, 12), "call_3", "found 3 results"), + `{"timestamp":"broken`, + ) + + // A fork of session A copies its history before continuing. + writeJSONL(t, filepath.Join(day, "rollout-2026-09-10T11-00-00-"+codexSessionB+".jsonl"), + codexLine(codexTS(11, 0, 0), "session_meta", obj{"id": codexSessionB, "cwd": "/tmp/codex", "cli_version": "0.41.0", "model_provider": "openai"}), + codexLine(codexTS(10, 0, 1), "event_msg", obj{"type": "user_message", "message": "fix the tests"}), + codexTokens(codexTS(10, 0, 5), first, first), + codexLine(codexTS(11, 0, 1), "turn_context", obj{"cwd": "/tmp/codex", "model": "gpt-5", "approval_policy": "never", "sandbox_policy": "read-only"}), + codexLine(codexTS(11, 0, 5), "event_msg", obj{"type": "user_message", "message": "and run the linter"}), + codexTokens(codexTS(11, 0, 6), codexUsageObj(3000, 900, 350, 90, 3350), codexUsageObj(500, 100, 50, 0, 550)), + ) + + // An old rollout: no session id in the metadata, no user_message events + // and no turn context. + writeJSONL(t, filepath.Join(home, "archived_sessions", "rollout-2025-08-01T09-00-00-"+codexSessionC+".jsonl"), + codexLine("2025-08-01T09:00:00.000Z", "session_meta", obj{"cwd": "/tmp/old", "cli_version": "0.20.0"}), + codexUserItem("2025-08-01T09:00:01.000Z", "old prompt"), + codexTokens("2025-08-01T09:00:02.000Z", nil, codexUsageObj(10, 0, 5, 0, 15)), + ) + return home +} + +func TestParseCodexRollouts(t *testing.T) { + home := writeCodexFixtures(t) + t.Setenv("CODEX_HOME", home) + + res, err := ParseCodexRollouts(CodexSessionRoots(), BackfillOptions{}) + require.NoError(t, err) + + assert.Equal(t, 3, res.Files) + assert.Equal(t, 1, res.BadLines) + require.Len(t, res.Sessions, 3) + + a := sessionByID(t, res, codexSessionA) + starts := eventsOfType(a, AICodeEventConversationStarts) + require.Len(t, starts, 1) + assert.Equal(t, "gpt-5-codex", starts[0].Model) + assert.Equal(t, "openai", starts[0].Provider) + assert.Equal(t, "medium", starts[0].ReasoningEffort) + assert.Equal(t, "on-request", starts[0].ApprovalPolicy) + assert.Equal(t, "workspace-write", starts[0].SandboxPolicy) + + prompts := eventsOfType(a, AICodeEventUserPrompt) + require.Len(t, prompts, 1, "environment context and the duplicated response item are not prompts") + assert.Equal(t, "fix the tests", prompts[0].Prompt) + assert.Equal(t, 13, *prompts[0].PromptLength) + + tokens := eventsOfType(a, AICodeEventSSEEvent) + require.Len(t, tokens, 2, "null info and re-emitted totals are skipped") + for _, e := range tokens { + assert.Equal(t, codexResponseCompleted, e.EventKind) + assert.Equal(t, "gpt-5-codex", e.Model) + } + assert.Equal(t, []int{1000, 800, 100, 40}, []int{*tokens[0].InputTokens, *tokens[0].CacheReadTokens, *tokens[0].OutputTokens, *tokens[0].ReasoningTokens}) + assert.Equal(t, []int{1500, 1000, 200, 50}, []int{*tokens[1].InputTokens, *tokens[1].CacheReadTokens, *tokens[1].OutputTokens, *tokens[1].ReasoningTokens}) + + tools := eventsOfType(a, AICodeEventToolResult) + require.Len(t, tools, 3) + assert.Equal(t, "shell", tools[0].ToolName) + assert.Equal(t, "call_1", tools[0].CallID) + assert.Equal(t, false, *tools[0].Success) + assert.Equal(t, 2500, *tools[0].DurationMs) + assert.Equal(t, "/tmp/codex", tools[0].ToolArguments["workdir"]) + assert.Equal(t, "apply_patch", tools[1].ToolName) + assert.Equal(t, true, *tools[1].Success) + assert.Equal(t, 300, *tools[1].DurationMs) + assert.Nil(t, tools[1].ToolArguments, "large patches are not uploaded") + assert.Equal(t, "docs__search", tools[2].ToolName) + assert.Nil(t, tools[2].Success) + assert.Nil(t, tools[2].DurationMs) + + for _, e := range a.Events { + assert.Equal(t, AICodeClientCodex, e.ClientType) + assert.Equal(t, codexSessionA, e.SessionID) + assert.Equal(t, codexSessionA, e.ConversationID) + assert.Equal(t, "/tmp/codex", e.Pwd) + assert.Equal(t, "0.40.0", e.AppVersion) + } + + b := sessionByID(t, res, codexSessionB) + bPrompts := eventsOfType(b, AICodeEventUserPrompt) + require.Len(t, bPrompts, 1, "copied history belongs to the original session") + assert.Equal(t, "and run the linter", bPrompts[0].Prompt) + bTokens := eventsOfType(b, AICodeEventSSEEvent) + require.Len(t, bTokens, 1) + assert.Equal(t, 500, *bTokens[0].InputTokens) + assert.Equal(t, "gpt-5", bTokens[0].Model) + assert.Equal(t, "read-only", eventsOfType(b, AICodeEventConversationStarts)[0].SandboxPolicy) + + c := sessionByID(t, res, codexSessionC) + cPrompts := eventsOfType(c, AICodeEventUserPrompt) + require.Len(t, cPrompts, 1) + assert.Equal(t, "old prompt", cPrompts[0].Prompt) + cTokens := eventsOfType(c, AICodeEventSSEEvent) + require.Len(t, cTokens, 1) + assert.Equal(t, codexFallbackModel, cTokens[0].Model) + assert.Equal(t, "/tmp/old", cTokens[0].Pwd) +} + +func TestParseCodexRolloutsIsDeterministic(t *testing.T) { + home := writeCodexFixtures(t) + roots := []string{filepath.Join(home, "sessions"), filepath.Join(home, "archived_sessions")} + + first, err := ParseCodexRollouts(roots, BackfillOptions{}) + require.NoError(t, err) + second, err := ParseCodexRollouts(roots, BackfillOptions{}) + require.NoError(t, err) + + ids := func(res *BackfillParseResult) map[string]bool { + out := map[string]bool{} + for _, s := range res.Sessions { + for _, e := range s.Events { + out[e.EventID] = true + } + } + return out + } + assert.Equal(t, ids(first), ids(second)) +} + +func TestCodexToolOutcome(t *testing.T) { + tests := []struct { + name string + output string + wantSuccess *bool + wantDuration *int + }{ + {name: "json metadata", output: `"{\"output\":\"ok\",\"metadata\":{\"exit_code\":0,\"duration_seconds\":0.25}}"`, wantSuccess: boolRef(true), wantDuration: intRef(250)}, + {name: "text form", output: `"Exit code: 2\nWall time: 1.5 seconds\nOutput:\nboom"`, wantSuccess: boolRef(false), wantDuration: intRef(1500)}, + {name: "plain output", output: `"hello"`}, + {name: "object output", output: `{"metadata":{"exit_code":0}}`, wantSuccess: boolRef(true)}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + success, duration := codexToolOutcome(json.RawMessage(tt.output)) + assert.Equal(t, tt.wantSuccess, success) + assert.Equal(t, tt.wantDuration, duration) + }) + } +} + +func TestCodexSessionRootsDefault(t *testing.T) { + home := t.TempDir() + t.Setenv("HOME", home) + t.Setenv("CODEX_HOME", "") + + assert.Equal(t, []string{ + filepath.Join(home, ".codex", "sessions"), + filepath.Join(home, ".codex", "archived_sessions"), + }, CodexSessionRoots()) +} diff --git a/model/aicode_backfill_common.go b/model/aicode_backfill_common.go index eaaf4b0..69ed6f2 100644 --- a/model/aicode_backfill_common.go +++ b/model/aicode_backfill_common.go @@ -5,6 +5,7 @@ import ( "bytes" "crypto/sha256" "encoding/hex" + "encoding/json" "errors" "io" "io/fs" @@ -14,6 +15,15 @@ import ( "sort" "strings" "time" + "unicode/utf8" +) + +// Text limits for backfilled events. They keep a 500-event batch well under +// the server's body limit even when prompts contain pasted files. +const ( + backfillMaxPromptBytes = 16 << 10 + backfillMaxErrorBytes = 1 << 10 + backfillMaxToolArgsBytes = 2 << 10 ) // backfillSettleWindow keeps sessions that were active very recently out of @@ -227,12 +237,13 @@ func SelectBackfillSessions(sessions []*BackfillSession, opts BackfillOptions) ( return kept, active } -// PackBackfillBatches groups sessions into upload requests of at most -// maxEvents events and maxCompleted completed sessions. Whole sessions are -// kept together; only a session larger than maxEvents is split, and a -// session is listed in completedSessionIds only on the request carrying its -// last event, so the server builds its summary once all events are stored. -func PackBackfillBatches(clientType string, sessions []*BackfillSession, maxEvents, maxCompleted int) []AICodeBackfillRequest { +// PackBackfillBatches groups sessions into upload requests bounded by +// maxEvents events, maxBytes of encoded events and maxCompleted completed +// sessions. Whole sessions are kept together when they fit; a larger session +// is split, and a session is listed in completedSessionIds only on the +// request carrying its last event, so the server builds its summary once all +// of its events are stored. +func PackBackfillBatches(clientType string, sessions []*BackfillSession, maxEvents, maxBytes, maxCompleted int) []AICodeBackfillRequest { ordered := make([]*BackfillSession, len(sessions)) copy(ordered, sessions) sort.SliceStable(ordered, func(i, j int) bool { @@ -244,31 +255,35 @@ func PackBackfillBatches(clientType string, sessions []*BackfillSession, maxEven var batches []AICodeBackfillRequest cur := AICodeBackfillRequest{ClientType: clientType} + curBytes := 0 flush := func() { if len(cur.Events) > 0 || len(cur.CompletedSessionIDs) > 0 { batches = append(batches, cur) } cur = AICodeBackfillRequest{ClientType: clientType} + curBytes = 0 } for _, s := range ordered { - events := s.Events + sizes := make([]int, len(s.Events)) + sessionBytes := 0 + for i := range s.Events { + sizes[i] = backfillEventSize(&s.Events[i]) + sessionBytes += sizes[i] + } + + // Start a fresh request rather than split a session that would fit + // into one on its own. if len(cur.CompletedSessionIDs) >= maxCompleted || - (len(cur.Events) > 0 && len(cur.Events)+len(events) > maxEvents) { + (len(cur.Events) > 0 && (len(cur.Events)+len(s.Events) > maxEvents || curBytes+sessionBytes > maxBytes)) { flush() } - for len(events) > 0 { - room := maxEvents - len(cur.Events) - if room <= 0 { - flush() - room = maxEvents - } - n := min(room, len(events)) - cur.Events = append(cur.Events, events[:n]...) - events = events[n:] - if len(events) > 0 { + for i, e := range s.Events { + if len(cur.Events) > 0 && (len(cur.Events) >= maxEvents || curBytes+sizes[i] > maxBytes) { flush() } + cur.Events = append(cur.Events, e) + curBytes += sizes[i] } cur.CompletedSessionIDs = append(cur.CompletedSessionIDs, s.ID) } @@ -276,6 +291,14 @@ func PackBackfillBatches(clientType string, sessions []*BackfillSession, maxEven return batches } +func backfillEventSize(e *AICodeBackfillEvent) int { + buf, err := json.Marshal(e) + if err != nil { + return 0 + } + return len(buf) + 1 +} + // finalizeBackfillSession sorts a session's events and sets its time range. func finalizeBackfillSession(s *BackfillSession) { sort.SliceStable(s.Events, func(i, j int) bool { @@ -294,6 +317,19 @@ func finalizeBackfillSession(s *BackfillSession) { } } +// capBackfillText truncates s to at most maxBytes without splitting a UTF-8 +// character. +func capBackfillText(s string, maxBytes int) string { + if len(s) <= maxBytes { + return s + } + cut := maxBytes + for cut > 0 && !utf8.RuneStart(s[cut]) { + cut-- + } + return s[:cut] +} + func intRef(v int) *int { return &v } func boolRef(v bool) *bool { return &v } diff --git a/model/aicode_backfill_common_test.go b/model/aicode_backfill_common_test.go index 0422449..1bf7f23 100644 --- a/model/aicode_backfill_common_test.go +++ b/model/aicode_backfill_common_test.go @@ -39,7 +39,7 @@ func TestPackBackfillBatches(t *testing.T) { testBackfillSession("big", base.Add(time.Hour), 12), } - batches := PackBackfillBatches(AICodeClientClaudeCode, sessions, 5, 50) + batches := PackBackfillBatches(AICodeClientClaudeCode, sessions, 5, AICodeBackfillMaxBatchBytes, 50) // a (4) fits alone; big (12) is split 5+5+2; c (3) joins big's last chunk. require.Len(t, batches, 4) @@ -64,7 +64,7 @@ func TestPackBackfillBatchesRespectsCompletedCap(t *testing.T) { sessions = append(sessions, testBackfillSession(fmt.Sprintf("s%d", i), base.Add(time.Duration(i)*time.Minute), 1)) } - batches := PackBackfillBatches(AICodeClientCodex, sessions, 500, 2) + batches := PackBackfillBatches(AICodeClientCodex, sessions, 500, AICodeBackfillMaxBatchBytes, 2) require.Len(t, batches, 3) assert.Equal(t, []string{"s0", "s1"}, batches[0].CompletedSessionIDs) @@ -72,6 +72,23 @@ func TestPackBackfillBatchesRespectsCompletedCap(t *testing.T) { assert.Equal(t, []string{"s4"}, batches[2].CompletedSessionIDs) } +func TestPackBackfillBatchesRespectsByteBudget(t *testing.T) { + base := time.Date(2026, 9, 1, 10, 0, 0, 0, time.UTC) + s := testBackfillSession("prompts", base, 6) + for i := range s.Events { + s.Events[i].Prompt = strings.Repeat("p", 1000) + } + eventSize := backfillEventSize(&s.Events[0]) + + batches := PackBackfillBatches(AICodeClientClaudeCode, []*BackfillSession{s}, 500, eventSize*2, 50) + + require.Len(t, batches, 3) + for _, b := range batches { + assert.Len(t, b.Events, 2) + } + assert.Equal(t, []string{"prompts"}, batches[2].CompletedSessionIDs) +} + func TestSelectBackfillSessions(t *testing.T) { now := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) old := testBackfillSession("old", now.AddDate(0, 0, -20), 2) diff --git a/model/aicode_backfill_types.go b/model/aicode_backfill_types.go index 5ee007e..0e17502 100644 --- a/model/aicode_backfill_types.go +++ b/model/aicode_backfill_types.go @@ -17,9 +17,11 @@ const ( AICodeBackfillStatusArchived = "archived" ) -// Server limits for one backfill request. +// Server limits for one backfill request. The body limit is 16 MiB; batches +// are packed to half of it to leave room for the JSON envelope. const ( AICodeBackfillMaxEvents = 500 + AICodeBackfillMaxBatchBytes = 8 << 20 AICodeBackfillMaxCompleted = 50 AICodeBackfillMaxSessionIDs = 1000 ) From 874078c4e4b27c00bbe0fa0f42b3ace18ae96264 Mon Sep 17 00:00:00 2001 From: Claude Date: Mon, 5 Oct 2026 11:35:22 +0000 Subject: [PATCH 3/3] feat(commands): add `shelltime cc backfill` and `shelltime codex backfill` Upload past Claude Code and Codex usage from local transcripts that live OTEL tracking missed (before `cc install` or while the daemon was down). The command asks the server which sessions it already has, skips those tracked live, archived or already uploaded, and uploads the rest in batches with retries on server errors. Sessions still running are held back. When it finishes it asks the server to refresh activity data and caches for the uploaded range. Flags: --since/--until, --dry-run (per-day preview), --no-prompts and --ai-summary (opt-in, uses AI credits). Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LMCwtubXzBYhQTnX44vhzF --- README.md | 18 ++ commands/aicode_backfill.go | 403 +++++++++++++++++++++++++++++++ commands/aicode_backfill_test.go | 332 +++++++++++++++++++++++++ commands/cc.go | 1 + commands/codex.go | 1 + 5 files changed, 755 insertions(+) create mode 100644 commands/aicode_backfill.go create mode 100644 commands/aicode_backfill_test.go diff --git a/README.md b/README.md index 26c65a3..ac3031b 100644 --- a/README.md +++ b/README.md @@ -94,8 +94,10 @@ shelltime codex install | `shelltime cc install` | Install Claude Code OTEL configuration into `~/.claude/settings.json` | | `shelltime cc uninstall` | Remove Claude Code OTEL configuration from `~/.claude/settings.json` | | `shelltime cc statusline` | Emit statusline JSON for Claude Code | +| `shelltime cc backfill` | Upload past Claude Code usage from local transcripts | | `shelltime codex install` | Add ShellTime OTEL config to `~/.codex/config.toml` | | `shelltime codex uninstall` | Remove ShellTime OTEL config from `~/.codex/config.toml` | +| `shelltime codex backfill` | Upload past Codex usage from local session files | ### Environment helpers @@ -189,6 +191,22 @@ Quota sync requires both a ShellTime login (`shelltime auth`) and a ChatGPT-auth Codex decides which windows are present. ShellTime displays the windows returned by Codex instead of assuming that every account has a fixed 5-hour window. +## Backfilling AI Usage + +Live tracking only records sessions that run while the OTEL configuration is installed and the daemon is running. To upload earlier sessions from the transcripts Claude Code and Codex keep on disk: + +```bash +shelltime cc backfill --dry-run # show what would be uploaded +shelltime cc backfill # upload Claude Code sessions +shelltime codex backfill # upload Codex sessions +``` + +- Claude Code transcripts are read from `~/.claude/projects` and `~/.config/claude/projects`, or the directories in `CLAUDE_CONFIG_DIR`. Codex sessions are read from `~/.codex/sessions` and `~/.codex/archived_sessions`, or `CODEX_HOME`. +- Prompts, token usage, models and tool calls are uploaded as if they had been tracked live; the server adds costs. Lines of code, commits and active time are not in the transcripts. +- Sessions the server already has from live tracking are skipped, as are sessions still running. Running the command again only uploads what is missing. +- Flags: `--since` / `--until` (`YYYY-MM-DD`) limit the range, `--no-prompts` uploads prompt lengths without the text, and `--ai-summary` also generates AI session summaries, which use your monthly AI credits. +- Claude Code deletes transcripts after 30 days by default (`cleanupPeriodDays`), so only recent history may be available. + ## Security and Privacy - **Data masking** redacts sensitive command content before it leaves your machine. diff --git a/commands/aicode_backfill.go b/commands/aicode_backfill.go new file mode 100644 index 0000000..7a38d29 --- /dev/null +++ b/commands/aicode_backfill.go @@ -0,0 +1,403 @@ +package commands + +import ( + "context" + "errors" + "fmt" + "log/slog" + "net/http" + "os" + "sort" + "strconv" + "time" + + "github.com/gookit/color" + "github.com/malamtime/cli/model" + "github.com/olekukonko/tablewriter" + "github.com/urfave/cli/v2" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/trace" +) + +var CCBackfillCommand = &cli.Command{ + Name: "backfill", + Usage: "Upload Claude Code usage from local transcripts that live tracking missed", + Flags: aiCodeBackfillFlags(), + Action: commandCCBackfill, +} + +var CodexBackfillCommand = &cli.Command{ + Name: "backfill", + Usage: "Upload Codex usage from local session files that live tracking missed", + Flags: aiCodeBackfillFlags(), + Action: commandCodexBackfill, +} + +// aiCodeBackfillSource describes where a client keeps its transcripts and +// how to read them. +type aiCodeBackfillSource struct { + clientType string + name string + roots func() []string + parse func(roots []string, opts model.BackfillOptions) (*model.BackfillParseResult, error) +} + +func commandCCBackfill(c *cli.Context) error { + return runAICodeBackfill(c, aiCodeBackfillSource{ + clientType: model.AICodeClientClaudeCode, + name: "Claude Code", + roots: model.ClaudeProjectRoots, + parse: model.ParseClaudeTranscripts, + }) +} + +func commandCodexBackfill(c *cli.Context) error { + return runAICodeBackfill(c, aiCodeBackfillSource{ + clientType: model.AICodeClientCodex, + name: "Codex", + roots: model.CodexSessionRoots, + parse: model.ParseCodexRollouts, + }) +} + +func aiCodeBackfillFlags() []cli.Flag { + return []cli.Flag{ + &cli.StringFlag{ + Name: "since", + Usage: "only upload sessions started on or after this day (YYYY-MM-DD, local time)", + }, + &cli.StringFlag{ + Name: "until", + Usage: "only upload sessions started on or before this day (YYYY-MM-DD, local time)", + }, + &cli.BoolFlag{ + Name: "dry-run", + Usage: "show what would be uploaded without uploading anything", + }, + &cli.BoolFlag{ + Name: "no-prompts", + Usage: "upload prompt lengths but not the prompt text", + }, + &cli.BoolFlag{ + Name: "ai-summary", + Usage: "also generate AI summaries for uploaded sessions (uses your monthly AI credits)", + }, + } +} + +// backfillBackoff is how long to wait before each retry of a failed upload. +var backfillBackoff = []time.Duration{time.Second, 2 * time.Second, 4 * time.Second, 8 * time.Second, 16 * time.Second} + +// backfillPlan sorts local sessions by what the server already has. +type backfillPlan struct { + upload []*model.BackfillSession + done []*model.BackfillSession + live []*model.BackfillSession + archived []*model.BackfillSession +} + +func runAICodeBackfill(c *cli.Context, src aiCodeBackfillSource) error { + ctx, span := commandTracer.Start(c.Context, "aicode.backfill", trace.WithSpanKind(trace.SpanKindClient)) + defer span.End() + span.SetAttributes(attribute.String("clientType", src.clientType)) + + SetupLogger(os.ExpandEnv("$HOME/" + model.COMMAND_BASE_STORAGE_FOLDER)) + + opts, err := backfillOptionsFromFlags(c) + if err != nil { + return err + } + + cfg, err := configService.ReadConfigFile(ctx) + if err != nil { + slog.Error("failed to read config file", slog.Any("err", err)) + return err + } + if cfg.Token == "" { + return errors.New("no ShellTime token configured, run `shelltime init` first") + } + endpoint := model.Endpoint{Token: cfg.Token, APIEndpoint: cfg.APIEndpoint} + + color.Yellow.Printf("Scanning %s transcripts...\n", src.name) + parsed, err := src.parse(src.roots(), opts) + if err != nil { + return fmt.Errorf("failed to read %s transcripts: %w", src.name, err) + } + sessions, active := model.SelectBackfillSessions(parsed.Sessions, opts) + span.SetAttributes(attribute.Int("files", parsed.Files), attribute.Int("sessions", len(sessions))) + if parsed.BadLines > 0 { + slog.Warn("skipped unreadable transcript lines", slog.Int("count", parsed.BadLines)) + } + if len(sessions) == 0 { + fmt.Printf("No finished %s sessions found in %d files.\n", src.name, parsed.Files) + return nil + } + + statuses, err := fetchBackfillStatuses(ctx, endpoint, src.clientType, sessions) + if err != nil { + return backfillServerError(err) + } + plan := classifyBackfillSessions(sessions, statuses) + printBackfillPlan(plan, active) + + if c.Bool("dry-run") { + printBackfillDays(src.clientType, plan.upload) + color.Green.Println("Dry run: nothing was uploaded.") + return nil + } + if len(plan.upload) == 0 { + color.Green.Println("Nothing new to upload.") + return nil + } + + batches := model.PackBackfillBatches(src.clientType, plan.upload, + model.AICodeBackfillMaxEvents, model.AICodeBackfillMaxBatchBytes, model.AICodeBackfillMaxCompleted) + accepted := 0 + skippedByServer := map[string]bool{} + for i := range batches { + batch := batches[i] + batch.AISummary = c.Bool("ai-summary") + var resp *model.AICodeBackfillResponse + err := sendWithRetry(ctx, func() error { + var err error + resp, err = model.SendAICodeBackfill(ctx, endpoint, batch) + return err + }) + if err != nil { + fmt.Println() + return backfillServerError(err) + } + accepted += resp.Accepted + for _, s := range resp.SkippedSessions { + skippedByServer[s.SessionID] = true + } + fmt.Printf("\rUploading... batch %d/%d", i+1, len(batches)) + } + fmt.Println() + + from, to := backfillRange(plan.upload) + err = sendWithRetry(ctx, func() error { + return model.CompleteAICodeBackfill(ctx, endpoint, model.AICodeBackfillCompleteRequest{ + ClientType: src.clientType, + From: from, + To: to, + }) + }) + if err != nil { + // The data is stored; only the activity refresh and cache flush failed. + color.Yellow.Printf("Uploaded, but refreshing dashboards failed: %v\n", err) + } + + color.Green.Printf("Uploaded %d events from %d sessions.\n", accepted, len(plan.upload)-len(skippedByServer)) + if len(skippedByServer) > 0 { + fmt.Printf("%d sessions were skipped because the server already tracks them.\n", len(skippedByServer)) + } + if c.Bool("ai-summary") { + fmt.Println("AI summaries are being generated in the background.") + } + fmt.Println("Note: free plans show the last 7 days of AI coding history.") + return nil +} + +func backfillOptionsFromFlags(c *cli.Context) (model.BackfillOptions, error) { + opts := model.BackfillOptions{NoPrompts: c.Bool("no-prompts"), Now: time.Now()} + if v := c.String("since"); v != "" { + day, err := time.ParseInLocation(time.DateOnly, v, time.Local) + if err != nil { + return opts, fmt.Errorf("invalid --since %q, expected YYYY-MM-DD", v) + } + opts.Since = day + } + if v := c.String("until"); v != "" { + day, err := time.ParseInLocation(time.DateOnly, v, time.Local) + if err != nil { + return opts, fmt.Errorf("invalid --until %q, expected YYYY-MM-DD", v) + } + // Inclusive of the whole day. + opts.Until = day.AddDate(0, 0, 1) + } + if !opts.Since.IsZero() && !opts.Until.IsZero() && !opts.Since.Before(opts.Until) { + return opts, errors.New("--since must not be after --until") + } + return opts, nil +} + +func fetchBackfillStatuses(ctx context.Context, endpoint model.Endpoint, clientType string, sessions []*model.BackfillSession) (map[string]model.AICodeBackfillSessionStatus, error) { + statuses := map[string]model.AICodeBackfillSessionStatus{} + for start := 0; start < len(sessions); start += model.AICodeBackfillMaxSessionIDs { + end := min(start+model.AICodeBackfillMaxSessionIDs, len(sessions)) + ids := make([]string, 0, end-start) + for _, s := range sessions[start:end] { + ids = append(ids, s.ID) + } + var chunk []model.AICodeBackfillSessionStatus + err := sendWithRetry(ctx, func() error { + var err error + chunk, err = model.FetchAICodeBackfillStatus(ctx, endpoint, clientType, ids) + return err + }) + if err != nil { + return nil, err + } + for _, st := range chunk { + statuses[st.SessionID] = st + } + } + return statuses, nil +} + +func classifyBackfillSessions(sessions []*model.BackfillSession, statuses map[string]model.AICodeBackfillSessionStatus) backfillPlan { + var plan backfillPlan + for _, s := range sessions { + st, known := statuses[s.ID] + switch { + case !known: + plan.upload = append(plan.upload, s) + case st.Status == model.AICodeBackfillStatusLive: + plan.live = append(plan.live, s) + case st.Status == model.AICodeBackfillStatusArchived: + plan.archived = append(plan.archived, s) + case st.Status == model.AICodeBackfillStatusBackfilled && st.EventCount >= len(s.Events): + plan.done = append(plan.done, s) + default: + // Partly uploaded before; re-sending is idempotent. + plan.upload = append(plan.upload, s) + } + } + return plan +} + +func sendWithRetry(ctx context.Context, fn func() error) error { + err := fn() + for _, wait := range backfillBackoff { + if err == nil || !isRetryableBackfillError(err) { + return err + } + slog.Warn("backfill request failed, retrying", slog.Any("err", err), slog.Duration("wait", wait)) + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(wait): + } + err = fn() + } + return err +} + +func isRetryableBackfillError(err error) bool { + if errors.Is(err, context.Canceled) { + return false + } + var statusErr *model.HTTPStatusError + if errors.As(err, &statusErr) { + return statusErr.StatusCode >= http.StatusInternalServerError || statusErr.StatusCode == http.StatusTooManyRequests + } + // Network errors and timeouts. + return true +} + +func backfillServerError(err error) error { + var statusErr *model.HTTPStatusError + if errors.As(err, &statusErr) && statusErr.StatusCode == http.StatusNotFound { + return errors.New("the ShellTime server does not support backfill yet, please try again after it is updated") + } + return fmt.Errorf("backfill failed: %w", err) +} + +func backfillRange(sessions []*model.BackfillSession) (time.Time, time.Time) { + var from, to time.Time + for _, s := range sessions { + if from.IsZero() || s.Start.Before(from) { + from = s.Start + } + if s.End.After(to) { + to = s.End + } + } + return from, to +} + +func countBackfillEvents(sessions []*model.BackfillSession) int { + n := 0 + for _, s := range sessions { + n += len(s.Events) + } + return n +} + +func printBackfillPlan(plan backfillPlan, active int) { + w := tablewriter.NewWriter(os.Stdout) + w.Header([]string{"SESSIONS", "COUNT", "EVENTS"}) + rows := []struct { + label string + sessions []*model.BackfillSession + }{ + {"to upload", plan.upload}, + {"already uploaded", plan.done}, + {"tracked live (skipped)", plan.live}, + {"archived on server (skipped)", plan.archived}, + } + for _, r := range rows { + w.Append([]string{r.label, strconv.Itoa(len(r.sessions)), strconv.Itoa(countBackfillEvents(r.sessions))}) + } + if active > 0 { + w.Append([]string{"still active (skipped)", strconv.Itoa(active), "-"}) + } + w.Render() +} + +type backfillDay struct { + sessions, prompts, requests, tokens int +} + +// printBackfillDays shows per-day totals of what would be uploaded. Tokens +// add up input, output and cache tokens; Codex input already includes cached +// tokens, so cache reads are only added for Claude Code. +func printBackfillDays(clientType string, sessions []*model.BackfillSession) { + if len(sessions) == 0 { + return + } + days := map[string]*backfillDay{} + for _, s := range sessions { + key := s.Start.Local().Format(time.DateOnly) + d := days[key] + if d == nil { + d = &backfillDay{} + days[key] = d + } + d.sessions++ + for _, e := range s.Events { + switch e.EventType { + case model.AICodeEventUserPrompt: + d.prompts++ + case model.AICodeEventApiRequest, model.AICodeEventSSEEvent: + d.requests++ + d.tokens += derefInt(e.InputTokens) + derefInt(e.OutputTokens) + derefInt(e.CacheCreationTokens) + if clientType == model.AICodeClientClaudeCode { + d.tokens += derefInt(e.CacheReadTokens) + } + } + } + } + + keys := make([]string, 0, len(days)) + for k := range days { + keys = append(keys, k) + } + sort.Strings(keys) + + w := tablewriter.NewWriter(os.Stdout) + w.Header([]string{"DAY", "SESSIONS", "PROMPTS", "REQUESTS", "TOKENS"}) + for _, k := range keys { + d := days[k] + w.Append([]string{k, strconv.Itoa(d.sessions), strconv.Itoa(d.prompts), strconv.Itoa(d.requests), strconv.Itoa(d.tokens)}) + } + w.Render() +} + +func derefInt(v *int) int { + if v == nil { + return 0 + } + return *v +} diff --git a/commands/aicode_backfill_test.go b/commands/aicode_backfill_test.go new file mode 100644 index 0000000..58d4c0e --- /dev/null +++ b/commands/aicode_backfill_test.go @@ -0,0 +1,332 @@ +package commands + +import ( + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "sync" + "testing" + "time" + + "github.com/malamtime/cli/model" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" + "github.com/urfave/cli/v2" +) + +// backfillServer fakes the server's backfill endpoints. +type backfillServer struct { + t *testing.T + mu sync.Mutex + statuses map[string]model.AICodeBackfillSessionStatus + // failures are returned, in order, before backfill uploads succeed. + failures []int + statusCode int + statusCalls int + uploads []model.AICodeBackfillRequest + uploadCalls int + completeCalls []model.AICodeBackfillCompleteRequest + server *httptest.Server +} + +func newBackfillServer(t *testing.T) *backfillServer { + bs := &backfillServer{t: t, statuses: map[string]model.AICodeBackfillSessionStatus{}} + bs.server = httptest.NewServer(http.HandlerFunc(bs.handle)) + t.Cleanup(bs.server.Close) + return bs +} + +func (bs *backfillServer) handle(w http.ResponseWriter, r *http.Request) { + bs.mu.Lock() + defer bs.mu.Unlock() + assert.Equal(bs.t, "CLI test-token", r.Header.Get("Authorization")) + + switch r.URL.Path { + case "/api/v1/cc/backfill/sessions": + bs.statusCalls++ + if bs.statusCode != 0 { + w.WriteHeader(bs.statusCode) + _, _ = w.Write([]byte("404 page not found")) + return + } + var req model.AICodeBackfillSessionsRequest + require.NoError(bs.t, json.NewDecoder(r.Body).Decode(&req)) + resp := model.AICodeBackfillSessionsResponse{Sessions: []model.AICodeBackfillSessionStatus{}} + for _, id := range req.SessionIDs { + if st, ok := bs.statuses[id]; ok { + resp.Sessions = append(resp.Sessions, st) + } + } + _ = json.NewEncoder(w).Encode(resp) + case "/api/v1/cc/backfill": + bs.uploadCalls++ + if len(bs.failures) > 0 { + code := bs.failures[0] + bs.failures = bs.failures[1:] + w.WriteHeader(code) + _, _ = fmt.Fprintf(w, `{"code":%d,"error":"backfill rejected"}`, code) + return + } + var req model.AICodeBackfillRequest + require.NoError(bs.t, json.NewDecoder(r.Body).Decode(&req)) + bs.uploads = append(bs.uploads, req) + _ = json.NewEncoder(w).Encode(model.AICodeBackfillResponse{ + Success: true, + Accepted: len(req.Events), + Summarized: len(req.CompletedSessionIDs), + }) + case "/api/v1/cc/backfill/complete": + var req model.AICodeBackfillCompleteRequest + require.NoError(bs.t, json.NewDecoder(r.Body).Decode(&req)) + bs.completeCalls = append(bs.completeCalls, req) + _, _ = w.Write([]byte(`{"success":true}`)) + default: + bs.t.Errorf("unexpected request to %s", r.URL.Path) + w.WriteHeader(http.StatusNotFound) + } +} + +func setupBackfillTest(t *testing.T, token string) *backfillServer { + t.Helper() + setupCCTest(t) + bs := newBackfillServer(t) + + original := configService + mc := model.NewMockConfigService(t) + mc.On("ReadConfigFile", mock.Anything).Return(model.ShellTimeConfig{Token: token, APIEndpoint: bs.server.URL}, nil).Maybe() + configService = mc + t.Cleanup(func() { configService = original }) + + originalBackoff := backfillBackoff + backfillBackoff = []time.Duration{time.Millisecond, time.Millisecond, time.Millisecond} + t.Cleanup(func() { backfillBackoff = originalBackoff }) + return bs +} + +func writeLines(t *testing.T, path string, lines ...map[string]any) { + t.Helper() + require.NoError(t, os.MkdirAll(filepath.Dir(path), 0o755)) + var b strings.Builder + for _, l := range lines { + buf, err := json.Marshal(l) + require.NoError(t, err) + b.Write(buf) + b.WriteString("\n") + } + require.NoError(t, os.WriteFile(path, []byte(b.String()), 0o644)) +} + +// writeClaudeSession writes a finished Claude Code session with one prompt +// and the given number of API responses. +func writeClaudeSession(t *testing.T, projects, sessionID string, start time.Time, responses int) { + t.Helper() + ts := func(d time.Duration) string { return start.Add(d).UTC().Format(time.RFC3339Nano) } + lines := []map[string]any{{ + "type": "user", "uuid": sessionID + "-u", "sessionId": sessionID, "timestamp": ts(0), "cwd": "/tmp/p", + "origin": map[string]any{"kind": "human"}, + "message": map[string]any{"role": "user", "content": "prompt for " + sessionID}, + }} + for i := 0; i < responses; i++ { + lines = append(lines, map[string]any{ + "type": "assistant", "uuid": sessionID + "-a" + string(rune('0'+i)), "sessionId": sessionID, + "timestamp": ts(time.Duration(i+1) * time.Second), "requestId": sessionID + "-req" + string(rune('0'+i)), + "message": map[string]any{ + "id": sessionID + "-msg" + string(rune('0'+i)), "model": "claude-sonnet-4-5", + "content": []any{map[string]any{"type": "text", "text": "ok"}}, + "usage": map[string]any{"input_tokens": 10, "output_tokens": 20, "cache_read_input_tokens": 100}, + }, + }) + } + writeLines(t, filepath.Join(projects, "-tmp-p", sessionID+".jsonl"), lines...) +} + +func runBackfill(t *testing.T, command *cli.Command, group string, args ...string) error { + t.Helper() + app := &cli.App{Name: "t", Commands: []*cli.Command{{Name: group, Subcommands: []*cli.Command{command}}}} + return app.Run(append([]string{"t", group, "backfill"}, args...)) +} + +func TestCCBackfill_UploadsOnlyNewSessions(t *testing.T) { + bs := setupBackfillTest(t, "test-token") + root := t.TempDir() + projects := filepath.Join(root, "projects") + start := time.Now().AddDate(0, 0, -3) + writeClaudeSession(t, projects, "sess-new", start, 2) + writeClaudeSession(t, projects, "sess-live", start.Add(time.Hour), 1) + writeClaudeSession(t, projects, "sess-done", start.Add(2*time.Hour), 1) + writeClaudeSession(t, projects, "sess-partial", start.Add(3*time.Hour), 2) + writeClaudeSession(t, projects, "sess-active", time.Now().Add(-5*time.Minute), 1) + t.Setenv("CLAUDE_CONFIG_DIR", root) + + bs.statuses["sess-live"] = model.AICodeBackfillSessionStatus{SessionID: "sess-live", Status: model.AICodeBackfillStatusLive} + bs.statuses["sess-done"] = model.AICodeBackfillSessionStatus{SessionID: "sess-done", Status: model.AICodeBackfillStatusBackfilled, EventCount: 2} + bs.statuses["sess-partial"] = model.AICodeBackfillSessionStatus{SessionID: "sess-partial", Status: model.AICodeBackfillStatusBackfilled, EventCount: 1} + + require.NoError(t, runBackfill(t, CCBackfillCommand, "cc")) + + require.Len(t, bs.uploads, 1) + upload := bs.uploads[0] + assert.Equal(t, model.AICodeClientClaudeCode, upload.ClientType) + assert.False(t, upload.AISummary) + assert.Equal(t, []string{"sess-new", "sess-partial"}, upload.CompletedSessionIDs) + assert.Len(t, upload.Events, 6) + for _, e := range upload.Events { + assert.Contains(t, []string{"sess-new", "sess-partial"}, e.SessionID) + assert.Equal(t, model.AICodeClientClaudeCode, e.ClientType) + } + + require.Len(t, bs.completeCalls, 1) + assert.Equal(t, model.AICodeClientClaudeCode, bs.completeCalls[0].ClientType) + assert.False(t, bs.completeCalls[0].From.After(bs.completeCalls[0].To)) +} + +func TestCCBackfill_SplitsLargeBackfills(t *testing.T) { + bs := setupBackfillTest(t, "test-token") + root := t.TempDir() + start := time.Now().AddDate(0, 0, -5) + for i := 0; i < model.AICodeBackfillMaxCompleted+5; i++ { + writeClaudeSession(t, filepath.Join(root, "projects"), "sess-"+strings.Repeat("x", i+1), start.Add(time.Duration(i)*time.Minute), 1) + } + t.Setenv("CLAUDE_CONFIG_DIR", root) + + require.NoError(t, runBackfill(t, CCBackfillCommand, "cc", "--ai-summary")) + + require.Len(t, bs.uploads, 2) + assert.Len(t, bs.uploads[0].CompletedSessionIDs, model.AICodeBackfillMaxCompleted) + assert.Len(t, bs.uploads[1].CompletedSessionIDs, 5) + for _, u := range bs.uploads { + assert.True(t, u.AISummary) + assert.LessOrEqual(t, len(u.Events), model.AICodeBackfillMaxEvents) + } +} + +func TestCCBackfill_DryRunOnlyChecksStatus(t *testing.T) { + bs := setupBackfillTest(t, "test-token") + root := t.TempDir() + writeClaudeSession(t, filepath.Join(root, "projects"), "sess-1", time.Now().AddDate(0, 0, -2), 1) + t.Setenv("CLAUDE_CONFIG_DIR", root) + + require.NoError(t, runBackfill(t, CCBackfillCommand, "cc", "--dry-run")) + + assert.Equal(t, 1, bs.statusCalls) + assert.Equal(t, 0, bs.uploadCalls) + assert.Empty(t, bs.completeCalls) +} + +func TestCCBackfill_RetriesServerErrors(t *testing.T) { + bs := setupBackfillTest(t, "test-token") + root := t.TempDir() + writeClaudeSession(t, filepath.Join(root, "projects"), "sess-1", time.Now().AddDate(0, 0, -2), 1) + t.Setenv("CLAUDE_CONFIG_DIR", root) + bs.failures = []int{http.StatusServiceUnavailable, http.StatusTooManyRequests} + + require.NoError(t, runBackfill(t, CCBackfillCommand, "cc")) + + assert.Equal(t, 3, bs.uploadCalls) + assert.Len(t, bs.uploads, 1) + assert.Len(t, bs.completeCalls, 1) +} + +func TestCCBackfill_AbortsOnClientError(t *testing.T) { + bs := setupBackfillTest(t, "test-token") + root := t.TempDir() + writeClaudeSession(t, filepath.Join(root, "projects"), "sess-1", time.Now().AddDate(0, 0, -2), 1) + t.Setenv("CLAUDE_CONFIG_DIR", root) + bs.failures = []int{http.StatusBadRequest} + + err := runBackfill(t, CCBackfillCommand, "cc") + + require.Error(t, err) + assert.Contains(t, err.Error(), "backfill rejected") + assert.Equal(t, 1, bs.uploadCalls) + assert.Empty(t, bs.completeCalls) +} + +func TestCCBackfill_ExplainsOldServer(t *testing.T) { + bs := setupBackfillTest(t, "test-token") + root := t.TempDir() + writeClaudeSession(t, filepath.Join(root, "projects"), "sess-1", time.Now().AddDate(0, 0, -2), 1) + t.Setenv("CLAUDE_CONFIG_DIR", root) + bs.statusCode = http.StatusNotFound + + err := runBackfill(t, CCBackfillCommand, "cc") + + require.Error(t, err) + assert.Contains(t, err.Error(), "does not support backfill yet") + assert.Equal(t, 1, bs.statusCalls, "client errors are not retried") +} + +func TestCCBackfill_RequiresToken(t *testing.T) { + setupBackfillTest(t, "") + + err := runBackfill(t, CCBackfillCommand, "cc") + + require.Error(t, err) + assert.Contains(t, err.Error(), "shelltime init") +} + +func TestCodexBackfill_UploadsSessions(t *testing.T) { + bs := setupBackfillTest(t, "test-token") + home := t.TempDir() + sessionID := "44444444-4444-4444-8444-444444444444" + start := time.Now().AddDate(0, 0, -1).UTC() + ts := func(d time.Duration) string { return start.Add(d).Format(time.RFC3339Nano) } + usage := map[string]any{"input_tokens": 1000, "cached_input_tokens": 600, "output_tokens": 50, "reasoning_output_tokens": 10, "total_tokens": 1050} + writeLines(t, filepath.Join(home, "sessions", "2026", "10", "04", "rollout-2026-10-04T10-00-00-"+sessionID+".jsonl"), + map[string]any{"timestamp": ts(0), "type": "session_meta", "payload": map[string]any{"id": sessionID, "cwd": "/tmp/c", "cli_version": "0.42.0"}}, + map[string]any{"timestamp": ts(time.Second), "type": "turn_context", "payload": map[string]any{"model": "gpt-5-codex"}}, + map[string]any{"timestamp": ts(2 * time.Second), "type": "event_msg", "payload": map[string]any{"type": "user_message", "message": "hello"}}, + map[string]any{"timestamp": ts(3 * time.Second), "type": "event_msg", "payload": map[string]any{"type": "token_count", "info": map[string]any{"total_token_usage": usage, "last_token_usage": usage}}}, + ) + t.Setenv("CODEX_HOME", home) + + require.NoError(t, runBackfill(t, CodexBackfillCommand, "codex", "--no-prompts")) + + require.Len(t, bs.uploads, 1) + upload := bs.uploads[0] + assert.Equal(t, model.AICodeClientCodex, upload.ClientType) + assert.Equal(t, []string{sessionID}, upload.CompletedSessionIDs) + types := map[string]model.AICodeBackfillEvent{} + for _, e := range upload.Events { + types[e.EventType] = e + } + assert.Contains(t, types, model.AICodeEventConversationStarts) + assert.Empty(t, types[model.AICodeEventUserPrompt].Prompt) + assert.Equal(t, 5, *types[model.AICodeEventUserPrompt].PromptLength) + sse := types[model.AICodeEventSSEEvent] + assert.Equal(t, "response.completed", sse.EventKind) + assert.Equal(t, 1000, *sse.InputTokens) + assert.Equal(t, 600, *sse.CacheReadTokens) + require.Len(t, bs.completeCalls, 1) + assert.Equal(t, model.AICodeClientCodex, bs.completeCalls[0].ClientType) +} + +func TestBackfillOptionsFromFlags(t *testing.T) { + parse := func(args ...string) (model.BackfillOptions, error) { + var opts model.BackfillOptions + var err error + app := &cli.App{Flags: aiCodeBackfillFlags(), Action: func(c *cli.Context) error { + opts, err = backfillOptionsFromFlags(c) + return nil + }} + require.NoError(t, app.Run(append([]string{"t"}, args...))) + return opts, err + } + + opts, err := parse("--since", "2026-09-01", "--until", "2026-09-30", "--no-prompts") + require.NoError(t, err) + assert.Equal(t, time.Date(2026, 9, 1, 0, 0, 0, 0, time.Local), opts.Since) + assert.Equal(t, time.Date(2026, 10, 1, 0, 0, 0, 0, time.Local), opts.Until, "until includes the whole day") + assert.True(t, opts.NoPrompts) + + _, err = parse("--since", "09/01/2026") + assert.ErrorContains(t, err, "invalid --since") + + _, err = parse("--since", "2026-09-30", "--until", "2026-09-01") + assert.ErrorContains(t, err, "must not be after") +} diff --git a/commands/cc.go b/commands/cc.go index 6001f63..7cea073 100644 --- a/commands/cc.go +++ b/commands/cc.go @@ -13,6 +13,7 @@ var CCCommand = &cli.Command{ CCInstallCommand, CCUninstallCommand, CCStatuslineCommand, + CCBackfillCommand, }, } diff --git a/commands/codex.go b/commands/codex.go index 3e4f580..1647242 100644 --- a/commands/codex.go +++ b/commands/codex.go @@ -12,6 +12,7 @@ var CodexCommand = &cli.Command{ Subcommands: []*cli.Command{ CodexInstallCommand, CodexUninstallCommand, + CodexBackfillCommand, }, }