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 {