Skip to content
Closed
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
2 changes: 1 addition & 1 deletion packages/clickhouse/pkg/hoststats/hoststats.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ type SandboxHostStat struct {
CgroupCPUUserUsec uint64 `ch:"cgroup_cpu_user_usec"` // cumulative, microseconds
CgroupCPUSystemUsec uint64 `ch:"cgroup_cpu_system_usec"` // cumulative, microseconds
CgroupMemoryUsage uint64 `ch:"cgroup_memory_usage_bytes"` // current, bytes
CgroupMemoryPeak uint64 `ch:"cgroup_memory_peak_bytes"` // interval peak, bytes (reset after each sample)
CgroupMemoryPeak uint64 `ch:"cgroup_memory_peak_bytes"` // interval peak, bytes; the sandbox lifetime peak on hosts whose kernel cannot reset memory.peak (Linux < 6.12)

// Pre-computed deltas between consecutive samples.
DeltaCgroupCPUUsageUsec uint64 `ch:"delta_cgroup_cpu_usage_usec"`
Expand Down
64 changes: 58 additions & 6 deletions packages/orchestrator/pkg/sandbox/cgroup/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@ import (
"slices"
"strconv"
"strings"
"sync/atomic"
"syscall"
"time"

"go.uber.org/zap"
Expand All @@ -33,14 +35,29 @@ const (
cgroupKillPollInterval = 100 * time.Millisecond
)

// memoryPeakResetUnsupported latches once the kernel rejects a memory.peak
// reset write: reset support is per-kernel (Linux 6.12+), so one failure
// proves it for every sandbox on this host. See resetMemoryPeak.
var memoryPeakResetUnsupported atomic.Bool

// writeMemoryPeakReset is injectable for tests: a regular file accepts the
// write, so a pre-6.12 kernel's EINVAL cannot be provoked without a real cgroup.
var writeMemoryPeakReset = func(memoryPeakFile *os.File) error {
_, err := memoryPeakFile.WriteString("0")

return err
}

// Stats contains resource usage statistics from a cgroup
type Stats struct {
CPUUsageUsec uint64 // microseconds
CPUUserUsec uint64 // microseconds
CPUSystemUsec uint64 // microseconds

MemoryUsageBytes uint64 // bytes
MemoryPeakBytes uint64 // bytes, reset after each GetStats() call
// MemoryPeakBytes is the peak since the previous GetStats() call — or the
// cgroup's lifetime peak on kernels without memory.peak reset (Linux < 6.12).
MemoryPeakBytes uint64 // bytes
}

// CgroupHandle represents a created cgroup for a sandbox.
Expand Down Expand Up @@ -391,7 +408,10 @@ func (m *managerImpl) Create(ctx context.Context, cgroupName string) (*CgroupHan
return nil, fmt.Errorf("failed to open cgroup directory: %w", err)
}

// O_RDWR FD must stay open for per-FD peak reset across GetStats() calls
// O_RDWR FD must stay open for per-FD peak reset across GetStats() calls.
// A successful open says nothing about reset support: pre-6.12 kernels
// expose memory.peak as 0444, but CAP_DAC_OVERRIDE lets root open it O_RDWR
// anyway, and only the write itself then fails. See resetMemoryPeak.
memPeakPath := filepath.Join(cgroupPath, "memory.peak")
memoryPeakFile, peakErr := os.OpenFile(memPeakPath, os.O_RDWR, 0)
if peakErr != nil {
Expand All @@ -415,7 +435,7 @@ func (m *managerImpl) Create(ctx context.Context, cgroupName string) (*CgroupHan
zap.String("cgroup_name", cgroupName),
zap.String("path", cgroupPath),
zap.Int("fd", handle.GetFD()),
zap.Bool("peak_reset_available", memoryPeakFile != nil))
zap.Bool("memory_peak_open", memoryPeakFile != nil))

return handle, nil
}
Expand Down Expand Up @@ -527,12 +547,44 @@ func (m *managerImpl) readAndResetMemoryPeak(ctx context.Context, memoryPeakFile
return 0, fmt.Errorf("failed to parse memory.peak value %q: %w", strings.TrimSpace(string(buf[:n])), parseErr)
}

// Reset per-FD peak for next interval
if _, err := memoryPeakFile.WriteString("0"); err != nil {
resetMemoryPeak(ctx, memoryPeakFile)

return peakBytes, nil
}

// resetMemoryPeak resets the per-FD peak so the next read covers one interval.
// Failure is never fatal — the peak already read stays valid, just widens to
// span more than one interval. Pre-6.12 kernels have no memory.peak write
// handler and fail every write with EINVAL; that is permanent, so the first
// such failure latches memoryPeakResetUnsupported (logged once, never
// retried). Other errors are transient and retried on the next sample.
func resetMemoryPeak(ctx context.Context, memoryPeakFile *os.File) {
if memoryPeakResetUnsupported.Load() {
return
}

err := writeMemoryPeakReset(memoryPeakFile)
if err == nil {
return
}

if !peakResetUnsupported(err) {
logger.L().Warn(ctx, "failed to reset memory.peak, interval peak semantics degraded", zap.Error(err))

return
}

return peakBytes, nil
if memoryPeakResetUnsupported.CompareAndSwap(false, true) {
logger.L().Warn(ctx, "memory.peak reset unsupported (requires kernel 6.12+); reporting lifetime peak instead of interval peak", zap.Error(err))
}
}

// peakResetUnsupported reports whether err means the kernel does not implement
// the memory.peak reset write at all, as opposed to a transient write failure.
func peakResetUnsupported(err error) bool {
return errors.Is(err, syscall.EINVAL) ||
errors.Is(err, syscall.ENOTSUP) ||
errors.Is(err, syscall.EOPNOTSUPP)
}

// cgroupPath returns the filesystem path for a sandbox's cgroup
Expand Down
220 changes: 220 additions & 0 deletions packages/orchestrator/pkg/sandbox/cgroup/peak_reset_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,220 @@
//go:build linux

package cgroup

import (
"errors"
"fmt"
"os"
"path/filepath"
"sync/atomic"
"syscall"
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.uber.org/zap/zapcore"
"go.uber.org/zap/zaptest/observer"

"github.com/e2b-dev/infra/packages/shared/pkg/logger"
)

const (
fakeCPUUsageUsec uint64 = 123456789
fakeCPUUserUsec uint64 = 100000000
fakeCPUSystemUsec uint64 = 23456789
fakeMemoryCurrent uint64 = 536870912
fakeMemoryPeak uint64 = 805306368
)

// peakWriteError wraps errno the way os.File.WriteString does — inside a
// *os.PathError — so tests exercise the errors.Is unwrap.
func peakWriteError(errno syscall.Errno) error {
return &os.PathError{Op: "write", Path: "/sys/fs/cgroup/e2b/sbx-test/memory.peak", Err: errno}
}

// stubPeakReset stubs the reset write; the process-wide latch outlives any
// single test, so it is cleared before and after.
func stubPeakReset(t *testing.T, writeErr error) *atomic.Int64 {
t.Helper()

calls := &atomic.Int64{}

original := writeMemoryPeakReset
memoryPeakResetUnsupported.Store(false)
writeMemoryPeakReset = func(*os.File) error {
calls.Add(1)

return writeErr
}

t.Cleanup(func() {
writeMemoryPeakReset = original
memoryPeakResetUnsupported.Store(false)
})

return calls
}

func observeWarnings(t *testing.T) *observer.ObservedLogs {
t.Helper()

core, logs := observer.New(zapcore.WarnLevel)
t.Cleanup(logger.ReplaceGlobals(t.Context(), logger.NewTracedLoggerFromCore(core)))

return logs
}

func fakeCgroupDir(t *testing.T) (string, *os.File) {
t.Helper()

cgroupPath := t.TempDir()

for name, content := range map[string]string{
"cpu.stat": fmt.Sprintf("usage_usec %d\nuser_usec %d\nsystem_usec %d\n", fakeCPUUsageUsec, fakeCPUUserUsec, fakeCPUSystemUsec),
"memory.current": fmt.Sprintf("%d\n", fakeMemoryCurrent),
"memory.peak": fmt.Sprintf("%d\n", fakeMemoryPeak),
} {
require.NoError(t, os.WriteFile(filepath.Join(cgroupPath, name), []byte(content), 0o644))
}

memoryPeakFile, err := os.OpenFile(filepath.Join(cgroupPath, "memory.peak"), os.O_RDWR, 0)
require.NoError(t, err)
t.Cleanup(func() { _ = memoryPeakFile.Close() })

return cgroupPath, memoryPeakFile
}

func assertParsedStats(t *testing.T, stats *Stats) {
t.Helper()

require.NotNil(t, stats)
assert.Equal(t, fakeCPUUsageUsec, stats.CPUUsageUsec)
assert.Equal(t, fakeCPUUserUsec, stats.CPUUserUsec)
assert.Equal(t, fakeCPUSystemUsec, stats.CPUSystemUsec)
assert.Equal(t, fakeMemoryCurrent, stats.MemoryUsageBytes)
assert.Equal(t, fakeMemoryPeak, stats.MemoryPeakBytes, "the peak read must be reported regardless of reset support")
}

//nolint:paralleltest // mutates the process-wide reset latch and the global logger
func TestGetStatsPeakResetUnsupportedLatchesAfterOneWarning(t *testing.T) {
calls := stubPeakReset(t, peakWriteError(syscall.EINVAL))
logs := observeWarnings(t)

cgroupPath, memoryPeakFile := fakeCgroupDir(t)
mgr := &managerImpl{}

for range 3 {
stats, err := mgr.getStatsForPath(t.Context(), cgroupPath, memoryPeakFile)
require.NoError(t, err)
assertParsedStats(t, stats)
}

assert.True(t, memoryPeakResetUnsupported.Load(), "EINVAL must latch the reset as unsupported")
assert.Equal(t, int64(1), calls.Load(), "reset must be attempted once, then skipped entirely")

warnings := logs.FilterLevelExact(zapcore.WarnLevel).All()
require.Len(t, warnings, 1, "unsupported reset must be logged exactly once per process")
assert.Contains(t, warnings[0].Message, "kernel 6.12+")

// The latch is per-kernel, so it must hold for a second sandbox's FD too.
otherPath, otherPeakFile := fakeCgroupDir(t)
stats, err := mgr.getStatsForPath(t.Context(), otherPath, otherPeakFile)
require.NoError(t, err)
assertParsedStats(t, stats)
assert.Equal(t, int64(1), calls.Load(), "latch must suppress the reset for every sandbox")
assert.Len(t, logs.FilterLevelExact(zapcore.WarnLevel).All(), 1)
}

//nolint:paralleltest // mutates the process-wide reset latch and the global logger
func TestGetStatsPeakResetTransientErrorDoesNotLatch(t *testing.T) {
calls := stubPeakReset(t, peakWriteError(syscall.EIO))
logs := observeWarnings(t)

cgroupPath, memoryPeakFile := fakeCgroupDir(t)
mgr := &managerImpl{}

for range 3 {
stats, err := mgr.getStatsForPath(t.Context(), cgroupPath, memoryPeakFile)
require.NoError(t, err)
assertParsedStats(t, stats)
}

assert.False(t, memoryPeakResetUnsupported.Load(), "a transient error must not latch")
assert.Equal(t, int64(3), calls.Load(), "a transient error must be retried on every sample")

warnings := logs.FilterLevelExact(zapcore.WarnLevel).All()
require.Len(t, warnings, 3, "transient failures keep warning, as before")
assert.Contains(t, warnings[0].Message, "interval peak semantics degraded")
}

//nolint:paralleltest // mutates the process-wide reset latch and the global logger
func TestGetStatsPeakResetSucceeds(t *testing.T) {
calls := stubPeakReset(t, nil)
logs := observeWarnings(t)

cgroupPath, memoryPeakFile := fakeCgroupDir(t)
mgr := &managerImpl{}

for range 3 {
stats, err := mgr.getStatsForPath(t.Context(), cgroupPath, memoryPeakFile)
require.NoError(t, err)
assertParsedStats(t, stats)
}

assert.False(t, memoryPeakResetUnsupported.Load())
assert.Equal(t, int64(3), calls.Load(), "a supported reset runs on every sample")
assert.Empty(t, logs.FilterLevelExact(zapcore.WarnLevel).All())
}

func TestPeakResetUnsupported(t *testing.T) {
t.Parallel()

for _, tc := range []struct {
name string
err error
want bool
}{
{
name: "nil",
},
{
name: "EINVAL",
err: syscall.EINVAL,
want: true,
},
{
name: "EINVAL wrapped by os.File.WriteString",
err: peakWriteError(syscall.EINVAL),
want: true,
},
{
name: "ENOTSUP",
err: syscall.ENOTSUP,
want: true,
},
{
name: "EOPNOTSUPP",
err: syscall.EOPNOTSUPP,
want: true,
},
{
name: "EIO",
err: syscall.EIO,
},
{
name: "EBADF",
err: syscall.EBADF,
},
{
name: "not a syscall error",
err: errors.New("boom"),
},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()

assert.Equal(t, tc.want, peakResetUnsupported(tc.err))
})
}
}
Loading