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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 4 additions & 3 deletions agent/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
9 changes: 6 additions & 3 deletions agent/go/internal/agent/agent_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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"))
})
Expand Down Expand Up @@ -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(
Expand Down
30 changes: 13 additions & 17 deletions agent/go/internal/command/process.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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()
Expand All @@ -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()
Expand All @@ -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
Expand Down
26 changes: 21 additions & 5 deletions agent/go/internal/command/process_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand All @@ -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() {
Expand Down Expand Up @@ -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)
Expand Down
5 changes: 1 addition & 4 deletions agent/go/internal/interrupts/node_restart.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@ package interrupts

import (
"context"
"errors"
"fmt"
"syscall"
"time"
Expand Down Expand Up @@ -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 {
Expand Down
55 changes: 46 additions & 9 deletions agent/go/internal/step/regular_step_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,13 @@ package step
import (
"bytes"
"context"
"errors"
"fmt"
"io"
"os"
"path/filepath"
"strconv"
"strings"
"sync"
"syscall"
"time"

Expand Down Expand Up @@ -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))
})
})

Expand Down Expand Up @@ -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("/"),
Expand Down Expand Up @@ -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,
Expand Down
8 changes: 5 additions & 3 deletions docs/contributing/ci-test-pools.md
Original file line number Diff line number Diff line change
Expand Up @@ -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` /
`<stage>_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.
`<stage>_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,
Expand Down
43 changes: 43 additions & 0 deletions k8s-tests/operator-agent/sigterm_grace/assert.yaml
Original file line number Diff line number Diff line change
@@ -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
Loading
Loading