feat: address worker capacity review findings with safe migrations and recovery reporting

This commit is contained in:
Matthew Meszaros
2026-09-17 06:24:14 -07:00
parent bab9f86727
commit d6025ea4c0
21 changed files with 337 additions and 120 deletions
+1
View File
@@ -61,6 +61,7 @@ export interface AdminWorkerEmail {
risk_evaluated_at?: string | null;
warmup_health?: string; // worst warmup health_state, "" if not in a pool
blocked_until?: string | null;
created_at: string;
}
export interface AdminWorkerEmailsResult {
+9
View File
@@ -17,4 +17,13 @@ func TestWorkerCapacityTarget(t *testing.T) {
if got := workerCapacityTarget(); got != 100 {
t.Fatalf("invalid capacity = %.0f, want safe default 100", got)
}
for _, invalid := range []string{"NaN", "+Inf", "100000000"} {
t.Run(invalid, func(t *testing.T) {
t.Setenv("WARMBLY_WORKER_CAPACITY", invalid)
if got := workerCapacityTarget(); got != 100 {
t.Fatalf("capacity %q = %v, want safe default 100", invalid, got)
}
})
}
}
+2 -1
View File
@@ -3,6 +3,7 @@ package main
import (
"context"
"log"
"math"
"os"
"os/signal"
"strconv"
@@ -266,7 +267,7 @@ func workerCapacityTarget() float64 {
return defaultWorkerCapacityTarget
}
n, err := strconv.ParseFloat(raw, 64)
if err != nil || n <= 0 {
if err != nil || n <= 0 || math.IsNaN(n) || math.IsInf(n, 0) || n > 99_999_999.99 {
return defaultWorkerCapacityTarget
}
return n
@@ -136,10 +136,12 @@ Rotation (`internal/app/worker/rotation.go`) is deliberately reluctant. Mailbox
| Urgency | Trigger | Residency floor | Destination bar |
|---|---|---|---|
| Immediate | Worker inactive or not heartbeating | none | anything eligible |
| Elevated | Mailbox is on another workspace's reserved worker | 6 hours | anything eligible |
| Immediate | Worker inactive, not heartbeating, or externally marked blocked or quarantined | none | anything eligible |
| Elevated | Worker externally marked throttled, or mailbox is on another workspace's reserved worker | 6 hours | anything eligible |
| Opportunistic | Worker over 85% utilization, isolated-egress drift, or placement concentration | 72 hours | must score materially better |
The externally managed health labels describe machine or operator state. Mailbox delivery outcomes do not set them.
### Isolated egress
Plans that reserve egress bind an organization to a worker (`dedicated_worker_assignments`). It is a strong placement preference, not a pin: the worker carries no marking, so if it dies the organization's mailboxes place normally instead of stranding, and the rotation loop pulls them onto a replacement once one exists. Mailboxes belonging to other tenants drift off a reserved worker on the same loop.
@@ -43,7 +43,7 @@ X-Warmbly-Event: <event key>
{
"event": "worker.offline",
"title": "Worker stopped responding",
"summary": "No healthy replacement of the same tier was available…",
"summary": "Some affected mailboxes were moved; the remaining mailboxes are paused.",
"severity": "urgent",
"fields": [{ "label": "Worker", "value": "…" }],
"link": "https://app.example.com",
@@ -64,7 +64,7 @@ A chat webhook URL is a bearer credential: anyone holding it can post to that ch
| Payment failed | Stripe could not collect an invoice |
| Workspace created | A new organization is created |
| User registered | A new account finishes signing up |
| Worker went offline | A worker stops heartbeating; the alert reports whether affected mailboxes were reassigned or paused |
| Worker went offline | A worker stops heartbeating; the alert reports complete, partial, or failed recovery with moved and paused mailbox counts |
| Warmup ban appealed | A blocked mailbox asks to rejoin the warmup pool |
| Workspace risk escalated | Risk scoring moves a workspace into a different posture |
+7 -2
View File
@@ -11,6 +11,7 @@ import (
"github.com/warmbly/warmbly/internal/api/middleware"
"github.com/warmbly/warmbly/internal/errx"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/utils/paging"
)
// User Management Handlers
@@ -297,10 +298,14 @@ func (h *Handler) AdminGetWorkerEmails(c *gin.Context) {
return
}
cursor := parseCursor(c.Query("cursor"))
beforeAt, beforeID, cursorErr := paging.DecodeTimeCursor(c.Query("cursor"))
if cursorErr != nil {
errx.JSON(c, cursorErr)
return
}
limit := parseLimit(c.Query("limit"), 50)
emails, pagination, xerr := h.AdminService.GetWorkerEmails(c.Request.Context(), workerID, cursor, limit)
emails, pagination, xerr := h.AdminService.GetWorkerEmails(c.Request.Context(), workerID, beforeAt, beforeID, limit)
if xerr != nil {
errx.JSON(c, xerr)
return
+9 -9
View File
@@ -185,11 +185,11 @@ func (h *Handler) FleetHeartbeat(c *gin.Context) {
func heartbeatAddress(reported, observed string) string {
reported = strings.TrimSpace(reported)
observed = strings.TrimSpace(observed)
if publicIPv4(observed) {
return observed
if normalized, ok := normalizedPublicIPv4(observed); ok {
return normalized
}
if publicIPv4(reported) {
return reported
if normalized, ok := normalizedPublicIPv4(reported); ok {
return normalized
}
if observed != "" {
return observed
@@ -197,21 +197,21 @@ func heartbeatAddress(reported, observed string) string {
return reported
}
func publicIPv4(raw string) bool {
func normalizedPublicIPv4(raw string) (string, bool) {
ip, err := netip.ParseAddr(raw)
if err != nil {
return false
return "", false
}
ip = ip.Unmap()
if !ip.Is4() || !ip.IsGlobalUnicast() || ip.IsPrivate() {
return false
return "", false
}
for _, prefix := range nonPublicIPv4Prefixes {
if prefix.Contains(ip) {
return false
return "", false
}
}
return true
return ip.String(), true
}
var nonPublicIPv4Prefixes = []netip.Prefix{
+1
View File
@@ -18,6 +18,7 @@ func TestHeartbeatAddressPrefersBackendObservedPublicIPv4(t *testing.T) {
{"observed public address replaces stale report", "198.51.100.8", "8.8.8.8", "8.8.8.8"},
{"private proxy hop keeps reported public address", "1.1.1.1", "10.0.0.4", "1.1.1.1"},
{"carrier grade nat is not a public address", "1.1.1.1", "100.64.2.3", "1.1.1.1"},
{"mapped public IPv4 is canonicalized", "1.1.1.1", "::ffff:8.8.8.8", "8.8.8.8"},
{"private direct address is still useful", "", "10.0.0.4", "10.0.0.4"},
{"public IPv6 does not replace requested IPv4", "1.1.1.1", "2001:4860:4860::8888", "1.1.1.1"},
}
+3 -3
View File
@@ -32,7 +32,7 @@ type AdminService interface {
ListWorkers(ctx context.Context, cursor *uuid.UUID, limit int) (*models.AdminWorkersResult, *errx.Error)
GetWorkerDetail(ctx context.Context, workerID uuid.UUID) (*models.AdminWorkerDetail, *errx.Error)
UpdateWorker(ctx context.Context, adminID, workerID uuid.UUID, update *models.AdminUpdateWorker, ipAddress, userAgent string) *errx.Error
GetWorkerEmails(ctx context.Context, workerID uuid.UUID, cursor *uuid.UUID, limit int) ([]models.AdminWorkerEmail, *models.Pagination, *errx.Error)
GetWorkerEmails(ctx context.Context, workerID uuid.UUID, beforeAt time.Time, beforeID uuid.UUID, limit int) ([]models.AdminWorkerEmail, *models.Pagination, *errx.Error)
GetWorkerStats(ctx context.Context, workerID uuid.UUID) (*models.WorkerStats, *errx.Error)
ReassignEmails(ctx context.Context, adminID uuid.UUID, emailIDs []uuid.UUID, newWorkerID uuid.UUID, ipAddress, userAgent string) *errx.Error
@@ -316,8 +316,8 @@ func (s *adminService) UpdateWorker(ctx context.Context, adminID, workerID uuid.
return nil
}
func (s *adminService) GetWorkerEmails(ctx context.Context, workerID uuid.UUID, cursor *uuid.UUID, limit int) ([]models.AdminWorkerEmail, *models.Pagination, *errx.Error) {
emails, pagination, err := s.repo.GetWorkerEmails(ctx, workerID, cursor, limit)
func (s *adminService) GetWorkerEmails(ctx context.Context, workerID uuid.UUID, beforeAt time.Time, beforeID uuid.UUID, limit int) ([]models.AdminWorkerEmail, *models.Pagination, *errx.Error) {
emails, pagination, err := s.repo.GetWorkerEmails(ctx, workerID, beforeAt, beforeID, limit)
if err != nil {
errs.CaptureException(err)
return nil, nil, errx.New(errx.Internal, "failed to get worker emails")
+138 -66
View File
@@ -3,7 +3,9 @@ package jobs
import (
"context"
"fmt"
"sort"
"strconv"
"strings"
"time"
"github.com/google/uuid"
@@ -26,10 +28,47 @@ type OperatorNotifier interface {
NotifyOperator(key, title, summary string, fields map[string]string)
}
// notifyOperatorWorkerDown alerts the operator that a worker is gone. It shares
// the same once-per-incident SetNX guard shape as the tenant notice, under its
// own key so the two audiences are independent.
func (s *JobsService) notifyOperatorWorkerDown(ctx context.Context, workerID uuid.UUID, mailboxes int, reassigned bool) {
type workerRecoveryOutcome struct {
Total int
Reassigned int
Stranded int
FailureReasons map[string]int
}
func (o workerRecoveryOutcome) state() string {
switch {
case o.Stranded == 0:
return "complete"
case o.Reassigned > 0:
return "partial"
default:
return "failed"
}
}
func (o workerRecoveryOutcome) failureSummary() string {
if len(o.FailureReasons) == 0 {
return "none"
}
keys := make([]string, 0, len(o.FailureReasons))
for reason := range o.FailureReasons {
keys = append(keys, reason)
}
sort.Strings(keys)
parts := make([]string, 0, len(keys))
for _, reason := range keys {
parts = append(parts, fmt.Sprintf("%s: %d", strings.ReplaceAll(reason, "_", " "), o.FailureReasons[reason]))
}
return strings.Join(parts, ", ")
}
type orgRecoveryOutcome struct {
Reassigned int
Stranded int
}
// notifyOperatorWorkerDown alerts the operator with the recovery result.
func (s *JobsService) notifyOperatorWorkerDown(ctx context.Context, workerID uuid.UUID, outcome workerRecoveryOutcome) {
if s.OpsNotifier == nil || s.Cache == nil {
return
}
@@ -37,18 +76,23 @@ func (s *JobsService) notifyOperatorWorkerDown(ctx context.Context, workerID uui
if err != nil || !ok {
return
}
outcome := "Mailboxes were moved to a healthy worker automatically."
if !reassigned {
outcome = "No healthy replacement had room, so sending from those mailboxes is paused."
summary := "All affected mailboxes were moved to eligible live workers."
if outcome.state() == "partial" {
summary = "Some affected mailboxes were moved; the remaining mailboxes are paused."
} else if outcome.state() == "failed" {
summary = "Recovery failed; affected mailboxes remain paused."
}
s.OpsNotifier.NotifyOperator(
"worker.offline",
"Worker stopped responding",
outcome,
summary,
map[string]string{
"Worker": workerID.String(),
"Mailboxes": strconv.Itoa(mailboxes),
"Reassigned": map[bool]string{true: "yes", false: "no"}[reassigned],
"Worker": workerID.String(),
"Outcome": outcome.state(),
"Mailboxes": strconv.Itoa(outcome.Total),
"Reassigned": strconv.Itoa(outcome.Reassigned),
"Stranded": strconv.Itoa(outcome.Stranded),
"Failure reasons": outcome.failureSummary(),
},
)
}
@@ -58,31 +102,36 @@ func (s *JobsService) notifyOperatorWorkerDown(ctx context.Context, workerID uui
// mailboxes for many orgs, and one org's alert must not suppress another's if
// reassignment completes across multiple scans. The shared group key coalesces
// an org's recipients into one email with everyone in To.
func (s *JobsService) notifyWorkerDown(ctx context.Context, workerID uuid.UUID, orgs map[uuid.UUID]int, reassigned bool) {
func (s *JobsService) notifyWorkerDown(ctx context.Context, workerID uuid.UUID, orgs map[uuid.UUID]orgRecoveryOutcome) {
if s.Notifier == nil || len(orgs) == 0 {
return
}
for orgID, n := range orgs {
for orgID, outcome := range orgs {
key := "worker:downnotify:" + workerID.String() + ":" + orgID.String()
ok, err := s.Cache.SetNX(ctx, key, "1", 6*time.Hour).Result()
if err != nil || !ok {
continue
}
noun := fmt.Sprintf("%d of your mailboxes were", n)
if n == 1 {
total := outcome.Reassigned + outcome.Stranded
noun := fmt.Sprintf("%d of your mailboxes were", total)
if total == 1 {
noun = "One of your mailboxes was"
}
body := noun + " on a sending worker that stopped responding. "
if reassigned {
switch {
case outcome.Stranded == 0:
body += "They were moved to a healthy worker automatically; no action is needed."
if n == 1 {
if total == 1 {
body = noun + " on a sending worker that stopped responding. It was moved to a healthy worker automatically; no action is needed."
}
} else {
case outcome.Reassigned == 0:
body += "Sending from them is paused until a replacement worker is available."
default:
body += fmt.Sprintf("%d were moved automatically; sending from the remaining %d is paused.", outcome.Reassigned, outcome.Stranded)
}
meta := map[string]any{"reassigned": outcome.Reassigned, "stranded": outcome.Stranded}
s.Notifier.NotifyOrg(ctx, orgID, models.PermManageEmails, uuid.Nil, models.NotifWorkerDowntime,
"Sending worker went offline", body, "/app/emails", nil,
"Sending worker went offline", body, "/app/emails", meta,
"worker_down:"+workerID.String())
}
}
@@ -141,13 +190,28 @@ func (s *JobsService) detectDeadWorkers(ctx context.Context) {
}
// Score each mailbox independently so recovery does not create a hotspot.
reassigned := 0
affectedOrgs := map[uuid.UUID]int{}
outcome := workerRecoveryOutcome{Total: len(accountIDs), FailureReasons: map[string]int{}}
affectedOrgs := map[uuid.UUID]orgRecoveryOutcome{}
destinations := map[uuid.UUID]struct{}{}
recordFailure := func(reason string, orgID *uuid.UUID) {
outcome.Stranded++
outcome.FailureReasons[reason]++
if orgID != nil {
orgOutcome := affectedOrgs[*orgID]
orgOutcome.Stranded++
affectedOrgs[*orgID] = orgOutcome
}
}
for _, accountID := range accountIDs {
account, aerr := s.EmailRepository.GetByID(ctx, accountID)
if aerr != nil || account == nil || account.OrganizationID == nil {
if aerr != nil || account == nil {
log.Warn().Err(aerr).Str("account_id", accountID.String()).Msg("dead worker reassign: mailbox owner unavailable")
recordFailure("mailbox_lookup_failed", nil)
continue
}
if account.OrganizationID == nil {
log.Warn().Str("account_id", accountID.String()).Msg("dead worker reassign: mailbox has no organization")
recordFailure("organization_missing", nil)
continue
}
@@ -162,33 +226,52 @@ func (s *JobsService) detectDeadWorkers(ctx context.Context) {
})
if rerr != nil {
log.Warn().Err(rerr).Str("account_id", accountID.String()).Msg("dead worker reassign: placement failed")
recordFailure("placement_failed", account.OrganizationID)
continue
}
if result != nil {
target = result.Worker
}
} else {
target, _ = s.findHealthyWorker(ctx, w)
var ferr error
target, ferr = s.findHealthyWorker(ctx, w)
if ferr != nil {
log.Warn().Err(ferr).Str("account_id", accountID.String()).Msg("dead worker reassign: fallback lookup failed")
recordFailure("replacement_lookup_failed", account.OrganizationID)
continue
}
}
if target == nil {
recordFailure("no_eligible_worker", account.OrganizationID)
continue
}
if s.AssignmentService != nil {
if err := s.AssignmentService.MoveMailbox(ctx, accountID, &w.ID, target.ID); err != nil {
log.Error().Err(err).Str("account_id", accountID.String()).Msg("failed to reassign email account")
recordFailure("move_failed", account.OrganizationID)
continue
}
} else if err := s.WorkerRepo.UpdateEmailAccountWorker(ctx, accountID, target.ID); err != nil {
log.Error().Err(err).Str("account_id", accountID.String()).Msg("failed to reassign email account")
recordFailure("move_failed", account.OrganizationID)
continue
} else {
_ = s.WorkerRepo.DecrementAccountCount(ctx, w.ID)
_ = s.WorkerRepo.IncrementAccountCount(ctx, target.ID)
weight := workerapp.MailboxWeight(account.Provider, account.Warmup != nil)
if err := s.WorkerRepo.AddLoadScore(ctx, w.ID, -weight); err != nil {
log.Warn().Err(err).Str("worker_id", w.ID.String()).Msg("dead worker reassign: source load update failed")
}
if err := s.WorkerRepo.AddLoadScore(ctx, target.ID, weight); err != nil {
log.Warn().Err(err).Str("worker_id", target.ID.String()).Msg("dead worker reassign: destination load update failed")
}
}
reassigned++
affectedOrgs[*account.OrganizationID]++
outcome.Reassigned++
orgOutcome := affectedOrgs[*account.OrganizationID]
orgOutcome.Reassigned++
affectedOrgs[*account.OrganizationID] = orgOutcome
destinations[target.ID] = struct{}{}
// The backend's worker reconciler loads the account onto its new
@@ -197,43 +280,43 @@ func (s *JobsService) detectDeadWorkers(ctx context.Context) {
// the worker reject it and log an error per mailbox.
}
if reassigned > 0 {
if outcome.Reassigned > 0 {
log.Info().
Str("dead_worker", w.ID.String()).
Int("destinations", len(destinations)).
Int("reassigned", reassigned).
Int("reassigned", outcome.Reassigned).
Int("stranded", outcome.Stranded).
Msg("email accounts reassigned from dead worker")
// Record in admin_audit_log so the dashboard's audit viewer
// shows when and where the fleet auto-reassigned. uuid.Nil for
// admin_user_id signals "system action".
if s.AdminRepo != nil {
_ = s.AdminRepo.CreateAuditLog(ctx, &models.AdminAuditLog{
ID: uuid.New(),
AdminUserID: uuid.Nil,
Action: "auto_reassign",
TargetType: "worker",
TargetID: w.ID,
Details: map[string]any{
"replacement_workers": len(destinations),
"accounts_reassigned": reassigned,
"reason": "heartbeat_expired",
},
IPAddress: "",
UserAgent: "system",
CreatedAt: time.Now(),
})
}
s.notifyWorkerDown(ctx, w.ID, affectedOrgs, true)
s.notifyOperatorWorkerDown(ctx, w.ID, reassigned, true)
} else {
log.Warn().Str("worker_id", w.ID.String()).Msg("no healthy replacement worker found")
s.notifyWorkerDown(ctx, w.ID, s.accountOrgs(ctx, accountIDs), false)
s.notifyOperatorWorkerDown(ctx, w.ID, len(accountIDs), false)
log.Warn().Str("worker_id", w.ID.String()).Int("stranded", outcome.Stranded).Msg("dead worker recovery failed")
}
if reassigned == len(accountIDs) {
// uuid.Nil identifies this as a system action in the operator audit.
if s.AdminRepo != nil {
_ = s.AdminRepo.CreateAuditLog(ctx, &models.AdminAuditLog{
ID: uuid.New(),
AdminUserID: uuid.Nil,
Action: "auto_reassign",
TargetType: "worker",
TargetID: w.ID,
Details: map[string]any{
"outcome": outcome.state(),
"replacement_workers": len(destinations),
"accounts_reassigned": outcome.Reassigned,
"accounts_stranded": outcome.Stranded,
"failure_reasons": outcome.FailureReasons,
"reason": "heartbeat_expired",
},
IPAddress: "",
UserAgent: "system",
CreatedAt: time.Now(),
})
}
s.notifyWorkerDown(ctx, w.ID, affectedOrgs)
s.notifyOperatorWorkerDown(ctx, w.ID, outcome)
if outcome.Reassigned == len(accountIDs) {
s.deactivateIfLongDead(ctx, w)
}
}
@@ -268,17 +351,6 @@ func (s *JobsService) deactivateIfLongDead(ctx context.Context, w models.Worker)
log.Info().Str("worker_id", w.ID.String()).Msg("dead worker deactivated")
}
// accountOrgs resolves which orgs own the given accounts (org -> count).
func (s *JobsService) accountOrgs(ctx context.Context, accountIDs []uuid.UUID) map[uuid.UUID]int {
orgs := map[uuid.UUID]int{}
for _, id := range accountIDs {
if account, err := s.EmailRepository.GetByID(ctx, id); err == nil && account != nil && account.OrganizationID != nil {
orgs[*account.OrganizationID]++
}
}
return orgs
}
func (s *JobsService) findHealthyWorker(ctx context.Context, deadWorker models.Worker) (*models.Worker, error) {
workers, err := s.WorkerRepo.ListPlaceableWorkers(ctx)
if err != nil {
+34
View File
@@ -0,0 +1,34 @@
package jobs
import "testing"
func TestWorkerRecoveryOutcomeState(t *testing.T) {
tests := []struct {
name string
outcome workerRecoveryOutcome
want string
}{
{name: "complete", outcome: workerRecoveryOutcome{Reassigned: 3}, want: "complete"},
{name: "partial", outcome: workerRecoveryOutcome{Reassigned: 2, Stranded: 1}, want: "partial"},
{name: "failed", outcome: workerRecoveryOutcome{Stranded: 3}, want: "failed"},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := tt.outcome.state(); got != tt.want {
t.Fatalf("state() = %q, want %q", got, tt.want)
}
})
}
}
func TestWorkerRecoveryOutcomeFailureSummary(t *testing.T) {
outcome := workerRecoveryOutcome{FailureReasons: map[string]int{
"no_eligible_worker": 2,
"move_failed": 1,
}}
if got, want := outcome.failureSummary(), "move failed: 1, no eligible worker: 2"; got != want {
t.Fatalf("failureSummary() = %q, want %q", got, want)
}
}
@@ -0,0 +1,7 @@
BEGIN;
ALTER TABLE fleet_nodes
DROP CONSTRAINT IF EXISTS fleet_nodes_capacity_target_positive,
DROP COLUMN IF EXISTS capacity_target;
COMMIT;
@@ -0,0 +1,8 @@
BEGIN;
ALTER TABLE fleet_nodes
ADD COLUMN capacity_target numeric(10,2) NOT NULL DEFAULT 100,
ADD CONSTRAINT fleet_nodes_capacity_target_positive
CHECK (capacity_target > 0) NOT VALID;
COMMIT;
@@ -50,8 +50,4 @@ CREATE UNIQUE INDEX worker_capacity_view_pk ON public.worker_capacity_view USING
REFRESH MATERIALIZED VIEW public.worker_capacity_view;
ALTER TABLE fleet_nodes
DROP CONSTRAINT IF EXISTS fleet_nodes_capacity_target_positive,
DROP COLUMN IF EXISTS capacity_target;
COMMIT;
@@ -1,14 +1,9 @@
-- Give each worker an assigned-mailbox planning target. Provider send limits
-- remain attached to each mailbox, tenant, or API project.
-- Switch worker load to assigned-mailbox counts and rebuild its capacity view.
BEGIN;
DROP MATERIALIZED VIEW IF EXISTS worker_capacity_view;
ALTER TABLE fleet_nodes
ADD COLUMN capacity_target numeric(10,2) NOT NULL DEFAULT 100,
ADD CONSTRAINT fleet_nodes_capacity_target_positive CHECK (capacity_target > 0);
-- Load is now the number of assigned mailboxes, independent of provider.
UPDATE workers w
SET load_score = (
@@ -0,0 +1,8 @@
BEGIN;
ALTER TABLE fleet_nodes
DROP CONSTRAINT IF EXISTS fleet_nodes_capacity_target_positive,
ADD CONSTRAINT fleet_nodes_capacity_target_positive
CHECK (capacity_target > 0) NOT VALID;
COMMIT;
@@ -0,0 +1,2 @@
ALTER TABLE fleet_nodes
VALIDATE CONSTRAINT fleet_nodes_capacity_target_positive;
+1
View File
@@ -193,6 +193,7 @@ type AdminWorkerEmail struct {
RiskEvaluatedAt *time.Time `json:"risk_evaluated_at,omitempty"`
WarmupHealth string `json:"warmup_health,omitempty"` // worst warmup health_state, "" if not in a pool
BlockedUntil *time.Time `json:"blocked_until,omitempty"`
CreatedAt time.Time `json:"created_at"`
}
// ReassignEmailsRequest represents the request to reassign emails
+8 -12
View File
@@ -37,7 +37,7 @@ type AdminRepository interface {
ListWorkers(ctx context.Context, cursor *uuid.UUID, limit int) (*models.AdminWorkersResult, error)
GetWorkerDetail(ctx context.Context, workerID uuid.UUID) (*models.AdminWorkerDetail, error)
UpdateWorker(ctx context.Context, workerID uuid.UUID, update *models.AdminUpdateWorker) error
GetWorkerEmails(ctx context.Context, workerID uuid.UUID, cursor *uuid.UUID, limit int) ([]models.AdminWorkerEmail, *models.Pagination, error)
GetWorkerEmails(ctx context.Context, workerID uuid.UUID, beforeAt time.Time, beforeID uuid.UUID, limit int) ([]models.AdminWorkerEmail, *models.Pagination, error)
GetWorkerStats(ctx context.Context, workerID uuid.UUID) (*models.WorkerStats, error)
ReassignEmails(ctx context.Context, emailIDs []uuid.UUID, newWorkerID uuid.UUID) error
@@ -858,20 +858,16 @@ func (r *adminRepository) UpdateWorker(ctx context.Context, workerID uuid.UUID,
}
// GetWorkerEmails gets emails connected to a worker
func (r *adminRepository) GetWorkerEmails(ctx context.Context, workerID uuid.UUID, cursor *uuid.UUID, limit int) ([]models.AdminWorkerEmail, *models.Pagination, error) {
func (r *adminRepository) GetWorkerEmails(ctx context.Context, workerID uuid.UUID, beforeAt time.Time, beforeID uuid.UUID, limit int) ([]models.AdminWorkerEmail, *models.Pagination, error) {
if limit <= 0 || limit > 100 {
limit = 50
}
args := []interface{}{workerID, limit + 1}
whereClause := "WHERE ea.worker_id = $1"
if cursor != nil {
whereClause += ` AND (ea.created_at, ea.id) < (
SELECT cursor_ea.created_at, cursor_ea.id
FROM email_accounts cursor_ea
WHERE cursor_ea.id = $3
)`
args = append(args, *cursor)
if beforeID != uuid.Nil {
whereClause += ` AND (ea.created_at, ea.id) < ($3, $4)`
args = append(args, beforeAt, beforeID)
}
// Health lives on warmup_pool_participants, one row per mailbox. The CASE
@@ -883,7 +879,7 @@ func (r *adminRepository) GetWorkerEmails(ctx context.Context, workerID uuid.UUI
COALESCE(ea.risk_band, 'clean'::email_risk_band)::text,
ea.risk_evaluated_at,
COALESCE(wh.health_state, '')::text,
wh.blocked_until
wh.blocked_until, ea.created_at
FROM email_accounts ea
LEFT JOIN LATERAL (
SELECT health_state, blocked_until
@@ -916,7 +912,7 @@ func (r *adminRepository) GetWorkerEmails(ctx context.Context, workerID uuid.UUI
err := rows.Scan(
&e.ID, &e.Email, &e.UserID, &e.OrganizationID,
&e.Status, &e.Provider, &e.WarmupEnabled, &e.LastSyncedAt,
&e.RiskBand, &e.RiskEvaluatedAt, &e.WarmupHealth, &e.BlockedUntil,
&e.RiskBand, &e.RiskEvaluatedAt, &e.WarmupHealth, &e.BlockedUntil, &e.CreatedAt,
)
if err != nil {
return nil, nil, err
@@ -930,7 +926,7 @@ func (r *adminRepository) GetWorkerEmails(ctx context.Context, workerID uuid.UUI
if len(emails) > limit {
emails = emails[:limit]
pagination.NextCursor = paging.UUIDString(emails[limit-1].ID)
pagination.NextCursor = paging.EncodeTime(emails[limit-1].CreatedAt, emails[limit-1].ID)
}
return emails, pagination, nil
+16 -13
View File
@@ -171,10 +171,21 @@ func (r *workerRepository) CountOrgMailboxes(ctx context.Context, orgID uuid.UUI
}
const mailboxPlacementStateSelect = `
WITH assigned_mailboxes AS MATERIALIZED (
SELECT worker_id, organization_id, provider
FROM email_accounts
WHERE worker_id IS NOT NULL
WITH shared_workers AS MATERIALIZED (
SELECT live_worker.id
FROM workers live_worker
JOIN fleet_nodes live_node ON live_node.id = live_worker.id
LEFT JOIN dedicated_worker_assignments live_reservation
ON live_reservation.worker_id = live_worker.id AND live_reservation.released_at IS NULL
WHERE live_worker.health_state IN ('healthy', 'watch')
AND live_node.active
AND live_node.last_seen_at > now() - $1::interval
AND live_reservation.worker_id IS NULL
),
assigned_mailboxes AS MATERIALIZED (
SELECT ea.worker_id, ea.organization_id, ea.provider
FROM email_accounts ea
JOIN shared_workers shared ON shared.id = ea.worker_id
),
worker_totals AS (
SELECT worker_id, count(*) AS total_count
@@ -202,15 +213,7 @@ const mailboxPlacementStateSelect = `
GROUP BY provider
),
fleet AS (
SELECT count(*) AS live_count
FROM workers live_worker
JOIN fleet_nodes live_node ON live_node.id = live_worker.id
LEFT JOIN dedicated_worker_assignments live_reservation
ON live_reservation.worker_id = live_worker.id AND live_reservation.released_at IS NULL
WHERE live_worker.health_state IN ('healthy', 'watch')
AND live_node.active
AND live_node.last_seen_at > now() - $1::interval
AND live_reservation.worker_id IS NULL
SELECT count(*) AS live_count FROM shared_workers
)
SELECT ea.id, ea.organization_id, ea.provider::text, (ea.warmup IS NOT NULL),
ea.worker_id, ea.worker_assigned_at,
@@ -146,3 +146,79 @@ func TestAdminWorkerQueriesLive(t *testing.T) {
}
_, _ = pool.Exec(ctx, `DELETE FROM fleet_nodes WHERE id = $1`, id)
}
func TestMailboxPlacementStateExcludesDedicatedWorkerMailboxesLive(t *testing.T) {
pool := placementLiveDB(t)
ctx := context.Background()
tx, err := pool.Begin(ctx)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = tx.Rollback(ctx) })
userID := uuid.New()
orgID := uuid.New()
subscriptionID := uuid.New()
sharedWorkerID := uuid.New()
secondSharedWorkerID := uuid.New()
dedicatedWorkerID := uuid.New()
sharedMailboxID := uuid.New()
dedicatedMailboxID := uuid.New()
if _, err := tx.Exec(ctx, `
INSERT INTO users (id, first_name, last_name, email)
VALUES ($1, 'Placement', 'Test', $2)`, userID, userID.String()+"@example.test"); err != nil {
t.Fatal(err)
}
if _, err := tx.Exec(ctx, `
INSERT INTO organizations (id, name, owner_user_id)
VALUES ($1, 'Placement Test', $2)`, orgID, userID); err != nil {
t.Fatal(err)
}
if _, err := tx.Exec(ctx, `
INSERT INTO subscriptions (id, user_id, organization_id, plan_id, stripe_customer_id, status)
VALUES ($1, $2, $3, '00000000-0000-0000-0000-000000000001', $4, 'active')`,
subscriptionID, userID, orgID, "test_"+subscriptionID.String()); err != nil {
t.Fatal(err)
}
if _, err := tx.Exec(ctx, `
INSERT INTO fleet_nodes (id, role, active, last_seen_at)
VALUES ($1, 'worker', true, now()), ($2, 'worker', true, now()), ($3, 'worker', true, now())`,
sharedWorkerID, secondSharedWorkerID, dedicatedWorkerID); err != nil {
t.Fatal(err)
}
if _, err := tx.Exec(ctx, `
INSERT INTO workers (id) VALUES ($1), ($2), ($3)`,
sharedWorkerID, secondSharedWorkerID, dedicatedWorkerID); err != nil {
t.Fatal(err)
}
if _, err := tx.Exec(ctx, `
INSERT INTO dedicated_worker_assignments (worker_id, organization_id, subscription_id)
VALUES ($1, $2, $3)`, dedicatedWorkerID, orgID, subscriptionID); err != nil {
t.Fatal(err)
}
if _, err := tx.Exec(ctx, `
INSERT INTO email_accounts (
id, user_id, organization_id, worker_id, email, name,
signature_plain, signature_html, provider, warmup_tag, worker_assigned_at
) VALUES
($1, $2, $3, $4, $5, 'Shared', '', '', 'gmail', '', now()),
($6, $2, $3, $7, $8, 'Dedicated', '', '', 'gmail', '', now())`,
sharedMailboxID, userID, orgID, sharedWorkerID, sharedMailboxID.String()+"@example.test",
dedicatedMailboxID, dedicatedWorkerID, dedicatedMailboxID.String()+"@example.test"); err != nil {
t.Fatal(err)
}
state, err := scanMailboxPlacementState(tx.QueryRow(ctx,
mailboxPlacementStateSelect+` WHERE ea.id = $2`, WorkerLivenessWindow, sharedMailboxID))
if err != nil {
t.Fatal(err)
}
if state.LiveWorkerCount != 2 {
t.Fatalf("LiveWorkerCount = %d, want 2 shared workers", state.LiveWorkerCount)
}
if state.OrgTotalMailboxes != 1 || state.ProviderTotalMailboxes != 1 {
t.Fatalf("concentration totals = org %d, provider %d; dedicated mailbox must be excluded",
state.OrgTotalMailboxes, state.ProviderTotalMailboxes)
}
}