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
2 changes: 1 addition & 1 deletion packages/clickhouse/pkg/hoststats/hoststats.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,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"` // lifetime peak, bytes
}

// Delivery is the interface for delivering host stats to storage backend
Expand Down
39 changes: 17 additions & 22 deletions packages/orchestrator/internal/sandbox/cgroup/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ type Stats struct {
CPUSystemUsec uint64 // microseconds

MemoryUsageBytes uint64 // bytes
MemoryPeakBytes uint64 // bytes, reset after each GetStats() call
MemoryPeakBytes uint64 // bytes, lifetime peak
}

// CgroupHandle represents a created cgroup for a sandbox.
Expand All @@ -42,8 +42,8 @@ type Stats struct {
// Lifecycle: Create → GetFD → cmd.Start() → ReleaseCgroupFD → GetStats (repeatedly) → Remove
//
// The caller MUST call ReleaseCgroupFD() right after cmd.Start() (regardless of
// whether Start succeeded or failed). Remove() only closes the memory.peak FD
// and deletes the cgroup directory — it does not release the cgroup directory FD.
// whether Start succeeded or failed). Remove() closes the memory.peak FD and
// deletes the cgroup directory — it does not release the cgroup directory FD.
type CgroupHandle struct {
sandboxID string
path string
Expand All @@ -68,8 +68,8 @@ func (h *CgroupHandle) GetFD() int {
// Call this after cmd.Start() — the kernel has already placed the process in
// the cgroup atomically during clone, so the directory FD is no longer needed.
//
// The memory.peak FD is intentionally kept open because the per-FD reset
// mechanism requires the same FD for the lifetime of stats collection.
// The memory.peak FD is intentionally kept open for the lifetime of stats
// collection. It will be needed for per-FD reset when kernel 6.12+ is available.
// That FD is closed later by Remove().
//
// Safe to call multiple times.
Expand Down Expand Up @@ -205,12 +205,12 @@ func (m *managerImpl) Create(ctx context.Context, sandboxID string) (*CgroupHand
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
// TODO: Change to os.O_RDWR when per-FD peak reset is re-enabled.
memPeakPath := filepath.Join(cgroupPath, "memory.peak")
memoryPeakFile, peakErr := os.OpenFile(memPeakPath, os.O_RDWR, 0)
memoryPeakFile, peakErr := os.OpenFile(memPeakPath, os.O_RDONLY, 0)
if peakErr != nil {
// Not fatal — memory.peak may not exist on older kernels
logger.L().Debug(ctx, "failed to open memory.peak for reset (will track lifetime peak)",
logger.L().Debug(ctx, "failed to open memory.peak",
Comment thread
cursor[bot] marked this conversation as resolved.
logger.WithSandboxID(sandboxID),
zap.String("path", memPeakPath),
zap.Error(peakErr))
Expand All @@ -228,8 +228,7 @@ func (m *managerImpl) Create(ctx context.Context, sandboxID string) (*CgroupHand
logger.L().Debug(ctx, "created cgroup for sandbox",
logger.WithSandboxID(sandboxID),
zap.String("path", cgroupPath),
zap.Int("fd", handle.GetFD()),
zap.Bool("peak_reset_available", memoryPeakFile != nil))
zap.Int("fd", handle.GetFD()))

return handle, nil
}
Expand Down Expand Up @@ -271,7 +270,7 @@ func (m *managerImpl) getStatsForPath(ctx context.Context, cgroupPath string, me
stats.MemoryUsageBytes, _ = strconv.ParseUint(strings.TrimSpace(string(memData)), 10, 64)

if memoryPeakFile != nil {
peakBytes, err := m.readAndResetMemoryPeak(ctx, memoryPeakFile)
peakBytes, err := m.readMemoryPeak(memoryPeakFile)
if err != nil {
logger.L().Debug(ctx, "failed to read memory.peak", zap.Error(err))
} else {
Expand All @@ -282,13 +281,12 @@ func (m *managerImpl) getStatsForPath(ctx context.Context, cgroupPath string, me
return stats, nil
}

// readAndResetMemoryPeak reads the current peak memory value and resets it for the next interval.
// It uses the persistent FD kept open in CgroupHandle for per-FD reset tracking.
// The cgroups v2 kernel interface works as follows:
// - Read requires file position 0 (seq_file), so we seek before reading.
// - Write resets the per-FD peak to current memory usage. The kernel ignores
// both the written content and the file offset, so no seek before write is needed.
func (m *managerImpl) readAndResetMemoryPeak(ctx context.Context, memoryPeakFile *os.File) (uint64, error) {
// readMemoryPeak reads the current peak memory value from the persistent FD.
// Read requires file position 0 (seq_file), so we seek before reading.
//
// TODO: When per-FD peak reset is available-enabled, this function
// should also write to the FD to reset the peak for the next interval.
func (m *managerImpl) readMemoryPeak(memoryPeakFile *os.File) (uint64, error) {
if _, err := memoryPeakFile.Seek(0, io.SeekStart); err != nil {
return 0, fmt.Errorf("failed to seek memory.peak for read: %w", err)
}
Expand All @@ -301,10 +299,7 @@ func (m *managerImpl) readAndResetMemoryPeak(ctx context.Context, memoryPeakFile

peakBytes, _ := strconv.ParseUint(strings.TrimSpace(string(buf[:n])), 10, 64)

// Reset per-FD peak for next interval
if _, err := memoryPeakFile.WriteString("0"); err != nil {
logger.L().Debug(ctx, "failed to reset memory.peak", zap.Error(err))
}
// TODO:The write for Per-FD peak reset introduced in kernel 6.12 belongs here.

return peakBytes, nil
}
Expand Down
25 changes: 14 additions & 11 deletions packages/orchestrator/internal/sandbox/cgroup/manager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -325,7 +325,7 @@ burst_usec 0`
assert.Equal(t, uint64(0), stats.MemoryPeakBytes, "MemoryPeakBytes should be 0 without peak FD")
}

func TestCgroupHandlePeakReset(t *testing.T) {
func TestCgroupHandleLifetimePeak(t *testing.T) {
t.Parallel()

if os.Geteuid() != 0 {
Expand All @@ -339,14 +339,14 @@ func TestCgroupHandlePeakReset(t *testing.T) {
err = mgr.Initialize(ctx)
require.NoError(t, err)

testSandboxID := "test-peak-reset"
testSandboxID := "test-lifetime-peak"

handle, err := mgr.Create(ctx, testSandboxID)
require.NoError(t, err)
defer handle.Remove(ctx)
defer handle.ReleaseCgroupFD()

// Allocate memory gradually so we can sample the peak reset behavior
// Allocate memory gradually to observe lifetime peak behavior
cmd := exec.CommandContext(ctx, "bash", "-c",
"x=''; for i in {1..10}; do x=$x$(head -c 5M /dev/zero | tr '\\0' 'x'); sleep 0.5; done; sleep 5")
cmd.SysProcAttr = &syscall.SysProcAttr{
Expand All @@ -365,30 +365,33 @@ func TestCgroupHandlePeakReset(t *testing.T) {
require.NoError(t, err)
peak1 := stats1.MemoryPeakBytes
require.Positive(t, peak1, "First peak should be non-zero")
assert.GreaterOrEqual(t, peak1, stats1.MemoryUsageBytes,
"Peak should be >= current memory")
t.Logf("First sample - peak: %d bytes, current: %d bytes", peak1, stats1.MemoryUsageBytes)

// Peak should represent interval peak (since last GetStats), not lifetime
time.Sleep(2 * time.Second)
stats2, err := handle.GetStats(ctx)
require.NoError(t, err)
peak2 := stats2.MemoryPeakBytes
require.Positive(t, peak2, "Second peak should be non-zero")
t.Logf("Second sample - peak: %d bytes, current: %d bytes", peak2, stats2.MemoryUsageBytes)

assert.GreaterOrEqual(t, peak2, peak1,
"Lifetime peak should be monotonically non-decreasing")
assert.GreaterOrEqual(t, peak2, stats2.MemoryUsageBytes,
"Peak memory should be >= current memory within the interval")
"Peak should be >= current memory")
t.Logf("Second sample - peak: %d bytes, current: %d bytes", peak2, stats2.MemoryUsageBytes)

time.Sleep(2 * time.Second)
stats3, err := handle.GetStats(ctx)
require.NoError(t, err)
peak3 := stats3.MemoryPeakBytes
require.Positive(t, peak3, "Third peak should be non-zero")
t.Logf("Third sample - peak: %d bytes, current: %d bytes", peak3, stats3.MemoryUsageBytes)

assert.GreaterOrEqual(t, peak3, peak2,
"Lifetime peak should be monotonically non-decreasing")
assert.GreaterOrEqual(t, peak3, stats3.MemoryUsageBytes,
"Peak memory should be >= current memory within the interval")
"Peak should be >= current memory")
t.Logf("Third sample - peak: %d bytes, current: %d bytes", peak3, stats3.MemoryUsageBytes)

t.Logf("Reset test complete - peaks tracked per interval: %d, %d, %d bytes",
t.Logf("Lifetime peak test complete - peaks: %d, %d, %d bytes",
peak1, peak2, peak3)

cmd.Process.Kill()
Expand Down
Loading