internal/queue: shard task workload across 4 queues

- Modify queue logic to shard task workload across 4 queues
- Update Terraform to create new queue resources
- Queue count determined via Terraform env variable

Run `terraform apply` to push resource changes.

Change-Id: I839ebd7cf3f915960a2b1fd1c7583438a001643b
Reviewed-on: https://go-review.googlesource.com/c/pkgsite-metrics/+/698995
Auto-Submit: Ethan Lee <ethanalee@google.com>
LUCI-TryBot-Result: Go LUCI <golang-scoped@luci-project-accounts.iam.gserviceaccount.com>
Reviewed-by: Markus Kusano <kusano@google.com>
diff --git a/internal/config/config.go b/internal/config/config.go
index 674f5f9..6e09702 100644
--- a/internal/config/config.go
+++ b/internal/config/config.go
@@ -55,10 +55,7 @@
 	// BigQueryDataset is the BigQuery dataset to write results to.
 	BigQueryDataset string
 
-	// QueueName is the name of the Cloud Tasks queue.
-	QueueName string
-
-	// QueueURL is the URL that the Cloud Tasks queue should send requests to.
+	// QueueURL is the URL that the Cloud Tasks queues should send requests to.
 	// It should be used when the worker is not on AppEngine.
 	QueueURL string
 
@@ -106,6 +103,9 @@
 
 	// ProxyURL is the url for the Go module proxy.
 	ProxyURL string
+
+	// Cloud Task Queues
+	QueueNames []string
 }
 
 // Init resolves all configuration values provided by the config package. It
@@ -126,7 +126,7 @@
 		LocationID:            "us-central1",
 		StaticPath:            ts,
 		BigQueryDataset:       GetEnv("GO_ECOSYSTEM_BIGQUERY_DATASET", "disable"),
-		QueueName:             os.Getenv("GO_ECOSYSTEM_QUEUE_NAME"),
+		QueueNames:            strings.Split(os.Getenv("GO_ECOSYSTEM_QUEUE_NAMES"), ","),
 		QueueURL:              os.Getenv("GO_ECOSYSTEM_QUEUE_URL"),
 		VulnDBBucketProjectID: os.Getenv("GO_ECOSYSTEM_VULNDB_BUCKET_PROJECT"),
 		BinaryBucket:          os.Getenv("GO_ECOSYSTEM_BINARY_BUCKET"),
@@ -139,6 +139,7 @@
 		PkgsiteDBSecret:       os.Getenv("GO_ECOSYSTEM_PKGSITE_DB_SECRET"),
 		ProxyURL:              GetEnv("GO_MODULE_PROXY_URL", "https://proxy.golang.org"),
 	}
