Skip to content
Merged
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
18 changes: 18 additions & 0 deletions internal/jobs/queue_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
package jobs

import "testing"

// TestQueueReconcileConst guards the queue name so a typo in the
// periodic-job closures doesn't silently route reconcilers back to the
// default queue, reintroducing the starvation bug.
//
// Background: prior to fix/reconcile-queue, all periodic jobs landed on
// river.QueueDefault. A weekly_digest fan-out (1 row per team) accumulated
// 232K available jobs and pinned all 5 worker slots, so the
// deploy_status_reconcile job (queued every 30s) never ran. Customers saw
// status="building" indefinitely while their pods were already Ready.
func TestQueueReconcileConst(t *testing.T) {
if queueReconcile != "reconcile" {
t.Errorf("queueReconcile = %q; want %q. River queue names are referenced as plain strings from periodic-job InsertOpts; renaming this without updating the matching QueueConfig entry in workers.go reintroduces the starvation bug.", queueReconcile, "reconcile")
}
}
25 changes: 22 additions & 3 deletions internal/jobs/workers.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,14 @@ import (
"instant.dev/worker/internal/provisioner"
)

// queueReconcile is the dedicated queue for fast periodic reconcilers
// (deploy-status every 30s, custom-domain every 5min). Isolated from the
// default queue so a fan-out backlog on bulk jobs (weekly_digest, etc.)
// cannot starve them — the previous symptom was deploy status staying in
// "building" indefinitely while 200K weekly_digest rows occupied every
// worker slot.
const queueReconcile = "reconcile"

// Workers wraps a running River client.
type Workers struct {
client *river.Client[pgx.Tx]
Expand Down Expand Up @@ -174,29 +182,40 @@ func StartWorkers(ctx context.Context, db *sql.DB, rdb *redis.Client, cfg *confi
// customDomainReconcileInterval in custom_domain_reconcile.go.
// RunOnStart=true so a worker restart immediately picks up domains
// that became verifiable while we were down.
// Routed to the "reconcile" queue so a backlog on the default queue
// (e.g. a weekly_digest fan-out) cannot starve it.
river.NewPeriodicJob(
river.PeriodicInterval(customDomainReconcileInterval),
func() (river.JobArgs, *river.InsertOpts) {
return CustomDomainReconcileArgs{}, nil
return CustomDomainReconcileArgs{}, &river.InsertOpts{Queue: queueReconcile}
},
&river.PeriodicJobOpts{RunOnStart: true},
),
// Deploy-status reconciler runs every 30s — see
// deployStatusReconcileInterval in deploy_status_reconcile.go.
// RunOnStart=true so a worker restart immediately reconciles any
// deployments stuck in "building" or "deploying" from the last cycle.
// On the "reconcile" queue for the same starvation-protection reason.
river.NewPeriodicJob(
river.PeriodicInterval(deployStatusReconcileInterval),
func() (river.JobArgs, *river.InsertOpts) {
return DeployStatusReconcileArgs{}, nil
return DeployStatusReconcileArgs{}, &river.InsertOpts{Queue: queueReconcile}
},
&river.PeriodicJobOpts{RunOnStart: true},
),
}

riverClient, err := river.NewClient(riverpgxv5.New(pool), &river.Config{
Queues: map[string]river.QueueConfig{
// Bulk email + heavyweight periodics live on the default queue.
// A fan-out (one row per team) can blow this to 100K+ rows; the
// reconcile queue below guarantees small-but-critical periodic
// jobs always have worker capacity.
river.QueueDefault: {MaxWorkers: 5},
// Reserved for fast, frequent reconcilers (deploy-status every 30s,
// custom-domain every 5min). 2 workers is enough because each
// invocation does one batched DB query + per-row k8s GETs.
queueReconcile: {MaxWorkers: 2},
},
Workers: workers,
PeriodicJobs: periodicJobs,
Expand All @@ -217,7 +236,7 @@ func StartWorkers(ctx context.Context, db *sql.DB, rdb *redis.Client, cfg *confi
}

slog.Info("jobs.workers.started",
"queues", fmt.Sprintf("%v", []string{river.QueueDefault}),
"queues", fmt.Sprintf("%v", []string{river.QueueDefault, queueReconcile}),
"max_workers", 5,
)

Expand Down