Skip to content
Open
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
73 changes: 73 additions & 0 deletions backend/plugins/jenkins/e2e/builds_unfinished_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
@@ -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,
Original file line number Diff line number Diff line change
@@ -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,
54 changes: 54 additions & 0 deletions backend/plugins/jenkins/e2e/stages_collector_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
68 changes: 50 additions & 18 deletions backend/plugins/jenkins/tasks/build_collector.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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{
Expand All @@ -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"`
Expand All @@ -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) {
Expand All @@ -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 {
Expand Down Expand Up @@ -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)
Expand Down
15 changes: 13 additions & 2 deletions backend/plugins/jenkins/tasks/stage_collector.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import (
"net/http"
"net/url"
"reflect"
"time"

"github.com/apache/devlake/core/dal"
"github.com/apache/devlake/core/errors"
Expand All @@ -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,
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down