package main import ( "context" "errors" "fmt" "log/slog" "net/http" "os" "slices" "strconv" "strings" "time" "github.com/nats-io/nats.go" "github.com/nats-io/nats.go/jetstream" "golang.org/x/sync/errgroup" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" "k8s.io/client-go/tools/leaderelection" "k8s.io/client-go/tools/leaderelection/resourcelock" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/assignmentqueue" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/backendpool" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/controller" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/giteaactions" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/opensandboxbackend" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/podbackend" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/runnerbootstrap" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/runnerfacade" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/taskassignment" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/taskscheduler" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/taskworker" ) type runComponent func(context.Context) error func (function runComponent) Run(ctx context.Context) error { return function(ctx) } type controllerConfig struct { Components controller.Selection TrustDomain, WorkloadAPIAddr string GiteaURL, GiteaUUID, GiteaToken string NATSURL, NATSProducerUser, NATSProducerPassword string NATSWorkerUser, NATSWorkerPassword, NATSCA, Stream, SubjectBase string FacadeListen, FacadeURL, FacadeSPIFFEID string CapabilityKey []byte PodNamespace, PodImage, PodServiceAccount, SPIRECluster, SPIREClass string SPIREAgentID string PodExecutorUID, PodCapacity int OpenSandboxURL, OpenSandboxAPIKey, OpenSandboxPool, VMRunnerLabel string VMTimeout, VMCapacity int } func runController(ctx context.Context) error { config, err := loadControllerConfig() if err != nil { return err } if !slices.Contains(config.Components, controller.Scheduler) { return errors.New("split worker deployment is not yet safe: scheduler/facade must be enabled with workers") } if !slices.Contains(config.Components, controller.PodWorker) && !slices.Contains(config.Components, controller.VMWorker) { return errors.New("scheduler requires at least one local backend worker") } producerConnection, err := connectNATS(config.NATSURL, config.NATSProducerUser, config.NATSProducerPassword, config.NATSCA, "gitea-dynamic-runner-producer") if err != nil { return fmt.Errorf("connect NATS producer: %w", err) } defer producerConnection.Close() producerJS, err := jetstream.New(producerConnection) if err != nil { return fmt.Errorf("open producer JetStream: %w", err) } 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 { return err } registry := runnerfacade.NewRegistry() podPool := backendpool.New(config.PodCapacity) vmPool := backendpool.New(config.VMCapacity) var podExecutorBackend *podbackend.Backend var vmExecutorBackend *opensandboxbackend.Backend giteaClient := giteaactions.NewClient(giteaactions.DefaultHTTPClient(), config.GiteaURL, config.GiteaUUID, config.GiteaToken) facade := &runnerfacade.Facade{ Registry: registry, Capabilities: capabilities, Upstream: giteaClient, OnTerminal: func(ctx context.Context, assignment taskassignment.Assignment) error { switch assignment.Backend { case taskassignment.BackendPod: if podExecutorBackend == nil { return errors.New("Pod lifecycle backend is not configured") } if err := podExecutorBackend.MarkTerminal(ctx, assignment.ID); err != nil { return err } podPool.Release(assignment.ID) case taskassignment.BackendVM: if vmExecutorBackend == nil { return errors.New("VM lifecycle backend is not configured") } if err := vmExecutorBackend.MarkTerminal(ctx, assignment.ID); err != nil { return err } vmPool.Release(assignment.ID) } return nil }, } bootstrap := runnerbootstrap.Bootstrap{ Capabilities: capabilities, FacadeURL: config.FacadeURL, FacadeSPIFFEID: config.FacadeSPIFFEID, WorkloadAPIAddr: config.WorkloadAPIAddr, } labels := []string{"self-hosted"} if slices.Contains(config.Components, controller.PodWorker) { labels = append(labels, string(taskassignment.BackendPod)) } if slices.Contains(config.Components, controller.VMWorker) { labels = append(labels, config.VMRunnerLabel) } poller := taskscheduler.Poller{ Client: giteaClient, Scheduler: &taskscheduler.Scheduler{TrustDomain: config.TrustDomain, Dispatcher: assignmentqueue.Publisher{ JetStream: producerJS, SubjectBase: config.SubjectBase, }}, Config: taskscheduler.PollerConfig{Version: "gitea-dynamic-runner/0.4", Labels: labels, Capacity: config.PodCapacity + config.VMCapacity}, OnError: func(err error) { slog.Error("scheduler error", "component", "scheduler", "error", err) }, } kubernetesConfig, err := rest.InClusterConfig() if err != nil { return fmt.Errorf("load leader election Kubernetes config: %w", err) } kubernetesClient, err := kubernetes.NewForConfig(kubernetesConfig) if err != nil { return fmt.Errorf("create leader election Kubernetes client: %w", err) } leaderIdentity := strings.TrimSpace(os.Getenv("HOSTNAME")) if leaderIdentity == "" { return errors.New("HOSTNAME is required for scheduler leader election") } facadeServer := runnerfacade.Server{ Facade: facade, ListenAddress: config.FacadeListen, TrustDomain: config.TrustDomain, WorkloadAPIAddr: config.WorkloadAPIAddr, UpstreamURL: config.GiteaURL, } components := controller.Registry{ controller.Scheduler: runComponent(func(ctx context.Context) error { group, groupContext := errgroup.WithContext(ctx) group.Go(func() error { return facadeServer.Run(groupContext) }) group.Go(func() error { return runSchedulerLeader(groupContext, kubernetesClient, config.PodNamespace, leaderIdentity, poller.Run) }) return group.Wait() }), } if slices.Contains(config.Components, controller.PodWorker) { client, err := podbackend.NewInClusterClient() if err != nil { return err } backend := podbackend.Backend{API: client, Config: podbackend.Config{ Namespace: config.PodNamespace, Image: config.PodImage, ServiceAccount: config.PodServiceAccount, ExecutorArgs: []string{"executor"}, TrustDomain: config.TrustDomain, SPIRECluster: config.SPIRECluster, SPIREClass: config.SPIREClass, SPIREAgentID: config.SPIREAgentID, ExecutorUID: config.PodExecutorUID, }} podExecutorBackend = &backend assignments, err := backend.RecoverAssignments(ctx) if err != nil { return err } for _, assignment := range assignments { podPool.Restore(assignment.ID) if err := registry.RecoverClaimed(assignment); err != nil { return fmt.Errorf("recover Pod facade claim %s: %w", assignment.ID, err) } } component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendPod, config.PodCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap, OnEvent: workerEventLogger(taskassignment.BackendPod)}, registry, podPool) if err != nil { return err } lifecycle := podbackend.Lifecycle{Backend: backend, OnError: func(err error) { slog.Error("backend lifecycle error", "component", "lifecycle", "backend", taskassignment.BackendPod, "error", err) }} components[controller.PodWorker] = runComponent(func(ctx context.Context) error { group, groupContext := errgroup.WithContext(ctx) group.Go(func() error { return component.Run(groupContext) }) group.Go(func() error { return lifecycle.Run(groupContext) }) return group.Wait() }) } if slices.Contains(config.Components, controller.VMWorker) { lifecycle := opensandboxbackend.NewLifecycleClient(config.OpenSandboxURL, config.OpenSandboxAPIKey, &http.Client{Timeout: 60 * time.Second}) backend := opensandboxbackend.Backend{Lifecycle: lifecycle, Config: opensandboxbackend.Config{ Pool: config.OpenSandboxPool, Timeout: config.VMTimeout, Entrypoint: []string{"/usr/local/bin/gitea-dynamic-runner", "executor"}, Env: map[string]string{"SPIFFE_ENDPOINT_SOCKET": config.WorkloadAPIAddr}, }} vmExecutorBackend = &backend assignments, err := backend.RecoverAssignments(ctx, config.TrustDomain) if err != nil { return err } for _, assignment := range assignments { vmPool.Restore(assignment.ID) if err := registry.RecoverClaimed(assignment); err != nil { return fmt.Errorf("recover VM facade claim %s: %w", assignment.ID, err) } } component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendVM, config.VMCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap, OnEvent: workerEventLogger(taskassignment.BackendVM)}, registry, vmPool) if err != nil { return err } lifecycleReconciler := opensandboxbackend.LifecycleReconciler{Backend: backend, OnError: func(err error) { slog.Error("backend lifecycle error", "component", "lifecycle", "backend", taskassignment.BackendVM, "error", err) }} components[controller.VMWorker] = runComponent(func(ctx context.Context) error { group, groupContext := errgroup.WithContext(ctx) group.Go(func() error { return component.Run(groupContext) }) group.Go(func() error { return lifecycleReconciler.Run(groupContext) }) return group.Wait() }) } return controller.Run(ctx, config.Components, components) } func runSchedulerLeader(ctx context.Context, client kubernetes.Interface, namespace, identity string, run func(context.Context) error) error { if client == nil || namespace == "" || identity == "" || run == nil { return errors.New("leader election client, namespace, identity, and scheduler are required") } electionContext, cancel := context.WithCancel(ctx) defer cancel() result := make(chan error, 1) lock := &resourcelock.LeaseLock{ LeaseMeta: metav1.ObjectMeta{Name: "dynamic-runner-scheduler", Namespace: namespace}, Client: client.CoordinationV1(), LockConfig: resourcelock.ResourceLockConfig{ Identity: identity, }, } elector, err := leaderelection.NewLeaderElector(leaderelection.LeaderElectionConfig{ Lock: lock, LeaseDuration: 15 * time.Second, RenewDeadline: 10 * time.Second, RetryPeriod: 2 * time.Second, ReleaseOnCancel: true, Callbacks: leaderelection.LeaderCallbacks{ OnStartedLeading: func(leaderContext context.Context) { result <- run(leaderContext) cancel() }, OnStoppedLeading: func() { if ctx.Err() == nil { select { case result <- errors.New("scheduler leadership lost"): default: } } }, }, }) if err != nil { return fmt.Errorf("configure scheduler leader election: %w", err) } go elector.Run(electionContext) select { case err := <-result: return err case <-ctx.Done(): return nil } } 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, admission assignmentqueue.Admission) (controller.Component, error) { consumer, err := assignmentqueue.OpenConsumer(ctx, js, config.Stream, config.SubjectBase, backend, capacity) if err != nil { return nil, err } return assignmentqueue.ConsumerComponent{ Consumer: consumer, Capacity: capacity, Processor: assignmentqueue.Processor{TrustDomain: config.TrustDomain, Accepter: accepter, Claims: claims, Admission: admission, OnEvent: func(event assignmentqueue.Event) { slog.Info("assignment transition", "component", "worker", "event", event.Name, "backend", event.Backend, "assignment", event.AssignmentID, "retry_delay", event.RetryDelay) }}, OnError: func(err error) { slog.Error("assignment processing error", "component", "worker", "backend", backend, "error", err) }, }, nil } func workerEventLogger(backend taskassignment.Backend) func(taskworker.Event) { return func(event taskworker.Event) { slog.Info("executor transition", "component", "worker", "event", event.Name, "backend", backend, "assignment", event.AssignmentID, "executor", event.Executor, "phase", event.Phase) } } func loadControllerConfig() (controllerConfig, error) { selection, err := controller.ParseSelection(os.Getenv("COMPONENTS")) if err != nil { return controllerConfig{}, err } read := func(name string) (string, error) { path := os.Getenv(name) if path == "" { return "", fmt.Errorf("%s is required", name) } value, err := os.ReadFile(path) if err != nil { return "", fmt.Errorf("read %s: %w", name, err) } return strings.TrimSpace(string(value)), nil } uuid, err := read("GITEA_RUNNER_UUID_FILE") if err != nil { return controllerConfig{}, err } token, err := read("GITEA_RUNNER_TOKEN_FILE") if err != nil { return controllerConfig{}, err } 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 } capabilityKey, err := read("RUNNER_FACADE_CAPABILITY_KEY_FILE") if err != nil { return controllerConfig{}, err } 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"), 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"), SPIRECluster: env("SPIRE_CLUSTER", "homelab"), SPIREClass: env("SPIRE_CLASS", "spire-mgmt-spire"), SPIREAgentID: os.Getenv("SPIRE_AGENT_ID"), PodExecutorUID: envInt("POD_EXECUTOR_UID", 2000), PodCapacity: envInt("POD_CAPACITY", 4), OpenSandboxURL: os.Getenv("OPENSANDBOX_API"), OpenSandboxPool: env("OPENSANDBOX_POOL", "ci-vm"), VMRunnerLabel: env("VM_RUNNER_LABEL", "vm"), VMTimeout: envInt("VM_TIMEOUT_SECONDS", 14400), VMCapacity: envInt("VM_CAPACITY", 1), } if config.WorkloadAPIAddr == "" || config.FacadeURL == "" || config.FacadeSPIFFEID == "" { return controllerConfig{}, errors.New("SPIFFE_ENDPOINT_SOCKET, RUNNER_FACADE_URL, and RUNNER_FACADE_SPIFFE_ID are required") } if slices.Contains(selection, controller.PodWorker) && config.PodImage == "" { return controllerConfig{}, errors.New("POD_EXECUTOR_IMAGE is required for pod-worker") } if slices.Contains(selection, controller.PodWorker) && config.SPIREAgentID == "" { return controllerConfig{}, errors.New("SPIRE_AGENT_ID is required for pod-worker") } if slices.Contains(selection, controller.VMWorker) { if config.VMRunnerLabel != "vm" && config.VMRunnerLabel != "vm-dev" { return controllerConfig{}, errors.New("VM_RUNNER_LABEL must be vm or vm-dev") } if config.OpenSandboxURL == "" { return controllerConfig{}, errors.New("OPENSANDBOX_API is required for vm-worker") } config.OpenSandboxAPIKey, err = read("OPENSANDBOX_API_KEY_FILE") if err != nil { return controllerConfig{}, err } } return config, nil } func env(name, fallback string) string { if value := strings.TrimSpace(os.Getenv(name)); value != "" { return value } return fallback } func envInt(name string, fallback int) int { value := strings.TrimSpace(os.Getenv(name)) if value == "" { return fallback } parsed, err := strconv.Atoi(value) if err != nil || parsed < 1 { return fallback } return parsed }