From 37c4de7626e4b93141ff8292e66abf01ff8dc266 Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Sun, 17 May 2026 14:38:27 -0700 Subject: [PATCH 01/10] fix(envd): self-heal MMDS routing on /init MMDS lookup failure When a user-installed PREROUTING/OUTPUT NAT rule in the same netns shadows the MMDS route (e.g. `iptables -t nat -A PREROUTING -d 169.254.169.254 -p tcp --dport 80 -j REDIRECT`), envd's resume-time token-hash lookup fails and /init returns 401. The user's app keeps working because their app needs the redirect. On MMDS lookup failure only (no happy-path cost), install a private nat chain (E2B_MMDS) with a RETURN for 169.254.169.254:80 and re-pin its jump at position 1 in PREROUTING and OUTPUT. Retry the MMDS lookup once. Self-healing rather than always-on: every /init that hits the slow path puts our jump back at position 1, so subsequent resumes recover even if the customer template still installs the conflicting rule. --- packages/envd/internal/api/init.go | 7 ++++ .../envd/internal/host/mmds_route_linux.go | 32 +++++++++++++++++++ .../envd/internal/host/mmds_route_other.go | 7 ++++ 3 files changed, 46 insertions(+) create mode 100644 packages/envd/internal/host/mmds_route_linux.go create mode 100644 packages/envd/internal/host/mmds_route_other.go diff --git a/packages/envd/internal/api/init.go b/packages/envd/internal/api/init.go index 4d0081b0a2..1e1967201e 100644 --- a/packages/envd/internal/api/init.go +++ b/packages/envd/internal/api/init.go @@ -71,6 +71,13 @@ func (a *API) checkMMDSHash(ctx context.Context, requestToken *SecureToken) (boo } mmdsHash, err := a.mmdsClient.GetAccessTokenHash(ctx) + if err != nil { + // Self-heal: a user-installed PREROUTING/OUTPUT redirect on 169.254.169.254:80 + // in the same netns can shadow our route. Reinstall our private chain + // at position 1 and retry once. + host.PinMMDSRoute(ctx) + mmdsHash, err = a.mmdsClient.GetAccessTokenHash(ctx) + } if err != nil { return false, false } diff --git a/packages/envd/internal/host/mmds_route_linux.go b/packages/envd/internal/host/mmds_route_linux.go new file mode 100644 index 0000000000..efc107f119 --- /dev/null +++ b/packages/envd/internal/host/mmds_route_linux.go @@ -0,0 +1,32 @@ +//go:build linux + +package host + +import ( + "context" + "os/exec" +) + +// PinMMDSRoute installs a private nat chain (E2B_MMDS) that returns early +// for MMDS traffic (169.254.169.254:80) and pins the jump at position 1 in +// PREROUTING and OUTPUT, removing any prior copies first. Idempotent. +// +// Intended for the self-heal path: only called when a real MMDS lookup +// fails, on the assumption that user iptables in the same netns clobbered +// our route. Re-running puts our jump back at the front of both chains. +func PinMMDSRoute(ctx context.Context) { + commands := [][]string{ + {"-N", "E2B_MMDS"}, + {"-F", "E2B_MMDS"}, + {"-A", "E2B_MMDS", "-d", "169.254.169.254", "-p", "tcp", "--dport", "80", "-j", "RETURN"}, + {"-D", "PREROUTING", "-j", "E2B_MMDS"}, + {"-I", "PREROUTING", "1", "-j", "E2B_MMDS"}, + {"-D", "OUTPUT", "-j", "E2B_MMDS"}, + {"-I", "OUTPUT", "1", "-j", "E2B_MMDS"}, + } + for _, args := range commands { + // -N fails when chain exists; -D fails when jump is absent. Both + // are expected on the first / clean-state run; swallow errors. + _ = exec.CommandContext(ctx, "iptables", append([]string{"-t", "nat"}, args...)...).Run() + } +} diff --git a/packages/envd/internal/host/mmds_route_other.go b/packages/envd/internal/host/mmds_route_other.go new file mode 100644 index 0000000000..ebaf7fc374 --- /dev/null +++ b/packages/envd/internal/host/mmds_route_other.go @@ -0,0 +1,7 @@ +//go:build !linux + +package host + +import "context" + +func PinMMDSRoute(_ context.Context) {} From b425dcfdbf0c176e9e42594e0a4b84e20eed466e Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Sun, 17 May 2026 14:45:21 -0700 Subject: [PATCH 02/10] use direct PREROUTING/OUTPUT RETURN rules for MMDS pin Switch from a private E2B_MMDS chain (jump + RETURN) to inserting the RETURN rule directly into nat PREROUTING/OUTPUT at position 1, matching the rules known to work. The chain approach was not equivalent: RETURN from a user-defined chain falls through to the next caller rule, so customer DNAT/REDIRECT rules at PREROUTING[2:] would still match. --- .../envd/internal/host/mmds_route_linux.go | 27 +++++++------------ 1 file changed, 10 insertions(+), 17 deletions(-) diff --git a/packages/envd/internal/host/mmds_route_linux.go b/packages/envd/internal/host/mmds_route_linux.go index efc107f119..101d764cc8 100644 --- a/packages/envd/internal/host/mmds_route_linux.go +++ b/packages/envd/internal/host/mmds_route_linux.go @@ -7,26 +7,19 @@ import ( "os/exec" ) -// PinMMDSRoute installs a private nat chain (E2B_MMDS) that returns early -// for MMDS traffic (169.254.169.254:80) and pins the jump at position 1 in -// PREROUTING and OUTPUT, removing any prior copies first. Idempotent. +// PinMMDSRoute pins a RETURN rule for MMDS traffic (169.254.169.254:80) at +// position 1 of nat PREROUTING and OUTPUT. Idempotent: each run deletes any +// existing copy of the rule first, then re-inserts at position 1, so user +// rules added above ours get pushed down. // // Intended for the self-heal path: only called when a real MMDS lookup // fails, on the assumption that user iptables in the same netns clobbered -// our route. Re-running puts our jump back at the front of both chains. +// our route. func PinMMDSRoute(ctx context.Context) { - commands := [][]string{ - {"-N", "E2B_MMDS"}, - {"-F", "E2B_MMDS"}, - {"-A", "E2B_MMDS", "-d", "169.254.169.254", "-p", "tcp", "--dport", "80", "-j", "RETURN"}, - {"-D", "PREROUTING", "-j", "E2B_MMDS"}, - {"-I", "PREROUTING", "1", "-j", "E2B_MMDS"}, - {"-D", "OUTPUT", "-j", "E2B_MMDS"}, - {"-I", "OUTPUT", "1", "-j", "E2B_MMDS"}, - } - for _, args := range commands { - // -N fails when chain exists; -D fails when jump is absent. Both - // are expected on the first / clean-state run; swallow errors. - _ = exec.CommandContext(ctx, "iptables", append([]string{"-t", "nat"}, args...)...).Run() + rule := []string{"-d", "169.254.169.254", "-p", "tcp", "--dport", "80", "-j", "RETURN"} + for _, chain := range []string{"PREROUTING", "OUTPUT"} { + // -D fails when the rule is absent; expected on first run. Swallow. + _ = exec.CommandContext(ctx, "iptables", append([]string{"-t", "nat", "-D", chain}, rule...)...).Run() + _ = exec.CommandContext(ctx, "iptables", append([]string{"-t", "nat", "-I", chain, "1"}, rule...)...).Run() } } From f3796b6e0e952ce2c9e8654bded10485b422da11 Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Sun, 17 May 2026 14:52:22 -0700 Subject: [PATCH 03/10] add -w to iptables for xtables lock contention Gemini review: a user iptables process can hold the xtables lock when our self-heal fires; -w 5 makes us wait up to 5s instead of failing immediately on EAGAIN. --- packages/envd/internal/host/mmds_route_linux.go | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) diff --git a/packages/envd/internal/host/mmds_route_linux.go b/packages/envd/internal/host/mmds_route_linux.go index 101d764cc8..abf76808fc 100644 --- a/packages/envd/internal/host/mmds_route_linux.go +++ b/packages/envd/internal/host/mmds_route_linux.go @@ -19,7 +19,15 @@ func PinMMDSRoute(ctx context.Context) { rule := []string{"-d", "169.254.169.254", "-p", "tcp", "--dport", "80", "-j", "RETURN"} for _, chain := range []string{"PREROUTING", "OUTPUT"} { // -D fails when the rule is absent; expected on first run. Swallow. - _ = exec.CommandContext(ctx, "iptables", append([]string{"-t", "nat", "-D", chain}, rule...)...).Run() - _ = exec.CommandContext(ctx, "iptables", append([]string{"-t", "nat", "-I", chain, "1"}, rule...)...).Run() + run(ctx, append([]string{"-D", chain}, rule...)...) + run(ctx, append([]string{"-I", chain, "1"}, rule...)...) } } + +// run executes iptables in the nat table with -w to wait for the xtables +// lock (a user iptables process may race us). Errors are intentionally +// swallowed; this is best-effort self-heal. +func run(ctx context.Context, args ...string) { + full := append([]string{"-w", "5", "-t", "nat"}, args...) + _ = exec.CommandContext(ctx, "iptables", full...).Run() +} From b139b814d5e72e1cc2b500764f9001fad8f843c8 Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Sun, 17 May 2026 15:18:10 -0700 Subject: [PATCH 04/10] bump envd to 0.5.24 + correct self-heal comment - pkg/version.go: 0.5.23 -> 0.5.24 (new behavioral path on /init). - init.go: comment referenced the private chain that was removed in b425dcfd; describe the direct PREROUTING/OUTPUT re-pin instead. --- packages/envd/internal/api/init.go | 7 ++++--- packages/envd/pkg/version.go | 2 +- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/packages/envd/internal/api/init.go b/packages/envd/internal/api/init.go index 1e1967201e..d267554dd7 100644 --- a/packages/envd/internal/api/init.go +++ b/packages/envd/internal/api/init.go @@ -72,9 +72,10 @@ func (a *API) checkMMDSHash(ctx context.Context, requestToken *SecureToken) (boo mmdsHash, err := a.mmdsClient.GetAccessTokenHash(ctx) if err != nil { - // Self-heal: a user-installed PREROUTING/OUTPUT redirect on 169.254.169.254:80 - // in the same netns can shadow our route. Reinstall our private chain - // at position 1 and retry once. + // Self-heal: a user-installed PREROUTING/OUTPUT redirect on + // 169.254.169.254:80 in the same netns can shadow our route. + // Re-pin our RETURN rule at position 1 of nat PREROUTING and + // OUTPUT, then retry once. host.PinMMDSRoute(ctx) mmdsHash, err = a.mmdsClient.GetAccessTokenHash(ctx) } diff --git a/packages/envd/pkg/version.go b/packages/envd/pkg/version.go index 88bc32c769..4cfbb08b29 100644 --- a/packages/envd/pkg/version.go +++ b/packages/envd/pkg/version.go @@ -1,3 +1,3 @@ package pkg -const Version = "0.5.23" +const Version = "0.5.24" From 361598df98ec5524b22cd04d44aff563b6775ea7 Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Sun, 17 May 2026 15:20:19 -0700 Subject: [PATCH 05/10] surface iptables -I failures from MMDS pin Previously the self-heal swallowed every iptables error, which hid real breakage (e.g. iptables binary missing, broken iptables module load, xtables lock contention beyond -w). Now PinMMDSRoute returns the first -I failure with the iptables stderr attached, and the caller logs at warn. -D failures stay silent: "rule absent" is the expected first-run case. --- packages/envd/internal/api/init.go | 4 ++- .../envd/internal/host/mmds_route_linux.go | 31 +++++++++++++------ .../envd/internal/host/mmds_route_other.go | 2 +- 3 files changed, 25 insertions(+), 12 deletions(-) diff --git a/packages/envd/internal/api/init.go b/packages/envd/internal/api/init.go index d267554dd7..db6f4c44f9 100644 --- a/packages/envd/internal/api/init.go +++ b/packages/envd/internal/api/init.go @@ -76,7 +76,9 @@ func (a *API) checkMMDSHash(ctx context.Context, requestToken *SecureToken) (boo // 169.254.169.254:80 in the same netns can shadow our route. // Re-pin our RETURN rule at position 1 of nat PREROUTING and // OUTPUT, then retry once. - host.PinMMDSRoute(ctx) + if pinErr := host.PinMMDSRoute(ctx); pinErr != nil { + a.logger.Warn().Err(pinErr).Msg("failed to pin MMDS iptables route") + } mmdsHash, err = a.mmdsClient.GetAccessTokenHash(ctx) } if err != nil { diff --git a/packages/envd/internal/host/mmds_route_linux.go b/packages/envd/internal/host/mmds_route_linux.go index abf76808fc..5810e4fcb9 100644 --- a/packages/envd/internal/host/mmds_route_linux.go +++ b/packages/envd/internal/host/mmds_route_linux.go @@ -4,6 +4,7 @@ package host import ( "context" + "fmt" "os/exec" ) @@ -14,20 +15,30 @@ import ( // // Intended for the self-heal path: only called when a real MMDS lookup // fails, on the assumption that user iptables in the same netns clobbered -// our route. -func PinMMDSRoute(ctx context.Context) { +// our route. Returns the first -I failure (if any); -D failures are +// expected (rule absent on first run) and silently swallowed. +func PinMMDSRoute(ctx context.Context) error { rule := []string{"-d", "169.254.169.254", "-p", "tcp", "--dport", "80", "-j", "RETURN"} for _, chain := range []string{"PREROUTING", "OUTPUT"} { - // -D fails when the rule is absent; expected on first run. Swallow. - run(ctx, append([]string{"-D", chain}, rule...)...) - run(ctx, append([]string{"-I", chain, "1"}, rule...)...) + // -D fails when the rule is absent (exit 1, expected on first run); + // nothing actionable to log. + _ = iptables(ctx, append([]string{"-D", chain}, rule...)...) + if err := iptables(ctx, append([]string{"-I", chain, "1"}, rule...)...); err != nil { + return fmt.Errorf("iptables -I nat %s: %w", chain, err) + } } + + return nil } -// run executes iptables in the nat table with -w to wait for the xtables -// lock (a user iptables process may race us). Errors are intentionally -// swallowed; this is best-effort self-heal. -func run(ctx context.Context, args ...string) { +// iptables runs `iptables -w 5 -t nat ...`. -w waits up to 5s for the +// xtables lock (a user iptables process may race us). +func iptables(ctx context.Context, args ...string) error { full := append([]string{"-w", "5", "-t", "nat"}, args...) - _ = exec.CommandContext(ctx, "iptables", full...).Run() + out, err := exec.CommandContext(ctx, "iptables", full...).CombinedOutput() + if err != nil { + return fmt.Errorf("%w: %s", err, out) + } + + return nil } diff --git a/packages/envd/internal/host/mmds_route_other.go b/packages/envd/internal/host/mmds_route_other.go index ebaf7fc374..b009d37228 100644 --- a/packages/envd/internal/host/mmds_route_other.go +++ b/packages/envd/internal/host/mmds_route_other.go @@ -4,4 +4,4 @@ package host import "context" -func PinMMDSRoute(_ context.Context) {} +func PinMMDSRoute(_ context.Context) error { return nil } From f7256d4d9079d4765e6fe6ff38d011769b57b3f3 Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Sun, 17 May 2026 15:24:47 -0700 Subject: [PATCH 06/10] rate-limit MMDS pin failure warn to 1/10s Drop the in-pin cooldown (we want self-heal to fire on every retry in case user rules keep landing back on top) and instead floor the warn log to once per 10s with a suppressed-since-last counter, mirroring internal/logs/exporter/rate_limited_logger.go. --- packages/envd/internal/api/init.go | 31 ++++++++++++++++++- .../envd/internal/host/mmds_route_linux.go | 5 ++- 2 files changed, 32 insertions(+), 4 deletions(-) diff --git a/packages/envd/internal/api/init.go b/packages/envd/internal/api/init.go index db6f4c44f9..c3b5c9c35d 100644 --- a/packages/envd/internal/api/init.go +++ b/packages/envd/internal/api/init.go @@ -10,6 +10,7 @@ import ( "net/netip" "os/exec" "strings" + "sync/atomic" "time" "github.com/awnumar/memguard" @@ -22,6 +23,34 @@ import ( "github.com/e2b-dev/infra/packages/shared/pkg/keys" ) +// rateLimitedWarn floors a recurring warn to at most one line per `floor`, +// counting how many were suppressed in between. /init is hammered by the +// orchestrator's infinite retry loop, so a persistent failure here would +// otherwise flood the log. +type rateLimitedWarn struct { + floor time.Duration + lastLogged atomic.Pointer[time.Time] + suppressed atomic.Int64 +} + +func (r *rateLimitedWarn) log(logger *zerolog.Logger, err error, msg string) { + last := r.lastLogged.Load() + if last != nil && time.Since(*last) <= r.floor { + r.suppressed.Add(1) + + return + } + now := time.Now() + if !r.lastLogged.CompareAndSwap(last, &now) { + r.suppressed.Add(1) + + return + } + logger.Warn().Err(err).Int64("suppressed", r.suppressed.Swap(0)).Msg(msg) +} + +var pinMMDSWarn = &rateLimitedWarn{floor: 10 * time.Second} + var ( ErrAccessTokenMismatch = errors.New("access token validation failed") ErrAccessTokenResetNotAuthorized = errors.New("access token reset not authorized") @@ -77,7 +106,7 @@ func (a *API) checkMMDSHash(ctx context.Context, requestToken *SecureToken) (boo // Re-pin our RETURN rule at position 1 of nat PREROUTING and // OUTPUT, then retry once. if pinErr := host.PinMMDSRoute(ctx); pinErr != nil { - a.logger.Warn().Err(pinErr).Msg("failed to pin MMDS iptables route") + pinMMDSWarn.log(a.logger, pinErr, "failed to pin MMDS iptables route") } mmdsHash, err = a.mmdsClient.GetAccessTokenHash(ctx) } diff --git a/packages/envd/internal/host/mmds_route_linux.go b/packages/envd/internal/host/mmds_route_linux.go index 5810e4fcb9..fd0cac7560 100644 --- a/packages/envd/internal/host/mmds_route_linux.go +++ b/packages/envd/internal/host/mmds_route_linux.go @@ -14,9 +14,8 @@ import ( // rules added above ours get pushed down. // // Intended for the self-heal path: only called when a real MMDS lookup -// fails, on the assumption that user iptables in the same netns clobbered -// our route. Returns the first -I failure (if any); -D failures are -// expected (rule absent on first run) and silently swallowed. +// fails. Returns the first -I failure (if any); -D failures are expected +// (rule absent on first run) and silently swallowed. func PinMMDSRoute(ctx context.Context) error { rule := []string{"-d", "169.254.169.254", "-p", "tcp", "--dport", "80", "-j", "RETURN"} for _, chain := range []string{"PREROUTING", "OUTPUT"} { From bd2a1278fd00d9684b21cfa86c16b55202fbb9f1 Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Sun, 17 May 2026 15:29:47 -0700 Subject: [PATCH 07/10] reuse shared rate limiter; bound MMDS poller per-tick - Extract the gate from internal/logs/exporter/rate_limited_logger.go into internal/logs/ratelimit.Limiter (logger-agnostic; returns (allowed, suppressedSinceLast)). Exporter wraps it for log.Printf, init.go uses it directly with zerolog. - Add Timeout to PollForMMDSOpts' http.Client (matches the existing mmdsAccessTokenClient 10s). Without it a -j DROP on MMDS would hang each tick on the TCP handshake until the parent ctx expired. --- packages/envd/internal/api/init.go | 36 ++++------------- packages/envd/internal/host/mmds.go | 8 +++- .../logs/exporter/rate_limited_logger.go | 24 ++++-------- .../envd/internal/logs/ratelimit/ratelimit.go | 39 +++++++++++++++++++ 4 files changed, 60 insertions(+), 47 deletions(-) create mode 100644 packages/envd/internal/logs/ratelimit/ratelimit.go diff --git a/packages/envd/internal/api/init.go b/packages/envd/internal/api/init.go index c3b5c9c35d..2efab03fd5 100644 --- a/packages/envd/internal/api/init.go +++ b/packages/envd/internal/api/init.go @@ -10,7 +10,6 @@ import ( "net/netip" "os/exec" "strings" - "sync/atomic" "time" "github.com/awnumar/memguard" @@ -20,36 +19,13 @@ import ( "github.com/e2b-dev/infra/packages/envd/internal/host" "github.com/e2b-dev/infra/packages/envd/internal/logs" + "github.com/e2b-dev/infra/packages/envd/internal/logs/ratelimit" "github.com/e2b-dev/infra/packages/shared/pkg/keys" ) -// rateLimitedWarn floors a recurring warn to at most one line per `floor`, -// counting how many were suppressed in between. /init is hammered by the -// orchestrator's infinite retry loop, so a persistent failure here would -// otherwise flood the log. -type rateLimitedWarn struct { - floor time.Duration - lastLogged atomic.Pointer[time.Time] - suppressed atomic.Int64 -} - -func (r *rateLimitedWarn) log(logger *zerolog.Logger, err error, msg string) { - last := r.lastLogged.Load() - if last != nil && time.Since(*last) <= r.floor { - r.suppressed.Add(1) - - return - } - now := time.Now() - if !r.lastLogged.CompareAndSwap(last, &now) { - r.suppressed.Add(1) - - return - } - logger.Warn().Err(err).Int64("suppressed", r.suppressed.Swap(0)).Msg(msg) -} - -var pinMMDSWarn = &rateLimitedWarn{floor: 10 * time.Second} +// /init is hammered by the orchestrator's infinite retry loop, so a +// persistent pin failure would otherwise flood the log. +var pinMMDSWarnLimit = ratelimit.New(10 * time.Second) var ( ErrAccessTokenMismatch = errors.New("access token validation failed") @@ -106,7 +82,9 @@ func (a *API) checkMMDSHash(ctx context.Context, requestToken *SecureToken) (boo // Re-pin our RETURN rule at position 1 of nat PREROUTING and // OUTPUT, then retry once. if pinErr := host.PinMMDSRoute(ctx); pinErr != nil { - pinMMDSWarn.log(a.logger, pinErr, "failed to pin MMDS iptables route") + if ok, suppressed := pinMMDSWarnLimit.Allow(); ok { + a.logger.Warn().Err(pinErr).Int64("suppressed", suppressed).Msg("failed to pin MMDS iptables route") + } } mmdsHash, err = a.mmdsClient.GetAccessTokenHash(ctx) } diff --git a/packages/envd/internal/host/mmds.go b/packages/envd/internal/host/mmds.go index e24cc6ba92..aa29cb85a3 100644 --- a/packages/envd/internal/host/mmds.go +++ b/packages/envd/internal/host/mmds.go @@ -130,7 +130,13 @@ func GetAccessTokenHashFromMMDS(ctx context.Context) (string, error) { } func PollForMMDSOpts(ctx context.Context, mmdsChan chan<- *MMDSOpts, envVars *utils.Map[string, string]) { - httpClient := &http.Client{Transport: &http.Transport{DisableKeepAlives: true}} + // Match mmdsAccessTokenClient: bound any single tick (e.g. -j DROP on + // MMDS would otherwise hang on the TCP handshake) and avoid keepalive + // so a broken intermediate doesn't poison a kept-open connection. + httpClient := &http.Client{ + Timeout: mmdsAccessTokenRequestClientTimeout, + Transport: &http.Transport{DisableKeepAlives: true}, + } defer httpClient.CloseIdleConnections() var lastErr error diff --git a/packages/envd/internal/logs/exporter/rate_limited_logger.go b/packages/envd/internal/logs/exporter/rate_limited_logger.go index 22500e2948..1845d10320 100644 --- a/packages/envd/internal/logs/exporter/rate_limited_logger.go +++ b/packages/envd/internal/logs/exporter/rate_limited_logger.go @@ -2,34 +2,24 @@ package exporter import ( "log" - "sync/atomic" "time" + + "github.com/e2b-dev/infra/packages/envd/internal/logs/ratelimit" ) type rateLimitedLogger struct { - floor time.Duration - format string - lastLogged atomic.Pointer[time.Time] - suppressed atomic.Int64 + limit *ratelimit.Limiter + format string } func newRateLimitedLogger(floor time.Duration, format string) *rateLimitedLogger { - return &rateLimitedLogger{floor: floor, format: format} + return &rateLimitedLogger{limit: ratelimit.New(floor), format: format} } func (r *rateLimitedLogger) log(args ...any) { - last := r.lastLogged.Load() - if last != nil && time.Since(*last) <= r.floor { - r.suppressed.Add(1) - - return - } - now := time.Now() - if !r.lastLogged.CompareAndSwap(last, &now) { - r.suppressed.Add(1) - + ok, suppressed := r.limit.Allow() + if !ok { return } - suppressed := r.suppressed.Swap(0) log.Printf(r.format+" (%d suppressed since last log)", append(args, suppressed)...) } diff --git a/packages/envd/internal/logs/ratelimit/ratelimit.go b/packages/envd/internal/logs/ratelimit/ratelimit.go new file mode 100644 index 0000000000..78e36b9b28 --- /dev/null +++ b/packages/envd/internal/logs/ratelimit/ratelimit.go @@ -0,0 +1,39 @@ +package ratelimit + +import ( + "sync/atomic" + "time" +) + +// Limiter gates a recurring log to at most one emit per `floor`, counting +// suppressed attempts in between. The caller decides how to format/emit; +// this type only owns the gating decision. +type Limiter struct { + floor time.Duration + lastLogged atomic.Pointer[time.Time] + suppressed atomic.Int64 +} + +func New(floor time.Duration) *Limiter { + return &Limiter{floor: floor} +} + +// Allow returns (true, suppressedSinceLast) when the caller should emit a +// log line; false otherwise. On true the caller should include +// `suppressedSinceLast` in the emitted message. +func (r *Limiter) Allow() (bool, int64) { + last := r.lastLogged.Load() + if last != nil && time.Since(*last) <= r.floor { + r.suppressed.Add(1) + + return false, 0 + } + now := time.Now() + if !r.lastLogged.CompareAndSwap(last, &now) { + r.suppressed.Add(1) + + return false, 0 + } + + return true, r.suppressed.Swap(0) +} From c6ee5e0bd228050f04c2ca187d3f25b45f99c1e4 Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Mon, 18 May 2026 12:12:29 -0700 Subject: [PATCH 08/10] serialize MMDS pin to avoid parallel iptables runs Coalesce concurrent /init self-heal calls with an atomic.Bool TryLock so we don't fire iptables mutations in parallel against the same nat table. Same pattern as isMountingNFS. --- packages/envd/internal/host/mmds_route_linux.go | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/packages/envd/internal/host/mmds_route_linux.go b/packages/envd/internal/host/mmds_route_linux.go index fd0cac7560..fe24c62b50 100644 --- a/packages/envd/internal/host/mmds_route_linux.go +++ b/packages/envd/internal/host/mmds_route_linux.go @@ -6,17 +6,29 @@ import ( "context" "fmt" "os/exec" + "sync/atomic" ) +// pinMMDSInflight prevents concurrent /init retries from running iptables +// in parallel against the same nat table. +var pinMMDSInflight atomic.Bool + // PinMMDSRoute pins a RETURN rule for MMDS traffic (169.254.169.254:80) at // position 1 of nat PREROUTING and OUTPUT. Idempotent: each run deletes any // existing copy of the rule first, then re-inserts at position 1, so user // rules added above ours get pushed down. // // Intended for the self-heal path: only called when a real MMDS lookup -// fails. Returns the first -I failure (if any); -D failures are expected -// (rule absent on first run) and silently swallowed. +// fails. Concurrent callers are coalesced — only one runs at a time, the +// rest return nil immediately. Returns the first -I failure (if any); +// -D failures are expected (rule absent on first run) and silently +// swallowed. func PinMMDSRoute(ctx context.Context) error { + if !pinMMDSInflight.CompareAndSwap(false, true) { + return nil + } + defer pinMMDSInflight.Store(false) + rule := []string{"-d", "169.254.169.254", "-p", "tcp", "--dport", "80", "-j", "RETURN"} for _, chain := range []string{"PREROUTING", "OUTPUT"} { // -D fails when the rule is absent (exit 1, expected on first run); From 9232385bc3f4960cb54748af2b01e4588a76c5f6 Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Mon, 18 May 2026 12:15:38 -0700 Subject: [PATCH 09/10] use semaphore.Weighted for MMDS pin coalescing --- .../envd/internal/host/mmds_route_linux.go | 21 ++++++++++--------- 1 file changed, 11 insertions(+), 10 deletions(-) diff --git a/packages/envd/internal/host/mmds_route_linux.go b/packages/envd/internal/host/mmds_route_linux.go index fe24c62b50..170e877247 100644 --- a/packages/envd/internal/host/mmds_route_linux.go +++ b/packages/envd/internal/host/mmds_route_linux.go @@ -6,12 +6,13 @@ import ( "context" "fmt" "os/exec" - "sync/atomic" + + "golang.org/x/sync/semaphore" ) -// pinMMDSInflight prevents concurrent /init retries from running iptables -// in parallel against the same nat table. -var pinMMDSInflight atomic.Bool +// pinMMDSSem serializes self-heal calls so concurrent /init retries don't +// run iptables in parallel against the same nat table. +var pinMMDSSem = semaphore.NewWeighted(1) // PinMMDSRoute pins a RETURN rule for MMDS traffic (169.254.169.254:80) at // position 1 of nat PREROUTING and OUTPUT. Idempotent: each run deletes any @@ -19,15 +20,15 @@ var pinMMDSInflight atomic.Bool // rules added above ours get pushed down. // // Intended for the self-heal path: only called when a real MMDS lookup -// fails. Concurrent callers are coalesced — only one runs at a time, the -// rest return nil immediately. Returns the first -I failure (if any); -// -D failures are expected (rule absent on first run) and silently -// swallowed. +// fails. Concurrent callers are coalesced via a semaphore — only one runs +// at a time, the rest return nil immediately. Returns the first -I failure +// (if any); -D failures are expected (rule absent on first run) and +// silently swallowed. func PinMMDSRoute(ctx context.Context) error { - if !pinMMDSInflight.CompareAndSwap(false, true) { + if !pinMMDSSem.TryAcquire(1) { return nil } - defer pinMMDSInflight.Store(false) + defer pinMMDSSem.Release(1) rule := []string{"-d", "169.254.169.254", "-p", "tcp", "--dport", "80", "-j", "RETURN"} for _, chain := range []string{"PREROUTING", "OUTPUT"} { From 6589bfb62ae944138e11ec5726035b0fb82011ad Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Petr=20Van=C4=9Bk?= Date: Tue, 19 May 2026 09:21:13 +0200 Subject: [PATCH 10/10] chore(envd): bump version to 0.5.25 --- packages/envd/pkg/version.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/envd/pkg/version.go b/packages/envd/pkg/version.go index 4cfbb08b29..48aeaf4b70 100644 --- a/packages/envd/pkg/version.go +++ b/packages/envd/pkg/version.go @@ -1,3 +1,3 @@ package pkg -const Version = "0.5.24" +const Version = "0.5.25"