diff --git a/agent/README.md b/agent/README.md index 381af2d6..868cfd41 100644 --- a/agent/README.md +++ b/agent/README.md @@ -62,9 +62,10 @@ agent MODE COPY_DIR agent interrupt COPY_DIR INTERRUPT_DATA ``` -SIGTERM cancels the active step or interrupt and prevents later steps from -starting. A failed operation or runtime error exits with status 1; malformed -arguments exit with status 2. +SIGTERM lets the active step or interrupt operation run to completion and +prevents later ones from starting, so a package's `gracefulShutdown` is the +time its script has to finish. A failed operation or runtime error exits with +status 1; malformed arguments exit with status 2. `agent --version` prints the semantic version embedded in the binary at build time, falling back to the embedded Git SHA or `unknown` when build metadata is diff --git a/agent/go/internal/agent/agent_test.go b/agent/go/internal/agent/agent_test.go index d6ca7680..f367371a 100644 --- a/agent/go/internal/agent/agent_test.go +++ b/agent/go/internal/agent/agent_test.go @@ -338,7 +338,7 @@ var _ = Describe("Agent.Run", func() { Expect(stderr.String()).To(ContainSubstring("agent stopped after receiving a termination signal")) }) - It("cancels an active step and does not start the next step", func() { + It("lets an active step finish and does not start the next step", func() { root := GinkgoT().TempDir() dataDir := GinkgoT().TempDir() writeCancellationPackageFixture(dataDir) @@ -387,7 +387,10 @@ var _ = Describe("Agent.Run", func() { var exitCode ExitCode Eventually(exitCodes).WithTimeout(5 * time.Second).Should(Receive(&exitCode)) Expect(exitCode).To(Equal(ExitFailure)) - Expect(filepath.Join(root, "package", "first-finished")).NotTo(BeAnExistingFile()) + // Cancellation landed while the first step was asleep. It must have been + // allowed to finish: killing it here is the gracefulShutdown regression + // in #672. Only the step that had not started yet is refused. + Expect(filepath.Join(root, "package", "first-finished")).To(BeAnExistingFile()) Expect(filepath.Join(root, "package", "second-started")).NotTo(BeAnExistingFile()) Expect(stderr.String()).To(ContainSubstring("agent stopped after receiving a termination signal")) }) @@ -549,7 +552,7 @@ func writeCancellationPackageFixture(directory string) { stepsDir := filepath.Join(directory, "skyhook_dir") Expect(os.WriteFile( filepath.Join(stepsDir, "apply"), - []byte("#!/bin/sh\n: > \"$NODEWRIGHT_AGENT_TEST_MARKER_DIR/first-started\"\nsleep 3600\n: > \"$NODEWRIGHT_AGENT_TEST_MARKER_DIR/first-finished\"\n"), + []byte("#!/bin/sh\n: > \"$NODEWRIGHT_AGENT_TEST_MARKER_DIR/first-started\"\nsleep 1\n: > \"$NODEWRIGHT_AGENT_TEST_MARKER_DIR/first-finished\"\n"), 0o700, )).To(Succeed()) Expect(os.WriteFile( diff --git a/agent/go/internal/command/process.go b/agent/go/internal/command/process.go index ab7f646d..98ba1fed 100644 --- a/agent/go/internal/command/process.go +++ b/agent/go/internal/command/process.go @@ -30,7 +30,7 @@ import ( ) // processWaitDelay bounds cleanup when a descendant keeps an inherited pipe -// open after the command exits or context cancellation begins. +// open after the command exits. const processWaitDelay = 2 * time.Second func validateRun(ctx context.Context, command Command) error { @@ -94,7 +94,16 @@ func executeCommand( return Result{}, fmt.Errorf("running command %q: %w", command.Executable, err) } - process := exec.CommandContext(ctx, executable, command.Arguments...) + // Deliberately exec.Command, not exec.CommandContext. The context is + // cancelled on SIGTERM (cmd/agent/main.go), and binding the child to it + // would SIGKILL a step mid-host-mutation the moment the pod is told to + // stop, defeating the package's gracefulShutdown, which is documented as + // the time a script has to finish. Cancellation still refuses to *start* a + // command (the ctx.Err() check above, and the checks between steps), so + // SIGTERM means "finish what is running, then stop", as in the Python + // agent. A hung child is bounded by the kubelet at the end of the grace + // period, not here. + process := exec.Command(executable, command.Arguments...) process.Args[0] = command.Executable process.Dir = workingDirectory process.Env = process.Environ() @@ -107,18 +116,6 @@ func executeCommand( if chroot != "" { process.SysProcAttr.Chroot = chroot } - process.Cancel = func() error { - if process.Process == nil { - return os.ErrProcessDone - } - if err := syscall.Kill(-process.Process.Pid, syscall.SIGKILL); err != nil { - if errors.Is(err, syscall.ESRCH) { - return os.ErrProcessDone - } - return fmt.Errorf("killing process group %d: %w", process.Process.Pid, err) - } - return nil - } process.WaitDelay = processWaitDelay runErr := process.Run() @@ -131,9 +128,8 @@ func executeCommand( if runErr == nil { return result, nil } - if err := ctx.Err(); err != nil { - return result, fmt.Errorf("running command %q: %w", command.Executable, err) - } + // The child is never terminated on cancellation, so a failure here is its + // own and is reported as such rather than attributed to the context. var exitErr *exec.ExitError if errors.As(runErr, &exitErr) { return result, nil diff --git a/agent/go/internal/command/process_test.go b/agent/go/internal/command/process_test.go index 16dd755c..c20b45da 100644 --- a/agent/go/internal/command/process_test.go +++ b/agent/go/internal/command/process_test.go @@ -51,8 +51,8 @@ var _ = Describe("command runner process execution", func() { Expect(result.Signal).To(Equal(os.Signal(syscall.SIGTERM))) }) - It("cancels the running process group through context", func() { - ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + It("lets a running process finish after cancellation and refuses to start another", func() { + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() ready := make(chan struct{}) go func() { @@ -63,14 +63,22 @@ var _ = Describe("command runner process execution", func() { } }() + // The helper reports readiness and then keeps running, so cancelling on + // readiness lands mid-run. The process must complete on its own and + // report its real exit status; it must not be killed. + start := time.Now() result, err := NewRunner().Run(ctx, helperCommand( - "wait", + "sleep", "300", WithStdout(&readinessWriter{ready: ready}), )) + Expect(err).NotTo(HaveOccurred()) + Expect(result).To(Equal(Result{ExitCode: 0})) + Expect(time.Since(start)).To(BeNumerically(">=", 300*time.Millisecond)) + + // Only a command that has not yet started is refused. + _, err = NewRunner().Run(ctx, helperCommand("exit", "0")) Expect(errors.Is(err, context.Canceled)).To(BeTrue()) - Expect(result.ExitCode).To(Equal(SignalExitCode)) - Expect(result.Signal).To(Equal(os.Signal(syscall.SIGKILL))) }) It("propagates output writer failures", func() { @@ -167,6 +175,14 @@ func runCommandTestHelper() bool { case "wait": _, _ = io.WriteString(os.Stdout, "ready") time.Sleep(time.Hour) + case "sleep": + milliseconds, err := strconv.Atoi(values[0]) + if err != nil { + os.Exit(2) + } + _, _ = io.WriteString(os.Stdout, "ready") + time.Sleep(time.Duration(milliseconds) * time.Millisecond) + os.Exit(0) default: _, _ = fmt.Fprintf(os.Stderr, "unknown helper action %q\n", action) os.Exit(2) diff --git a/agent/go/internal/interrupts/node_restart.go b/agent/go/internal/interrupts/node_restart.go index 7f9561ae..a8647d3d 100644 --- a/agent/go/internal/interrupts/node_restart.go +++ b/agent/go/internal/interrupts/node_restart.go @@ -18,7 +18,6 @@ package interrupts import ( "context" - "errors" "fmt" "syscall" "time" @@ -51,9 +50,7 @@ func (n NodeRestart) Run(ctx context.Context, config execution.Config) (executio command.WithStderr(config.Stderr()), ) result, runErr := runner.Run(ctx, cmd) - if nodeRestartCompleted(result) && (runErr == nil || - errors.Is(runErr, context.Canceled) || - errors.Is(runErr, context.DeadlineExceeded)) { + if runErr == nil && nodeRestartCompleted(result) { return execution.StatusSuccess, nil } if runErr != nil { diff --git a/agent/go/internal/step/regular_step_test.go b/agent/go/internal/step/regular_step_test.go index fbbbc171..fa25c6dc 100644 --- a/agent/go/internal/step/regular_step_test.go +++ b/agent/go/internal/step/regular_step_test.go @@ -21,13 +21,13 @@ package step import ( "bytes" "context" - "errors" "fmt" "io" "os" "path/filepath" "strconv" "strings" + "sync" "syscall" "time" @@ -301,20 +301,39 @@ var _ = Describe("RegularStep.Run", func() { ))) }) - It("cancels an in-flight step", func() { + It("lets an in-flight step finish after cancellation", func() { stepRoot, skyhookDir, executable := prepareStepTestExecutable() value := NewRegularStep( filepath.Base(executable), WithOnHost(false), - WithArguments([]string{"-test.run=^TestStep$", "--", "wait"}), + WithArguments([]string{"-test.run=^TestStep$", "--", "sleep", "300"}), ) - ctx, cancel := context.WithCancel(context.Background()) - time.AfterFunc(100*time.Millisecond, cancel) - - status, err := value.Run(ctx, newStepRunConfig(stepRoot, skyhookDir)) + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + ready := make(chan struct{}) + go func() { + select { + case <-ready: + cancel() + case <-ctx.Done(): + } + }() + + // Cancelling on the helper's readiness write, rather than on a timer, + // guarantees the cancellation lands while the step is running and not + // before it starts, where runStep's ctx.Err() check would refuse it. + // The step must complete on its own and report its real outcome; only + // a step that has not started yet is refused (covered by runSteps tests). + start := time.Now() + status, err := value.Run(ctx, newStepRunConfig( + stepRoot, + skyhookDir, + execution.WithRunOutput(&readinessWriter{ready: ready}, io.Discard), + )) - Expect(status).To(Equal(execution.StatusFailed)) - Expect(errors.Is(err, context.Canceled)).To(BeTrue()) + Expect(err).NotTo(HaveOccurred()) + Expect(status).To(Equal(execution.StatusSuccess)) + Expect(time.Since(start)).To(BeNumerically(">=", 300*time.Millisecond)) }) }) @@ -357,6 +376,16 @@ var _ = Describe("RegularStep.WithVersions", func() { }) }) +type readinessWriter struct { + ready chan struct{} + once sync.Once +} + +func (writer *readinessWriter) Write(data []byte) (int, error) { + writer.once.Do(func() { close(writer.ready) }) + return len(data), nil +} + func newStepRunConfig(stepRoot, skyhookDir string, options ...execution.Option) execution.Config { options = append([]execution.Option{ execution.WithRootMount("/"), @@ -418,6 +447,14 @@ func runStepTestHelper() bool { time.Sleep(time.Hour) case "wait": time.Sleep(time.Hour) + case "sleep": + milliseconds, err := strconv.Atoi(values[0]) + if err != nil { + os.Exit(2) + } + _, _ = io.WriteString(os.Stdout, "ready") + time.Sleep(time.Duration(milliseconds) * time.Millisecond) + os.Exit(0) case "versions": _, _ = fmt.Fprintf( os.Stdout, diff --git a/docs/contributing/ci-test-pools.md b/docs/contributing/ci-test-pools.md index 0f559cc9..cb5010d3 100644 --- a/docs/contributing/ci-test-pools.md +++ b/docs/contributing/ci-test-pools.md @@ -55,9 +55,11 @@ the `agentless` package image, so the only one that proves a package's scripts r It covers `apply`, `config`, `interrupt`, `post-interrupt` and `uninstall`, each with its `-check` counterpart, plus log retention, `SKYHOOK_AGENT_WRITE_LOGS=false`, and the `check_results` / -`_ALL_CHECKED` summary artifacts. `upgrade` is **not** covered: the `shellscript` package -these scenarios use declares no `upgrade` mode in any published version, so an upgrade scenario -needs a package that supports one. +`_ALL_CHECKED` summary artifacts. `sigterm_grace` covers pod teardown mid-step: it deletes +the package pod while `apply.sh` is running and asserts the step finishes inside the package's +`gracefulShutdown` rather than being killed, and that the retry then skips it. `upgrade` is **not** +covered: the `shellscript` package these scenarios use declares no `upgrade` mode in any published +version, so an upgrade scenario needs a package that supports one. `agent-ci.yaml` already runs it when the *agent* changes, against a freshly built agent. This row covers the other direction: the operator is what builds the pod the agent runs in — its args, diff --git a/k8s-tests/operator-agent/sigterm_grace/assert.yaml b/k8s-tests/operator-agent/sigterm_grace/assert.yaml new file mode 100644 index 00000000..886f67f0 --- /dev/null +++ b/k8s-tests/operator-agent/sigterm_grace/assert.yaml @@ -0,0 +1,43 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + + +apiVersion: v1 +kind: Node +metadata: + labels: + nodewright.nvidia.com/test-node: skyhooke2e + nodewright.nvidia.com/status_sigterm-grace-agent-operator: complete + annotations: + nodewright.nvidia.com/status_sigterm-grace-agent-operator: complete +status: + (conditions[?type == 'nodewright.nvidia.com/sigterm-grace-agent-operator/NotReady']): + - reason: "Complete" + status: "False" + (conditions[?type == 'nodewright.nvidia.com/sigterm-grace-agent-operator/Erroring']): + - reason: "Not Erroring" + status: "False" +--- +apiVersion: nodewright.nvidia.com/v1alpha1 +kind: NodeWright +metadata: + name: sigterm-grace-agent-operator +status: + status: complete + nodeStatus: + # grab values should be one and is complete + (values(@)): + - complete diff --git a/k8s-tests/operator-agent/sigterm_grace/chainsaw-test.yaml b/k8s-tests/operator-agent/sigterm_grace/chainsaw-test.yaml new file mode 100644 index 00000000..5dfed03c --- /dev/null +++ b/k8s-tests/operator-agent/sigterm_grace/chainsaw-test.yaml @@ -0,0 +1,102 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +# yaml-language-server: $schema=https://raw.githubusercontent.com/kyverno/chainsaw/main/.schemas/json/test-chainsaw-v1alpha1.json +apiVersion: chainsaw.kyverno.io/v1alpha1 +kind: Test +metadata: + name: sigterm-grace-agent-operator +spec: + timeouts: + assert: 240s + # The step sleeps for 60s and the second poll below waits for it to end, so + # the suite's usual 90s exec budget leaves no room for a slow scheduler. The + # retry counts below are sized so a poll that never matches still runs out + # inside this budget: each check_node.sh iteration is a kubectl exec plus a 1s + # sleep, and if chainsaw kills the op first the Data:/Check: diagnostic is lost. + exec: 150s + steps: + - try: + - script: + content: | + ## remove annotation from last run + ../../../operator/bin/nodewright reset sigterm-grace-agent-operator --confirm 2>/dev/null || true + - script: + content: | + ## reinstall the debug pod in case it was deleted + ../setup.sh kind-worker setup + - script: + content: | + ## clean up from prior runs + ../check_node.sh kind-worker "rm -rf /var/log/skyhook/sigterm-grace-agent-operator || true" ".*" 2 + ../check_node.sh kind-worker "rm -rf /var/lib/skyhook/sigterm-grace-agent-operator || true" ".*" 2 + - apply: + file: nodewright.yaml + - script: + content: | + set -e + ## Wait until apply.sh is genuinely mid-run before pulling the pod. Deleting it + ## any earlier tests scheduling, not signal handling. + ../check_node.sh kind-worker "cat /var/lib/skyhook/sigterm-grace-agent-operator/progress" "^start" 60 + - script: + content: | + set -e + ## This is the regression under test (#672). The pod receives SIGTERM while + ## apply.sh is asleep, with gracefulShutdown well above the remaining sleep. + ## The agent must let the step finish: the Python agent always did, and the Go + ## agent used to SIGKILL the step's process group the instant SIGTERM arrived, + ## which made gracefulShutdown a no-op. + ## + ## Delete only a pod whose apply container is running right now, and fail if there + ## is none. A label-scoped delete with --ignore-not-found would also accept a pod + ## that had already finished apply, or no pod at all, and the start/end counts + ## below would then pass without a SIGTERM ever reaching the step. The apply + ## stage is an init container named -apply (job_builder.go), so its + ## running state is what identifies the pod. + pod=$(kubectl -n nodewright get pod \ + -l nodewright.nvidia.com/name=sigterm-grace-agent-operator,nodewright.nvidia.com/stage=apply \ + -o jsonpath='{range .items[*]}{.metadata.name}{" "}{.status.initContainerStatuses[?(@.name=="sigterm-grace-agent-operator-apply")].state.running.startedAt}{"\n"}{end}' \ + | awk 'NF == 2 { print $1; exit }') + if [ -z "$pod" ]; then + echo "no pod with a running sigterm-grace-agent-operator-apply container" + kubectl -n nodewright get pod -l nodewright.nvidia.com/name=sigterm-grace-agent-operator -o wide + exit 1 + fi + ## --wait=false because the pod legitimately outlives the delete call for the + ## rest of the step. + kubectl -n nodewright delete pod "$pod" --wait=false + - script: + content: | + set -e + ## The step wrote "start" before the delete; only a step that survived the + ## SIGTERM writes "end". + progress=/var/lib/skyhook/sigterm-grace-agent-operator/progress + ../check_node.sh kind-worker "cat $progress" "^end" 75 + ## "end" alone does not prove survival: a killed step leaves no flag, so the + ## retry pod re-runs apply.sh from scratch and appends a second "start" and + ## then its own "end". Exactly one of each is what pins the first attempt as + ## the one that finished. grep -c prints the count, so match it as a whole line. + ../check_node.sh kind-worker "grep -c '^start' $progress" "^1$" 2 + ../check_node.sh kind-worker "grep -c '^end' $progress" "^1$" 2 + - assert: + ## The retry pod finds apply.sh's completion flag, skips it and finishes the + ## remaining stages, so the package converges. A Go agent that had killed the + ## step would instead re-run it from scratch here, and the absence of "end" + ## above would already have failed the test. + file: assert.yaml + - finally: + - delete: + file: nodewright.yaml diff --git a/k8s-tests/operator-agent/sigterm_grace/nodewright.yaml b/k8s-tests/operator-agent/sigterm_grace/nodewright.yaml new file mode 100644 index 00000000..b9db8200 --- /dev/null +++ b/k8s-tests/operator-agent/sigterm_grace/nodewright.yaml @@ -0,0 +1,55 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +apiVersion: nodewright.nvidia.com/v1alpha1 +kind: NodeWright +metadata: + labels: + app.kubernetes.io/part-of: skyhook-operator + app.kubernetes.io/created-by: skyhook-operator + name: sigterm-grace-agent-operator +spec: + nodeSelectors: + matchLabels: + nodewright.nvidia.com/test-node: skyhooke2e + packages: + sigterm-grace-agent-operator: + version: "1.1.1" + image: ghcr.io/nvidia/skyhook-packages/shellscript + # Longer than the step below runs. The test deletes the package pod while + # the step is mid-sleep; this is the window the step must be allowed to + # finish within. See docs/user-guide/custom-resource.md#gracefulshutdown. + gracefulShutdown: 2m + configMap: + # Progress markers go to the package's state root rather than + # $SKYHOOK_DIR: the copy dir is per-attempt, and the retry after the + # pod delete must read what the killed attempt wrote. + apply.sh: | + #!/bin/bash + progress=/var/lib/skyhook/sigterm-grace-agent-operator/progress + mkdir -p "$(dirname "$progress")" + echo "start $(date +%s)" >> "$progress" + sleep 60 + echo "end $(date +%s)" >> "$progress" + apply_check.sh: | + #!/bin/bash + echo "apply checked" + config.sh: | + #!/bin/bash + echo "configuring" + config_check.sh: | + #!/bin/bash + echo "config checked"