+
 	if OnCloudRun() {
 		sa, err := gceMetadata(ctx, "instance/service-accounts/default/email")
 		if err != nil {
diff --git a/internal/queue/queue.go b/internal/queue/queue.go
index c41f0b5..940d092 100644
--- a/internal/queue/queue.go
+++ b/internal/queue/queue.go
@@ -13,6 +13,7 @@
 	"errors"
 	"fmt"
 	"io"
+	"math/rand/v2"
 	"strings"
 	"time"
 
@@ -50,19 +51,19 @@
 	if err != nil {
 		return nil, err
 	}
-	g, err := newGCP(cfg, client, cfg.QueueName)
+	g, err := newGCP(cfg, client, cfg.QueueNames)
 	if err != nil {
 		return nil, err
 	}
-	log.Infof(ctx, "enqueuing at %s with queueURL=%q", g.queueName, g.queueURL)
+	log.Infof(ctx, "enqueuing to %d queues with URL=%q", len(g.queueNames), g.queueURL)
 	return g, nil
 }
 
 // GCP provides a Queue implementation backed by the Google Cloud Tasks API.
 type GCP struct {
-	client    *cloudtasks.Client
-	queueName string // full GCP name of the queue
-	queueURL  string // non-AppEngine URL to post tasks to
+	client     *cloudtasks.Client
+	queueNames []string // full GCP name of the queue
+	queueURL   string   // non-AppEngine URL to post tasks to
 	// token holds information that lets the task queue construct an authorized request to the worker.
 	// Since the worker sits behind the IAP, the queue needs an identity token that includes the
 	// identity of a service account that has access, and the client ID for the IAP.
@@ -73,10 +74,10 @@
 // newGCP returns a new Queue that can be used to enqueue tasks using the
 // cloud tasks API.  The given queueID should be the name of the queue in the
 // cloud tasks console.
-func newGCP(cfg *config.Config, client *cloudtasks.Client, queueID string) (_ *GCP, err error) {
-	defer derrors.Wrap(&err, "newGCP(cfg, client, %q)", queueID)
-	if queueID == "" {
-		return nil, errors.New("empty queueID")
+func newGCP(cfg *config.Config, client *cloudtasks.Client, queueIDs []string) (_ *GCP, err error) {
+	defer derrors.Wrap(&err, "newGCP(cfg, client, %q)", queueIDs)
+	if len(queueIDs) == 0 {
+		return nil, errors.New("must provide at least one queueID")
 	}
 	if cfg.ProjectID == "" {
 		return nil, errors.New("empty ProjectID")
@@ -90,10 +91,14 @@
 	if cfg.ServiceAccount == "" {
 		return nil, errors.New("empty ServiceAccount")
 	}
+	var queueNames []string
+	for _, id := range queueIDs {
+		queueNames = append(queueNames, fmt.Sprintf("projects/%s/locations/%s/queues/%s", cfg.ProjectID, cfg.LocationID, id))
+	}
 	return &GCP{
-		client:    client,
-		queueName: fmt.Sprintf("projects/%s/locations/%s/queues/%s", cfg.ProjectID, cfg.LocationID, queueID),
-		queueURL:  cfg.QueueURL,
+		client:     client,
+		queueNames: queueNames,
+		queueURL:   cfg.QueueURL,
 		token: &taskspb.HttpRequest_OidcToken{
 			OidcToken: &taskspb.OidcToken{
 				ServiceAccountEmail: cfg.ServiceAccount,
@@ -171,9 +176,12 @@
 		relativeURI += "?" + params
 	}
 
+	queueIndex := rand.IntN(len(q.queueNames))
+	queueName := q.queueNames[queueIndex]
+
 	taskID := newTaskID(opts.Namespace, task)
 	taskpb := &taskspb.Task{
-		Name:             fmt.Sprintf("%s/tasks/%s", q.queueName, taskID),
+		Name:             fmt.Sprintf("%s/tasks/%s", queueName, taskID),
 		DispatchDeadline: durationpb.New(maxCloudTasksTimeout),
 		MessageType: &taskspb.Task_HttpRequest{
 			HttpRequest: &taskspb.HttpRequest{
@@ -184,7 +192,7 @@
 		},
 	}
 	req := &taskspb.CreateTaskRequest{
-		Parent: q.queueName,
+		Parent: queueName,
 		Task:   taskpb,
 	}
 	// If suffix is non-empty, append it to the task name.
diff --git a/internal/queue/queue_test.go b/internal/queue/queue_test.go
index 911f043..bb5bbd0 100644
--- a/internal/queue/queue_test.go
+++ b/internal/queue/queue_test.go
@@ -53,8 +53,45 @@
 		QueueURL:       "http://1.2.3.4:8000",
 		ServiceAccount: "sa",
 	}
+	queueIDs := []string{"queueID-0", "queueID-1", "queueID-2", "queueID-3"}
+	var possibleQueueNames []string
+	for _, qID := range queueIDs {
+		possibleQueueNames = append(possibleQueueNames, "projects/Project/locations/us-central1/queues/"+qID)
+	}
+
+	gcp, err := newGCP(&cfg, nil, queueIDs)
+	if err != nil {
+		t.Fatal(err)
+	}
+
+	opts := &Options{
+		Namespace:      "test",
+		TaskNameSuffix: "suf",
+	}
+	sreq := &testTask{
+		name:   "name",
+		path:   "mod@v1.2.3",
+		params: "importedby=0&mode=test&insecure=true",
+	}
+
+	got, err := gcp.newTaskRequest(sreq, opts)
+	if err != nil {
+		t.Fatal(err)
+	}
+
+	validParent := false
+	for _, pqn := range possibleQueueNames {
+		if got.Parent == pqn {
+			validParent = true
+			break
+		}
+	}
+	if !validParent {
+		t.Errorf("got.Parent = %q, want one of %v", got.Parent, possibleQueueNames)
+	}
+
 	want := &taskspb.CreateTaskRequest{
-		Parent: "projects/Project/locations/us-central1/queues/queueID",
+		Parent: got.Parent,
 		Task: &taskspb.Task{
 			DispatchDeadline: durationpb.New(maxCloudTasksTimeout),
 			MessageType: &taskspb.Task_HttpRequest{
@@ -70,35 +107,33 @@
 			},
 		},
 	}
-	gcp, err := newGCP(&cfg, nil, "queueID")
-	if err != nil {
-		t.Fatal(err)
-	}
-	opts := &Options{
-		Namespace:      "test",
-		TaskNameSuffix: "suf",
-	}
-	sreq := &testTask{
-		name:   "name",
-		path:   "mod@v1.2.3",
-		params: "importedby=0&mode=test&insecure=true",
-	}
-	got, err := gcp.newTaskRequest(sreq, opts)
-	if err != nil {
-		t.Fatal(err)
-	}
 	want.Task.Name = got.Task.Name
+
 	if diff := cmp.Diff(want, got, protocmp.Transform()); diff != "" {
 		t.Errorf("mismatch (-want, +got):\n%s", diff)
 	}
 
 	opts.DisableProxyFetch = true
 	want.Task.MessageType.(*taskspb.Task_HttpRequest).HttpRequest.Url += "&proxyfetch=off"
+
 	got, err = gcp.newTaskRequest(sreq, opts)
 	if err != nil {
 		t.Fatal(err)
 	}
 	want.Task.Name = got.Task.Name
+
+	validParent = false
+	for _, pqn := range possibleQueueNames {
+		if got.Parent == pqn {
+			validParent = true
+			break
+		}
+	}
+	if !validParent {
+		t.Errorf("got.Parent = %q, want one of %v", got.Parent, possibleQueueNames)
+	}
+	want.Parent = got.Parent
+
 	if diff := cmp.Diff(want, got, protocmp.Transform()); diff != "" {
 		t.Errorf("mismatch (-want, +got):\n%s", diff)
 	}
diff --git a/terraform/environment/worker.tf b/terraform/environment/worker.tf
index 4a774d0..a534ddd 100644
--- a/terraform/environment/worker.tf
+++ b/terraform/environment/worker.tf
@@ -102,6 +102,10 @@
           name  = "GO_ECOSYSTEM_WORKER_USE_PROFILER"
           value = var.use_profiler
         }
+        env {
+          name = "GO_ECOSYSTEM_QUEUE_NAMES"
+          value = join(",", [for queue in google_cloud_tasks_queue.worker_queues : queue.name])
+        }
         resources {
           limits = {
             "cpu"    = "8000m"
@@ -113,10 +117,6 @@
           value = local.worker_url
         }
         env {
-          name  = "GO_ECOSYSTEM_QUEUE_NAME"
-          value = "${var.env}-worker-tasks"
-        }
-        env {
           name = "GITHUB_ACCESS_TOKEN"
           value_from {
             secret_key_ref {
@@ -207,14 +207,15 @@
 ################################################################
 # Other components.
 
-resource "google_cloud_tasks_queue" "worker_tasks" {
-  name     = "${var.env}-worker-tasks"
-  location = var.region
+resource "google_cloud_tasks_queue" "worker_queues" {
+  count = var.env == "prod" ? 1 : 4
   project  = var.project
+  name     = "${var.env}-worker-queue-${count.index}"
+  location = var.region
 
   rate_limits {
-    max_concurrent_dispatches = 200
     max_dispatches_per_second = 500
+    max_concurrent_dispatches = 5000
   }
 
   retry_config {