diff --git a/lib/network/bridge_linux.go b/lib/network/bridge_linux.go index 2e5bfc986..7d43ac8da 100644 --- a/lib/network/bridge_linux.go +++ b/lib/network/bridge_linux.go @@ -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") @@ -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) @@ -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}, diff --git a/lib/network/bridge_linux_test.go b/lib/network/bridge_linux_test.go index 3c86fd1b3..2cf427b85 100644 --- a/lib/network/bridge_linux_test.go +++ b/lib/network/bridge_linux_test.go @@ -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 diff --git a/lib/network/manager.go b/lib/network/manager.go index bdfacad9e..8a9560f1a 100644 --- a/lib/network/manager.go +++ b/lib/network/manager.go @@ -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 } @@ -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. diff --git a/lib/resources/README.md b/lib/resources/README.md index 1ca94319e..8e1360cc8 100644 --- a/lib/resources/README.md +++ b/lib/resources/README.md @@ -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 diff --git a/lib/resources/network_linux.go b/lib/resources/network_linux.go index 0ec292a0a..be7396f84 100644 --- a/lib/resources/network_linux.go +++ b/lib/resources/network_linux.go @@ -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 diff --git a/lib/resources/network_linux_test.go b/lib/resources/network_linux_test.go new file mode 100644 index 000000000..d7d32f8fd --- /dev/null +++ b/lib/resources/network_linux_test.go @@ -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) + }) + } +} diff --git a/lib/resources/resource.go b/lib/resources/resource.go index f8f624e93..cf90a2e60 100644 --- a/lib/resources/resource.go +++ b/lib/resources/resource.go @@ -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 ) @@ -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 } @@ -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 } @@ -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) }