feat: assemble Go scheduler and backend workers

This commit is contained in:
2026-09-20 20:24:45 +00:00
parent 27599631f1
commit e271537760
8 changed files with 355 additions and 35 deletions
+240
View File
@@ -0,0 +1,240 @@
package main
import (
"context"
"errors"
"fmt"
"log"
"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"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/assignmentqueue"
"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, NATSUser, NATSPassword, NATSCA, Stream, SubjectBase string
FacadeListen, FacadeURL, FacadeSPIFFEID string
CapabilityKey []byte
PodNamespace, PodImage, PodServiceAccount, SPIRECluster, SPIREClass string
PodExecutorUID, PodCapacity int
OpenSandboxURL, OpenSandboxAPIKey, OpenSandboxPool 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")
}
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...)
if err != nil {
return fmt.Errorf("connect NATS: %w", err)
}
defer natsConnection.Close()
js, err := jetstream.New(natsConnection)
if err != nil {
return fmt.Errorf("open JetStream: %w", err)
}
if _, err := js.Stream(ctx, config.Stream); err != nil {
return fmt.Errorf("open assignment stream %s: %w", config.Stream, err)
}
capabilities, err := runnerfacade.NewCapabilities(config.CapabilityKey)
if err != nil {
return err
}
registry := runnerfacade.NewRegistry()
giteaClient := giteaactions.NewClient(giteaactions.DefaultHTTPClient(), config.GiteaURL, config.GiteaUUID, config.GiteaToken)
facade := &runnerfacade.Facade{Registry: registry, Capabilities: capabilities, Upstream: giteaClient}
bootstrap := runnerbootstrap.Bootstrap{
Capabilities: capabilities, FacadeURL: config.FacadeURL, FacadeSPIFFEID: config.FacadeSPIFFEID,
}
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, string(taskassignment.BackendVM))
}
poller := taskscheduler.Poller{
Client: giteaClient,
Scheduler: &taskscheduler.Scheduler{TrustDomain: config.TrustDomain, Dispatcher: assignmentqueue.Publisher{
JetStream: js, SubjectBase: config.SubjectBase,
}},
Config: taskscheduler.PollerConfig{Version: "gitea-dynamic-runner/0.4", Labels: labels},
OnError: func(err error) { log.Printf("scheduler: %v", err) },
}
facadeServer := runnerfacade.Server{
Facade: facade, ListenAddress: config.FacadeListen, TrustDomain: config.TrustDomain,
WorkloadAPIAddr: config.WorkloadAPIAddr,
}
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 poller.Run(groupContext) })
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, ExecutorUID: config.PodExecutorUID,
}}
component, err := workerComponent(ctx, js, config, taskassignment.BackendPod, config.PodCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry)
if err != nil {
return err
}
components[controller.PodWorker] = component
}
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},
}}
component, err := workerComponent(ctx, js, config, taskassignment.BackendVM, config.VMCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry)
if err != nil {
return err
}
components[controller.VMWorker] = component
}
return controller.Run(ctx, config.Components, components)
}
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 {
return nil, err
}
return assignmentqueue.ConsumerComponent{
Consumer: consumer, Capacity: capacity,
Processor: assignmentqueue.Processor{TrustDomain: config.TrustDomain, Accepter: accepter, Claims: claims},
OnError: func(err error) { log.Printf("%s worker: %v", backend, err) },
}, nil
}
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
}
natsPassword, err := read("NATS_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"), NATSUser: env("NATS_USER", "ci-worker"), NATSPassword: natsPassword,
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"), PodExecutorUID: envInt("POD_EXECUTOR_UID", 2000), PodCapacity: envInt("POD_CAPACITY", 4),
OpenSandboxURL: os.Getenv("OPENSANDBOX_API"), OpenSandboxPool: env("OPENSANDBOX_POOL", "ci-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.VMWorker) {
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
}
@@ -0,0 +1,59 @@
package main
import (
"os"
"path/filepath"
"testing"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/controller"
)
func secretFile(t *testing.T, name, value string) string {
t.Helper()
path := filepath.Join(t.TempDir(), name)
if err := os.WriteFile(path, []byte(value+"\n"), 0o600); err != nil {
t.Fatal(err)
}
return path
}
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("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")
t.Setenv("RUNNER_FACADE_SPIFFE_ID", "spiffe://ddupan.top/ns/gitea-actions/sa/gitea-dynamic-runner")
t.Setenv("POD_EXECUTOR_IMAGE", "zot.ddupan.top/ci/gitea-runner@sha256:abc")
config, err := loadControllerConfig()
if err != nil {
t.Fatal(err)
}
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" {
t.Fatal("file secrets were not loaded")
}
if string(config.CapabilityKey) != "0123456789abcdef0123456789abcdef" || config.PodExecutorUID != 2000 {
t.Fatalf("config = %#v", config)
}
}
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("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")
t.Setenv("RUNNER_FACADE_SPIFFE_ID", "spiffe://ddupan.top/controller")
t.Setenv("OPENSANDBOX_API", "http://opensandbox.internal")
if _, err := loadControllerConfig(); err == nil {
t.Fatal("expected missing OpenSandbox API key file error")
}
}
+10 -4
View File
@@ -19,12 +19,18 @@ func main() {
}
func run() error {
if len(os.Args) != 2 {
return errors.New("usage: gitea-dynamic-runner executor")
if len(os.Args) > 2 {
return errors.New("usage: gitea-dynamic-runner [controller|executor]")
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
switch os.Args[1] {
command := "controller"
if len(os.Args) == 2 {
command = os.Args[1]
}
switch command {
case "controller":
return runController(ctx)
case "executor":
config, err := runnerbootstrap.ExecutorConfigFromEnvironment()
if err != nil {
@@ -32,6 +38,6 @@ func run() error {
}
return runnerbootstrap.RunExecutor(ctx, config)
default:
return fmt.Errorf("unknown command %q", os.Args[1])
return fmt.Errorf("unknown command %q", command)
}
}