diff --git a/README.md b/README.md index 41eaff3..afa813e 100644 --- a/README.md +++ b/README.md @@ -60,7 +60,8 @@ credential 都从挂载文件读取,不接受明文环境变量: - `GITEA_RUNNER_UUID_FILE`、`GITEA_RUNNER_TOKEN_FILE`:scheduler 的常驻 RunnerService registration;该 credential 不下发给 executor。 -- `NATS_PASSWORD_FILE`:assignment stream 的连接密码。 +- `NATS_PRODUCER_PASSWORD_FILE`、`NATS_WORKER_PASSWORD_FILE`:分别使用现有最小权限的 + `ci-producer` publish 连接和 `ci-worker` pull/ACK 连接,controller 不合并权限。 - `RUNNER_FACADE_CAPABILITY_KEY_FILE`:至少 32 字节的 controller HMAC key。 - `OPENSANDBOX_API_KEY_FILE`:仅启用 `vm-worker` 时读取。 diff --git a/cmd/gitea-dynamic-runner/controller.go b/cmd/gitea-dynamic-runner/controller.go index 0607b69..2e2b1ee 100644 --- a/cmd/gitea-dynamic-runner/controller.go +++ b/cmd/gitea-dynamic-runner/controller.go @@ -36,7 +36,8 @@ type controllerConfig struct { Components controller.Selection TrustDomain, WorkloadAPIAddr string GiteaURL, GiteaUUID, GiteaToken string - NATSURL, NATSUser, NATSPassword, NATSCA, Stream, SubjectBase string + NATSURL, NATSProducerUser, NATSProducerPassword string + NATSWorkerUser, NATSWorkerPassword, NATSCA, Stream, SubjectBase string FacadeListen, FacadeURL, FacadeSPIFFEID string CapabilityKey []byte PodNamespace, PodImage, PodServiceAccount, SPIRECluster, SPIREClass string @@ -57,25 +58,27 @@ func runController(ctx context.Context) error { return errors.New("scheduler requires at least one local backend worker") } - natsOptions := []nats.Option{nats.Name("gitea-dynamic-runner")} - if config.NATSUser != "" || config.NATSPassword != "" { - natsOptions = append(natsOptions, nats.UserInfo(config.NATSUser, config.NATSPassword)) - } - if config.NATSCA != "" { - natsOptions = append(natsOptions, nats.RootCAs(config.NATSCA)) - } - natsConnection, err := nats.Connect(config.NATSURL, natsOptions...) + producerConnection, err := connectNATS(config.NATSURL, config.NATSProducerUser, config.NATSProducerPassword, config.NATSCA, "gitea-dynamic-runner-producer") if err != nil { - return fmt.Errorf("connect NATS: %w", err) + return fmt.Errorf("connect NATS producer: %w", err) } - defer natsConnection.Close() - js, err := jetstream.New(natsConnection) + defer producerConnection.Close() + producerJS, err := jetstream.New(producerConnection) if err != nil { - return fmt.Errorf("open JetStream: %w", err) + return fmt.Errorf("open producer JetStream: %w", err) } - if _, err := js.Stream(ctx, config.Stream); err != nil { + if _, err := producerJS.Stream(ctx, config.Stream); err != nil { return fmt.Errorf("open assignment stream %s: %w", config.Stream, err) } + workerConnection, err := connectNATS(config.NATSURL, config.NATSWorkerUser, config.NATSWorkerPassword, config.NATSCA, "gitea-dynamic-runner-worker") + if err != nil { + return fmt.Errorf("connect NATS worker: %w", err) + } + defer workerConnection.Close() + workerJS, err := jetstream.New(workerConnection) + if err != nil { + return fmt.Errorf("open worker JetStream: %w", err) + } capabilities, err := runnerfacade.NewCapabilities(config.CapabilityKey) if err != nil { @@ -98,7 +101,7 @@ func runController(ctx context.Context) error { poller := taskscheduler.Poller{ Client: giteaClient, Scheduler: &taskscheduler.Scheduler{TrustDomain: config.TrustDomain, Dispatcher: assignmentqueue.Publisher{ - JetStream: js, SubjectBase: config.SubjectBase, + JetStream: producerJS, SubjectBase: config.SubjectBase, }}, Config: taskscheduler.PollerConfig{Version: "gitea-dynamic-runner/0.4", Labels: labels}, OnError: func(err error) { log.Printf("scheduler: %v", err) }, @@ -126,7 +129,7 @@ func runController(ctx context.Context) error { ExecutorArgs: []string{"executor"}, TrustDomain: config.TrustDomain, SPIRECluster: config.SPIRECluster, SPIREClass: config.SPIREClass, ExecutorUID: config.PodExecutorUID, }} - component, err := workerComponent(ctx, js, config, taskassignment.BackendPod, config.PodCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry) + component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendPod, config.PodCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry) if err != nil { return err } @@ -139,7 +142,7 @@ func runController(ctx context.Context) error { Entrypoint: []string{"/usr/local/bin/gitea-dynamic-runner", "executor"}, Env: map[string]string{"SPIFFE_ENDPOINT_SOCKET": config.WorkloadAPIAddr}, }} - component, err := workerComponent(ctx, js, config, taskassignment.BackendVM, config.VMCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry) + component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendVM, config.VMCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry) if err != nil { return err } @@ -148,6 +151,14 @@ func runController(ctx context.Context) error { return controller.Run(ctx, config.Components, components) } +func connectNATS(server, user, password, caFile, clientName string) (*nats.Conn, error) { + options := []nats.Option{nats.Name(clientName), nats.UserInfo(user, password)} + if caFile != "" { + options = append(options, nats.RootCAs(caFile)) + } + return nats.Connect(server, options...) +} + func workerComponent(ctx context.Context, js jetstream.JetStream, config controllerConfig, backend taskassignment.Backend, capacity int, accepter assignmentqueue.Accepter, claims assignmentqueue.Claims) (controller.Component, error) { consumer, err := assignmentqueue.OpenConsumer(ctx, js, config.Stream, config.SubjectBase, backend, capacity) if err != nil { @@ -184,7 +195,11 @@ func loadControllerConfig() (controllerConfig, error) { if err != nil { return controllerConfig{}, err } - natsPassword, err := read("NATS_PASSWORD_FILE") + producerPassword, err := read("NATS_PRODUCER_PASSWORD_FILE") + if err != nil { + return controllerConfig{}, err + } + workerPassword, err := read("NATS_WORKER_PASSWORD_FILE") if err != nil { return controllerConfig{}, err } @@ -195,7 +210,9 @@ func loadControllerConfig() (controllerConfig, error) { config := controllerConfig{ Components: selection, TrustDomain: env("TRUST_DOMAIN", "ddupan.top"), WorkloadAPIAddr: os.Getenv("SPIFFE_ENDPOINT_SOCKET"), GiteaURL: env("GITEA_INSTANCE_URL", "https://git.ddupan.top"), GiteaUUID: uuid, GiteaToken: token, - NATSURL: env("NATS_URL", "tls://nats.ad.ddupan.top:4222"), NATSUser: env("NATS_USER", "ci-worker"), NATSPassword: natsPassword, + NATSURL: env("NATS_URL", "tls://nats.ad.ddupan.top:4222"), + NATSProducerUser: env("NATS_PRODUCER_USER", "ci-producer"), NATSProducerPassword: producerPassword, + NATSWorkerUser: env("NATS_WORKER_USER", "ci-worker"), NATSWorkerPassword: workerPassword, NATSCA: os.Getenv("NATS_CA_FILE"), Stream: env("NATS_STREAM", "CI_RUNNER"), SubjectBase: env("NATS_SUBJECT_BASE", "ci.runner"), FacadeListen: env("RUNNER_FACADE_LISTEN", ":8443"), FacadeURL: os.Getenv("RUNNER_FACADE_URL"), FacadeSPIFFEID: os.Getenv("RUNNER_FACADE_SPIFFE_ID"), CapabilityKey: []byte(capabilityKey), PodNamespace: env("POD_NAMESPACE", "gitea-actions"), PodImage: os.Getenv("POD_EXECUTOR_IMAGE"), PodServiceAccount: env("POD_SERVICE_ACCOUNT", "gitea-task-executor"), diff --git a/cmd/gitea-dynamic-runner/controller_test.go b/cmd/gitea-dynamic-runner/controller_test.go index 0e14844..edb0d4d 100644 --- a/cmd/gitea-dynamic-runner/controller_test.go +++ b/cmd/gitea-dynamic-runner/controller_test.go @@ -21,7 +21,8 @@ func TestLoadControllerConfigUsesFileSecrets(t *testing.T) { t.Setenv("COMPONENTS", "scheduler,pod-worker") t.Setenv("GITEA_RUNNER_UUID_FILE", secretFile(t, "uuid", "scheduler-uuid")) t.Setenv("GITEA_RUNNER_TOKEN_FILE", secretFile(t, "token", "scheduler-token")) - t.Setenv("NATS_PASSWORD_FILE", secretFile(t, "nats", "nats-password")) + t.Setenv("NATS_PRODUCER_PASSWORD_FILE", secretFile(t, "nats-producer", "producer-password")) + t.Setenv("NATS_WORKER_PASSWORD_FILE", secretFile(t, "nats-worker", "worker-password")) t.Setenv("RUNNER_FACADE_CAPABILITY_KEY_FILE", secretFile(t, "capability", "0123456789abcdef0123456789abcdef")) t.Setenv("SPIFFE_ENDPOINT_SOCKET", "unix:///run/spire/agent-sockets/spire-agent.sock") t.Setenv("RUNNER_FACADE_URL", "https://gitea-runner-facade.gitea-actions.svc:8443") @@ -35,7 +36,7 @@ func TestLoadControllerConfigUsesFileSecrets(t *testing.T) { if len(config.Components) != 2 || config.Components[0] != controller.Scheduler || config.Components[1] != controller.PodWorker { t.Fatalf("components = %#v", config.Components) } - if config.GiteaUUID != "scheduler-uuid" || config.GiteaToken != "scheduler-token" || config.NATSPassword != "nats-password" { + if config.GiteaUUID != "scheduler-uuid" || config.GiteaToken != "scheduler-token" || config.NATSProducerPassword != "producer-password" || config.NATSWorkerPassword != "worker-password" { t.Fatal("file secrets were not loaded") } if string(config.CapabilityKey) != "0123456789abcdef0123456789abcdef" || config.PodExecutorUID != 2000 { @@ -47,7 +48,8 @@ func TestLoadControllerConfigRequiresOpenSandboxSecretOnlyForVM(t *testing.T) { t.Setenv("COMPONENTS", "scheduler,vm-worker") t.Setenv("GITEA_RUNNER_UUID_FILE", secretFile(t, "uuid", "uuid")) t.Setenv("GITEA_RUNNER_TOKEN_FILE", secretFile(t, "token", "token")) - t.Setenv("NATS_PASSWORD_FILE", secretFile(t, "nats", "password")) + t.Setenv("NATS_PRODUCER_PASSWORD_FILE", secretFile(t, "nats-producer", "producer-password")) + t.Setenv("NATS_WORKER_PASSWORD_FILE", secretFile(t, "nats-worker", "worker-password")) t.Setenv("RUNNER_FACADE_CAPABILITY_KEY_FILE", secretFile(t, "capability", "0123456789abcdef0123456789abcdef")) t.Setenv("SPIFFE_ENDPOINT_SOCKET", "unix:///run/spire/agent-sockets/spire-agent.sock") t.Setenv("RUNNER_FACADE_URL", "https://facade:8443")