From 39d809c20603c9d8e3e8d68f073471e4bdda0d6e Mon Sep 17 00:00:00 2001 From: Abhishek Kumar Singh Date: Wed, 6 May 2026 15:58:41 +0530 Subject: [PATCH 1/5] Add configurable affinity_mode for egress pod selection MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The current StartEgressAffinity always scores idle pods at 0.5 and busy pods at 1.0. Combined with MaximumAffinity=1 in the psrpc client, this means the first pod to accept any job wins all subsequent jobs until its CPU budget is exhausted — a staircase pattern rather than even spread. The existing code comment acknowledges this is intentional for mixed fleets ("avoids having many instances with one track request each, taking availability from room composite"). However for a TrackEgress-only fleet the packing strategy provides no benefit and causes sequential scale-out delays. This commit adds a configurable affinity_mode field to ServiceConfig: pack (default) — current behaviour, unchanged spread — CPU-proportional scoring; idle pods score 1.0 and win immediately via MaximumAffinity short-circuit, busy pods score proportionally so the least-loaded pod wins after ShortCircuitTimeout. Best for single-type (TrackEgress-only) fleets. type_aware — RoomComposite/Web requests prefer idle pods (1.0 idle / 0.5 busy); Track/Participant requests spread by CPU load. Best of both worlds for mixed fleets. Default is "pack" so all existing deployments are unaffected. Co-Authored-By: Claude Sonnet 4.6 (cherry picked from commit e3ab4765602737fd33135ca9540b86db692bc8a2) --- pkg/config/service.go | 6 ++++ pkg/server/server_rpc.go | 39 +++++++++++++++++----- pkg/server/server_rpc_test.go | 62 +++++++++++++++++++++++++++++++++++ pkg/stats/monitor.go | 12 +++++++ 4 files changed, 111 insertions(+), 8 deletions(-) create mode 100644 pkg/server/server_rpc_test.go diff --git a/pkg/config/service.go b/pkg/config/service.go index 174c4adaa..a91ef7a9d 100644 --- a/pkg/config/service.go +++ b/pkg/config/service.go @@ -68,6 +68,12 @@ type ServiceConfig struct { PrometheusPort int `yaml:"prometheus_port"` // prometheus handler port DebugHandlerPort int `yaml:"debug_handler_port"` // egress debug handler port + // AffinityMode controls job distribution across pods. + // "pack" — busy pods win (default; preserves upstream behaviour) + // "spread" — CPU-proportional; use for TrackEgress-only fleets + // "type_aware" — RoomComposite/Web prefer idle pods; Track/Participant spread by load + AffinityMode string `yaml:"affinity_mode"` + *CPUCostConfig `yaml:"cpu_cost"` // CPU costs for the different egress types } diff --git a/pkg/server/server_rpc.go b/pkg/server/server_rpc.go index 3910dde2f..af4b5ef0e 100644 --- a/pkg/server/server_rpc.go +++ b/pkg/server/server_rpc.go @@ -207,18 +207,41 @@ func (s *Server) processEnded(req *rpc.StartEgressRequest, info *livekit.EgressI func (s *Server) StartEgressAffinity(_ context.Context, req *rpc.StartEgressRequest) float32 { if s.IsDisabled() || !s.monitor.CanAcceptRequest(req) { - // cannot accept return -1 } - if s.activeRequests.Load() == 0 { - // group multiple track and track composite requests. - // if this instance is idle and another is already handling some, the request will go to that server. - // this avoids having many instances with one track request each, taking availability from room composite. - return 0.5 + switch s.conf.AffinityMode { + + case "spread": + // Idle pods return 1.0 → MaximumAffinity short-circuit fires immediately. + // Busy pods return proportional score → least loaded wins after ShortCircuitTimeout. + return s.monitor.AvailableCPUFraction() + + case "type_aware": + if isHeavyEgressRequest(req) { + // RoomComposite/Web need ~4 CPUs. Strongly prefer idle pods. + if s.activeRequests.Load() == 0 { + return 1.0 + } + return 0.5 + } + return s.monitor.AvailableCPUFraction() + + default: // "pack" or empty — upstream behaviour unchanged + if s.activeRequests.Load() == 0 { + return 0.5 + } + return 1 + } +} + +func isHeavyEgressRequest(req *rpc.StartEgressRequest) bool { + switch req.Request.(type) { + case *rpc.StartEgressRequest_RoomComposite, + *rpc.StartEgressRequest_Web: + return true } - // already handling a request and has available cpu - return 1 + return false } func (s *Server) ListActiveEgress(ctx context.Context, _ *rpc.ListActiveEgressRequest) (*rpc.ListActiveEgressResponse, error) { diff --git a/pkg/server/server_rpc_test.go b/pkg/server/server_rpc_test.go new file mode 100644 index 000000000..9040b9962 --- /dev/null +++ b/pkg/server/server_rpc_test.go @@ -0,0 +1,62 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package server + +import ( + "testing" + + "github.com/stretchr/testify/require" + + "github.com/livekit/protocol/livekit" + "github.com/livekit/protocol/rpc" +) + +func TestIsHeavyEgressRequest(t *testing.T) { + for _, tc := range []struct { + name string + req *rpc.StartEgressRequest + expected bool + }{ + { + name: "RoomComposite is heavy", + req: &rpc.StartEgressRequest{Request: &rpc.StartEgressRequest_RoomComposite{RoomComposite: &livekit.RoomCompositeEgressRequest{}}}, + expected: true, + }, + { + name: "Web is heavy", + req: &rpc.StartEgressRequest{Request: &rpc.StartEgressRequest_Web{Web: &livekit.WebEgressRequest{}}}, + expected: true, + }, + { + name: "Track is not heavy", + req: &rpc.StartEgressRequest{Request: &rpc.StartEgressRequest_Track{Track: &livekit.TrackEgressRequest{}}}, + expected: false, + }, + { + name: "TrackComposite is not heavy", + req: &rpc.StartEgressRequest{Request: &rpc.StartEgressRequest_TrackComposite{TrackComposite: &livekit.TrackCompositeEgressRequest{}}}, + expected: false, + }, + { + name: "Participant is not heavy", + req: &rpc.StartEgressRequest{Request: &rpc.StartEgressRequest_Participant{Participant: &livekit.ParticipantEgressRequest{}}}, + expected: false, + }, + } { + t.Run(tc.name, func(t *testing.T) { + require.Equal(t, tc.expected, isHeavyEgressRequest(tc.req)) + }) + } +} diff --git a/pkg/stats/monitor.go b/pkg/stats/monitor.go index c8d4af432..9f8e9e25d 100644 --- a/pkg/stats/monitor.go +++ b/pkg/stats/monitor.go @@ -604,6 +604,18 @@ func (m *Monitor) GetAvailableCPU() float64 { return available } +// AvailableCPUFraction returns the fraction of CPU budget remaining (0.0–1.0). +// Returns 1.0 when the pod is idle. Used by StartEgressAffinity for spread/type_aware modes. +func (m *Monitor) AvailableCPUFraction() float32 { + m.mu.Lock() + defer m.mu.Unlock() + total, available, _, _ := m.getCPUUsageLocked() + if total == 0 { + return 1.0 + } + return float32(available / total) +} + func (m *Monitor) getCPUUsageLocked() (total, available, pending, used float64) { total = m.cpuStats.NumCPU() if m.requests.Load() == 0 { From 66c9e0748e4d32d63c322896deb759658f18175c Mon Sep 17 00:00:00 2001 From: Abhishek Kumar Singh Date: Thu, 7 May 2026 15:46:25 +0530 Subject: [PATCH 2/5] fix(spread): pending-counter + jitter to fix burst skew in spread/type_aware modes Addresses both root causes of the 24/51-job-on-one-pod skew observed in the 2026-05-07 load test: Cause A (strict > tie-break): psrpc's ShortCircuitTimeout means the first replier wins when all idle pods return the same score. Fixed by subtracting rand.Float32()*0.001 jitter so idle pods produce distinct scores, making the strict-> comparison effectively random among equally-idle peers. Cause B (m.requests.Inc lag): StartEgressAffinity is called before StartEgress, so the winning pod's m.requests counter stays 0 across an entire 200ms burst window and all callers see score 1.0. Fixed by a pendingClaims atomic.Int32 that increments at affinity time and decrements at StartEgress accept (consumePendingClaim). A 2s self-decay timer guards against claims that are never fulfilled. A CAS loop in consumePendingClaim ensures exactly one decrement fires per increment even when StartEgress and the timer race. New monitor helper AvailableCPUFractionWithPending deducts pendingSlots*TrackCpuCost from the available budget so the score decreases with each in-flight claim. Image: asia-south1-docker.pkg.dev/avian-pulsar-430509-f6/r41-livekit/egress:v1.12.0-r41.2 Co-Authored-By: Claude Sonnet 4.6 (cherry picked from commit ded32ad00dd1072fdfd2ce2c9ecea13f6d2bfde1) --- pkg/server/server.go | 1 + pkg/server/server_rpc.go | 42 ++++++++++++++++++++++++++---- pkg/server/server_rpc_test.go | 48 +++++++++++++++++++++++++++++++++++ pkg/stats/monitor.go | 20 +++++++++++++++ 4 files changed, 106 insertions(+), 5 deletions(-) diff --git a/pkg/server/server.go b/pkg/server/server.go index a6cb24b95..8bf95fd1d 100644 --- a/pkg/server/server.go +++ b/pkg/server/server.go @@ -57,6 +57,7 @@ type Server struct { ioClient info.SessionReporter activeRequests atomic.Int32 + pendingClaims atomic.Int32 // claimed-but-not-yet-accepted requests (spread/type_aware modes) terminating core.Fuse shutdown core.Fuse } diff --git a/pkg/server/server_rpc.go b/pkg/server/server_rpc.go index af4b5ef0e..5864a1ec0 100644 --- a/pkg/server/server_rpc.go +++ b/pkg/server/server_rpc.go @@ -16,6 +16,7 @@ package server import ( "context" + "math/rand" "net/http" "os" "os/exec" @@ -44,8 +45,25 @@ var ( tracer = otel.Tracer("github.com/livekit/egress/pkg/server") ) +// consumePendingClaim decrements pendingClaims by 1, floored at 0. +// Called both from StartEgress (claim accepted) and the 2s self-decay timer +// set in StartEgressAffinity. The CompareAndSwap loop ensures exactly one +// decrement fires per increment even when both callers race. +func (s *Server) consumePendingClaim() { + for { + n := s.pendingClaims.Load() + if n <= 0 { + return + } + if s.pendingClaims.CompareAndSwap(n, n-1) { + return + } + } +} + func (s *Server) StartEgress(ctx context.Context, req *rpc.StartEgressRequest) (*livekit.EgressInfo, error) { s.activeRequests.Inc() + s.consumePendingClaim() // hand slot from pending to m.requests ctx, span := tracer.Start(ctx, "Service.StartEgress") defer span.End() @@ -213,19 +231,33 @@ func (s *Server) StartEgressAffinity(_ context.Context, req *rpc.StartEgressRequ switch s.conf.AffinityMode { case "spread": - // Idle pods return 1.0 → MaximumAffinity short-circuit fires immediately. - // Busy pods return proportional score → least loaded wins after ShortCircuitTimeout. - return s.monitor.AvailableCPUFraction() + // Increment pendingClaims so subsequent affinity calls on this pod see a + // lower score before m.requests.Inc fires in StartEgress. The 2s self-decay + // guards against claims that are never accepted (lost RPCs, rejected requests). + // consumePendingClaim uses a CAS loop so exactly one of StartEgress or the + // timer decrements the counter — no double-decrement. + s.pendingClaims.Inc() + time.AfterFunc(2*time.Second, s.consumePendingClaim) + pending := s.pendingClaims.Load() + // Subtract rand jitter ≤ 0.001 to break equal-score ties in psrpc's strict-> + // comparison (first-replier-wins). Idle pods land in [0.999, 1.0]; busy pods + // score ≤ 0.96 (4 CPU, TrackCpuCost=0.15) and never beat idle pods. + return s.monitor.AvailableCPUFractionWithPending(pending) - rand.Float32()*0.001 case "type_aware": if isHeavyEgressRequest(req) { // RoomComposite/Web need ~4 CPUs. Strongly prefer idle pods. + // Jitter breaks first-replier ties among equally-idle pods. if s.activeRequests.Load() == 0 { - return 1.0 + return 1.0 - rand.Float32()*0.001 } return 0.5 } - return s.monitor.AvailableCPUFraction() + // Light requests use the same pending-counter + jitter logic as spread. + s.pendingClaims.Inc() + time.AfterFunc(2*time.Second, s.consumePendingClaim) + pending := s.pendingClaims.Load() + return s.monitor.AvailableCPUFractionWithPending(pending) - rand.Float32()*0.001 default: // "pack" or empty — upstream behaviour unchanged if s.activeRequests.Load() == 0 { diff --git a/pkg/server/server_rpc_test.go b/pkg/server/server_rpc_test.go index 9040b9962..9a90d44eb 100644 --- a/pkg/server/server_rpc_test.go +++ b/pkg/server/server_rpc_test.go @@ -15,6 +15,7 @@ package server import ( + "sync" "testing" "github.com/stretchr/testify/require" @@ -23,6 +24,53 @@ import ( "github.com/livekit/protocol/rpc" ) +// TestConsumePendingClaim_FloorAtZero verifies the CAS loop never decrements below zero. +func TestConsumePendingClaim_FloorAtZero(t *testing.T) { + s := &Server{} + + // Counter starts at 0; extra consume calls must be no-ops. + s.consumePendingClaim() + s.consumePendingClaim() + require.Equal(t, int32(0), s.pendingClaims.Load()) + + // One increment; one consume; counter back to 0. + s.pendingClaims.Inc() + require.Equal(t, int32(1), s.pendingClaims.Load()) + s.consumePendingClaim() + require.Equal(t, int32(0), s.pendingClaims.Load()) + + // Second consume after counter is already 0 is a no-op. + s.consumePendingClaim() + require.Equal(t, int32(0), s.pendingClaims.Load()) +} + +// TestConsumePendingClaim_NoDoubleDecrement verifies that concurrent consumptions +// of the same claim (StartEgress + self-decay timer racing) each fire exactly once +// and together decrement the counter by exactly N, not 2N. +func TestConsumePendingClaim_NoDoubleDecrement(t *testing.T) { + const claims = 50 + s := &Server{} + + for i := 0; i < claims; i++ { + s.pendingClaims.Inc() + } + require.Equal(t, int32(claims), s.pendingClaims.Load()) + + // Simulate StartEgress and self-decay racing concurrently for all claims. + // Both sides try to decrement; together they must not go below zero. + var wg sync.WaitGroup + for i := 0; i < claims*2; i++ { // 2× concurrency, one genuine + one spurious per claim + wg.Add(1) + go func() { + defer wg.Done() + s.consumePendingClaim() + }() + } + wg.Wait() + + require.Equal(t, int32(0), s.pendingClaims.Load(), "counter must be exactly 0, not negative") +} + func TestIsHeavyEgressRequest(t *testing.T) { for _, tc := range []struct { name string diff --git a/pkg/stats/monitor.go b/pkg/stats/monitor.go index 9f8e9e25d..410f25a27 100644 --- a/pkg/stats/monitor.go +++ b/pkg/stats/monitor.go @@ -616,6 +616,26 @@ func (m *Monitor) AvailableCPUFraction() float32 { return float32(available / total) } +// AvailableCPUFractionWithPending is like AvailableCPUFraction but deducts +// pendingSlots extra track-cost units from the available budget. Used by +// StartEgressAffinity in spread/type_aware modes to account for claimed-but-not-yet-accepted +// requests during burst windows (where m.requests.Inc has not yet fired). +func (m *Monitor) AvailableCPUFractionWithPending(pendingSlots int32) float32 { + m.mu.Lock() + defer m.mu.Unlock() + total, available, _, _ := m.getCPUUsageLocked() + if total == 0 { + return 1.0 + } + if pendingSlots > 0 { + available -= float64(pendingSlots) * m.cpuCostConfig.TrackCpuCost + if available < 0 { + available = 0 + } + } + return float32(available / total) +} + func (m *Monitor) getCPUUsageLocked() (total, available, pending, used float64) { total = m.cpuStats.NumCPU() if m.requests.Load() == 0 { From adb48d38b4f4b2d668ad7d23eabf25fa6c92e79c Mon Sep 17 00:00:00 2001 From: Abhishek Kumar Singh Date: Mon, 18 May 2026 11:30:46 +0530 Subject: [PATCH 3/5] fix(spread): idle pod returns 1.0 affinity to trigger ShortCircuitTimeout MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit In spread mode, an idle pod computed AvailableCPUFraction = 2.4/4.0 = 0.6. Since psrpc's MaximumAffinity is 1.0, the ShortCircuitTimeout (500ms fast path) never fired — every dispatch waited the full AffinityTimeout even when idle pods were ready. An idle pod (activeRequests=0, pendingClaims<=1) now returns 1.0 plus tiny jitter, triggering the 500ms short-circuit. Busy pods still return a proportional fraction. Also fix jitter direction bug in type_aware heavy-request path: 1.0 - jitter landed below MaximumAffinity; changed to 1.0 + jitter. (cherry picked from commit d2fc658d2e94db663e68d2f1f43b236da877d5ab) --- pkg/config/service.go | 14 +++++++++++ pkg/server/server_rpc.go | 51 +++++++++++++++++++++++++++++++++------- 2 files changed, 57 insertions(+), 8 deletions(-) diff --git a/pkg/config/service.go b/pkg/config/service.go index a91ef7a9d..186db36d3 100644 --- a/pkg/config/service.go +++ b/pkg/config/service.go @@ -74,6 +74,20 @@ type ServiceConfig struct { // "type_aware" — RoomComposite/Web prefer idle pods; Track/Participant spread by load AffinityMode string `yaml:"affinity_mode"` + // SoftRejectFloor is the affinity score returned when CanAcceptRequest is + // false but the pod has not yet reached MaxActiveRequests. This prevents + // the pod from going completely silent during transient CPU spikes (e.g. + // Chrome cold-start), keeping it in the dispatcher's selection pool as a + // last resort. Set to 0 (default) to disable and preserve upstream -1 behaviour. + // Recommended production value: 0.01 + SoftRejectFloor float32 `yaml:"soft_reject_floor"` + + // MaxActiveRequests is the hard capacity limit for this pod. When the pod + // has this many active recording jobs, it returns -1 (hard reject) even if + // SoftRejectFloor is set. Set to 0 to disable this guard. + // Recommended value: 16 (= interviews_per_pod × 4 tracks per interview) + MaxActiveRequests int32 `yaml:"max_active_requests"` + *CPUCostConfig `yaml:"cpu_cost"` // CPU costs for the different egress types } diff --git a/pkg/server/server_rpc.go b/pkg/server/server_rpc.go index 5864a1ec0..e2a1a0406 100644 --- a/pkg/server/server_rpc.go +++ b/pkg/server/server_rpc.go @@ -224,8 +224,12 @@ func (s *Server) processEnded(req *rpc.StartEgressRequest, info *livekit.EgressI } func (s *Server) StartEgressAffinity(_ context.Context, req *rpc.StartEgressRequest) float32 { - if s.IsDisabled() || !s.monitor.CanAcceptRequest(req) { - return -1 + if s.IsDisabled() { + return -1 // pod is shutting down — always hard reject + } + + if !s.monitor.CanAcceptRequest(req) { + return s.softRejectScore() } switch s.conf.AffinityMode { @@ -239,24 +243,36 @@ func (s *Server) StartEgressAffinity(_ context.Context, req *rpc.StartEgressRequ s.pendingClaims.Inc() time.AfterFunc(2*time.Second, s.consumePendingClaim) pending := s.pendingClaims.Load() - // Subtract rand jitter ≤ 0.001 to break equal-score ties in psrpc's strict-> - // comparison (first-replier-wins). Idle pods land in [0.999, 1.0]; busy pods - // score ≤ 0.96 (4 CPU, TrackCpuCost=0.15) and never beat idle pods. + + // An idle pod returns ≥ 1.0 to trigger psrpc's ShortCircuitTimeout (500ms + // fast path). Without this, AvailableCPUFractionWithPending returns ~0.6 for + // an idle pod (2.4/4.0), which never crosses MaximumAffinity=1.0, so every + // dispatch waits the full AffinityTimeout even when pods are completely free. + // Jitter is added (not subtracted) so the score stays ≥ 1.0. + if s.activeRequests.Load() == 0 && pending <= 1 { + return 1.0 + rand.Float32()*0.001 + } + return s.monitor.AvailableCPUFractionWithPending(pending) - rand.Float32()*0.001 case "type_aware": if isHeavyEgressRequest(req) { // RoomComposite/Web need ~4 CPUs. Strongly prefer idle pods. - // Jitter breaks first-replier ties among equally-idle pods. + // Jitter added (not subtracted) so idle score stays ≥ MaximumAffinity=1.0. if s.activeRequests.Load() == 0 { - return 1.0 - rand.Float32()*0.001 + return 1.0 + rand.Float32()*0.001 } return 0.5 } - // Light requests use the same pending-counter + jitter logic as spread. + // Light requests use the same pending-counter + idle-1.0 logic as spread. s.pendingClaims.Inc() time.AfterFunc(2*time.Second, s.consumePendingClaim) pending := s.pendingClaims.Load() + + if s.activeRequests.Load() == 0 && pending <= 1 { + return 1.0 + rand.Float32()*0.001 + } + return s.monitor.AvailableCPUFractionWithPending(pending) - rand.Float32()*0.001 default: // "pack" or empty — upstream behaviour unchanged @@ -267,6 +283,25 @@ func (s *Server) StartEgressAffinity(_ context.Context, req *rpc.StartEgressRequ } } +// softRejectScore returns the score to use when CanAcceptRequest is false. +// If SoftRejectFloor is non-zero and the pod is below its hard capacity limit, +// returns the floor so the pod still participates in dispatcher selection as a +// last resort. Otherwise returns -1 (silent / hard reject). +func (s *Server) softRejectScore() float32 { + floor := s.conf.SoftRejectFloor + + if floor <= 0 { + return -1 + } + + maxActive := s.conf.MaxActiveRequests + if maxActive > 0 && s.activeRequests.Load() >= maxActive { + return -1 + } + + return floor +} + func isHeavyEgressRequest(req *rpc.StartEgressRequest) bool { switch req.Request.(type) { case *rpc.StartEgressRequest_RoomComposite, From 9529af722ad57332f3709261e6a585ed6786b21b Mon Sep 17 00:00:00 2001 From: Abhishek Kumar Singh Date: Mon, 18 May 2026 11:30:53 +0530 Subject: [PATCH 4/5] feat(egress): add soft_reject_floor to prevent all-pods-silent drops When all egress pods are simultaneously over their CPU budget (Chrome cold-start storm), every pod returns -1 and the dispatcher gets zero bids. The job is permanently dropped. New config fields soft_reject_floor (float, default 0) and max_active_requests (int, default 0) allow a pod to return a small positive score instead of -1 when CanAcceptRequest is false but the pod is below its design capacity. The dispatcher can then select this pod as a last resort instead of dropping the job. Guarded by MaxActiveRequests so genuinely full pods still hard-reject. Also add unit tests for softRejectScore helper. (cherry picked from commit acf8a3da0946099162a7f58c8bf3e1988c78c64c) --- pkg/server/server_rpc_test.go | 32 ++++++++++++++++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/pkg/server/server_rpc_test.go b/pkg/server/server_rpc_test.go index 9a90d44eb..9847e8899 100644 --- a/pkg/server/server_rpc_test.go +++ b/pkg/server/server_rpc_test.go @@ -22,6 +22,8 @@ import ( "github.com/livekit/protocol/livekit" "github.com/livekit/protocol/rpc" + + "github.com/livekit/egress/pkg/config" ) // TestConsumePendingClaim_FloorAtZero verifies the CAS loop never decrements below zero. @@ -71,6 +73,36 @@ func TestConsumePendingClaim_NoDoubleDecrement(t *testing.T) { require.Equal(t, int32(0), s.pendingClaims.Load(), "counter must be exactly 0, not negative") } +// --- Tests for softRejectScore (TASK-03) --- + +func TestSoftRejectFloorDisabled(t *testing.T) { + // SoftRejectFloor=0 → feature disabled → always return -1 + s := &Server{conf: &config.ServiceConfig{SoftRejectFloor: 0}} + require.Equal(t, float32(-1), s.softRejectScore()) +} + +func TestSoftRejectFloorReturnedWhenBelowMax(t *testing.T) { + // Floor set, activeRequests < MaxActiveRequests → return floor + s := &Server{conf: &config.ServiceConfig{SoftRejectFloor: 0.01, MaxActiveRequests: 16}} + s.activeRequests.Store(8) + require.Equal(t, float32(0.01), s.softRejectScore()) +} + +func TestSoftRejectFloorHardRejectWhenAtMax(t *testing.T) { + // Floor set, activeRequests >= MaxActiveRequests → -1 (genuinely full) + s := &Server{conf: &config.ServiceConfig{SoftRejectFloor: 0.01, MaxActiveRequests: 16}} + s.activeRequests.Store(16) + require.Equal(t, float32(-1), s.softRejectScore()) +} + +func TestSoftRejectFloorNoGuardWhenMaxIsZero(t *testing.T) { + // MaxActiveRequests=0 → guard disabled → return floor regardless of load + s := &Server{conf: &config.ServiceConfig{SoftRejectFloor: 0.01, MaxActiveRequests: 0}} + s.activeRequests.Store(20) + require.Equal(t, float32(0.01), s.softRejectScore()) +} + + func TestIsHeavyEgressRequest(t *testing.T) { for _, tc := range []struct { name string From bad7fa92e07a89a35800966acfb2d05daf424766 Mon Sep 17 00:00:00 2001 From: Abhishek Kumar Singh Date: Mon, 18 May 2026 12:16:49 +0530 Subject: [PATCH 5/5] revert: remove unused type_aware changes and CI workflow MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit type_aware mode is not used in production (affinity_mode: spread in both config blocks). Reverts idle-1.0 and jitter-direction changes in that branch to keep the diff minimal and production-relevant only. build-push-ar.yml removed — images are built and pushed manually. (cherry picked from commit 3c412585eff7bd4c9d6e24c4e3fa234cf4bf1da1) --- pkg/server/server_rpc.go | 10 ++-------- 1 file changed, 2 insertions(+), 8 deletions(-) diff --git a/pkg/server/server_rpc.go b/pkg/server/server_rpc.go index e2a1a0406..aa9829b5b 100644 --- a/pkg/server/server_rpc.go +++ b/pkg/server/server_rpc.go @@ -258,21 +258,15 @@ func (s *Server) StartEgressAffinity(_ context.Context, req *rpc.StartEgressRequ case "type_aware": if isHeavyEgressRequest(req) { // RoomComposite/Web need ~4 CPUs. Strongly prefer idle pods. - // Jitter added (not subtracted) so idle score stays ≥ MaximumAffinity=1.0. if s.activeRequests.Load() == 0 { - return 1.0 + rand.Float32()*0.001 + return 1.0 - rand.Float32()*0.001 } return 0.5 } - // Light requests use the same pending-counter + idle-1.0 logic as spread. + // Light requests use the same pending-counter + jitter logic as spread. s.pendingClaims.Inc() time.AfterFunc(2*time.Second, s.consumePendingClaim) pending := s.pendingClaims.Load() - - if s.activeRequests.Load() == 0 && pending <= 1 { - return 1.0 + rand.Float32()*0.001 - } - return s.monitor.AvailableCPUFractionWithPending(pending) - rand.Float32()*0.001 default: // "pack" or empty — upstream behaviour unchanged