mirror of
https://github.com/warmbly/warmbly.git
synced 2026-09-12 08:04:41 +00:00
feat: make worker capacity a placement target rather than a hard gate, so Eligible refuses only on health and an over-target worker costs enough score to lose to anything with room instead of returning nil and dropping assignment into selectFallback, score projected utilization including the incoming mailbox's own weight, and penalise the 454/421 auth pressure the capacity view already collected and threw away
This commit is contained in:
@@ -312,7 +312,9 @@ In production, workers are treated as individually addressable executors:
|
||||
- worker events are delivered through worker-specific Kafka topics
|
||||
- the platform can rebalance or migrate accounts between workers, reluctantly
|
||||
|
||||
Placement is a score, never a filter (`internal/app/worker/placement.go`). Hard constraints cover only whether the work can be done: heartbeating, health in `healthy`/`watch`, and enough capacity headroom for the mailbox's weight. Everything else is a preference term: capacity headroom, incumbency (weighted highest), region match, tenant blast radius, per-provider crowding on one address, and foreign tenants for orgs entitled to isolated egress.
|
||||
Placement is a score, never a filter (`internal/app/worker/placement.go`). Hard constraints cover only whether the work can be done: heartbeating and health in `healthy`/`watch`. Everything else is a preference term: capacity headroom (projected, so the incoming mailbox's own weight counts), incumbency (weighted highest), region match, tenant blast radius, per-provider crowding on one address, auth pressure from `454`/`421` throttles on that address, and foreign tenants for orgs entitled to isolated egress.
|
||||
|
||||
Capacity is a target, not a ceiling, and nothing refuses a placement for being over it. Over-target costs more score than any bonus a candidate can earn, so a worker with room always wins when one exists; when none does, the least-overloaded wins with every other preference still applied. Do not put capacity back into `Eligible`: `base_capacity` is a flat `16` for every worker regardless of the machine, so refusing on it refuses on a guess, and it refused precisely when the fleet was full, dropping assignment into `selectFallback` (first healthy worker, no region, no blast radius, no provider crowding).
|
||||
|
||||
Capacity is one number for every worker in cold-mailbox equivalents, because each mailbox declares its own cost through `MailboxWeight`: `smtp_imap` = 1.0, `gmail`/`outlook` = 0.05, warmup-only = 0.4. Those are the `email_provider` enum values as stored; do not invent provider strings for them.
|
||||
|
||||
|
||||
@@ -114,7 +114,7 @@ So the levers invert: IP *stability* per mailbox beats IP diversity, and a migra
|
||||
|
||||
### Choosing a worker
|
||||
|
||||
Placement scores every live worker and takes the best (`internal/app/worker/placement.go`). Hard constraints are only about whether the work can be done at all: the worker has to be heartbeating, in `healthy` or `watch`, and have capacity headroom for the mailbox's weight. Everything else is a preference:
|
||||
Placement scores every live worker and takes the best (`internal/app/worker/placement.go`). Hard constraints are only about whether the work can be done at all: the worker has to be heartbeating and in `healthy` or `watch`. Everything else is a preference:
|
||||
|
||||
| Term | Why |
|
||||
|---|---|
|
||||
@@ -123,9 +123,12 @@ Placement scores every live worker and takes the best (`internal/app/worker/plac
|
||||
| Region match | Sign-ins from where the provider expects them raise fewer challenges |
|
||||
| Tenant blast radius | Spread one customer across workers so a single failure does not stop their sending |
|
||||
| Provider crowding | Many accounts of one provider signing in from one address is what earns a per-IP throttle |
|
||||
| Auth pressure | `454`/`421` throttles on a worker's address mean the provider is already pushing back on sign-ins from there |
|
||||
| Foreign tenants | Only for organizations entitled to isolated egress |
|
||||
|
||||
Capacity is one number for every worker, in cold-mailbox equivalents, because each mailbox already declares its own cost: an `smtp_imap` mailbox weighs `1.0`, a Gmail or Outlook API mailbox `0.05`, and a warmup-only assignment `0.4`.
|
||||
Capacity is one number for every worker, in cold-mailbox equivalents, because each mailbox already declares its own cost: an `smtp_imap` mailbox weighs `1.0`, a Gmail or Outlook API mailbox `0.05`, and a warmup-only assignment `0.4`. The mailbox being placed counts against the target too, so the same nearly-full worker can be under target for a Gmail mailbox and over it for an `smtp_imap` one.
|
||||
|
||||
**Capacity is a target, not a ceiling.** Nothing refuses a placement for being over it. Going over costs enough score that a worker with room wins whenever one exists, even against an incumbent; once no worker has room, the least-overloaded one still wins with region, blast radius and provider crowding all applied. This is deliberate: the target is a flat `16` for every worker regardless of the machine, so refusing on it refuses on a guess, and refusing exactly when the fleet is full is when good placement matters most.
|
||||
|
||||
### Moving a mailbox
|
||||
|
||||
|
||||
@@ -91,8 +91,9 @@ type PlacementResult struct {
|
||||
Worker *models.Worker
|
||||
Score float64
|
||||
IncumbentScore float64
|
||||
// IncumbentEligible is false when the current worker could not host the
|
||||
// mailbox at all, in which case IncumbentScore is meaningless.
|
||||
// IncumbentEligible is false when the current worker is too unhealthy to
|
||||
// host the mailbox at all, in which case IncumbentScore is meaningless.
|
||||
// Being over capacity does not clear it; that shows up in IncumbentScore.
|
||||
IncumbentEligible bool
|
||||
// Mandated is set when the worker was chosen by an entitlement rather than
|
||||
// by scoring, which today means an isolated-egress reservation. Callers
|
||||
|
||||
@@ -42,6 +42,9 @@ type WorkerCapacityRow struct {
|
||||
// dimensionless except Effective (mailbox-equivalents) and Load (sum of
|
||||
// mailbox weights). Utilization is Load/Effective and is the value the
|
||||
// scheduler sorts on when picking the next worker.
|
||||
//
|
||||
// Effective is a target, not a ceiling. Nothing refuses a placement for being
|
||||
// over it; see Score in placement.go for what being over it costs.
|
||||
type Capacity struct {
|
||||
Base float64
|
||||
HealthMul float64
|
||||
@@ -49,6 +52,10 @@ type Capacity struct {
|
||||
Effective float64
|
||||
Load float64
|
||||
Utilization float64
|
||||
|
||||
// AuthPressure is how hard this worker's address is being pushed back on
|
||||
// by mailbox providers, in [0, 1]. See computeAuthPressure.
|
||||
AuthPressure float64
|
||||
}
|
||||
|
||||
// ComputeCapacity is the placement math: Base * Health * Age, floored at
|
||||
@@ -69,9 +76,37 @@ func ComputeCapacity(row WorkerCapacityRow) Capacity {
|
||||
if c.Effective > 0 {
|
||||
c.Utilization = c.Load / c.Effective
|
||||
}
|
||||
c.AuthPressure = computeAuthPressure(row.AuthErrors1h, row.SendsAttempted1h)
|
||||
return c
|
||||
}
|
||||
|
||||
// authPressureFullScale is the auth-error rate that saturates the pressure
|
||||
// signal. Auth errors are the one thing in worker_health_samples that is
|
||||
// genuinely about the worker's own address rather than about the mailbox or
|
||||
// its contact list: a per-IP authentication throttle (454 4.7.0) or per-IP
|
||||
// rate limit (421 4.7.28) is the provider saying this client is signing in
|
||||
// too much. Bounces and complaints follow the mailbox and say nothing about
|
||||
// where it is running from.
|
||||
//
|
||||
// 5% is already a bad hour, so that is full scale.
|
||||
const authPressureFullScale = 0.05
|
||||
|
||||
// computeAuthPressure scales the last hour's auth-error rate into [0, 1].
|
||||
//
|
||||
// The denominator floors at 1 rather than at attempts, because auth errors
|
||||
// also arrive from sync on a worker that sent nothing. Errors with no attempts
|
||||
// is the worst case there is, and it saturates, which is the right answer.
|
||||
func computeAuthPressure(authErrors, sendsAttempted int64) float64 {
|
||||
if authErrors <= 0 {
|
||||
return 0
|
||||
}
|
||||
denom := float64(sendsAttempted)
|
||||
if denom < 1 {
|
||||
denom = 1
|
||||
}
|
||||
return clampUnit((float64(authErrors) / denom) / authPressureFullScale)
|
||||
}
|
||||
|
||||
// clampUnit pins x into [0, 1]. Negatives are surprisingly easy to feed
|
||||
// in - PostgreSQL's NULLIF/divide-by-zero handling can leak NaN through
|
||||
// in pathological cases - so we coerce defensively here.
|
||||
|
||||
@@ -198,3 +198,29 @@ func TestComputeCapacity_HealthStatesArePassedThrough(t *testing.T) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestComputeAuthPressure(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
authErrors int64
|
||||
attempted int64
|
||||
want float64
|
||||
}{
|
||||
{"clean worker", 0, 500, 0},
|
||||
{"one error in a thousand", 1, 1000, 0.02},
|
||||
{"one percent", 10, 1000, 0.2},
|
||||
{"full scale at five percent", 50, 1000, 1},
|
||||
{"saturates past full scale", 500, 1000, 1},
|
||||
// Sync auth failures arrive on a worker that sent nothing; the
|
||||
// denominator floor makes that saturate rather than divide by zero.
|
||||
{"errors with no attempts", 3, 0, 1},
|
||||
}
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
got := computeAuthPressure(tc.authErrors, tc.attempted)
|
||||
if math.Abs(got-tc.want) > 1e-9 {
|
||||
t.Fatalf("got %v want %v", got, tc.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -41,6 +41,9 @@ const (
|
||||
weightIsolation = 1.5
|
||||
weightBlastRadius = 0.8
|
||||
weightProviderLoad = 0.7
|
||||
weightAuthPressure = 1.2
|
||||
weightOverTarget = 3.0
|
||||
weightOverload = 4.0
|
||||
)
|
||||
|
||||
// providerSoftCap is how many mailboxes of ONE provider a single worker is
|
||||
@@ -91,16 +94,19 @@ type PlacementRequest struct {
|
||||
IsolatedEgress bool
|
||||
}
|
||||
|
||||
// Eligible reports whether a candidate may host the mailbox at all. These are
|
||||
// the only hard constraints left: the worker has to be able to do the work.
|
||||
// Everything else is a preference expressed in the score.
|
||||
// Eligible reports whether a candidate may host the mailbox at all. Health is
|
||||
// the only hard constraint: capacity is a preference in the score, because
|
||||
// base_capacity is a flat 16 for every worker and refusing on it refused on a
|
||||
// guess. It also refused where placement mattered most - a full fleet returned
|
||||
// nil and fell through to selectFallback, which reads none of region, blast
|
||||
// radius or provider crowding.
|
||||
func (c PlacementCandidate) Eligible(req PlacementRequest) bool {
|
||||
switch c.Health {
|
||||
case models.WorkerHealthHealthy, models.WorkerHealthWatch:
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
return c.Capacity.Effective-c.Capacity.Load >= req.Weight
|
||||
}
|
||||
|
||||
// Score ranks an eligible candidate. Higher is better. The terms are additive
|
||||
@@ -109,12 +115,29 @@ func (c PlacementCandidate) Eligible(req PlacementRequest) bool {
|
||||
func (c PlacementCandidate) Score(req PlacementRequest) float64 {
|
||||
var score float64
|
||||
|
||||
// Capacity: prefer the worker with the most room, so the fleet fills evenly.
|
||||
utilization := c.Capacity.Utilization
|
||||
if utilization > 1 {
|
||||
utilization = 1
|
||||
// Capacity: prefer the worker with the most room, so the fleet fills
|
||||
// evenly. Projected, so the incoming mailbox's own weight counts - Eligible
|
||||
// used to be the only thing that read req.Weight.
|
||||
projected := c.Capacity.Utilization
|
||||
if c.Capacity.Effective > 0 {
|
||||
projected = (c.Capacity.Load + req.Weight) / c.Capacity.Effective
|
||||
}
|
||||
score += weightHeadroom * (1 - utilization)
|
||||
if projected <= 1 {
|
||||
score += weightHeadroom * (1 - projected)
|
||||
} else {
|
||||
// A soft wall: the flat cost exceeds the largest bonus a candidate can
|
||||
// earn (incumbency 2.0 + region 0.6), so anything with room wins when
|
||||
// it exists, and the ramp still ranks the overloaded against each
|
||||
// other when nothing does.
|
||||
over := projected - 1
|
||||
score -= weightOverTarget + weightOverload*over
|
||||
}
|
||||
|
||||
// Provider pushback: 454/421 auth throttles mean this address is signing
|
||||
// in too much, so steer new mailboxes away. Below incumbency on purpose -
|
||||
// evacuating a throttled worker is rotation's call, behind a residency
|
||||
// floor, not a per-mailbox sign-in challenge paid here.
|
||||
score -= weightAuthPressure * c.Capacity.AuthPressure
|
||||
|
||||
// Stickiness: the incumbent wins ties and most non-ties. A mailbox that
|
||||
// stays put keeps presenting the same client IP to its provider.
|
||||
|
||||
@@ -20,15 +20,17 @@ func candidate(id uuid.UUID, effective, load float64) PlacementCandidate {
|
||||
return c
|
||||
}
|
||||
|
||||
func TestEligibleRequiresHealthAndHeadroom(t *testing.T) {
|
||||
func TestEligibleRequiresHealthOnly(t *testing.T) {
|
||||
id := uuid.New()
|
||||
req := PlacementRequest{Weight: 1.0}
|
||||
|
||||
if c := candidate(id, 16, 4); !c.Eligible(req) {
|
||||
t.Fatal("healthy worker with headroom should be eligible")
|
||||
}
|
||||
if c := candidate(id, 16, 15.5); c.Eligible(req) {
|
||||
t.Fatal("worker with less headroom than the mailbox weight should not be eligible")
|
||||
// Capacity is a preference, not a fence: an over-target worker is still a
|
||||
// legal home, it just scores badly. See TestOverTargetWorkerIsLastResort.
|
||||
if c := candidate(id, 16, 40); !c.Eligible(req) {
|
||||
t.Fatal("being over the capacity target must not refuse a placement")
|
||||
}
|
||||
|
||||
for _, state := range []models.WorkerHealthState{
|
||||
@@ -154,14 +156,77 @@ func TestRegionMatchIsAPreferenceNotARequirement(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestSelectPlacementReturnsNilWhenNothingFits(t *testing.T) {
|
||||
full := candidate(uuid.New(), 16, 16)
|
||||
if got := SelectPlacement([]PlacementCandidate{full}, PlacementRequest{Weight: 1.0}); got != nil {
|
||||
t.Fatal("expected no placement when every worker is full")
|
||||
func TestSelectPlacementReturnsNilOnlyWhenNothingIsHealthy(t *testing.T) {
|
||||
sick := candidate(uuid.New(), 16, 0)
|
||||
sick.Health = models.WorkerHealthQuarantined
|
||||
if got := SelectPlacement([]PlacementCandidate{sick}, PlacementRequest{Weight: 1.0}); got != nil {
|
||||
t.Fatal("expected no placement when every worker is unhealthy")
|
||||
}
|
||||
if got := SelectPlacement(nil, PlacementRequest{Weight: 1.0}); got != nil {
|
||||
t.Fatal("expected no placement from an empty fleet")
|
||||
}
|
||||
|
||||
// A full fleet is not an empty one. Returning nil here used to drop
|
||||
// assignment into selectFallback, which ignores every preference term.
|
||||
full := candidate(uuid.New(), 16, 16)
|
||||
if got := SelectPlacement([]PlacementCandidate{full}, PlacementRequest{Weight: 1.0}); got == nil {
|
||||
t.Fatal("a full but healthy fleet must still place the mailbox")
|
||||
}
|
||||
}
|
||||
|
||||
func TestOverTargetWorkerIsLastResort(t *testing.T) {
|
||||
over, room := uuid.New(), uuid.New()
|
||||
|
||||
// The over-target worker is also the incumbent and matches the region, so
|
||||
// it collects every bonus available. It must still lose to spare capacity.
|
||||
a := candidate(over, 16, 18)
|
||||
a.Region = "eu"
|
||||
b := candidate(room, 16, 15)
|
||||
|
||||
got := SelectPlacement([]PlacementCandidate{a, b}, PlacementRequest{
|
||||
Weight: 1.0, CurrentWorkerID: &over, Region: "eu",
|
||||
})
|
||||
if got == nil || got.WorkerID != room {
|
||||
t.Fatal("a worker with room must beat an over-target one holding every bonus")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAmongOverloadedWorkersTheLeastOverloadedWins(t *testing.T) {
|
||||
bad, worse := uuid.New(), uuid.New()
|
||||
|
||||
got := SelectPlacement([]PlacementCandidate{
|
||||
candidate(worse, 16, 40),
|
||||
candidate(bad, 16, 18),
|
||||
}, PlacementRequest{Weight: 1.0})
|
||||
if got == nil || got.WorkerID != bad {
|
||||
t.Fatal("with no room anywhere, the least-overloaded worker should win")
|
||||
}
|
||||
}
|
||||
|
||||
func TestMailboxWeightCountsAgainstTheTarget(t *testing.T) {
|
||||
id := uuid.New()
|
||||
c := candidate(id, 16, 15.5)
|
||||
|
||||
// The same worker is under target for a Gmail mailbox and over it for an
|
||||
// smtp_imap one. Only the projected utilization can tell them apart.
|
||||
light := c.Score(PlacementRequest{Weight: MailboxWeight("gmail", false)})
|
||||
heavy := c.Score(PlacementRequest{Weight: MailboxWeight("smtp_imap", false)})
|
||||
if !(light > 0 && heavy < -weightOverTarget+1) {
|
||||
t.Fatalf("weight should change the verdict: light=%.3f heavy=%.3f", light, heavy)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAuthPressureSteersAwayFromAThrottledAddress(t *testing.T) {
|
||||
throttled, clean := uuid.New(), uuid.New()
|
||||
|
||||
a := candidate(throttled, 16, 8)
|
||||
a.Capacity.AuthPressure = 1
|
||||
b := candidate(clean, 16, 8)
|
||||
|
||||
got := SelectPlacement([]PlacementCandidate{a, b}, PlacementRequest{Weight: 1.0})
|
||||
if got == nil || got.WorkerID != clean {
|
||||
t.Fatal("placement should avoid an address the provider is auth-throttling")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSelectPlacementIsDeterministicOnTies(t *testing.T) {
|
||||
|
||||
Reference in New Issue
Block a user