diff --git a/CMakeLists.txt b/CMakeLists.txt index 2254bb4b..c6dd0699 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -45,6 +45,16 @@ option(STORE_USE_REDIS "build mooncake store with redis" OFF) if (STORE_USE_REDIS) add_compile_definitions(STORE_USE_REDIS) endif() +option(STORE_USE_K8S_LEASE "build mooncake store with K8s Lease leader election" OFF) +if (STORE_USE_K8S_LEASE) + if (STORE_USE_ETCD) + message(FATAL_ERROR "STORE_USE_K8S_LEASE and STORE_USE_ETCD cannot be enabled together because both build Go c-shared HA backends.") + endif() + if (USE_ETCD AND NOT USE_ETCD_LEGACY) + message(FATAL_ERROR "STORE_USE_K8S_LEASE cannot be enabled with non-legacy USE_ETCD because both build Go c-shared libraries in the same process.") + endif() + add_compile_definitions(STORE_USE_K8S_LEASE) +endif() option(STORE_USE_JEMALLOC "Use jemalloc in mooncake store master" OFF) diff --git a/mooncake-common/CMakeLists.txt b/mooncake-common/CMakeLists.txt index 9f07f7d8..dade5b98 100644 --- a/mooncake-common/CMakeLists.txt +++ b/mooncake-common/CMakeLists.txt @@ -2,6 +2,10 @@ if ((USE_ETCD AND NOT USE_ETCD_LEGACY) OR STORE_USE_ETCD) add_subdirectory(etcd) endif() +if (STORE_USE_K8S_LEASE) + add_subdirectory(k8s-lease) +endif() + include_directories(${CMAKE_CURRENT_SOURCE_DIR}/include) add_subdirectory(src) diff --git a/mooncake-common/k8s-lease/CMakeLists.txt b/mooncake-common/k8s-lease/CMakeLists.txt new file mode 100644 index 00000000..7ccc2ac6 --- /dev/null +++ b/mooncake-common/k8s-lease/CMakeLists.txt @@ -0,0 +1,20 @@ +add_custom_command( + OUTPUT ${CMAKE_CURRENT_BINARY_DIR}/libk8s_lease_wrapper.so + COMMAND bash -c "go mod tidy" && bash -c "go build -buildmode=c-shared -o ${CMAKE_CURRENT_BINARY_DIR}/libk8s_lease_wrapper.so k8s_lease_wrapper.go" && cp ${CMAKE_CURRENT_BINARY_DIR}/libk8s_lease_wrapper.h ${CMAKE_CURRENT_SOURCE_DIR} + WORKING_DIRECTORY ${CMAKE_CURRENT_SOURCE_DIR} + COMMENT "Building K8s Lease Go shared library" + DEPENDS k8s_lease_wrapper.go +) + +set(K8S_LEASE_WRAPPER_INCLUDE ${CMAKE_CURRENT_BINARY_DIR}/libk8s_lease_wrapper.h) +set(K8S_LEASE_WRAPPER_LIB ${CMAKE_CURRENT_BINARY_DIR}/libk8s_lease_wrapper.so) + +add_custom_target( + build_k8s_lease_wrapper + DEPENDS ${K8S_LEASE_WRAPPER_LIB} +) + +install( + FILES ${K8S_LEASE_WRAPPER_LIB} + DESTINATION lib +) diff --git a/mooncake-common/k8s-lease/cmd/envtest-server/main.go b/mooncake-common/k8s-lease/cmd/envtest-server/main.go new file mode 100644 index 00000000..efda29de --- /dev/null +++ b/mooncake-common/k8s-lease/cmd/envtest-server/main.go @@ -0,0 +1,61 @@ +// envtest-server starts a real kube-apiserver + etcd via envtest, writes the +// KUBECONFIG path to stdout, and blocks until SIGTERM or SIGINT. This lets +// C++ tests launch it as a subprocess and talk to a real K8s API without a +// full cluster. +package main + +import ( + "fmt" + "os" + "os/signal" + "path/filepath" + "syscall" + + "k8s.io/client-go/tools/clientcmd" + clientcmdapi "k8s.io/client-go/tools/clientcmd/api" + "sigs.k8s.io/controller-runtime/pkg/envtest" +) + +func main() { + env := &envtest.Environment{} + + cfg, err := env.Start() + if err != nil { + fmt.Fprintf(os.Stderr, "envtest start failed: %v\n", err) + os.Exit(1) + } + + // Write a KUBECONFIG file that points at the envtest kube-apiserver. + kubeconfigPath := filepath.Join(os.TempDir(), fmt.Sprintf("envtest-kubeconfig-%d", os.Getpid())) + kubeconfig := clientcmdapi.NewConfig() + kubeconfig.Clusters["envtest"] = &clientcmdapi.Cluster{ + Server: cfg.Host, + CertificateAuthorityData: cfg.CAData, + } + kubeconfig.AuthInfos["envtest"] = &clientcmdapi.AuthInfo{ + ClientCertificateData: cfg.CertData, + ClientKeyData: cfg.KeyData, + } + kubeconfig.Contexts["envtest"] = &clientcmdapi.Context{ + Cluster: "envtest", + AuthInfo: "envtest", + } + kubeconfig.CurrentContext = "envtest" + + if err := clientcmd.WriteToFile(*kubeconfig, kubeconfigPath); err != nil { + fmt.Fprintf(os.Stderr, "failed to write kubeconfig: %v\n", err) + env.Stop() + os.Exit(1) + } + + // Print the kubeconfig path — the parent process reads this from stdout. + fmt.Println(kubeconfigPath) + + // Block until SIGTERM or SIGINT. + sigCh := make(chan os.Signal, 1) + signal.Notify(sigCh, syscall.SIGTERM, syscall.SIGINT) + <-sigCh + + os.Remove(kubeconfigPath) + env.Stop() +} diff --git a/mooncake-common/k8s-lease/go.mod b/mooncake-common/k8s-lease/go.mod new file mode 100644 index 00000000..4bfc2052 --- /dev/null +++ b/mooncake-common/k8s-lease/go.mod @@ -0,0 +1,60 @@ +module github.com/kvcache-ai/Mooncake/mooncake-common/k8s-lease + +go 1.24.0 + +require ( + k8s.io/api v0.34.3 + k8s.io/apimachinery v0.34.3 + k8s.io/client-go v0.34.3 + k8s.io/utils v0.0.0-20251002143259-bc988d571ff4 + sigs.k8s.io/controller-runtime v0.22.5 +) + +require ( + github.com/beorn7/perks v1.0.1 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect + github.com/emicklei/go-restful/v3 v3.12.2 // indirect + github.com/evanphx/json-patch/v5 v5.9.11 // indirect + github.com/fxamacker/cbor/v2 v2.9.0 // indirect + github.com/go-logr/logr v1.4.3 // indirect + github.com/go-openapi/jsonpointer v0.21.0 // indirect + github.com/go-openapi/jsonreference v0.20.2 // indirect + github.com/go-openapi/swag v0.23.0 // indirect + github.com/gogo/protobuf v1.3.2 // indirect + github.com/google/gnostic-models v0.7.0 // indirect + github.com/google/go-cmp v0.7.0 // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/josharian/intern v1.0.0 // indirect + github.com/json-iterator/go v1.1.12 // indirect + github.com/mailru/easyjson v0.7.7 // indirect + github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect + github.com/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee // indirect + github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect + github.com/pmezard/go-difflib v1.0.0 // indirect + github.com/prometheus/client_golang v1.23.2 // indirect + github.com/prometheus/client_model v0.6.2 // indirect + github.com/prometheus/common v0.66.1 // indirect + github.com/prometheus/procfs v0.16.1 // indirect + github.com/spf13/pflag v1.0.9 // indirect + github.com/x448/float16 v0.8.4 // indirect + go.yaml.in/yaml/v2 v2.4.3 // indirect + go.yaml.in/yaml/v3 v3.0.4 // indirect + golang.org/x/net v0.47.0 // indirect + golang.org/x/oauth2 v0.30.0 // indirect + golang.org/x/sys v0.38.0 // indirect + golang.org/x/term v0.37.0 // indirect + golang.org/x/text v0.31.0 // indirect + golang.org/x/time v0.9.0 // indirect + google.golang.org/protobuf v1.36.8 // indirect + gopkg.in/evanphx/json-patch.v4 v4.13.0 // indirect + gopkg.in/inf.v0 v0.9.1 // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect + k8s.io/apiextensions-apiserver v0.34.3 // indirect + k8s.io/klog/v2 v2.130.1 // indirect + k8s.io/kube-openapi v0.0.0-20250910181357-589584f1c912 // indirect + sigs.k8s.io/json v0.0.0-20250730193827-2d320260d730 // indirect + sigs.k8s.io/randfill v1.0.0 // indirect + sigs.k8s.io/structured-merge-diff/v6 v6.3.2-0.20260122202528-d9cc6641c482 // indirect + sigs.k8s.io/yaml v1.6.0 // indirect +) diff --git a/mooncake-common/k8s-lease/k8s_lease_integration_test.go b/mooncake-common/k8s-lease/k8s_lease_integration_test.go new file mode 100644 index 00000000..5db713fe --- /dev/null +++ b/mooncake-common/k8s-lease/k8s_lease_integration_test.go @@ -0,0 +1,571 @@ +//go:build integration + +package main + +import ( + "context" + "fmt" + "os" + "sync" + "testing" + "time" + + 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" + + "sigs.k8s.io/controller-runtime/pkg/envtest" +) + +var ( + testEnv *envtest.Environment + testConfig *rest.Config +) + +type electionStateNoRelease struct { + cancel context.CancelFunc + elected chan struct{} + lost chan struct{} +} + +func TestMain(m *testing.M) { + testEnv = &envtest.Environment{} + + var err error + testConfig, err = testEnv.Start() + if err != nil { + fmt.Fprintf(os.Stderr, "failed to start envtest: %v\n", err) + os.Exit(1) + } + + // Set up global client for the wrapper + client, err := kubernetes.NewForConfig(testConfig) + if err != nil { + fmt.Fprintf(os.Stderr, "failed to create clientset: %v\n", err) + testEnv.Stop() + os.Exit(1) + } + clientMutex.Lock() + globalClient = client + clientMutex.Unlock() + + code := m.Run() + + testEnv.Stop() + os.Exit(code) +} + +func runElectionWithoutRelease(namespace, leaseName, identity string, + leaseDurationSec, renewDeadlineSec, retryPeriodSec int) (*electionStateNoRelease, error) { + if err := ensureClientInitialized(); err != nil { + return nil, err + } + + ctx, cancel := context.WithCancel(context.Background()) + state := &electionStateNoRelease{ + cancel: cancel, + elected: make(chan struct{}), + lost: make(chan struct{}), + } + + lock := &resourcelock.LeaseLock{ + LeaseMeta: metav1.ObjectMeta{ + Name: leaseName, + Namespace: namespace, + }, + Client: globalClient.CoordinationV1(), + LockConfig: resourcelock.ResourceLockConfig{ + Identity: identity, + }, + } + + le, err := leaderelection.NewLeaderElector(leaderelection.LeaderElectionConfig{ + Lock: lock, + LeaseDuration: time.Duration(leaseDurationSec) * time.Second, + RenewDeadline: time.Duration(renewDeadlineSec) * time.Second, + RetryPeriod: time.Duration(retryPeriodSec) * time.Second, + ReleaseOnCancel: false, + Callbacks: leaderelection.LeaderCallbacks{ + OnStartedLeading: func(ctx context.Context) { + close(state.elected) + <-ctx.Done() + }, + OnStoppedLeading: func() { + close(state.lost) + }, + }, + }) + if err != nil { + cancel() + return nil, fmt.Errorf("failed to create leader elector: %w", err) + } + + go le.Run(ctx) + return state, nil +} + +// TestSingleLeaderElection verifies a single candidate becomes leader. +func TestSingleLeaderElection(t *testing.T) { + ns := "default" + lease := "single-election-test" + identity := "node-1:8080" + + err := runElection(ns, lease, identity, 5, 4, 1) + if err != nil { + t.Fatalf("runElection failed: %v", err) + } + + // Wait for elected + key := electionKey(ns, lease) + electionMutex.Lock() + state := elections[key] + electionMutex.Unlock() + + select { + case <-state.elected: + // success + case <-time.After(15 * time.Second): + t.Fatal("timed out waiting for election") + } + + // Verify holder via getHolder + holder, transitions, err := getHolder(ns, lease) + if err != nil { + t.Fatalf("getHolder failed: %v", err) + } + if holder != identity { + t.Errorf("expected holder %q, got %q", identity, holder) + } + // First election — transitions should be 0 or 1 + if transitions < 0 { + t.Errorf("expected non-negative transitions, got %d", transitions) + } + + // Cancel the election + electionMutex.Lock() + state = elections[key] + electionMutex.Unlock() + state.cancel() + + select { + case <-state.lost: + // success + case <-time.After(10 * time.Second): + t.Fatal("timed out waiting for election loss after cancel") + } +} + +// TestLeaderEpoch verifies leaseTransitions increments across elections. +func TestLeaderEpoch(t *testing.T) { + ns := "default" + lease := "epoch-test" + + // First election + err := runElection(ns, lease, "node-epoch-1:8080", 5, 4, 1) + if err != nil { + t.Fatalf("first runElection failed: %v", err) + } + + key := electionKey(ns, lease) + electionMutex.Lock() + state1 := elections[key] + electionMutex.Unlock() + + select { + case <-state1.elected: + case <-time.After(15 * time.Second): + t.Fatal("timed out on first election") + } + + _, trans1, _ := getHolder(ns, lease) + + // Cancel first election and wait for loss + state1.cancel() + select { + case <-state1.lost: + case <-time.After(10 * time.Second): + t.Fatal("timed out waiting for first election loss") + } + + // Wait for lease to expire / be released + time.Sleep(2 * time.Second) + + // Second election + err = runElection(ns, lease, "node-epoch-2:8080", 5, 4, 1) + if err != nil { + t.Fatalf("second runElection failed: %v", err) + } + + electionMutex.Lock() + state2 := elections[key] + electionMutex.Unlock() + + select { + case <-state2.elected: + case <-time.After(15 * time.Second): + t.Fatal("timed out on second election") + } + + _, trans2, _ := getHolder(ns, lease) + if trans2 <= trans1 { + t.Errorf("expected transitions to increment: first=%d, second=%d", trans1, trans2) + } + + state2.cancel() + select { + case <-state2.lost: + case <-time.After(10 * time.Second): + t.Fatal("timed out waiting for second election loss") + } +} + +// TestSequentialLeadershipHandoff tests that a second candidate can acquire +// leadership after the first one releases it. +func TestSequentialLeadershipHandoff(t *testing.T) { + ns := "default" + lease := "two-candidate-test" + + err1 := runElection(ns, lease, "candidate-a:8080", 5, 4, 1) + if err1 != nil { + t.Fatalf("first runElection failed: %v", err1) + } + + key := electionKey(ns, lease) + electionMutex.Lock() + stateA := elections[key] + electionMutex.Unlock() + + // Wait for first candidate to win + select { + case <-stateA.elected: + case <-time.After(15 * time.Second): + t.Fatal("timed out waiting for first candidate") + } + + // Verify holder is candidate-a + holder, _, err := getHolder(ns, lease) + if err != nil { + t.Fatalf("getHolder failed: %v", err) + } + if holder != "candidate-a:8080" { + t.Errorf("expected candidate-a, got %q", holder) + } + + // Cancel candidate-a + stateA.cancel() + select { + case <-stateA.lost: + case <-time.After(10 * time.Second): + t.Fatal("timed out waiting for candidate-a loss") + } + + // Wait for lease to expire + time.Sleep(2 * time.Second) + + // Start candidate-b + err2 := runElection(ns, lease, "candidate-b:8080", 5, 4, 1) + if err2 != nil { + t.Fatalf("second runElection failed: %v", err2) + } + + electionMutex.Lock() + stateB := elections[key] + electionMutex.Unlock() + + select { + case <-stateB.elected: + case <-time.After(15 * time.Second): + t.Fatal("timed out waiting for candidate-b") + } + + holder, _, err = getHolder(ns, lease) + if err != nil { + t.Fatalf("getHolder after takeover failed: %v", err) + } + if holder != "candidate-b:8080" { + t.Errorf("expected candidate-b, got %q", holder) + } + + stateB.cancel() + select { + case <-stateB.lost: + case <-time.After(10 * time.Second): + t.Fatal("timed out waiting for candidate-b loss") + } +} + +// TestConcurrentCandidateElection starts two candidates simultaneously and +// verifies that exactly one wins leadership. +func TestConcurrentCandidateElection(t *testing.T) { + ns := "default" + lease := "concurrent-election-test" + + type result struct { + identity string + elected bool + } + + candidates := []string{"candidate-a:8080", "candidate-b:8080"} + results := make(chan result, len(candidates)) + + lock := func(identity string) *resourcelock.LeaseLock { + return &resourcelock.LeaseLock{ + LeaseMeta: metav1.ObjectMeta{ + Name: lease, + Namespace: ns, + }, + Client: globalClient.CoordinationV1(), + LockConfig: resourcelock.ResourceLockConfig{ + Identity: identity, + }, + } + } + + var wg sync.WaitGroup + for _, id := range candidates { + wg.Add(1) + go func(identity string) { + defer wg.Done() + + // Short timeout: enough for one to acquire, but the loser + // times out before the winner's lease could expire. + ctx, cancel := context.WithTimeout(context.Background(), 8*time.Second) + defer cancel() + + elected := make(chan struct{}) + le, err := leaderelection.NewLeaderElector(leaderelection.LeaderElectionConfig{ + Lock: lock(identity), + LeaseDuration: 5 * time.Second, + RenewDeadline: 3 * time.Second, + RetryPeriod: 1 * time.Second, + ReleaseOnCancel: true, + Callbacks: leaderelection.LeaderCallbacks{ + OnStartedLeading: func(ctx context.Context) { + close(elected) + <-ctx.Done() + }, + OnStoppedLeading: func() {}, + }, + }) + if err != nil { + t.Errorf("NewLeaderElector(%s): %v", identity, err) + return + } + + go le.Run(ctx) + + select { + case <-elected: + results <- result{identity, true} + // Keep holding until context expires (8s total). + // Winner does NOT release early, so loser cannot + // re-acquire within its own 8s window. + <-ctx.Done() + case <-ctx.Done(): + results <- result{identity, false} + } + }(id) + } + + wg.Wait() + close(results) + + winners := 0 + for r := range results { + if r.elected { + winners++ + t.Logf("winner: %s", r.identity) + } + } + + if winners != 1 { + t.Fatalf("expected exactly 1 winner, got %d", winners) + } +} + +// TestCancelElection tests that cancelling an election makes WaitLost return. +func TestCancelElection(t *testing.T) { + ns := "default" + lease := "cancel-test" + + err := runElection(ns, lease, "cancel-node:8080", 5, 4, 1) + if err != nil { + t.Fatalf("runElection failed: %v", err) + } + + key := electionKey(ns, lease) + electionMutex.Lock() + state := elections[key] + electionMutex.Unlock() + + // Wait for elected + select { + case <-state.elected: + case <-time.After(15 * time.Second): + t.Fatal("timed out waiting for election") + } + + // Cancel + state.cancel() + + // WaitLost should return promptly + select { + case <-state.lost: + // success + case <-time.After(10 * time.Second): + t.Fatal("WaitLost did not return after cancel") + } +} + +// TestGetHolderDuringElection verifies getHolder works while election is active. +func TestGetHolderDuringElection(t *testing.T) { + ns := "default" + lease := "active-get-holder-test" + identity := "active-node:8080" + + err := runElection(ns, lease, identity, 5, 4, 1) + if err != nil { + t.Fatalf("runElection failed: %v", err) + } + + key := electionKey(ns, lease) + electionMutex.Lock() + state := elections[key] + electionMutex.Unlock() + + select { + case <-state.elected: + case <-time.After(15 * time.Second): + t.Fatal("timed out waiting for election") + } + + // Concurrent getHolder calls during active election + var wg sync.WaitGroup + for i := 0; i < 5; i++ { + wg.Add(1) + go func() { + defer wg.Done() + holder, _, err := getHolder(ns, lease) + if err != nil { + t.Errorf("getHolder during election failed: %v", err) + return + } + if holder != identity { + t.Errorf("expected %q, got %q", identity, holder) + } + }() + } + wg.Wait() + + state.cancel() + <-state.lost +} + +// TestGetHolderReturnsEmptyAfterLeaderDeath verifies that after a leader stops +// renewing its lease without releasing it, getHolder returns an empty holder +// once the lease expires. This is the integration-level counterpart to the +// unit test TestGetHolderReturnsEmptyForExpiredLease. +func TestGetHolderReturnsEmptyAfterLeaderDeath(t *testing.T) { + ns := "default" + lease := "expired-leader-test" + identity := "doomed-leader:8080" + + // Acquire leadership without ReleaseOnCancel so canceling simulates a dead + // leader that stops renewing and leaves the old holder until expiry. + state, err := runElectionWithoutRelease(ns, lease, identity, 5, 4, 1) + if err != nil { + t.Fatalf("runElection failed: %v", err) + } + + select { + case <-state.elected: + case <-time.After(15 * time.Second): + t.Fatal("timed out waiting for election") + } + + // Verify holder while active. + holder, _, err := getHolder(ns, lease) + if err != nil { + t.Fatalf("getHolder (active) failed: %v", err) + } + if holder != identity { + t.Fatalf("expected active holder %q, got %q", identity, holder) + } + + // Simulate leader death: stop renewing without explicitly releasing. + state.cancel() + select { + case <-state.lost: + case <-time.After(10 * time.Second): + t.Fatal("timed out waiting for loss") + } + + // Wait for the lease to expire (leaseDuration=5s, add margin). + time.Sleep(7 * time.Second) + + // After expiry, getHolder must return empty holder so that the + // supervisor will attempt acquisition. + holder, _, err = getHolder(ns, lease) + if err != nil { + t.Fatalf("getHolder (expired) failed: %v", err) + } + if holder != "" { + t.Errorf("expected empty holder after lease expiry, got %q", holder) + } +} + +// TestFailoverAfterLeaderDeath verifies that a new candidate can acquire +// leadership after the previous leader dies and its lease expires. +func TestFailoverAfterLeaderDeath(t *testing.T) { + ns := "default" + lease := "failover-test" + + // First leader acquires without ReleaseOnCancel so canceling leaves the + // old holder in place until the lease naturally expires. + state1, err := runElectionWithoutRelease(ns, lease, "leader-1:8080", 5, 4, 1) + if err != nil { + t.Fatalf("first runElection failed: %v", err) + } + + select { + case <-state1.elected: + case <-time.After(15 * time.Second): + t.Fatal("timed out waiting for first election") + } + + // Simulate crash: cancel without release, wait for expiry. + state1.cancel() + <-state1.lost + time.Sleep(7 * time.Second) + + // Second candidate should be able to acquire. + err = runElection(ns, lease, "leader-2:8080", 5, 4, 1) + if err != nil { + t.Fatalf("second runElection failed: %v", err) + } + + key := electionKey(ns, lease) + electionMutex.Lock() + state2 := elections[key] + electionMutex.Unlock() + + select { + case <-state2.elected: + // success — failover worked + case <-time.After(15 * time.Second): + t.Fatal("second candidate failed to acquire after leader death") + } + + holder, _, err := getHolder(ns, lease) + if err != nil { + t.Fatalf("getHolder after failover failed: %v", err) + } + if holder != "leader-2:8080" { + t.Errorf("expected new leader %q, got %q", "leader-2:8080", holder) + } + + state2.cancel() + <-state2.lost +} diff --git a/mooncake-common/k8s-lease/k8s_lease_wrapper.go b/mooncake-common/k8s-lease/k8s_lease_wrapper.go new file mode 100644 index 00000000..72f90731 --- /dev/null +++ b/mooncake-common/k8s-lease/k8s_lease_wrapper.go @@ -0,0 +1,489 @@ +package main + +/* +#include +#include +#include + +// Trampoline to invoke C/C++ callback safely from Go via cgo. +typedef void (*holder_change_cb_t)(void* ctx, + const char* holder, size_t holderSize, + int64_t leaseTransitions); + +static inline void call_holder_change_cb(holder_change_cb_t func, void* ctx, + const char* holder, size_t holderSize, + int64_t leaseTransitions) { + func(ctx, holder, holderSize, leaseTransitions); +} +*/ +import "C" + +import ( + "context" + "fmt" + "os" + "sync" + "time" + "unsafe" + + coordinationv1 "k8s.io/api/coordination/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/watch" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/rest" + "k8s.io/client-go/tools/clientcmd" + "k8s.io/client-go/tools/leaderelection" + "k8s.io/client-go/tools/leaderelection/resourcelock" +) + +// electionState holds the runtime state for a single leader election. +type electionState struct { + cancel context.CancelFunc + elected chan struct{} // closed when OnStartedLeading fires + lost chan struct{} // closed when OnStoppedLeading fires + err error // set before lost is closed, if any + transitions int64 // set before elected is closed +} + +// watchState holds the runtime state for a single Lease watch. +type watchState struct { + cancel context.CancelFunc +} + +var ( + globalClient kubernetes.Interface + clientMutex sync.Mutex + initClientFn = initClient + + elections = make(map[string]*electionState) + electionMutex sync.Mutex + + watches = make(map[string]*watchState) + watchMutex sync.Mutex +) + +func electionKey(namespace, leaseName string) string { + return namespace + "/" + leaseName +} + +func ensureClientInitialized() error { + clientMutex.Lock() + initialized := globalClient != nil + clientMutex.Unlock() + if initialized { + return nil + } + return initClientFn() +} + +// initClient creates the K8s clientset from in-cluster config or KUBECONFIG. +func initClient() error { + clientMutex.Lock() + defer clientMutex.Unlock() + if globalClient != nil { + return nil + } + + config, err := rest.InClusterConfig() + if err != nil { + // Fall back to KUBECONFIG + kubeconfig := os.Getenv("KUBECONFIG") + if kubeconfig == "" { + home := os.Getenv("HOME") + if home != "" { + kubeconfig = home + "/.kube/config" + } + } + config, err = clientcmd.BuildConfigFromFlags("", kubeconfig) + if err != nil { + return fmt.Errorf("failed to build k8s config: %w", err) + } + } + + client, err := kubernetes.NewForConfig(config) + if err != nil { + return fmt.Errorf("failed to create k8s clientset: %w", err) + } + globalClient = client + return nil +} + +// runElection starts a leader election goroutine for the given namespace/leaseName. +func runElection(namespace, leaseName, identity string, + leaseDurationSec, renewDeadlineSec, retryPeriodSec int) error { + if err := ensureClientInitialized(); err != nil { + return err + } + + key := electionKey(namespace, leaseName) + + electionMutex.Lock() + if _, exists := elections[key]; exists { + electionMutex.Unlock() + return fmt.Errorf("election already running for %s", key) + } + + ctx, cancel := context.WithCancel(context.Background()) + state := &electionState{ + cancel: cancel, + elected: make(chan struct{}), + lost: make(chan struct{}), + } + elections[key] = state + electionMutex.Unlock() + + lock := &resourcelock.LeaseLock{ + LeaseMeta: metav1.ObjectMeta{ + Name: leaseName, + Namespace: namespace, + }, + Client: globalClient.CoordinationV1(), + LockConfig: resourcelock.ResourceLockConfig{ + Identity: identity, + }, + } + + le, err := leaderelection.NewLeaderElector(leaderelection.LeaderElectionConfig{ + Lock: lock, + LeaseDuration: time.Duration(leaseDurationSec) * time.Second, + RenewDeadline: time.Duration(renewDeadlineSec) * time.Second, + RetryPeriod: time.Duration(retryPeriodSec) * time.Second, + ReleaseOnCancel: true, + Callbacks: leaderelection.LeaderCallbacks{ + OnStartedLeading: func(ctx context.Context) { + _, transitions, err := getHolder(namespace, leaseName) + if err == nil { + state.transitions = transitions + } + close(state.elected) + // Block until context is cancelled (leadership lost or explicit cancel) + <-ctx.Done() + }, + OnStoppedLeading: func() { + close(state.lost) + // Auto-cleanup: remove from map so the same key can be reused. + electionMutex.Lock() + if elections[key] == state { + delete(elections, key) + } + electionMutex.Unlock() + }, + }, + }) + if err != nil { + electionMutex.Lock() + delete(elections, key) + electionMutex.Unlock() + cancel() + return fmt.Errorf("failed to create leader elector: %w", err) + } + + go le.Run(ctx) + return nil +} + +// getHolder reads the current Lease holder identity and transitions. +func getHolder(namespace, leaseName string) (string, int64, error) { + if err := ensureClientInitialized(); err != nil { + return "", 0, err + } + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + + lease, err := globalClient.CoordinationV1().Leases(namespace).Get(ctx, leaseName, metav1.GetOptions{}) + if err != nil { + return "", 0, fmt.Errorf("failed to get lease: %w", err) + } + + holder := "" + if lease.Spec.HolderIdentity != nil { + holder = *lease.Spec.HolderIdentity + } + transitions := int64(0) + if lease.Spec.LeaseTransitions != nil { + transitions = int64(*lease.Spec.LeaseTransitions) + } + + // Treat expired leases as having no holder so that the C++ supervisor + // will attempt acquisition instead of going to standby. + if holder != "" && lease.Spec.RenewTime != nil && lease.Spec.LeaseDurationSeconds != nil { + expiry := lease.Spec.RenewTime.Time.Add(time.Duration(*lease.Spec.LeaseDurationSeconds) * time.Second) + if time.Now().After(expiry) { + holder = "" + } + } + return holder, transitions, nil +} + +//export K8sLeaseInit +func K8sLeaseInit(errMsg **C.char) C.int { + if err := ensureClientInitialized(); err != nil { + *errMsg = C.CString(err.Error()) + return -1 + } + return 0 +} + +//export K8sLeaseRunElection +func K8sLeaseRunElection( + ns, leaseName, identity *C.char, + leaseDurationSec, renewDeadlineSec, retryPeriodSec C.int, + errMsg **C.char, +) C.int { + nsStr := C.GoString(ns) + ln := C.GoString(leaseName) + id := C.GoString(identity) + + err := runElection(nsStr, ln, id, + int(leaseDurationSec), int(renewDeadlineSec), int(retryPeriodSec)) + if err != nil { + *errMsg = C.CString(err.Error()) + return -1 + } + return 0 +} + +//export K8sLeaseWaitElected +func K8sLeaseWaitElected( + ns, leaseName *C.char, + timeoutSec C.int, + leaseTransitions *C.longlong, + errMsg **C.char, +) C.int { + key := electionKey(C.GoString(ns), C.GoString(leaseName)) + + electionMutex.Lock() + state, exists := elections[key] + electionMutex.Unlock() + + if !exists { + *errMsg = C.CString("no election running for " + key) + return -1 + } + + timeout := time.Duration(timeoutSec) * time.Second + + // Wait for elected, lost, or timeout + select { + case <-state.elected: + *leaseTransitions = C.longlong(state.transitions) + return 0 + case <-state.lost: + *errMsg = C.CString("election lost before becoming leader") + return -1 + case <-time.After(timeout): + state.cancel() + <-state.lost + *errMsg = C.CString("election timed out after " + fmt.Sprintf("%d", int(timeoutSec)) + "s") + return -1 + } +} + +//export K8sLeaseWaitLost +func K8sLeaseWaitLost( + ns, leaseName *C.char, + errMsg **C.char, +) C.int { + key := electionKey(C.GoString(ns), C.GoString(leaseName)) + + electionMutex.Lock() + state, exists := elections[key] + electionMutex.Unlock() + + if !exists { + // Already cleaned up by OnStoppedLeading — election is over. + return 0 + } + + <-state.lost + + if state.err != nil { + *errMsg = C.CString(state.err.Error()) + return -1 + } + return 0 +} + +//export K8sLeaseCancelElection +func K8sLeaseCancelElection( + ns, leaseName *C.char, + errMsg **C.char, +) C.int { + key := electionKey(C.GoString(ns), C.GoString(leaseName)) + + electionMutex.Lock() + state, exists := elections[key] + electionMutex.Unlock() + + if !exists { + // Idempotent — no error if no election + return 0 + } + + state.cancel() + return 0 +} + +//export K8sLeaseGetHolder +func K8sLeaseGetHolder( + ns, leaseName *C.char, + holderIdentity **C.char, + leaseTransitions *C.longlong, + errMsg **C.char, +) C.int { + nsStr := C.GoString(ns) + ln := C.GoString(leaseName) + + holder, transitions, err := getHolder(nsStr, ln) + if err != nil { + if apierrors.IsNotFound(err) { + *holderIdentity = nil + *leaseTransitions = 0 + return 1 + } + errStr := err.Error() + *errMsg = C.CString(errStr) + return -1 + } + + if holder == "" { + *holderIdentity = nil + } else { + *holderIdentity = C.CString(holder) + } + *leaseTransitions = C.longlong(transitions) + return 0 +} + +//export K8sLeaseWatchHolder +func K8sLeaseWatchHolder( + ns, leaseName *C.char, + callbackCtx unsafe.Pointer, + callbackFunc C.holder_change_cb_t, + errMsg **C.char, +) C.int { + nsStr := C.GoString(ns) + ln := C.GoString(leaseName) + key := electionKey(nsStr, ln) + + if callbackFunc == nil { + *errMsg = C.CString("callback function is nil") + return -1 + } + if err := ensureClientInitialized(); err != nil { + *errMsg = C.CString(err.Error()) + return -1 + } + + watchMutex.Lock() + if _, exists := watches[key]; exists { + watchMutex.Unlock() + *errMsg = C.CString("watch already running for " + key) + return -1 + } + + ctx, cancel := context.WithCancel(context.Background()) + watches[key] = &watchState{cancel: cancel} + watchMutex.Unlock() + + go func() { + defer func() { + watchMutex.Lock() + delete(watches, key) + watchMutex.Unlock() + }() + + for { + select { + case <-ctx.Done(): + return + default: + } + + watcher, err := globalClient.CoordinationV1().Leases(nsStr).Watch(ctx, metav1.ListOptions{ + FieldSelector: "metadata.name=" + ln, + }) + if err != nil { + select { + case <-ctx.Done(): + return + default: + time.Sleep(time.Second) + continue + } + } + + for event := range watcher.ResultChan() { + select { + case <-ctx.Done(): + watcher.Stop() + return + default: + } + + if event.Type == watch.Modified || event.Type == watch.Added { + lease, ok := event.Object.(*coordinationv1.Lease) + if !ok { + continue + } + holder := "" + if lease.Spec.HolderIdentity != nil { + holder = *lease.Spec.HolderIdentity + } + transitions := int64(0) + if lease.Spec.LeaseTransitions != nil { + transitions = int64(*lease.Spec.LeaseTransitions) + } + + var holderPtr *C.char + var holderSize C.size_t + if holder != "" { + holderPtr = C.CString(holder) + holderSize = C.size_t(len(holder)) + } + + C.call_holder_change_cb(callbackFunc, callbackCtx, + holderPtr, holderSize, C.int64_t(transitions)) + + if holderPtr != nil { + C.free(unsafe.Pointer(holderPtr)) + } + } + } + + // Watch channel closed — retry unless cancelled + select { + case <-ctx.Done(): + return + default: + time.Sleep(time.Second) + } + } + }() + + return 0 +} + +//export K8sLeaseCancelWatch +func K8sLeaseCancelWatch( + ns, leaseName *C.char, + errMsg **C.char, +) C.int { + key := electionKey(C.GoString(ns), C.GoString(leaseName)) + + watchMutex.Lock() + state, exists := watches[key] + watchMutex.Unlock() + + if !exists { + // Idempotent + return 0 + } + + state.cancel() + return 0 +} + +func main() {} diff --git a/mooncake-common/k8s-lease/k8s_lease_wrapper_test.go b/mooncake-common/k8s-lease/k8s_lease_wrapper_test.go new file mode 100644 index 00000000..705b579a --- /dev/null +++ b/mooncake-common/k8s-lease/k8s_lease_wrapper_test.go @@ -0,0 +1,293 @@ +package main + +import ( + "context" + "fmt" + "testing" + "time" + + coordinationv1 "k8s.io/api/coordination/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/kubernetes/fake" + "k8s.io/utils/ptr" +) + +// swapClient replaces globalClient and returns the old one. +func swapClient(newClient kubernetes.Interface) kubernetes.Interface { + clientMutex.Lock() + defer clientMutex.Unlock() + old := globalClient + globalClient = newClient + return old +} + +// TestGetHolderWithFakeClient tests getHolder using a fake K8s clientset. +func TestGetHolderWithFakeClient(t *testing.T) { + holderID := "node-1:8080" + transitions := int32(3) + lease := &coordinationv1.Lease{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-lease", + Namespace: "default", + }, + Spec: coordinationv1.LeaseSpec{ + HolderIdentity: &holderID, + LeaseTransitions: &transitions, + }, + } + + fakeClient := fake.NewSimpleClientset(lease) + old := swapClient(fakeClient) + defer swapClient(old) + + holder, trans, err := getHolder("default", "test-lease") + if err != nil { + t.Fatalf("getHolder failed: %v", err) + } + if holder != holderID { + t.Errorf("expected holder %q, got %q", holderID, holder) + } + if trans != int64(transitions) { + t.Errorf("expected transitions %d, got %d", transitions, trans) + } +} + +// TestGetHolderNotFound tests getHolder when the Lease does not exist. +func TestGetHolderNotFound(t *testing.T) { + fakeClient := fake.NewSimpleClientset() + old := swapClient(fakeClient) + defer swapClient(old) + + _, _, err := getHolder("default", "nonexistent") + if err == nil { + t.Fatal("expected error for nonexistent lease, got nil") + } +} + +// TestGetHolderEmptyIdentity tests getHolder when holder is nil. +func TestGetHolderEmptyIdentity(t *testing.T) { + lease := &coordinationv1.Lease{ + ObjectMeta: metav1.ObjectMeta{ + Name: "empty-lease", + Namespace: "default", + }, + Spec: coordinationv1.LeaseSpec{}, + } + fakeClient := fake.NewSimpleClientset(lease) + old := swapClient(fakeClient) + defer swapClient(old) + + holder, trans, err := getHolder("default", "empty-lease") + if err != nil { + t.Fatalf("getHolder failed: %v", err) + } + if holder != "" { + t.Errorf("expected empty holder, got %q", holder) + } + if trans != 0 { + t.Errorf("expected 0 transitions, got %d", trans) + } +} + +// TestGetHolderReturnsEmptyForExpiredLease verifies that getHolder treats a +// lease whose renewTime + leaseDuration is in the past as having no holder. +// This is critical for failover: when a leader pod dies without releasing the +// lease, standbys must see an empty holder so the supervisor attempts +// acquisition instead of looping in standby. +func TestGetHolderReturnsEmptyForExpiredLease(t *testing.T) { + holderID := "dead-leader:8080" + leaseDuration := int32(5) + transitions := int32(2) + expiredRenewTime := metav1.NewMicroTime(time.Now().Add(-10 * time.Second)) + + lease := &coordinationv1.Lease{ + ObjectMeta: metav1.ObjectMeta{ + Name: "expired-lease", + Namespace: "default", + }, + Spec: coordinationv1.LeaseSpec{ + HolderIdentity: &holderID, + LeaseDurationSeconds: &leaseDuration, + LeaseTransitions: &transitions, + RenewTime: &expiredRenewTime, + }, + } + + fakeClient := fake.NewSimpleClientset(lease) + old := swapClient(fakeClient) + defer swapClient(old) + + holder, trans, err := getHolder("default", "expired-lease") + if err != nil { + t.Fatalf("getHolder failed: %v", err) + } + if holder != "" { + t.Errorf("expected empty holder for expired lease, got %q", holder) + } + // Transitions should still be reported even for expired leases. + if trans != int64(transitions) { + t.Errorf("expected transitions %d, got %d", transitions, trans) + } +} + +// TestGetHolderReturnsHolderForActiveLease verifies that getHolder returns the +// holder identity when the lease is still active (renewTime + leaseDuration is +// in the future). +func TestGetHolderReturnsHolderForActiveLease(t *testing.T) { + holderID := "active-leader:8080" + leaseDuration := int32(15) + transitions := int32(1) + recentRenewTime := metav1.NewMicroTime(time.Now()) + + lease := &coordinationv1.Lease{ + ObjectMeta: metav1.ObjectMeta{ + Name: "active-lease", + Namespace: "default", + }, + Spec: coordinationv1.LeaseSpec{ + HolderIdentity: &holderID, + LeaseDurationSeconds: &leaseDuration, + LeaseTransitions: &transitions, + RenewTime: &recentRenewTime, + }, + } + + fakeClient := fake.NewSimpleClientset(lease) + old := swapClient(fakeClient) + defer swapClient(old) + + holder, trans, err := getHolder("default", "active-lease") + if err != nil { + t.Fatalf("getHolder failed: %v", err) + } + if holder != holderID { + t.Errorf("expected holder %q, got %q", holderID, holder) + } + if trans != int64(transitions) { + t.Errorf("expected transitions %d, got %d", transitions, trans) + } +} + +// TestElectionKeyFormat tests the election key construction. +func TestElectionKeyFormat(t *testing.T) { + tests := []struct { + ns, name, want string + }{ + {"default", "leader", "default/leader"}, + {"kube-system", "my-lock", "kube-system/my-lock"}, + {"", "bare", "/bare"}, + } + for _, tc := range tests { + got := electionKey(tc.ns, tc.name) + if got != tc.want { + t.Errorf("electionKey(%q, %q) = %q, want %q", tc.ns, tc.name, got, tc.want) + } + } +} + +// TestLeaseCRUDWithFakeClient tests basic Lease CRUD via the K8s API. +func TestLeaseCRUDWithFakeClient(t *testing.T) { + fakeClient := fake.NewSimpleClientset() + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + + holderID := "node-a:9090" + transitions := int32(0) + lease := &coordinationv1.Lease{ + ObjectMeta: metav1.ObjectMeta{ + Name: "crud-test", + Namespace: "default", + }, + Spec: coordinationv1.LeaseSpec{ + HolderIdentity: &holderID, + LeaseTransitions: &transitions, + }, + } + created, err := fakeClient.CoordinationV1().Leases("default").Create(ctx, lease, metav1.CreateOptions{}) + if err != nil { + t.Fatalf("create lease failed: %v", err) + } + if *created.Spec.HolderIdentity != holderID { + t.Errorf("created holder = %q, want %q", *created.Spec.HolderIdentity, holderID) + } + + newHolder := "node-b:9090" + newTransitions := int32(1) + created.Spec.HolderIdentity = &newHolder + created.Spec.LeaseTransitions = &newTransitions + updated, err := fakeClient.CoordinationV1().Leases("default").Update(ctx, created, metav1.UpdateOptions{}) + if err != nil { + t.Fatalf("update lease failed: %v", err) + } + if *updated.Spec.HolderIdentity != newHolder { + t.Errorf("updated holder = %q, want %q", *updated.Spec.HolderIdentity, newHolder) + } + if *updated.Spec.LeaseTransitions != newTransitions { + t.Errorf("updated transitions = %d, want %d", *updated.Spec.LeaseTransitions, newTransitions) + } + + got, err := fakeClient.CoordinationV1().Leases("default").Get(ctx, "crud-test", metav1.GetOptions{}) + if err != nil { + t.Fatalf("get lease failed: %v", err) + } + if *got.Spec.HolderIdentity != newHolder { + t.Errorf("got holder = %q, want %q", *got.Spec.HolderIdentity, newHolder) + } + + err = fakeClient.CoordinationV1().Leases("default").Delete(ctx, "crud-test", metav1.DeleteOptions{}) + if err != nil { + t.Fatalf("delete lease failed: %v", err) + } + + _, err = fakeClient.CoordinationV1().Leases("default").Get(ctx, "crud-test", metav1.GetOptions{}) + if err == nil { + t.Fatal("expected error after delete, got nil") + } +} + +// TestConcurrentGetHolder tests concurrent calls to getHolder. +func TestConcurrentGetHolder(t *testing.T) { + holderID := "concurrent-node:8080" + lease := &coordinationv1.Lease{ + ObjectMeta: metav1.ObjectMeta{ + Name: "concurrent-lease", + Namespace: "default", + }, + Spec: coordinationv1.LeaseSpec{ + HolderIdentity: &holderID, + LeaseTransitions: ptr.To(int32(5)), + }, + } + fakeClient := fake.NewSimpleClientset(lease) + old := swapClient(fakeClient) + defer swapClient(old) + + const n = 10 + errCh := make(chan error, n) + for i := 0; i < n; i++ { + go func() { + holder, trans, err := getHolder("default", "concurrent-lease") + if err != nil { + errCh <- err + return + } + if holder != holderID { + errCh <- fmt.Errorf("expected holder %q, got %q", holderID, holder) + return + } + if trans != 5 { + errCh <- fmt.Errorf("expected transitions 5, got %d", trans) + return + } + errCh <- nil + }() + } + + for i := 0; i < n; i++ { + if err := <-errCh; err != nil { + t.Fatalf("concurrent getHolder failed: %v", err) + } + } +}