feat: add filtering by initiator to provisioner job listing in the CLI (#20137)

Relates to https://github.com/coder/internal/issues/934

This PR provides a mechanism to filter provisioner jobs according to who
initiated the job.
This will be used to find pending prebuild jobs when prebuilds have
overwhelmed the provisioner job queue. They can then be canceled.

If prebuilds are overwhelming provisioners, the following steps will be
taken:

```bash
# pause prebuild reconciliation to limit provisioner queue pollution:
coder prebuilds pause 
# cancel pending provisioner jobs to clear the queue
coder provisioner jobs list --initiator="prebuilds" --status="pending" | jq ... | xargs -n1 -I{} coder provisioner jobs cancel {}
# push a fixed template and wait for the import to complete
coder templates push ... # push a fixed template
# resume prebuild reconciliation
coder prebuilds resume
```

This interface differs somewhat from what was specified in the issue,
but still provides a mechanism that addresses the issue. The original
proposal was made by myself and this simpler implementation makes sense.
I might add a `--search` parameter in a follow-up if there is appetite
for it.

Potential follow ups:
* Support for this usage: `coder provisioner jobs list --search
"initiator:prebuilds status:pending"`
* Adding the same parameters to `coder provisioner jobs cancel` as a
convenience feature so that operators don't have to pipe through `jq`
and `xargs`
This commit is contained in:
Sas Swart
2025-10-06 08:56:43 +00:00
committed by GitHub
parent c1357d4e27
commit d17dd5d787
23 changed files with 395 additions and 39 deletions
+11
View File
@@ -3744,6 +3744,13 @@ const docTemplate = `{
"description": "Provisioner tags to filter by (JSON of the form {'tag1':'value1','tag2':'value2'})",
"name": "tags",
"in": "query"
},
{
"type": "string",
"format": "uuid",
"description": "Filter results by initiator",
"name": "initiator",
"in": "query"
}
],
"responses": {
@@ -15974,6 +15981,10 @@ const docTemplate = `{
"type": "string",
"format": "uuid"
},
"initiator_id": {
"type": "string",
"format": "uuid"
},
"input": {
"$ref": "#/definitions/codersdk.ProvisionerJobInput"
},
+11
View File
@@ -3299,6 +3299,13 @@
"description": "Provisioner tags to filter by (JSON of the form {'tag1':'value1','tag2':'value2'})",
"name": "tags",
"in": "query"
},
{
"type": "string",
"format": "uuid",
"description": "Filter results by initiator",
"name": "initiator",
"in": "query"
}
],
"responses": {
@@ -14532,6 +14539,10 @@
"type": "string",
"format": "uuid"
},
"initiator_id": {
"type": "string",
"format": "uuid"
},
"input": {
"$ref": "#/definitions/codersdk.ProvisionerJobInput"
},
+2
View File
@@ -2484,10 +2484,12 @@ func (s *MethodTestSuite) TestExtraMethods() {
ds, err := db.GetProvisionerJobsByOrganizationAndStatusWithQueuePositionAndProvisioner(context.Background(), database.GetProvisionerJobsByOrganizationAndStatusWithQueuePositionAndProvisionerParams{
OrganizationID: org.ID,
InitiatorID: uuid.Nil,
})
s.NoError(err, "get provisioner jobs by org")
check.Args(database.GetProvisionerJobsByOrganizationAndStatusWithQueuePositionAndProvisionerParams{
OrganizationID: org.ID,
InitiatorID: uuid.Nil,
}).Asserts(j1, policy.ActionRead, j2, policy.ActionRead).Returns(ds)
}))
}
+4 -1
View File
@@ -9834,6 +9834,7 @@ WHERE
AND (COALESCE(array_length($2::uuid[], 1), 0) = 0 OR pj.id = ANY($2::uuid[]))
AND (COALESCE(array_length($3::provisioner_job_status[], 1), 0) = 0 OR pj.job_status = ANY($3::provisioner_job_status[]))
AND ($4::tagset = 'null'::tagset OR provisioner_tagset_contains(pj.tags::tagset, $4::tagset))
AND ($5::uuid = '00000000-0000-0000-0000-000000000000'::uuid OR pj.initiator_id = $5::uuid)
GROUP BY
pj.id,
qp.queue_position,
@@ -9849,7 +9850,7 @@ GROUP BY
ORDER BY
pj.created_at DESC
LIMIT
$5::int
$6::int
`
type GetProvisionerJobsByOrganizationAndStatusWithQueuePositionAndProvisionerParams struct {
@@ -9857,6 +9858,7 @@ type GetProvisionerJobsByOrganizationAndStatusWithQueuePositionAndProvisionerPar
IDs []uuid.UUID `db:"ids" json:"ids"`
Status []ProvisionerJobStatus `db:"status" json:"status"`
Tags StringMap `db:"tags" json:"tags"`
InitiatorID uuid.UUID `db:"initiator_id" json:"initiator_id"`
Limit sql.NullInt32 `db:"limit" json:"limit"`
}
@@ -9881,6 +9883,7 @@ func (q *sqlQuerier) GetProvisionerJobsByOrganizationAndStatusWithQueuePositionA
pq.Array(arg.IDs),
pq.Array(arg.Status),
arg.Tags,
arg.InitiatorID,
arg.Limit,
)
if err != nil {
@@ -224,6 +224,7 @@ WHERE
AND (COALESCE(array_length(@ids::uuid[], 1), 0) = 0 OR pj.id = ANY(@ids::uuid[]))
AND (COALESCE(array_length(@status::provisioner_job_status[], 1), 0) = 0 OR pj.job_status = ANY(@status::provisioner_job_status[]))
AND (@tags::tagset = 'null'::tagset OR provisioner_tagset_contains(pj.tags::tagset, @tags::tagset))
AND (@initiator_id::uuid = '00000000-0000-0000-0000-000000000000'::uuid OR pj.initiator_id = @initiator_id::uuid)
GROUP BY
pj.id,
qp.queue_position,
+4
View File
@@ -76,6 +76,7 @@ func (api *API) provisionerJob(rw http.ResponseWriter, r *http.Request) {
// @Param ids query []string false "Filter results by job IDs" format(uuid)
// @Param status query codersdk.ProvisionerJobStatus false "Filter results by status" enums(pending,running,succeeded,canceling,canceled,failed)
// @Param tags query object false "Provisioner tags to filter by (JSON of the form {'tag1':'value1','tag2':'value2'})"
// @Param initiator query string false "Filter results by initiator" format(uuid)
// @Success 200 {array} codersdk.ProvisionerJob
// @Router /organizations/{organization}/provisionerjobs [get]
func (api *API) provisionerJobs(rw http.ResponseWriter, r *http.Request) {
@@ -110,6 +111,7 @@ func (api *API) handleAuthAndFetchProvisionerJobs(rw http.ResponseWriter, r *htt
ids = p.UUIDs(qp, nil, "ids")
}
tags := p.JSONStringMap(qp, database.StringMap{}, "tags")
initiatorID := p.UUID(qp, uuid.Nil, "initiator_id")
p.ErrorExcessParams(qp)
if len(p.Errors) > 0 {
httpapi.Write(ctx, rw, http.StatusBadRequest, codersdk.Response{
@@ -125,6 +127,7 @@ func (api *API) handleAuthAndFetchProvisionerJobs(rw http.ResponseWriter, r *htt
Limit: sql.NullInt32{Int32: limit, Valid: limit > 0},
IDs: ids,
Tags: tags,
InitiatorID: initiatorID,
})
if err != nil {
if httpapi.Is404Error(err) {
@@ -355,6 +358,7 @@ func convertProvisionerJob(pj database.GetProvisionerJobsByIDsWithQueuePositionR
job := codersdk.ProvisionerJob{
ID: provisionerJob.ID,
OrganizationID: provisionerJob.OrganizationID,
InitiatorID: provisionerJob.InitiatorID,
CreatedAt: provisionerJob.CreatedAt,
Type: codersdk.ProvisionerJobType(provisionerJob.Type),
Error: provisionerJob.Error.String,
+102
View File
@@ -58,6 +58,8 @@ func TestProvisionerJobs(t *testing.T) {
StartedAt: sql.NullTime{Time: dbtime.Now(), Valid: true},
Type: database.ProvisionerJobTypeWorkspaceBuild,
Input: json.RawMessage(`{"workspace_build_id":"` + wbID.String() + `"}`),
InitiatorID: member.ID,
Tags: database.StringMap{"initiatorTest": "true"},
})
dbgen.WorkspaceBuild(t, db, database.WorkspaceBuild{
ID: wbID,
@@ -71,6 +73,7 @@ func TestProvisionerJobs(t *testing.T) {
dbgen.ProvisionerJob(t, db, nil, database.ProvisionerJob{
OrganizationID: owner.OrganizationID,
Tags: database.StringMap{"count": strconv.Itoa(i)},
InitiatorID: owner.UserID,
})
}
@@ -165,6 +168,94 @@ func TestProvisionerJobs(t *testing.T) {
require.Len(t, jobs, 1)
})
t.Run("Initiator", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitMedium)
jobs, err := templateAdminClient.OrganizationProvisionerJobs(ctx, owner.OrganizationID, &codersdk.OrganizationProvisionerJobsOptions{
InitiatorID: &member.ID,
})
require.NoError(t, err)
require.GreaterOrEqual(t, len(jobs), 1)
require.Equal(t, member.ID, jobs[0].InitiatorID)
})
t.Run("InitiatorWithOtherFilters", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitMedium)
// Test filtering by initiator ID combined with status filter
jobs, err := templateAdminClient.OrganizationProvisionerJobs(ctx, owner.OrganizationID, &codersdk.OrganizationProvisionerJobsOptions{
InitiatorID: &owner.UserID,
Status: []codersdk.ProvisionerJobStatus{codersdk.ProvisionerJobSucceeded},
})
require.NoError(t, err)
// Verify all returned jobs have the correct initiator and status
for _, job := range jobs {
require.Equal(t, owner.UserID, job.InitiatorID)
require.Equal(t, codersdk.ProvisionerJobSucceeded, job.Status)
}
})
t.Run("InitiatorWithLimit", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitMedium)
// Test filtering by initiator ID with limit
jobs, err := templateAdminClient.OrganizationProvisionerJobs(ctx, owner.OrganizationID, &codersdk.OrganizationProvisionerJobsOptions{
InitiatorID: &owner.UserID,
Limit: 1,
})
require.NoError(t, err)
require.Len(t, jobs, 1)
// Verify the returned job has the correct initiator
require.Equal(t, owner.UserID, jobs[0].InitiatorID)
})
t.Run("InitiatorWithTags", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitMedium)
// Test filtering by initiator ID combined with tags
jobs, err := templateAdminClient.OrganizationProvisionerJobs(ctx, owner.OrganizationID, &codersdk.OrganizationProvisionerJobsOptions{
InitiatorID: &member.ID,
Tags: map[string]string{"initiatorTest": "true"},
})
require.NoError(t, err)
require.Len(t, jobs, 1)
// Verify the returned job has the correct initiator and tags
require.Equal(t, member.ID, jobs[0].InitiatorID)
require.Equal(t, "true", jobs[0].Tags["initiatorTest"])
})
t.Run("InitiatorNotFound", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitMedium)
// Test with non-existent initiator ID
nonExistentID := uuid.New()
jobs, err := templateAdminClient.OrganizationProvisionerJobs(ctx, owner.OrganizationID, &codersdk.OrganizationProvisionerJobsOptions{
InitiatorID: &nonExistentID,
})
require.NoError(t, err)
require.Len(t, jobs, 0)
})
t.Run("InitiatorNil", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitMedium)
// Test with nil initiator ID (should return all jobs)
jobs, err := templateAdminClient.OrganizationProvisionerJobs(ctx, owner.OrganizationID, &codersdk.OrganizationProvisionerJobsOptions{
InitiatorID: nil,
})
require.NoError(t, err)
require.GreaterOrEqual(t, len(jobs), 50) // Should return all jobs (up to default limit)
})
t.Run("Limit", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitMedium)
@@ -185,6 +276,17 @@ func TestProvisionerJobs(t *testing.T) {
require.Error(t, err)
require.Len(t, jobs, 0)
})
t.Run("MemberDeniedWithInitiator", func(t *testing.T) {
t.Parallel()
ctx := testutil.Context(t, testutil.WaitMedium)
// Member should not be able to access jobs even with initiator filter
jobs, err := memberClient.OrganizationProvisionerJobs(ctx, owner.OrganizationID, &codersdk.OrganizationProvisionerJobsOptions{
InitiatorID: &member.ID,
})
require.Error(t, err)
require.Len(t, jobs, 0)
})
})
// Ensures that when a provisioner job is in the succeeded state,