diff --git a/internal/orchestrator/orchestrator.go b/internal/orchestrator/orchestrator.go index 1cf2e79..d201935 100644 --- a/internal/orchestrator/orchestrator.go +++ b/internal/orchestrator/orchestrator.go @@ -9,6 +9,7 @@ import ( "fmt" "io" "log/slog" + "sync/atomic" "time" "golang.org/x/sync/errgroup" @@ -368,33 +369,42 @@ func (ro *RunOrchestrator) createSecrets( ) (int, error) { cfg := ro.config - secretsCreated := 0 + g, gctx := errgroup.WithContext(ctx) + g.SetLimit(10) + + var secretsCreated atomic.Int32 for i := range plans { - secretName := plans[i].VMName + "-cloudinit" - secretLabels := map[string]string{ - constants.LabelAppName: plans[i].VMSpec.Labels[constants.LabelAppName], - constants.LabelManagedBy: constants.ManagedByValue, - constants.LabelComponent: plans[i].Component, - constants.LabelRunID: runID, - } - if err := resources.CreateCloudInitSecret(ctx, ro.client, secretName, - cfg.Namespace, plans[i].VMSpec.CloudInitUserdata, secretLabels); err != nil { - return 0, fmt.Errorf("creating cloud-init secret for %q: %w", plans[i].VMName, err) - } - plans[i].VMSpec.CloudInitSecretName = secretName - secretsCreated++ - ro.logger.Info("secret created", - slog.String("secret_name", secretName), - slog.String("namespace", cfg.Namespace)) - - _, _ = ro.auditor.RecordResource(ctx, execID, audit.ResourceRecord{ - ResourceType: "Secret", - ResourceName: secretName, - Namespace: cfg.Namespace, + g.Go(func() error { + secretName := plans[i].VMName + "-cloudinit" + secretLabels := map[string]string{ + constants.LabelAppName: plans[i].VMSpec.Labels[constants.LabelAppName], + constants.LabelManagedBy: constants.ManagedByValue, + constants.LabelComponent: plans[i].Component, + constants.LabelRunID: runID, + } + if err := resources.CreateCloudInitSecret(gctx, ro.client, secretName, + cfg.Namespace, plans[i].VMSpec.CloudInitUserdata, secretLabels); err != nil { + return fmt.Errorf("creating cloud-init secret for %q: %w", plans[i].VMName, err) + } + plans[i].VMSpec.CloudInitSecretName = secretName + secretsCreated.Add(1) + ro.logger.Info("secret created", + slog.String("secret_name", secretName), + slog.String("namespace", cfg.Namespace)) + + _, _ = ro.auditor.RecordResource(ctx, execID, audit.ResourceRecord{ + ResourceType: "Secret", + ResourceName: secretName, + Namespace: cfg.Namespace, + }) + return nil }) } - return secretsCreated, nil + if err := g.Wait(); err != nil { + return 0, err + } + return int(secretsCreated.Load()), nil } func (ro *RunOrchestrator) createVMs( diff --git a/internal/orchestrator/orchestrator_test.go b/internal/orchestrator/orchestrator_test.go index adc2faa..c773a56 100644 --- a/internal/orchestrator/orchestrator_test.go +++ b/internal/orchestrator/orchestrator_test.go @@ -248,6 +248,25 @@ var _ = Describe("RunOrchestrator", func() { Expect(err).NotTo(HaveOccurred()) } }) + + It("should create multiple secrets concurrently", func() { + logger := logging.NewLogger(buf, false) + ro := orchestrator.NewRunOrchestrator(logger, c, cfg, auditor, buf) + + result, err := ro.Run(ctx, 0, "test-run-id", []string{"cpu"}, 5) + Expect(err).NotTo(HaveOccurred()) + Expect(result.SecretCount).To(Equal(5)) + + for i := range 5 { + secret := &corev1.Secret{} + err := c.Get(ctx, client.ObjectKey{ + Name: fmt.Sprintf("virtwork-cpu-%d-cloudinit", i), + Namespace: constants.DefaultNamespace, + }, secret) + Expect(err).NotTo(HaveOccurred()) + Expect(secret.Labels[constants.LabelRunID]).To(Equal("test-run-id")) + } + }) }) Context("workloads with DataVolumes", func() {