diff --git a/cmd/unbounded-net-controller/detail_cache.go b/cmd/unbounded-net-controller/detail_cache.go index ccd64fe2c..b53f69d2a 100644 --- a/cmd/unbounded-net-controller/detail_cache.go +++ b/cmd/unbounded-net-controller/detail_cache.go @@ -17,7 +17,12 @@ 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 +} // 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 +54,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 +72,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{}, errors.New("legacy detail base changed or expired") + } + } + 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_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) }