diff --git a/docs/adr/60951-track-gh-aw-logs-download-telemetry.md b/docs/adr/60951-track-gh-aw-logs-download-telemetry.md new file mode 100644 index 00000000000..ef9ed2fe44b --- /dev/null +++ b/docs/adr/60951-track-gh-aw-logs-download-telemetry.md @@ -0,0 +1,46 @@ +# ADR-60951: Track gh aw logs download telemetry + +**Date**: 2026-09-15 +**Status**: Draft +**Deciders**: gh-aw maintainers + +--- + +### Context + +This pull request changes the `gh aw logs` pipeline to capture artifact download duration and downloaded size for each workflow run, persist those values into cached JSONL/log schemas, and print an end-of-run download summary across single-target, multi-target, and stdin-driven entry points. The PR description explains that the command previously did not show how long artifact downloads took or how much data they transferred, which made slow runs and GitHub API usage harder to diagnose. The implementation also derives an approximate per-run GitHub API cost from existing rate-limit reports rather than introducing a separate telemetry source. The architectural question is whether `gh aw logs` should treat download-performance telemetry as first-class run metadata and surface it consistently through its reporting and cache formats. + +### Decision + +We will record per-run artifact download duration and artifact size inside the `gh aw logs` workflow-run model, propagate those fields into persisted `RunData`/schema outputs, and render an aggregate end-of-run summary for real downloads. We will measure download duration around the existing artifact-download path, compute size from the downloaded artifact directory, and exclude cache hits from the aggregate timing and size metrics so the summary reflects actual transfer work performed during the invocation. We will also estimate GitHub API cost per downloaded run from existing rate-limit snapshots instead of adding a new API accounting mechanism. We chose this because the PR evidence shows the main problem is lack of visibility into download cost, and extending the current logs/reporting pipeline solves that with minimal new architecture. + +### Alternatives Considered + +#### Alternative 1: Keep download telemetry out of the run model and rely on ad hoc debug logging + +The team could have added temporary or verbose-only log lines around artifact downloads without changing `WorkflowRun`, `RunData`, or the JSON schemas. This was considered because it would be a smaller code change and avoid widening cached output contracts. It was not chosen because the PR explicitly updates cached JSONL persistence and schemas, indicating the telemetry needs to survive beyond a single terminal session and be available for downstream analysis. + +#### Alternative 2: Report only aggregate command-level download timing + +Another option was to compute one total download duration and size for the full command without attaching telemetry to each run. This was considered because it would still improve operator visibility while avoiding new per-run fields. It was not chosen because the PR adds `DownloadDuration` and `DownloadSizeBytes` directly to `WorkflowRun` and `RunData`, showing the intended design is per-run observability that can be aggregated later and reused from cache. + +### Consequences + +#### Positive +- `gh aw logs` gains concrete visibility into artifact transfer cost, making slow or heavy runs easier to diagnose. +- The telemetry is available both in terminal summaries and persisted JSON/JSONL outputs, enabling later analysis and cache round-tripping. +- Reusing existing rate-limit snapshots provides a lightweight API-cost estimate without introducing a separate tracking subsystem. + +#### Negative +- The run/report model and published schemas gain additional fields that maintainers must preserve or evolve carefully. +- Size measurement adds extra filesystem work after downloads and may produce partial observability when directory sizing fails. +- The API cost figure is only an estimate and can understate true usage during rate-limit window resets. + +#### Neutral +- Cache hits continue to produce zero download metrics and are intentionally excluded from aggregate download summaries. +- The implementation factors artifact downloading into a helper function to isolate timing and sizing behavior from the rest of run processing. +- All three logs entry points now call the same end-of-run summary renderer, increasing consistency across invocation modes. + +--- + +*ADR created by [adr-writer agent]. Review and finalize before changing status from Draft to Accepted.* diff --git a/pkg/cli/logs_models.go b/pkg/cli/logs_models.go index 1b8ae8f5091..f96880ac7cb 100644 --- a/pkg/cli/logs_models.go +++ b/pkg/cli/logs_models.go @@ -97,6 +97,13 @@ type WorkflowRun struct { EffectiveTokens int // Cost-normalized token count computed from per-model multipliers AvgTimeBetweenTurns time.Duration // Average time between consecutive LLM API calls (from per-turn timestamps when available) LogsPath string + // DownloadDuration is the wall-clock time spent downloading this run's artifacts + // from GitHub. It is zero for runs served from the on-disk cache (no download + // was performed). + DownloadDuration time.Duration + // DownloadSizeBytes is the total on-disk size of the artifacts downloaded for + // this run, used to estimate average/maximum transfer volume across a batch. + DownloadSizeBytes int64 } // LogMetrics represents extracted metrics from log files diff --git a/pkg/cli/logs_multi.go b/pkg/cli/logs_multi.go index 78a356cacf4..dff4afe081c 100644 --- a/pkg/cli/logs_multi.go +++ b/pkg/cli/logs_multi.go @@ -164,6 +164,9 @@ func DownloadWorkflowLogsForTargets( //nolint:largefunc // Keeps shared collecti if err := prepareCachedLogsJSONL(&opts); err != nil { return err } + if opts.collectionStats == nil { + opts.collectionStats = &logsCollectionStats{} + } defer func() { err = errors.Join(err, finalizeCachedLogsJSONL(opts.cachedJSONLWriter, opts.cachedJSONLSourcePaths, opts.cachedJSONLWildcard, opts.StartDate, opts.EndDate)) }() @@ -180,6 +183,7 @@ func DownloadWorkflowLogsForTargets( //nolint:largefunc // Keeps shared collecti finishGitHubAPIRateLimitReports(activeCtx, allAPIRateLimits, opts.JSONOutput) cacheGitHubAPIRateLimitReports(opts.cachedJSONLWriter, allAPIRateLimits...) apiRateLimit, apiRateLimits := partitionGitHubAPIRateLimitReports(allAPIRateLimits) + renderLogsDownloadStatsSummary(opts.collectionStats, allAPIRateLimits...) if len(processedRuns) == 0 { if len(allErrors) > 0 { return errors.Join(allErrors...) diff --git a/pkg/cli/logs_orchestrator.go b/pkg/cli/logs_orchestrator.go index 30112992e4b..06fba6534d3 100644 --- a/pkg/cli/logs_orchestrator.go +++ b/pkg/cli/logs_orchestrator.go @@ -304,6 +304,9 @@ func DownloadWorkflowLogs(ctx context.Context, opts LogsDownloadOptions) (err er if err := prepareCachedLogsJSONL(&opts); err != nil { return err } + if opts.collectionStats == nil { + opts.collectionStats = &logsCollectionStats{} + } defer func() { err = errors.Join(err, finalizeCachedLogsJSONL(opts.cachedJSONLWriter, opts.cachedJSONLSourcePaths, opts.cachedJSONLWildcard, opts.StartDate, opts.EndDate)) }() @@ -314,6 +317,7 @@ func DownloadWorkflowLogs(ctx context.Context, opts LogsDownloadOptions) (err er } renderLogsCollectionStats(opts.collectionStats) finishGitHubAPIRateLimitReport(ctx, apiRateLimit, opts.JSONOutput) + renderLogsDownloadStatsSummary(opts.collectionStats, apiRateLimit) cacheGitHubAPIRateLimitReports(opts.cachedJSONLWriter, apiRateLimit) if handled, err := handleEmptyProcessedRuns(result.processedRuns, opts, result.timeoutReached, result.storageLimitReached, result.continuation, nil, apiRateLimit, nil); handled || err != nil { logsOrchestratorLog.Printf("No processed runs to render (timeoutReached=%v, err=%v)", result.timeoutReached, err) diff --git a/pkg/cli/logs_orchestrator_download.go b/pkg/cli/logs_orchestrator_download.go index 7e582e70ddc..e05360b431f 100644 --- a/pkg/cli/logs_orchestrator_download.go +++ b/pkg/cli/logs_orchestrator_download.go @@ -39,6 +39,15 @@ type logsCollectionStats struct { discoveredRuns atomic.Int64 downloadedReports atomic.Int64 cachedReports atomic.Int64 + // downloadCount, totalDownloadNanos, maxDownloadNanos, totalDownloadBytes, and + // maxDownloadBytes track artifact-download timing and size across runs that were + // actually downloaded this invocation (not served from the on-disk cache), so an + // avg/max summary can be rendered at the end of the run. + downloadCount atomic.Int64 + totalDownloadNanos atomic.Int64 + maxDownloadNanos atomic.Int64 + totalDownloadBytes atomic.Int64 + maxDownloadBytes atomic.Int64 } func (s *logsCollectionStats) recordDiscovered(count int) { @@ -57,6 +66,36 @@ func (s *logsCollectionStats) recordResult(result DownloadResult) { } if !result.Skipped && result.Error == nil { s.downloadedReports.Add(1) + if result.Run.DownloadDuration > 0 { + s.recordDownloadStats(result.Run.DownloadDuration, result.Run.DownloadSizeBytes) + } + } +} + +// recordDownloadStats accumulates per-run download duration and size so that +// renderLogsDownloadStatsSummary can report avg/max values at the end of the run. +func (s *logsCollectionStats) recordDownloadStats(duration time.Duration, sizeBytes int64) { + if s == nil { + return + } + s.downloadCount.Add(1) + s.totalDownloadNanos.Add(duration.Nanoseconds()) + s.totalDownloadBytes.Add(sizeBytes) + atomicMaxInt64(&s.maxDownloadNanos, duration.Nanoseconds()) + atomicMaxInt64(&s.maxDownloadBytes, sizeBytes) +} + +// atomicMaxInt64 atomically sets *addr to value if value is greater than the +// current contents, using a compare-and-swap retry loop. +func atomicMaxInt64(addr *atomic.Int64, value int64) { + for { + current := addr.Load() + if value <= current { + return + } + if addr.CompareAndSwap(current, value) { + return + } } } @@ -70,6 +109,73 @@ func renderLogsCollectionStats(stats *logsCollectionStats) { ))) } +// gitHubAPIRateLimitCostEstimate sums the core GitHub API requests consumed +// across one or more rate-limit reports (Start/End snapshots taken around the +// logs command). Returns ok=false when no populated report is available. +func gitHubAPIRateLimitCostEstimate(reports []*GitHubAPIRateLimitReport) (int, bool) { + var total int + var found bool + for _, report := range reports { + if report == nil || report.Start == nil || report.End == nil { + continue + } + diff := report.End.Used - report.Start.Used + if diff < 0 || report.End.Reset != report.Start.Reset { + // The rate-limit window reset mid-run (Used wrapped back down, or the + // reset timestamp itself moved even though Used happened to still be + // >= Start.Used); fall back to the ending value as a lower-bound + // approximation rather than mixing counters from different windows. + diff = report.End.Used + } + total += diff + found = true + } + return total, found +} + +// formatDownloadByteSize renders a byte count in a compact human-readable form +// (B, KB, MB, GB) for the end-of-run download stats summary. +func formatDownloadByteSize(bytes int64) string { + const unit = 1024 + if bytes < unit { + return fmt.Sprintf("%dB", bytes) + } + div, exp := int64(unit), 0 + for n := bytes / unit; n >= unit; n /= unit { + div *= unit + exp++ + } + return fmt.Sprintf("%.1f%ciB", float64(bytes)/float64(div), "KMGTPE"[exp]) +} + +// renderLogsDownloadStatsSummary prints an end-of-run informational line +// summarizing per-run artifact download duration and size (avg/max across runs +// actually downloaded this invocation), plus an estimated GitHub API rate-limit +// cost per run derived from the provided rate-limit reports. It is a no-op when +// no runs were downloaded or stats were not tracked. +func renderLogsDownloadStatsSummary(stats *logsCollectionStats, reports ...*GitHubAPIRateLimitReport) { + if stats == nil { + return + } + count := stats.downloadCount.Load() + if count == 0 { + return + } + avgDuration := time.Duration(stats.totalDownloadNanos.Load() / count) + maxDuration := time.Duration(stats.maxDownloadNanos.Load()) + avgSize := stats.totalDownloadBytes.Load() / count + maxSize := stats.maxDownloadBytes.Load() + msg := fmt.Sprintf( + "Download stats: avg %s (max %s) per run; avg size %s (max %s) per run", + avgDuration.Round(time.Millisecond), maxDuration.Round(time.Millisecond), + formatDownloadByteSize(avgSize), formatDownloadByteSize(maxSize), + ) + if calls, ok := gitHubAPIRateLimitCostEstimate(reports); ok { + msg += fmt.Sprintf("; GitHub API cost estimate: ~%.1f requests/run", float64(calls)/float64(count)) + } + fmt.Fprintln(os.Stderr, console.FormatInfoMessage(msg)) +} + type processWorkflowRunBatchOptions struct { count int outputDir string diff --git a/pkg/cli/logs_orchestrator_stdin.go b/pkg/cli/logs_orchestrator_stdin.go index 599de1041d7..fe306644ccb 100644 --- a/pkg/cli/logs_orchestrator_stdin.go +++ b/pkg/cli/logs_orchestrator_stdin.go @@ -280,6 +280,7 @@ func DownloadWorkflowLogsFromStdin(ctx context.Context, opts StdinLogsOptions) ( finishGitHubAPIRateLimitReports(ctx, allAPIRateLimits, opts.JSONOutput) cacheGitHubAPIRateLimitReports(cachedJSONLWriter, allAPIRateLimits...) apiRateLimit, apiRateLimits := partitionGitHubAPIRateLimitReports(allAPIRateLimits) + renderLogsDownloadStatsSummary(collectionStats, allAPIRateLimits...) if opts.JSONOutput { logsData := buildLogsData([]ProcessedRun{}, opts.OutputDir, nil) logsData.GitHubAPIRateLimit = populatedGitHubAPIRateLimitReport(apiRateLimit) @@ -307,6 +308,7 @@ func DownloadWorkflowLogsFromStdin(ctx context.Context, opts StdinLogsOptions) ( finishGitHubAPIRateLimitReports(ctx, allAPIRateLimits, opts.JSONOutput) cacheGitHubAPIRateLimitReports(cachedJSONLWriter, allAPIRateLimits...) apiRateLimit, apiRateLimits := partitionGitHubAPIRateLimitReports(allAPIRateLimits) + renderLogsDownloadStatsSummary(collectionStats, allAPIRateLimits...) return renderLogsOutput(processedRuns, renderLogsOutputOptions{ outputDir: opts.OutputDir, summaryFile: opts.SummaryFile, diff --git a/pkg/cli/logs_orchestrator_unit_test.go b/pkg/cli/logs_orchestrator_unit_test.go index b0fdf8a3bda..b46f9e74503 100644 --- a/pkg/cli/logs_orchestrator_unit_test.go +++ b/pkg/cli/logs_orchestrator_unit_test.go @@ -862,6 +862,90 @@ func TestLogsCollectionStatsReportsDiscoveredDownloadedAndCached(t *testing.T) { assert.Contains(t, stderr, "Runs: 4 discovered; reports: 1 downloaded, 2 skipped because cached analyses were reused") } +func TestLogsCollectionStatsRecordsDownloadDurationAndSize(t *testing.T) { + stats := &logsCollectionStats{} + stats.recordResult(DownloadResult{RunAnalysis: RunAnalysis{Run: WorkflowRun{DownloadDuration: 2 * time.Second, DownloadSizeBytes: 1000}}}) + stats.recordResult(DownloadResult{RunAnalysis: RunAnalysis{Run: WorkflowRun{DownloadDuration: 4 * time.Second, DownloadSizeBytes: 3000}}}) + // A cached hit must not contribute to download duration/size stats. + stats.recordResult(DownloadResult{Cached: true, RunAnalysis: RunAnalysis{Run: WorkflowRun{DownloadDuration: 10 * time.Second, DownloadSizeBytes: 999_999}}}) + + assert.Equal(t, int64(2), stats.downloadCount.Load()) + assert.Equal(t, (2*time.Second + 4*time.Second).Nanoseconds(), stats.totalDownloadNanos.Load()) + assert.Equal(t, (4 * time.Second).Nanoseconds(), stats.maxDownloadNanos.Load()) + assert.Equal(t, int64(4000), stats.totalDownloadBytes.Load()) + assert.Equal(t, int64(3000), stats.maxDownloadBytes.Load()) +} + +func TestRenderLogsDownloadStatsSummaryReportsAvgMaxAndRateLimitCost(t *testing.T) { + stats := &logsCollectionStats{} + stats.recordResult(DownloadResult{RunAnalysis: RunAnalysis{Run: WorkflowRun{DownloadDuration: 2 * time.Second, DownloadSizeBytes: 1024}}}) + stats.recordResult(DownloadResult{RunAnalysis: RunAnalysis{Run: WorkflowRun{DownloadDuration: 6 * time.Second, DownloadSizeBytes: 3072}}}) + + report := &GitHubAPIRateLimitReport{ + Start: &GitHubAPIRateLimitState{Used: 10}, + End: &GitHubAPIRateLimitState{Used: 30}, + } + + _, stderr := captureOutput(t, func() error { + renderLogsDownloadStatsSummary(stats, report) + return nil + }) + + assert.Contains(t, stderr, "Download stats: avg 4s (max 6s) per run") + assert.Contains(t, stderr, "avg size 2.0KiB (max 3.0KiB) per run") + assert.Contains(t, stderr, "GitHub API cost estimate: ~10.0 requests/run") +} + +func TestRenderLogsDownloadStatsSummaryNoOpWithoutDownloads(t *testing.T) { + stats := &logsCollectionStats{} + stats.recordResult(DownloadResult{Cached: true}) + + _, stderr := captureOutput(t, func() error { + renderLogsDownloadStatsSummary(stats) + return nil + }) + + assert.Empty(t, stderr) +} + +func TestGitHubAPIRateLimitCostEstimateHandlesWindowResetAndMissingReports(t *testing.T) { + calls, ok := gitHubAPIRateLimitCostEstimate(nil) + assert.False(t, ok) + assert.Zero(t, calls) + + calls, ok = gitHubAPIRateLimitCostEstimate([]*GitHubAPIRateLimitReport{nil, {}}) + assert.False(t, ok) + assert.Zero(t, calls) + + // A window reset mid-run (End.Used < Start.Used) falls back to End.Used as a + // lower-bound approximation instead of a negative cost. + calls, ok = gitHubAPIRateLimitCostEstimate([]*GitHubAPIRateLimitReport{ + {Start: &GitHubAPIRateLimitState{Used: 4900}, End: &GitHubAPIRateLimitState{Used: 5}}, + }) + assert.True(t, ok) + assert.Equal(t, 5, calls) + + calls, ok = gitHubAPIRateLimitCostEstimate([]*GitHubAPIRateLimitReport{ + {Start: &GitHubAPIRateLimitState{Used: 10}, End: &GitHubAPIRateLimitState{Used: 25}}, + {Start: &GitHubAPIRateLimitState{Used: 0}, End: &GitHubAPIRateLimitState{Used: 5}}, + }) + assert.True(t, ok) + assert.Equal(t, 20, calls) + + // A window reset can also happen while End.Used is still >= Start.Used (the + // window reset and then accumulated enough new usage to exceed the old + // value by coincidence). The differing Reset timestamps must still select + // the End.Used fallback instead of silently mixing counters across windows. + calls, ok = gitHubAPIRateLimitCostEstimate([]*GitHubAPIRateLimitReport{ + { + Start: &GitHubAPIRateLimitState{Used: 4990, Reset: 1000}, + End: &GitHubAPIRateLimitState{Used: 5000, Reset: 2000}, + }, + }) + assert.True(t, ok) + assert.Equal(t, 5000, calls) +} + // TestDownloadWorkflowLogsReportsCollectionStatsForJSONLAndDiskCacheHits verifies // that the single-target DownloadWorkflowLogs entry point wires --cached-jsonl // discovery/download results through to the rendered collection-stats summary, @@ -945,6 +1029,133 @@ func TestDownloadWorkflowLogsReportsCollectionStatsForJSONLAndDiskCacheHits(t *t assert.Contains(t, stderr, "Runs: 2 discovered; reports: 0 downloaded, 2 skipped because cached analyses were reused") } +// TestDownloadWorkflowLogsRendersDownloadStatsForFreshDownloadWithoutCachedJSONL +// verifies that the "Download stats: ..." summary is rendered by the +// single-target DownloadWorkflowLogs entry point when no --cached-jsonl option +// is set at all -- the common case for a plain `gh aw logs` invocation. It +// drives a real (non-cached) run through prepareRunDownload / +// downloadAndTimeRunArtifacts using a fake `gh` binary on PATH so the actual +// download and stats-recording code paths execute unmocked, instead of only +// exercising logsCollectionStats/renderLogsDownloadStatsSummary directly. +func TestDownloadWorkflowLogsRendersDownloadStatsForFreshDownloadWithoutCachedJSONL(t *testing.T) { + const runID int64 = 303 + outputDir := t.TempDir() + + fakeBinDir := t.TempDir() + fakeGH := filepath.Join(fakeBinDir, "gh") + fakeGHScript := "#!/bin/sh\n" + + "if [ \"$1\" = \"api\" ]; then\n" + + " case \"$*\" in\n" + + " *artifacts*) printf '%s\\n' \"usage\" ;;\n" + + fmt.Sprintf(" *) printf '%%s\\n' '{\"id\":%d,\"status\":\"completed\",\"conclusion\":\"success\",\"repository\":{\"full_name\":\"owner/repo\"}}' ;;\n", runID) + + " esac\n" + + " exit 0\n" + + "fi\n" + + "if [ \"$1\" = \"run\" ] && [ \"$2\" = \"download\" ]; then\n" + + " dir=\"\"\n" + + " while [ $# -gt 0 ]; do\n" + + " if [ \"$1\" = \"--dir\" ]; then dir=\"$2\"; shift 2; continue; fi\n" + + " shift\n" + + " done\n" + + " mkdir -p \"$dir\"\n" + + " printf '%s' '{\"engine_id\":\"claude\"}' > \"$dir/aw_info.json\"\n" + + " printf '%s' '{\"total_tokens\":100}' > \"$dir/usage.jsonl\"\n" + + " exit 0\n" + + "fi\n" + + "exit 1\n" + require.NoError(t, os.WriteFile(fakeGH, []byte(fakeGHScript), 0o755)) + t.Setenv("PATH", fakeBinDir+string(os.PathListSeparator)+os.Getenv("PATH")) + + originalFetch := logsFetchWorkflowRunBatch + t.Cleanup(func() { logsFetchWorkflowRunBatch = originalFetch }) + batchCalls := 0 + logsFetchWorkflowRunBatch = func(_ context.Context, _ LogsDownloadOptions, _ string, _ int, _ bool) (workflowRunBatch, error) { + batchCalls++ + if batchCalls > 1 { + return workflowRunBatch{}, nil + } + return workflowRunBatch{ + runs: []WorkflowRun{{DatabaseID: runID, Repository: "owner/repo"}}, + totalFetched: 1, + batchSize: 1, + oldestFetchedCreatedAt: time.Now(), + }, nil + } + + _, stderr := captureOutput(t, func() error { + return DownloadWorkflowLogs(context.Background(), LogsDownloadOptions{ + Count: 1, + OutputDir: outputDir, + SummaryFile: "summary.json", + ArtifactSets: []string{"usage"}, + SuppressRender: true, + }) + }) + + assert.Contains(t, stderr, "Runs: 1 discovered; reports: 1 downloaded, 0 skipped because cached analyses were reused") + assert.Contains(t, stderr, "Download stats: avg") + assert.Contains(t, stderr, "avg size") +} + +// TestDownloadAndTimeRunArtifactsExcludesPreexistingBytes verifies that +// DownloadSizeBytes reflects only the bytes added by this invocation's +// download, not the whole run directory. A prior incremental/cache pass (or +// locally generated metadata already on disk) must not inflate the reported +// size, and must not make it nonzero when nothing new was actually +// transferred. +func TestDownloadAndTimeRunArtifactsExcludesPreexistingBytes(t *testing.T) { + const runID int64 = 505 + runOutputDir := t.TempDir() + + // Simulate leftover bytes from an earlier pass plus locally generated + // metadata that already exist on disk before this download runs. + preexisting := make([]byte, 5000) + require.NoError(t, os.WriteFile(filepath.Join(runOutputDir, "leftover.txt"), preexisting, 0o600)) + + const artifactPayloadSize = 123 + fakeBinDir := t.TempDir() + fakeGH := filepath.Join(fakeBinDir, "gh") + fakeGHScript := "#!/bin/sh\n" + + "if [ \"$1\" = \"api\" ]; then\n" + + " case \"$*\" in\n" + + " *artifacts*) printf '%s\\n' \"usage\" ;;\n" + + fmt.Sprintf(" *) printf '%%s\\n' '{\"id\":%d,\"status\":\"completed\",\"conclusion\":\"success\",\"repository\":{\"full_name\":\"owner/repo\"}}' ;;\n", runID) + + " esac\n" + + " exit 0\n" + + "fi\n" + + "if [ \"$1\" = \"run\" ] && [ \"$2\" = \"download\" ]; then\n" + + " dir=\"\"\n" + + " while [ $# -gt 0 ]; do\n" + + " if [ \"$1\" = \"--dir\" ]; then dir=\"$2\"; shift 2; continue; fi\n" + + " shift\n" + + " done\n" + + " mkdir -p \"$dir\"\n" + + " printf '%s' '{\"engine_id\":\"claude\"}' > \"$dir/aw_info.json\"\n" + + fmt.Sprintf(" head -c %d /dev/zero > \"$dir/usage.jsonl\"\n", artifactPayloadSize) + + " exit 0\n" + + "fi\n" + + "exit 1\n" + require.NoError(t, os.WriteFile(fakeGH, []byte(fakeGHScript), 0o755)) + t.Setenv("PATH", fakeBinDir+string(os.PathListSeparator)+os.Getenv("PATH")) + + result := &DownloadResult{RunAnalysis: RunAnalysis{Run: WorkflowRun{DatabaseID: runID, Repository: "owner/repo"}}} + params := concurrentRunDownloadParams{ + outputDir: filepath.Dir(runOutputDir), + artifactFilter: []string{"usage"}, + dlOwner: "owner", + dlRepo: "repo", + } + downloadAndTimeRunArtifacts(context.Background(), result.Run, runOutputDir, params, params, result) + + assert.Greater(t, result.Run.DownloadDuration, time.Duration(0)) + // The preexisting leftover.txt (5000 bytes) must not be counted: a naive + // whole-directory measurement would report at least 5000 bytes, but the + // actual artifact payload downloaded here is much smaller. + assert.Positive(t, result.Run.DownloadSizeBytes) + assert.Less(t, result.Run.DownloadSizeBytes, int64(len(preexisting)), + "DownloadSizeBytes must exclude the preexisting leftover.txt bytes already on disk before this download") +} + // TestDownloadWorkflowLogsFromStdinFiltersCachedJSONLByDateRange verifies that // --stdin honors --start-date/--end-date by pruning out-of-range cached run // records, mirroring the discovery-mode behavior in DownloadWorkflowLogs. diff --git a/pkg/cli/logs_report.go b/pkg/cli/logs_report.go index 47388710b35..f8593c9b1c1 100644 --- a/pkg/cli/logs_report.go +++ b/pkg/cli/logs_report.go @@ -202,7 +202,18 @@ type RunData struct { Experiments *ExperimentData `json:"experiments,omitempty" console:"-"` // A/B experiment assignments for this run Graders *GradersData `json:"graders,omitempty" console:"-"` // Deterministic grader results for this run SafeOutputs []CreatedItemReport `json:"safe_outputs,omitempty" console:"-"` // Entities affected by safe-output handlers - awInfo *AwInfo + // DownloadDurationMS is the wall-clock time (milliseconds) spent by `gh aw logs` + // downloading this run's artifacts from GitHub. Zero means no download duration + // was recorded for this invocation (for example, an on-disk cache hit); when the + // record comes from a cached JSONL file, this preserves whatever value was + // measured the first time the run was downloaded, since cached JSONL records are + // carried through unchanged rather than re-measured. + DownloadDurationMS int64 `json:"download_duration_ms,omitempty" console:"-"` + // DownloadSizeBytes is the total on-disk size of the artifacts downloaded for + // this run. See DownloadDurationMS for how this value behaves for cache hits and + // cached JSONL records. + DownloadSizeBytes int64 `json:"download_size_bytes,omitempty" console:"-"` + awInfo *AwInfo } // logsAggregate accumulates cross-run totals while runs are converted to RunData. @@ -440,6 +451,8 @@ func applyAwInfoToRunData(runData *RunData, awInfo *AwInfo) { } func applyGitHubMetadataToRunData(runData *RunData, run WorkflowRun) { + runData.DownloadDurationMS = run.DownloadDuration.Milliseconds() + runData.DownloadSizeBytes = run.DownloadSizeBytes if run.Repository != "" { runData.Repository = run.Repository } diff --git a/pkg/cli/logs_run_processor.go b/pkg/cli/logs_run_processor.go index 430881199d4..89521ae8ca6 100644 --- a/pkg/cli/logs_run_processor.go +++ b/pkg/cli/logs_run_processor.go @@ -368,43 +368,10 @@ func processSingleRunDownload( result, ok, err := prepareRunDownload(ctx, run, runOutputDir, perRunParams, params.storageLimit) if err != nil { handleArtifactDownloadError(result, err, params.verbose) + } else if !ok { + downloadAndTimeRunArtifacts(ctx, run, runOutputDir, perRunParams, params, result) } else { - if !ok { - writeWorkflowRunFolderLocation(run.DatabaseID, runOutputDir) - logsOrchestratorLog.Printf("Downloading artifacts for run %d: owner=%s, repo=%s", run.DatabaseID, perRunParams.dlOwner, perRunParams.dlRepo) - err := params.storageLimit.runDownloadDeferredReserved(ctx, runOutputDir, func() error { - if err := waitForConfiguredRateLimit(ctx, params.verbose, params.maxGitHubAPIRateLimit, logsRunPreflightAPIReserve, params.rateLimitState); err != nil { - return err - } - - if err := os.MkdirAll(runOutputDir, constants.DirPermSensitive); err != nil { - return fmt.Errorf("failed to create run output directory: %w", err) - } - if metadata, err := fetchAndCacheWorkflowRunMetadata(ctx, run, runOutputDir, perRunParams.dlOwner, perRunParams.dlRepo, perRunParams.dlHost, params.verbose); err != nil { - logsOrchestratorLog.Printf("Failed to fetch workflow run metadata for run %d: %v", run.DatabaseID, err) - } else { - applyWorkflowRunMetadata(&result.Run, metadata) - } - if err := downloadRunArtifacts(ctx, downloadArtifactsOptions{runID: run.DatabaseID, outputDir: runOutputDir, verbose: params.verbose, owner: perRunParams.dlOwner, repo: perRunParams.dlRepo, hostname: perRunParams.dlHost, artifactFilter: params.artifactFilter}); err != nil { - return err - } - // When evals are requested but not found in the usage artifact (older runs - // that predate the conclusion-job copy), fall back to the dedicated evals - // artifact so those runs are not silently skipped. This applies both when - // --evals is set and when --artifacts evals was explicitly listed. - if params.evalsArtifactRequested && !runHasEvals(runOutputDir, params.verbose) { - tryDownloadEvalsArtifactFallback(ctx, run.DatabaseID, runOutputDir, perRunParams) - } - analyzeRunArtifacts(ctx, result, runOutputDir, params.verbose, params.artifactFilter) - return nil - }) - - if err != nil { - handleArtifactDownloadError(result, err, params.verbose) - } - } else { - logsOrchestratorLog.Printf("Cache hit for run %d, using cached summary", run.DatabaseID) - } + logsOrchestratorLog.Printf("Cache hit for run %d, using cached summary", run.DatabaseID) } completed := completedCount.Add(1) @@ -414,6 +381,75 @@ func processSingleRunDownload( return *result, nil } +// downloadAndTimeRunArtifacts performs the actual artifact download for a run +// that was not served from cache, recording wall-clock duration and resulting +// on-disk size onto result.Run so that end-of-run download stats can be +// aggregated across all runs processed in this invocation. +func downloadAndTimeRunArtifacts( + ctx context.Context, + run WorkflowRun, + runOutputDir string, + perRunParams concurrentRunDownloadParams, + params concurrentRunDownloadParams, + result *DownloadResult, +) { + writeWorkflowRunFolderLocation(run.DatabaseID, runOutputDir) + logsOrchestratorLog.Printf("Downloading artifacts for run %d: owner=%s, repo=%s", run.DatabaseID, perRunParams.dlOwner, perRunParams.dlRepo) + // sizeBefore/downloadStart bracket only the actual GitHub download calls + // (metadata fetch, artifact download, evals fallback) rather than the + // preceding rate-limit wait or the trailing artifact analysis, so a long + // quota wait or CPU-heavy analysis is never misreported as download + // latency, and pre-existing bytes on disk (from an earlier incremental or + // cache pass) are excluded from the reported size. + var downloadStart time.Time + var sizeBefore int64 + err := params.storageLimit.runDownloadDeferredReserved(ctx, runOutputDir, func() error { + if err := waitForConfiguredRateLimit(ctx, params.verbose, params.maxGitHubAPIRateLimit, logsRunPreflightAPIReserve, params.rateLimitState); err != nil { + return err + } + + if err := os.MkdirAll(runOutputDir, constants.DirPermSensitive); err != nil { + return fmt.Errorf("failed to create run output directory: %w", err) + } + if size, sizeErr := logsDirectorySize(runOutputDir); sizeErr == nil { + sizeBefore = size + } else { + logsOrchestratorLog.Printf("failed to compute pre-download size for run %d: %v", run.DatabaseID, sizeErr) + } + downloadStart = time.Now() + if metadata, err := fetchAndCacheWorkflowRunMetadata(ctx, run, runOutputDir, perRunParams.dlOwner, perRunParams.dlRepo, perRunParams.dlHost, params.verbose); err != nil { + logsOrchestratorLog.Printf("Failed to fetch workflow run metadata for run %d: %v", run.DatabaseID, err) + } else { + applyWorkflowRunMetadata(&result.Run, metadata) + } + if err := downloadRunArtifacts(ctx, downloadArtifactsOptions{runID: run.DatabaseID, outputDir: runOutputDir, verbose: params.verbose, owner: perRunParams.dlOwner, repo: perRunParams.dlRepo, hostname: perRunParams.dlHost, artifactFilter: params.artifactFilter}); err != nil { + return err + } + // When evals are requested but not found in the usage artifact (older runs + // that predate the conclusion-job copy), fall back to the dedicated evals + // artifact so those runs are not silently skipped. This applies both when + // --evals is set and when --artifacts evals was explicitly listed. + if params.evalsArtifactRequested && !runHasEvals(runOutputDir, params.verbose) { + tryDownloadEvalsArtifactFallback(ctx, run.DatabaseID, runOutputDir, perRunParams) + } + result.Run.DownloadDuration = time.Since(downloadStart) + if size, sizeErr := logsDirectorySize(runOutputDir); sizeErr == nil { + if size > sizeBefore { + result.Run.DownloadSizeBytes = size - sizeBefore + } + } else { + logsOrchestratorLog.Printf("failed to compute download size for run %d: %v", run.DatabaseID, sizeErr) + } + analyzeRunArtifacts(ctx, result, runOutputDir, params.verbose, params.artifactFilter) + return nil + }) + + if err != nil { + handleArtifactDownloadError(result, err, params.verbose) + return + } +} + func writeWorkflowRunFolderLocation(runID int64, runOutputDir string) { runFolder, err := filepath.Abs(runOutputDir) if err != nil { diff --git a/schemas/logs-jsonl.schema.json b/schemas/logs-jsonl.schema.json index 56710de4e9d..bf1586613e5 100644 --- a/schemas/logs-jsonl.schema.json +++ b/schemas/logs-jsonl.schema.json @@ -983,6 +983,12 @@ "additionalProperties": false } }, + "download_duration_ms": { + "type": "integer" + }, + "download_size_bytes": { + "type": "integer" + }, "engine_version": { "type": "string" }, diff --git a/schemas/logs.schema.json b/schemas/logs.schema.json index fdee5eca964..c8b9e9fa519 100644 --- a/schemas/logs.schema.json +++ b/schemas/logs.schema.json @@ -1101,6 +1101,12 @@ ], "additionalProperties": false } + }, + "download_duration_ms": { + "type": "integer" + }, + "download_size_bytes": { + "type": "integer" } }, "required": [