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
44 changes: 23 additions & 21 deletions lib/network/bridge_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -667,42 +667,39 @@ func (m *manager) removeRateLimit(tapName string) error {
func (m *manager) setupBridgeHTB(ctx context.Context, bridgeName string, capacityBps int64) error {
log := logger.FromContext(ctx)

if capacityBps <= 0 {
log.DebugContext(ctx, "skipping HTB setup - no capacity configured", "bridge", bridgeName)
return nil
}

// Check if HTB qdisc already exists
checkCmd := exec.Command("tc", "qdisc", "show", "dev", bridgeName)
checkCmd.SysProcAttr = &syscall.SysProcAttr{
AmbientCaps: []uintptr{unix.CAP_NET_ADMIN},
}
output, err := checkCmd.Output()
if err == nil && strings.Contains(string(output), "htb") {
log.InfoContext(ctx, "HTB qdisc ready", "bridge", bridgeName, "status", "existing")
if err != nil || !strings.Contains(string(output), "htb") {
// Add the scheduler without replacing existing per-VM classes and filters.
cmd := exec.Command("tc", "qdisc", "add", "dev", bridgeName, "root",
"handle", htbRootHandle, "htb")
cmd.SysProcAttr = &syscall.SysProcAttr{
AmbientCaps: []uintptr{unix.CAP_NET_ADMIN},
}
if output, err := cmd.CombinedOutput(); err != nil {
return fmt.Errorf("tc qdisc add htb: %w (output: %s)", err, string(output))
}
}

if capacityBps <= 0 {
log.InfoContext(ctx, "HTB qdisc ready", "bridge", bridgeName, "capacity", "unknown")
return nil
}

rateStr := formatTcRate(capacityBps)

// 1. Add root HTB qdisc (no default - all traffic must be classified)
cmd := exec.Command("tc", "qdisc", "add", "dev", bridgeName, "root",
"handle", htbRootHandle, "htb")
cmd.SysProcAttr = &syscall.SysProcAttr{
AmbientCaps: []uintptr{unix.CAP_NET_ADMIN},
}
if output, err := cmd.CombinedOutput(); err != nil {
return fmt.Errorf("tc qdisc add htb: %w (output: %s)", err, string(output))
}

// 2. Add root class for total capacity
cmd = exec.Command("tc", "class", "add", "dev", bridgeName, "parent", htbRootHandle,
// Ensure the shared parent exists even after an uncapped startup.
cmd := exec.Command("tc", "class", "replace", "dev", bridgeName, "parent", htbRootHandle,
"classid", htbRootClassID, "htb", "rate", rateStr)
cmd.SysProcAttr = &syscall.SysProcAttr{
AmbientCaps: []uintptr{unix.CAP_NET_ADMIN},
}
if output, err := cmd.CombinedOutput(); err != nil {
return fmt.Errorf("tc class add root: %w (output: %s)", err, string(output))
return fmt.Errorf("tc class replace root: %w (output: %s)", err, string(output))
}

log.InfoContext(ctx, "HTB qdisc ready", "bridge", bridgeName, "capacity", rateStr, "status", "configured")
Expand All @@ -723,6 +720,11 @@ func (m *manager) addVMClass(ctx context.Context, bridgeName, tapName string, ra
}
ceilStr := formatTcRate(ceilBps)

parent := htbRootClassID
if m.uncappedUploads {
parent = htbRootHandle
}

// Start with derived class ID, probe linearly on collision.
classIDVal := deriveClassIDVal(tapName)

Expand All @@ -733,7 +735,7 @@ func (m *manager) addVMClass(ctx context.Context, bridgeName, tapName string, ra
fullClassID := fmt.Sprintf("1:%s", classID)

// Try tc class add (NOT replace) so we detect collisions.
cmd := exec.Command("tc", "class", "add", "dev", bridgeName, "parent", htbRootClassID,
cmd := exec.Command("tc", "class", "add", "dev", bridgeName, "parent", parent,
"classid", fullClassID, "htb", "rate", rateStr, "ceil", ceilStr, "prio", "1")
cmd.SysProcAttr = &syscall.SysProcAttr{
AmbientCaps: []uintptr{unix.CAP_NET_ADMIN},
Expand Down
69 changes: 69 additions & 0 deletions lib/network/bridge_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,80 @@
package network

import (
"os"
"os/exec"
"testing"

"github.com/kernel/hypeman/cmd/api/config"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/vishvananda/netlink"
)

func TestHTBUnknownCapacity(t *testing.T) {
if os.Geteuid() != 0 {
t.Skip("requires root and an HTB-capable kernel")
}
if os.Getenv("HYPEMAN_TEST_HTB_NAMESPACE") != "1" {
cmd := exec.CommandContext(t.Context(), "unshare", "--net", os.Args[0], "-test.run=^TestHTBUnknownCapacity$", "-test.v")
cmd.Env = append(os.Environ(), "HYPEMAN_TEST_HTB_NAMESPACE=1")
output, err := cmd.CombinedOutput()
require.NoError(t, err, "%s", output)
t.Logf("%s", output)
return
}

ctx := t.Context()
bridge := &netlink.Bridge{LinkAttrs: netlink.LinkAttrs{Name: "htb-test"}}
require.NoError(t, netlink.LinkAdd(bridge))
tap := &netlink.Dummy{LinkAttrs: netlink.LinkAttrs{Name: "htb-tap"}}
require.NoError(t, netlink.LinkAdd(tap))
tapLink, err := netlink.LinkByName(tap.Name)
require.NoError(t, err)
cfg := &config.Config{Network: config.NetworkConfig{BridgeName: bridge.Name}}
show := func(t *testing.T, kind string) string {
output, err := exec.Command("tc", kind, "show", "dev", bridge.Name).CombinedOutput()
require.NoError(t, err, "%s", output)
return string(output)
}

for _, tt := range []struct {
name string
capacity int64
rootRate string
}{
{"unknown", 0, ""},
{"known-after-unknown", 125_000_000, "1Gbit"},
{"unknown-after-known", 0, ""},
{"updated-known", 250_000_000, "2Gbit"},
} {
t.Run(tt.name, func(t *testing.T) {
m := &manager{config: cfg}
require.NoError(t, m.SetupHTB(ctx, tt.capacity))
classID, err := m.addVMClass(ctx, bridge.Name, tap.Name, 37_500_000, 150_000_000)
require.NoError(t, err)
parent := "root"
if tt.capacity > 0 {
parent = "parent 1:1"
assert.Contains(t, show(t, "class"), "class htb 1:1 root rate "+tt.rootRate)
}
classes := show(t, "class")
assert.Contains(t, classes, "class htb 1:"+classID+" "+parent)
assert.Contains(t, classes, "rate 300Mbit ceil 1200Mbit")
if tt.name == "unknown" {
assert.Len(t, parseBridgeClasses(classes), 1)
}

filters := show(t, "filter")
m = &manager{config: cfg}
require.NoError(t, m.SetupHTB(ctx, tt.capacity))
assert.Equal(t, classes, show(t, "class"))
assert.Equal(t, filters, show(t, "filter"))
require.NoError(t, m.removeVMClass(ctx, bridge.Name, tapLink.Attrs().Index))
})
}
}

func TestParseBridgeFilters(t *testing.T) {
output := `filter parent 1: protocol all pref 1 basic chain 0
filter parent 1: protocol all pref 1 basic chain 0 handle 0x1 flowid 1:a3f2
Expand Down
9 changes: 8 additions & 1 deletion lib/network/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@ type manager struct {
defaultNetwork *Network
pendingAllocations map[string]pendingAllocation
tcMu sync.Mutex // Serializes shared bridge tc mutations.
uncappedUploads bool // Protected by tcMu.
metrics *Metrics
}

Expand Down Expand Up @@ -201,7 +202,13 @@ func (m *manager) getDefaultNetwork(ctx context.Context) (*Network, error) {
// SetupHTB initializes HTB qdisc on the bridge for upload fair sharing.
// capacityBps is the total network capacity in bytes per second.
func (m *manager) SetupHTB(ctx context.Context, capacityBps int64) error {
return m.setupBridgeHTB(ctx, m.config.Network.BridgeName, capacityBps)
m.tcMu.Lock()
defer m.tcMu.Unlock()
if err := m.setupBridgeHTB(ctx, m.config.Network.BridgeName, capacityBps); err != nil {
return err
}
m.uncappedUploads = capacityBps <= 0
return nil
}

// GetUploadBurstMultiplier returns the configured multiplier for upload burst ceiling.
Expand Down
3 changes: 3 additions & 0 deletions lib/resources/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,9 @@ Bidirectional rate limiting with separate download and upload controls:

**Capacity tracking:**
- Uses max(download, upload) per instance since they share physical link
- Failed capacity discovery logs a warning and disables host network admission enforcement; explicit per-instance rate limits remain unchanged.
- With unknown capacity, per-VM upload classes attach directly to the HTB scheduler, without a shared host-capacity parent or a guessed bandwidth limit. Existing classes are preserved during startup; recreated classes use the current capacity mode.
- `/resources` reports `source: "unknown"` in this case. Zero capacity, effective limit, and availability are placeholders, not enforced limits; allocation tracking remains active.

### Disk I/O

Expand Down
6 changes: 2 additions & 4 deletions lib/resources/network_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,14 +37,12 @@ func NewNetworkResource(ctx context.Context, cfg *config.Config, instLister Inst
// Auto-detect from uplink interface
uplink, err := getUplinkInterface(cfg.Network.UplinkInterface)
if err != nil {
// No uplink found - network limiting disabled
log.WarnContext(ctx, "no uplink interface found, network limiting disabled", "error", err)
log.WarnContext(ctx, "no uplink interface found, network admission limit disabled", "error", err)
capacity = 0
} else {
speed, err := getInterfaceSpeed(uplink)
if err != nil || speed <= 0 {
// Speed detection failed - network limiting disabled
log.WarnContext(ctx, "failed to detect interface speed, network limiting disabled", "interface", uplink, "error", err, "speed", speed)
log.WarnContext(ctx, "failed to detect interface speed, network admission limit disabled", "interface", uplink, "error", err, "speed", speed)
capacity = 0
} else {
// speed is in Mbps, convert to bytes/sec
Expand Down
49 changes: 49 additions & 0 deletions lib/resources/network_linux_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
//go:build linux

package resources

import (
"context"
"testing"

"github.com/kernel/hypeman/cmd/api/config"
"github.com/kernel/hypeman/lib/paths"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

func TestNetworkAdmissionCapacity(t *testing.T) {
for _, tt := range []struct {
configured string
source SourceType
wantErr bool
}{
{"", SourceUnknown, false},
{"1Gbps", SourceConfigured, true},
{"0Gbps", SourceConfigured, true},
} {
t.Run(string(tt.source)+tt.configured, func(t *testing.T) {
cfg := &config.Config{Capacity: config.CapacityConfig{Network: tt.configured},
Network: config.NetworkConfig{UplinkInterface: "missing-test-interface"},
Oversubscription: config.OversubscriptionConfig{Network: 1}}
lister := &mockInstanceLister{allocations: []InstanceAllocation{
{State: "Running", NetworkDownloadBps: 50_000_000},
}}
network, err := NewNetworkResource(context.Background(), cfg, lister)
require.NoError(t, err)
mgr := NewManager(cfg, paths.New(t.TempDir()))
mgr.SetInstanceLister(lister)
mgr.resources[ResourceNetwork] = network
err = mgr.ReserveAllocation(context.Background(), "test", 0, 0, 125_000_001, 0, 0, 0, false)
if tt.wantErr {
require.ErrorContains(t, err, "insufficient network bandwidth")
} else {
require.NoError(t, err)
}
status, err := mgr.GetStatus(context.Background(), ResourceNetwork)
require.NoError(t, err)
assert.Equal(t, tt.source, status.Source)
assert.Equal(t, int64(50_000_000), status.Allocated)
})
}
}
7 changes: 6 additions & 1 deletion lib/resources/resource.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ const (
type SourceType string

const (
SourceUnknown SourceType = "unknown" // Capacity unavailable; network admission limit disabled
SourceDetected SourceType = "detected" // Auto-detected from host hardware
SourceConfigured SourceType = "configured" // Explicitly configured by operator
)
Expand Down Expand Up @@ -349,6 +350,8 @@ func (m *Manager) GetStatus(ctx context.Context, rt ResourceType) (*ResourceStat
if rt == ResourceNetwork {
if m.cfg.Capacity.Network != "" {
status.Source = SourceConfigured
} else if status.Capacity == 0 {
status.Source = SourceUnknown
} else {
status.Source = SourceDetected
}
Expand Down Expand Up @@ -604,6 +607,8 @@ func (m *Manager) admissionStatusLocked(rt ResourceType, visibleAllocated int64,
status.Allocated = visibleAllocated + pending.NetworkBps
if m.cfg.Capacity.Network != "" {
status.Source = SourceConfigured
} else if status.Capacity == 0 {
status.Source = SourceUnknown
} else {
status.Source = SourceDetected
}
Expand Down Expand Up @@ -686,7 +691,7 @@ func (m *Manager) validateAllocationLocked(ctx context.Context, excludeID string
if err != nil {
return fmt.Errorf("check network capacity: %w", err)
}
if req.NetworkBps > status.Available {
if status.Source != SourceUnknown && req.NetworkBps > status.Available {
return fmt.Errorf("insufficient network bandwidth: requested %s/s, but only %s/s available (currently allocated: %s/s, effective limit: %s/s with %.1fx oversubscription)",
datasize.ByteSize(req.NetworkBps).HR(), datasize.ByteSize(status.Available).HR(), datasize.ByteSize(status.Allocated).HR(), datasize.ByteSize(status.EffectiveLimit).HR(), status.OversubRatio)
}
Expand Down
Loading