From 6afcb434612c34da49193fb2ee6a567903f01040 Mon Sep 17 00:00:00 2001 From: "Patrick W. Healy" Date: Thu, 17 Sep 2026 14:02:22 +0000 Subject: [PATCH] net controller: keep routine caches thin and preserve exact summary freshness Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: d2243398-6c36-4c3d-969e-7ed7bfb5b459 --- .../cluster_status.go | 9 +- cmd/unbounded-net-controller/detail_cache.go | 32 ++- cmd/unbounded-net-controller/detail_legacy.go | 151 +++++++++++++ .../detail_legacy_test.go | 204 +++++++++++++++++ .../detail_lifecycle.go | 4 + .../detail_requests.go | 7 +- .../detail_responses.go | 47 +--- .../detail_transport_test.go | 4 +- .../legacy_ingestion_test.go | 113 ++++++++++ .../node_detail_observer.go | 33 +++ .../node_detail_observer_test.go | 167 ++++++++++++++ cmd/unbounded-net-controller/node_legacy.go | 61 ++++++ .../node_legacy_test.go | 207 ++++++++++++++++++ cmd/unbounded-net-controller/node_overview.go | 19 +- .../node_overview_test.go | 7 +- cmd/unbounded-net-controller/node_status.go | 128 ++++++++--- .../overview_diagnostics_test.go | 148 +++++++++++++ cmd/unbounded-net-controller/server.go | 17 +- .../status_overview_ingestion_test.go | 199 +++++++++++++++++ cmd/unbounded-net-controller/status_proto.go | 11 +- cmd/unbounded-net-controller/status_types.go | 60 ++--- .../summary_metadata_test.go | 62 ++++++ .../summary_receipt_time_test.go | 102 +++++++++ 23 files changed, 1669 insertions(+), 123 deletions(-) create mode 100644 cmd/unbounded-net-controller/detail_legacy.go create mode 100644 cmd/unbounded-net-controller/detail_legacy_test.go create mode 100644 cmd/unbounded-net-controller/legacy_ingestion_test.go create mode 100644 cmd/unbounded-net-controller/node_detail_observer.go create mode 100644 cmd/unbounded-net-controller/node_detail_observer_test.go create mode 100644 cmd/unbounded-net-controller/node_legacy.go create mode 100644 cmd/unbounded-net-controller/node_legacy_test.go create mode 100644 cmd/unbounded-net-controller/overview_diagnostics_test.go create mode 100644 cmd/unbounded-net-controller/status_overview_ingestion_test.go create mode 100644 cmd/unbounded-net-controller/summary_metadata_test.go create mode 100644 cmd/unbounded-net-controller/summary_receipt_time_test.go diff --git a/cmd/unbounded-net-controller/cluster_status.go b/cmd/unbounded-net-controller/cluster_status.go index 56beaff1d..33891d570 100644 --- a/cmd/unbounded-net-controller/cluster_status.go +++ b/cmd/unbounded-net-controller/cluster_status.go @@ -18,6 +18,7 @@ import ( "k8s.io/klog/v2" "github.com/Azure/unbounded/internal/net/controller" + statuspkg "github.com/Azure/unbounded/internal/net/status" statusv1alpha1 "github.com/Azure/unbounded/internal/net/status/v1alpha1" "github.com/Azure/unbounded/internal/version" ) @@ -822,12 +823,8 @@ func collectClusterProblems(status *ClusterStatusResponse) []StatusProblem { } if overview := status.NodeOverviews[node.NodeInfo.Name]; overview != nil { - if overview.RouteMismatch { - appendProblem("node", nodeName, "Route next-hop mismatches (expected vs present)") - } - - if unhealthy := overview.PeerCount - overview.HealthyPeers; unhealthy > 0 { - appendProblem("node", nodeName, fmt.Sprintf("%d peers are not healthy", unhealthy)) + for _, message := range statuspkg.OverviewDiagnosticMessages(*overview, node.NodeInfo.ProviderID) { + appendProblem("node", nodeName, message) } continue diff --git a/cmd/unbounded-net-controller/detail_cache.go b/cmd/unbounded-net-controller/detail_cache.go index ccd64fe2c..7afc26bd4 100644 --- a/cmd/unbounded-net-controller/detail_cache.go +++ b/cmd/unbounded-net-controller/detail_cache.go @@ -17,7 +17,14 @@ import ( // nodeDetailSnapshot carries immutable details, separate from routine status. // Status and all its nested data must remain read-only, including for callers // retaining a returned snapshot after its cache entry expires. -type nodeDetailSnapshot = statusv1alpha1.NodeDetailSnapshot +type nodeDetailSnapshot struct { + statusv1alpha1.NodeDetailSnapshot + + legacyRevision uint64 + peerIdentity *peerIdentityDigest +} + +var errLegacyDetailBaseUnavailable = errors.New("legacy detail base changed or expired") // nodeDetailCache is a leader-local, TTL-only store. It owns no second result // history or per-entry timers. TTL bounds retention time, not peak memory. @@ -49,6 +56,10 @@ func newNodeDetailCache(ttl time.Duration) (*nodeDetailCache, error) { // maps, and pointers remain shared and must not be mutated by the caller. // Only Store renews the receipt-based TTL; request validation belongs upstream. func (c *nodeDetailCache) Store(nodeName, requestID string, collectedAt time.Time, status *NodeStatusResponse) (nodeDetailSnapshot, error) { + return c.store(nodeName, requestID, collectedAt, status, 0, nil, nil) +} + +func (c *nodeDetailCache) store(nodeName, requestID string, collectedAt time.Time, status *NodeStatusResponse, revision uint64, identity *peerIdentityDigest, expected *NodeStatusResponse) (nodeDetailSnapshot, error) { if nodeName == "" { return nodeDetailSnapshot{}, errors.New("node detail cache requires a node name") } @@ -63,13 +74,20 @@ func (c *nodeDetailCache) Store(nodeName, requestID string, collectedAt time.Tim defer c.mu.Unlock() now := c.clock.Now() + if expected != nil { + previous, ok := c.entries[nodeName] + if !ok || previous.Status != expected || !now.Before(previous.ExpiresAt) { + return nodeDetailSnapshot{}, errLegacyDetailBaseUnavailable + } + } + snapshot := nodeDetailSnapshot{ - NodeName: nodeName, - RequestID: requestID, - CollectedAt: collectedAt, - ReceivedAt: now, - ExpiresAt: now.Add(c.ttl), - Status: &statusCopy, + NodeDetailSnapshot: statusv1alpha1.NodeDetailSnapshot{ + NodeName: nodeName, RequestID: requestID, CollectedAt: collectedAt, + ReceivedAt: now, ExpiresAt: now.Add(c.ttl), Status: &statusCopy, + }, + legacyRevision: revision, + peerIdentity: identity, } c.entries[nodeName] = snapshot c.notify() diff --git a/cmd/unbounded-net-controller/detail_legacy.go b/cmd/unbounded-net-controller/detail_legacy.go new file mode 100644 index 000000000..3870ab818 --- /dev/null +++ b/cmd/unbounded-net-controller/detail_legacy.go @@ -0,0 +1,151 @@ +// Copyright (c) Microsoft Corporation. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "crypto/rand" + "errors" + "time" + + statusv1alpha1 "github.com/Azure/unbounded/internal/net/status/v1alpha1" +) + +// ObserveLegacy stores actual legacy details, not a summary or source update. +// A continuous stream reuses the completed cache association rather than making +// a request record per publication. It never fulfills an unrelated pending +// request. A non-nil base requires the exact, unexpired delta base at commit. +func (m *nodeDetailRequests) ObserveLegacy(nodeName string, status *NodeStatusResponse, revision uint64, identity *peerIdentityDigest, base *NodeStatusResponse) error { + m.mu.Lock() + defer m.mu.Unlock() + + m.expireLocked(time.Now()) + + if m.ctx.Err() != nil || m.closed { + return errors.New("detail request leader is unavailable") + } + + if status == nil || status.NodeInfo.Name != nodeName || revision == 0 { + return errors.New("legacy details require a matching node name and wire revision") + } + + uid, err := m.hooks.Resolve(nodeName) + if err != nil || uid == "" { + m.forgetLocked(nodeName) + + return errors.New("legacy detail node identity is unavailable") + } + + var request *nodeDetailRequest + + if snapshot, ok := m.cache.Get(nodeName); ok { + existing := m.requests[snapshot.RequestID] + if existing != nil && existing.uid == uid && existing.state == statusv1alpha1.NodeDetailComplete { + request = existing + } else if existing != nil && existing.uid != uid { + m.invalidateLocked(existing) + } + } + + if request == nil { + request = &nodeDetailRequest{ + nodeName: nodeName, uid: uid, state: statusv1alpha1.NodeDetailComplete, + command: statusv1alpha1.DetailRequest{RequestID: rand.Text()}, + cancel: func() {}, + } + } + + snapshot, err := m.cache.store(nodeName, request.command.RequestID, status.Timestamp, status, revision, identity, base) + if err != nil { + return err + } + + request.wakeAt = snapshot.ExpiresAt + m.requests[request.command.RequestID] = request + m.notify() + + return nil +} + +// LegacyBase returns only the current wire base. An on-demand snapshot cannot +// substitute for it, even if its payload happens to look identical. +func (m *nodeDetailRequests) LegacyBase(nodeName string, revision uint64) (*NodeStatusResponse, *peerIdentityDigest, bool) { + m.mu.Lock() + defer m.mu.Unlock() + + if m.ctx.Err() != nil || m.closed { + return nil, nil, false + } + + snapshot, ok := m.cache.Get(nodeName) + if !ok || snapshot.legacyRevision == 0 || snapshot.legacyRevision != revision { + return nil, nil, false + } + + request := m.requests[snapshot.RequestID] + + uid, err := m.hooks.Resolve(nodeName) + if request == nil || err != nil || request.uid != uid { + m.forgetLocked(nodeName) + + return nil, nil, false + } + + return snapshot.Status, snapshot.peerIdentity, true +} + +// Forget drops the current node's request state and cache ownership. +func (m *nodeDetailRequests) Forget(nodeName string) { + m.mu.Lock() + defer m.mu.Unlock() + + m.forgetLocked(nodeName) +} + +func (m *nodeDetailRequests) forgetLocked(nodeName string) { + for _, request := range m.requests { + if request.nodeName == nodeName { + m.invalidateLocked(request) + } + } + + m.cache.Delete(nodeName) +} + +// CompleteFailure terminates only a correlated live request. A failed refresh +// leaves an older, still-valid cached result alone. +func (m *nodeDetailRequests) CompleteFailure(nodeName, requestID, message string) error { + m.mu.Lock() + defer m.mu.Unlock() + + m.expireLocked(time.Now()) + + request := m.requests[requestID] + if m.ctx.Err() != nil || m.closed || request == nil || request.nodeName != nodeName || message == "" { + return errors.New("detail failure does not match an available request") + } + + if uid, err := m.hooks.Resolve(nodeName); err != nil || uid != request.uid { + m.invalidateLocked(request) + + return errors.New("detail failure node was deleted or replaced") + } + + if request.state == statusv1alpha1.NodeDetailComplete || request.state == statusv1alpha1.NodeDetailUnavailable { + return nil + } + + if request.state != statusv1alpha1.NodeDetailPending { + return errors.New("detail request is no longer pending") + } + + request.cancel() + request.state = statusv1alpha1.NodeDetailUnavailable + request.message = message + request.poll = false + request.wakeAt = time.Now().Add(m.timeout) + delete(m.active, nodeName) + m.notify() + + return nil +} diff --git a/cmd/unbounded-net-controller/detail_legacy_test.go b/cmd/unbounded-net-controller/detail_legacy_test.go new file mode 100644 index 000000000..aa56150bd --- /dev/null +++ b/cmd/unbounded-net-controller/detail_legacy_test.go @@ -0,0 +1,204 @@ +// Copyright (c) Microsoft Corporation. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "testing" + "testing/synctest" + "time" + + "k8s.io/apimachinery/pkg/types" + + statusv1alpha1 "github.com/Azure/unbounded/internal/net/status/v1alpha1" +) + +func TestObserveLegacyReusesAssociation(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + manager := testDetailRequests(t, nodeDetailRequestHooks{}) + + status := testDetailStatus() + if err := manager.ObserveLegacy("node", status, 1, nil, nil); err != nil { + t.Fatal(err) + } + + initial := manager.Request("node", false) + if initial.State != statusv1alpha1.NodeDetailComplete { + t.Fatal("legacy data was not reusable") + } + + time.Sleep(time.Second) + + for revision := uint64(2); revision <= 100; revision++ { + if err := manager.ObserveLegacy("node", status, revision, nil, nil); err != nil { + t.Fatal(err) + } + } + + updated := manager.Result("node", initial.RequestID) + if updated.State != statusv1alpha1.NodeDetailComplete || updated.Details == nil || + updated.Details.ReceivedAt != initial.Details.ReceivedAt.Add(time.Second) { + t.Fatal("legacy updates replaced the completed association or failed to refresh TTL") + } + + manager.mu.Lock() + count := len(manager.requests) + manager.mu.Unlock() + + if count != 1 { + t.Fatalf("publications accumulated %d request records", count) + } + + fresh := manager.Request("node", true) + if err := manager.ObserveLegacy("node", status, 101, nil, nil); err != nil { + t.Fatal(err) + } + + if result := manager.Result("node", fresh.RequestID); result.State != statusv1alpha1.NodeDetailPending { + t.Fatal("uncorrelated legacy publication fulfilled a fresh request") + } + + if err := manager.Complete("node", fresh.RequestID, status); err != nil { + t.Fatal(err) + } + + if err := manager.ObserveLegacy("node", status, 102, nil, nil); err != nil { + t.Fatal(err) + } + + if result := manager.Result("node", fresh.RequestID); result.State != statusv1alpha1.NodeDetailComplete || result.Details == nil { + t.Fatal("continuous legacy publication immediately expired a requested result") + } + }) +} + +func TestLegacyBaseExpiryReplacementAndIdentity(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + uid := types.UID("uid") + manager := testDetailRequests(t, nodeDetailRequestHooks{ + Resolve: func(string) (types.UID, error) { return uid, nil }, + }) + + memo := &peerIdentityDigest{1} + if err := manager.ObserveLegacy("node", testDetailStatus(), 5, memo, nil); err != nil { + t.Fatal(err) + } + + base, identity, ok := manager.LegacyBase("node", 5) + if !ok || identity != memo { + t.Fatal("wire base did not retain its validation memo") + } + + if _, _, ok := manager.LegacyBase("node", 4); ok { + t.Fatal("wrong wire revision accepted") + } + + request := manager.Request("node", true) + if err := manager.Complete("node", request.RequestID, testDetailStatus()); err != nil { + t.Fatal(err) + } + + if _, _, ok := manager.LegacyBase("node", 5); ok { + t.Fatal("one-shot details masqueraded as legacy wire base") + } + + if err := manager.ObserveLegacy("node", testDetailStatus(), 6, memo, base); err == nil { + t.Fatal("a delta overwrote a newer one-shot result") + } + + if err := manager.ObserveLegacy("node", testDetailStatus(), 7, memo, nil); err != nil { + t.Fatal(err) + } + + time.Sleep(manager.cache.ttl) + synctest.Wait() + assertNodeDetailEntries(t, manager.cache, 0) + + if _, _, ok := manager.LegacyBase("node", 7); ok { + t.Fatal("expired wire base accepted") + } + + if err := manager.ObserveLegacy("node", testDetailStatus(), 8, memo, nil); err != nil { + t.Fatal(err) + } + + synctest.Wait() + + uid = "replacement" + + if _, _, ok := manager.LegacyBase("node", 8); ok { + t.Fatal("replaced node reused old details") + } + + assertNodeDetailEntries(t, manager.cache, 0) + + uid = "" + + if err := manager.ObserveLegacy("node", testDetailStatus(), 9, nil, nil); err == nil { + t.Fatal("unknown node UID accepted") + } + + if err := manager.ObserveLegacy("node", nil, 9, nil, nil); err == nil { + t.Fatal("nil legacy details accepted") + } + + manager.Close() + + if err := manager.ObserveLegacy("node", testDetailStatus(), 9, nil, nil); err == nil { + t.Fatal("closed manager accepted legacy details") + } + }) +} + +func TestCompleteDetailFailurePreservesCachedResult(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + manager := testDetailRequests(t, nodeDetailRequestHooks{}) + if err := manager.ObserveLegacy("node", testDetailStatus(), 1, nil, nil); err != nil { + t.Fatal(err) + } + + cached := manager.Request("node", false) + refresh := manager.Request("node", true) + + for _, wrong := range []struct{ node, id, message string }{ + {"other", refresh.RequestID, "failed"}, + {"node", "unknown", "failed"}, + {"node", refresh.RequestID, ""}, + } { + if err := manager.CompleteFailure(wrong.node, wrong.id, wrong.message); err == nil { + t.Fatal("uncorrelated or empty failure accepted") + } + } + + for range 2 { + if err := manager.CompleteFailure("node", refresh.RequestID, "collection failed"); err != nil { + t.Fatal(err) + } + } + + result := manager.Result("node", refresh.RequestID) + if result.State != statusv1alpha1.NodeDetailUnavailable || result.Error != "collection failed" || result.Details != nil { + t.Fatal("failure was not explicit") + } + + previous := manager.Result("node", cached.RequestID) + if previous.Details == nil || *previous.Details != *cached.Details { + t.Fatal("failed refresh modified valid old details") + } + + if _, ok := manager.Pending("node"); ok { + t.Fatal("failed request retained a command") + } + + request := manager.Request("node", true) + time.Sleep(manager.timeout) + synctest.Wait() + + if err := manager.CompleteFailure("node", request.RequestID, "late"); err == nil { + t.Fatal("late failure revived an expired request") + } + + manager.Forget("node") + assertNodeDetailEntries(t, manager.cache, 0) + }) +} diff --git a/cmd/unbounded-net-controller/detail_lifecycle.go b/cmd/unbounded-net-controller/detail_lifecycle.go index ce2f49b9f..bba467143 100644 --- a/cmd/unbounded-net-controller/detail_lifecycle.go +++ b/cmd/unbounded-net-controller/detail_lifecycle.go @@ -107,6 +107,10 @@ func (h *healthState) startDetailRequests(ctx context.Context, nodeInformer cach h.detailRequests = manager h.detailMu.Unlock() + if h.statusCache != nil { + h.statusCache.BindDetails(manager) + } + go func() { <-manager.done diff --git a/cmd/unbounded-net-controller/detail_requests.go b/cmd/unbounded-net-controller/detail_requests.go index 0477b9a81..62863cbae 100644 --- a/cmd/unbounded-net-controller/detail_requests.go +++ b/cmd/unbounded-net-controller/detail_requests.go @@ -12,6 +12,7 @@ import ( "time" "k8s.io/apimachinery/pkg/types" + "k8s.io/klog/v2" statusv1alpha1 "github.com/Azure/unbounded/internal/net/status/v1alpha1" ) @@ -55,7 +56,7 @@ type nodeDetailRequests struct { } // newNodeDetailRequests takes exclusive lifecycle ownership of cache. All -// snapshots must enter through Complete so their node UID binding is known. +// snapshots enter through Complete or ObserveLegacy so their UID binding is known. func newNodeDetailRequests(ctx context.Context, cache *nodeDetailCache, timeout time.Duration, hooks nodeDetailRequestHooks) (*nodeDetailRequests, error) { if cache == nil || timeout <= 0 || hooks.Resolve == nil { return nil, errors.New("detail requests require a cache, positive timeout, and node resolver") @@ -263,7 +264,8 @@ func (m *nodeDetailRequests) resultLocked(request *nodeDetailRequest) statusv1al } if request.state == statusv1alpha1.NodeDetailComplete { if snapshot, ok := m.cache.Get(request.nodeName); ok && snapshot.RequestID == request.command.RequestID { - result.Details = &snapshot + details := snapshot.NodeDetailSnapshot + result.Details = &details } else { result.State = statusv1alpha1.NodeDetailExpired result.Error = "details expired or were replaced" @@ -358,6 +360,7 @@ func (m *nodeDetailRequests) run() { defer close(cacheDone) if err := m.cache.Run(m.ctx); err != nil { + klog.Errorf("Node detail cache expiry loop failed: %v", err) m.cancel() } }() diff --git a/cmd/unbounded-net-controller/detail_responses.go b/cmd/unbounded-net-controller/detail_responses.go index 9b4e3a8a7..aaf206fd8 100644 --- a/cmd/unbounded-net-controller/detail_responses.go +++ b/cmd/unbounded-net-controller/detail_responses.go @@ -3,51 +3,6 @@ package main -import ( - "errors" - "time" - - statusv1alpha1 "github.com/Azure/unbounded/internal/net/status/v1alpha1" -) - -// Fail records a correlated collection failure without deleting an older, -// still-valid snapshot that a viewer may display during a failed refresh. -func (m *nodeDetailRequests) Fail(nodeName, requestID, reason string) error { - m.mu.Lock() - defer m.mu.Unlock() - - m.expireLocked(time.Now()) - - request := m.requests[requestID] - if reason == "" || m.ctx.Err() != nil || m.closed || request == nil || request.nodeName != nodeName { - return errors.New("detail failure does not match an available request") - } - - if uid, err := m.hooks.Resolve(nodeName); err != nil || uid != request.uid { - m.invalidateLocked(request) - return errors.New("detail request node was deleted or replaced") - } - - if request.state == statusv1alpha1.NodeDetailComplete || - (request.state == statusv1alpha1.NodeDetailUnavailable && request.message == reason) { - return nil - } - - if request.state != statusv1alpha1.NodeDetailPending { - return errors.New("detail request is no longer pending") - } - - request.cancel() - request.state = statusv1alpha1.NodeDetailUnavailable - request.message = reason - request.poll = false - request.wakeAt = time.Now().Add(m.timeout) - delete(m.active, nodeName) - m.notify() - - return nil -} - func handleNodeDetailResponse(health *healthState, nodeName, requestID string, status *NodeStatusResponse, failure string) NodeStatusPushAck { ack := NodeStatusPushAck{Status: "error", DetailRequestID: requestID, SummarySupported: true} @@ -64,7 +19,7 @@ func handleNodeDetailResponse(health *healthState, nodeName, requestID string, s var err error if failure != "" { - err = manager.Fail(nodeName, requestID, failure) + err = manager.CompleteFailure(nodeName, requestID, failure) } else { err = manager.Complete(nodeName, requestID, status) } diff --git a/cmd/unbounded-net-controller/detail_transport_test.go b/cmd/unbounded-net-controller/detail_transport_test.go index 205dddbf6..61ae87132 100644 --- a/cmd/unbounded-net-controller/detail_transport_test.go +++ b/cmd/unbounded-net-controller/detail_transport_test.go @@ -321,11 +321,11 @@ func TestDetailFailureKeepsPreviousValidSnapshot(t *testing.T) { before, _ := manager.cache.Get("node") refresh := manager.Request("node", true) - if err := manager.Fail("other", refresh.RequestID, "bad"); err == nil { + if err := manager.CompleteFailure("other", refresh.RequestID, "bad"); err == nil { t.Fatal("failure for a different node was accepted") } - if err := manager.Fail("node", refresh.RequestID, "response exceeds the transport frame limit"); err != nil { + if err := manager.CompleteFailure("node", refresh.RequestID, "response exceeds the transport frame limit"); err != nil { t.Fatal(err) } diff --git a/cmd/unbounded-net-controller/legacy_ingestion_test.go b/cmd/unbounded-net-controller/legacy_ingestion_test.go new file mode 100644 index 000000000..3ea80e336 --- /dev/null +++ b/cmd/unbounded-net-controller/legacy_ingestion_test.go @@ -0,0 +1,113 @@ +// Copyright (c) Microsoft Corporation. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "errors" + "net/http" + "testing" + + "google.golang.org/protobuf/proto" + "k8s.io/apimachinery/pkg/types" + + statusproto "github.com/Azure/unbounded/internal/net/status/proto" +) + +func TestLegacyIngestionPropagatesStorageFailure(t *testing.T) { + message := &statusproto.NodeStatusMessage{ + Type: "node_status_full", NodeName: "node", + Status: &statusproto.NodeStatusFull{ + NodeInfo: &statusproto.NodeInfo{Name: "node"}, + Peers: []*statusproto.PeerStatus{{Name: "peer"}}, + }, + } + + binary, err := proto.Marshal(message) + if err != nil { + t.Fatal(err) + } + + decoded, err := decodeProtoWSMessage(binary) + if err != nil { + t.Fatal(err) + } + + for _, failure := range []string{"none", "identity", "closed"} { + t.Run(failure, func(t *testing.T) { + for _, transport := range []string{"raw-http", "json-http", "protobuf-http", "json-ws", "protobuf-ws"} { + t.Run(transport, func(t *testing.T) { + manager := testDetailRequests(t, nodeDetailRequestHooks{ + Resolve: func(string) (types.UID, error) { + if failure == "identity" { + return "", errors.New("node identity unavailable") + } + + return "uid", nil + }, + }) + cache := NewNodeStatusCache() + cache.BindDetails(manager) + health := &healthState{statusCache: cache} + + if failure == "closed" { + manager.Close() + } + + var ( + ack NodeStatusPushAck + code int + ackType string + pushErr error + ) + + switch transport { + case "raw-http": + ack, code, pushErr = handleStatusPushRequest(health, []byte(`{"nodeInfo":{"name":"node"},"peers":[{"name":"peer"}]}`)) + case "json-http": + ack, code, pushErr = handleStatusPushRequest(health, []byte(`{"mode":"full","nodeName":"node","status":{"nodeInfo":{"name":"node"},"peers":[{"name":"peer"}]}}`)) + case "protobuf-http": + ack, code, pushErr = handleProtoPushRequest(health, binary, "push") + case "json-ws": + ackType, ack = handleNodeStatusWSMessage(health, []byte(`{"type":"node_status_full","nodeName":"node","status":{"nodeInfo":{"name":"node"},"peers":[{"name":"peer"}]}}`)) + case "protobuf-ws": + ackType, ack = handleProtoWSMessage(health, decoded, "ws") + } + + if failure != "none" { + if code != 0 && (code != http.StatusServiceUnavailable || pushErr == nil) { + t.Fatalf("HTTP storage failure was hidden: code=%d ack=%+v err=%v", code, ack, pushErr) + } + + if code == 0 && (ackType != "node_status_resync" || ack.Status != "resync_required" || ack.Reason == "") { + t.Fatalf("WebSocket storage failure was hidden: type=%s ack=%+v", ackType, ack) + } + + if ack.Revision != 0 || cache.Len() != 0 { + t.Fatal("failed storage advanced publication state") + } + + assertNodeDetailEntries(t, manager.cache, 0) + + return + } + + if pushErr != nil || ack.Status != "ok" || ack.Revision != 1 || + (code != 0 && code != http.StatusOK) || (code == 0 && ackType != "node_status_ack") { + t.Fatalf("successful storage was rejected: code=%d type=%s ack=%+v err=%v", code, ackType, ack, pushErr) + } + + entry, ok := cache.Get("node") + if !ok || entry.Overview == nil || entry.Overview.PeerCount != 1 || len(entry.Status.Peers) != 0 { + t.Fatal("routine cache retained full peers or lost overview counts") + } + + snapshot, ok := manager.cache.Get("node") + if !ok || len(snapshot.Status.Peers) != 1 { + t.Fatal("accepted publication lost expiring details") + } + }) + } + }) + } +} diff --git a/cmd/unbounded-net-controller/node_detail_observer.go b/cmd/unbounded-net-controller/node_detail_observer.go new file mode 100644 index 000000000..29fbee53f --- /dev/null +++ b/cmd/unbounded-net-controller/node_detail_observer.go @@ -0,0 +1,33 @@ +// Copyright (c) Microsoft Corporation. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import "k8s.io/klog/v2" + +// ObserveLegacyDetails bridges accepted legacy publications into the detail API +// while preserving existing full-cache and callback behavior during migration. +// It does not replay existing entries or refresh TTL on reads/source changes. +func (c *NodeStatusCache) ObserveLegacyDetails(manager *nodeDetailRequests) { + c.mu.Lock() + defer c.mu.Unlock() + + c.legacyObserver = manager +} + +// The cache lock preserves publication order against concurrent full/delta +// updates. The observer resolves informer identity but performs no network I/O. +func (c *NodeStatusCache) observeLegacyLocked(nodeName string, entry *CachedNodeStatus) { + if c.legacyObserver == nil { + return + } + + status := *entry.Status + if status.NodeInfo.Name == "" { + status.NodeInfo.Name = nodeName + } + + if err := c.legacyObserver.ObserveLegacy(nodeName, &status, entry.Revision, entry.peerIdentity, nil); err != nil { + klog.Errorf("Observing legacy node %q details failed: %v", nodeName, err) + } +} diff --git a/cmd/unbounded-net-controller/node_detail_observer_test.go b/cmd/unbounded-net-controller/node_detail_observer_test.go new file mode 100644 index 000000000..3230b2e29 --- /dev/null +++ b/cmd/unbounded-net-controller/node_detail_observer_test.go @@ -0,0 +1,167 @@ +// Copyright (c) Microsoft Corporation. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "context" + "errors" + "net/http" + "net/http/httptest" + "strings" + "sync/atomic" + "testing" + "testing/synctest" + "time" + + statusv1alpha1 "github.com/Azure/unbounded/internal/net/status/v1alpha1" +) + +func TestLegacyObserverBridgeAuthenticatedPublishAndCachedDetails(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + var pulls atomic.Int32 + + manager := testDetailRequests(t, nodeDetailRequestHooks{ + Pull: func(context.Context, string) (*NodeStatusResponse, error) { + pulls.Add(1) + + return nil, errors.New("node HTTP endpoint is unreachable") + }, + }) + health := &healthState{ + statusCache: NewNodeStatusCache(), detailRequests: manager, + nodeTokenVerifier: fakeServiceAccountTokenVerifier{}, + } + health.isLeader.Store(true) + health.statusCache.ObserveLegacyDetails(manager) + + fullCallbacks := 0 + + health.statusCache.SetOnChange(func(_ string, status *NodeStatusResponse) { + fullCallbacks++ + + if len(status.Peers) == 0 { + t.Error("bridge replaced the existing full callback payload") + } + }) + + issuer := testTokenIssuer(t) + token := testNodeToken(t, issuer) + mux := http.NewServeMux() + registerPushHandlers(mux, health, nil, make(chan struct{}, maxConcurrentNodeWS), issuer) + registerStatusHandlers(mux, health, false, nil, nil, nil) + + fullBody := `{"mode":"full","nodeName":"node-a","status":{"nodeInfo":{"name":"node-a"},"peers":[{"name":"peer-1"}]}}` + + publish := func(body, bearer string) int { + request := httptest.NewRequest(http.MethodPost, "/status/push", strings.NewReader(body)) + request.Header.Set("Authorization", "Bearer "+bearer) + + response := httptest.NewRecorder() + mux.ServeHTTP(response, request) + + return response.Code + } + if code := publish(fullBody, "invalid"); code != http.StatusUnauthorized { + t.Fatalf("unauthenticated publication returned %d", code) + } + + assertNodeDetailEntries(t, manager.cache, 0) + + if code := publish(fullBody, token); code != http.StatusOK { + t.Fatalf("full publication returned %d", code) + } + + response, first := serveDetailRequest(t, mux, http.MethodPost, "/status/node/node-a/details", "{}") + if response.Code != http.StatusOK || first.State != statusv1alpha1.NodeDetailComplete || + first.Details == nil || first.Details.Status.Peers[0].Name != "peer-1" || pulls.Load() != 0 { + t.Fatal("cached authenticated details required a direct HTTP pull") + } + + time.Sleep(time.Second) + + deltaBody := `{"mode":"delta","nodeName":"node-a","baseRevision":1,"delta":{"peers":[{"name":"peer-2"},{"name":"peer-3"}]}}` + if code := publish(deltaBody, token); code != http.StatusOK { + t.Fatalf("delta publication returned %d", code) + } + + current := manager.Request("node-a", false) + if current.RequestID != first.RequestID || current.Details == nil || len(current.Details.Status.Peers) != 2 || + !current.Details.ExpiresAt.Equal(first.Details.ExpiresAt.Add(time.Second)) || pulls.Load() != 0 { + t.Fatal("delta did not refresh the existing cache association") + } + + legacy, ok := health.statusCache.Get("node-a") + if !ok || legacy.Overview != nil || len(legacy.Status.Peers) != 2 || fullCallbacks != 2 { + t.Fatal("bridge changed existing full-cache or callback behavior") + } + + health.statusCache.UpdateSource("node-a", "ws") + + if _, err := health.statusCache.StoreOverview("summary-node", statusv1alpha1.NodeStatusOverview{ + NodeInfo: NodeInfo{Name: "summary-node"}, + }, "push"); err != nil { + t.Fatal(err) + } + + after, _ := manager.cache.Get("node-a") + if after.ExpiresAt != current.Details.ExpiresAt || fullCallbacks != 3 { + t.Fatal("source/summary change renewed the detail TTL or changed callbacks") + } + + fresh := manager.Request("node-a", true) + + if code := publish(fullBody, token); code != http.StatusOK { + t.Fatalf("next full publication returned %d", code) + } + + pending := manager.Result("node-a", fresh.RequestID) + if pending.State != statusv1alpha1.NodeDetailPending || pending.RequestID != fresh.RequestID || pending.Deadline != fresh.Deadline { + t.Fatal("legacy publication reset or completed a pending correlated refresh") + } + + time.Sleep(manager.cache.ttl) + synctest.Wait() + assertNodeDetailEntries(t, manager.cache, 0) + + legacy, ok = health.statusCache.Get("node-a") + + if !ok || len(legacy.Status.Peers) != 1 || legacy.Overview != nil { + t.Fatal("bridge prematurely removed the legacy full base") + } + }) +} + +func TestLegacyObserverBridgeFailurePreservesOldStore(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + manager := testDetailRequests(t, nodeDetailRequestHooks{}) + cache := NewNodeStatusCache() + cache.ObserveLegacyDetails(manager) + + status := NodeStatusResponse{Peers: []WireGuardPeerStatus{{Name: "peer"}}} + cache.StoreFull("node", status, "ws") + + details := manager.Request("node", false) + + if details.Details == nil || details.Details.Status.NodeInfo.Name != "node" || + status.NodeInfo.Name != "" || cache.entries["node"].Status.NodeInfo.Name != "" { + t.Fatal("bridge failed to normalize only its own shallow snapshot") + } + + manager.Close() + + if revision := cache.StoreFull("node", status, "ws"); revision != 2 { + t.Fatal("observer failure changed old cache revision behavior") + } + + if len(cache.entries["node"].Status.Peers) != 1 { + t.Fatal("observer failure discarded the old full payload") + } + + cache.ObserveLegacyDetails(nil) + + if revision := cache.StoreFull("node", status, "ws"); revision != 3 { + t.Fatal("detaching the bridge changed legacy storage") + } + }) +} diff --git a/cmd/unbounded-net-controller/node_legacy.go b/cmd/unbounded-net-controller/node_legacy.go new file mode 100644 index 000000000..349bf628b --- /dev/null +++ b/cmd/unbounded-net-controller/node_legacy.go @@ -0,0 +1,61 @@ +// Copyright (c) Microsoft Corporation. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "time" + + statuspkg "github.com/Azure/unbounded/internal/net/status" +) + +// BindDetails switches routine storage to overview-only for this process. +// Existing full entries lose their heavy ownership; subsequent legacy deltas +// require an unexpired TTL base. A closed manager stays bound and fails closed. +func (c *NodeStatusCache) BindDetails(manager *nodeDetailRequests) { + if manager == nil { + panic("node status requires a non-nil detail manager") + } + + c.mu.Lock() + defer c.mu.Unlock() + + c.details = manager + c.legacyObserver = nil + + for name, previous := range c.entries { + entry := *previous + if entry.Overview == nil { + overview := statuspkg.OverviewFromStatus(entry.Status, time.Now()) + entry.Overview = &overview + } + + metadata := statuspkg.OverviewMetadata(*entry.Overview) + entry.Status = &metadata + entry.peerIdentity = nil + c.entries[name] = &entry + } +} + +func (c *NodeStatusCache) legacyEntryLocked(nodeName string, status *NodeStatusResponse, revision uint64, identity *peerIdentityDigest, source string, base *NodeStatusResponse) (*CachedNodeStatus, error) { + entry := &CachedNodeStatus{ + Status: status, Revision: revision, Source: source, + ReceivedAt: time.Now(), peerIdentity: identity, legacy: true, + } + if c.details == nil { + return entry, nil + } + + if err := c.details.ObserveLegacy(nodeName, status, revision, identity, base); err != nil { + return nil, err + } + + overview := statuspkg.OverviewFromStatus(status, entry.ReceivedAt) + overview.StatusSource = source + metadata := statuspkg.OverviewMetadata(overview) + entry.Status = &metadata + entry.Overview = &overview + entry.peerIdentity = nil + + return entry, nil +} diff --git a/cmd/unbounded-net-controller/node_legacy_test.go b/cmd/unbounded-net-controller/node_legacy_test.go new file mode 100644 index 000000000..93af5dd15 --- /dev/null +++ b/cmd/unbounded-net-controller/node_legacy_test.go @@ -0,0 +1,207 @@ +// Copyright (c) Microsoft Corporation. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "encoding/json" + "testing" + "testing/synctest" + "time" + + statuspkg "github.com/Azure/unbounded/internal/net/status" + statusv1alpha1 "github.com/Azure/unbounded/internal/net/status/v1alpha1" +) + +func retentionFixture(peers int) NodeStatusResponse { + status := protoToNodeStatus(measurementTestStatus(peers)) + status.NodeInfo.Name = "node" + status.RoutingTable.Routes = []RouteEntry{{Destination: "10.0.0.0/8", NextHops: []NextHop{{Device: "wg0"}}}} + status.BpfEntries = []BpfEntry{{CIDR: "10.0.0.0/8", Node: "private-detail-marker"}} + + return status +} + +func assertThinStatus(t *testing.T, entry *CachedNodeStatus, peerCount int) { + t.Helper() + + if entry == nil || entry.Overview == nil || entry.Overview.PeerCount != peerCount || + len(entry.Status.Peers) != 0 || len(entry.Status.RoutingTable.Routes) != 0 || + len(entry.Status.BpfEntries) != 0 || entry.peerIdentity != nil { + t.Fatal("routine cache retained details/memo or lost observed facts") + } +} + +func TestBoundNodeCacheRetainsOverviewOnly(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + manager := testDetailRequests(t, nodeDetailRequestHooks{}) + cache := NewNodeStatusCache() + status := retentionFixture(1024) + cache.StoreFull("node", status, "ws") + cache.BindDetails(manager) + assertThinStatus(t, cache.entries["node"], 1024) + + if _, conflict, err := cache.ApplyDelta("node", 1, map[string]json.RawMessage{"statusSource": []byte(`"ws"`)}, "ws"); err != nil || !conflict { + t.Fatal("pre-binding payload remained usable as a wire base") + } + + fullCallbacks := 0 + overviewCallbacks := 0 + + cache.SetOnChange(func(string, *NodeStatusResponse) { fullCallbacks++ }) + cache.SetOnOverviewChange(func(_ string, overview statusv1alpha1.NodeStatusOverview) { + overviewCallbacks++ + + if overview.PeerCount != 1024 || overview.RouteCount != 1 { + t.Error("callback lost observed counts") + } + }) + + revision, err := cache.StoreFullChecked("node", status, "ws") + if err != nil { + t.Fatal(err) + } + + assertThinStatus(t, cache.entries["node"], 1024) + + base, _, ok := manager.LegacyBase("node", revision) + if !ok || &base.Peers[0] != &status.Peers[0] { + t.Fatal("TTL wire base missing or deeply copied") + } + + time.Sleep(time.Second) + + measurements, err := statuspkg.PeerMeasurementsToProto(status.Peers) + if err != nil { + t.Fatal(err) + } + + revision, conflict, err := cache.ApplyParsedDelta("node", revision, parsedDelta{peerMeasurements: measurements}, "ws") + if err != nil || conflict { + t.Fatalf("measurement delta: conflict=%v error=%v", conflict, err) + } + + assertThinStatus(t, cache.entries["node"], 1024) + + if _, memo, ok := manager.LegacyBase("node", revision); !ok || memo == nil { + t.Fatal("measurement identity memo was not kept with details") + } + + snapshot, _ := manager.cache.Get("node") + + time.Sleep(time.Second) + cache.Get("node") + cache.GetAll() + cache.UpdateSource("node", "push") + after, _ := manager.cache.Get("node") + + if after.ExpiresAt != snapshot.ExpiresAt || fullCallbacks != 0 || overviewCallbacks != 3 { + t.Fatal("routine read/source update refreshed TTL or sent full data") + } + + time.Sleep(9 * time.Second) + synctest.Wait() + assertNodeDetailEntries(t, manager.cache, 0) + assertThinStatus(t, cache.entries["node"], 1024) + + if _, conflict, err := cache.ApplyParsedDelta("node", revision, parsedDelta{peerMeasurements: measurements}, "ws"); err != nil || !conflict { + t.Fatal("expired legacy base did not require resync") + } + }) +} + +func TestBoundNodeCacheSummaryDeletionAndClose(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + manager := testDetailRequests(t, nodeDetailRequestHooks{}) + cache := NewNodeStatusCache() + cache.BindDetails(manager) + + status := retentionFixture(3) + cache.StoreFull("node", status, "ws") + before, _ := manager.cache.Get("node") + + time.Sleep(time.Second) + + if _, err := cache.StoreOverview("node", statuspkg.OverviewFromStatus(&status, time.Now()), "push"); err != nil { + t.Fatal(err) + } + + after, _ := manager.cache.Get("node") + if after.ExpiresAt != before.ExpiresAt { + t.Fatal("summary publication refreshed details") + } + + cache.CleanupStaleEntries(map[string]bool{}) + assertNodeDetailEntries(t, manager.cache, 0) + cache.StoreFull("node", status, "ws") + cache.Delete("node") + assertNodeDetailEntries(t, manager.cache, 0) + manager.Close() + + if _, err := cache.StoreFullChecked("node", status, "ws"); err == nil { + t.Fatal("closed lifecycle silently accepted full data") + } + + if cache.Len() != 0 { + t.Fatal("closed lifecycle repopulated routine cache") + } + }) +} + +func TestBoundNodeCacheRejectsReplacedDeltaBase(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + manager := testDetailRequests(t, nodeDetailRequestHooks{}) + cache := NewNodeStatusCache() + cache.BindDetails(manager) + + status := retentionFixture(3) + revision := cache.StoreFull("node", status, "ws") + previous := cache.entries["node"] + base, _, _ := manager.LegacyBase("node", revision) + request := manager.Request("node", true) + + if err := manager.Complete("node", request.RequestID, &status); err != nil { + t.Fatal(err) + } + + if _, conflict, err := cache.commitParsedDeltaBase("node", previous, &status, nil, "ws", base); err != nil || !conflict { + t.Fatal("stale delta replaced a newer one-shot result") + } + + if result := manager.Result("node", request.RequestID); result.State != statusv1alpha1.NodeDetailComplete { + t.Fatal("delta invalidated requested result") + } + }) +} + +func TestBoundNodeCacheDisablesLegacyBridgeObserver(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + manager := testDetailRequests(t, nodeDetailRequestHooks{}) + cache := NewNodeStatusCache() + cache.ObserveLegacyDetails(manager) + + status := retentionFixture(3) + revision := cache.StoreFull("node", status, "ws") + cache.BindDetails(manager) + + if cache.legacyObserver != nil { + t.Fatal("thin binding left the compatibility observer active") + } + + measurements, err := statuspkg.PeerMeasurementsToProto(status.Peers) + if err != nil { + t.Fatal(err) + } + + if _, conflict, err := cache.ApplyParsedDelta("node", revision, parsedDelta{peerMeasurements: measurements}, "ws"); err != nil || conflict { + t.Fatalf("bridge transition lost its valid TTL base: conflict=%v error=%v", conflict, err) + } + + snapshot, ok := manager.cache.Get("node") + if !ok || len(snapshot.Status.Peers) != 3 { + t.Fatal("compatibility observer replaced detailed data with thin metadata") + } + + assertThinStatus(t, cache.entries["node"], 3) + }) +} diff --git a/cmd/unbounded-net-controller/node_overview.go b/cmd/unbounded-net-controller/node_overview.go index fe5f79225..0b9a14540 100644 --- a/cmd/unbounded-net-controller/node_overview.go +++ b/cmd/unbounded-net-controller/node_overview.go @@ -18,7 +18,9 @@ func (c *NodeStatusCache) StoreOverview(nodeName string, overview statusv1alpha1 } if overview.PeerCount < 0 || overview.HealthyPeers < 0 || - overview.HealthyPeers > overview.PeerCount || overview.RouteCount < 0 { + overview.HealthyPeers > overview.PeerCount || overview.RouteCount < 0 || + overview.RouteMismatchCount < 0 || overview.UnhealthyPeerLinks < 0 || + (overview.RouteMismatchCount > 0 && !overview.RouteMismatch) { return 0, fmt.Errorf("summary contains invalid observed counts") } @@ -38,20 +40,31 @@ func (c *NodeStatusCache) StoreOverview(nodeName string, overview statusv1alpha1 revision = previous.Revision + 1 } - c.entries[nodeName] = &CachedNodeStatus{ + entry := &CachedNodeStatus{ Status: &metadata, Overview: &overview, Source: source, Revision: revision, ReceivedAt: time.Now(), } + c.entries[nodeName] = entry fn := c.onOverviewChange c.mu.Unlock() if fn != nil { - fn(nodeName, overview) + fn(nodeName, entry.overviewForNotification()) } return revision, nil } +// Notifications use controller receipt time, matching full cluster rebuilds, +// without changing the node's wire metadata or renewing it on source changes. +func (entry *CachedNodeStatus) overviewForNotification() statusv1alpha1.NodeStatusOverview { + overview := *entry.Overview + received := entry.ReceivedAt + overview.LastPushTime = &received + + return overview +} + // SetOnOverviewChange registers the summary-only cache mutation callback. func (c *NodeStatusCache) SetOnOverviewChange(fn func(string, statusv1alpha1.NodeStatusOverview)) { c.mu.Lock() diff --git a/cmd/unbounded-net-controller/node_overview_test.go b/cmd/unbounded-net-controller/node_overview_test.go index 7b5734891..065f81ac7 100644 --- a/cmd/unbounded-net-controller/node_overview_test.go +++ b/cmd/unbounded-net-controller/node_overview_test.go @@ -6,6 +6,7 @@ package main import ( "bytes" "encoding/json" + "reflect" "strings" "sync" "testing" @@ -86,6 +87,9 @@ func TestNodeOverviewCacheRejectsInvalidFacts(t *testing.T) { {HealthyPeers: -1}, {PeerCount: 1, HealthyPeers: 2}, {RouteCount: -1}, + {RouteMismatchCount: -1}, + {UnhealthyPeerLinks: -1}, + {RouteMismatchCount: 1}, } { cache := NewNodeStatusCache() if _, err := cache.StoreOverview("node", overview, "ws"); err == nil || cache.Len() != 0 { @@ -107,6 +111,7 @@ func TestClusterOverviewPreservesCountsAndEnrichment(t *testing.T) { overview := statusv1alpha1.NodeStatusOverview{ NodeInfo: NodeInfo{Name: "node", SiteName: "site", WireGuard: &WireGuardStatusInfo{Interface: "wg0"}}, StatusSource: "ws", PeerCount: 20, HealthyPeers: 17, RouteCount: 30, RouteMismatch: true, + RouteMismatchCount: 2, UnhealthyPeerLinks: 3, } c.PatchOverview("node", overview) snapshot := c.Get() @@ -130,7 +135,7 @@ func TestClusterOverviewPreservesCountsAndEnrichment(t *testing.T) { overview.NodeErrors = []NodeError{{Type: "cni", Message: "blocked"}} c.PatchOverview("node", overview) - if buildClusterSummary(snapshot).NodeSummaries[0] != row { + if !reflect.DeepEqual(buildClusterSummary(snapshot).NodeSummaries[0], row) { t.Fatal("patching changed a previously returned snapshot") } diff --git a/cmd/unbounded-net-controller/node_status.go b/cmd/unbounded-net-controller/node_status.go index 3dbb37b88..d23b4e373 100644 --- a/cmd/unbounded-net-controller/node_status.go +++ b/cmd/unbounded-net-controller/node_status.go @@ -6,11 +6,15 @@ package main import ( "context" "encoding/json" + "errors" "fmt" "net/http" "sync" "time" + "k8s.io/klog/v2" + + statuspkg "github.com/Azure/unbounded/internal/net/status" statusproto "github.com/Azure/unbounded/internal/net/status/proto" statusv1alpha1 "github.com/Azure/unbounded/internal/net/status/v1alpha1" ) @@ -24,6 +28,7 @@ type CachedNodeStatus struct { Overview *statusv1alpha1.NodeStatusOverview peerIdentity *peerIdentityDigest + legacy bool } // NodeStatusCache is a thread-safe cache of node status data pushed from node agents. @@ -32,6 +37,8 @@ type NodeStatusCache struct { entries map[string]*CachedNodeStatus onChange func(nodeName string, status *NodeStatusResponse) onOverviewChange func(nodeName string, overview statusv1alpha1.NodeStatusOverview) + legacyObserver *nodeDetailRequests + details *nodeDetailRequests } // NewNodeStatusCache creates an empty NodeStatusCache. @@ -49,33 +56,60 @@ func (c *NodeStatusCache) Len() int { // StoreFull stores a full node status payload and returns the new revision. func (c *NodeStatusCache) StoreFull(nodeName string, status NodeStatusResponse, source string) uint64 { + revision, err := c.StoreFullChecked(nodeName, status, source) + if err != nil { + klog.Errorf("Store node %q status failed: %v", nodeName, err) + } + + return revision +} + +// StoreFullChecked exposes identity/lifecycle failures to ingestion handlers. +func (c *NodeStatusCache) StoreFullChecked(nodeName string, status NodeStatusResponse, source string) (uint64, error) { if source == "" { source = "push" } + if nodeName == "" || (status.NodeInfo.Name != "" && status.NodeInfo.Name != nodeName) { + return 0, fmt.Errorf("legacy status identity does not match node %q", nodeName) + } + c.mu.Lock() + if c.details != nil && status.NodeInfo.Name == "" { + status.NodeInfo.Name = nodeName + } + prevRevision := uint64(0) if existing, ok := c.entries[nodeName]; ok { prevRevision = existing.Revision } revision := prevRevision + 1 - c.entries[nodeName] = &CachedNodeStatus{ - Status: &status, - ReceivedAt: time.Now(), - Source: source, - Revision: revision, + + entry, err := c.legacyEntryLocked(nodeName, &status, revision, nil, source, nil) + if err != nil { + c.mu.Unlock() + + return 0, err } + + c.entries[nodeName] = entry + c.observeLegacyLocked(nodeName, entry) + fn := c.onChange - statusPtr := c.entries[nodeName].Status + overviewFn := c.onOverviewChange c.mu.Unlock() - if fn != nil { - fn(nodeName, statusPtr) + if entry.Overview != nil { + if overviewFn != nil { + overviewFn(nodeName, entry.overviewForNotification()) + } + } else if fn != nil { + fn(nodeName, entry.Status) } - return revision + return revision, nil } // parsedDelta holds pre-deserialized delta fields parsed outside the lock. @@ -233,7 +267,7 @@ func (c *NodeStatusCache) applyParsedDelta(nodeName string, baseRevision uint64, return 0, true, nil } - if entry.Overview != nil || (pd.peerMeasurements != nil && baseRevision == 0) || (baseRevision != 0 && entry.Revision != baseRevision) { + if (entry.Overview != nil && !entry.legacy) || (pd.peerMeasurements != nil && baseRevision == 0) || (baseRevision != 0 && entry.Revision != baseRevision) { rev := entry.Revision c.mu.RUnlock() @@ -245,6 +279,17 @@ func (c *NodeStatusCache) applyParsedDelta(nodeName string, baseRevision uint64, prevRevision := entry.Revision peerIdentity := entry.peerIdentity + if c.details != nil { + var exists bool + + prevStatus, peerIdentity, exists = c.details.LegacyBase(nodeName, prevRevision) + if !exists { + c.mu.RUnlock() + + return prevRevision, true, nil + } + } + c.mu.RUnlock() // Phase 2: Copy and merge OUTSIDE the lock. This is the expensive @@ -337,10 +382,14 @@ func (c *NodeStatusCache) applyParsedDelta(nodeName string, baseRevision uint64, merged.NodeInfo.Name = nodeName } - return c.commitParsedDelta(nodeName, entry, &merged, peerIdentity, source) + return c.commitParsedDeltaBase(nodeName, entry, &merged, peerIdentity, source, prevStatus) } func (c *NodeStatusCache) commitParsedDelta(nodeName string, previous *CachedNodeStatus, merged *NodeStatusResponse, peerIdentity *peerIdentityDigest, source string) (uint64, bool, error) { + return c.commitParsedDeltaBase(nodeName, previous, merged, peerIdentity, source, previous.Status) +} + +func (c *NodeStatusCache) commitParsedDeltaBase(nodeName string, previous *CachedNodeStatus, merged *NodeStatusResponse, peerIdentity *peerIdentityDigest, source string, base *NodeStatusResponse) (uint64, bool, error) { // Phase 3: Write lock for the brief pointer swap. c.mu.Lock() // Deletion and recreation can reuse a revision. The identity memo and @@ -361,19 +410,31 @@ func (c *NodeStatusCache) commitParsedDelta(nodeName string, previous *CachedNod } revision := entry.Revision + 1 - c.entries[nodeName] = &CachedNodeStatus{ - Status: merged, - ReceivedAt: time.Now(), - Source: source, - Revision: revision, - peerIdentity: peerIdentity, + + next, err := c.legacyEntryLocked(nodeName, merged, revision, peerIdentity, source, base) + if err != nil { + c.mu.Unlock() + + if errors.Is(err, errLegacyDetailBaseUnavailable) { + return entry.Revision, true, nil + } + + return entry.Revision, false, err } + + c.entries[nodeName] = next + c.observeLegacyLocked(nodeName, next) + fn := c.onChange - mergedPtr := c.entries[nodeName].Status + overviewFn := c.onOverviewChange c.mu.Unlock() - if fn != nil { - fn(nodeName, mergedPtr) + if next.Overview != nil { + if overviewFn != nil { + overviewFn(nodeName, next.overviewForNotification()) + } + } else if fn != nil { + fn(nodeName, next.Status) } return revision, false, nil @@ -431,6 +492,10 @@ func (c *NodeStatusCache) Delete(nodeName string) { defer c.mu.Unlock() delete(c.entries, nodeName) + + if c.details != nil { + c.details.Forget(nodeName) + } } // UpdateSource updates the cached status source for a node without changing @@ -461,16 +526,25 @@ func (c *NodeStatusCache) UpdateSourceIf(nodeName, expectedSource, source string updated := *entry updated.Source = source + + if entry.Overview != nil { + overview := *entry.Overview + overview.StatusSource = source + metadata := statuspkg.OverviewMetadata(overview) + updated.Overview = &overview + updated.Status = &metadata + } + c.entries[nodeName] = &updated fn := c.onChange overviewFn := c.onOverviewChange - statusCopy := entry.Status + statusCopy := updated.Status c.mu.Unlock() - if entry.Overview != nil && overviewFn != nil { - overview := *entry.Overview - overview.StatusSource = source - overviewFn(nodeName, overview) + if updated.Overview != nil { + if overviewFn != nil { + overviewFn(nodeName, updated.overviewForNotification()) + } } else if fn != nil { fn(nodeName, statusCopy) } @@ -486,6 +560,10 @@ func (c *NodeStatusCache) CleanupStaleEntries(validNodes map[string]bool) { for name := range c.entries { if !validNodes[name] { delete(c.entries, name) + + if c.details != nil { + c.details.Forget(name) + } } } } diff --git a/cmd/unbounded-net-controller/overview_diagnostics_test.go b/cmd/unbounded-net-controller/overview_diagnostics_test.go new file mode 100644 index 000000000..fc055c31d --- /dev/null +++ b/cmd/unbounded-net-controller/overview_diagnostics_test.go @@ -0,0 +1,148 @@ +// Copyright (c) Microsoft Corporation. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "encoding/json" + "slices" + "testing" + "time" + + statuspkg "github.com/Azure/unbounded/internal/net/status" + statusproto "github.com/Azure/unbounded/internal/net/status/proto" + statusv1alpha1 "github.com/Azure/unbounded/internal/net/status/v1alpha1" +) + +func TestOverviewDiagnosticMessagesMatchFullProblems(t *testing.T) { + now := time.Now() + yes := true + node := &NodeStatusResponse{ + NodeInfo: NodeInfo{Name: "node", K8sReady: "Ready", ProviderID: "azure://vm", WireGuard: &WireGuardStatusInfo{Interface: "wg0"}}, + Peers: []WireGuardPeerStatus{ + {Name: "same", Tunnel: PeerTunnelStatus{Protocol: "IPIP"}, HealthCheck: &HealthCheckPeerStatus{Status: " UP "}}, + {Name: "same", Tunnel: PeerTunnelStatus{LastHandshake: now}, HealthCheck: &HealthCheckPeerStatus{Status: " down "}}, + {Name: "same", HealthCheck: &HealthCheckPeerStatus{Enabled: true, Status: "DOWN"}}, + }, + RoutingTable: RoutingTableInfo{Routes: []RouteEntry{{ + NextHops: []NextHop{{Expected: &yes}, {Present: &yes}, {Expected: &yes}}, + }}}, + } + + fullProblems := collectClusterProblems(&ClusterStatusResponse{Nodes: []*NodeStatusResponse{node}}) + if len(fullProblems) != 1 || len(fullProblems[0].Errors) != 3 { + t.Fatalf("legacy diagnostic fixture changed: %+v", fullProblems) + } + + overview := statuspkg.OverviewFromStatus(node, now) + overview.NodeInfo.ProviderID = "" + messages := statuspkg.OverviewDiagnosticMessages(overview, node.NodeInfo.ProviderID) + slices.Sort(messages) + + want := slices.Clone(fullProblems[0].Errors) + slices.Sort(want) + + if !slices.Equal(messages, want) || overview.RouteMismatchCount != 3 || overview.UnhealthyPeerLinks != 2 { + t.Fatalf("summary diagnostics=%v, legacy=%v, overview=%+v", messages, want, overview) + } + + if otherCloud := statuspkg.OverviewDiagnosticMessages(overview, "other://vm"); len(otherCloud) != 2 { + t.Fatalf("IPIP warning applied outside Azure: %v", otherCloud) + } + + metadata := statuspkg.OverviewMetadata(overview) + metadata.NodeInfo.ProviderID = node.NodeInfo.ProviderID + cluster := &ClusterStatusResponse{ + Nodes: []*NodeStatusResponse{&metadata}, + NodeOverviews: map[string]*statusv1alpha1.NodeStatusOverview{"node": &overview}, + } + + problems := collectClusterProblems(cluster) + if len(problems) != 1 || !slices.Equal(problems[0].Errors, fullProblems[0].Errors) { + t.Fatalf("overview problem pipeline=%+v, legacy=%+v", problems, fullProblems) + } + + metadata.NodeInfo.ProviderID = "other://vm" + + if problems = collectClusterProblems(cluster); len(problems) != 1 || len(problems[0].Errors) != 2 { + t.Fatalf("controller ignored enriched cloud identity: %+v", problems) + } + + overview.RouteMismatch = false + overview.RouteMismatchCount = 0 + overview.UnhealthyPeerLinks = 0 + + if problems = collectClusterProblems(cluster); len(problems) != 0 { + t.Fatalf("overview peer counts incorrectly replaced diagnostic link health: %+v", problems) + } + + overview.RouteMismatchCount = 3 + overview.UnhealthyPeerLinks = 2 + + overview.UsesIPIP = false + if noIPIP := statuspkg.OverviewDiagnosticMessages(overview, node.NodeInfo.ProviderID); len(noIPIP) != 2 { + t.Fatalf("Azure warning applied without IPIP: %v", noIPIP) + } +} + +func TestProtoOverviewPreservesDiagnosticFacts(t *testing.T) { + got := protoToNodeOverview(&statusproto.NodeStatusOverview{ + NodeInfo: &statusproto.NodeInfo{Name: "node", ProviderId: "azure://vm"}, + PeerCount: 4, HealthyPeers: 1, RouteCount: 2, RouteMismatch: true, + RouteMismatchCount: 3, UnhealthyPeerLinks: 2, UsesIpip: true, + }) + if got.RouteMismatchCount != 3 || got.UnhealthyPeerLinks != 2 || !got.UsesIPIP || + got.PeerCount != 4 || got.HealthyPeers != 1 || !got.RouteMismatch || got.NodeInfo.ProviderID != "azure://vm" { + t.Fatalf("overview converter lost diagnostics: %+v", got) + } +} + +func TestViewerSummaryPreservesInterfaceOnlineIndependentOfCNI(t *testing.T) { + for _, native := range []bool{false, true} { + for _, tc := range []struct { + name string + status *WireGuardStatusInfo + online bool + }{ + {"missing", nil, false}, + {"public key without interface", &WireGuardStatusInfo{PublicKey: "key"}, false}, + {"interface with CNI failure", &WireGuardStatusInfo{Interface: "wg0"}, true}, + } { + t.Run(tc.name, func(t *testing.T) { + node := &NodeStatusResponse{ + NodeInfo: NodeInfo{Name: "node", WireGuard: tc.status}, + NodeErrors: []NodeError{{Type: "cni", Message: "blocked"}}, + } + + cluster := &ClusterStatusResponse{Nodes: []*NodeStatusResponse{node}} + if native { + cluster.NodeOverviews = map[string]*statusv1alpha1.NodeStatusOverview{ + "node": {NodeInfo: node.NodeInfo, NodeErrors: node.NodeErrors}, + } + } + + summary := buildClusterSummary(cluster).NodeSummaries[0] + if summary.WireGuardOnline != tc.online || summary.ErrorCount != 1 || + summary.FirstError != "blocked" || summary.CniStatus != "Errors" { + t.Fatalf("interface/CNI facts conflated: %+v", summary) + } + + data, err := json.Marshal(summary) + if err != nil { + t.Fatal(err) + } + + var decoded struct { + Online *bool `json:"wireGuardOnline"` + } + if err := json.Unmarshal(data, &decoded); err != nil { + t.Fatal(err) + } + + if decoded.Online == nil || *decoded.Online != tc.online { + t.Fatalf("explicit online/offline fact omitted or changed: %s", data) + } + }) + } + } +} diff --git a/cmd/unbounded-net-controller/server.go b/cmd/unbounded-net-controller/server.go index dace866f1..0c19a1e84 100644 --- a/cmd/unbounded-net-controller/server.go +++ b/cmd/unbounded-net-controller/server.go @@ -1272,7 +1272,11 @@ func handleStatusPushRequestWithSource(health *healthState, bodyBytes []byte, so return NodeStatusPushAck{}, http.StatusBadRequest, fmt.Errorf("nodeInfo.name is required") } - ack.Revision = health.statusCache.StoreFull(nodeStatus.NodeInfo.Name, nodeStatus, source) + ack.Revision, err = health.statusCache.StoreFullChecked(nodeStatus.NodeInfo.Name, nodeStatus, source) + if err != nil { + return NodeStatusPushAck{}, http.StatusServiceUnavailable, fmt.Errorf("failed to store full status: %w", err) + } + klog.V(5).Infof("Received full status push from node %s", nodeStatus.NodeInfo.Name) return ack, http.StatusOK, nil @@ -1318,7 +1322,11 @@ func handleStatusPushRequestWithSource(health *healthState, bodyBytes []byte, so envelope.Status.NodeInfo.Name = nodeName } - ack.Revision = health.statusCache.StoreFull(nodeName, *envelope.Status, source) + ack.Revision, err = health.statusCache.StoreFullChecked(nodeName, *envelope.Status, source) + if err != nil { + return NodeStatusPushAck{}, http.StatusServiceUnavailable, fmt.Errorf("failed to store full status: %w", err) + } + klog.V(5).Infof("Received full status push from node %s", nodeName) return ack, http.StatusOK, nil @@ -1409,7 +1417,10 @@ func handleNodeStatusWSMessageWithSource(health *healthState, data []byte, sourc message.Status.NodeInfo.Name = nodeName } - rev := health.statusCache.StoreFull(nodeName, *message.Status, source) + rev, err := health.statusCache.StoreFullChecked(nodeName, *message.Status, source) + if err != nil { + return "node_status_resync", NodeStatusPushAck{Status: "resync_required", Reason: err.Error()} + } return "node_status_ack", NodeStatusPushAck{Status: "ok", Revision: rev} case "node_status_delta": diff --git a/cmd/unbounded-net-controller/status_overview_ingestion_test.go b/cmd/unbounded-net-controller/status_overview_ingestion_test.go new file mode 100644 index 000000000..05e49ee46 --- /dev/null +++ b/cmd/unbounded-net-controller/status_overview_ingestion_test.go @@ -0,0 +1,199 @@ +// Copyright (c) Microsoft Corporation. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "encoding/json" + "net/http" + "testing" + "time" + + "google.golang.org/protobuf/proto" + + statusproto "github.com/Azure/unbounded/internal/net/status/proto" + statusv1alpha1 "github.com/Azure/unbounded/internal/net/status/v1alpha1" +) + +func submitOverview(t *testing.T, channel string, health *healthState, message *statusproto.NodeStatusMessage) NodeStatusPushAck { + t.Helper() + + var ( + data []byte + err error + ) + if channel == "proto-http" || channel == "proto-ws" { + data, err = proto.Marshal(message) + } else { + envelope := NodeStatusWSMessage{ + Type: message.Type, NodeName: message.NodeName, DetailRequestID: message.DetailRequestId, + DetailError: message.DetailError, + } + if message.Summary != nil { + overview := protoToNodeOverview(message.Summary) + envelope.Summary = &overview + } + + if message.Status != nil { + status := protoToNodeStatus(message.Status) + envelope.Status = &status + } + + if message.Delta != nil { + envelope.Delta = map[string]json.RawMessage{"timestamp": json.RawMessage(`null`)} + } + + data, err = json.Marshal(envelope) + } + + if err != nil { + t.Fatal(err) + } + + switch channel { + case "proto-http": + ack, code, err := handleProtoPushRequest(health, data, "push") + if err != nil || code != http.StatusOK { + return NodeStatusPushAck{Status: "rejected"} + } + + return ack + case "json-http": + ack, code, err := handleStatusPushRequestWithSource(health, data, "push") + if err != nil || code != http.StatusOK { + return NodeStatusPushAck{Status: "rejected"} + } + + return ack + case "proto-ws": + decoded, err := decodeProtoWSMessage(data) + if err != nil { + return NodeStatusPushAck{Status: "rejected"} + } + + _, ack := handleProtoWSMessage(health, decoded, "ws") + + return ack + case "json-ws": + _, ack := handleNodeStatusWSMessageWithSource(health, data, "ws") + return ack + default: + t.Fatalf("unknown test channel %s", channel) + return NodeStatusPushAck{} + } +} + +func TestOverviewIngestionAllChannels(t *testing.T) { + for _, channel := range []string{"proto-http", "json-http", "proto-ws", "json-ws"} { + t.Run(channel, func(t *testing.T) { + health := &healthState{statusCache: NewNodeStatusCache()} + message := &statusproto.NodeStatusMessage{ + Type: statusv1alpha1.NodeStatusSummaryType, NodeName: "node", + Summary: &statusproto.NodeStatusOverview{ + TimestampUnixNs: 1000, LastPushTimeUnixNs: 2000, + NodeInfo: &statusproto.NodeInfo{ + Name: "node", SiteName: "site", WireGuard: &statusproto.WireGuardStatusInfo{Interface: "wg0"}, + }, + PeerCount: 10, HealthyPeers: 8, RouteCount: 20, RouteMismatch: true, + NodeErrors: []*statusproto.NodeError{{Type: "cni", Message: "blocked"}}, + HealthCheck: &statusproto.HealthCheckStatus{Summary: "not healthy"}, + NodePodInfo: &statusproto.NodePodInfo{PodName: "agent"}, + }, + } + + ack := submitOverview(t, channel, health, message) + if ack.Status != "ok" || ack.Revision != 1 || !ack.SummarySupported { + t.Fatalf("unexpected ACK: %+v", ack) + } + + entry, ok := health.statusCache.Get("node") + if !ok || entry.Overview == nil { + t.Fatal("overview was not stored") + } + + overview := entry.Overview + if overview.PeerCount != 10 || overview.HealthyPeers != 8 || overview.RouteCount != 20 || !overview.RouteMismatch || + overview.NodeErrors[0].Message != "blocked" || overview.NodeInfo.SiteName != "site" || + overview.NodeInfo.WireGuard.Interface != "wg0" || overview.NodePodInfo.PodName != "agent" || + overview.HealthCheck.Summary != "not healthy" || + !overview.Timestamp.Equal(time.Unix(0, 1000)) || !overview.LastPushTime.Equal(time.Unix(0, 2000)) { + t.Fatalf("overview changed during ingestion: %+v", overview) + } + + if entry.Status.Peers != nil || entry.Status.RoutingTable.Routes != nil || entry.Status.BpfEntries != nil { + t.Fatal("summary retained diagnostic arrays") + } + + if next := submitOverview(t, channel, health, message); next.Revision != 2 { + t.Fatal("summary resync did not advance routine revision") + } + }) + } +} + +func TestOverviewIngestionRejectsInvalidEnvelopes(t *testing.T) { + for _, channel := range []string{"proto-http", "json-http", "proto-ws", "json-ws"} { + for _, tc := range []struct { + name string + mutate func(*statusproto.NodeStatusMessage) + }{ + {"missing", func(m *statusproto.NodeStatusMessage) { m.Summary = nil }}, + {"identity mismatch", func(m *statusproto.NodeStatusMessage) { m.Summary.NodeInfo.Name = "other" }}, + {"negative counts", func(m *statusproto.NodeStatusMessage) { m.Summary.PeerCount = -1 }}, + {"impossible counts", func(m *statusproto.NodeStatusMessage) { m.Summary.HealthyPeers = 1 }}, + {"full mixed with summary", func(m *statusproto.NodeStatusMessage) { m.Status = &statusproto.NodeStatusFull{} }}, + {"delta mixed with summary", func(m *statusproto.NodeStatusMessage) { m.Delta = &statusproto.NodeStatusDelta{} }}, + {"detail correlation on summary", func(m *statusproto.NodeStatusMessage) { m.DetailRequestId = "request" }}, + {"collection error on summary", func(m *statusproto.NodeStatusMessage) { m.DetailError = "failed" }}, + {"summary in full", func(m *statusproto.NodeStatusMessage) { m.Type = "node_status_full" }}, + } { + t.Run(channel+"/"+tc.name, func(t *testing.T) { + health := &healthState{statusCache: NewNodeStatusCache()} + message := &statusproto.NodeStatusMessage{ + Type: statusv1alpha1.NodeStatusSummaryType, NodeName: "node", + Summary: &statusproto.NodeStatusOverview{NodeInfo: &statusproto.NodeInfo{Name: "node"}}, + } + tc.mutate(message) + + ack := submitOverview(t, channel, health, message) + if ack.Status == "ok" || health.statusCache.Len() != 0 { + t.Fatalf("invalid summary accepted: %+v", ack) + } + }) + } + } +} + +func TestOverviewIdentityRejectsDuplicateAndConflictingFields(t *testing.T) { + for _, data := range []string{ + `{"nodeName":"node","summary":{"nodeInfo":{"name":"other"}}}`, + `{"nodeName":"node","summary":{"nodeInfo":{"name":"other"}},"Summary":null}`, + `{"summary":{"nodeInfo":{"name":"other","Name":"node"}}}`, + `{"summary":{"nodeInfo":{"name":"other"},"NodeInfo":{"name":"node"}}}`, + } { + if _, err := extractNodeNameFromWSMessage([]byte(data)); err == nil { + t.Fatalf("ambiguous summary identity accepted: %s", data) + } + } + + name, err := extractNodeNameFromWSMessage([]byte(`{"summary":{"nodeInfo":{"name":"node"}}}`)) + if err != nil || name != "node" { + t.Fatalf("summary-only identity lost: %q %v", name, err) + } +} + +func TestOverviewCapabilityProtoAck(t *testing.T) { + data, err := marshalProtoAck("node_status_ack", NodeStatusPushAck{Status: "ok", Revision: 7}) + if err != nil { + t.Fatal(err) + } + + var ack statusproto.NodeStatusAck + if err := proto.Unmarshal(data, &ack); err != nil { + t.Fatal(err) + } + + if !ack.SummarySupported || !ack.PeerMeasurements || ack.Revision != 7 { + t.Fatalf("ACK lost capability or revision: %v", &ack) + } +} diff --git a/cmd/unbounded-net-controller/status_proto.go b/cmd/unbounded-net-controller/status_proto.go index 635ebc611..f2229c1eb 100644 --- a/cmd/unbounded-net-controller/status_proto.go +++ b/cmd/unbounded-net-controller/status_proto.go @@ -5,6 +5,7 @@ package main import ( "fmt" + "net/http" "time" "google.golang.org/protobuf/proto" @@ -531,7 +532,10 @@ func handleProtoWSMessage(health *healthState, decoded *decodedProtoWSMessage, s status.NodeInfo.Name = nodeName } - rev := health.statusCache.StoreFull(nodeName, status, source) + rev, err := health.statusCache.StoreFullChecked(nodeName, status, source) + if err != nil { + return "node_status_resync", NodeStatusPushAck{Status: "resync_required", Reason: err.Error()} + } return "node_status_ack", NodeStatusPushAck{Status: "ok", Revision: rev} case "node_status_delta": @@ -624,7 +628,10 @@ func handleProtoPushRequest(health *healthState, bodyBytes []byte, source string status.NodeInfo.Name = nodeName } - ack.Revision = health.statusCache.StoreFull(nodeName, status, source) + ack.Revision, err = health.statusCache.StoreFullChecked(nodeName, status, source) + if err != nil { + return NodeStatusPushAck{}, http.StatusServiceUnavailable, fmt.Errorf("failed to store full status: %w", err) + } return ack, 200, nil case "node_status_delta": diff --git a/cmd/unbounded-net-controller/status_types.go b/cmd/unbounded-net-controller/status_types.go index db908722c..c0c3ac236 100644 --- a/cmd/unbounded-net-controller/status_types.go +++ b/cmd/unbounded-net-controller/status_types.go @@ -6,6 +6,7 @@ package main import ( "bytes" "encoding/json" + "reflect" "sort" "time" @@ -166,20 +167,23 @@ type ClusterSummary struct { // NodeSummary is a compact per-node summary for use in ClusterSummary. type NodeSummary struct { - Name string `json:"name"` - SiteName string `json:"siteName,omitempty"` - IsGateway bool `json:"isGateway,omitempty"` - K8sReady string `json:"k8sReady,omitempty"` - StatusSource string `json:"statusSource,omitempty"` - CniStatus string `json:"cniStatus,omitempty"` - CniTone string `json:"cniTone,omitempty"` - ErrorCount int `json:"errorCount,omitempty"` - FirstError string `json:"firstError,omitempty"` - PeerCount int `json:"peerCount,omitempty"` - HealthyPeers int `json:"healthyPeers,omitempty"` - RouteCount int `json:"routeCount,omitempty"` - RouteMismatch bool `json:"routeMismatch,omitempty"` - FetchError string `json:"fetchError,omitempty"` + NodeInfo *NodeInfo `json:"nodeInfo,omitempty"` + LastPushTime *time.Time `json:"lastPushTime,omitempty"` + Name string `json:"name"` + SiteName string `json:"siteName,omitempty"` + IsGateway bool `json:"isGateway,omitempty"` + K8sReady string `json:"k8sReady,omitempty"` + StatusSource string `json:"statusSource,omitempty"` + CniStatus string `json:"cniStatus,omitempty"` + CniTone string `json:"cniTone,omitempty"` + ErrorCount int `json:"errorCount,omitempty"` + FirstError string `json:"firstError,omitempty"` + PeerCount int `json:"peerCount,omitempty"` + HealthyPeers int `json:"healthyPeers,omitempty"` + RouteCount int `json:"routeCount,omitempty"` + RouteMismatch bool `json:"routeMismatch,omitempty"` + FetchError string `json:"fetchError,omitempty"` + WireGuardOnline bool `json:"wireGuardOnline"` } // buildClusterSummary extracts a ClusterSummary from a full ClusterStatusResponse. @@ -197,18 +201,22 @@ func buildClusterSummary(status *ClusterStatusResponse) *ClusterSummary { overview = &projected } + info := node.NodeInfo ns := NodeSummary{ - Name: node.NodeInfo.Name, - SiteName: node.NodeInfo.SiteName, - IsGateway: node.NodeInfo.IsGateway, - K8sReady: node.NodeInfo.K8sReady, - StatusSource: node.StatusSource, - PeerCount: overview.PeerCount, - HealthyPeers: overview.HealthyPeers, - RouteCount: overview.RouteCount, - RouteMismatch: overview.RouteMismatch, - FetchError: node.FetchError, - ErrorCount: len(node.NodeErrors), + NodeInfo: &info, + LastPushTime: node.LastPushTime, + Name: node.NodeInfo.Name, + SiteName: node.NodeInfo.SiteName, + IsGateway: node.NodeInfo.IsGateway, + K8sReady: node.NodeInfo.K8sReady, + StatusSource: node.StatusSource, + PeerCount: overview.PeerCount, + HealthyPeers: overview.HealthyPeers, + RouteCount: overview.RouteCount, + RouteMismatch: overview.RouteMismatch, + FetchError: node.FetchError, + ErrorCount: len(node.NodeErrors), + WireGuardOnline: node.NodeInfo.WireGuard != nil && node.NodeInfo.WireGuard.Interface != "", } // Include first error message so the frontend can show it inline @@ -424,7 +432,7 @@ func computeClusterSummaryDelta(prev, curr *ClusterSummary) *ClusterSummaryDelta for _, ns := range curr.NodeSummaries { old, existed := prevByName[ns.Name] - if !existed || ns != old { + if !existed || !reflect.DeepEqual(ns, old) { delta.NodeSummaries = append(delta.NodeSummaries, ns) changed = true } diff --git a/cmd/unbounded-net-controller/summary_metadata_test.go b/cmd/unbounded-net-controller/summary_metadata_test.go new file mode 100644 index 000000000..a37b55165 --- /dev/null +++ b/cmd/unbounded-net-controller/summary_metadata_test.go @@ -0,0 +1,62 @@ +// Copyright (c) Microsoft Corporation. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "encoding/json" + "reflect" + "strings" + "testing" + "time" + + statusv1alpha1 "github.com/Azure/unbounded/internal/net/status/v1alpha1" +) + +func TestSummaryPreservesNodeInfoAndFreshnessWithoutDetails(t *testing.T) { + received := time.Now() + node := retentionFixture(10000) + node.LastPushTime = &received + node.NodeInfo.ProviderID = "azure://vm" + node.NodeInfo.InternalIPs = []string{"192.0.2.1"} + node.NodeInfo.K8sLabels = map[string]string{"node.kubernetes.io/instance-type": "test"} + node.NodeInfo.K8sUpdatedAt = &received + node.NodeInfo.BuildInfo = &statusv1alpha1.BuildInfo{Version: "test"} + cluster := &ClusterStatusResponse{Nodes: []*NodeStatusResponse{&node}} + summary := buildClusterSummary(cluster) + + row := summary.NodeSummaries[0] + if row.NodeInfo == &node.NodeInfo || !reflect.DeepEqual(row.NodeInfo, &node.NodeInfo) || + row.LastPushTime == nil || !row.LastPushTime.Equal(received) { + t.Fatal("summary lost metadata/freshness or retained the full node allocation") + } + + data, err := json.Marshal(summary) + if err != nil { + t.Fatal(err) + } + + for _, forbidden := range []string{`"peers"`, `"routingTable"`, `"bpfEntries"`, "private-detail-marker"} { + if strings.Contains(string(data), forbidden) { + t.Fatalf("summary retained diagnostic data: %s", forbidden) + } + } + + if delta := computeClusterSummaryDelta(summary, buildClusterSummary(cluster)); delta != nil { + t.Fatalf("equal metadata copies emitted a spurious update: %+v", delta) + } + + node.NodeInfo = NodeInfo{Name: "node", ProviderID: "azure://replacement"} + + next := buildClusterSummary(cluster) + if delta := computeClusterSummaryDelta(summary, next); delta == nil || len(delta.NodeSummaries) != 1 { + t.Fatal("metadata-only update was omitted") + } + + later := received.Add(time.Second) + node.LastPushTime = &later + + if delta := computeClusterSummaryDelta(next, buildClusterSummary(cluster)); delta == nil || len(delta.NodeSummaries) != 1 { + t.Fatal("freshness-only update was omitted") + } +} diff --git a/cmd/unbounded-net-controller/summary_receipt_time_test.go b/cmd/unbounded-net-controller/summary_receipt_time_test.go new file mode 100644 index 000000000..57c81e430 --- /dev/null +++ b/cmd/unbounded-net-controller/summary_receipt_time_test.go @@ -0,0 +1,102 @@ +// Copyright (c) Microsoft Corporation. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "testing" + "testing/synctest" + "time" + + statusv1alpha1 "github.com/Azure/unbounded/internal/net/status/v1alpha1" +) + +func TestSummaryNotificationsPreserveReceiptTime(t *testing.T) { + for _, mode := range []string{"summary", "legacy-full", "legacy-delta"} { + t.Run(mode, func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + cache := NewNodeStatusCache() + cache.BindDetails(testDetailRequests(t, nodeDetailRequestHooks{})) + + cluster := NewClusterStatusCache(&healthState{}) + cluster.status = &ClusterStatusResponse{ + Nodes: []*NodeStatusResponse{{NodeInfo: NodeInfo{Name: "node"}}}, + } + cluster.nodeIndex["node"] = 0 + cache.SetOnOverviewChange(cluster.PatchOverview) + + claimedTime := time.Unix(1, 0) + + var revision uint64 + + publish := func() { + t.Helper() + + var err error + + switch { + case mode == "summary": + revision, err = cache.StoreOverview("node", statusv1alpha1.NodeStatusOverview{ + NodeInfo: NodeInfo{Name: "node"}, LastPushTime: &claimedTime, + }, "ws") + case revision == 0 || mode == "legacy-full": + status := retentionFixture(1) + status.LastPushTime = &claimedTime + revision, err = cache.StoreFullChecked("node", status, "ws") + default: + now := time.Now() + + var resync bool + + revision, resync, err = cache.ApplyParsedDelta("node", revision, parsedDelta{timestamp: &now}, "ws") + if resync { + t.Fatal("unexpected resync") + } + } + + if err != nil { + t.Fatal(err) + } + } + assertReceipt := func(want time.Time) *ClusterSummary { + t.Helper() + + summary := buildClusterSummary(cluster.Get()) + + got := summary.NodeSummaries[0].LastPushTime + if got == nil || !got.Equal(want) { + t.Fatalf("summary receipt time = %v, want %v", got, want) + } + + return summary + } + + publish() + + entry, _ := cache.Get("node") + + first := assertReceipt(entry.ReceivedAt) + if !entry.Overview.LastPushTime.Equal(claimedTime) { + t.Fatal("notification changed the stored wire metadata") + } + + time.Sleep(time.Second) + cache.UpdateSource("node", "apiserver-ws") + assertReceipt(entry.ReceivedAt) + time.Sleep(time.Second) + publish() + + nextEntry, _ := cache.Get("node") + + next := assertReceipt(nextEntry.ReceivedAt) + if !first.NodeSummaries[0].LastPushTime.Equal(entry.ReceivedAt) { + t.Fatal("new notification mutated an earlier snapshot") + } + + if delta := computeClusterSummaryDelta(first, next); delta == nil || len(delta.NodeSummaries) != 1 { + t.Fatal("receipt-only update did not reach the summary delta") + } + }) + }) + } +}