diff --git a/cmd/gitea-dynamic-runner/controller.go b/cmd/gitea-dynamic-runner/controller.go index 8b243f0..252382c 100644 --- a/cmd/gitea-dynamic-runner/controller.go +++ b/cmd/gitea-dynamic-runner/controller.go @@ -15,6 +15,11 @@ import ( "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/controller" @@ -133,6 +138,18 @@ func runController(ctx context.Context) error { Config: taskscheduler.PollerConfig{Version: "gitea-dynamic-runner/0.4", Labels: labels}, OnError: func(err error) { log.Printf("scheduler: %v", 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, @@ -141,7 +158,9 @@ func runController(ctx context.Context) error { 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) }) + group.Go(func() error { + return runSchedulerLeader(groupContext, kubernetesClient, config.PodNamespace, leaderIdentity, poller.Run) + }) return group.Wait() }), } @@ -211,6 +230,50 @@ func runController(ctx context.Context) error { 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 != "" { diff --git a/go.mod b/go.mod index aa2b32b..941ad92 100644 --- a/go.mod +++ b/go.mod @@ -45,6 +45,7 @@ require ( github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/nats-io/nkeys v0.4.16 // indirect github.com/nats-io/nuid v1.0.1 // indirect + github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect github.com/sirupsen/logrus v1.10.2 // indirect github.com/x448/float16 v0.8.4 // indirect go.yaml.in/yaml/v2 v2.4.4 // indirect diff --git a/go.sum b/go.sum index 62546f1..e61ea34 100644 --- a/go.sum +++ b/go.sum @@ -112,6 +112,8 @@ go.opentelemetry.io/otel/sdk/metric v1.39.0 h1:cXMVVFVgsIf2YL6QkRF4Urbr/aMInf+2W go.opentelemetry.io/otel/sdk/metric v1.39.0/go.mod h1:xq9HEVH7qeX69/JnwEfp6fVq5wosJsY1mt4lLfYdVew= go.opentelemetry.io/otel/trace v1.39.0 h1:2d2vfpEDmCJ5zVYz7ijaJdOF59xLomrvj7bjt6/qCJI= go.opentelemetry.io/otel/trace v1.39.0/go.mod h1:88w4/PnZSazkGzz/w84VHpQafiU4EtqqlVdxWy+rNOA= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ= go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ= go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw=