diff --git a/backend/plugins/github/tasks/cicd_run_collector.go b/backend/plugins/github/tasks/cicd_run_collector.go index fe19da43eaf..91e3f87d788 100644 --- a/backend/plugins/github/tasks/cicd_run_collector.go +++ b/backend/plugins/github/tasks/cicd_run_collector.go @@ -48,6 +48,11 @@ const FILTERED_SEARCH_CAP = 1000 // githubTimeLayout is the ISO8601 format GitHub expects in the `created` filter. const githubTimeLayout = "2006-01-02T15:04:05Z" +// maxUnfinishedRunLookback bounds how far back the low-water-mark query can reach +// for unfinished runs. GitHub Actions cancels workflows after 35 days max; 14 days +// safely catches long runs while preventing zombie runs from pinning the window forever. +const maxUnfinishedRunLookback = 14 * 24 * time.Hour + // TimeWindow is an inclusive-both-ends range for the `/actions/runs` `created=..` query. type TimeWindow struct { From time.Time @@ -125,6 +130,21 @@ func CollectRuns(taskCtx plugin.SubTaskContext) errors.Error { sinceSource = "epoch_fullsync" } + if manager.IsIncremental() { + oldestUnfinished, err := loadOldestUnfinishedRunCreatedAt(taskCtx, data.Options.ConnectionId, data.Options.GithubId, maxUnfinishedRunLookback) + if err != nil { + return err + } + + if oldestUnfinished != nil && oldestUnfinished.Before(windowStart) { + logger.Info("cicd_run_collector: moving windowStart back from %s to %s to recheck unfinished runs", + windowStart.Format(githubTimeLayout), + oldestUnfinished.Format(githubTimeLayout)) + windowStart = oldestUnfinished.Truncate(time.Second) + sinceSource = "tool_runs_unfinished_lwm" + } + } + logger.Info("cicd_run_collector: collecting workflow runs in [%s, %s] (incremental=%v, since_source=%s)", windowStart.Format(githubTimeLayout), until.Format(githubTimeLayout), @@ -166,6 +186,32 @@ func loadLatestRunUpdatedAt(taskCtx plugin.SubTaskContext, connectionId uint64, return &fallback, nil } +func loadOldestUnfinishedRunCreatedAt(taskCtx plugin.SubTaskContext, connectionId uint64, repoId int, maxLookback time.Duration) (*time.Time, errors.Error) { + db := taskCtx.GetDal() + oldest := &models.GithubRun{} + clauses := []dal.Clause{ + dal.Where("connection_id = ? AND repo_id = ? AND status != ? AND github_created_at IS NOT NULL", connectionId, repoId, "completed"), + dal.Orderby("github_created_at ASC"), + dal.Limit(1), + } + if maxLookback > 0 { + cutoff := time.Now().Add(-maxLookback) + clauses = append(clauses, dal.Where("github_created_at >= ?", cutoff)) + } + err := db.First(oldest, clauses...) + if err != nil { + if db.IsErrorNotFound(err) { + return nil, nil + } + return nil, err + } + if oldest.GithubCreatedAt == nil { + return nil, nil + } + t := oldest.GithubCreatedAt.UTC() + return &t, nil +} + // buildRunsQuery assembles the filtered-mode query for a single leaf TimeWindow. // Shared between registerCollectorForLeafWindows and tests. func buildRunsQuery(reqData *helper.RequestData) (url.Values, errors.Error) { @@ -209,19 +255,17 @@ func registerCollectorForLeafWindows( if len(body.WorkflowRuns) == 0 { return nil, nil } - // Range is already bounded in filtered mode; only keep completed runs. - filtered := make([]json.RawMessage, 0, len(body.WorkflowRuns)) + // Range is already bounded in filtered mode; persist all runs so in-progress + // and queued runs are tracked in _tool_github_runs and updated upon completion. + rawRuns := make([]json.RawMessage, 0, len(body.WorkflowRuns)) for _, run := range body.WorkflowRuns { - if run.Status != "completed" { - continue - } runJSON, err := json.Marshal(run) if err != nil { return nil, errors.Convert(err) } - filtered = append(filtered, json.RawMessage(runJSON)) + rawRuns = append(rawRuns, json.RawMessage(runJSON)) } - return filtered, nil + return rawRuns, nil }, }) } diff --git a/backend/plugins/github/tasks/cicd_run_collector_test.go b/backend/plugins/github/tasks/cicd_run_collector_test.go index c131c6cc6b4..deb73524427 100644 --- a/backend/plugins/github/tasks/cicd_run_collector_test.go +++ b/backend/plugins/github/tasks/cicd_run_collector_test.go @@ -92,6 +92,58 @@ func TestCicdRunLoadLatestRunUpdatedAt_PropagatesError(t *testing.T) { mockDal.AssertExpectations(t) } +func TestCicdRunLoadOldestUnfinishedRunCreatedAt_ReturnsOldestTimestamp(t *testing.T) { + connectionId := uint64(1) + repoId := 101 + oldestTs := time.Date(2025, 4, 1, 10, 11, 12, 0, time.UTC) + + mockDal := new(mockdal.Dal) + mockDal.On("First", mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything).Run(func(args mock.Arguments) { + dst := args.Get(0).(*models.GithubRun) + dst.ConnectionId = connectionId + dst.RepoId = repoId + dst.Status = "in_progress" + dst.GithubCreatedAt = &oldestTs + }).Return(nil).Once() + + ctx := unithelper.DummySubTaskContext(mockDal) + created, err := loadOldestUnfinishedRunCreatedAt(ctx, connectionId, repoId, 14*24*time.Hour) + + require.Nil(t, err) + require.NotNil(t, created) + assert.True(t, created.Equal(oldestTs)) + mockDal.AssertExpectations(t) +} + +func TestCicdRunLoadOldestUnfinishedRunCreatedAt_NotFoundReturnsNil(t *testing.T) { + mockDal := new(mockdal.Dal) + notFoundErr := errors.Default.New("record not found") + mockDal.On("First", mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return(notFoundErr).Once() + mockDal.On("IsErrorNotFound", notFoundErr).Return(true).Once() + + ctx := unithelper.DummySubTaskContext(mockDal) + created, err := loadOldestUnfinishedRunCreatedAt(ctx, 1, 101, 14*24*time.Hour) + + require.Nil(t, err) + assert.Nil(t, created) + mockDal.AssertExpectations(t) +} + +func TestCicdRunLoadOldestUnfinishedRunCreatedAt_PropagatesError(t *testing.T) { + mockDal := new(mockdal.Dal) + dbErr := errors.Default.New("db unavailable") + mockDal.On("First", mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return(dbErr).Once() + mockDal.On("IsErrorNotFound", dbErr).Return(false).Once() + + ctx := unithelper.DummySubTaskContext(mockDal) + created, err := loadOldestUnfinishedRunCreatedAt(ctx, 1, 101, 14*24*time.Hour) + + assert.Nil(t, created) + require.NotNil(t, err) + assert.Contains(t, err.Error(), "db unavailable") + mockDal.AssertExpectations(t) +} + // newTestBuilder constructs a leafWindowBuilder with a stubbed probe for unit testing. func newTestBuilder(probe probeFunc) *leafWindowBuilder { mockDal := new(mockdal.Dal)