Skip to content
Open
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
41 changes: 39 additions & 2 deletions reaper.go
Original file line number Diff line number Diff line change
Expand Up @@ -154,7 +154,9 @@ func (r *reaperSpawner) cleanupLocked() error {
// lookupContainer returns a DockerContainer type with the reaper container in the case
// it's found in the running state, and including the labels for sessionID, reaper, and ryuk.
// It will perform a retry with exponential backoff to allow for the container to be started and
// avoid potential false negatives.
// avoid potential false negatives. A container found in a stopped state is removed and
// reported as errReaperNotFound, so the caller creates a new reaper instead of waiting
// on a container that will never become ready.
func (r *reaperSpawner) lookupContainer(ctx context.Context, sessionID string) (*DockerContainer, error) {
dockerClient, err := NewDockerClientWithOpts(ctx)
if err != nil {
Expand Down Expand Up @@ -194,6 +196,26 @@ func (r *reaperSpawner) lookupContainer(ctx context.Context, sessionID string) (
return nil, fmt.Errorf("found %d reaper containers for session ID %q", len(resp.Items), sessionID)
}

switch state := resp.Items[0].State; state {
case container.StateRunning:
// Continue below and return the container for reuse.
case container.StateCreated, container.StateRestarting:
// The container is on its way up, retry until it is running.
return nil, fmt.Errorf("container not running: state %s", state)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
default:
// Exited, dead, paused or removing: the reaper shuts itself down
// once it has had no clients for its reconnection timeout, so a
// stopped container will never become ready again. Remove what is
// left of it so a new reaper can be created under the same name.
// Auto-removed containers may already be gone or being removed,
// which is fine: creating the new reaper retries on a name conflict.
if _, err := dockerClient.ContainerRemove(ctx, resp.Items[0].ID, client.ContainerRemoveOptions{Force: true}); !isCleanupSafe(err) {
return nil, fmt.Errorf("remove stopped container: %w", err)
}

return nil, backoff.Permanent(errReaperNotFound)
}

r, err := provider.ContainerFromType(ctx, resp.Items[0])
if err != nil {
return nil, fmt.Errorf("from docker: %w", err)
Expand Down Expand Up @@ -316,7 +338,7 @@ func (r *reaperSpawner) reuseOrCreate(ctx context.Context, sessionID string, pro

// Look for an existing reaper created in the same test session but in a
// different test process execution e.g. when running tests in parallel.
container, err := r.lookupContainer(context.Background(), sessionID)
container, err := r.lookupContainer(ctx, sessionID)
if err != nil {
if !errors.Is(err, errReaperNotFound) {
return nil, fmt.Errorf("look up container: %w", err)
Expand Down Expand Up @@ -344,13 +366,28 @@ func (r *reaperSpawner) reuseOrCreate(ctx context.Context, sessionID string, pro
func (r *reaperSpawner) fromContainer(ctx context.Context, sessionID string, provider ReaperProvider, dockerContainer *DockerContainer) (*Reaper, error) {
log.Printf("⏳ Waiting for Reaper %q to be ready", dockerContainer.ID[:8])

// The reaper might have terminated between being looked up and now, e.g.
// because it reached its reconnection timeout with no clients. Waiting on
// a stopped container would take the full startup timeout, so check first
// and report not-found, which triggers a retry that recreates the reaper.
if err := r.isRunning(ctx, dockerContainer); err != nil {
return nil, err
}

// Reusing an existing container so we determine the port from the container's exposed ports.
if err := wait.ForAll(
wait.ForLog("Started"),
wait.ForExposedPort().
WithPollInterval(100*time.Millisecond).
SkipInternalCheck(),
).WaitUntilReady(ctx, dockerContainer); err != nil {
// The reaper can also terminate while we wait for it. Include the
// not-running error, if any, so the retry recreates the reaper
// instead of treating the wait failure as permanent.
if stateErr := r.isRunning(ctx, dockerContainer); stateErr != nil {
err = errors.Join(err, stateErr)
}

return nil, fmt.Errorf("wait for reaper %s: %w", dockerContainer.ID[:8], err)
}

Expand Down
73 changes: 73 additions & 0 deletions reaper_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import (
"github.com/containerd/errdefs"
"github.com/moby/moby/api/types/container"
"github.com/moby/moby/api/types/network"
"github.com/moby/moby/client"
"github.com/stretchr/testify/require"

"github.com/testcontainers/testcontainers-go/internal/config"
Expand Down Expand Up @@ -439,6 +440,73 @@ func Test_RecreateReaperIfTerminated(t *testing.T) {
require.NoError(t, err, "connecting to Reaper should be successful")
}

// Test_RecreateReaperIfStopped tests that a reaper container which still exists
// but is no longer running, e.g. it shut down after its reconnection timeout
// with no clients but was not removed yet, is replaced instead of being waited
// on until the startup timeout expires.
func Test_RecreateReaperIfStopped(t *testing.T) {
reaperDisable(t, false)

SkipIfProviderIsNotHealthy(t)

ctx := context.Background()

provider, err := NewDockerProvider()
require.NoError(t, err)

// Create a stopped container that the lookup identifies as the session's
// reaper: same name and labels, but exited and not auto-removed.
require.NoError(t, provider.PullImage(ctx, alpineImage))

labels := core.DefaultLabels(testSessionID)
labels[core.LabelReaper] = "true"
labels[core.LabelRyuk] = "true"
delete(labels, core.LabelReap)

cli := provider.Client()
created, err := cli.ContainerCreate(ctx, client.ContainerCreateOptions{
Config: &container.Config{
Image: alpineImage,
Cmd: []string{"true"},
Labels: labels,
},
Name: reaperContainerNameFromSessionID(testSessionID),
})
require.NoError(t, err)
t.Cleanup(func() {
if _, err := cli.ContainerRemove(context.Background(), created.ID, client.ContainerRemoveOptions{Force: true}); err != nil && !errdefs.IsNotFound(err) {
require.NoError(t, err)
}
})

_, err = cli.ContainerStart(ctx, created.ID, client.ContainerStartOptions{})
require.NoError(t, err)

require.Eventually(t, func() bool {
inspect, err := cli.ContainerInspect(ctx, created.ID, client.ContainerInspectOptions{})
return err == nil && !inspect.Container.State.Running
}, time.Second*10, time.Millisecond*100, "stopped reaper container should have exited")

// The stopped container must be replaced by a fresh reaper well within
// the startup timeout that waiting on it for readiness would burn.
timeout, cancel := context.WithTimeout(ctx, time.Second*30)
defer cancel()

spawner := &reaperSpawner{}
reaper, err := spawner.reaper(context.WithValue(timeout, core.DockerHostContextKey, provider.host), testSessionID, provider)
cleanupReaper(t, reaper, spawner)
require.NoError(t, err, "creating the Reaper should not error")
require.NotEqual(t, created.ID, reaper.container.GetContainerID(), "expected a new reaper container")

// The stopped container was removed to free up the reaper name.
_, err = cli.ContainerInspect(ctx, created.ID, client.ContainerInspectOptions{})
require.True(t, errdefs.IsNotFound(err), "stopped reaper container should have been removed, got: %v", err)

termSignal, err := reaper.Connect()
cleanupTermSignal(t, termSignal)
require.NoError(t, err, "connecting to Reaper should be successful")
}

func TestReaper_reuseItFromOtherTestProgramUsingDocker(t *testing.T) {
reaperDisable(t, false)

Expand Down Expand Up @@ -633,6 +701,11 @@ func TestSpawnerRetryError(t *testing.T) {
err: fmt.Errorf("foo: %w", context.Canceled),
permanent: false,
},
{
name: "wait failure on a stopped container",
err: fmt.Errorf("foo: %w", errors.Join(errors.New("container exited with code 0"), errdefs.ErrNotFound.WithMessage("container state: exited"))),
permanent: false,
},
{
name: "random error",
err: errors.New("some random error"),
Expand Down
Loading