Files
gitea-dynamic-runner/cmd/gitea-dynamic-runner/controller.go
T
panxiao81 38e8d59541
test / go (pull_request) Successful in 3m10s
test / shell (pull_request) Failing after 10m4s
test / python (pull_request) Failing after 10m4s
feat: 增加 Runner 生命周期结构化日志
2026-09-21 09:31:46 +00:00

406 lines
17 KiB
Go

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
}