From 31088ca8f25f85c3ae1d1db46b5ef3936d8854e0 Mon Sep 17 00:00:00 2001 From: Hank Mueller Date: Mon, 17 Aug 2026 04:24:13 +0000 Subject: [PATCH] flock-agent: GC orphaned allocations; retry birdc on socket-not-ready (#5) flock-agent: GC orphaned allocations; retry birdc on socket-not-ready --- pkg/agent/anycast_linux.go | 10 ++++ pkg/agent/bird.go | 48 ++++++++++++--- pkg/agent/handlers.go | 31 ++++++++++ pkg/agent/orphan_test.go | 119 +++++++++++++++++++++++++++++++++++++ pkg/agent/runtime_linux.go | 47 +++++++++++++++ 5 files changed, 247 insertions(+), 8 deletions(-) create mode 100644 pkg/agent/orphan_test.go diff --git a/pkg/agent/anycast_linux.go b/pkg/agent/anycast_linux.go index 53e4fd4..2f912e8 100644 --- a/pkg/agent/anycast_linux.go +++ b/pkg/agent/anycast_linux.go @@ -172,6 +172,16 @@ func (r *AnycastReconciler) renderBird(desired map[string]anycastTarget) { if a.State != StateCommitted { continue } + // Defense-in-depth: skip IPs for pods that are no longer on this + // node. The startup and periodic GC clean the store, but this guard + // prevents advertising a stale route in the window before GC runs + // and on any reconcile pass for IPs that outlast a GC tick. + if r.Pods != nil { + pod, ok := r.Pods.Get(a.Namespace, a.PodName) + if !ok || (a.OwnerUID != "" && string(pod.UID) != a.OwnerUID) { + continue + } + } if ip := net.ParseIP(a.IP6); ip != nil && !ipInAny(ip, nodeV6) { add(ip) } diff --git a/pkg/agent/bird.go b/pkg/agent/bird.go index 08cc040..c714b35 100644 --- a/pkg/agent/bird.go +++ b/pkg/agent/bird.go @@ -8,6 +8,7 @@ import ( "os" "os/exec" "path/filepath" + "strings" "sync" "time" @@ -137,6 +138,15 @@ func (b *BirdManager) scheduleReload() { }) } +// birdcMaxAttempts and birdcRetryDelay bound the startup-socket retry loop. +// BIRD may not have bound its control socket yet when flock-agent first +// tries to configure it; 20 × 500 ms = 10 s covers the typical BIRD +// startup window without blocking indefinitely. +const ( + birdcMaxAttempts = 20 + birdcRetryDelay = 500 * time.Millisecond +) + func (b *BirdManager) reload() { birdctl := b.BirdctlPath if birdctl == "" { @@ -146,18 +156,40 @@ func (b *BirdManager) reload() { if socket == "" { socket = "/run/flock/bird.ctl" } - cmd := exec.Command(birdctl, "-s", socket, "configure") - out, err := cmd.CombinedOutput() - if err != nil { - // First-run case: bird may not be ready yet — retry on next change. - if errors.Is(err, exec.ErrNotFound) || os.IsNotExist(err) { - b.Logger.Warn("birdc not available", "err", err) + + for attempt := 1; attempt <= birdcMaxAttempts; attempt++ { + if attempt > 1 { + time.Sleep(birdcRetryDelay) + } + // Socket absent → BIRD hasn't bound it yet; retry. + if _, err := os.Stat(socket); os.IsNotExist(err) { + b.Logger.Debug("birdc socket not ready, retrying", + "attempt", attempt, "socket", socket) + continue + } + cmd := exec.Command(birdctl, "-s", socket, "configure") + out, err := cmd.CombinedOutput() + if err == nil { + b.Logger.Info("birdc configure ok", "out", string(out)) return } - b.Logger.Warn("birdc reload failed", "err", err, "out", string(out)) + if errors.Is(err, exec.ErrNotFound) { + b.Logger.Warn("birdc not found", "err", err) + return + } + outStr := string(out) + // "Unable to connect" means BIRD exists but isn't listening yet. + if strings.Contains(outStr, "Unable to connect") { + b.Logger.Debug("birdc not ready, retrying", + "attempt", attempt, "err", err) + continue + } + // Any other failure (syntax error, etc.) is not retriable. + b.Logger.Warn("birdc reload failed", "err", err, "out", outStr) return } - b.Logger.Info("birdc configure ok", "out", string(out)) + b.Logger.Error("birdc configure gave up after retries", + "socket", socket, "attempts", birdcMaxAttempts) } // SummaryRoutes installs blackhole kernel routes for each NodeConfig CIDR. diff --git a/pkg/agent/handlers.go b/pkg/agent/handlers.go index 61ad351..b0727a4 100644 --- a/pkg/agent/handlers.go +++ b/pkg/agent/handlers.go @@ -330,3 +330,34 @@ func anycastStrings(ips []net.IP) []string { } return out } + +// orphanedCommitted returns committed allocations whose owner pod is no +// longer running on this node. An allocation is orphaned when: +// - lookupUID cannot find the pod by namespace+name (pod deleted or +// rescheduled to another node), OR +// - the pod is found but its UID doesn't match the allocation's OwnerUID +// (pod was replaced — name reuse after deletion). +// +// Allocations with an empty OwnerUID are only considered orphaned if the +// pod is absent; an empty UID prevents a false-positive on legacy entries +// that predate the UID field. +// +// lookupUID returns the current pod UID and found=true when the pod is on +// this node; found=false when absent. Callers plug in the live PodCache. +func orphanedCommitted(allocations []Allocation, lookupUID func(ns, name string) (uid string, found bool)) []Allocation { + var out []Allocation + for _, a := range allocations { + if a.State != StateCommitted { + continue + } + uid, found := lookupUID(a.Namespace, a.PodName) + if !found { + out = append(out, a) + continue + } + if a.OwnerUID != "" && uid != a.OwnerUID { + out = append(out, a) + } + } + return out +} diff --git a/pkg/agent/orphan_test.go b/pkg/agent/orphan_test.go new file mode 100644 index 0000000..4d7b24b --- /dev/null +++ b/pkg/agent/orphan_test.go @@ -0,0 +1,119 @@ +package agent + +import "testing" + +// lookupFixed returns a lookupUID func that maps pod names to UIDs. +// An empty UID in the map means "pod found but no UID" (legacy entry). +// A missing key means "pod not found on this node". +func lookupFixed(m map[string]string) func(ns, name string) (string, bool) { + return func(_, name string) (string, bool) { + uid, found := m[name] + return uid, found + } +} + +func TestOrphanedCommitted_PodAbsent(t *testing.T) { + allocs := []Allocation{{ + ContainerID: "c1", Namespace: "ns", PodName: "gitea-0", + OwnerUID: "uid-old", State: StateCommitted, IP6: "2001:db8::1", + }} + // Pod not on this node at all. + got := orphanedCommitted(allocs, lookupFixed(map[string]string{})) + if len(got) != 1 || got[0].ContainerID != "c1" { + t.Fatalf("expected 1 orphan, got %v", got) + } +} + +func TestOrphanedCommitted_UIDMismatch(t *testing.T) { + // Pod name re-used: informer has a newer pod with a different UID. + allocs := []Allocation{{ + ContainerID: "c1", Namespace: "ns", PodName: "gitea-0", + OwnerUID: "uid-old", State: StateCommitted, IP6: "2001:db8::1", + }} + got := orphanedCommitted(allocs, lookupFixed(map[string]string{ + "gitea-0": "uid-new", + })) + if len(got) != 1 || got[0].ContainerID != "c1" { + t.Fatalf("expected 1 orphan on UID mismatch, got %v", got) + } +} + +func TestOrphanedCommitted_LivePod(t *testing.T) { + // Matching UID — pod is live; must not be reported as orphan. + allocs := []Allocation{{ + ContainerID: "c1", Namespace: "ns", PodName: "gitea-0", + OwnerUID: "uid-live", State: StateCommitted, IP6: "2001:db8::1", + }} + got := orphanedCommitted(allocs, lookupFixed(map[string]string{ + "gitea-0": "uid-live", + })) + if len(got) != 0 { + t.Fatalf("live pod must not be orphaned, got %v", got) + } +} + +func TestOrphanedCommitted_PendingSkipped(t *testing.T) { + // Pending entries are excluded (handled by separate startup GC). + allocs := []Allocation{{ + ContainerID: "c1", Namespace: "ns", PodName: "pod-a", + OwnerUID: "uid-x", State: StatePending, IP6: "2001:db8::1", + }} + got := orphanedCommitted(allocs, lookupFixed(map[string]string{})) + if len(got) != 0 { + t.Fatalf("pending must be skipped, got %v", got) + } +} + +func TestOrphanedCommitted_EmptyOwnerUID_PodFound(t *testing.T) { + // Legacy allocation with empty OwnerUID: if the pod is found, do NOT + // treat it as orphaned (can't verify ownership without UID). + allocs := []Allocation{{ + ContainerID: "c1", Namespace: "ns", PodName: "gitea-0", + OwnerUID: "", State: StateCommitted, IP6: "2001:db8::1", + }} + got := orphanedCommitted(allocs, lookupFixed(map[string]string{ + "gitea-0": "uid-any", + })) + if len(got) != 0 { + t.Fatalf("empty OwnerUID with found pod must not be orphaned, got %v", got) + } +} + +func TestOrphanedCommitted_EmptyOwnerUID_PodAbsent(t *testing.T) { + // Legacy allocation with empty OwnerUID: if the pod is absent, it IS + // orphaned (pod is gone from this node regardless of UID). + allocs := []Allocation{{ + ContainerID: "c1", Namespace: "ns", PodName: "gitea-0", + OwnerUID: "", State: StateCommitted, IP6: "2001:db8::1", + }} + got := orphanedCommitted(allocs, lookupFixed(map[string]string{})) + if len(got) != 1 { + t.Fatalf("empty OwnerUID with absent pod must be orphaned, got %v", got) + } +} + +func TestOrphanedCommitted_Mixed(t *testing.T) { + // Three committed allocations: one live, one absent, one UID-mismatched. + allocs := []Allocation{ + {ContainerID: "live", Namespace: "ns", PodName: "pod-live", + OwnerUID: "uid-live", State: StateCommitted}, + {ContainerID: "absent", Namespace: "ns", PodName: "pod-absent", + OwnerUID: "uid-gone", State: StateCommitted}, + {ContainerID: "replaced", Namespace: "ns", PodName: "pod-replaced", + OwnerUID: "uid-old", State: StateCommitted}, + } + got := orphanedCommitted(allocs, lookupFixed(map[string]string{ + "pod-live": "uid-live", + "pod-replaced": "uid-new", + })) + if len(got) != 2 { + t.Fatalf("expected 2 orphans, got %d: %v", len(got), got) + } + seen := map[string]bool{} + for _, a := range got { + seen[a.ContainerID] = true + } + if !seen["absent"] || !seen["replaced"] { + t.Fatalf("wrong orphan set: %v", got) + } +} diff --git a/pkg/agent/runtime_linux.go b/pkg/agent/runtime_linux.go index 73ec4fd..af709ee 100644 --- a/pkg/agent/runtime_linux.go +++ b/pkg/agent/runtime_linux.go @@ -82,6 +82,34 @@ func (s *Server) configureRuntime(ctx context.Context) error { return fmt.Errorf("pod informer: %w", err) } + // Startup orphan GC: the pod informer is now fully synced. Walk all + // committed allocations and release any whose owner pod is absent from + // this node. This catches ungraceful evictions where CNI DEL never ran + // (TaintManagerEviction path) and prevents stale public /128s from + // suppressing the live pod's BGP advertisement after rescheduling. + gcOrphans := func(label string) int { + orphans := orphanedCommitted(s.Store.Snapshot(), func(ns, name string) (string, bool) { + pod, ok := pods.Get(ns, name) + if !ok { + return "", false + } + return string(pod.UID), true + }) + for _, a := range orphans { + s.Logger.Info(label, + "container_id", a.ContainerID, + "pod", a.Namespace+"/"+a.PodName, + "ip6", a.IP6, + "ip4", a.IP4, + ) + _ = Teardown(a.ContainerID, net.ParseIP(a.IP6), net.ParseIP(a.IP4)) + _ = s.Store.Delete(a.ContainerID) + ipam.Release(net.ParseIP(a.IP6), net.ParseIP(a.IP4)) + } + return len(orphans) + } + gcOrphans("GC orphaned committed allocation (startup)") + // Keep NetworkUnavailable=False so the node.kubernetes.io/network- // unavailable taint never gets re-applied. Calico's calico-node sets // it on shutdown; without an owner replacing it, kubelet's controller @@ -132,6 +160,25 @@ func (s *Server) configureRuntime(ctx context.Context) error { } }() + // Periodic orphan GC: defense-in-depth against allocations that escape + // the startup sweep (e.g. a pod evicted while the agent is running and + // the CNI DEL is never delivered). Keeps the store and IPAM in sync + // with the live pod set without requiring a full agent restart. + go func() { + t := time.NewTicker(60 * time.Second) + defer t.Stop() + for { + select { + case <-ctx.Done(): + return + case <-t.C: + if n := gcOrphans("GC orphaned committed allocation (periodic)"); n > 0 { + anycast.Trigger() + } + } + } + }() + // NetworkPolicy enforcement. world := netpol.NewWorld(s.Logger) if err := world.Start(ctx, s.restCfg); err != nil {