diff --git a/pkg/election/election.go b/pkg/election/election.go index d019be8c..67b20752 100644 --- a/pkg/election/election.go +++ b/pkg/election/election.go @@ -59,7 +59,7 @@ func NewManager(config *kubevip.Config, k8sClientset, rwClientset *kubernetes.Cl func RunOrDie(ctx context.Context, run *RunConfig, c *kubevip.Config) error { switch c.LeaderElectionType { case "kubernetes", "": - runKubernetesLeaderElectionOrDie(ctx, run) + return runKubernetesLeaderElectionOrDie(ctx, run) case "etcd": if err := runEtcdLeaderElectionOrDie(ctx, run); err != nil { return err @@ -71,20 +71,25 @@ func RunOrDie(ctx context.Context, run *RunConfig, c *kubevip.Config) error { return nil } -func runKubernetesLeaderElectionOrDie(ctx context.Context, run *RunConfig) { +func runKubernetesLeaderElectionOrDie(ctx context.Context, run *RunConfig) error { + annotations, err := kubevip.WithLeaseVIPs(run.LeaseAnnotations, run.Config.InstanceName, run.Config.RoutingProtocol, run.VIPs) + if err != nil { + return err + } + leaseClient := run.Mgr.KubernetesClient.CoordinationV1().Leases(run.LeaseID.Namespace()) // we use the Lease lock type since edits to Leases are less common // and fewer objects in the cluster watch "all Leases". - lock := &resourcelock.LeaseLock{ + baseLock := &resourcelock.LeaseLock{ LeaseMeta: metav1.ObjectMeta{ - Name: run.LeaseID.Name(), - Namespace: run.LeaseID.Namespace(), - Annotations: run.LeaseAnnotations, + Name: run.LeaseID.Name(), + Namespace: run.LeaseID.Namespace(), }, Client: run.Mgr.KubernetesClient.CoordinationV1(), LockConfig: resourcelock.ResourceLockConfig{ Identity: run.Config.NodeName, }, } + lock := newAnnotatedLeaseLock(baseLock, leaseClient, run.LeaseID.Name(), annotations) // start the leader election code loop leaderelection.RunOrDie(ctx, leaderelection.LeaderElectionConfig{ @@ -105,6 +110,7 @@ func runKubernetesLeaderElectionOrDie(ctx context.Context, run *RunConfig) { OnNewLeader: run.OnNewLeader, }, }) + return nil } func runEtcdLeaderElectionOrDie(ctx context.Context, run *RunConfig) error { @@ -135,6 +141,7 @@ type RunConfig struct { LeaseID lease.ID Mgr *Manager LeaseAnnotations map[string]string + VIPs []string // onStartedLeading is called when this member starts leading. OnStartedLeading func(context.Context) diff --git a/pkg/election/lease_lock.go b/pkg/election/lease_lock.go new file mode 100644 index 00000000..78c1a4c0 --- /dev/null +++ b/pkg/election/lease_lock.go @@ -0,0 +1,92 @@ +package election + +import ( + "context" + + log "log/slog" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + coordinationv1client "k8s.io/client-go/kubernetes/typed/coordination/v1" + "k8s.io/client-go/tools/leaderelection/resourcelock" + "k8s.io/client-go/util/retry" +) + +type annotatedLeaseLock struct { + resourcelock.Interface + leases coordinationv1client.LeaseInterface + name string + annotations map[string]string +} + +func newAnnotatedLeaseLock(lock resourcelock.Interface, leases coordinationv1client.LeaseInterface, + name string, annotations map[string]string) resourcelock.Interface { + return &annotatedLeaseLock{Interface: lock, leases: leases, name: name, annotations: annotations} +} + +func (lock *annotatedLeaseLock) Create(ctx context.Context, record resourcelock.LeaderElectionRecord) error { + if err := lock.Interface.Create(ctx, record); err != nil { + return err + } + lock.ensure(ctx, record) + return nil +} + +func (lock *annotatedLeaseLock) Update(ctx context.Context, record resourcelock.LeaderElectionRecord) error { + if err := lock.Interface.Update(ctx, record); err != nil { + return err + } + lock.ensure(ctx, record) + return nil +} + +// ensure applies the configured annotations once this process holds the lease. Failures are +// logged rather than returned: the lease write already succeeded, so reporting an error would +// make the elector stand down while it still holds the lease. +func (lock *annotatedLeaseLock) ensure(ctx context.Context, record resourcelock.LeaderElectionRecord) { + if record.HolderIdentity != lock.Identity() { + return + } + changed, err := lock.ensureAnnotations(ctx) + if err != nil { + log.Warn("failed to annotate lease", "lease", lock.name, "err", err) + return + } + if !changed { + return + } + // Annotating out of band bumps the resourceVersion, so refresh the wrapped lock's + // cached lease or its next optimistic Update conflicts. + if _, _, err := lock.Interface.Get(ctx); err != nil { + log.Warn("failed to refresh lease after annotating", "lease", lock.name, "err", err) + } +} + +func (lock *annotatedLeaseLock) ensureAnnotations(ctx context.Context) (bool, error) { + changed := false + err := retry.RetryOnConflict(retry.DefaultRetry, func() error { + resource, err := lock.leases.Get(ctx, lock.name, metav1.GetOptions{}) + if err != nil { + return err + } + if resource.Annotations == nil { + resource.Annotations = make(map[string]string, len(lock.annotations)) + } + resourceChanged := false + for key, value := range lock.annotations { + if resource.Annotations[key] == value { + continue + } + resource.Annotations[key] = value + resourceChanged = true + } + if !resourceChanged { + return nil + } + _, err = lock.leases.Update(ctx, resource, metav1.UpdateOptions{}) + if err == nil { + changed = true + } + return err + }) + return changed, err +} diff --git a/pkg/election/lease_lock_test.go b/pkg/election/lease_lock_test.go new file mode 100644 index 00000000..f0db898e --- /dev/null +++ b/pkg/election/lease_lock_test.go @@ -0,0 +1,177 @@ +package election + +import ( + "context" + "fmt" + "testing" + + "github.com/kube-vip/kube-vip/pkg/kubevip" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/client-go/kubernetes/fake" + k8stesting "k8s.io/client-go/testing" + "k8s.io/client-go/tools/leaderelection/resourcelock" +) + +func TestAnnotatedLeaseLockPersistsAnnotationsOnCreateAndUpdate(t *testing.T) { + client := fake.NewSimpleClientset() + leaseClient := client.CoordinationV1().Leases("default") + base := &resourcelock.LeaseLock{ + LeaseMeta: metav1.ObjectMeta{Name: "lease", Namespace: "default"}, + Client: client.CoordinationV1(), + LockConfig: resourcelock.ResourceLockConfig{ + Identity: "node-a", + }, + } + annotations, err := kubevip.WithLeaseVIPs(map[string]string{"example.test/preserved": "true"}, + "release_a", 248, []string{"192.0.2.10"}) + if err != nil { + t.Fatalf("WithLeaseVIPs() error = %v", err) + } + lock := newAnnotatedLeaseLock(base, leaseClient, "lease", annotations) + record := resourcelock.LeaderElectionRecord{HolderIdentity: "node-a"} + if err := lock.Create(context.Background(), record); err != nil { + t.Fatalf("Create() error = %v", err) + } + if err := lock.Update(context.Background(), record); err != nil { + t.Fatalf("Update() error = %v", err) + } + + resource, err := leaseClient.Get(context.Background(), "lease", metav1.GetOptions{}) + if err != nil { + t.Fatalf("get Lease: %v", err) + } + value, err := kubevip.ParseLeaseVIPs(resource.Annotations[kubevip.LeaseVIPs]) + if err != nil { + t.Fatalf("ParseLeaseVIPs() error = %v", err) + } + if value.InstanceName != "release_a" || value.IFAProto != 248 || len(value.VIPs) != 1 || + value.VIPs[0] != (kubevip.LeaseVIP{Index: 0, Value: "192.0.2.10", Kind: kubevip.LeaseVIPKindAddress}) { + t.Fatalf("Lease VIP metadata = %+v", value) + } + if resource.Annotations["example.test/preserved"] != "true" { + t.Fatal("Lease update dropped a configured annotation") + } +} + +// A failed annotation write must not be reported to the leader elector: the lease itself +// was already written, and an error makes the elector stand down while it still holds it. +func TestAnnotatedLeaseLockAnnotationFailureDoesNotSurfaceToElector(t *testing.T) { + client := fake.NewSimpleClientset() + base := &resourcelock.LeaseLock{ + LeaseMeta: metav1.ObjectMeta{Name: "lease", Namespace: "default"}, + Client: client.CoordinationV1(), + LockConfig: resourcelock.ResourceLockConfig{Identity: "node-a"}, + } + + failing := fake.NewSimpleClientset() + failing.PrependReactor("get", "leases", func(k8stesting.Action) (bool, runtime.Object, error) { + return true, nil, fmt.Errorf("annotation backend unavailable") + }) + + annotations, err := kubevip.WithLeaseVIPs(nil, "release_a", 248, []string{"192.0.2.10"}) + if err != nil { + t.Fatal(err) + } + lock := newAnnotatedLeaseLock(base, failing.CoordinationV1().Leases("default"), "lease", annotations) + record := resourcelock.LeaderElectionRecord{HolderIdentity: "node-a"} + + if err := lock.Create(context.Background(), record); err != nil { + t.Fatalf("Create() error = %v, want nil so the elector keeps the lease", err) + } + if err := lock.Update(context.Background(), record); err != nil { + t.Fatalf("Update() error = %v, want nil so the elector keeps the lease", err) + } + + if _, err := client.CoordinationV1().Leases("default").Get(context.Background(), "lease", + metav1.GetOptions{}); err != nil { + t.Fatalf("wrapped lock did not write the lease: %v", err) + } +} + +func TestAnnotatedLeaseLockFollowerDoesNotOverwriteAnnotations(t *testing.T) { + client := fake.NewSimpleClientset() + leaseClient := client.CoordinationV1().Leases("default") + newBase := func(identity string) *resourcelock.LeaseLock { + return &resourcelock.LeaseLock{ + LeaseMeta: metav1.ObjectMeta{Name: "lease", Namespace: "default"}, + Client: client.CoordinationV1(), + LockConfig: resourcelock.ResourceLockConfig{Identity: identity}, + } + } + ownerBase := newBase("node-a") + followerBase := newBase("node-b") + active, err := kubevip.WithLeaseVIPs(nil, "release_a", 248, []string{"192.0.2.10"}) + if err != nil { + t.Fatal(err) + } + creator := newAnnotatedLeaseLock(ownerBase, leaseClient, "lease", active) + if err := creator.Create(context.Background(), resourcelock.LeaderElectionRecord{HolderIdentity: "node-a"}); err != nil { + t.Fatalf("Create() error = %v", err) + } + + follower, err := kubevip.WithLeaseVIPs(nil, "release_b", 249, []string{"192.0.2.20"}) + if err != nil { + t.Fatal(err) + } + observer := newAnnotatedLeaseLock(followerBase, leaseClient, "lease", follower) + if _, _, err := observer.Get(context.Background()); err != nil { + t.Fatalf("Get() error = %v", err) + } + resource, err := leaseClient.Get(context.Background(), "lease", metav1.GetOptions{}) + if err != nil { + t.Fatal(err) + } + metadata, err := kubevip.ParseLeaseVIPs(resource.Annotations[kubevip.LeaseVIPs]) + if err != nil { + t.Fatal(err) + } + if metadata.InstanceName != "release_a" || metadata.IFAProto != 248 { + t.Fatalf("follower overwrote active metadata: %+v", metadata) + } +} + +func TestAnnotatedLeaseLockReleaseDoesNotOverwriteSuccessorMetadata(t *testing.T) { + client := fake.NewSimpleClientset() + leaseClient := client.CoordinationV1().Leases("default") + newLock := func(identity, instanceName string, protocol int, vip string) resourcelock.Interface { + base := &resourcelock.LeaseLock{ + LeaseMeta: metav1.ObjectMeta{Name: "lease", Namespace: "default"}, + Client: client.CoordinationV1(), + LockConfig: resourcelock.ResourceLockConfig{Identity: identity}, + } + annotations, err := kubevip.WithLeaseVIPs(nil, instanceName, protocol, []string{vip}) + if err != nil { + t.Fatal(err) + } + return newAnnotatedLeaseLock(base, leaseClient, "lease", annotations) + } + + first := newLock("node-a", "release_a", 248, "192.0.2.10") + if err := first.Create(context.Background(), resourcelock.LeaderElectionRecord{HolderIdentity: "node-a"}); err != nil { + t.Fatal(err) + } + if err := first.Update(context.Background(), resourcelock.LeaderElectionRecord{}); err != nil { + t.Fatal(err) + } + + second := newLock("node-b", "release_b", 249, "192.0.2.20") + if _, _, err := second.Get(context.Background()); err != nil { + t.Fatal(err) + } + if err := second.Update(context.Background(), resourcelock.LeaderElectionRecord{HolderIdentity: "node-b"}); err != nil { + t.Fatal(err) + } + + resource, err := leaseClient.Get(context.Background(), "lease", metav1.GetOptions{}) + if err != nil { + t.Fatal(err) + } + metadata, err := kubevip.ParseLeaseVIPs(resource.Annotations[kubevip.LeaseVIPs]) + if err != nil { + t.Fatal(err) + } + if metadata.InstanceName != "release_b" || metadata.IFAProto != 249 || metadata.VIPs[0].Value != "192.0.2.20" { + t.Fatalf("successor metadata = %+v", metadata) + } +} diff --git a/pkg/kubevip/annotations.go b/pkg/kubevip/annotations.go index 92c39fac..284e5ac2 100644 --- a/pkg/kubevip/annotations.go +++ b/pkg/kubevip/annotations.go @@ -77,6 +77,9 @@ const ( // Name of the service lease object ServiceLease = "kube-vip.io/leaseName" + // Versioned kube-vip ownership metadata stored on Kubernetes election Leases + LeaseVIPs = "kube-vip.io/lease-vips" + // Forces kube-vip to use per service election for this particular service ForcePerServiceElection = "kube-vip.io/forcePerServiceElection" diff --git a/pkg/kubevip/lease_annotations.go b/pkg/kubevip/lease_annotations.go new file mode 100644 index 00000000..c53a05b5 --- /dev/null +++ b/pkg/kubevip/lease_annotations.go @@ -0,0 +1,134 @@ +package kubevip + +import ( + "encoding/json" + "fmt" + "net/netip" + "slices" + "strings" +) + +const LeaseVIPsVersion = "v1" + +// LeaseVIPKind distinguishes literal addresses from names that resolve to one. +type LeaseVIPKind string + +const ( + LeaseVIPKindAddress LeaseVIPKind = "address" + LeaseVIPKindName LeaseVIPKind = "name" +) + +type LeaseVIPsValue struct { + Version string `json:"version"` + InstanceName string `json:"instance_name"` + IFAProto int `json:"ifa_proto"` + VIPs []LeaseVIP `json:"vips"` +} + +type LeaseVIP struct { + Index int `json:"index"` + Value string `json:"value"` + Kind LeaseVIPKind `json:"kind"` +} + +func WithLeaseVIPs(annotations map[string]string, instanceName string, ifaProto int, vips []string) (map[string]string, error) { + result := make(map[string]string, len(annotations)+1) + for key, value := range annotations { + result[key] = value + } + + encoded, err := json.Marshal(LeaseVIPsValue{ + Version: LeaseVIPsVersion, + InstanceName: instanceName, + IFAProto: ifaProto, + VIPs: normalizeLeaseVIPs(vips), + }) + if err != nil { + return nil, fmt.Errorf("encode %s annotation: %w", LeaseVIPs, err) + } + result[LeaseVIPs] = string(encoded) + return result, nil +} + +func ParseLeaseVIPs(value string) (LeaseVIPsValue, error) { + var parsed LeaseVIPsValue + if err := json.Unmarshal([]byte(value), &parsed); err != nil { + return LeaseVIPsValue{}, fmt.Errorf("decode %s annotation: %w", LeaseVIPs, err) + } + if parsed.Version != LeaseVIPsVersion { + return LeaseVIPsValue{}, fmt.Errorf("unsupported %s annotation version %q", LeaseVIPs, parsed.Version) + } + for index, vip := range parsed.VIPs { + if vip.Index != index { + return LeaseVIPsValue{}, fmt.Errorf("invalid %s VIP index %d at position %d", LeaseVIPs, vip.Index, index) + } + switch vip.Kind { + case LeaseVIPKindAddress, LeaseVIPKindName: + default: + return LeaseVIPsValue{}, fmt.Errorf("invalid %s VIP kind %q at index %d", LeaseVIPs, vip.Kind, vip.Index) + } + } + return parsed, nil +} + +func normalizeLeaseVIPs(values []string) []LeaseVIP { + unique := make(map[string]struct{}, len(values)) + addresses := make([]string, 0, len(values)) + for _, value := range values { + for candidate := range strings.SplitSeq(value, ",") { + candidate = strings.TrimSpace(candidate) + if candidate == "" { + continue + } + if _, exists := unique[candidate]; exists { + continue + } + unique[candidate] = struct{}{} + addresses = append(addresses, candidate) + } + } + // Sorting keeps the annotation byte-identical however callers happen to order VIPs. + slices.SortFunc(addresses, compareLeaseVIPs) + + result := make([]LeaseVIP, 0, len(addresses)) + for _, address := range addresses { + kind := LeaseVIPKindName + if _, isAddress := leaseVIPAddress(address); isAddress { + kind = LeaseVIPKindAddress + } + result = append(result, LeaseVIP{Index: len(result), Value: address, Kind: kind}) + } + return result +} + +// compareLeaseVIPs orders addresses numerically and ahead of names, which keeps VIPs like +// 10.0.0.2 and 10.0.0.10 in the order an operator expects. Values that are not addresses, +// such as DNS records, are kept and ordered lexically. +func compareLeaseVIPs(a, b string) int { + addressA, isAddressA := leaseVIPAddress(a) + addressB, isAddressB := leaseVIPAddress(b) + switch { + case isAddressA && isAddressB: + if order := addressA.Compare(addressB); order != 0 { + return order + } + // Distinct spellings of one address still need a stable order. + return strings.Compare(a, b) + case isAddressA: + return -1 + case isAddressB: + return 1 + default: + return strings.Compare(a, b) + } +} + +func leaseVIPAddress(value string) (netip.Addr, bool) { + if address, err := netip.ParseAddr(value); err == nil { + return address.Unmap(), true + } + if prefix, err := netip.ParsePrefix(value); err == nil { + return prefix.Addr().Unmap(), true + } + return netip.Addr{}, false +} diff --git a/pkg/kubevip/lease_annotations_test.go b/pkg/kubevip/lease_annotations_test.go new file mode 100644 index 00000000..b66006f2 --- /dev/null +++ b/pkg/kubevip/lease_annotations_test.go @@ -0,0 +1,94 @@ +package kubevip + +import ( + "slices" + "testing" +) + +func TestWithLeaseVIPsEncodesVersionedInstanceOwnership(t *testing.T) { + base := map[string]string{"example.test/preserved": "true", LeaseVIPs: "stale"} + annotations, err := WithLeaseVIPs(base, "release_a", 248, []string{ + "2001:db8::10/128", "192.0.2.10", "192.0.2.10/32", "api.example.test", + }) + if err != nil { + t.Fatalf("WithLeaseVIPs() error = %v", err) + } + if annotations["example.test/preserved"] != "true" { + t.Fatal("WithLeaseVIPs() dropped an existing annotation") + } + if base[LeaseVIPs] != "stale" { + t.Fatal("WithLeaseVIPs() mutated the input annotations") + } + + value, err := ParseLeaseVIPs(annotations[LeaseVIPs]) + if err != nil { + t.Fatalf("ParseLeaseVIPs() error = %v", err) + } + if value.Version != LeaseVIPsVersion || value.InstanceName != "release_a" || value.IFAProto != 248 { + t.Fatalf("Lease VIP metadata = %+v", value) + } + // Values are stored verbatim so DNS records survive alongside addresses. + want := []LeaseVIP{ + {Index: 0, Value: "192.0.2.10", Kind: LeaseVIPKindAddress}, + {Index: 1, Value: "192.0.2.10/32", Kind: LeaseVIPKindAddress}, + {Index: 2, Value: "2001:db8::10/128", Kind: LeaseVIPKindAddress}, + {Index: 3, Value: "api.example.test", Kind: LeaseVIPKindName}, + } + if !slices.Equal(value.VIPs, want) { + t.Fatalf("Lease VIPs = %v, want %v", value.VIPs, want) + } +} + +// The annotation is rewritten whenever a node starts campaigning, so the encoding +// has to be stable even when callers collect the same VIPs in a different order. +func TestWithLeaseVIPsIsIndependentOfInputOrder(t *testing.T) { + first, err := WithLeaseVIPs(nil, "release_a", 248, []string{ + "2001:db8::10", "192.0.2.10", "10.0.0.2", "10.0.0.10", + }) + if err != nil { + t.Fatalf("WithLeaseVIPs() error = %v", err) + } + second, err := WithLeaseVIPs(nil, "release_a", 248, []string{ + "10.0.0.10", "192.0.2.10", "2001:db8::10", "10.0.0.2", + }) + if err != nil { + t.Fatalf("WithLeaseVIPs() error = %v", err) + } + if first[LeaseVIPs] != second[LeaseVIPs] { + t.Fatalf("annotation changed with input order:\n%s\n%s", first[LeaseVIPs], second[LeaseVIPs]) + } + + value, err := ParseLeaseVIPs(first[LeaseVIPs]) + if err != nil { + t.Fatalf("ParseLeaseVIPs() error = %v", err) + } + want := []string{"10.0.0.2", "10.0.0.10", "192.0.2.10", "2001:db8::10"} + if len(value.VIPs) != len(want) { + t.Fatalf("Lease VIPs = %v, want %v", value.VIPs, want) + } + for index, address := range want { + if value.VIPs[index] != (LeaseVIP{Index: index, Value: address, Kind: LeaseVIPKindAddress}) { + t.Fatalf("Lease VIPs = %v, want %v", value.VIPs, want) + } + } +} + +func TestParseLeaseVIPsRejectsUnknownVersion(t *testing.T) { + if _, err := ParseLeaseVIPs(`{"version":"v2","instance_name":"release_a","ifa_proto":248,"vips":[]}`); err == nil { + t.Fatal("ParseLeaseVIPs() accepted an unknown version") + } +} + +func TestParseLeaseVIPsRejectsUnknownKind(t *testing.T) { + if _, err := ParseLeaseVIPs( + `{"version":"v1","instance_name":"release_a","ifa_proto":248,"vips":[{"index":0,"value":"192.0.2.10","kind":"cidr"}]}`, + ); err == nil { + t.Fatal("ParseLeaseVIPs() accepted an unknown VIP kind") + } +} + +func TestParseLeaseVIPsRejectsOutOfOrderIndexes(t *testing.T) { + if _, err := ParseLeaseVIPs(`{"version":"v1","instance_name":"release_a","ifa_proto":248,"vips":[{"index":1,"value":"192.0.2.10"}]}`); err == nil { + t.Fatal("ParseLeaseVIPs() accepted an out-of-order VIP index") + } +}