diff --git a/backend/plugins/jenkins/e2e/builds_unfinished_test.go b/backend/plugins/jenkins/e2e/builds_unfinished_test.go new file mode 100644 index 00000000000..299a346ee6c --- /dev/null +++ b/backend/plugins/jenkins/e2e/builds_unfinished_test.go @@ -0,0 +1,73 @@ +/* +Licensed to the Apache Software Foundation (ASF) under one or more +contributor license agreements. See the NOTICE file distributed with +this work for additional information regarding copyright ownership. +The ASF licenses this file to You under the Apache License, Version 2.0 +(the "License"); you may not use this file except in compliance with +the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package e2e + +import ( + "testing" + + "github.com/apache/devlake/core/dal" + "github.com/apache/devlake/helpers/e2ehelper" + "github.com/apache/devlake/plugins/jenkins/impl" + "github.com/apache/devlake/plugins/jenkins/models" + "github.com/apache/devlake/plugins/jenkins/tasks" + "github.com/stretchr/testify/assert" +) + +var singleJobOptions = &tasks.JenkinsOptions{ + ConnectionId: 1, + JobName: `devlake`, + JobFullName: `devlake`, + JobPath: `job/`, +} + +// A single job's builds still running at the last collection are re-collected until they +// finish: the incremental list only returns builds started since that collection. +func TestJenkinsUnfinishedBuilds(t *testing.T) { + var jenkins impl.Jenkins + dataflowTester := e2ehelper.NewDataFlowTester(t, "jenkins", jenkins) + + // devlake#1 finished, devlake#2 running, other#7 running but in another job + dataflowTester.FlushTabler(&models.JenkinsBuild{}) + dataflowTester.ImportCsvIntoTabler("./raw_tables/_tool_jenkins_builds_unfinished.csv", models.JenkinsBuild{}) + + var builds []tasks.SimpleBuild + err := dataflowTester.Dal.All(&builds, tasks.UnfinishedBuilds(singleJobOptions)...) + assert.Nil(t, err) + assert.Equal(t, []tasks.SimpleBuild{{Number: "2", FullName: "devlake#2"}}, builds) +} + +// A build listed while running and re-collected once finished ends up with its final state. +func TestJenkinsUnfinishedBuildExtraction(t *testing.T) { + var jenkins impl.Jenkins + dataflowTester := e2ehelper.NewDataFlowTester(t, "jenkins", jenkins) + + // row 1: devlake#5 from the list, running; row 2: its detail, re-collected once finished + dataflowTester.ImportCsvIntoRawTable("./raw_tables/_raw_jenkins_api_builds_unfinished.csv", "_raw_jenkins_api_builds") + dataflowTester.FlushTabler(&models.JenkinsBuild{}) + dataflowTester.FlushTabler(&models.JenkinsBuildCommit{}) + dataflowTester.Subtask(tasks.ExtractApiBuildsMeta, &tasks.JenkinsTaskData{Options: singleJobOptions}) + + var build models.JenkinsBuild + err := dataflowTester.Dal.First(&build, dal.Where("full_name = ?", "devlake#5")) + assert.Nil(t, err) + assert.False(t, build.Building) + assert.Equal(t, "SUCCESS", build.Result) + assert.Equal(t, float64(90000000), build.Duration) + assert.Equal(t, "devlake", build.JobName) + assert.Equal(t, "job/", build.JobPath) +} diff --git a/backend/plugins/jenkins/e2e/raw_tables/_raw_jenkins_api_builds_unfinished.csv b/backend/plugins/jenkins/e2e/raw_tables/_raw_jenkins_api_builds_unfinished.csv new file mode 100644 index 00000000000..6e6c94ce594 --- /dev/null +++ b/backend/plugins/jenkins/e2e/raw_tables/_raw_jenkins_api_builds_unfinished.csv @@ -0,0 +1,3 @@ +id,params,data,url,input,created_at +1,"{""ConnectionId"":1,""FullName"":""devlake""}","{""_class"":""org.jenkinsci.plugins.workflow.job.WorkflowRun"",""number"":5,""timestamp"":1767052800000,""duration"":0,""estimatedDuration"":600000,""building"":true,""result"":null,""actions"":[]}",https://jenkins/job/devlake/api/json,null,2026-01-01 00:30:00 +2,"{""ConnectionId"":1,""FullName"":""devlake""}","{""_class"":""org.jenkinsci.plugins.workflow.job.WorkflowRun"",""number"":5,""timestamp"":1767052800000,""duration"":90000000,""estimatedDuration"":600000,""building"":false,""result"":""SUCCESS"",""actions"":[]}",https://jenkins/job/devlake/5/api/json,"{""Number"":""5"",""FullName"":""devlake#5""}",2026-01-02 00:30:00 diff --git a/backend/plugins/jenkins/e2e/raw_tables/_tool_jenkins_builds_finished_since.csv b/backend/plugins/jenkins/e2e/raw_tables/_tool_jenkins_builds_finished_since.csv new file mode 100644 index 00000000000..17646e10644 --- /dev/null +++ b/backend/plugins/jenkins/e2e/raw_tables/_tool_jenkins_builds_finished_since.csv @@ -0,0 +1,4 @@ +connection_id,full_name,job_name,job_path,duration,estimated_duration,number,result,timestamp,start_time,class,_raw_data_params,_raw_data_table,_raw_data_id,_raw_data_remark +1,devlake#1,devlake,job/devlake/,3600000,0,1,SUCCESS,1767052800000,2025-12-30T00:00:00.000+00:00,WorkflowRun,,,0, +1,devlake#2,devlake,job/devlake/,345600000,0,2,SUCCESS,1767052800000,2025-12-30T00:00:00.000+00:00,WorkflowRun,,,0, +1,devlake#3,devlake,job/devlake/,600000,0,3,SUCCESS,1767312000000,2026-01-02T00:00:00.000+00:00,WorkflowRun,,,0, diff --git a/backend/plugins/jenkins/e2e/raw_tables/_tool_jenkins_builds_unfinished.csv b/backend/plugins/jenkins/e2e/raw_tables/_tool_jenkins_builds_unfinished.csv new file mode 100644 index 00000000000..0128c7c5ef6 --- /dev/null +++ b/backend/plugins/jenkins/e2e/raw_tables/_tool_jenkins_builds_unfinished.csv @@ -0,0 +1,4 @@ +connection_id,full_name,job_name,job_path,duration,estimated_duration,number,result,timestamp,start_time,class,building,_raw_data_params,_raw_data_table,_raw_data_id,_raw_data_remark +1,devlake#1,devlake,job/,3600000,0,1,SUCCESS,1767052800000,2025-12-30T00:00:00.000+00:00,WorkflowRun,0,,,0, +1,devlake#2,devlake,job/,0,0,2,,1767139200000,2025-12-31T00:00:00.000+00:00,WorkflowRun,1,,,0, +1,other#7,other,job/,0,0,7,,1767139200000,2025-12-31T00:00:00.000+00:00,WorkflowRun,1,,,0, diff --git a/backend/plugins/jenkins/e2e/stages_collector_test.go b/backend/plugins/jenkins/e2e/stages_collector_test.go new file mode 100644 index 00000000000..177c9eda574 --- /dev/null +++ b/backend/plugins/jenkins/e2e/stages_collector_test.go @@ -0,0 +1,54 @@ +/* +Licensed to the Apache Software Foundation (ASF) under one or more +contributor license agreements. See the NOTICE file distributed with +this work for additional information regarding copyright ownership. +The ASF licenses this file to You under the Apache License, Version 2.0 +(the "License"); you may not use this file except in compliance with +the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package e2e + +import ( + "testing" + "time" + + "github.com/apache/devlake/core/dal" + "github.com/apache/devlake/helpers/e2ehelper" + "github.com/apache/devlake/plugins/jenkins/impl" + "github.com/apache/devlake/plugins/jenkins/models" + "github.com/apache/devlake/plugins/jenkins/tasks" + "github.com/stretchr/testify/assert" +) + +// The incremental stage collection must include builds that started before the +// last collection but finished after it (e.g. a build that waited days on an +// input step): they reach _tool_jenkins_builds only once finished. +func TestJenkinsStagesFinishedSince(t *testing.T) { + var jenkins impl.Jenkins + dataflowTester := e2ehelper.NewDataFlowTester(t, "jenkins", jenkins) + + // devlake#1: started 3 days before, ran 1 hour -> collected by an earlier run + // devlake#2: started 3 days before, ran 4 days -> finished after the last run + // devlake#3: started after the last run, ran 10 min + dataflowTester.FlushTabler(&models.JenkinsBuild{}) + dataflowTester.ImportCsvIntoTabler("./raw_tables/_tool_jenkins_builds_finished_since.csv", models.JenkinsBuild{}) + + lastCollection := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + var builds []string + err := dataflowTester.Dal.Pluck("tjb.full_name", &builds, + dal.From("_tool_jenkins_builds as tjb"), + tasks.FinishedSince(lastCollection), + dal.Orderby("tjb.full_name"), + ) + assert.Nil(t, err) + assert.Equal(t, []string{"devlake#2", "devlake#3"}, builds) +} diff --git a/backend/plugins/jenkins/tasks/build_collector.go b/backend/plugins/jenkins/tasks/build_collector.go index b447bed7bcc..f180daae3a6 100644 --- a/backend/plugins/jenkins/tasks/build_collector.go +++ b/backend/plugins/jenkins/tasks/build_collector.go @@ -34,6 +34,20 @@ import ( const RAW_BUILD_TABLE = "jenkins_api_builds" +// buildFields is the tree of fields requested for every build, as a list item or on its own +const buildFields = "timestamp,number,duration,building,estimatedDuration,fullDisplayName,result,actions[lastBuiltRevision[SHA1,branch[name]],remoteUrls,mercurialRevisionNumber,causes[*]],changeSet[kind,revisions[revision]]" + +// UnfinishedBuilds selects the builds of a single job that were still running when last +// collected, so their final state can be collected once they finish. +func UnfinishedBuilds(options *JenkinsOptions) []dal.Clause { + return []dal.Clause{ + dal.Select("tjb.number,tjb.full_name"), + dal.From("_tool_jenkins_builds as tjb"), + dal.Where(`tjb.connection_id = ? and tjb.job_path = ? and tjb.job_name = ? and tjb.building = ?`, + options.ConnectionId, options.JobPath, options.JobName, true), + } +} + var CollectApiBuildsMeta = plugin.SubTaskMeta{ Name: "collectApiBuilds", EntryPoint: CollectApiBuilds, @@ -72,6 +86,7 @@ func CollectApiBuilds(taskCtx plugin.SubTaskContext) errors.Error { func collectSingleJobApiBuilds(taskCtx plugin.SubTaskContext) errors.Error { // The API input is defined in the plugin's task definition, be that the UI or advanced blueprint. data := taskCtx.GetData().(*JenkinsTaskData) + db := taskCtx.GetDal() collector, err := helper.NewStatefulApiCollectorForFinalizableEntity(helper.FinalizableApiCollectorArgs{ RawDataSubTaskArgs: helper.RawDataSubTaskArgs{ Params: JenkinsApiParams{ @@ -89,12 +104,15 @@ func collectSingleJobApiBuilds(taskCtx plugin.SubTaskContext) errors.Error { UrlTemplate: fmt.Sprintf("%sjob/%s/api/json", data.Options.JobPath, data.Options.JobName), Query: func(reqData *helper.RequestData, createdAfter *time.Time) (url.Values, errors.Error) { query := url.Values{} - treeValue := fmt.Sprintf( - "allBuilds[timestamp,number,duration,building,estimatedDuration,fullDisplayName,result,actions[lastBuiltRevision[SHA1,branch[name]],remoteUrls,mercurialRevisionNumber,causes[*]],changeSet[kind,revisions[revision]]]{%d,%d}", - reqData.Pager.Skip, reqData.Pager.Skip+reqData.Pager.Size) + treeValue := fmt.Sprintf("allBuilds[%s]{%d,%d}", + buildFields, reqData.Pager.Skip, reqData.Pager.Skip+reqData.Pager.Size) query.Set("tree", treeValue) return query, nil }, + // Running builds are kept too: the list only returns builds that started since + // the last collection, so a build still running now would never be listed + // again. Stored as building, CollectUnfinishedDetails re-collects it until it + // finishes. ResponseParser: func(res *http.Response) ([]json.RawMessage, errors.Error) { var data struct { Builds []json.RawMessage `json:"allBuilds"` @@ -103,20 +121,7 @@ func collectSingleJobApiBuilds(taskCtx plugin.SubTaskContext) errors.Error { if err != nil { return nil, err } - - builds := make([]json.RawMessage, 0, len(data.Builds)) - for _, build := range data.Builds { - var buildObj map[string]interface{} - err := json.Unmarshal(build, &buildObj) - if err != nil { - return nil, errors.Convert(err) - } - if buildObj["result"] != nil { - builds = append(builds, build) - } - } - - return builds, nil + return data.Builds, nil }, }, GetCreated: func(item json.RawMessage) (time.Time, errors.Error) { @@ -130,6 +135,33 @@ func collectSingleJobApiBuilds(taskCtx plugin.SubTaskContext) errors.Error { return time.Unix(seconds, nanos), nil }, }, + CollectUnfinishedDetails: &helper.FinalizableApiCollectorDetailArgs{ + BuildInputIterator: func() (helper.Iterator, errors.Error) { + cursor, err := db.Cursor(UnfinishedBuilds(data.Options)...) + if err != nil { + return nil, err + } + return helper.NewDalCursorIterator(db, cursor, reflect.TypeOf(SimpleBuild{})) + }, + FinalizableApiCollectorCommonArgs: helper.FinalizableApiCollectorCommonArgs{ + UrlTemplate: fmt.Sprintf("%sjob/%s/{{ .Input.Number }}/api/json", data.Options.JobPath, data.Options.JobName), + Query: func(reqData *helper.RequestData, createdAfter *time.Time) (url.Values, errors.Error) { + return url.Values{"tree": {buildFields}}, nil + }, + // A build deleted while running (e.g. by hand) is skipped instead of failing the task + AfterResponse: func(res *http.Response) errors.Error { + if res.StatusCode == http.StatusNotFound { + return helper.ErrIgnoreAndContinue + } + return nil + }, + ResponseParser: func(res *http.Response) ([]json.RawMessage, errors.Error) { + var build json.RawMessage + err := helper.UnmarshalResponse(res, &build) + return []json.RawMessage{build}, err + }, + }, + }, }) if err != nil { @@ -189,7 +221,7 @@ func collectMultiBranchJobApiBuilds(taskCtx plugin.SubTaskContext) errors.Error UrlTemplate: "{{ .Input.Path }}api/json", Query: func(reqData *helper.RequestData) (url.Values, errors.Error) { query := url.Values{} - treeValue := "allBuilds[timestamp,number,duration,building,estimatedDuration,fullDisplayName,result,actions[lastBuiltRevision[SHA1,branch[name]],remoteUrls,mercurialRevisionNumber,causes[*]],changeSet[kind,revisions[revision]]]" + treeValue := fmt.Sprintf("allBuilds[%s]", buildFields) query.Set("tree", treeValue) logger.Debug("Query: %v", query) diff --git a/backend/plugins/jenkins/tasks/stage_collector.go b/backend/plugins/jenkins/tasks/stage_collector.go index 0398023d7af..8f7c9c04314 100644 --- a/backend/plugins/jenkins/tasks/stage_collector.go +++ b/backend/plugins/jenkins/tasks/stage_collector.go @@ -23,6 +23,7 @@ import ( "net/http" "net/url" "reflect" + "time" "github.com/apache/devlake/core/dal" "github.com/apache/devlake/core/errors" @@ -32,6 +33,16 @@ import ( const RAW_STAGE_TABLE = "jenkins_api_stages" +// FinishedSince selects the builds (aliased tjb) that finished at or after since. +// Builds are stored only once they have a result, so a build that started before +// the last collection and finished after it (e.g. one waiting days on an input step) +// appears in _tool_jenkins_builds only now; filtering on start_time would skip its +// stages forever. timestamp and duration are both in milliseconds, so the sum is +// plain arithmetic on every supported database. +func FinishedSince(since time.Time) dal.Clause { + return dal.Where(`tjb.timestamp + tjb.duration >= ?`, since.UnixMilli()) +} + var CollectApiStagesMeta = plugin.SubTaskMeta{ Name: "collectApiStages", EntryPoint: CollectApiStages, @@ -79,7 +90,7 @@ func collectSingleBuildApiStages(taskCtx plugin.SubTaskContext) errors.Error { data.Options.ConnectionId, data.Options.JobPath, data.Options.JobName, "WorkflowRun"), } if apiCollector.IsIncremental() && apiCollector.GetSince() != nil { - clauses = append(clauses, dal.Where(`tjb.start_time >= ?`, apiCollector.GetSince())) + clauses = append(clauses, FinishedSince(*apiCollector.GetSince())) } cursor, err := db.Cursor(clauses...) if err != nil { @@ -146,7 +157,7 @@ func collectMultiBranchBuildApiStages(taskCtx plugin.SubTaskContext) errors.Erro data.Options.ConnectionId, fmt.Sprintf("%s%%", data.Options.JobFullName), "WorkflowRun"), } if apiCollector.IsIncremental() && apiCollector.GetSince() != nil { - clauses = append(clauses, dal.Where(`tjb.start_time >= ?`, apiCollector.GetSince())) + clauses = append(clauses, FinishedSince(*apiCollector.GetSince())) } cursor, err := db.Cursor(clauses...) if err != nil {