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
180 changes: 168 additions & 12 deletions cgroup/cgroup.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,8 @@ const (
maxName = "cpu.max"
burstName = "cpu.max.burst"
statName = "cpu.stat"
cpusetName = "cpuset.cpus"
cpusetEffName = "cpuset.cpus.effective"

removeWait = time.Second
removePollInterval = 10 * time.Millisecond
Expand All @@ -56,16 +58,18 @@ type Knobs struct {
QuotaUs int64
PeriodUs int64
BurstUs int64
CPUSet string
}

// ResolveKnobs applies the Guaranteed-at-N defaults: weight = vCPU count, quota = vCPU count x period, burst = 0.
// ResolveKnobs applies the Guaranteed-at-N defaults: weight = vCPU count, quota = vCPU count x period, burst = 0, no placement.
func ResolveKnobs(cfg *types.Config) Knobs {
period := cmp.Or(cfg.CPUPeriodUs, int64(DefaultPeriodUs))
return Knobs{
Weight: cmp.Or(cfg.CPUWeight, cfg.CPU),
QuotaUs: cmp.Or(cfg.CPUQuotaUs, int64(cfg.CPU)*period),
PeriodUs: period,
BurstUs: cfg.CPUBurstUs,
CPUSet: cfg.CPUSetCPUs,
}
}

Expand All @@ -83,27 +87,68 @@ func (k Knobs) Validate() error {
if k.BurstUs < 0 || k.BurstUs > k.QuotaUs {
return fmt.Errorf("--cpu-burst-us must be 0..quota (%d), got %d", k.QuotaUs, k.BurstUs)
}
if _, err := ParseCPUList(k.CPUSet); err != nil {
return fmt.Errorf("--cpuset-cpus: %w", err)
}
return nil
}

// EffectiveCPUs resolves the cpu set a VM may run on — explicit placement wins over the machine fence; nil means all cores.
func EffectiveCPUs(placement, fence string) []int {
cpus, err := ParseCPUList(cmp.Or(placement, fence))
if err != nil {
return nil
}
return cpus
}

// ParseCPUList parses a kernel cpu-list ("0-14", "0,2-4"); empty input is nil (no placement).
func ParseCPUList(s string) ([]int, error) {
if s == "" {
return nil, nil
}
var cpus []int
for part := range strings.SplitSeq(s, ",") {
lo, hi, ok := strings.Cut(strings.TrimSpace(part), "-")
start, err := strconv.Atoi(lo)
if err != nil || start < 0 {
return nil, fmt.Errorf("invalid cpu list %q", s)
}
end := start
if ok {
if end, err = strconv.Atoi(hi); err != nil || end < start {
return nil, fmt.Errorf("invalid cpu list %q", s)
}
}
for c := start; c <= end; c++ {
cpus = append(cpus, c)
}
}
slices.Sort(cpus)
return slices.Compact(cpus), nil
}

// ScopeDir returns vmID's scope directory under parentDir.
func ScopeDir(parentDir, vmID string) string {
return filepath.Join(parentDir, scopePrefix+vmID+scopeSuffix)
}

// Prepare creates or reconfigures vmID's scope and returns its opened directory for CLONE_INTO_CGROUP; idempotent, so a relaunch reuses a scope its dying predecessor still occupies.
func Prepare(parentDir, vmID string, k Knobs) (*os.File, error) {
func Prepare(parentDir, fence, vmID string, k Knobs) (*os.File, error) {
if vmID == "" {
return nil, errors.New("cgroup scope: empty vm id")
}
if err := ensureParent(parentDir); err != nil {
if err := ensureParent(parentDir, fence, k.CPUSet != ""); err != nil {
return nil, err
}
dir := ScopeDir(parentDir, vmID)
mkErr := os.Mkdir(dir, 0o750)
if mkErr != nil && !errors.Is(mkErr, fs.ErrExist) {
return nil, fmt.Errorf("create scope: %w", mkErr)
}
if err := placeScope(parentDir, dir, k.CPUSet); err != nil {
return nil, err
}
if err := writeControl(dir, weightName, strconv.Itoa(k.Weight)); err != nil {
return nil, err
}
Expand Down Expand Up @@ -199,35 +244,146 @@ func parseStat(data string) map[string]int64 {
return stat
}

// ensureParent enables cpu at every ancestor, not just the leaf — cgroup v2 subtree delegation is hierarchical.
func ensureParent(parentDir string) error {
// ensureParent enables the needed controllers at every ancestor, not just the leaf — cgroup v2 subtree delegation is hierarchical.
func ensureParent(parentDir, fence string, placed bool) error {
rel, err := filepath.Rel(Root, parentDir)
if err != nil || rel == "." || strings.HasPrefix(rel, "..") {
return fmt.Errorf("cgroup parent %q must be under %s", parentDir, Root)
}
if err := utils.EnsureDirs(parentDir); err != nil {
return err
}
if err := enableCPU(Root); err != nil {
ctrls := []string{"cpu"}
if fence != "" || placed {
ctrls = append(ctrls, "cpuset")
}
if err := forEachLevel(rel, func(dir string) error { return enableControllers(dir, ctrls) }); err != nil {
return err
}
return reconcileFence(parentDir, fence)
}

// reconcileFence converges the parent's cpuset.cpus to the configured fence; a cleared config resets a stale fence once, and the cpuset controller is never disabled.
func reconcileFence(parentDir, fence string) error {
current, err := readControl(parentDir, cpusetName)
if fence == "" {
if err == nil && current != "" {
return writeControl(parentDir, cpusetName, "\n")
}
return nil
}
// The kernel echoes cpu lists canonicalized ("0-3,4-7" reads back "0-7"): compare parsed sets, or a non-canonical config string re-runs the scan and the serialized cpuset write on every launch.
if cur, curErr := ParseCPUList(current); curErr == nil {
if want, wantErr := ParseCPUList(fence); wantErr == nil && slices.Equal(cur, want) {
return nil
}
}
if err := checkSubset(fence, filepath.Dir(parentDir), "cgroup_cpus fence"); err != nil {
return err
}
if err := checkScopePlacements(parentDir, fence); err != nil {
return err
}
return writeControl(parentDir, cpusetName, fence)
}

// placeScope applies a per-VM placement; the kernel silently degrades ungrantable requests to the parent set, so the subset check is cocoon's.
func placeScope(parentDir, dir, cpuset string) error {
if cpuset == "" {
return nil
}
if current, err := readControl(dir, cpusetName); err == nil {
if cur, curErr := ParseCPUList(current); curErr == nil {
if want, wantErr := ParseCPUList(cpuset); wantErr == nil && slices.Equal(cur, want) {
return nil
}
}
}
if err := checkSubset(cpuset, parentDir, "--cpuset-cpus"); err != nil {
return err
}
return writeControl(dir, cpusetName, cpuset)
}

func checkSubset(cpuset, dir, what string) error {
want, err := ParseCPUList(cpuset)
if err != nil {
return fmt.Errorf("%s: %w", what, err)
}
effRaw, err := readControl(dir, cpusetEffName)
if err != nil {
return fmt.Errorf("read effective cpuset: %w", err)
}
eff, err := ParseCPUList(effRaw)
if err != nil {
return fmt.Errorf("parse %s: %w", filepath.Join(dir, cpusetEffName), err)
}
for _, c := range want {
if !slices.Contains(eff, c) {
return fmt.Errorf("%s %q: cpu %d not in %s effective set %q", what, cpuset, c, dir, effRaw)
}
}
return nil
}

// checkScopePlacements refuses a fence shrink that would invalidate a live VM's explicit placement.
func checkScopePlacements(parentDir, fence string) error {
ids, err := ListScopeVMIDs(parentDir)
if err != nil {
return err
}
fenceCPUs, _ := ParseCPUList(fence)
for _, id := range ids {
placement, err := readControl(ScopeDir(parentDir, id), cpusetName)
if err != nil || placement == "" {
continue
}
cpus, err := ParseCPUList(placement)
if err != nil {
continue
}
for _, c := range cpus {
if !slices.Contains(fenceCPUs, c) {
return fmt.Errorf("fence %q excludes cpu %d used by vm %s placement %q; stop that VM or widen the fence", fence, c, id, placement)
}
}
}
return nil
}

func forEachLevel(rel string, fn func(dir string) error) error {
if err := fn(Root); err != nil {
return err
}
dir := Root
for part := range strings.SplitSeq(rel, string(filepath.Separator)) {
dir = filepath.Join(dir, part)
if err := enableCPU(dir); err != nil {
if err := fn(dir); err != nil {
return err
}
}
return nil
}

// enableCPU reads before writing: subtree_control writes take the kernel's hierarchy-wide cgroup_mutex, so steady-state launches must not contend on a no-op write.
func enableCPU(dir string) error {
path := filepath.Join(dir, subtreeControlName)
if data, err := os.ReadFile(path); err == nil && slices.Contains(strings.Fields(string(data)), "cpu") { //nolint:gosec // fixed name under the config-derived parent
// enableControllers reads before writing: subtree_control writes take the kernel's hierarchy-wide cgroup_mutex, so steady-state launches must not contend on a no-op write. Missing controllers go in one combined write.
func enableControllers(dir string, ctrls []string) error {
have, _ := readControl(dir, subtreeControlName)
enabled := strings.Fields(have)
var missing []string
for _, ctrl := range ctrls {
if !slices.Contains(enabled, ctrl) {
missing = append(missing, "+"+ctrl)
}
}
if len(missing) == 0 {
return nil
}
return writeControl(dir, subtreeControlName, "+cpu")
return writeControl(dir, subtreeControlName, strings.Join(missing, " "))
}

func readControl(dir, name string) (string, error) {
data, err := os.ReadFile(filepath.Join(dir, name)) //nolint:gosec // fixed name under the config-derived parent
return strings.TrimSpace(string(data)), err
}

func writeControl(dir, name, value string) error {
Expand Down
126 changes: 126 additions & 0 deletions cgroup/cgroup_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"os"
"path/filepath"
"slices"
"strings"
"testing"

"github.com/cocoonstack/cocoon/types"
Expand Down Expand Up @@ -109,3 +110,128 @@ func TestParseStat(t *testing.T) {
t.Errorf("got %v", stat)
}
}

func TestParseCPUList(t *testing.T) {
tests := []struct {
in string
want []int
wantErr bool
}{
{in: "", want: nil},
{in: "3", want: []int{3}},
{in: "0-3", want: []int{0, 1, 2, 3}},
{in: "0,2-4,7", want: []int{0, 2, 3, 4, 7}},
{in: "4-2", wantErr: true},
{in: "a-b", wantErr: true},
{in: "-1", wantErr: true},
{in: "1,,2", wantErr: true},
}
for _, tt := range tests {
t.Run(tt.in, func(t *testing.T) {
got, err := ParseCPUList(tt.in)
if (err != nil) != tt.wantErr {
t.Fatalf("err = %v, wantErr %v", err, tt.wantErr)
}
if !tt.wantErr && !slices.Equal(got, tt.want) {
t.Errorf("got %v, want %v", got, tt.want)
}
})
}
}

func TestKnobsValidateRejectsBadCPUSet(t *testing.T) {
k := Knobs{Weight: 1, QuotaUs: 100000, PeriodUs: 100000, CPUSet: "9-1"}
if err := k.Validate(); err == nil {
t.Error("want error for invalid cpu list")
}
}

func TestCheckSubset(t *testing.T) {
dir := t.TempDir()
if err := os.WriteFile(filepath.Join(dir, "cpuset.cpus.effective"), []byte("0-14\n"), 0o600); err != nil {
t.Fatalf("setup: %v", err)
}
if err := checkSubset("2-4", dir, "fence"); err != nil {
t.Errorf("subset rejected: %v", err)
}
if err := checkSubset("14-15", dir, "fence"); err == nil {
t.Error("want error: cpu 15 outside effective 0-14")
}
}

func TestCheckScopePlacements(t *testing.T) {
parent := t.TempDir()
mk := func(id, placement string) {
dir := ScopeDir(parent, id)
if err := os.Mkdir(dir, 0o755); err != nil {
t.Fatalf("setup: %v", err)
}
if err := os.WriteFile(filepath.Join(dir, "cpuset.cpus"), []byte(placement+"\n"), 0o600); err != nil {
t.Fatalf("setup: %v", err)
}
}
mk("A", "")
mk("B", "2-3")

if err := checkScopePlacements(parent, "0-7"); err != nil {
t.Errorf("fence covering placements rejected: %v", err)
}
if err := checkScopePlacements(parent, "0-2"); err == nil {
t.Error("want error: fence 0-2 excludes cpu 3 used by B")
}
}

func TestReconcileFenceClearsStaleValue(t *testing.T) {
parent := t.TempDir()
path := filepath.Join(parent, "cpuset.cpus")
if err := os.WriteFile(path, []byte("0-14\n"), 0o600); err != nil {
t.Fatalf("setup: %v", err)
}
if err := reconcileFence(parent, ""); err != nil {
t.Fatalf("reconcile: %v", err)
}
data, err := os.ReadFile(path)
if err != nil || strings.TrimSpace(string(data)) != "" {
t.Errorf("stale fence not cleared: %q err=%v", data, err)
}
if err := reconcileFence(parent, ""); err != nil {
t.Errorf("steady-state empty reconcile: %v", err)
}
}

func TestReconcileFenceCanonicalEquality(t *testing.T) {
parent := t.TempDir()
if err := os.WriteFile(filepath.Join(parent, "cpuset.cpus"), []byte("0-7\n"), 0o600); err != nil {
t.Fatalf("setup: %v", err)
}
// No cpuset.cpus.effective fixture exists: reaching checkSubset would fail, so success proves the parsed-set gate short-circuited.
if err := reconcileFence(parent, "0-3,4-7"); err != nil {
t.Errorf("canonically-equal fence rewrote: %v", err)
}
}

func TestPlaceScopeReadGate(t *testing.T) {
parent := t.TempDir()
dir := ScopeDir(parent, "X")
if err := os.Mkdir(dir, 0o755); err != nil {
t.Fatalf("setup: %v", err)
}
if err := os.WriteFile(filepath.Join(dir, "cpuset.cpus"), []byte("2-3\n"), 0o600); err != nil {
t.Fatalf("setup: %v", err)
}
if err := placeScope(parent, dir, "2,3"); err != nil {
t.Errorf("equal placement rewrote: %v", err)
}
}

func TestEffectiveCPUs(t *testing.T) {
if got := EffectiveCPUs("2-3", "0-14"); !slices.Equal(got, []int{2, 3}) {
t.Errorf("placement wins: got %v", got)
}
if got := EffectiveCPUs("", "0-1"); !slices.Equal(got, []int{0, 1}) {
t.Errorf("fence fallback: got %v", got)
}
if EffectiveCPUs("", "") != nil {
t.Error("no constraint: want nil")
}
}
Loading