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
6 changes: 4 additions & 2 deletions pkg/container/docker/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -1638,9 +1638,11 @@ func mergeEnvVars(base, extra map[string]string) map[string]string {
// the workload and returns the host-side port it is bound on.
func (c *Client) setupIngressContainer(ctx context.Context, containerName string, upstreamPort int, attachStdio bool,
externalEndpointsConfig map[string]*network.EndpointSettings, networkPermissions *permissions.NetworkPermissions) (int, error) {
squidPort, err := networking.FindOrUsePort(upstreamPort + 1)
// A random port avoids every same-image workload racing to bind the same
// preferred port when starting concurrently (see #6063).
squidPort, err := networking.FindOrUsePort(0)
if err != nil {
return 0, fmt.Errorf("failed to find or use port %d: %w", squidPort, err)
return 0, fmt.Errorf("failed to find an available port: %w", err)
}
squidExposedPorts := map[string]struct{}{
fmt.Sprintf("%d/tcp", squidPort): {},
Expand Down
4 changes: 3 additions & 1 deletion pkg/container/docker/envoy.go
Original file line number Diff line number Diff line change
Expand Up @@ -886,7 +886,9 @@ func (e *envoyProxy) SetupIngress(ctx context.Context, spec proxySpec, _ egressR

var ingressPort int
if spec.TransportType != "stdio" && spec.UpstreamPort > 0 {
port, err := networking.FindOrUsePort(spec.UpstreamPort + 1)
// A random port avoids every same-image workload racing to bind the
// same preferred port when starting concurrently (see #6063).
port, err := networking.FindOrUsePort(0)
if err != nil {
return 0, fmt.Errorf("failed to find ingress port: %w", err)
}
Expand Down
66 changes: 66 additions & 0 deletions pkg/container/docker/envoy_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,10 +6,12 @@ package docker
import (
"context"
"encoding/json"
"fmt"
"os"
"os/exec"
"regexp"
"strings"
"sync"
"testing"
"time"

Expand Down Expand Up @@ -821,6 +823,70 @@ func TestEnvoyProxy_SetupOrchestration(t *testing.T) {
}
}

// TestEnvoyProxy_SetupIngress_ConcurrentSameUpstreamPort guards against a
// deterministic ingress port collision: two workloads built from the same
// image share the same fixed UpstreamPort, so deriving the ingress port from
// UpstreamPort+1 made every concurrently-starting instance request the
// identical host port. The fix picks a random port instead, so concurrent
// SetupIngress calls sharing UpstreamPort must land on different ports.
func TestEnvoyProxy_SetupIngress_ConcurrentSameUpstreamPort(t *testing.T) {
t.Parallel()

const upstreamPort = 8080
const workloadCount = 4

ports := make([]int, workloadCount)
errs := make([]error, workloadCount)

var wg sync.WaitGroup
for i := range workloadCount {
wg.Add(1)
go func(idx int) {
defer wg.Done()
api := &fakeDockerAPI{
createFunc: func(_ context.Context, _ *container.Config, _ *container.HostConfig,
_ *network.NetworkingConfig, _ *v1.Platform, _ string) (container.CreateResponse, error) {
return container.CreateResponse{ID: "cid-envoy"}, nil
},
startFunc: func(_ context.Context, _ string, _ mobyclient.ContainerStartOptions) error { return nil },
}
e := &envoyProxy{client: &Client{
api: api,
imageManager: &fakeImageManager{availableImages: map[string]struct{}{defaultEnvoyImage: {}}},
}}
spec := proxySpec{
WorkloadName: fmt.Sprintf("app-%d", idx),
TransportType: "streamable-http",
UpstreamPort: upstreamPort,
GatewayIP: dockerDefaultBridgeGatewayIP,
Endpoints: map[string]*network.EndpointSettings{},
}
port, err := e.SetupIngress(t.Context(), spec, egressResult{})
ports[idx] = port
errs[idx] = err
}(i)
}

done := make(chan struct{})
go func() {
wg.Wait()
close(done)
}()
select {
case <-done:
case <-time.After(10 * time.Second):
t.Fatal("timeout waiting for concurrent SetupIngress calls")
}

seen := make(map[int]bool, workloadCount)
for i, port := range ports {
require.NoError(t, errs[i])
assert.NotEqual(t, upstreamPort+1, port, "ingress port must not be derived from UpstreamPort+1")
assert.False(t, seen[port], "port %d was reserved by more than one concurrent workload", port)
seen[port] = true
}
}

// TestListenerTimeoutsDisabledForStreaming guards against Envoy's 15s default
// RouteAction.timeout (a hard total-response cap) and 5m stream_idle_timeout
// truncating long-lived MCP streams (SSE / streamable-http). Both the ingress
Expand Down
Loading