feat: 通过 Lease 单例调度任务
test / python (pull_request) Successful in 31s
test / shell (pull_request) Successful in 34s
test / go (pull_request) Successful in 3m38s

This commit is contained in:
2026-09-21 06:35:06 +00:00
parent 5e94182308
commit 2a23f0c62e
3 changed files with 67 additions and 1 deletions
+64 -1
View File
@@ -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 != "" {