mirror of
https://github.com/vmware-tanzu/velero.git
synced 2026-08-15 19:56:06 +00:00
Fix hook timeout review feedback
Signed-off-by: chlins <chlins.zhang@gmail.com>
This commit is contained in:
@@ -168,7 +168,7 @@ func (e *defaultPodCommandExecutor) ExecutePodCommand(log logrus.FieldLogger, it
|
||||
Stderr: &stderr,
|
||||
}
|
||||
|
||||
// The timeout drives the context so the exec stream is actually cancelled, rather than
|
||||
// The timeout drives the context so the exec stream is actually canceled, rather than
|
||||
// being left running on the API server after this function has returned.
|
||||
ctx, cancel := context.WithTimeout(context.Background(), localHook.Timeout.Duration)
|
||||
defer cancel()
|
||||
@@ -178,17 +178,15 @@ func (e *defaultPodCommandExecutor) ExecutePodCommand(log logrus.FieldLogger, it
|
||||
errCh := make(chan error, 1)
|
||||
|
||||
go func() {
|
||||
errCh <- executor.StreamWithContext(ctx, streamOptions)
|
||||
streamErr := executor.StreamWithContext(ctx, streamOptions)
|
||||
// Inspect the local context as soon as the stream returns. Otherwise a stream error
|
||||
// completed before the deadline could be misclassified if this goroutine sends its
|
||||
// result before the caller is scheduled to receive it.
|
||||
errCh <- normalizeExecHookError(streamErr, ctx.Err(), localHook.Timeout.Duration)
|
||||
}()
|
||||
|
||||
select {
|
||||
case err = <-errCh:
|
||||
// On a timeout the stream returns because the context expired, so both this case
|
||||
// and ctx.Done() are ready and the select picks one at random. Report the timeout
|
||||
// either way instead of surfacing the context error only some of the time.
|
||||
if errors.Is(ctx.Err(), context.DeadlineExceeded) {
|
||||
return errors.Errorf("timed out after %v", localHook.Timeout.Duration)
|
||||
}
|
||||
case <-ctx.Done():
|
||||
return errors.Errorf("timed out after %v", localHook.Timeout.Duration)
|
||||
}
|
||||
@@ -199,6 +197,14 @@ func (e *defaultPodCommandExecutor) ExecutePodCommand(log logrus.FieldLogger, it
|
||||
return err
|
||||
}
|
||||
|
||||
func normalizeExecHookError(streamErr, contextErr error, timeout time.Duration) error {
|
||||
if errors.Is(contextErr, context.DeadlineExceeded) {
|
||||
return errors.Errorf("timed out after %v", timeout)
|
||||
}
|
||||
|
||||
return streamErr
|
||||
}
|
||||
|
||||
func ensureContainerExists(pod *corev1api.Pod, container string) error {
|
||||
existsAsMainContainer := slices.ContainsFunc(pod.Spec.Containers, func(c corev1api.Container) bool {
|
||||
return c.Name == container
|
||||
|
||||
@@ -177,6 +177,15 @@ func TestExecutePodCommand(t *testing.T) {
|
||||
hookError: errors.New("hook error"),
|
||||
expectedError: "hook error",
|
||||
},
|
||||
{
|
||||
name: "stream deadline exceeded before local timeout",
|
||||
command: []string{"some", "command"},
|
||||
expectedContainerName: "foo",
|
||||
expectedErrorMode: v1.HookErrorModeFail,
|
||||
expectedTimeout: defaultTimeout,
|
||||
hookError: context.DeadlineExceeded,
|
||||
expectedError: context.DeadlineExceeded.Error(),
|
||||
},
|
||||
{
|
||||
// Timeouts from pod annotations go through time.ParseDuration, which accepts
|
||||
// negative values. Without clamping, the hook would run with no timeout at all.
|
||||
@@ -264,6 +273,54 @@ func TestExecutePodCommand(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeExecHookError(t *testing.T) {
|
||||
hookErr := errors.New("hook error")
|
||||
tests := []struct {
|
||||
name string
|
||||
streamErr error
|
||||
contextErr error
|
||||
expectedError string
|
||||
preserveStreamErr bool
|
||||
}{
|
||||
{
|
||||
name: "local context deadline exceeded",
|
||||
streamErr: context.DeadlineExceeded,
|
||||
contextErr: context.DeadlineExceeded,
|
||||
expectedError: "timed out after 30s",
|
||||
},
|
||||
{
|
||||
name: "stream deadline exceeded before local timeout",
|
||||
streamErr: context.DeadlineExceeded,
|
||||
expectedError: context.DeadlineExceeded.Error(),
|
||||
preserveStreamErr: true,
|
||||
},
|
||||
{
|
||||
name: "ordinary hook error",
|
||||
streamErr: hookErr,
|
||||
expectedError: hookErr.Error(),
|
||||
preserveStreamErr: true,
|
||||
},
|
||||
{
|
||||
name: "no errors",
|
||||
},
|
||||
}
|
||||
|
||||
for _, test := range tests {
|
||||
t.Run(test.name, func(t *testing.T) {
|
||||
err := normalizeExecHookError(test.streamErr, test.contextErr, defaultTimeout)
|
||||
if test.expectedError == "" {
|
||||
require.NoError(t, err)
|
||||
return
|
||||
}
|
||||
|
||||
require.EqualError(t, err, test.expectedError)
|
||||
if test.preserveStreamErr && err != test.streamErr {
|
||||
t.Fatalf("expected stream error to be returned unchanged")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestEnsureContainerExists(t *testing.T) {
|
||||
pod := &corev1api.Pod{
|
||||
Spec: corev1api.PodSpec{
|
||||
|
||||
@@ -1,9 +1,26 @@
|
||||
/*
|
||||
Copyright 2026 the Velero contributors.
|
||||
|
||||
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 podexec
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/url"
|
||||
"runtime"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -22,23 +39,38 @@ const timeoutTestPodJSON = `{
|
||||
"spec": {"containers": [{"name": "container-1"}]}
|
||||
}`
|
||||
|
||||
// contextAwareExecutor returns once its context is cancelled, like the SPDY executor does.
|
||||
// contextAwareExecutor returns once its context is canceled, like the SPDY executor does.
|
||||
type contextAwareExecutor struct {
|
||||
cancelled chan struct{}
|
||||
cancelledOnce bool
|
||||
canceled chan struct{}
|
||||
canceledOnce bool
|
||||
}
|
||||
|
||||
func (e *contextAwareExecutor) Stream(options remotecommand.StreamOptions) error { return nil }
|
||||
|
||||
func (e *contextAwareExecutor) StreamWithContext(ctx context.Context, options remotecommand.StreamOptions) error {
|
||||
<-ctx.Done()
|
||||
if !e.cancelledOnce {
|
||||
e.cancelledOnce = true
|
||||
close(e.cancelled)
|
||||
if !e.canceledOnce {
|
||||
e.canceledOnce = true
|
||||
close(e.canceled)
|
||||
}
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
// contextIgnoringExecutor lets the outer timeout path return before the stream does.
|
||||
// Once released, the stream goroutine can only exit if its result channel is buffered.
|
||||
type contextIgnoringExecutor struct {
|
||||
release <-chan struct{}
|
||||
returned *sync.WaitGroup
|
||||
}
|
||||
|
||||
func (e *contextIgnoringExecutor) Stream(options remotecommand.StreamOptions) error { return nil }
|
||||
|
||||
func (e *contextIgnoringExecutor) StreamWithContext(ctx context.Context, options remotecommand.StreamOptions) error {
|
||||
defer e.returned.Done()
|
||||
<-e.release
|
||||
return nil
|
||||
}
|
||||
|
||||
func newTimeoutTestExecutor(t *testing.T, exec remotecommand.Executor) (*defaultPodCommandExecutor, map[string]any) {
|
||||
t.Helper()
|
||||
|
||||
@@ -70,10 +102,10 @@ func timeoutTestHook(timeout time.Duration) *v1.ExecHook {
|
||||
}
|
||||
}
|
||||
|
||||
// A hook that times out must have its exec stream cancelled, otherwise the command keeps
|
||||
// A hook that times out must have its exec stream canceled, otherwise the command keeps
|
||||
// running on the API server after ExecutePodCommand has returned.
|
||||
func TestExecutePodCommandCancelsStreamOnTimeout(t *testing.T) {
|
||||
exec := &contextAwareExecutor{cancelled: make(chan struct{})}
|
||||
exec := &contextAwareExecutor{canceled: make(chan struct{})}
|
||||
podCommandExecutor, pod := newTimeoutTestExecutor(t, exec)
|
||||
|
||||
err := podCommandExecutor.ExecutePodCommand(velerotest.NewLogger(), pod, "ns", "pod-1", "hookName", timeoutTestHook(100*time.Millisecond))
|
||||
@@ -82,26 +114,32 @@ func TestExecutePodCommandCancelsStreamOnTimeout(t *testing.T) {
|
||||
}
|
||||
|
||||
select {
|
||||
case <-exec.cancelled:
|
||||
case <-exec.canceled:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("stream was not cancelled after the hook timed out")
|
||||
t.Fatal("stream was not canceled after the hook timed out")
|
||||
}
|
||||
}
|
||||
|
||||
// When the stream returns because the context expired, both select cases are ready and one
|
||||
// is picked at random, so the reported error must not depend on which one wins.
|
||||
func TestExecutePodCommandTimeoutErrorIsDeterministic(t *testing.T) {
|
||||
const rounds = 50
|
||||
const (
|
||||
rounds = 50
|
||||
expectedError = "timed out after 1ms"
|
||||
)
|
||||
|
||||
messages := map[string]int{}
|
||||
for range rounds {
|
||||
exec := &contextAwareExecutor{cancelled: make(chan struct{})}
|
||||
exec := &contextAwareExecutor{canceled: make(chan struct{})}
|
||||
podCommandExecutor, pod := newTimeoutTestExecutor(t, exec)
|
||||
|
||||
err := podCommandExecutor.ExecutePodCommand(velerotest.NewLogger(), pod, "ns", "pod-1", "hookName", timeoutTestHook(time.Millisecond))
|
||||
if err == nil {
|
||||
t.Fatal("expected a timeout error")
|
||||
}
|
||||
if err.Error() != expectedError {
|
||||
t.Fatalf("expected %q, got %q", expectedError, err)
|
||||
}
|
||||
messages[err.Error()]++
|
||||
}
|
||||
|
||||
@@ -117,8 +155,11 @@ func TestExecutePodCommandDoesNotLeakOnTimeout(t *testing.T) {
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
before := runtime.NumGoroutine()
|
||||
|
||||
release := make(chan struct{})
|
||||
returned := &sync.WaitGroup{}
|
||||
for range rounds {
|
||||
exec := &contextAwareExecutor{cancelled: make(chan struct{})}
|
||||
returned.Add(1)
|
||||
exec := &contextIgnoringExecutor{release: release, returned: returned}
|
||||
podCommandExecutor, pod := newTimeoutTestExecutor(t, exec)
|
||||
|
||||
if err := podCommandExecutor.ExecutePodCommand(velerotest.NewLogger(), pod, "ns", "pod-1", "hookName", timeoutTestHook(50*time.Millisecond)); err == nil {
|
||||
@@ -126,7 +167,11 @@ func TestExecutePodCommandDoesNotLeakOnTimeout(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
time.Sleep(time.Second)
|
||||
// Every ExecutePodCommand call has already taken the timeout path. Releasing the
|
||||
// streams now forces their goroutines to send into an errCh with no receiver.
|
||||
close(release)
|
||||
returned.Wait()
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
runtime.GC()
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user