diff --git a/docs/architecture/interrupt-flow.md b/docs/architecture/interrupt-flow.md index 4be725da8..fa1d8307e 100644 --- a/docs/architecture/interrupt-flow.md +++ b/docs/architecture/interrupt-flow.md @@ -191,6 +191,46 @@ a pod that cannot finish terminating — a stuck finalizer, an unresponsive kubelet — blocks the node's interrupt until it is cleared. Set `spec.drainConfig.timeout` to bound that wait; there is no default. +### DrainBlocked Condition + +A `DrainBlocked` condition is set on the NodeWright whenever one or more selected +nodes have a drain that cannot currently make progress. It names the blocked +nodes and, where the blocker is a PodDisruptionBudget, includes the apiserver's +own message verbatim: + +```yaml +- type: DrainBlocked + status: "True" + reason: PodDisruptionBudget + message: "2/5 nodes blocked draining (node-a, node-c); default/web-0 on node-a: + The disruption budget web-pdb needs 3 healthy pods and has 3 currently" +``` + +`reason` is one of `PodDisruptionBudget`, `UnmanagedPod`, `EmptyDirData`, or +`MultipleCauses` when more than one kind of blocker is present across the +affected nodes. The condition clears automatically once every previously +blocked node has drained — no action is required beyond removing the +underlying blocker. + +`DrainBlocked` is independent of the `Blocked` condition (which is reserved for +an uninstalled dependency): a NodeWright can be both dependency-blocked and +drain-blocked at the same time, so the two conditions never share a type. + +A PodDisruptionBudget rejection is treated as a self-resolving wait state, not +a reconcile error: it no longer aborts the reconcile pass for the remaining +nodes. However, this also means it no longer puts the NodeWright into +controller-runtime's exponential backoff. With a PDB at zero allowed +disruptions and `spec.drainConfig.timeout` unset, the operator currently +retries the eviction every 2 seconds indefinitely, where it previously backed +off toward roughly 1000 seconds — a real increase in eviction-API load. +Per-node throttling for the drain-blocked case specifically is tracked in +[#632](https://github.com/NVIDIA/nodewright/issues/632) and not yet shipped. + +Unmanaged pods (`force: false`) and `emptyDir` pods (`deleteEmptyDirData: +false`) are also reported in `DrainBlocked`, though — unlike a PDB rejection — +these were already wait states before this condition existed; `DrainBlocked` +only makes them visible without reading operator logs. + ### Recovering From a Drain Timeout When `spec.drainConfig.timeout` expires, the operator records a `DrainTimeout` @@ -198,6 +238,13 @@ warning event, marks the node and NodeWright `erroring`, and leaves the node cordoned. The operator stops issuing further evict/delete actions while the blocking condition remains, so package stages do not proceed on that node. +`DrainTimeout` and `DrainBlocked` answer different questions: `DrainBlocked` +names what is currently preventing progress and clears itself once the +blocker is gone, with no configured limit on how long that may take. +`DrainTimeout` only fires once `spec.drainConfig.timeout` is set and elapses — +which, for a PDB blocker, is now the only path that turns a stuck drain into +an `erroring` node; the PDB rejection itself no longer does so directly. + To recover, remove the underlying blocker first, such as a PDB with zero allowed disruptions, an unmanaged pod when `force: false`, or an `emptyDir` pod when `deleteEmptyDirData: false`. Then reset the failed rollout metadata: diff --git a/docs/architecture/operator-status.md b/docs/architecture/operator-status.md index 3a0e633aa..93c79b4b6 100644 --- a/docs/architecture/operator-status.md +++ b/docs/architecture/operator-status.md @@ -90,6 +90,7 @@ The operator also sets additional condition types that may be useful for trouble - `NodesIgnored`: selected nodes are skipped because they have the ignore label set - `ApplyPackage`: the controller is applying a package to a node - `DeploymentPolicyNotFound`: the referenced `DeploymentPolicy` is missing at reconcile time +- `DrainBlocked`: one or more selected nodes have a drain that cannot currently make progress. Reasons: `PodDisruptionBudget`, `UnmanagedPod`, `EmptyDirData`, or `MultipleCauses` when more than one kind of blocker is present. Independent of `Blocked` — a NodeWright can be both dependency-blocked and drain-blocked at once. These conditions complement, rather than replace, `.status.status` and `Ready`. diff --git a/operator/internal/controller/cluster_state_v2.go b/operator/internal/controller/cluster_state_v2.go index abb9d7697..fd5c2d25f 100644 --- a/operator/internal/controller/cluster_state_v2.go +++ b/operator/internal/controller/cluster_state_v2.go @@ -19,6 +19,7 @@ package controller import ( + "context" "fmt" "sort" "strings" @@ -402,6 +403,7 @@ type SkyhookNodes interface { IsPaused() bool HasUninstallWork() (bool, error) UpdateBlockedCondition() error + UpdateDrainBlockedCondition(ctx context.Context, logger logr.Logger) UpdateUninstallConditions() error UpdateNodeStateMalformedCondition() NodeCount() int @@ -626,6 +628,81 @@ func (s *skyhookNodes) UpdateBlockedCondition() error { return nil } +// UpdateDrainBlockedCondition rebuilds the DrainBlocked condition from each node's +// persisted drain-blocker annotation (SkyhookNode.DrainBlocked), rather than from a +// single pass's live findings. This makes the condition level-triggered, matching +// UpdateBlockedCondition: it is correct from persisted state alone regardless of +// whether this reconcile pass actually ran RunSkyhookPackages for this Skyhook (paused, +// disabled, complete Skyhooks skip it) or returned from it early (an error, or +// spec.serial stopping after the first node) — those nodes simply keep whatever was +// last recorded for them, rather than being wrongly treated as unblocked. +// +// A node whose annotation fails to parse is skipped for this computation, the same +// tolerance UpdateBlockedCondition applies to unreadable nodeState — this is a +// NodeStateMalformed-adjacent concern, but has no dedicated user-visible signal of +// its own; a node dropping out of this aggregate silently is an accepted limitation +// rather than a deliberate design, tracked for follow-up. +func (s *skyhookNodes) UpdateDrainBlockedCondition(ctx context.Context, logger logr.Logger) { + blocks := make([]wrapper.DrainBlockedNode, 0, len(s.nodes)) + for _, node := range s.nodes { + // A node with no runnable, interrupt-requiring package this pass will never + // reach EnsureNodeIsReadyForInterrupt, so nothing else clears its persisted + // drain-blocker annotation. Clear it here instead: this function is the one + // thing that runs on every pass regardless of paused/disabled/complete state + // or an error elsewhere in the reconcile, which is what keeps a removed or + // finished drain from reporting a blocker that no longer exists. + if !nodeNeedsInterruptDrain(ctx, node) { + if err := node.SetDrainBlocked(nil); err != nil { + logger.Error(err, "clearing stale drain blocked state", "node", node.GetNode().Name) + } + } + + blocked, err := node.DrainBlocked() + if err != nil { + continue + } + if len(blocked) == 0 { + continue + } + blocks = append(blocks, wrapper.DrainBlockedNode{ + NodeName: node.GetNode().Name, + Blocked: blocked, + }) + } + + if len(blocks) == 0 { + wrapper.RemoveSkyhookConditionTypes(s.skyhook, wrapper.SkyhookConditionDrainBlocked) + return + } + + // The message builder truncates detail lines past drainBlockedDetailLineLimit and + // points the reader at controller logs for the rest, mirroring + // updateTaintToleranceCondition's log-before-truncate pattern — so log the same + // persisted set the message is built from here, not a caller-local subset. + if totalBlockedPods := countBlockedPods(blocks); totalBlockedPods > wrapper.ReadyConditionNodeListLimit { + logger.Info("DrainBlocked condition message truncated; full blocked set", "nodewright", s.skyhook.Name, "drainBlocks", blocks) + } + + wrapper.AddSkyhookCondition(s.skyhook, metav1.Condition{ + Type: wrapper.SkyhookConditionDrainBlocked, + Status: metav1.ConditionTrue, + ObservedGeneration: s.skyhook.Generation, + LastTransitionTime: metav1.Now(), + Reason: wrapper.DrainBlockedConditionReason(blocks), + Message: wrapper.DrainBlockedConditionMessage(blocks, len(s.nodes)), + }) +} + +// countBlockedPods sums Blocked across every node, for deciding whether the DrainBlocked +// message will be truncated and the full set needs logging. +func countBlockedPods(blocks []wrapper.DrainBlockedNode) int { + total := 0 + for _, b := range blocks { + total += len(b.Blocked) + } + return total +} + // isPackageCompleteOnAllNodes reports whether the package has reached its // terminal-complete stage (per node.IsPackageComplete semantics) on every node // this Skyhook selects. Returns false when there are no selected nodes: with diff --git a/operator/internal/controller/mock/SkyhookNodes.go b/operator/internal/controller/mock/SkyhookNodes.go index b08f680c0..754f02eeb 100644 --- a/operator/internal/controller/mock/SkyhookNodes.go +++ b/operator/internal/controller/mock/SkyhookNodes.go @@ -23,6 +23,8 @@ package controller import ( + "context" + "github.com/NVIDIA/nodewright/operator/api/nodewright/v1alpha1" "github.com/NVIDIA/nodewright/operator/internal/wrapper" "github.com/go-logr/logr" @@ -35,10 +37,19 @@ func NewMockSkyhookNodes(t interface { mock.TestingT Cleanup(func()) }) *MockSkyhookNodes { + if helper, ok := t.(interface{ Helper() }); ok { + helper.Helper() + } + mock := &MockSkyhookNodes{} mock.Mock.Test(t) - t.Cleanup(func() { mock.AssertExpectations(t) }) + t.Cleanup(func() { + if helper, ok := t.(interface{ Helper() }); ok { + helper.Helper() + } + mock.AssertExpectations(t) + }) return mock } @@ -70,7 +81,7 @@ type MockSkyhookNodes_AddCompartment_Call struct { // AddCompartment is a helper method to define mock.On call // - name string // - compartment *wrapper.Compartment -func (_e *MockSkyhookNodes_Expecter) AddCompartment(name interface{}, compartment interface{}) *MockSkyhookNodes_AddCompartment_Call { +func (_e *MockSkyhookNodes_Expecter) AddCompartment(name any, compartment any) *MockSkyhookNodes_AddCompartment_Call { return &MockSkyhookNodes_AddCompartment_Call{Call: _e.mock.On("AddCompartment", name, compartment)} } @@ -127,7 +138,7 @@ type MockSkyhookNodes_AddCompartmentNode_Call struct { // AddCompartmentNode is a helper method to define mock.On call // - name string // - node wrapper.SkyhookNode -func (_e *MockSkyhookNodes_Expecter) AddCompartmentNode(name interface{}, node interface{}) *MockSkyhookNodes_AddCompartmentNode_Call { +func (_e *MockSkyhookNodes_Expecter) AddCompartmentNode(name any, node any) *MockSkyhookNodes_AddCompartmentNode_Call { return &MockSkyhookNodes_AddCompartmentNode_Call{Call: _e.mock.On("AddCompartmentNode", name, node)} } @@ -172,7 +183,7 @@ type MockSkyhookNodes_AddNode_Call struct { // AddNode is a helper method to define mock.On call // - node wrapper.SkyhookNode -func (_e *MockSkyhookNodes_Expecter) AddNode(node interface{}) *MockSkyhookNodes_AddNode_Call { +func (_e *MockSkyhookNodes_Expecter) AddNode(node any) *MockSkyhookNodes_AddNode_Call { return &MockSkyhookNodes_AddNode_Call{Call: _e.mock.On("AddNode", node)} } @@ -232,7 +243,7 @@ type MockSkyhookNodes_AssignNodeToCompartment_Call struct { // AssignNodeToCompartment is a helper method to define mock.On call // - node wrapper.SkyhookNode -func (_e *MockSkyhookNodes_Expecter) AssignNodeToCompartment(node interface{}) *MockSkyhookNodes_AssignNodeToCompartment_Call { +func (_e *MockSkyhookNodes_Expecter) AssignNodeToCompartment(node any) *MockSkyhookNodes_AssignNodeToCompartment_Call { return &MockSkyhookNodes_AssignNodeToCompartment_Call{Call: _e.mock.On("AssignNodeToCompartment", node)} } @@ -430,7 +441,7 @@ type MockSkyhookNodes_GetNode_Call struct { // GetNode is a helper method to define mock.On call // - name string -func (_e *MockSkyhookNodes_Expecter) GetNode(name interface{}) *MockSkyhookNodes_GetNode_Call { +func (_e *MockSkyhookNodes_Expecter) GetNode(name any) *MockSkyhookNodes_GetNode_Call { return &MockSkyhookNodes_GetNode_Call{Call: _e.mock.On("GetNode", name)} } @@ -802,7 +813,7 @@ type MockSkyhookNodes_Migrate_Call struct { // Migrate is a helper method to define mock.On call // - logger logr.Logger -func (_e *MockSkyhookNodes_Expecter) Migrate(logger interface{}) *MockSkyhookNodes_Migrate_Call { +func (_e *MockSkyhookNodes_Expecter) Migrate(logger any) *MockSkyhookNodes_Migrate_Call { return &MockSkyhookNodes_Migrate_Call{Call: _e.mock.On("Migrate", logger)} } @@ -919,7 +930,7 @@ type MockSkyhookNodes_SetStatus_Call struct { // SetStatus is a helper method to define mock.On call // - status v1alpha1.Status -func (_e *MockSkyhookNodes_Expecter) SetStatus(status interface{}) *MockSkyhookNodes_SetStatus_Call { +func (_e *MockSkyhookNodes_Expecter) SetStatus(status any) *MockSkyhookNodes_SetStatus_Call { return &MockSkyhookNodes_SetStatus_Call{Call: _e.mock.On("SetStatus", status)} } @@ -1058,7 +1069,7 @@ type MockSkyhookNodes_UpdateCondition_Call struct { // UpdateCondition is a helper method to define mock.On call // - logger logr.Logger -func (_e *MockSkyhookNodes_Expecter) UpdateCondition(logger interface{}) *MockSkyhookNodes_UpdateCondition_Call { +func (_e *MockSkyhookNodes_Expecter) UpdateCondition(logger any) *MockSkyhookNodes_UpdateCondition_Call { return &MockSkyhookNodes_UpdateCondition_Call{Call: _e.mock.On("UpdateCondition", logger)} } @@ -1085,6 +1096,52 @@ func (_c *MockSkyhookNodes_UpdateCondition_Call) RunAndReturn(run func(logger lo return _c } +// UpdateDrainBlockedCondition provides a mock function for the type MockSkyhookNodes +func (_mock *MockSkyhookNodes) UpdateDrainBlockedCondition(ctx context.Context, logger logr.Logger) { + _mock.Called(ctx, logger) + return +} + +// MockSkyhookNodes_UpdateDrainBlockedCondition_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'UpdateDrainBlockedCondition' +type MockSkyhookNodes_UpdateDrainBlockedCondition_Call struct { + *mock.Call +} + +// UpdateDrainBlockedCondition is a helper method to define mock.On call +// - ctx context.Context +// - logger logr.Logger +func (_e *MockSkyhookNodes_Expecter) UpdateDrainBlockedCondition(ctx any, logger any) *MockSkyhookNodes_UpdateDrainBlockedCondition_Call { + return &MockSkyhookNodes_UpdateDrainBlockedCondition_Call{Call: _e.mock.On("UpdateDrainBlockedCondition", ctx, logger)} +} + +func (_c *MockSkyhookNodes_UpdateDrainBlockedCondition_Call) Run(run func(ctx context.Context, logger logr.Logger)) *MockSkyhookNodes_UpdateDrainBlockedCondition_Call { + _c.Call.Run(func(args mock.Arguments) { + var arg0 context.Context + if args[0] != nil { + arg0 = args[0].(context.Context) + } + var arg1 logr.Logger + if args[1] != nil { + arg1 = args[1].(logr.Logger) + } + run( + arg0, + arg1, + ) + }) + return _c +} + +func (_c *MockSkyhookNodes_UpdateDrainBlockedCondition_Call) Return() *MockSkyhookNodes_UpdateDrainBlockedCondition_Call { + _c.Call.Return() + return _c +} + +func (_c *MockSkyhookNodes_UpdateDrainBlockedCondition_Call) RunAndReturn(run func(context.Context, logr.Logger)) *MockSkyhookNodes_UpdateDrainBlockedCondition_Call { + _c.Run(run) + return _c +} + // UpdateNodeStateMalformedCondition provides a mock function for the type MockSkyhookNodes func (_mock *MockSkyhookNodes) UpdateNodeStateMalformedCondition() { _mock.Called() diff --git a/operator/internal/controller/skyhook_controller.go b/operator/internal/controller/skyhook_controller.go index f50b7605d..a4fe71e11 100644 --- a/operator/internal/controller/skyhook_controller.go +++ b/operator/internal/controller/skyhook_controller.go @@ -613,6 +613,12 @@ func (r *SkyhookReconciler) refreshSkyhookConditions(ctx context.Context, cluste if err := skyhook.UpdateBlockedCondition(); err != nil { return fmt.Errorf("error updating blocked condition: %w", err) } + // DrainBlocked (PDB/unmanaged-pod/emptyDir drain blockers). Distinct from the + // NonInterruptPodsRunning-flavored Blocked condition r.updateDrainBlockedCondition + // below maintains — same name prefix, different condition type. Rebuilt from + // persisted per-node state so it stays correct on paused/disabled/complete/error/ + // serial-partial passes; see cluster_state_v2.go's UpdateDrainBlockedCondition. + skyhook.UpdateDrainBlockedCondition(ctx, log.FromContext(ctx)) if err := r.updateDrainBlockedCondition(ctx, skyhook); err != nil { return fmt.Errorf("error updating drain blocked condition: %w", err) } @@ -1428,6 +1434,29 @@ func (r *SkyhookReconciler) TrackReboots(ctx context.Context, clusterState *clus return updates, utilerrors.NewAggregate(errs) } +// saveThenWrap persists any in-memory node mutations (for example SetDrainBlocked) before an +// early error return, so they reach the apiserver instead of being dropped when the next +// pass rebuilds cluster state. Save failures are aggregated with the original error. +func (r *SkyhookReconciler) saveThenWrap(ctx context.Context, clusterState *clusterState, skyhook SkyhookNodes, err error) error { + if _, saveErrs := r.SaveNodesAndSkyhook(ctx, clusterState, skyhook); len(saveErrs) > 0 { + return utilerrors.NewAggregate(append(saveErrs, err)) + } + return err +} + +// filterApplicablePackages drops packages whose uninstall is in progress or already +// completed on this node. See shouldSkipApplyForUninstall for the exact rule. +func filterApplicablePackages(toRun []*v1alpha1.Package, nodeState v1alpha1.NodeState, beingDeleted bool) []*v1alpha1.Package { + filtered := make([]*v1alpha1.Package, 0, len(toRun)) + for _, pkg := range toRun { + if shouldSkipApplyForUninstall(pkg, nodeState, beingDeleted) { + continue + } + filtered = append(filtered, pkg) + } + return filtered +} + // RunSkyhookPackages runs all skyhook packages then saves and requeues if changes were made func (r *SkyhookReconciler) RunSkyhookPackages(ctx context.Context, clusterState *clusterState, nodePicker *NodePicker, skyhook SkyhookNodes) (*ctrl.Result, error) { @@ -1462,6 +1491,7 @@ func (r *SkyhookReconciler) RunSkyhookPackages(ctx context.Context, clusterState } selectedNode := nodePicker.SelectNodes(skyhook) + serialStop := false for _, node := range selectedNode { // Skip nodes that are waiting on higher-priority skyhooks @@ -1474,6 +1504,11 @@ func (r *SkyhookReconciler) RunSkyhookPackages(ctx context.Context, clusterState continue } + // The stale-annotation clear for nodes with no runnable interrupt-requiring + // package now lives in UpdateDrainBlockedCondition (see cluster_state_v2.go), + // which runs on every pass — paused, disabled, complete, and error exits + // included — rather than only the passes that reach this loop. + toRun, err := node.RunNext() if err != nil { return nil, fmt.Errorf("error getting next packages to run: %w", err) @@ -1496,14 +1531,7 @@ func (r *SkyhookReconciler) RunSkyhookPackages(ctx context.Context, clusterState return nil, fmt.Errorf("node %s: reading state while filtering runnable packages: %w", node.GetNode().Name, err) } - filtered := make([]*v1alpha1.Package, 0, len(toRun)) - for _, pkg := range toRun { - if shouldSkipApplyForUninstall(pkg, nodeState, beingDeleted) { - continue - } - filtered = append(filtered, pkg) - } - toRun = filtered + toRun = filterApplicablePackages(toRun, nodeState, beingDeleted) // prepend the uninstall packages so they are ran first. // filterUninstallForNode drops entries that aren't in this node's @@ -1517,8 +1545,10 @@ func (r *SkyhookReconciler) RunSkyhookPackages(ctx context.Context, clusterState ok, err := r.ProcessInterrupt(ctx, node, f, interrupt, interrupt != nil && f.Name == pack) if err != nil { - // TODO: error handle - return nil, fmt.Errorf("error processing if we should interrupt [%s:%s]: %w", f.Name, f.Version, err) + // ProcessInterrupt may have already recorded drain blockers in memory before + // erroring; saveThenWrap persists them so they are not lost. + return nil, r.saveThenWrap(ctx, clusterState, skyhook, + fmt.Errorf("error processing if we should interrupt [%s:%s]: %w", f.Name, f.Version, err)) } if !ok { requeue = true @@ -1527,16 +1557,30 @@ func (r *SkyhookReconciler) RunSkyhookPackages(ctx context.Context, clusterState err = r.ApplyPackage(ctx, logger, clusterState, node, f, interrupt != nil && f.Name == pack) if err != nil { - return nil, fmt.Errorf("error applying package [%s:%s]: %w", f.Name, f.Version, err) + return nil, r.saveThenWrap(ctx, clusterState, skyhook, + fmt.Errorf("error applying package [%s:%s]: %w", f.Name, f.Version, err)) } // process one package at a time if skyhook.GetSkyhook().Spec.Serial { - return &ctrl.Result{RequeueAfter: time.Second * 2}, nil + serialStop = true + break } } + + if serialStop { + break + } } + // Re-run here too (also done unconditionally in refreshSkyhookConditions) so the + // condition reflects this pass's fresh findings immediately rather than waiting one + // more reconcile for refreshSkyhookConditions to pick them up. Rebuilt from persisted + // per-node state (see UpdateDrainBlockedCondition), including its own truncation-log + // line, rather than from a local slice — that is what keeps it correct for nodes this + // pass skipped or never reached. + skyhook.UpdateDrainBlockedCondition(ctx, logger) + saved, errs := r.SaveNodesAndSkyhook(ctx, clusterState, skyhook) if len(errs) > 0 { return &ctrl.Result{}, utilerrors.NewAggregate(errs) @@ -1545,7 +1589,7 @@ func (r *SkyhookReconciler) RunSkyhookPackages(ctx context.Context, clusterState requeue = true } - if !skyhook.IsComplete() || requeue { + if serialStop || !skyhook.IsComplete() || requeue { return &ctrl.Result{RequeueAfter: time.Second * 2}, nil // not sure this is better then just requeue bool } @@ -2618,19 +2662,19 @@ func (r *SkyhookReconciler) HasRunningPackages(ctx context.Context, skyhookNode return false, nil } -func (r *SkyhookReconciler) DrainNode(ctx context.Context, skyhookNode wrapper.SkyhookNode, _package *v1alpha1.Package) (bool, error) { +func (r *SkyhookReconciler) DrainNode(ctx context.Context, skyhookNode wrapper.SkyhookNode, _package *v1alpha1.Package) (drain.DrainResult, error) { drained, err := r.IsDrained(ctx, skyhookNode) if err != nil { - return false, err + return drain.DrainResult{}, err } if drained { skyhookNode.ClearDrainStart() - return true, nil + return drain.DrainResult{Ready: true}, nil } drainStartedAt, err := skyhookNode.DrainStartedAt() if err != nil { - return false, fmt.Errorf("error reading drain start for node [%s]: %w", skyhookNode.GetNode().Name, err) + return drain.DrainResult{}, fmt.Errorf("error reading drain start for node [%s]: %w", skyhookNode.GetNode().Name, err) } drainConfig := skyhookNode.GetSkyhook().Spec.DrainConfig @@ -2655,18 +2699,25 @@ func (r *SkyhookReconciler) DrainNode(ctx context.Context, skyhookNode wrapper.S _package.Version, ) skyhookNode.SetStatus(v1alpha1.StatusErroring) - return false, nil + // Preserve whatever blockers were last recorded rather than clearing them: this is + // the moment a stuck drain becomes a user-visible DrainTimeout error, so the + // condition should keep explaining what was blocking it, not go silent. + lastBlocked, blockedErr := skyhookNode.DrainBlocked() + if blockedErr != nil { + return drain.DrainResult{}, blockedErr + } + return drain.DrainResult{Blocked: lastBlocked}, nil } pods, err := r.dal.GetPods(ctx, client.MatchingFields{ fieldSelectorNodeName: skyhookNode.GetNode().Name, }) if err != nil { - return false, err + return drain.DrainResult{}, err } if pods == nil || len(pods.Items) == 0 { - return true, nil + return drain.DrainResult{Ready: true}, nil } r.recorder.Eventf(skyhookNode.GetNode(), nil, EventTypeNormal, EventsReasonSkyhookDrain, "DrainNode", @@ -2680,18 +2731,35 @@ func (r *SkyhookReconciler) DrainNode(ctx context.Context, skyhookNode wrapper.S options := drain.OptionsFromConfig(skyhookNode.GetSkyhook().Spec.DrainConfig) options.PackageNamespace = r.opts.Namespace errs := make([]error, 0) + blocked := make([]drain.BlockedPod, 0) waitingForPods := false for _, pod := range pods.Items { decision := drain.DecidePod(&pod, options) switch decision.Action { case drain.ActionBlock: waitingForPods = true + if reason, ok := blockReasonFromDrainReason(decision.Reason); ok { + blocked = append(blocked, drain.BlockedPod{ + Namespace: pod.Namespace, + Name: pod.Name, + Reason: reason, + }) + } case drain.ActionEvict: waitingForPods = true eviction := policyv1.Eviction{DeleteOptions: options.EvictionDeleteOptions()} err := r.Client.SubResource("eviction").Create(ctx, &pod, &eviction) if err != nil { - errs = append(errs, fmt.Errorf("error evicting pod [%s:%s]: %w", pod.Namespace, pod.Name, err)) + if reason, detail, ok := classifyEvictionRejection(err); ok { + blocked = append(blocked, drain.BlockedPod{ + Namespace: pod.Namespace, + Name: pod.Name, + Reason: reason, + Detail: detail, + }) + } else { + errs = append(errs, fmt.Errorf("error evicting pod [%s:%s]: %w", pod.Namespace, pod.Name, err)) + } } case drain.ActionDelete: waitingForPods = true @@ -2703,10 +2771,44 @@ func (r *SkyhookReconciler) DrainNode(ctx context.Context, skyhookNode wrapper.S } if len(errs) > 0 { - return false, utilerrors.NewAggregate(errs) + return drain.DrainResult{Blocked: blocked}, utilerrors.NewAggregate(errs) + } + + return drain.DrainResult{Ready: !waitingForPods, Blocked: blocked}, nil +} + +// classifyEvictionRejection inspects a failed eviction create and reports whether it is a +// PodDisruptionBudget rejection (a self-resolving wait state) rather than a genuine error. +// The PDB cause message is copied verbatim — apiserver-generated prose, not a stable contract. +func classifyEvictionRejection(err error) (drain.BlockReason, string, bool) { + var statusErr *apierrors.StatusError + if !errors.As(err, &statusErr) || !apierrors.IsTooManyRequests(statusErr) { + return "", "", false + } + details := statusErr.ErrStatus.Details + if details == nil { + return "", "", false + } + for _, cause := range details.Causes { + if cause.Type == policyv1.DisruptionBudgetCause { + return drain.BlockReasonPodDisruptionBudget, cause.Message, true + } } + return "", "", false +} - return !waitingForPods, nil +// blockReasonFromDrainReason maps a drain.Decision reason to the DrainBlocked condition's +// taxonomy. ReasonTerminating is deliberately excluded: an already-terminating pod is not a +// blocker to report, just one drain is still waiting to finish evicting. +func blockReasonFromDrainReason(reason string) (drain.BlockReason, bool) { + switch reason { + case drain.ReasonUnmanaged: + return drain.BlockReasonUnmanagedPod, true + case drain.ReasonEmptyDir: + return drain.BlockReasonEmptyDirData, true + default: + return "", false + } } // Interrupt should not be called unless safe to do so, IE already cordoned and drained @@ -3433,15 +3535,31 @@ func (r *SkyhookReconciler) EnsureNodeIsReadyForInterrupt(ctx context.Context, s _package.Version, skyhookNode.GetSkyhook().Name, ) + // We have not reached DrainNode this pass, so we don't know whether any + // previously-recorded PDB/unmanaged/emptyDir blockers still apply. A stale + // blocker naming a pod that no longer holds the drain is worse than reporting + // none — non-interrupt work has its own accurate signal in + // updateDrainBlockedCondition's Blocked/NonInterruptPodsRunning condition. + if err := skyhookNode.SetDrainBlocked(nil); err != nil { + return false, err + } return false, nil } - ready, err := r.DrainNode(ctx, skyhookNode, _package) + result, err := r.DrainNode(ctx, skyhookNode, _package) + // Persist regardless of err: this is what makes DrainBlocked level-triggered rather + // than dependent on this pass reaching UpdateDrainBlockedCondition later. See + // SetDrainBlocked's doc comment. The mutation only reaches the apiserver once + // SaveNodesAndSkyhook runs — see RunSkyhookPackages' error-path handling for why + // callers here must not simply return before that happens. + if setErr := skyhookNode.SetDrainBlocked(result.Blocked); setErr != nil && err == nil { + err = setErr + } if err != nil { return false, fmt.Errorf("error draining node [%s]: %w", skyhookNode.GetNode().Name, err) } - return ready, nil + return result.Ready, nil } // ApplyPackage starts a pod on node for the package diff --git a/operator/internal/controller/skyhook_controller_test.go b/operator/internal/controller/skyhook_controller_test.go index 4c4193790..b4527b3f9 100644 --- a/operator/internal/controller/skyhook_controller_test.go +++ b/operator/internal/controller/skyhook_controller_test.go @@ -616,22 +616,22 @@ var _ = Describe("skyhook controller tests", func() { skyhookNode, err := wrapper.NewSkyhookNode(node, skyhook) Expect(err).ToNot(HaveOccurred()) - drained, err := r.DrainNode(ctx, skyhookNode, &v1alpha1.Package{ + result, err := r.DrainNode(ctx, skyhookNode, &v1alpha1.Package{ PackageRef: v1alpha1.PackageRef{Name: "pkg", Version: "1.0.0"}, }) Expect(err).ToNot(HaveOccurred()) - Expect(drained).To(BeFalse()) + Expect(result.Ready).To(BeFalse()) Expect(gracePeriodSeconds).To(Equal(int64(7))) deletedPod := &corev1.Pod{} err = testClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "workload"}, deletedPod) Expect(apierrors.IsNotFound(err)).To(BeTrue()) - drained, err = r.DrainNode(ctx, skyhookNode, &v1alpha1.Package{ + result, err = r.DrainNode(ctx, skyhookNode, &v1alpha1.Package{ PackageRef: v1alpha1.PackageRef{Name: "pkg", Version: "1.0.0"}, }) Expect(err).ToNot(HaveOccurred()) - Expect(drained).To(BeTrue()) + Expect(result.Ready).To(BeTrue()) }) It("should not report drained while an evicted pod is still terminating", func() { @@ -688,9 +688,9 @@ var _ = Describe("skyhook controller tests", func() { PackageRef: v1alpha1.PackageRef{Name: "pkg", Version: "1.0.0"}, } - drained, err := r.DrainNode(ctx, skyhookNode, _package) + result, err := r.DrainNode(ctx, skyhookNode, _package) Expect(err).ToNot(HaveOccurred()) - Expect(drained).To(BeFalse()) + Expect(result.Ready).To(BeFalse()) Expect(deleteCount).To(Equal(1)) terminating := &corev1.Pod{} @@ -701,18 +701,18 @@ var _ = Describe("skyhook controller tests", func() { Expect(err).ToNot(HaveOccurred()) Expect(isDrained).To(BeFalse()) - drained, err = r.DrainNode(ctx, skyhookNode, _package) + result, err = r.DrainNode(ctx, skyhookNode, _package) Expect(err).ToNot(HaveOccurred()) - Expect(drained).To(BeFalse()) + Expect(result.Ready).To(BeFalse()) Expect(deleteCount).To(Equal(1)) terminating.Finalizers = nil Expect(testClient.Update(ctx, terminating)).To(Succeed()) Expect(apierrors.IsNotFound(testClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "workload"}, &corev1.Pod{}))).To(BeTrue()) - drained, err = r.DrainNode(ctx, skyhookNode, _package) + result, err = r.DrainNode(ctx, skyhookNode, _package) Expect(err).ToNot(HaveOccurred()) - Expect(drained).To(BeTrue()) + Expect(result.Ready).To(BeTrue()) }) It("should wait without deleting unmanaged pods when force is false", func() { @@ -761,11 +761,11 @@ var _ = Describe("skyhook controller tests", func() { skyhookNode, err := wrapper.NewSkyhookNode(node, skyhook) Expect(err).ToNot(HaveOccurred()) - drained, err := r.DrainNode(ctx, skyhookNode, &v1alpha1.Package{ + result, err := r.DrainNode(ctx, skyhookNode, &v1alpha1.Package{ PackageRef: v1alpha1.PackageRef{Name: "pkg", Version: "1.0.0"}, }) Expect(err).ToNot(HaveOccurred()) - Expect(drained).To(BeFalse()) + Expect(result.Ready).To(BeFalse()) Expect(deleteCalled).To(BeFalse()) Expect(evictCalled).To(BeFalse()) Expect(skyhookNode.Status()).To(Equal(v1alpha1.StatusInProgress)) @@ -1110,11 +1110,11 @@ var _ = Describe("skyhook controller tests", func() { skyhookNode, err := wrapper.NewSkyhookNode(node, skyhook) Expect(err).ToNot(HaveOccurred()) - drained, err := r.DrainNode(ctx, skyhookNode, &v1alpha1.Package{ + result, err := r.DrainNode(ctx, skyhookNode, &v1alpha1.Package{ PackageRef: v1alpha1.PackageRef{Name: "pkg", Version: "1.0.0"}, }) Expect(err).ToNot(HaveOccurred()) - Expect(drained).To(BeFalse()) + Expect(result.Ready).To(BeFalse()) Expect(deleteCalled).To(BeFalse()) Expect(skyhookNode.Status()).To(Equal(v1alpha1.StatusErroring)) }) @@ -1174,11 +1174,11 @@ var _ = Describe("skyhook controller tests", func() { skyhookNode, err := wrapper.NewSkyhookNode(node, skyhook) Expect(err).ToNot(HaveOccurred()) - drained, err := r.DrainNode(ctx, skyhookNode, &v1alpha1.Package{ + result, err := r.DrainNode(ctx, skyhookNode, &v1alpha1.Package{ PackageRef: v1alpha1.PackageRef{Name: "pkg", Version: "1.0.0"}, }) Expect(err).ToNot(HaveOccurred()) - Expect(drained).To(BeFalse()) + Expect(result.Ready).To(BeFalse()) Expect(skyhookNode.Status()).To(Equal(v1alpha1.StatusErroring)) Eventually(recorder.Events).Should(Receive(ContainSubstring("Warning Drain drain timed out after [1s] for node [node-a] package [pkg:1.0.0] from [nodewright:drain-timeout]"))) Eventually(recorder.Events).Should(Receive(ContainSubstring("Warning Drain drain timed out after [1s] for node [node-a] package [pkg:1.0.0]"))) diff --git a/operator/internal/drain/drain.go b/operator/internal/drain/drain.go index a0a4d2291..e9d18bb8a 100644 --- a/operator/internal/drain/drain.go +++ b/operator/internal/drain/drain.go @@ -67,6 +67,42 @@ type Options struct { PackageNamespace string } +type BlockReason string + +const ( + BlockReasonPodDisruptionBudget BlockReason = "PodDisruptionBudget" + BlockReasonUnmanagedPod BlockReason = "UnmanagedPod" + BlockReasonEmptyDirData BlockReason = "EmptyDirData" +) + +// BlockedPod is one pod currently preventing a node's drain from completing, +// with enough context to render a DrainBlocked condition message. +type BlockedPod struct { + Namespace string `json:"namespace"` + Name string `json:"name"` + Reason BlockReason `json:"reason"` + // Detail is apiserver-generated prose (e.g. the PDB cause message) and is + // copied verbatim — it is not a stable contract, so never parsed. + Detail string `json:"detail,omitempty"` +} + +// DrainResult is what DrainNode reports back: whether the node is fully +// drained, and — if not — which pods are blocking it and why. +type DrainResult struct { + Ready bool + Blocked []BlockedPod +} + +// IsZero reports whether the result carries no information: DrainNode returns this +// alongside an error on paths that never got far enough to evaluate any pod (a +// GetPods failure, for instance), as distinct from a result that legitimately found +// nothing blocking. The caller uses this to avoid persisting a false "nothing +// blocking" over a real, previously-recorded blocker just because this pass +// couldn't look. +func (r DrainResult) IsZero() bool { + return !r.Ready && len(r.Blocked) == 0 +} + func DefaultOptions() Options { return Options{ DeleteEmptyDirData: true, diff --git a/operator/internal/wrapper/mock/SkyhookNode.go b/operator/internal/wrapper/mock/SkyhookNode.go index 4cadb293b..b7472980a 100644 --- a/operator/internal/wrapper/mock/SkyhookNode.go +++ b/operator/internal/wrapper/mock/SkyhookNode.go @@ -24,11 +24,12 @@ package wrapper import ( "github.com/NVIDIA/nodewright/operator/api/nodewright/v1alpha1" + "github.com/NVIDIA/nodewright/operator/internal/drain" "github.com/NVIDIA/nodewright/operator/internal/wrapper" "github.com/go-logr/logr" mock "github.com/stretchr/testify/mock" v10 "k8s.io/api/core/v1" - "k8s.io/apimachinery/pkg/apis/meta/v1" + v1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) // NewMockSkyhookNode creates a new instance of MockSkyhookNode. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. @@ -37,10 +38,19 @@ func NewMockSkyhookNode(t interface { mock.TestingT Cleanup(func()) }) *MockSkyhookNode { + if helper, ok := t.(interface{ Helper() }); ok { + helper.Helper() + } + mock := &MockSkyhookNode{} mock.Mock.Test(t) - t.Cleanup(func() { mock.AssertExpectations(t) }) + t.Cleanup(func() { + if helper, ok := t.(interface{ Helper() }); ok { + helper.Helper() + } + mock.AssertExpectations(t) + }) return mock } @@ -212,6 +222,61 @@ func (_c *MockSkyhookNode_Cordon_Call) RunAndReturn(run func() bool) *MockSkyhoo return _c } +// DrainBlocked provides a mock function for the type MockSkyhookNode +func (_mock *MockSkyhookNode) DrainBlocked() ([]drain.BlockedPod, error) { + ret := _mock.Called() + + if len(ret) == 0 { + panic("no return value specified for DrainBlocked") + } + + var r0 []drain.BlockedPod + var r1 error + if returnFunc, ok := ret.Get(0).(func() ([]drain.BlockedPod, error)); ok { + return returnFunc() + } + if returnFunc, ok := ret.Get(0).(func() []drain.BlockedPod); ok { + r0 = returnFunc() + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).([]drain.BlockedPod) + } + } + if returnFunc, ok := ret.Get(1).(func() error); ok { + r1 = returnFunc() + } else { + r1 = ret.Error(1) + } + return r0, r1 +} + +// MockSkyhookNode_DrainBlocked_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'DrainBlocked' +type MockSkyhookNode_DrainBlocked_Call struct { + *mock.Call +} + +// DrainBlocked is a helper method to define mock.On call +func (_e *MockSkyhookNode_Expecter) DrainBlocked() *MockSkyhookNode_DrainBlocked_Call { + return &MockSkyhookNode_DrainBlocked_Call{Call: _e.mock.On("DrainBlocked")} +} + +func (_c *MockSkyhookNode_DrainBlocked_Call) Run(run func()) *MockSkyhookNode_DrainBlocked_Call { + _c.Call.Run(func(args mock.Arguments) { + run() + }) + return _c +} + +func (_c *MockSkyhookNode_DrainBlocked_Call) Return(blockedPods []drain.BlockedPod, err error) *MockSkyhookNode_DrainBlocked_Call { + _c.Call.Return(blockedPods, err) + return _c +} + +func (_c *MockSkyhookNode_DrainBlocked_Call) RunAndReturn(run func() ([]drain.BlockedPod, error)) *MockSkyhookNode_DrainBlocked_Call { + _c.Call.Return(run) + return _c +} + // DrainStartedAt provides a mock function for the type MockSkyhookNode func (_mock *MockSkyhookNode) DrainStartedAt() (*v1.Time, error) { ret := _mock.Called() @@ -473,7 +538,7 @@ type MockSkyhookNode_HasInterrupt_Call struct { // HasInterrupt is a helper method to define mock.On call // - _package v1alpha1.Package -func (_e *MockSkyhookNode_Expecter) HasInterrupt(_package interface{}) *MockSkyhookNode_HasInterrupt_Call { +func (_e *MockSkyhookNode_Expecter) HasInterrupt(_package any) *MockSkyhookNode_HasInterrupt_Call { return &MockSkyhookNode_HasInterrupt_Call{Call: _e.mock.On("HasInterrupt", _package)} } @@ -612,7 +677,7 @@ type MockSkyhookNode_IsPackageComplete_Call struct { // IsPackageComplete is a helper method to define mock.On call // - _package v1alpha1.Package -func (_e *MockSkyhookNode_Expecter) IsPackageComplete(_package interface{}) *MockSkyhookNode_IsPackageComplete_Call { +func (_e *MockSkyhookNode_Expecter) IsPackageComplete(_package any) *MockSkyhookNode_IsPackageComplete_Call { return &MockSkyhookNode_IsPackageComplete_Call{Call: _e.mock.On("IsPackageComplete", _package)} } @@ -663,7 +728,7 @@ type MockSkyhookNode_Migrate_Call struct { // Migrate is a helper method to define mock.On call // - logger logr.Logger -func (_e *MockSkyhookNode_Expecter) Migrate(logger interface{}) *MockSkyhookNode_Migrate_Call { +func (_e *MockSkyhookNode_Expecter) Migrate(logger any) *MockSkyhookNode_Migrate_Call { return &MockSkyhookNode_Migrate_Call{Call: _e.mock.On("Migrate", logger)} } @@ -716,7 +781,7 @@ type MockSkyhookNode_NextStage_Call struct { // NextStage is a helper method to define mock.On call // - _package *v1alpha1.Package -func (_e *MockSkyhookNode_Expecter) NextStage(_package interface{}) *MockSkyhookNode_NextStage_Call { +func (_e *MockSkyhookNode_Expecter) NextStage(_package any) *MockSkyhookNode_NextStage_Call { return &MockSkyhookNode_NextStage_Call{Call: _e.mock.On("NextStage", _package)} } @@ -778,7 +843,7 @@ type MockSkyhookNode_PackageStatus_Call struct { // PackageStatus is a helper method to define mock.On call // - name string -func (_e *MockSkyhookNode_Expecter) PackageStatus(name interface{}) *MockSkyhookNode_PackageStatus_Call { +func (_e *MockSkyhookNode_Expecter) PackageStatus(name any) *MockSkyhookNode_PackageStatus_Call { return &MockSkyhookNode_PackageStatus_Call{Call: _e.mock.On("PackageStatus", name)} } @@ -961,7 +1026,7 @@ type MockSkyhookNode_RemoveState_Call struct { // RemoveState is a helper method to define mock.On call // - _package v1alpha1.PackageRef -func (_e *MockSkyhookNode_Expecter) RemoveState(_package interface{}) *MockSkyhookNode_RemoveState_Call { +func (_e *MockSkyhookNode_Expecter) RemoveState(_package any) *MockSkyhookNode_RemoveState_Call { return &MockSkyhookNode_RemoveState_Call{Call: _e.mock.On("RemoveState", _package)} } @@ -1001,7 +1066,7 @@ type MockSkyhookNode_RemoveTaint_Call struct { // RemoveTaint is a helper method to define mock.On call // - key string -func (_e *MockSkyhookNode_Expecter) RemoveTaint(key interface{}) *MockSkyhookNode_RemoveTaint_Call { +func (_e *MockSkyhookNode_Expecter) RemoveTaint(key any) *MockSkyhookNode_RemoveTaint_Call { return &MockSkyhookNode_RemoveTaint_Call{Call: _e.mock.On("RemoveTaint", key)} } @@ -1116,6 +1181,57 @@ func (_c *MockSkyhookNode_RunNext_Call) RunAndReturn(run func() ([]*v1alpha1.Pac return _c } +// SetDrainBlocked provides a mock function for the type MockSkyhookNode +func (_mock *MockSkyhookNode) SetDrainBlocked(blocked []drain.BlockedPod) error { + ret := _mock.Called(blocked) + + if len(ret) == 0 { + panic("no return value specified for SetDrainBlocked") + } + + var r0 error + if returnFunc, ok := ret.Get(0).(func([]drain.BlockedPod) error); ok { + r0 = returnFunc(blocked) + } else { + r0 = ret.Error(0) + } + return r0 +} + +// MockSkyhookNode_SetDrainBlocked_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'SetDrainBlocked' +type MockSkyhookNode_SetDrainBlocked_Call struct { + *mock.Call +} + +// SetDrainBlocked is a helper method to define mock.On call +// - blocked []drain.BlockedPod +func (_e *MockSkyhookNode_Expecter) SetDrainBlocked(blocked any) *MockSkyhookNode_SetDrainBlocked_Call { + return &MockSkyhookNode_SetDrainBlocked_Call{Call: _e.mock.On("SetDrainBlocked", blocked)} +} + +func (_c *MockSkyhookNode_SetDrainBlocked_Call) Run(run func(blocked []drain.BlockedPod)) *MockSkyhookNode_SetDrainBlocked_Call { + _c.Call.Run(func(args mock.Arguments) { + var arg0 []drain.BlockedPod + if args[0] != nil { + arg0 = args[0].([]drain.BlockedPod) + } + run( + arg0, + ) + }) + return _c +} + +func (_c *MockSkyhookNode_SetDrainBlocked_Call) Return(err error) *MockSkyhookNode_SetDrainBlocked_Call { + _c.Call.Return(err) + return _c +} + +func (_c *MockSkyhookNode_SetDrainBlocked_Call) RunAndReturn(run func([]drain.BlockedPod) error) *MockSkyhookNode_SetDrainBlocked_Call { + _c.Call.Return(run) + return _c +} + // SetState provides a mock function for the type MockSkyhookNode func (_mock *MockSkyhookNode) SetState(state v1alpha1.NodeState) error { ret := _mock.Called(state) @@ -1140,7 +1256,7 @@ type MockSkyhookNode_SetState_Call struct { // SetState is a helper method to define mock.On call // - state v1alpha1.NodeState -func (_e *MockSkyhookNode_Expecter) SetState(state interface{}) *MockSkyhookNode_SetState_Call { +func (_e *MockSkyhookNode_Expecter) SetState(state any) *MockSkyhookNode_SetState_Call { return &MockSkyhookNode_SetState_Call{Call: _e.mock.On("SetState", state)} } @@ -1180,7 +1296,7 @@ type MockSkyhookNode_SetStatus_Call struct { // SetStatus is a helper method to define mock.On call // - status v1alpha1.Status -func (_e *MockSkyhookNode_Expecter) SetStatus(status interface{}) *MockSkyhookNode_SetStatus_Call { +func (_e *MockSkyhookNode_Expecter) SetStatus(status any) *MockSkyhookNode_SetStatus_Call { return &MockSkyhookNode_SetStatus_Call{Call: _e.mock.On("SetStatus", status)} } @@ -1253,7 +1369,7 @@ type MockSkyhookNode_StartDrain_Call struct { // StartDrain is a helper method to define mock.On call // - startedAt v1.Time -func (_e *MockSkyhookNode_Expecter) StartDrain(startedAt interface{}) *MockSkyhookNode_StartDrain_Call { +func (_e *MockSkyhookNode_Expecter) StartDrain(startedAt any) *MockSkyhookNode_StartDrain_Call { return &MockSkyhookNode_StartDrain_Call{Call: _e.mock.On("StartDrain", startedAt)} } @@ -1392,7 +1508,7 @@ type MockSkyhookNode_Taint_Call struct { // Taint is a helper method to define mock.On call // - key string -func (_e *MockSkyhookNode_Expecter) Taint(key interface{}) *MockSkyhookNode_Taint_Call { +func (_e *MockSkyhookNode_Expecter) Taint(key any) *MockSkyhookNode_Taint_Call { return &MockSkyhookNode_Taint_Call{Call: _e.mock.On("Taint", key)} } @@ -1514,7 +1630,7 @@ type MockSkyhookNode_Upsert_Call struct { // - stage v1alpha1.Stage // - restarts int32 // - containerSHA string -func (_e *MockSkyhookNode_Expecter) Upsert(_package interface{}, image interface{}, state interface{}, stage interface{}, restarts interface{}, containerSHA interface{}) *MockSkyhookNode_Upsert_Call { +func (_e *MockSkyhookNode_Expecter) Upsert(_package any, image any, state any, stage any, restarts any, containerSHA any) *MockSkyhookNode_Upsert_Call { return &MockSkyhookNode_Upsert_Call{Call: _e.mock.On("Upsert", _package, image, state, stage, restarts, containerSHA)} } diff --git a/operator/internal/wrapper/mock/SkyhookNodeOnly.go b/operator/internal/wrapper/mock/SkyhookNodeOnly.go index 22344f013..6fa78cae8 100644 --- a/operator/internal/wrapper/mock/SkyhookNodeOnly.go +++ b/operator/internal/wrapper/mock/SkyhookNodeOnly.go @@ -24,10 +24,11 @@ package wrapper import ( "github.com/NVIDIA/nodewright/operator/api/nodewright/v1alpha1" + "github.com/NVIDIA/nodewright/operator/internal/drain" "github.com/go-logr/logr" mock "github.com/stretchr/testify/mock" v10 "k8s.io/api/core/v1" - "k8s.io/apimachinery/pkg/apis/meta/v1" + v1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) // NewMockSkyhookNodeOnly creates a new instance of MockSkyhookNodeOnly. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. @@ -36,10 +37,19 @@ func NewMockSkyhookNodeOnly(t interface { mock.TestingT Cleanup(func()) }) *MockSkyhookNodeOnly { + if helper, ok := t.(interface{ Helper() }); ok { + helper.Helper() + } + mock := &MockSkyhookNodeOnly{} mock.Mock.Test(t) - t.Cleanup(func() { mock.AssertExpectations(t) }) + t.Cleanup(func() { + if helper, ok := t.(interface{ Helper() }); ok { + helper.Helper() + } + mock.AssertExpectations(t) + }) return mock } @@ -178,6 +188,61 @@ func (_c *MockSkyhookNodeOnly_Cordon_Call) RunAndReturn(run func() bool) *MockSk return _c } +// DrainBlocked provides a mock function for the type MockSkyhookNodeOnly +func (_mock *MockSkyhookNodeOnly) DrainBlocked() ([]drain.BlockedPod, error) { + ret := _mock.Called() + + if len(ret) == 0 { + panic("no return value specified for DrainBlocked") + } + + var r0 []drain.BlockedPod + var r1 error + if returnFunc, ok := ret.Get(0).(func() ([]drain.BlockedPod, error)); ok { + return returnFunc() + } + if returnFunc, ok := ret.Get(0).(func() []drain.BlockedPod); ok { + r0 = returnFunc() + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).([]drain.BlockedPod) + } + } + if returnFunc, ok := ret.Get(1).(func() error); ok { + r1 = returnFunc() + } else { + r1 = ret.Error(1) + } + return r0, r1 +} + +// MockSkyhookNodeOnly_DrainBlocked_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'DrainBlocked' +type MockSkyhookNodeOnly_DrainBlocked_Call struct { + *mock.Call +} + +// DrainBlocked is a helper method to define mock.On call +func (_e *MockSkyhookNodeOnly_Expecter) DrainBlocked() *MockSkyhookNodeOnly_DrainBlocked_Call { + return &MockSkyhookNodeOnly_DrainBlocked_Call{Call: _e.mock.On("DrainBlocked")} +} + +func (_c *MockSkyhookNodeOnly_DrainBlocked_Call) Run(run func()) *MockSkyhookNodeOnly_DrainBlocked_Call { + _c.Call.Run(func(args mock.Arguments) { + run() + }) + return _c +} + +func (_c *MockSkyhookNodeOnly_DrainBlocked_Call) Return(blockedPods []drain.BlockedPod, err error) *MockSkyhookNodeOnly_DrainBlocked_Call { + _c.Call.Return(blockedPods, err) + return _c +} + +func (_c *MockSkyhookNodeOnly_DrainBlocked_Call) RunAndReturn(run func() ([]drain.BlockedPod, error)) *MockSkyhookNodeOnly_DrainBlocked_Call { + _c.Call.Return(run) + return _c +} + // DrainStartedAt provides a mock function for the type MockSkyhookNodeOnly func (_mock *MockSkyhookNodeOnly) DrainStartedAt() (*v1.Time, error) { ret := _mock.Called() @@ -347,7 +412,7 @@ type MockSkyhookNodeOnly_Migrate_Call struct { // Migrate is a helper method to define mock.On call // - logger logr.Logger -func (_e *MockSkyhookNodeOnly_Expecter) Migrate(logger interface{}) *MockSkyhookNodeOnly_Migrate_Call { +func (_e *MockSkyhookNodeOnly_Expecter) Migrate(logger any) *MockSkyhookNodeOnly_Migrate_Call { return &MockSkyhookNodeOnly_Migrate_Call{Call: _e.mock.On("Migrate", logger)} } @@ -409,7 +474,7 @@ type MockSkyhookNodeOnly_PackageStatus_Call struct { // PackageStatus is a helper method to define mock.On call // - name string -func (_e *MockSkyhookNodeOnly_Expecter) PackageStatus(name interface{}) *MockSkyhookNodeOnly_PackageStatus_Call { +func (_e *MockSkyhookNodeOnly_Expecter) PackageStatus(name any) *MockSkyhookNodeOnly_PackageStatus_Call { return &MockSkyhookNodeOnly_PackageStatus_Call{Call: _e.mock.On("PackageStatus", name)} } @@ -548,7 +613,7 @@ type MockSkyhookNodeOnly_RemoveState_Call struct { // RemoveState is a helper method to define mock.On call // - _package v1alpha1.PackageRef -func (_e *MockSkyhookNodeOnly_Expecter) RemoveState(_package interface{}) *MockSkyhookNodeOnly_RemoveState_Call { +func (_e *MockSkyhookNodeOnly_Expecter) RemoveState(_package any) *MockSkyhookNodeOnly_RemoveState_Call { return &MockSkyhookNodeOnly_RemoveState_Call{Call: _e.mock.On("RemoveState", _package)} } @@ -588,7 +653,7 @@ type MockSkyhookNodeOnly_RemoveTaint_Call struct { // RemoveTaint is a helper method to define mock.On call // - key string -func (_e *MockSkyhookNodeOnly_Expecter) RemoveTaint(key interface{}) *MockSkyhookNodeOnly_RemoveTaint_Call { +func (_e *MockSkyhookNodeOnly_Expecter) RemoveTaint(key any) *MockSkyhookNodeOnly_RemoveTaint_Call { return &MockSkyhookNodeOnly_RemoveTaint_Call{Call: _e.mock.On("RemoveTaint", key)} } @@ -648,6 +713,57 @@ func (_c *MockSkyhookNodeOnly_Reset_Call) RunAndReturn(run func()) *MockSkyhookN return _c } +// SetDrainBlocked provides a mock function for the type MockSkyhookNodeOnly +func (_mock *MockSkyhookNodeOnly) SetDrainBlocked(blocked []drain.BlockedPod) error { + ret := _mock.Called(blocked) + + if len(ret) == 0 { + panic("no return value specified for SetDrainBlocked") + } + + var r0 error + if returnFunc, ok := ret.Get(0).(func([]drain.BlockedPod) error); ok { + r0 = returnFunc(blocked) + } else { + r0 = ret.Error(0) + } + return r0 +} + +// MockSkyhookNodeOnly_SetDrainBlocked_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'SetDrainBlocked' +type MockSkyhookNodeOnly_SetDrainBlocked_Call struct { + *mock.Call +} + +// SetDrainBlocked is a helper method to define mock.On call +// - blocked []drain.BlockedPod +func (_e *MockSkyhookNodeOnly_Expecter) SetDrainBlocked(blocked any) *MockSkyhookNodeOnly_SetDrainBlocked_Call { + return &MockSkyhookNodeOnly_SetDrainBlocked_Call{Call: _e.mock.On("SetDrainBlocked", blocked)} +} + +func (_c *MockSkyhookNodeOnly_SetDrainBlocked_Call) Run(run func(blocked []drain.BlockedPod)) *MockSkyhookNodeOnly_SetDrainBlocked_Call { + _c.Call.Run(func(args mock.Arguments) { + var arg0 []drain.BlockedPod + if args[0] != nil { + arg0 = args[0].([]drain.BlockedPod) + } + run( + arg0, + ) + }) + return _c +} + +func (_c *MockSkyhookNodeOnly_SetDrainBlocked_Call) Return(err error) *MockSkyhookNodeOnly_SetDrainBlocked_Call { + _c.Call.Return(err) + return _c +} + +func (_c *MockSkyhookNodeOnly_SetDrainBlocked_Call) RunAndReturn(run func([]drain.BlockedPod) error) *MockSkyhookNodeOnly_SetDrainBlocked_Call { + _c.Call.Return(run) + return _c +} + // SetState provides a mock function for the type MockSkyhookNodeOnly func (_mock *MockSkyhookNodeOnly) SetState(state v1alpha1.NodeState) error { ret := _mock.Called(state) @@ -672,7 +788,7 @@ type MockSkyhookNodeOnly_SetState_Call struct { // SetState is a helper method to define mock.On call // - state v1alpha1.NodeState -func (_e *MockSkyhookNodeOnly_Expecter) SetState(state interface{}) *MockSkyhookNodeOnly_SetState_Call { +func (_e *MockSkyhookNodeOnly_Expecter) SetState(state any) *MockSkyhookNodeOnly_SetState_Call { return &MockSkyhookNodeOnly_SetState_Call{Call: _e.mock.On("SetState", state)} } @@ -712,7 +828,7 @@ type MockSkyhookNodeOnly_SetStatus_Call struct { // SetStatus is a helper method to define mock.On call // - status v1alpha1.Status -func (_e *MockSkyhookNodeOnly_Expecter) SetStatus(status interface{}) *MockSkyhookNodeOnly_SetStatus_Call { +func (_e *MockSkyhookNodeOnly_Expecter) SetStatus(status any) *MockSkyhookNodeOnly_SetStatus_Call { return &MockSkyhookNodeOnly_SetStatus_Call{Call: _e.mock.On("SetStatus", status)} } @@ -785,7 +901,7 @@ type MockSkyhookNodeOnly_StartDrain_Call struct { // StartDrain is a helper method to define mock.On call // - startedAt v1.Time -func (_e *MockSkyhookNodeOnly_Expecter) StartDrain(startedAt interface{}) *MockSkyhookNodeOnly_StartDrain_Call { +func (_e *MockSkyhookNodeOnly_Expecter) StartDrain(startedAt any) *MockSkyhookNodeOnly_StartDrain_Call { return &MockSkyhookNodeOnly_StartDrain_Call{Call: _e.mock.On("StartDrain", startedAt)} } @@ -924,7 +1040,7 @@ type MockSkyhookNodeOnly_Taint_Call struct { // Taint is a helper method to define mock.On call // - key string -func (_e *MockSkyhookNodeOnly_Expecter) Taint(key interface{}) *MockSkyhookNodeOnly_Taint_Call { +func (_e *MockSkyhookNodeOnly_Expecter) Taint(key any) *MockSkyhookNodeOnly_Taint_Call { return &MockSkyhookNodeOnly_Taint_Call{Call: _e.mock.On("Taint", key)} } @@ -1013,7 +1129,7 @@ type MockSkyhookNodeOnly_Upsert_Call struct { // - stage v1alpha1.Stage // - restarts int32 // - containerSHA string -func (_e *MockSkyhookNodeOnly_Expecter) Upsert(_package interface{}, image interface{}, state interface{}, stage interface{}, restarts interface{}, containerSHA interface{}) *MockSkyhookNodeOnly_Upsert_Call { +func (_e *MockSkyhookNodeOnly_Expecter) Upsert(_package any, image any, state any, stage any, restarts any, containerSHA any) *MockSkyhookNodeOnly_Upsert_Call { return &MockSkyhookNodeOnly_Upsert_Call{Call: _e.mock.On("Upsert", _package, image, state, stage, restarts, containerSHA)} } diff --git a/operator/internal/wrapper/node.go b/operator/internal/wrapper/node.go index 155f11e98..eb33df810 100644 --- a/operator/internal/wrapper/node.go +++ b/operator/internal/wrapper/node.go @@ -26,6 +26,7 @@ import ( "time" "github.com/NVIDIA/nodewright/operator/api/nodewright/v1alpha1" + "github.com/NVIDIA/nodewright/operator/internal/drain" "github.com/NVIDIA/nodewright/operator/internal/graph" "github.com/NVIDIA/nodewright/operator/internal/version" "github.com/go-logr/logr" @@ -113,6 +114,12 @@ type SkyhookNodeOnly interface { DrainStartedAt() (*metav1.Time, error) // ClearDrainStart removes the drain start marker for this Skyhook on this node. ClearDrainStart() + // SetDrainBlocked persists this pass's drain-blocker findings for this Skyhook on + // this node, replacing any previously recorded findings. Nil or empty clears it. + SetDrainBlocked(blocked []drain.BlockedPod) error + // DrainBlocked returns the drain-blocker findings last persisted for this Skyhook on + // this node. + DrainBlocked() ([]drain.BlockedPod, error) // Uncordon marks the node schedulable and removes this Skyhook's cordon annotation if present. Uncordon() // Reset clears Skyhook-related state and annotations so the node can be reconfigured from scratch. @@ -198,6 +205,10 @@ func (node *skyhookNode) drainStartAnnotationKey() string { return fmt.Sprintf("%s/drainStart_%s", v1alpha1.METADATA_PREFIX, node.skyhookName) } +func (node *skyhookNode) drainBlockedAnnotationKey() string { + return fmt.Sprintf("%s/drainBlocked_%s", v1alpha1.METADATA_PREFIX, node.skyhookName) +} + // GetSkyhook returns the Skyhook associated with this node, or nil if only a name was set. func (node *skyhookNode) GetSkyhook() *Skyhook { return node.skyhook @@ -602,6 +613,71 @@ func (node *skyhookNode) ClearDrainStart() { node.updated = true } +// SetDrainBlocked persists this pass's drain-blocker findings (PDB rejections, unmanaged +// pods, emptyDir pods) for this Skyhook on this node, replacing whatever was recorded +// previously. Nil or empty blocked clears the annotation. Unlike the per-package +// nodeState annotation, this always sets the whole value: there is exactly one DrainNode +// call's findings to record per pass, never several to merge. +// +// Persisting here — at the moment the caller learns the result — rather than accumulating +// in a caller-local slice is what lets the DrainBlocked condition be level-triggered: +// UpdateDrainBlockedCondition rebuilds it from this annotation on every reconcile pass, +// including passes that skip RunSkyhookPackages entirely (paused, disabled, complete +// Skyhooks) or that return from it early (an error, or spec.serial stopping after the +// first node). Without persisting immediately, all three of those paths would either +// never update the condition or wrongly clear it for nodes the pass never reached. +func (node *skyhookNode) SetDrainBlocked(blocked []drain.BlockedPod) error { + key := node.drainBlockedAnnotationKey() + + if len(blocked) == 0 { + if node.Annotations == nil { + return nil + } + if _, ok := node.Annotations[key]; !ok { + return nil + } + delete(node.Annotations, key) + node.updated = true + return nil + } + + data, err := json.Marshal(blocked) + if err != nil { + return fmt.Errorf("error marshalling drain blocked pods: %w", err) + } + + if node.Annotations == nil { + node.Annotations = map[string]string{} + } + + if existing, ok := node.Annotations[key]; !ok || existing != string(data) { + node.Annotations[key] = string(data) + node.updated = true + } + + return nil +} + +// DrainBlocked returns the drain-blocker findings last persisted for this Skyhook on this +// node by SetDrainBlocked, or nil if none are recorded. +func (node *skyhookNode) DrainBlocked() ([]drain.BlockedPod, error) { + if node.Annotations == nil { + return nil, nil + } + + value, ok := node.Annotations[node.drainBlockedAnnotationKey()] + if !ok { + return nil, nil + } + + var blocked []drain.BlockedPod + if err := json.Unmarshal([]byte(value), &blocked); err != nil { + return nil, fmt.Errorf("error unmarshalling drain blocked pods: %w", err) + } + + return blocked, nil +} + // Uncordon marks the node schedulable and removes this Skyhook's cordon annotation if present. func (node *skyhookNode) Uncordon() { @@ -647,6 +723,7 @@ func (node *skyhookNode) Reset() { delete(node.Annotations, cordonAnnotationKey(node.skyhookName)) delete(node.Annotations, node.drainStartAnnotationKey()) + delete(node.Annotations, node.drainBlockedAnnotationKey()) delete(node.Annotations, fmt.Sprintf("%s/nodeState_%s", v1alpha1.METADATA_PREFIX, node.skyhookName)) delete(node.Annotations, fmt.Sprintf("%s/status_%s", v1alpha1.METADATA_PREFIX, node.skyhookName)) delete(node.Annotations, fmt.Sprintf("%s/version_%s", v1alpha1.METADATA_PREFIX, node.skyhookName)) diff --git a/operator/internal/wrapper/skyhook_conditions.go b/operator/internal/wrapper/skyhook_conditions.go index 9509edc8a..9413e2b40 100644 --- a/operator/internal/wrapper/skyhook_conditions.go +++ b/operator/internal/wrapper/skyhook_conditions.go @@ -20,14 +20,18 @@ package wrapper import ( "fmt" + "sort" "strings" "github.com/NVIDIA/nodewright/operator/api/nodewright/v1alpha1" + "github.com/NVIDIA/nodewright/operator/internal/drain" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) const ( // ReadyConditionNodeListLimit caps condition message fan-out to avoid etcd object bloat and excess watch bandwidth on large rollouts. + // Also used as drainBlockedDetailLineLimit's value: the two independent condition-message + // fan-out caps must move together, or a future tuning pass silently desyncs them. ReadyConditionNodeListLimit = 10 SkyhookConditionReady = "Ready" @@ -40,6 +44,9 @@ const ( SkyhookConditionUninstallFailed = "UninstallFailed" SkyhookConditionNodeStateMalformed = "NodeStateMalformed" SkyhookConditionDeletionBlocked = "DeletionBlocked" + SkyhookConditionDrainBlocked = "DrainBlocked" + + drainBlockedReasonMultiple = "MultipleCauses" SkyhookReasonNonInterruptPodsRunning = "NonInterruptPodsRunning" @@ -321,3 +328,105 @@ func FormatNodeList(nodes []string) string { } return fmt.Sprintf(" (%s)", strings.Join(nodes, ", ")) } + +// DrainBlockedNode is one node's drain blockers for the DrainBlocked condition +// message builder below. +type DrainBlockedNode struct { + NodeName string + Blocked []drain.BlockedPod +} + +// DrainBlockedConditionReason picks the condition Reason from the set of block +// reasons observed this pass. MultipleCauses covers both "one node has two kinds +// of blocker" and "different nodes are blocked for different reasons". +func DrainBlockedConditionReason(nodes []DrainBlockedNode) string { + seen := make(map[drain.BlockReason]struct{}) + for _, n := range nodes { + for _, b := range n.Blocked { + seen[b.Reason] = struct{}{} + } + } + if len(seen) != 1 { + return drainBlockedReasonMultiple + } + // Exactly one reason observed: return it directly rather than hand-mapping + // against drain.BlockReason's constants. The old switch fell through to + // MultipleCauses for any drain.BlockReason it didn't explicitly list — so a + // new reason added to that type would silently mislabel a single-cause block + // as "more than one kind of blocker present," with no compiler error to catch it. + for reason := range seen { + return string(reason) + } + return drainBlockedReasonMultiple +} + +// drainBlockedDetailLineLimit caps the number of per-pod detail lines rendered into the +// DrainBlocked message, so it stays well inside .status.conditions[].message's 32768-byte +// apiserver limit even though each Detail carries apiserver-generated prose of unbounded +// length; a genuinely pathological Detail could still overrun this line-count cap, but PDB +// cause messages are short in practice. Shares ReadyConditionNodeListLimit's value rather +// than redeclaring it, since both exist to bound condition-message fan-out for the same +// reason. The full set is always available from nodes[].Blocked for logging by the caller; +// this function stays pure and does not log. +const drainBlockedDetailLineLimit = ReadyConditionNodeListLimit + +// DrainBlockedConditionMessage renders the aggregate DrainBlocked message: a +// "N/total nodes blocked draining (names)" summary line — following the same +// truncation idiom as the Ready condition — followed by one "/ on +// : " line per blocked pod that carries a Detail (PDB +// cases only; Detail is apiserver prose and is never altered). +// +// Node and pod order are sorted rather than taken from nodes/nodes[].Blocked as given: +// node order there comes from a compartment map and pod order from the informer store, +// neither of which is stable between otherwise-identical reconcile passes. An unsorted +// message reshuffles every pass and triggers a spurious status write each time — the same +// reason the Ready condition's node lists are sorted. +func DrainBlockedConditionMessage(nodes []DrainBlockedNode, totalSelected int) string { + sorted := make([]DrainBlockedNode, len(nodes)) + copy(sorted, nodes) + sort.Slice(sorted, func(i, j int) bool { return sorted[i].NodeName < sorted[j].NodeName }) + + names := make([]string, 0, len(sorted)) + for i := range sorted { + names = append(names, sorted[i].NodeName) + blocked := make([]drain.BlockedPod, len(sorted[i].Blocked)) + copy(blocked, sorted[i].Blocked) + sort.Slice(blocked, func(a, b int) bool { + if blocked[a].Namespace != blocked[b].Namespace { + return blocked[a].Namespace < blocked[b].Namespace + } + return blocked[a].Name < blocked[b].Name + }) + sorted[i].Blocked = blocked + } + // names is already in NodeName order: it's appended while iterating sorted, which was + // sorted by NodeName above. Re-sorting here would just prove that twice. + + lines := []string{fmt.Sprintf("%d/%d nodes blocked draining%s", len(sorted), totalSelected, FormatNodeList(names))} + + detailLines := 0 + truncated := false + for _, n := range sorted { + for _, b := range n.Blocked { + detail := b.Detail + if detail == "" { + // DrainNode creates unmanaged/emptyDir blockers with a Reason but no + // apiserver-generated Detail (that's PDB-only). Fall back to the reason + // so these blockers still surface which pod is holding drain, instead + // of being silently dropped from the message. + detail = string(b.Reason) + } + if detailLines >= drainBlockedDetailLineLimit { + truncated = true + continue + } + lines = append(lines, fmt.Sprintf("%s/%s on %s: %s", b.Namespace, b.Name, n.NodeName, detail)) + detailLines++ + } + } + if truncated { + lines = append(lines, "(additional detail truncated; see controller logs)") + } + + return strings.Join(lines, "; ") +}