Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 0 additions & 5 deletions pkg/cli/drain3_train.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
2 changes: 2 additions & 0 deletions pkg/cli/drain3_train_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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))
}

Expand Down
8 changes: 7 additions & 1 deletion pkg/cli/logs_cached_json.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
1 change: 1 addition & 0 deletions pkg/cli/logs_multi.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()))
}
Expand Down
48 changes: 48 additions & 0 deletions pkg/cli/logs_multi_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 })
Expand Down
1 change: 1 addition & 0 deletions pkg/cli/logs_orchestrator.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
40 changes: 40 additions & 0 deletions pkg/cli/logs_orchestrator_download.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"path/filepath"
"strings"
"sync"
"sync/atomic"
"time"

"github.com/github/gh-aw/pkg/console"
Expand All @@ -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
Expand All @@ -50,6 +86,7 @@ type processWorkflowRunBatchOptions struct {
cachedRuns cachedLogsRuns
cachedJSONLWriter *cachedLogsJSONLWriter
countLimit *logsCountLimit
collectionStats *logsCollectionStats
}

func prepareLogsDownload(ctx context.Context, opts LogsDownloadOptions) (logsDownloadRuntime, error) {
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions pkg/cli/logs_orchestrator_stdin.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down
1 change: 1 addition & 0 deletions pkg/cli/logs_orchestrator_types.go
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@ type LogsDownloadOptions struct {
inheritTimeoutContext bool
cachedJSONLWriter *cachedLogsJSONLWriter
cachedJSONLCache *cachedLogsJSONLCache
collectionStats *logsCollectionStats
}

type workflowLogsResult struct {
Expand Down
Loading
Loading