|
| 1 | +// Copyright (c) 2026, NVIDIA CORPORATION. All rights reserved. |
| 2 | +// |
| 3 | +// Licensed under the Apache License, Version 2.0 (the "License"); |
| 4 | +// you may not use this file except in compliance with the License. |
| 5 | +// You may obtain a copy of the License at |
| 6 | +// |
| 7 | +// http://www.apache.org/licenses/LICENSE-2.0 |
| 8 | +// |
| 9 | +// Unless required by applicable law or agreed to in writing, software |
| 10 | +// distributed under the License is distributed on an "AS IS" BASIS, |
| 11 | +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
| 12 | +// See the License for the specific language governing permissions and |
| 13 | +// limitations under the License. |
| 14 | + |
| 15 | +package monitor |
| 16 | + |
| 17 | +import ( |
| 18 | + "context" |
| 19 | + "errors" |
| 20 | + "fmt" |
| 21 | + "log/slog" |
| 22 | + "sort" |
| 23 | + "strings" |
| 24 | + "sync" |
| 25 | + "time" |
| 26 | + |
| 27 | + "github.com/nvidia/nvsentinel/commons/pkg/healthpub" |
| 28 | + pb "github.com/nvidia/nvsentinel/data-models/pkg/protos" |
| 29 | + "github.com/nvidia/nvsentinel/health-monitors/nic-health-monitor/pkg/checks" |
| 30 | +) |
| 31 | + |
| 32 | +// pollStall reports a poll that stops completing. When the host stalls the |
| 33 | +// sysfs reads a poll makes, the loop freezes and liveness restarts the |
| 34 | +// container, repeatedly, without anything saying that NIC state cannot be |
| 35 | +// observed. The watchdog publishes that as an event before liveness acts. |
| 36 | +// |
| 37 | +// One node-level event covers both loops: unhealthy once any category has |
| 38 | +// been in flight past the deadline, healthy once none is. Once a poll has |
| 39 | +// completed after start, a healthy baseline is published, which closes a |
| 40 | +// stall left open by a liveness restart. |
| 41 | +// |
| 42 | +// Only the watchdog publishes these events, so they leave in order and a |
| 43 | +// publish retrying against an unavailable platform connector never holds up |
| 44 | +// a polling loop. |
| 45 | +type pollStall struct { |
| 46 | + deadline time.Duration |
| 47 | + strategy pb.ProcessingStrategy |
| 48 | + now func() time.Time |
| 49 | + // waitingOnServer reports a publish waiting out the platform connector's |
| 50 | + // retry window; a poll blocked there is not a NIC stall. |
| 51 | + waitingOnServer func() bool |
| 52 | + |
| 53 | + mu sync.Mutex |
| 54 | + started map[string]time.Time |
| 55 | + completed bool |
| 56 | + // pending is the message of a stall not yet delivered. It outlives the |
| 57 | + // stall, so a stall that ends before a failed publish is retried is |
| 58 | + // still reported. |
| 59 | + pending string |
| 60 | + reported bool |
| 61 | + baselined bool |
| 62 | +} |
| 63 | + |
| 64 | +// EnablePollStallDetection turns on the poll stall watchdog. A deadline of |
| 65 | +// zero or less leaves it off. Call before the polling loops start. |
| 66 | +func (m *NICHealthMonitor) EnablePollStallDetection(deadline time.Duration, strategy pb.ProcessingStrategy) { |
| 67 | + if deadline <= 0 { |
| 68 | + return |
| 69 | + } |
| 70 | + |
| 71 | + m.stall = &pollStall{ |
| 72 | + deadline: deadline, |
| 73 | + strategy: strategy, |
| 74 | + now: time.Now, |
| 75 | + started: map[string]time.Time{}, |
| 76 | + |
| 77 | + waitingOnServer: m.WaitingOnServer, |
| 78 | + } |
| 79 | + |
| 80 | + slog.Info("NIC poll stall detection enabled", "deadline", deadline, "processing_strategy", strategy.String()) |
| 81 | +} |
| 82 | + |
| 83 | +// RunPollStallWatchdog checks for stalled polls until ctx is cancelled. It |
| 84 | +// returns at once when stall detection is off. |
| 85 | +func (m *NICHealthMonitor) RunPollStallWatchdog(ctx context.Context) error { |
| 86 | + if m.stall == nil { |
| 87 | + return nil |
| 88 | + } |
| 89 | + |
| 90 | + interval := min(time.Second, m.stall.deadline/4) |
| 91 | + ticker := time.NewTicker(interval) |
| 92 | + |
| 93 | + defer ticker.Stop() |
| 94 | + |
| 95 | + for { |
| 96 | + select { |
| 97 | + case <-ctx.Done(): |
| 98 | + return nil |
| 99 | + case <-ticker.C: |
| 100 | + m.checkPollStalls(ctx) |
| 101 | + } |
| 102 | + } |
| 103 | +} |
| 104 | + |
| 105 | +// beginPoll records that a poll of category is in flight. |
| 106 | +func (m *NICHealthMonitor) beginPoll(category string) { |
| 107 | + if m.stall == nil { |
| 108 | + return |
| 109 | + } |
| 110 | + |
| 111 | + m.stall.mu.Lock() |
| 112 | + m.stall.started[category] = m.stall.now() |
| 113 | + m.stall.mu.Unlock() |
| 114 | +} |
| 115 | + |
| 116 | +// endPoll records that a poll of category completed. The watchdog publishes |
| 117 | +// any recovery this causes. |
| 118 | +func (m *NICHealthMonitor) endPoll(category string) { |
| 119 | + s := m.stall |
| 120 | + if s == nil { |
| 121 | + return |
| 122 | + } |
| 123 | + |
| 124 | + s.mu.Lock() |
| 125 | + delete(s.started, category) |
| 126 | + s.completed = true |
| 127 | + s.mu.Unlock() |
| 128 | +} |
| 129 | + |
| 130 | +// checkPollStalls publishes the stall event once any category has been in |
| 131 | +// flight past the deadline, and the healthy event once none is and a poll |
| 132 | +// has completed. A failed publish is retried on the next check. Nothing is |
| 133 | +// checked while a publish waits on the platform connector. |
| 134 | +func (m *NICHealthMonitor) checkPollStalls(ctx context.Context) { |
| 135 | + s := m.stall |
| 136 | + |
| 137 | + if s.waitingOnServer() { |
| 138 | + return |
| 139 | + } |
| 140 | + |
| 141 | + s.mu.Lock() |
| 142 | + stalled := s.stalledLocked() |
| 143 | + |
| 144 | + if len(stalled) > 0 && !s.reported { |
| 145 | + s.pending = stallMessage(stalled) |
| 146 | + } |
| 147 | + s.mu.Unlock() |
| 148 | + |
| 149 | + if m.publishPendingStall(ctx) { |
| 150 | + m.publishRecovery(ctx) |
| 151 | + } |
| 152 | +} |
| 153 | + |
| 154 | +// publishPendingStall publishes the pending stall event, if any, and reports |
| 155 | +// whether none is left pending. |
| 156 | +func (m *NICHealthMonitor) publishPendingStall(ctx context.Context) bool { |
| 157 | + s := m.stall |
| 158 | + |
| 159 | + s.mu.Lock() |
| 160 | + pending := s.pending |
| 161 | + s.mu.Unlock() |
| 162 | + |
| 163 | + if pending == "" { |
| 164 | + return true |
| 165 | + } |
| 166 | + |
| 167 | + evt := checks.NewHealthEvent(m.nodeName, checks.PollStallCheckName, |
| 168 | + pending, nil, false, false, pb.RecommendedAction_NONE, s.strategy) |
| 169 | + if !m.publishStallEvent(ctx, evt) { |
| 170 | + return false |
| 171 | + } |
| 172 | + |
| 173 | + s.mu.Lock() |
| 174 | + s.pending = "" |
| 175 | + s.reported = true |
| 176 | + s.mu.Unlock() |
| 177 | + |
| 178 | + return true |
| 179 | +} |
| 180 | + |
| 181 | +// publishRecovery publishes the healthy event when a poll has completed, none |
| 182 | +// is stalled, and it either ends a reported stall or is the first since start. |
| 183 | +func (m *NICHealthMonitor) publishRecovery(ctx context.Context) { |
| 184 | + s := m.stall |
| 185 | + |
| 186 | + s.mu.Lock() |
| 187 | + needHealthy := s.completed && len(s.stalledLocked()) == 0 && (s.reported || !s.baselined) |
| 188 | + s.mu.Unlock() |
| 189 | + |
| 190 | + if !needHealthy { |
| 191 | + return |
| 192 | + } |
| 193 | + |
| 194 | + evt := checks.NewHealthEvent(m.nodeName, checks.PollStallCheckName, |
| 195 | + "NIC polls are completing", nil, false, true, pb.RecommendedAction_NONE, s.strategy) |
| 196 | + if !m.publishStallEvent(ctx, evt) { |
| 197 | + return |
| 198 | + } |
| 199 | + |
| 200 | + s.mu.Lock() |
| 201 | + s.reported = false |
| 202 | + s.baselined = true |
| 203 | + s.mu.Unlock() |
| 204 | +} |
| 205 | + |
| 206 | +// stalledLocked returns how long each category has been in flight, for those |
| 207 | +// past the deadline. The caller holds s.mu. |
| 208 | +func (s *pollStall) stalledLocked() map[string]time.Duration { |
| 209 | + now := s.now() |
| 210 | + stalled := map[string]time.Duration{} |
| 211 | + |
| 212 | + for category, start := range s.started { |
| 213 | + if elapsed := now.Sub(start); elapsed >= s.deadline { |
| 214 | + stalled[category] = elapsed |
| 215 | + } |
| 216 | + } |
| 217 | + |
| 218 | + return stalled |
| 219 | +} |
| 220 | + |
| 221 | +// stallMessage names the stalled categories in a stable order. |
| 222 | +func stallMessage(stalled map[string]time.Duration) string { |
| 223 | + categories := make([]string, 0, len(stalled)) |
| 224 | + for category := range stalled { |
| 225 | + categories = append(categories, category) |
| 226 | + } |
| 227 | + |
| 228 | + sort.Strings(categories) |
| 229 | + |
| 230 | + parts := make([]string, 0, len(categories)) |
| 231 | + for _, category := range categories { |
| 232 | + parts = append(parts, fmt.Sprintf("%s poll in flight for %s", category, stalled[category].Round(time.Second))) |
| 233 | + } |
| 234 | + |
| 235 | + return "NIC state cannot be observed: " + strings.Join(parts, ", ") |
| 236 | +} |
| 237 | + |
| 238 | +// publishStallEvent publishes one stall event and reports whether it is |
| 239 | +// settled. A permanent rejection counts as settled, as it does for check |
| 240 | +// events. |
| 241 | +func (m *NICHealthMonitor) publishStallEvent(ctx context.Context, evt *pb.HealthEvent) bool { |
| 242 | + batch := &pb.HealthEvents{Version: 1, Events: []*pb.HealthEvent{evt}} |
| 243 | + if err := m.pub.Publish(ctx, batch); err != nil { |
| 244 | + if errors.Is(err, healthpub.ErrPublishRejected) { |
| 245 | + slog.Error("Platform connector rejected NIC poll stall event for good; dropping it", |
| 246 | + "is_healthy", evt.IsHealthy, "error", err) |
| 247 | + |
| 248 | + return true |
| 249 | + } |
| 250 | + |
| 251 | + slog.Error("Failed to send NIC poll stall event", "is_healthy", evt.IsHealthy, "error", err) |
| 252 | + |
| 253 | + return false |
| 254 | + } |
| 255 | + |
| 256 | + m.logSentEvents(checks.PollStallCheckName, batch.Events) |
| 257 | + |
| 258 | + return true |
| 259 | +} |
0 commit comments