Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 3 additions & 6 deletions cmd/unbounded-net-controller/cluster_status.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)
Expand Down Expand Up @@ -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
Expand Down
32 changes: 25 additions & 7 deletions cmd/unbounded-net-controller/detail_cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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")
}
Expand All @@ -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()
Expand Down
151 changes: 151 additions & 0 deletions cmd/unbounded-net-controller/detail_legacy.go
Original file line number Diff line number Diff line change
@@ -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
}
Loading