From 94796630861d8fa4de4108c6a81febe3fb52ec7d Mon Sep 17 00:00:00 2001 From: sjmiller609 <7516283+sjmiller609@users.noreply.github.com> Date: Tue, 6 Oct 2026 22:32:21 +0000 Subject: [PATCH 1/3] Use 10Gbps when network capacity detection fails --- lib/resources/network_linux.go | 12 +++--- lib/resources/network_linux_test.go | 63 +++++++++++++++++++++++++++++ 2 files changed, 69 insertions(+), 6 deletions(-) create mode 100644 lib/resources/network_linux_test.go diff --git a/lib/resources/network_linux.go b/lib/resources/network_linux.go index 0ec292a0a..3dd1ce982 100644 --- a/lib/resources/network_linux.go +++ b/lib/resources/network_linux.go @@ -14,6 +14,8 @@ import ( "github.com/vishvananda/netlink" ) +const fallbackNetworkCapacity int64 = 10_000_000_000 / 8 + // NetworkResource implements Resource for network bandwidth discovery and tracking. type NetworkResource struct { capacity int64 // bytes per second @@ -37,15 +39,13 @@ 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) - capacity = 0 + log.WarnContext(ctx, "no uplink interface found, falling back to 10Gbps", "error", err) + capacity = fallbackNetworkCapacity } 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) - capacity = 0 + log.WarnContext(ctx, "failed to detect interface speed, falling back to 10Gbps", "interface", uplink, "error", err, "speed", speed) + capacity = fallbackNetworkCapacity } else { // speed is in Mbps, convert to bytes/sec capacity = speed * 1000 * 1000 / 8 diff --git a/lib/resources/network_linux_test.go b/lib/resources/network_linux_test.go new file mode 100644 index 000000000..a70fbd92e --- /dev/null +++ b/lib/resources/network_linux_test.go @@ -0,0 +1,63 @@ +//go:build linux + +package resources + +import ( + "bytes" + "context" + "encoding/json" + "log/slog" + "testing" + + "github.com/kernel/hypeman/cmd/api/config" + "github.com/kernel/hypeman/lib/logger" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestNewNetworkResourceCapacity(t *testing.T) { + for _, tt := range []struct { + name string + iface string + configured string + want int64 + warn bool + wantErr bool + }{ + {name: "missing interface", iface: "hypeman-test", want: 1_250_000_000, warn: true}, + {name: "loopback without speed", iface: "lo", want: 1_250_000_000, warn: true}, + {name: "configured capacity", iface: "hypeman-test", configured: "2Gbps", want: 250_000_000}, + {name: "invalid configured capacity", iface: "hypeman-test", configured: "invalid", wantErr: true}, + } { + t.Run(tt.name, func(t *testing.T) { + var output bytes.Buffer + ctx := logger.AddToContext(context.Background(), slog.New(slog.NewJSONHandler(&output, nil))) + cfg := &config.Config{ + Capacity: config.CapacityConfig{Network: tt.configured}, + Network: config.NetworkConfig{UplinkInterface: tt.iface}, + } + network, err := NewNetworkResource(ctx, cfg, nil) + if tt.wantErr { + require.ErrorContains(t, err, "parse network limit") + assert.Nil(t, network) + assert.Empty(t, output.String()) + return + } + require.NoError(t, err) + assert.Equal(t, tt.want, network.Capacity()) + if tt.warn { + var record struct { + Level string + Msg string + Interface string + } + require.NoError(t, json.Unmarshal(output.Bytes(), &record)) + assert.Equal(t, "WARN", record.Level) + assert.Contains(t, record.Msg, "falling back to 10Gbps") + assert.Equal(t, tt.iface, record.Interface) + } else { + assert.Empty(t, output.String()) + } + }) + } +} From 27a9ed2f8959c7661df0a3233a5c0be5650df62f Mon Sep 17 00:00:00 2001 From: sjmiller609 <7516283+sjmiller609@users.noreply.github.com> Date: Tue, 6 Oct 2026 22:51:08 +0000 Subject: [PATCH 2/3] Disable network admission when capacity is unknown --- lib/resources/README.md | 2 + lib/resources/network_linux.go | 10 ++--- lib/resources/network_linux_test.go | 64 +++++++++++------------------ lib/resources/resource.go | 7 +++- 4 files changed, 37 insertions(+), 46 deletions(-) diff --git a/lib/resources/README.md b/lib/resources/README.md index 1ca94319e..a3e7bcd3f 100644 --- a/lib/resources/README.md +++ b/lib/resources/README.md @@ -66,6 +66,8 @@ 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. +- `/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 3dd1ce982..be7396f84 100644 --- a/lib/resources/network_linux.go +++ b/lib/resources/network_linux.go @@ -14,8 +14,6 @@ import ( "github.com/vishvananda/netlink" ) -const fallbackNetworkCapacity int64 = 10_000_000_000 / 8 - // NetworkResource implements Resource for network bandwidth discovery and tracking. type NetworkResource struct { capacity int64 // bytes per second @@ -39,13 +37,13 @@ func NewNetworkResource(ctx context.Context, cfg *config.Config, instLister Inst // Auto-detect from uplink interface uplink, err := getUplinkInterface(cfg.Network.UplinkInterface) if err != nil { - log.WarnContext(ctx, "no uplink interface found, falling back to 10Gbps", "error", err) - capacity = fallbackNetworkCapacity + 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 { - log.WarnContext(ctx, "failed to detect interface speed, falling back to 10Gbps", "interface", uplink, "error", err, "speed", speed) - capacity = fallbackNetworkCapacity + 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 capacity = speed * 1000 * 1000 / 8 diff --git a/lib/resources/network_linux_test.go b/lib/resources/network_linux_test.go index a70fbd92e..d7d32f8fd 100644 --- a/lib/resources/network_linux_test.go +++ b/lib/resources/network_linux_test.go @@ -3,61 +3,47 @@ package resources import ( - "bytes" "context" - "encoding/json" - "log/slog" "testing" "github.com/kernel/hypeman/cmd/api/config" - "github.com/kernel/hypeman/lib/logger" + "github.com/kernel/hypeman/lib/paths" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) -func TestNewNetworkResourceCapacity(t *testing.T) { +func TestNetworkAdmissionCapacity(t *testing.T) { for _, tt := range []struct { - name string - iface string configured string - want int64 - warn bool + source SourceType wantErr bool }{ - {name: "missing interface", iface: "hypeman-test", want: 1_250_000_000, warn: true}, - {name: "loopback without speed", iface: "lo", want: 1_250_000_000, warn: true}, - {name: "configured capacity", iface: "hypeman-test", configured: "2Gbps", want: 250_000_000}, - {name: "invalid configured capacity", iface: "hypeman-test", configured: "invalid", wantErr: true}, + {"", SourceUnknown, false}, + {"1Gbps", SourceConfigured, true}, + {"0Gbps", SourceConfigured, true}, } { - t.Run(tt.name, func(t *testing.T) { - var output bytes.Buffer - ctx := logger.AddToContext(context.Background(), slog.New(slog.NewJSONHandler(&output, nil))) - cfg := &config.Config{ - Capacity: config.CapacityConfig{Network: tt.configured}, - Network: config.NetworkConfig{UplinkInterface: tt.iface}, - } - network, err := NewNetworkResource(ctx, cfg, nil) - if tt.wantErr { - require.ErrorContains(t, err, "parse network limit") - assert.Nil(t, network) - assert.Empty(t, output.String()) - return - } + 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) - assert.Equal(t, tt.want, network.Capacity()) - if tt.warn { - var record struct { - Level string - Msg string - Interface string - } - require.NoError(t, json.Unmarshal(output.Bytes(), &record)) - assert.Equal(t, "WARN", record.Level) - assert.Contains(t, record.Msg, "falling back to 10Gbps") - assert.Equal(t, tt.iface, record.Interface) + 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 { - assert.Empty(t, output.String()) + 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) } From 66b3203486a15ce596f4522a96badf2ed325de5d Mon Sep 17 00:00:00 2001 From: sjmiller609 <7516283+sjmiller609@users.noreply.github.com> Date: Tue, 6 Oct 2026 23:21:48 +0000 Subject: [PATCH 3/3] Preserve upload shaping when network capacity is unknown --- lib/network/bridge_linux.go | 44 ++++++++++---------- lib/network/bridge_linux_test.go | 69 ++++++++++++++++++++++++++++++++ lib/network/manager.go | 9 ++++- lib/resources/README.md | 1 + 4 files changed, 101 insertions(+), 22 deletions(-) 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 a3e7bcd3f..8e1360cc8 100644 --- a/lib/resources/README.md +++ b/lib/resources/README.md @@ -67,6 +67,7 @@ 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