diff --git a/pkg/cli/drain3_train.go b/pkg/cli/drain3_train.go index e02d7ea1a3e..44227da0f52 100644 --- a/pkg/cli/drain3_train.go +++ b/pkg/cli/drain3_train.go @@ -107,10 +107,5 @@ func writeDrain3Weights(coordinator *agentdrain.Coordinator, outputDir string) e return fmt.Errorf("log pattern training: write weights file: %w", err) } fmt.Fprintln(os.Stderr, console.FormatSuccessMessage("Log pattern weights written to: "+outputPath)) - fmt.Fprintln(os.Stderr, console.FormatInfoMessage( - "To embed these weights as default, copy the file and rebuild:\n"+ - " cp "+outputPath+" pkg/agentdrain/data/default_weights.json\n"+ - " make build", - )) return nil } diff --git a/pkg/cli/drain3_train_test.go b/pkg/cli/drain3_train_test.go index b006a60a0d8..ba597c64642 100644 --- a/pkg/cli/drain3_train_test.go +++ b/pkg/cli/drain3_train_test.go @@ -125,6 +125,8 @@ func TestTrainDrain3Weights_LoadsExistingWeights(t *testing.T) { }) assert.Contains(t, stderr, "Loaded log pattern weights from: "+weightsPath) + assert.NotContains(t, stderr, "To embed these weights as default") + assert.NotContains(t, stderr, "make build") assert.FileExists(t, filepath.Join(outputDir, drain3WeightsFilename)) } diff --git a/pkg/cli/logs_cached_json.go b/pkg/cli/logs_cached_json.go index 98507b8d359..ff2a2a3e35a 100644 --- a/pkg/cli/logs_cached_json.go +++ b/pkg/cli/logs_cached_json.go @@ -137,7 +137,13 @@ func loadCachedLogsJSONL(path string) (*cachedLogsJSONLCache, error) { } func prepareCachedLogsJSONL(opts *LogsDownloadOptions) error { - if opts.CachedJSONL == "" || opts.cachedJSONLWriter != nil { + if opts.CachedJSONL == "" { + return nil + } + if opts.collectionStats == nil { + opts.collectionStats = &logsCollectionStats{} + } + if opts.cachedJSONLWriter != nil { return nil } cache, err := loadCachedLogsJSONL(opts.CachedJSONL) diff --git a/pkg/cli/logs_multi.go b/pkg/cli/logs_multi.go index 8065235aa7c..3c9326b6d3e 100644 --- a/pkg/cli/logs_multi.go +++ b/pkg/cli/logs_multi.go @@ -165,6 +165,7 @@ func DownloadWorkflowLogsForTargets( //nolint:largefunc // Keeps shared collecti allAPIRateLimits := startGitHubAPIRateLimitReports(activeCtx, logsTargetRateLimitHosts(targets)) results := collectLogsTargets(activeCtx, opts, targets) processedRuns, continuations, timeoutReached, countLimitReached, storageLimitReached, allErrors := mergeLogsTargetResults(results, initialErrors) + renderLogsCollectionStats(opts.collectionStats) for _, err := range allErrors { fmt.Fprintln(os.Stderr, console.FormatWarningMessage("Skipping workflow target: "+err.Error())) } diff --git a/pkg/cli/logs_multi_test.go b/pkg/cli/logs_multi_test.go index 182ce2f0712..d00edb82871 100644 --- a/pkg/cli/logs_multi_test.go +++ b/pkg/cli/logs_multi_test.go @@ -120,6 +120,54 @@ func TestDownloadWorkflowLogsForTargetsConcurrentAndResilient(t *testing.T) { assert.Equal(t, int64(123), report.Continuations[0].BeforeRunID) } +// TestDownloadWorkflowLogsForTargetsReportsCollectionStatsAcrossTargets verifies +// that DownloadWorkflowLogsForTargets renders one combined collection-stats +// summary reflecting every target's recorded discovered/downloaded/cached +// counts, confirming the shared collectionStats pointer is wired through to +// each concurrent target and rendered once after they finish. +func TestDownloadWorkflowLogsForTargetsReportsCollectionStatsAcrossTargets(t *testing.T) { + original := collectWorkflowLogsForTarget + t.Cleanup(func() { collectWorkflowLogsForTarget = original }) + collectWorkflowLogsForTarget = func(_ context.Context, opts LogsDownloadOptions) (workflowLogsResult, error) { + opts.collectionStats.recordDiscovered(2) + opts.collectionStats.recordResult(DownloadResult{}) + opts.collectionStats.recordResult(DownloadResult{Cached: true}) + return workflowLogsResult{ + processedRuns: []ProcessedRun{{ + Run: WorkflowRun{ + DatabaseID: int64(len(opts.WorkflowName)), + WorkflowName: opts.WorkflowName, + CreatedAt: time.Now(), + LogsPath: filepath.Join(opts.OutputDir, "run-1"), + }, + }}, + }, nil + } + + tempDir := t.TempDir() + originalDir, err := os.Getwd() + require.NoError(t, err) + require.NoError(t, os.Chdir(tempDir)) + t.Cleanup(func() { _ = os.Chdir(originalDir) }) + + cachedJSONL := filepath.Join(tempDir, "logs.jsonl") + require.NoError(t, os.WriteFile(cachedJSONL, nil, 0o600)) + + _, stderr := captureOutput(t, func() error { + return DownloadWorkflowLogsForTargets(context.Background(), LogsDownloadOptions{ + OutputDir: filepath.Join(tempDir, "logs"), + SummaryFile: "summary.json", + SuppressRender: true, + CachedJSONL: cachedJSONL, + }, []logsWorkflowTarget{ + {workflowName: "alpha", repoOverride: "org/repo-a"}, + {workflowName: "bravo", repoOverride: "org/repo-b"}, + }, nil) + }) + + assert.Contains(t, stderr, "Runs: 4 discovered; reports: 2 downloaded, 2 skipped because cached analyses were reused") +} + func TestDownloadWorkflowLogsForTargetsReturnsErrorWhenAllFail(t *testing.T) { original := collectWorkflowLogsForTarget t.Cleanup(func() { collectWorkflowLogsForTarget = original }) diff --git a/pkg/cli/logs_orchestrator.go b/pkg/cli/logs_orchestrator.go index 1bddfa180a0..b78412d5ee3 100644 --- a/pkg/cli/logs_orchestrator.go +++ b/pkg/cli/logs_orchestrator.go @@ -310,6 +310,7 @@ func DownloadWorkflowLogs(ctx context.Context, opts LogsDownloadOptions) (err er if err != nil { return err } + renderLogsCollectionStats(opts.collectionStats) finishGitHubAPIRateLimitReport(ctx, apiRateLimit, opts.JSONOutput) cacheGitHubAPIRateLimitReports(opts.cachedJSONLWriter, apiRateLimit) if handled, err := handleEmptyProcessedRuns(result.processedRuns, opts, result.timeoutReached, result.storageLimitReached, result.continuation, nil, apiRateLimit, nil); handled || err != nil { diff --git a/pkg/cli/logs_orchestrator_download.go b/pkg/cli/logs_orchestrator_download.go index 25333c3ea66..169cbc0e61a 100644 --- a/pkg/cli/logs_orchestrator_download.go +++ b/pkg/cli/logs_orchestrator_download.go @@ -8,6 +8,7 @@ import ( "path/filepath" "strings" "sync" + "sync/atomic" "time" "github.com/github/gh-aw/pkg/console" @@ -34,6 +35,41 @@ type workflowRunBatch struct { oldestFetchedCreatedAt time.Time } +type logsCollectionStats struct { + discoveredRuns atomic.Int64 + downloadedReports atomic.Int64 + cachedReports atomic.Int64 +} + +func (s *logsCollectionStats) recordDiscovered(count int) { + if s != nil { + s.discoveredRuns.Add(int64(count)) + } +} + +func (s *logsCollectionStats) recordResult(result DownloadResult) { + if s == nil { + return + } + if result.Cached { + s.cachedReports.Add(1) + return + } + if !result.Skipped && result.Error == nil { + s.downloadedReports.Add(1) + } +} + +func renderLogsCollectionStats(stats *logsCollectionStats) { + if stats == nil { + return + } + fmt.Fprintln(os.Stderr, console.FormatInfoMessage(fmt.Sprintf( + "Runs: %d discovered; reports: %d downloaded, %d skipped because cached analyses were reused", + stats.discoveredRuns.Load(), stats.downloadedReports.Load(), stats.cachedReports.Load(), + ))) +} + type processWorkflowRunBatchOptions struct { count int outputDir string @@ -50,6 +86,7 @@ type processWorkflowRunBatchOptions struct { cachedRuns cachedLogsRuns cachedJSONLWriter *cachedLogsJSONLWriter countLimit *logsCountLimit + collectionStats *logsCollectionStats } func prepareLogsDownload(ctx context.Context, opts LogsDownloadOptions) (logsDownloadRuntime, error) { @@ -355,6 +392,7 @@ func fetchAndProcessLogsBatch(state *logsCollectionState, runtime logsDownloadRu if err != nil { return handleLogsBatchError(state, runtime.fetchAllInRange, opts.countLimit, err) } + opts.collectionStats.recordDiscovered(len(batch.runs)) if len(batch.runs) == 0 { cursor, shouldContinue, shouldStop := handleEmptyWorkflowRunBatch(batch, opts.Verbose) if shouldStop { @@ -384,6 +422,7 @@ func fetchAndProcessLogsBatch(state *logsCollectionState, runtime logsDownloadRu cachedRuns: runtime.cachedRuns, cachedJSONLWriter: opts.cachedJSONLWriter, countLimit: opts.countLimit, + collectionStats: opts.collectionStats, }) state.timeoutReached = state.timeoutReached || batchTimedOut logProcessedWorkflowRunBatch(opts, runtime.fetchAllInRange, state.iteration, batchProcessed, len(state.processedRuns), opts.Verbose) @@ -712,6 +751,7 @@ func (c *orderedLogsRunCollector) recordResult(index int, result DownloadResult) func (c *orderedLogsRunCollector) processReadyResult(index int) { result := c.pendingResults[index] + c.opts.collectionStats.recordResult(result) if c.processedBase+c.acceptedCount >= c.opts.count || c.opts.countLimit.isReached() { finalizeLogsRunDownload(c.opts.storageLimit, result) return diff --git a/pkg/cli/logs_orchestrator_stdin.go b/pkg/cli/logs_orchestrator_stdin.go index c67f5dc3e0f..304bf199d51 100644 --- a/pkg/cli/logs_orchestrator_stdin.go +++ b/pkg/cli/logs_orchestrator_stdin.go @@ -209,11 +209,14 @@ func DownloadWorkflowLogsFromStdin(ctx context.Context, opts StdinLogsOptions) ( gradersOnly: opts.GradersOnly, } downloadResults := downloadRunArtifactsConcurrent(ctx, runs, runArtifactsConcurrentOptions{outputDir: opts.OutputDir, verbose: opts.Verbose, maxRuns: len(runs), repoOverride: opts.RepoOverride, artifactFilter: artifactFilter, evalsOnly: opts.EvalsOnly, artifactSets: opts.ArtifactSets, storageLimit: storageLimit, cachedRuns: cachedRuns, filters: filters}) + collectionStats := &logsCollectionStats{} + collectionStats.recordDiscovered(len(runs)) // Process download results applying the same filters as DownloadWorkflowLogs. var processedRuns []ProcessedRun var storageLimitReached bool for _, result := range downloadResults { + collectionStats.recordResult(result) if result.CachedRun != nil { processedRuns = append(processedRuns, processedRunFromCachedData(*result.CachedRun)) continue @@ -269,6 +272,9 @@ func DownloadWorkflowLogsFromStdin(ctx context.Context, opts StdinLogsOptions) ( } finalizeLogsRunDownload(storageLimit, result) } + if opts.CachedJSONL != "" { + renderLogsCollectionStats(collectionStats) + } if len(processedRuns) == 0 { finishGitHubAPIRateLimitReports(ctx, allAPIRateLimits, opts.JSONOutput) diff --git a/pkg/cli/logs_orchestrator_types.go b/pkg/cli/logs_orchestrator_types.go index 5519f35ae8c..d96892310e4 100644 --- a/pkg/cli/logs_orchestrator_types.go +++ b/pkg/cli/logs_orchestrator_types.go @@ -69,6 +69,7 @@ type LogsDownloadOptions struct { inheritTimeoutContext bool cachedJSONLWriter *cachedLogsJSONLWriter cachedJSONLCache *cachedLogsJSONLCache + collectionStats *logsCollectionStats } type workflowLogsResult struct { diff --git a/pkg/cli/logs_orchestrator_unit_test.go b/pkg/cli/logs_orchestrator_unit_test.go index e31f5eb6380..77d81189d5b 100644 --- a/pkg/cli/logs_orchestrator_unit_test.go +++ b/pkg/cli/logs_orchestrator_unit_test.go @@ -7,6 +7,7 @@ import ( "fmt" "os" "path/filepath" + "strconv" "testing" "time" @@ -806,6 +807,105 @@ func TestCurrentLogsGuardrailStatusReportsAllBoundaries(t *testing.T) { assert.LessOrEqual(t, status.timeoutRemaining, time.Minute) } +func TestLogsCollectionStatsReportsDiscoveredDownloadedAndCached(t *testing.T) { + stats := &logsCollectionStats{} + stats.recordDiscovered(4) + stats.recordResult(DownloadResult{}) + stats.recordResult(DownloadResult{Cached: true}) + stats.recordResult(DownloadResult{Cached: true, CachedRun: &RunData{RunID: 42}}) + stats.recordResult(DownloadResult{Skipped: true}) + + _, stderr := captureOutput(t, func() error { + renderLogsCollectionStats(stats) + return nil + }) + + assert.Contains(t, stderr, "Runs: 4 discovered; reports: 1 downloaded, 2 skipped because cached analyses were reused") +} + +// TestDownloadWorkflowLogsReportsCollectionStatsForJSONLAndDiskCacheHits verifies +// that the single-target DownloadWorkflowLogs entry point wires --cached-jsonl +// discovery/download results through to the rendered collection-stats summary, +// counting both a JSONL cache hit and an on-disk (run_summary.json) cache hit +// toward the "skipped" total. This mirrors the entry-point-level check +// requested in review: mocking only the batch-discovery indirection point +// (logsFetchWorkflowRunBatch) so the real cache-lookup and stats-recording +// code in the download pipeline executes unmocked. +func TestDownloadWorkflowLogsReportsCollectionStatsForJSONLAndDiskCacheHits(t *testing.T) { + originalFetch := logsFetchWorkflowRunBatch + t.Cleanup(func() { logsFetchWorkflowRunBatch = originalFetch }) + + outputDir := t.TempDir() + + // Run 1: on-disk cache hit — a complete run_summary.json plus artifact marker + // already exists locally, so the download pipeline reuses it without any + // --cached-jsonl involvement. + const diskCachedRunID int64 = 101 + diskRunDir := filepath.Join(outputDir, fmt.Sprintf("run-%d", diskCachedRunID)) + require.NoError(t, os.MkdirAll(diskRunDir, 0o755)) + require.NoError(t, os.WriteFile( + filepath.Join(diskRunDir, runAPIResponseFileName), + []byte(fmt.Sprintf(`{"id":%d,"status":"completed","conclusion":"success"}`, diskCachedRunID)), + 0o600, + )) + require.NoError(t, saveRunSummary(diskRunDir, &RunSummary{ + CLIVersion: GetVersion(), + RunID: diskCachedRunID, + ProcessedAt: time.Now(), + RunAnalysis: RunAnalysis{ + Run: WorkflowRun{DatabaseID: diskCachedRunID, WorkflowName: "Disk Cached", Status: "completed", Conclusion: "success"}, + }, + }, false)) + require.NoError(t, markArtifactDownloaded(diskRunDir, constants.UsageArtifactName.String())) + + // Run 2: JSONL cache hit — the run is only ever known via --cached-jsonl. + const jsonlCachedRunID int64 = 202 + jsonlUpdatedAt := time.Now().Add(-time.Hour).Truncate(time.Second) + cachedJSONLPath := filepath.Join(outputDir, "cached-logs.jsonl") + cachedRecord := fmt.Sprintf( + `{"schema_version":2,"kind":"run","run":{"run_id":%d,"status":"completed","conclusion":"success","run_attempt":"1","updated_at":%q,"repository":"owner/repo"}}`+"\n", + jsonlCachedRunID, jsonlUpdatedAt.Format(time.RFC3339), + ) + require.NoError(t, os.WriteFile(cachedJSONLPath, []byte(cachedRecord), 0o600)) + + batchCalls := 0 + logsFetchWorkflowRunBatch = func(_ context.Context, _ LogsDownloadOptions, _ string, _ int, _ bool) (workflowRunBatch, error) { + batchCalls++ + if batchCalls > 1 { + return workflowRunBatch{}, nil + } + return workflowRunBatch{ + runs: []WorkflowRun{ + {DatabaseID: diskCachedRunID}, + { + DatabaseID: jsonlCachedRunID, + Status: "completed", + Conclusion: "success", + Attempt: 1, + UpdatedAt: jsonlUpdatedAt, + Repository: "owner/repo", + }, + }, + totalFetched: 2, + batchSize: 2, + oldestFetchedCreatedAt: time.Now(), + }, nil + } + + _, stderr := captureOutput(t, func() error { + return DownloadWorkflowLogs(context.Background(), LogsDownloadOptions{ + Count: 2, + OutputDir: outputDir, + SummaryFile: "summary.json", + CachedJSONL: cachedJSONLPath, + ArtifactSets: []string{"usage"}, + SuppressRender: true, + }) + }) + + assert.Contains(t, stderr, "Runs: 2 discovered; reports: 0 downloaded, 2 skipped because cached analyses were reused") +} + // TestDownloadWorkflowLogsFromStdinFiltersCachedJSONLByDateRange verifies that // --stdin honors --start-date/--end-date by pruning out-of-range cached run // records, mirroring the discovery-mode behavior in DownloadWorkflowLogs. @@ -831,3 +931,72 @@ func TestDownloadWorkflowLogsFromStdinFiltersCachedJSONLByDateRange(t *testing.T assert.NotContains(t, content, `"run_id":1`) assert.Contains(t, content, `"run_id":2`) } + +// TestDownloadWorkflowLogsFromStdinReportsCollectionStatsForJSONLAndDiskCacheHits +// verifies that the --stdin entry point (DownloadWorkflowLogsFromStdin) wires +// its collection stats through to the rendered summary, counting both a +// disk-cached run (existing run_summary.json) and a --cached-jsonl-cached run +// toward the "skipped" total. Run metadata is fetched through a fake `gh` +// binary on PATH so no live GitHub API access is required. +func TestDownloadWorkflowLogsFromStdinReportsCollectionStatsForJSONLAndDiskCacheHits(t *testing.T) { + outputDir := t.TempDir() + + const diskCachedRunID int64 = 101 + const jsonlCachedRunID int64 = 202 + jsonlUpdatedAt := time.Now().Add(-time.Hour).Truncate(time.Second) + + fakeBinDir := t.TempDir() + fakeGH := filepath.Join(fakeBinDir, "gh") + fakeGHScript := "#!/bin/sh\n" + + "case \"$*\" in\n" + + fmt.Sprintf(" *\"/runs/%d --jq\"*) cat <<'EOF'\n", diskCachedRunID) + + fmt.Sprintf(`{"databaseId":%d,"number":1,"htmlUrl":"https://github.com/owner/repo/actions/runs/%d","status":"completed","conclusion":"success","workflowName":"Disk Cached","createdAt":"2026-01-01T00:00:00Z","updatedAt":"2026-01-01T00:01:00Z","repository":"owner/repo"}`+"\n", diskCachedRunID, diskCachedRunID) + + "EOF\n" + + " ;;\n" + + fmt.Sprintf(" *\"/runs/%d --jq\"*) cat <<'EOF'\n", jsonlCachedRunID) + + fmt.Sprintf(`{"databaseId":%d,"number":2,"htmlUrl":"https://github.com/owner/repo/actions/runs/%d","status":"completed","conclusion":"success","workflowName":"JSONL Cached","attempt":1,"createdAt":"2026-01-01T00:00:00Z","updatedAt":%q,"repository":"owner/repo"}`+"\n", jsonlCachedRunID, jsonlCachedRunID, jsonlUpdatedAt.Format(time.RFC3339)) + + "EOF\n" + + " ;;\n" + + "esac\n" + require.NoError(t, os.WriteFile(fakeGH, []byte(fakeGHScript), 0o755)) + t.Setenv("PATH", fakeBinDir+string(os.PathListSeparator)+os.Getenv("PATH")) + + // Run 101: on-disk cache hit — a complete run_summary.json plus artifact + // marker already exists locally. + diskRunDir := filepath.Join(outputDir, fmt.Sprintf("run-%d", diskCachedRunID)) + require.NoError(t, os.MkdirAll(diskRunDir, 0o755)) + require.NoError(t, os.WriteFile( + filepath.Join(diskRunDir, runAPIResponseFileName), + []byte(fmt.Sprintf(`{"id":%d,"status":"completed","conclusion":"success","repository":{"full_name":"owner/repo"}}`, diskCachedRunID)), + 0o600, + )) + require.NoError(t, saveRunSummary(diskRunDir, &RunSummary{ + CLIVersion: GetVersion(), + RunID: diskCachedRunID, + ProcessedAt: time.Now(), + RunAnalysis: RunAnalysis{ + Run: WorkflowRun{DatabaseID: diskCachedRunID, WorkflowName: "Disk Cached", Status: "completed", Conclusion: "success"}, + }, + }, false)) + require.NoError(t, markArtifactDownloaded(diskRunDir, constants.UsageArtifactName.String())) + + // Run 202: JSONL cache hit — known only via --cached-jsonl. + cachedJSONLPath := filepath.Join(outputDir, "cached-logs.jsonl") + cachedRecord := fmt.Sprintf( + `{"schema_version":2,"kind":"run","run":{"run_id":%d,"status":"completed","conclusion":"success","run_attempt":"1","updated_at":%q,"repository":"owner/repo"}}`+"\n", + jsonlCachedRunID, jsonlUpdatedAt.Format(time.RFC3339), + ) + require.NoError(t, os.WriteFile(cachedJSONLPath, []byte(cachedRecord), 0o600)) + + _, stderr := captureOutput(t, func() error { + return DownloadWorkflowLogsFromStdin(context.Background(), StdinLogsOptions{ + RunURLs: []string{strconv.FormatInt(diskCachedRunID, 10), strconv.FormatInt(jsonlCachedRunID, 10)}, + OutputDir: outputDir, + RepoOverride: "owner/repo", + CachedJSONL: cachedJSONLPath, + ArtifactSets: []string{"usage"}, + }) + }) + + assert.Contains(t, stderr, "Runs: 2 discovered; reports: 0 downloaded, 2 skipped because cached analyses were reused") +}