mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
Increase actor termination grace period to 30m (#1565)
A one-minute grace period is nowhere near enough for an actor to suspend normally, so evictions were very likely to leave actors CRASHED instead of cleanly suspended. Previously there were two 1 minute timeouts. Waiting for in-flight RPCs to finish, and then another for SIGTERM. With such a short termination grace period the in-flight RPC could have eaten a significant portion of the timeout leaving very little for the actor. With a longer timeout this is less of a concern, and has now been collapsed into a single deadline for both operations. Fixes #1560 - [ x ] Tests pass - [ x ] Appropriate changes to documentation are included in the PR
This commit is contained in:
+39
-20
@@ -93,9 +93,11 @@ var (
|
||||
// ingress to.
|
||||
const actorHTTPUpstream = "http://" + ateomnet.ActorVethIP + ":80"
|
||||
|
||||
// Workers get a conservative shutdown period. This needs to be significantly less than the K8s
|
||||
// termination grace period for the ateom.
|
||||
const workloadGracePeriod = 1 * time.Minute
|
||||
// workloadGracePeriod is the whole budget for draining the worker on shutdown.
|
||||
// It needs to stay significantly less than the K8s termination grace period
|
||||
// for the ateom, so the escalation to SIGKILL happens here rather than as a
|
||||
// kubelet SIGKILL of ateom itself.
|
||||
const workloadGracePeriod = 30 * time.Minute
|
||||
|
||||
// resumeTimeout is the conservative ceiling for unpausing a paused sandbox.
|
||||
const resumeTimeout = 30 * time.Second
|
||||
@@ -492,15 +494,18 @@ func (s *AteomService) gracefulShutdown(ctx context.Context) {
|
||||
// a SIGTERM.
|
||||
s.cancelActiveRestoreOrRunRPC()
|
||||
|
||||
// One deadline covers the whole drain. Waiting for the lock and waiting out
|
||||
// SIGTERM below both run against it, so the two phases split a single grace
|
||||
// period rather than each getting one: an RPC that burns most of the budget
|
||||
// leaves the containers only the remainder, and the total stays bounded by
|
||||
// workloadGracePeriod however the time falls between them.
|
||||
deadline := time.Now().Add(workloadGracePeriod)
|
||||
|
||||
// Attempt to acquire the lock used to serialize ateom RPCs. This will wait for any
|
||||
// pending RPCs to finish (suspend, resume, etc...). After the RPCs finish there
|
||||
// should be no active session. The run / resume was cancelled and the
|
||||
// checkpoint / restore will stop the workload and clear the active session.
|
||||
//
|
||||
// In the worst case, these RPCs take almost the entire grace period and then
|
||||
// fail. We will then proceed to send SIGTERM to the containers and wait for
|
||||
// them to exit, potentially waiting for 2x the total grace period.
|
||||
lockCtx, lockCancel := context.WithTimeout(ctx, workloadGracePeriod)
|
||||
lockCtx, lockCancel := context.WithDeadline(ctx, deadline)
|
||||
defer lockCancel()
|
||||
|
||||
if !s.lock.LockContext(lockCtx) {
|
||||
@@ -521,7 +526,7 @@ func (s *AteomService) gracefulShutdown(ctx context.Context) {
|
||||
wg.Add(1)
|
||||
go func(containerName string) {
|
||||
defer wg.Done()
|
||||
if err := s.killContainer(ctx, session, containerName); err != nil {
|
||||
if err := killContainer(ctx, session.rcmd, containerName, deadline); err != nil {
|
||||
slog.WarnContext(ctx, "Failed to kill container during shutdown", slog.String("container", containerName), slog.Any("err", err))
|
||||
}
|
||||
}(name)
|
||||
@@ -531,23 +536,39 @@ func (s *AteomService) gracefulShutdown(ctx context.Context) {
|
||||
slog.InfoContext(ctx, "Shutting down")
|
||||
}
|
||||
|
||||
// killContainer stops a container by sending SIGTERM, waiting for the grace period,
|
||||
// and escalating to SIGKILL if necessary.
|
||||
func (s *AteomService) killContainer(ctx context.Context, session *workloadSession, name string) error {
|
||||
// containerKillTimeout bounds the post-SIGKILL wait, so a completely broken
|
||||
// gVisor cannot hold shutdown open indefinitely. It is deliberately not drawn
|
||||
// from the grace period: by this point the container has already had its
|
||||
// allowance and the deadline has passed. A var so tests can shorten it.
|
||||
var containerKillTimeout = 5 * time.Second
|
||||
|
||||
// containerRuntime is the slice of *runsc that graceful shutdown needs. Narrowed
|
||||
// to an interface so killContainer's SIGTERM-then-SIGKILL escalation can be
|
||||
// exercised without executing runsc.
|
||||
type containerRuntime interface {
|
||||
cmdKill(ctx context.Context, containerName, signal string) error
|
||||
cmdWait(ctx context.Context, containerName string) error
|
||||
}
|
||||
|
||||
// killContainer stops a container by sending SIGTERM, waiting until deadline, and
|
||||
// escalating to SIGKILL if necessary. deadline is the shared drain deadline, so a
|
||||
// caller that has already spent most of the grace period elsewhere leaves the
|
||||
// container only what is left of it.
|
||||
func killContainer(ctx context.Context, rcmd containerRuntime, name string, deadline time.Time) error {
|
||||
// Propagate SIGTERM to the application container so it can save state and close connections.
|
||||
// If the actor installed no SIGTERM handler it terminates immediately.
|
||||
slog.InfoContext(ctx, "Sending SIGTERM to container", slog.String("container", name))
|
||||
if err := session.rcmd.cmdKill(ctx, name, "SIGTERM"); err != nil {
|
||||
if err := rcmd.cmdKill(ctx, name, "SIGTERM"); err != nil {
|
||||
slog.ErrorContext(ctx, "Failed to propagate SIGTERM to container", slog.String("container", name), slog.Any("err", err))
|
||||
return fmt.Errorf("failed to propagate SIGTERM to container %q: %w", name, err)
|
||||
}
|
||||
|
||||
done := make(chan error, 1)
|
||||
go func() {
|
||||
done <- session.rcmd.cmdWait(ctx, name)
|
||||
done <- rcmd.cmdWait(ctx, name)
|
||||
}()
|
||||
|
||||
sigTermCtx, sigTermCtxCancel := context.WithTimeout(ctx, workloadGracePeriod)
|
||||
sigTermCtx, sigTermCtxCancel := context.WithDeadline(ctx, deadline)
|
||||
defer sigTermCtxCancel()
|
||||
|
||||
err := waitContainerStop(sigTermCtx, done)
|
||||
@@ -568,15 +589,13 @@ func (s *AteomService) killContainer(ctx context.Context, session *workloadSessi
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
// sigTermCtx timed out. Send SIGKILL.
|
||||
// sigTermCtx hit the drain deadline. Send SIGKILL.
|
||||
slog.WarnContext(ctx, "Grace period expired; killing container", slog.String("container", name))
|
||||
if err := session.rcmd.cmdKill(ctx, name, "SIGKILL"); err != nil {
|
||||
if err := rcmd.cmdKill(ctx, name, "SIGKILL"); err != nil {
|
||||
slog.WarnContext(ctx, "Failed to send SIGKILL to container (it might have already exited)", slog.String("container", name), slog.Any("err", err))
|
||||
}
|
||||
|
||||
// Block until the killed container actually exits, but set a short timeout (e.g. 5 seconds)
|
||||
// to avoid blocking indefinitely if gVisor is completely broken.
|
||||
killCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
|
||||
killCtx, cancel := context.WithTimeout(ctx, containerKillTimeout)
|
||||
defer cancel()
|
||||
|
||||
err = waitContainerStop(killCtx, done)
|
||||
|
||||
@@ -0,0 +1,195 @@
|
||||
//go:build linux
|
||||
|
||||
// Copyright 2026 Google LLC
|
||||
//
|
||||
// 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.
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"reflect"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// fakeRuntime stands in for *runsc. It records the signals killContainer
|
||||
// delivers and unblocks cmdWait when the container is configured to die on one
|
||||
// of them.
|
||||
type fakeRuntime struct {
|
||||
mu sync.Mutex
|
||||
signals []string
|
||||
|
||||
// exitOn is the signal that makes the container exit. Empty means it never
|
||||
// does, which is how the wedged-sandbox case is set up.
|
||||
exitOn string
|
||||
// waitErr is what cmdWait reports once it unblocks, standing in for a
|
||||
// `runsc wait` that failed rather than a container that stopped.
|
||||
waitErr error
|
||||
// killErr maps a signal to the error cmdKill returns for it.
|
||||
killErr map[string]error
|
||||
|
||||
exited chan struct{}
|
||||
exitedOnce sync.Once
|
||||
}
|
||||
|
||||
func newFakeRuntime(exitOn string, waitErr error, killErr map[string]error) *fakeRuntime {
|
||||
return &fakeRuntime{exitOn: exitOn, waitErr: waitErr, killErr: killErr, exited: make(chan struct{})}
|
||||
}
|
||||
|
||||
func (f *fakeRuntime) cmdKill(_ context.Context, _, signal string) error {
|
||||
f.mu.Lock()
|
||||
f.signals = append(f.signals, signal)
|
||||
f.mu.Unlock()
|
||||
|
||||
if f.exitOn == signal {
|
||||
f.exitedOnce.Do(func() { close(f.exited) })
|
||||
}
|
||||
return f.killErr[signal]
|
||||
}
|
||||
|
||||
func (f *fakeRuntime) cmdWait(ctx context.Context, _ string) error {
|
||||
select {
|
||||
case <-f.exited:
|
||||
return f.waitErr
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
}
|
||||
|
||||
func (f *fakeRuntime) sentSignals() []string {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
return append([]string(nil), f.signals...)
|
||||
}
|
||||
|
||||
// TestKillContainer covers the SIGTERM-then-SIGKILL escalation against the
|
||||
// shared drain deadline. The deadlines here are milliseconds rather than the
|
||||
// production grace period, which is what the package-level vars are for.
|
||||
func TestKillContainer(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
// exitOn, waitErr and killErr configure the fake runtime.
|
||||
exitOn string
|
||||
waitErr error
|
||||
killErr map[string]error
|
||||
// deadlineIn is the drain deadline relative to the start of the call. A
|
||||
// negative value is a deadline the caller has already spent elsewhere.
|
||||
deadlineIn time.Duration
|
||||
wantSignals []string
|
||||
wantErr string
|
||||
}{
|
||||
{
|
||||
name: "container exits on SIGTERM inside the deadline",
|
||||
exitOn: "SIGTERM",
|
||||
deadlineIn: time.Minute,
|
||||
wantSignals: []string{"SIGTERM"},
|
||||
},
|
||||
{
|
||||
name: "container ignoring SIGTERM is killed at the deadline",
|
||||
exitOn: "SIGKILL",
|
||||
deadlineIn: 20 * time.Millisecond,
|
||||
wantSignals: []string{"SIGTERM", "SIGKILL"},
|
||||
},
|
||||
{
|
||||
// The lock wait ate the whole grace period, so the container gets
|
||||
// SIGTERM and no time at all before the escalation.
|
||||
name: "deadline already spent leaves no grace",
|
||||
exitOn: "SIGKILL",
|
||||
deadlineIn: -time.Second,
|
||||
wantSignals: []string{"SIGTERM", "SIGKILL"},
|
||||
},
|
||||
{
|
||||
name: "container surviving SIGKILL is reported",
|
||||
deadlineIn: 20 * time.Millisecond,
|
||||
wantSignals: []string{"SIGTERM", "SIGKILL"},
|
||||
wantErr: "failed to exit even after SIGKILL",
|
||||
},
|
||||
{
|
||||
name: "undeliverable SIGTERM does not escalate",
|
||||
killErr: map[string]error{"SIGTERM": errors.New("sandbox gone")},
|
||||
deadlineIn: time.Minute,
|
||||
wantSignals: []string{"SIGTERM"},
|
||||
wantErr: "failed to propagate SIGTERM",
|
||||
},
|
||||
{
|
||||
// A failing `runsc wait` says nothing about the container, so the
|
||||
// caller is told rather than the container killed.
|
||||
name: "failed wait does not escalate",
|
||||
exitOn: "SIGTERM",
|
||||
waitErr: errors.New("runsc wait exploded"),
|
||||
deadlineIn: time.Minute,
|
||||
wantSignals: []string{"SIGTERM"},
|
||||
wantErr: "wait failed",
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range tests {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
origKillTimeout := containerKillTimeout
|
||||
containerKillTimeout = 20 * time.Millisecond
|
||||
t.Cleanup(func() { containerKillTimeout = origKillTimeout })
|
||||
|
||||
// A cancelled context is what releases a cmdWait the fake never
|
||||
// unblocks, so the wait goroutine does not outlive the test.
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.Cleanup(cancel)
|
||||
|
||||
f := newFakeRuntime(tc.exitOn, tc.waitErr, tc.killErr)
|
||||
err := killContainer(ctx, f, "counter", time.Now().Add(tc.deadlineIn))
|
||||
|
||||
switch {
|
||||
case tc.wantErr == "" && err != nil:
|
||||
t.Errorf("killContainer() = %v, want nil", err)
|
||||
case tc.wantErr != "" && err == nil:
|
||||
t.Errorf("killContainer() = nil, want error containing %q", tc.wantErr)
|
||||
case tc.wantErr != "" && !strings.Contains(err.Error(), tc.wantErr):
|
||||
t.Errorf("killContainer() = %v, want error containing %q", err, tc.wantErr)
|
||||
}
|
||||
|
||||
if got := f.sentSignals(); !reflect.DeepEqual(got, tc.wantSignals) {
|
||||
t.Errorf("signals delivered = %v, want %v", got, tc.wantSignals)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestKillContainerHonorsParentCancellation asserts that a cancelled shutdown
|
||||
// context stops the drain instead of escalating: ateom is going away anyway, and
|
||||
// the containers go down with the pod.
|
||||
func TestKillContainerHonorsParentCancellation(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
f := newFakeRuntime("", nil, nil)
|
||||
|
||||
// Cancel once SIGTERM has been delivered and the wait is under way.
|
||||
go func() {
|
||||
for {
|
||||
if len(f.sentSignals()) > 0 {
|
||||
cancel()
|
||||
return
|
||||
}
|
||||
time.Sleep(time.Millisecond)
|
||||
}
|
||||
}()
|
||||
|
||||
err := killContainer(ctx, f, "counter", time.Now().Add(time.Minute))
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Errorf("killContainer() = %v, want context.Canceled", err)
|
||||
}
|
||||
if got := f.sentSignals(); !reflect.DeepEqual(got, []string{"SIGTERM"}) {
|
||||
t.Errorf("signals delivered = %v, want [SIGTERM]", got)
|
||||
}
|
||||
}
|
||||
@@ -34,18 +34,22 @@ import (
|
||||
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/kata"
|
||||
)
|
||||
|
||||
// workloadGracePeriod is the whole budget for draining the worker on shutdown:
|
||||
// waiting for an in-flight RPC to release the lock and letting the guest
|
||||
// workloads handle SIGTERM both draw on it, and ateom escalates to SIGKILL once
|
||||
// it is gone. Matches ateom-gvisor, and is deliberately shorter than the pod's
|
||||
// own termination grace period — 3600s, set by
|
||||
// workerTerminationGracePeriodSeconds in cmd/atecontroller — so the escalation
|
||||
// happens here rather than as a kubelet SIGKILL of ateom itself.
|
||||
const workloadGracePeriod = 30 * time.Minute
|
||||
|
||||
// workloadKillTimeout bounds the post-SIGKILL wait. The VM teardown that
|
||||
// follows is what ultimately guarantees the workload is gone, so a wedged
|
||||
// kata-agent must not hold shutdown open past this.
|
||||
// A var so tests can shorten it.
|
||||
var workloadKillTimeout = 5 * time.Second
|
||||
|
||||
const (
|
||||
// workloadGracePeriod is how long a guest workload gets to handle SIGTERM and
|
||||
// exit on its own before ateom escalates to SIGKILL. Matches ateom-gvisor, and
|
||||
// is deliberately shorter than the pod's own termination grace period so the
|
||||
// escalation happens here rather than as a kubelet SIGKILL of ateom itself.
|
||||
workloadGracePeriod = 1 * time.Minute
|
||||
|
||||
// workloadKillTimeout bounds the post-SIGKILL wait. The VM teardown that
|
||||
// follows is what ultimately guarantees the workload is gone, so a wedged
|
||||
// kata-agent must not hold shutdown open past this.
|
||||
workloadKillTimeout = 5 * time.Second
|
||||
|
||||
// signalDeliveryTimeout bounds one SignalProcess round-trip. Delivering a signal
|
||||
// is a local ttrpc call that returns in microseconds; if it has not come back by
|
||||
// now the agent is not answering, and waiting longer will not change that.
|
||||
@@ -72,12 +76,17 @@ func (s *AteomService) gracefulShutdown(ctx context.Context) {
|
||||
// SIGTERM the guest it just produced is strictly worse than aborting it.
|
||||
s.cancelActiveRestoreOrRunRPC()
|
||||
|
||||
// One deadline covers the whole drain. Waiting for the lock and waiting out
|
||||
// SIGTERM below both run against it, so the two phases split a single grace
|
||||
// period rather than each getting one: an RPC that burns most of the budget
|
||||
// leaves the workloads only the remainder, and the total stays bounded by
|
||||
// workloadGracePeriod however the time falls between them.
|
||||
deadline := time.Now().Add(workloadGracePeriod)
|
||||
|
||||
// Wait for whatever still holds lock — a suspend, a resume — to finish, but
|
||||
// only for the grace period. In the worst case that RPC burns nearly all of it
|
||||
// and then fails, and the stop below spends another grace period on top; that
|
||||
// is bounded well inside the pod's own termination grace period, and is the
|
||||
// price of not truncating an RPC that may be saving the actor's state.
|
||||
lockCtx, lockCancel := context.WithTimeout(ctx, workloadGracePeriod)
|
||||
// not past the deadline. Letting it run that long is the price of not
|
||||
// truncating an RPC that may be saving the actor's state.
|
||||
lockCtx, lockCancel := context.WithDeadline(ctx, deadline)
|
||||
defer lockCancel()
|
||||
if !s.lock.LockContext(lockCtx) {
|
||||
slog.ErrorContext(ctx, "Failed to acquire lock during graceful shutdown; another RPC is still running")
|
||||
@@ -105,9 +114,20 @@ func (s *AteomService) gracefulShutdown(ctx context.Context) {
|
||||
return
|
||||
}
|
||||
|
||||
// Drain the actors concurrently, for the same reason their workloads are
|
||||
// drained concurrently below: they share one deadline, so in series the first
|
||||
// actor's wait would come out of every later one's allowance and the last
|
||||
// would be SIGKILLed with no grace at all.
|
||||
var wg sync.WaitGroup
|
||||
for _, t := range targets {
|
||||
gracefullyStopActor(ctx, t)
|
||||
wg.Add(1)
|
||||
go func(t drainTarget) {
|
||||
defer wg.Done()
|
||||
gracefullyStopActor(ctx, t, deadline)
|
||||
}(t)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
slog.InfoContext(ctx, "Shutting down")
|
||||
}
|
||||
|
||||
@@ -127,9 +147,18 @@ type drainTarget struct {
|
||||
workloadIDs []string
|
||||
}
|
||||
|
||||
// guestAgent is the slice of *kata.AgentClient the drain needs: signal a guest
|
||||
// process and wait for it to exit. Narrowed to an interface so
|
||||
// stopGuestWorkload's SIGTERM-then-SIGKILL escalation can be exercised without a
|
||||
// running VM.
|
||||
type guestAgent interface {
|
||||
SignalProcess(ctx context.Context, containerID, execID string, signal uint32) error
|
||||
WaitProcess(ctx context.Context, containerID, execID string) (int32, error)
|
||||
}
|
||||
|
||||
// gracefullyStopActor signals the actor's guest workloads with SIGTERM and waits
|
||||
// out the grace period, escalating to SIGKILL.
|
||||
func gracefullyStopActor(ctx context.Context, t drainTarget) {
|
||||
// until deadline, escalating to SIGKILL.
|
||||
func gracefullyStopActor(ctx context.Context, t drainTarget, deadline time.Time) {
|
||||
id := t.id
|
||||
|
||||
// Obtain a kata-agent client to signal the guest: reuse the log-forwarding
|
||||
@@ -148,15 +177,15 @@ func gracefullyStopActor(ctx context.Context, t drainTarget) {
|
||||
agent, dialed = a, a
|
||||
}
|
||||
|
||||
// Stop the workloads concurrently. Each one is entitled to the full grace
|
||||
// period, so stopping them in series would multiply it by the container
|
||||
// count and overrun the pod's own termination grace period.
|
||||
// Stop the workloads concurrently. They share one deadline, so stopping them
|
||||
// in series would spend the first workload's wait out of every later one's
|
||||
// allowance and leave the last with none.
|
||||
var wg sync.WaitGroup
|
||||
for _, wid := range t.workloadIDs {
|
||||
wg.Add(1)
|
||||
go func(wid string) {
|
||||
defer wg.Done()
|
||||
if err := stopGuestWorkload(ctx, agent, id, wid); err != nil {
|
||||
if err := stopGuestWorkload(ctx, agent, id, wid, deadline); err != nil {
|
||||
slog.WarnContext(ctx, "Failed to stop guest workload during shutdown", slog.String("id", id), slog.String("workload", wid), slog.Any("err", err))
|
||||
}
|
||||
}(wid)
|
||||
@@ -167,9 +196,11 @@ func gracefullyStopActor(ctx context.Context, t drainTarget) {
|
||||
}
|
||||
}
|
||||
|
||||
// stopGuestWorkload stops one guest workload, wait out workloadGracePeriod, then
|
||||
// escalate to SIGKILL and wait a bounded time for the kill to land.
|
||||
func stopGuestWorkload(ctx context.Context, agent *kata.AgentClient, id, wid string) error {
|
||||
// stopGuestWorkload stops one guest workload, waits until deadline, then
|
||||
// escalates to SIGKILL and waits a bounded time for the kill to land. deadline is
|
||||
// the shared drain deadline, so a caller that has already spent most of the grace
|
||||
// period elsewhere leaves the workload only what is left of it.
|
||||
func stopGuestWorkload(ctx context.Context, agent guestAgent, id, wid string, deadline time.Time) error {
|
||||
// Propagate SIGTERM so the actor can save state and close connections.
|
||||
// An actor that installed no handler terminates immediately.
|
||||
slog.InfoContext(ctx, "Sending SIGTERM to guest workload", slog.String("id", id), slog.String("workload", wid))
|
||||
@@ -190,7 +221,7 @@ func stopGuestWorkload(ctx context.Context, agent *kata.AgentClient, id, wid str
|
||||
done <- err
|
||||
}()
|
||||
|
||||
termCtx, termCancel := context.WithTimeout(ctx, workloadGracePeriod)
|
||||
termCtx, termCancel := context.WithDeadline(ctx, deadline)
|
||||
defer termCancel()
|
||||
err := waitWorkloadStop(termCtx, done)
|
||||
if err == nil {
|
||||
@@ -216,7 +247,9 @@ func stopGuestWorkload(ctx context.Context, agent *kata.AgentClient, id, wid str
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
slog.WarnContext(ctx, "Grace period expired; killing guest workload", slog.String("id", id), slog.String("workload", wid), slog.Duration("grace", workloadGracePeriod))
|
||||
// The deadline, not the configured grace period: the lock wait may have eaten
|
||||
// part of the budget before the workload ever saw SIGTERM.
|
||||
slog.WarnContext(ctx, "Grace period expired; killing guest workload", slog.String("id", id), slog.String("workload", wid), slog.Time("deadline", deadline))
|
||||
if err := signalWorkload(ctx, agent, wid, syscall.SIGKILL); err != nil {
|
||||
slog.WarnContext(ctx, "Failed to SIGKILL guest workload (it might have already exited)", slog.String("id", id), slog.String("workload", wid), slog.Any("err", err))
|
||||
}
|
||||
@@ -244,7 +277,7 @@ func stopGuestWorkload(ctx context.Context, agent *kata.AgentClient, id, wid str
|
||||
// land mid-drain — leaves the unix socket to CH perfectly healthy while the agent
|
||||
// never answers, and an unbounded call there would hang until the kubelet's
|
||||
// SIGKILL at the end of the pod's termination grace period.
|
||||
func signalWorkload(ctx context.Context, agent *kata.AgentClient, wid string, sig syscall.Signal) error {
|
||||
func signalWorkload(ctx context.Context, agent guestAgent, wid string, sig syscall.Signal) error {
|
||||
sigCtx, cancel := context.WithTimeout(ctx, signalDeliveryTimeout)
|
||||
defer cancel()
|
||||
return agent.SignalProcess(sigCtx, wid, wid, uint32(sig))
|
||||
|
||||
@@ -0,0 +1,197 @@
|
||||
//go:build linux
|
||||
|
||||
// Copyright 2026 Google LLC
|
||||
//
|
||||
// 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.
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"reflect"
|
||||
"strings"
|
||||
"sync"
|
||||
"syscall"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// fakeGuestAgent stands in for *kata.AgentClient. It records the signals
|
||||
// stopGuestWorkload delivers and unblocks WaitProcess when the workload is
|
||||
// configured to die on one of them.
|
||||
type fakeGuestAgent struct {
|
||||
mu sync.Mutex
|
||||
signals []syscall.Signal
|
||||
|
||||
// exitOn is the signal that makes the workload exit. Zero means it never
|
||||
// does, which is how the wedged-guest case is set up.
|
||||
exitOn syscall.Signal
|
||||
// waitErr is what WaitProcess reports once it unblocks, standing in for a
|
||||
// dead agent connection rather than a process that stopped.
|
||||
waitErr error
|
||||
// signalErr maps a signal to the error SignalProcess returns for it.
|
||||
signalErr map[syscall.Signal]error
|
||||
|
||||
exited chan struct{}
|
||||
exitedOnce sync.Once
|
||||
}
|
||||
|
||||
func newFakeGuestAgent(exitOn syscall.Signal, waitErr error, signalErr map[syscall.Signal]error) *fakeGuestAgent {
|
||||
return &fakeGuestAgent{exitOn: exitOn, waitErr: waitErr, signalErr: signalErr, exited: make(chan struct{})}
|
||||
}
|
||||
|
||||
func (f *fakeGuestAgent) SignalProcess(_ context.Context, _, _ string, signal uint32) error {
|
||||
sig := syscall.Signal(signal)
|
||||
|
||||
f.mu.Lock()
|
||||
f.signals = append(f.signals, sig)
|
||||
f.mu.Unlock()
|
||||
|
||||
if f.exitOn != 0 && f.exitOn == sig {
|
||||
f.exitedOnce.Do(func() { close(f.exited) })
|
||||
}
|
||||
return f.signalErr[sig]
|
||||
}
|
||||
|
||||
func (f *fakeGuestAgent) WaitProcess(ctx context.Context, _, _ string) (int32, error) {
|
||||
select {
|
||||
case <-f.exited:
|
||||
return 0, f.waitErr
|
||||
case <-ctx.Done():
|
||||
return 0, ctx.Err()
|
||||
}
|
||||
}
|
||||
|
||||
func (f *fakeGuestAgent) sentSignals() []syscall.Signal {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
return append([]syscall.Signal(nil), f.signals...)
|
||||
}
|
||||
|
||||
// TestStopGuestWorkload covers the SIGTERM-then-SIGKILL escalation against the
|
||||
// shared drain deadline. The deadlines here are milliseconds rather than the
|
||||
// production grace period, which is what the package-level vars are for.
|
||||
func TestStopGuestWorkload(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
// exitOn, waitErr and signalErr configure the fake agent.
|
||||
exitOn syscall.Signal
|
||||
waitErr error
|
||||
signalErr map[syscall.Signal]error
|
||||
// deadlineIn is the drain deadline relative to the start of the call. A
|
||||
// negative value is a deadline the caller has already spent elsewhere.
|
||||
deadlineIn time.Duration
|
||||
wantSignals []syscall.Signal
|
||||
wantErr string
|
||||
}{
|
||||
{
|
||||
name: "workload exits on SIGTERM inside the deadline",
|
||||
exitOn: syscall.SIGTERM,
|
||||
deadlineIn: time.Minute,
|
||||
wantSignals: []syscall.Signal{syscall.SIGTERM},
|
||||
},
|
||||
{
|
||||
name: "workload ignoring SIGTERM is killed at the deadline",
|
||||
exitOn: syscall.SIGKILL,
|
||||
deadlineIn: 20 * time.Millisecond,
|
||||
wantSignals: []syscall.Signal{syscall.SIGTERM, syscall.SIGKILL},
|
||||
},
|
||||
{
|
||||
// The lock wait ate the whole grace period, so the workload gets
|
||||
// SIGTERM and no time at all before the escalation.
|
||||
name: "deadline already spent leaves no grace",
|
||||
exitOn: syscall.SIGKILL,
|
||||
deadlineIn: -time.Second,
|
||||
wantSignals: []syscall.Signal{syscall.SIGTERM, syscall.SIGKILL},
|
||||
},
|
||||
{
|
||||
name: "workload surviving SIGKILL is reported",
|
||||
deadlineIn: 20 * time.Millisecond,
|
||||
wantSignals: []syscall.Signal{syscall.SIGTERM, syscall.SIGKILL},
|
||||
wantErr: "failed to exit even after SIGKILL",
|
||||
},
|
||||
{
|
||||
name: "undeliverable SIGTERM does not escalate",
|
||||
signalErr: map[syscall.Signal]error{syscall.SIGTERM: errors.New("ttrpc closed")},
|
||||
deadlineIn: time.Minute,
|
||||
wantSignals: []syscall.Signal{syscall.SIGTERM},
|
||||
wantErr: "while propagating SIGTERM to workload",
|
||||
},
|
||||
{
|
||||
// A WaitProcess that errors is a dead agent connection, not a bad
|
||||
// exit, so liveness is unknown and the caller is told rather than
|
||||
// the workload killed.
|
||||
name: "failed wait does not escalate",
|
||||
exitOn: syscall.SIGTERM,
|
||||
waitErr: errors.New("ttrpc: closed"),
|
||||
deadlineIn: time.Minute,
|
||||
wantSignals: []syscall.Signal{syscall.SIGTERM},
|
||||
wantErr: `while waiting for workload "counter" to exit`,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range tests {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
origKillTimeout := workloadKillTimeout
|
||||
workloadKillTimeout = 20 * time.Millisecond
|
||||
t.Cleanup(func() { workloadKillTimeout = origKillTimeout })
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.Cleanup(cancel)
|
||||
|
||||
f := newFakeGuestAgent(tc.exitOn, tc.waitErr, tc.signalErr)
|
||||
err := stopGuestWorkload(ctx, f, "actor-1", "counter", time.Now().Add(tc.deadlineIn))
|
||||
|
||||
switch {
|
||||
case tc.wantErr == "" && err != nil:
|
||||
t.Errorf("stopGuestWorkload() = %v, want nil", err)
|
||||
case tc.wantErr != "" && err == nil:
|
||||
t.Errorf("stopGuestWorkload() = nil, want error containing %q", tc.wantErr)
|
||||
case tc.wantErr != "" && !strings.Contains(err.Error(), tc.wantErr):
|
||||
t.Errorf("stopGuestWorkload() = %v, want error containing %q", err, tc.wantErr)
|
||||
}
|
||||
|
||||
if got := f.sentSignals(); !reflect.DeepEqual(got, tc.wantSignals) {
|
||||
t.Errorf("signals delivered = %v, want %v", got, tc.wantSignals)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestStopGuestWorkloadHonorsParentCancellation asserts that a cancelled shutdown
|
||||
// context stops the drain instead of escalating: ateom is going away anyway, and
|
||||
// the guest goes down with the VM.
|
||||
func TestStopGuestWorkloadHonorsParentCancellation(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
f := newFakeGuestAgent(0, nil, nil)
|
||||
|
||||
// Cancel once SIGTERM has been delivered and the wait is under way.
|
||||
go func() {
|
||||
for {
|
||||
if len(f.sentSignals()) > 0 {
|
||||
cancel()
|
||||
return
|
||||
}
|
||||
time.Sleep(time.Millisecond)
|
||||
}
|
||||
}()
|
||||
|
||||
err := stopGuestWorkload(ctx, f, "actor-1", "counter", time.Now().Add(time.Minute))
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Errorf("stopGuestWorkload() = %v, want context.Canceled", err)
|
||||
}
|
||||
if got := f.sentSignals(); !reflect.DeepEqual(got, []syscall.Signal{syscall.SIGTERM}) {
|
||||
t.Errorf("signals delivered = %v, want [SIGTERM]", got)
|
||||
}
|
||||
}
|
||||
@@ -139,91 +139,6 @@ func waitForWorkerRemoved(ctx context.Context, t *testing.T, clients *e2e.Client
|
||||
}
|
||||
}
|
||||
|
||||
// TestGracefulWorkerTerminationTimeout exercises the case where the workload
|
||||
// container hangs (exceeds the 1-minute workloadGracePeriod) during SIGTERM.
|
||||
// The ateom is expected to SIGKILL the container, letting the control plane
|
||||
// mark the worker removed and the actor CRASHED. Runs against both runtimes.
|
||||
func TestGracefulWorkerTerminationTimeout(t *testing.T) {
|
||||
nsObj := e2e.CreateNamespace(t)
|
||||
|
||||
ctx := context.Background()
|
||||
clients := e2e.GetClients()
|
||||
|
||||
at, err := createActorTemplate(ctx, t, clients, nsObj, ateapipb.SnapshotContentScope_SNAPSHOT_CONTENT_SCOPE_FULL, ateapipb.SnapshotContentScope_SNAPSHOT_CONTENT_SCOPE_FULL, ateapipb.ResumeSource_RESUME_SOURCE_COLD_BOOT)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to initialize ActorTemplate: %v", err)
|
||||
}
|
||||
|
||||
actorID := "graceful-term-timeout-" + nsObj.Name
|
||||
if _, err := clients.SubstrateAPI.CreateActor(ctx, &ateapipb.CreateActorRequest{
|
||||
Actor: &ateapipb.Actor{
|
||||
Metadata: &ateapipb.ResourceMetadata{Atespace: demoAtespace, Name: actorID},
|
||||
ActorTemplate: e2e.TemplateRef(at),
|
||||
},
|
||||
}); err != nil {
|
||||
t.Fatalf("failed to create Actor: %v", err)
|
||||
}
|
||||
defer func() {
|
||||
_, _ = clients.SubstrateAPI.DeleteActor(ctx, &ateapipb.DeleteActorRequest{
|
||||
Actor: &ateapipb.ObjectRef{Atespace: demoAtespace, Name: actorID},
|
||||
})
|
||||
}()
|
||||
|
||||
// Bring the actor up on a worker so it is bound to a pod.
|
||||
if _, err := e2e.ResumeActorAwaitCapacity(t, ctx, clients, &ateapipb.ResumeActorRequest{
|
||||
Actor: &ateapipb.ObjectRef{Atespace: demoAtespace, Name: actorID},
|
||||
}); err != nil {
|
||||
t.Fatalf("failed to resume Actor: %v", err)
|
||||
}
|
||||
waitForActorState(ctx, t, clients, actorID, ateapipb.ActorState_ACTOR_STATE_RUNNING)
|
||||
|
||||
// Set the sigterm sleep interval to 90 seconds (longer than the 1-minute grace period).
|
||||
if _, err := callActorPath(t, resources.ActorRef{Atespace: demoAtespace, Name: actorID}, "GET", "/set-sigterm-sleep?duration=90"); err != nil {
|
||||
t.Fatalf("failed to set sigterm sleep: %v", err)
|
||||
}
|
||||
|
||||
running, err := clients.SubstrateAPI.GetActor(ctx, &ateapipb.GetActorRequest{
|
||||
Actor: &ateapipb.ObjectRef{Atespace: demoAtespace, Name: actorID},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("failed to get running Actor: %v", err)
|
||||
}
|
||||
podNS := running.GetStatus().GetWorkerAssignment().GetWorkerNamespace()
|
||||
podName := running.GetStatus().GetWorkerAssignment().GetWorkerPod()
|
||||
if podNS == "" || podName == "" {
|
||||
t.Fatalf("running actor has no bound worker pod: ns=%q name=%q", podNS, podName)
|
||||
}
|
||||
t.Logf("Actor %q bound to worker pod %s/%s", actorID, podNS, podName)
|
||||
|
||||
// Evict the worker pod. The kubelet sends SIGTERM to ateom, which propagates
|
||||
// it into the sandbox; the container hangs, triggering the 1-minute timeout,
|
||||
// followed by SIGKILL by ateom.
|
||||
if err := clients.K8s.CoreV1().Pods(podNS).Delete(ctx, podName, metav1.DeleteOptions{}); err != nil {
|
||||
t.Fatalf("failed to delete worker pod %s/%s: %v", podNS, podName, err)
|
||||
}
|
||||
|
||||
// The worker record must eventually be removed once the pod is gone.
|
||||
// Since there is a 1-minute timeout + up to 5s SIGKILL wait, we need a
|
||||
// larger timeout (e.g. 120 seconds).
|
||||
if err := waitForWorkerRemoved(ctx, t, clients, podName, 120*time.Second); err != nil {
|
||||
t.Fatalf("worker %s not removed after pod deletion: %v", podName, err)
|
||||
}
|
||||
|
||||
// Verify the actor lands in ACTOR_STATE_CRASHED.
|
||||
waitForActorStateWithTimeout(ctx, t, clients, actorID, ateapipb.ActorState_ACTOR_STATE_CRASHED, 120*time.Second)
|
||||
|
||||
// Verify the pod assignment was cleared.
|
||||
actor, err := clients.SubstrateAPI.GetActor(ctx, &ateapipb.GetActorRequest{
|
||||
Actor: &ateapipb.ObjectRef{Atespace: demoAtespace, Name: actorID},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("failed to get actor: %v", err)
|
||||
}
|
||||
if pod := actor.GetStatus().GetWorkerAssignment().GetWorkerPod(); pod != "" {
|
||||
t.Errorf("actor still bound to worker pod %q, expected empty", pod)
|
||||
}
|
||||
}
|
||||
|
||||
// TestGracefulWorkerTerminationSuspend exercises the case where a worker pod is
|
||||
// deleted (evicted), and while the container is in its SIGTERM shutdown phase,
|
||||
// we initiate a suspend. Suspend should succeed.
|
||||
|
||||
Reference in New Issue
Block a user