From e7ea649dc222836f40ce6c5642cf01e70081da8b Mon Sep 17 00:00:00 2001 From: Jon Ayers Date: Mon, 9 Mar 2026 16:47:02 -0500 Subject: [PATCH] fix: optimize GetProvisionerJobsByIDsWithQueuePosition query (#22724) --- coderd/database/querier_test.go | 76 +++++++++++++++++++++ coderd/database/queries.sql.go | 27 +++++--- coderd/database/queries/provisionerjobs.sql | 27 +++++--- 3 files changed, 114 insertions(+), 16 deletions(-) diff --git a/coderd/database/querier_test.go b/coderd/database/querier_test.go index 428273bd1d..7b2420302c 100644 --- a/coderd/database/querier_test.go +++ b/coderd/database/querier_test.go @@ -3867,6 +3867,37 @@ func TestGetProvisionerJobsByIDsWithQueuePosition(t *testing.T) { queueSizes: nil, // TODO(yevhenii): should it be empty array instead? queuePositions: nil, }, + // Many daemons with identical tags should produce same results as one. + { + name: "duplicate-daemons-same-tags", + jobTags: []database.StringMap{ + {"a": "1"}, + {"a": "1", "b": "2"}, + }, + daemonTags: []database.StringMap{ + {"a": "1", "b": "2"}, + {"a": "1", "b": "2"}, + {"a": "1", "b": "2"}, + }, + queueSizes: []int64{2, 2}, + queuePositions: []int64{1, 2}, + }, + // Jobs that don't match any queried job's daemon should still + // have correct queue positions. + { + name: "irrelevant-daemons-filtered", + jobTags: []database.StringMap{ + {"a": "1"}, + {"x": "9"}, + }, + daemonTags: []database.StringMap{ + {"a": "1"}, + {"x": "9"}, + }, + queueSizes: []int64{1}, + queuePositions: []int64{1}, + skipJobIDs: map[int]struct{}{1: {}}, + }, } for _, tc := range testCases { @@ -4192,6 +4223,51 @@ func TestGetProvisionerJobsByIDsWithQueuePosition_OrderValidation(t *testing.T) assert.EqualValues(t, []int64{1, 2, 3, 4, 5, 6}, queuePositions, "expected queue positions to be set correctly") } +func TestGetProvisionerJobsByIDsWithQueuePosition_DuplicateDaemons(t *testing.T) { + t.Parallel() + db, _ := dbtestutil.NewDB(t) + now := dbtime.Now() + ctx := testutil.Context(t, testutil.WaitShort) + + // Create 3 pending jobs with the same tags. + jobs := make([]database.ProvisionerJob, 3) + for i := range jobs { + jobs[i] = dbgen.ProvisionerJob(t, db, nil, database.ProvisionerJob{ + CreatedAt: now.Add(-time.Duration(3-i) * time.Minute), + Tags: database.StringMap{"scope": "organization", "owner": ""}, + }) + } + + // Create 50 daemons with identical tags (simulates scale). + for i := range 50 { + dbgen.ProvisionerDaemon(t, db, database.ProvisionerDaemon{ + Name: fmt.Sprintf("daemon_%d", i), + Provisioners: []database.ProvisionerType{database.ProvisionerTypeEcho}, + Tags: database.StringMap{"scope": "organization", "owner": ""}, + }) + } + + jobIDs := make([]uuid.UUID, len(jobs)) + for i, j := range jobs { + jobIDs[i] = j.ID + } + + results, err := db.GetProvisionerJobsByIDsWithQueuePosition(ctx, + database.GetProvisionerJobsByIDsWithQueuePositionParams{ + IDs: jobIDs, + StaleIntervalMS: provisionerdserver.StaleInterval.Milliseconds(), + }) + require.NoError(t, err) + require.Len(t, results, 3) + + // All daemons have identical tags, so queue should be same as + // if there were just one daemon. + for i, r := range results { + assert.Equal(t, int64(3), r.QueueSize, "job %d queue size", i) + assert.Equal(t, int64(i+1), r.QueuePosition, "job %d queue position", i) + } +} + func TestGroupRemovalTrigger(t *testing.T) { t.Parallel() diff --git a/coderd/database/queries.sql.go b/coderd/database/queries.sql.go index 5c051af0c7..fe17470c9f 100644 --- a/coderd/database/queries.sql.go +++ b/coderd/database/queries.sql.go @@ -12654,7 +12654,7 @@ const getProvisionerJobsByIDsWithQueuePosition = `-- name: GetProvisionerJobsByI WITH filtered_provisioner_jobs AS ( -- Step 1: Filter provisioner_jobs SELECT - id, created_at + id, created_at, tags FROM provisioner_jobs WHERE @@ -12669,21 +12669,32 @@ pending_jobs AS ( WHERE job_status = 'pending' ), -online_provisioner_daemons AS ( - SELECT id, tags FROM provisioner_daemons pd - WHERE pd.last_seen_at IS NOT NULL AND pd.last_seen_at >= (NOW() - ($2::bigint || ' ms')::interval) +unique_daemon_tags AS ( + SELECT DISTINCT tags FROM provisioner_daemons pd + WHERE pd.last_seen_at IS NOT NULL + AND pd.last_seen_at >= (NOW() - ($2::bigint || ' ms')::interval) +), +relevant_daemon_tags AS ( + SELECT udt.tags + FROM unique_daemon_tags udt + WHERE EXISTS ( + SELECT 1 FROM filtered_provisioner_jobs fpj + WHERE provisioner_tagset_contains(udt.tags, fpj.tags) + ) ), ranked_jobs AS ( -- Step 3: Rank only pending jobs based on provisioner availability SELECT pj.id, pj.created_at, - ROW_NUMBER() OVER (PARTITION BY opd.id ORDER BY pj.initiator_id = 'c42fdf75-3097-471c-8c33-fb52454d81c0'::uuid ASC, pj.created_at ASC) AS queue_position, - COUNT(*) OVER (PARTITION BY opd.id) AS queue_size + ROW_NUMBER() OVER (PARTITION BY rdt.tags ORDER BY pj.initiator_id = 'c42fdf75-3097-471c-8c33-fb52454d81c0'::uuid ASC, pj.created_at ASC) AS queue_position, + COUNT(*) OVER (PARTITION BY rdt.tags) AS queue_size FROM pending_jobs pj - INNER JOIN online_provisioner_daemons opd - ON provisioner_tagset_contains(opd.tags, pj.tags) -- Join only on the small pending set + INNER JOIN + relevant_daemon_tags rdt + ON + provisioner_tagset_contains(rdt.tags, pj.tags) ), final_jobs AS ( -- Step 4: Compute best queue position and max queue size per job diff --git a/coderd/database/queries/provisionerjobs.sql b/coderd/database/queries/provisionerjobs.sql index 0f1b1db94d..a11c18e9ee 100644 --- a/coderd/database/queries/provisionerjobs.sql +++ b/coderd/database/queries/provisionerjobs.sql @@ -79,7 +79,7 @@ WHERE WITH filtered_provisioner_jobs AS ( -- Step 1: Filter provisioner_jobs SELECT - id, created_at + id, created_at, tags FROM provisioner_jobs WHERE @@ -94,21 +94,32 @@ pending_jobs AS ( WHERE job_status = 'pending' ), -online_provisioner_daemons AS ( - SELECT id, tags FROM provisioner_daemons pd - WHERE pd.last_seen_at IS NOT NULL AND pd.last_seen_at >= (NOW() - (@stale_interval_ms::bigint || ' ms')::interval) +unique_daemon_tags AS ( + SELECT DISTINCT tags FROM provisioner_daemons pd + WHERE pd.last_seen_at IS NOT NULL + AND pd.last_seen_at >= (NOW() - (@stale_interval_ms::bigint || ' ms')::interval) +), +relevant_daemon_tags AS ( + SELECT udt.tags + FROM unique_daemon_tags udt + WHERE EXISTS ( + SELECT 1 FROM filtered_provisioner_jobs fpj + WHERE provisioner_tagset_contains(udt.tags, fpj.tags) + ) ), ranked_jobs AS ( -- Step 3: Rank only pending jobs based on provisioner availability SELECT pj.id, pj.created_at, - ROW_NUMBER() OVER (PARTITION BY opd.id ORDER BY pj.initiator_id = 'c42fdf75-3097-471c-8c33-fb52454d81c0'::uuid ASC, pj.created_at ASC) AS queue_position, - COUNT(*) OVER (PARTITION BY opd.id) AS queue_size + ROW_NUMBER() OVER (PARTITION BY rdt.tags ORDER BY pj.initiator_id = 'c42fdf75-3097-471c-8c33-fb52454d81c0'::uuid ASC, pj.created_at ASC) AS queue_position, + COUNT(*) OVER (PARTITION BY rdt.tags) AS queue_size FROM pending_jobs pj - INNER JOIN online_provisioner_daemons opd - ON provisioner_tagset_contains(opd.tags, pj.tags) -- Join only on the small pending set + INNER JOIN + relevant_daemon_tags rdt + ON + provisioner_tagset_contains(rdt.tags, pj.tags) ), final_jobs AS ( -- Step 4: Compute best queue position and max queue size per job