From d6025ea4c0e6e26398f33ff07dfcb4d47a5e590d Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Thu, 17 Sep 2026 06:24:14 -0700 Subject: [PATCH] feat: address worker capacity review findings with safe migrations and recovery reporting --- admin/src/lib/api/models/admin.ts | 1 + cmd/worker/capacity_test.go | 9 + cmd/worker/main.go | 3 +- .../content/docs/development/architecture.mdx | 6 +- .../development/operator-notifications.mdx | 4 +- internal/api/handler/admin.go | 9 +- internal/api/handler/fleet_nodes.go | 18 +- internal/api/handler/fleet_nodes_test.go | 1 + internal/app/admin/service.go | 6 +- internal/app/consumer/dead_worker.go | 204 ++++++++++++------ internal/app/consumer/dead_worker_test.go | 34 +++ .../000169_worker_capacity_target.down.sql | 7 + .../000169_worker_capacity_target.up.sql | 8 + ...0170_worker_operational_capacity.down.sql} | 4 - ...000170_worker_operational_capacity.up.sql} | 7 +- ...1_validate_worker_capacity_target.down.sql | 8 + ...171_validate_worker_capacity_target.up.sql | 2 + internal/models/admin.go | 1 + internal/repository/pg_admin.go | 20 +- internal/repository/pg_worker_placement.go | 29 +-- .../repository/placement_queries_live_test.go | 76 +++++++ 21 files changed, 337 insertions(+), 120 deletions(-) create mode 100644 internal/app/consumer/dead_worker_test.go create mode 100644 internal/infrastructure/db/migrations/000169_worker_capacity_target.down.sql create mode 100644 internal/infrastructure/db/migrations/000169_worker_capacity_target.up.sql rename internal/infrastructure/db/migrations/{000169_worker_operational_capacity.down.sql => 000170_worker_operational_capacity.down.sql} (94%) rename internal/infrastructure/db/migrations/{000169_worker_operational_capacity.up.sql => 000170_worker_operational_capacity.up.sql} (88%) create mode 100644 internal/infrastructure/db/migrations/000171_validate_worker_capacity_target.down.sql create mode 100644 internal/infrastructure/db/migrations/000171_validate_worker_capacity_target.up.sql diff --git a/admin/src/lib/api/models/admin.ts b/admin/src/lib/api/models/admin.ts index 2192551ba..cb774a887 100644 --- a/admin/src/lib/api/models/admin.ts +++ b/admin/src/lib/api/models/admin.ts @@ -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 { diff --git a/cmd/worker/capacity_test.go b/cmd/worker/capacity_test.go index 800af8b65..ec402ac73 100644 --- a/cmd/worker/capacity_test.go +++ b/cmd/worker/capacity_test.go @@ -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) + } + }) + } } diff --git a/cmd/worker/main.go b/cmd/worker/main.go index a203b29ff..af453ee3a 100644 --- a/cmd/worker/main.go +++ b/cmd/worker/main.go @@ -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 diff --git a/docs/content/docs/development/architecture.mdx b/docs/content/docs/development/architecture.mdx index 1bdbbe181..1439f5458 100644 --- a/docs/content/docs/development/architecture.mdx +++ b/docs/content/docs/development/architecture.mdx @@ -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. diff --git a/docs/content/docs/development/operator-notifications.mdx b/docs/content/docs/development/operator-notifications.mdx index 741bd77f1..d6245cb29 100644 --- a/docs/content/docs/development/operator-notifications.mdx +++ b/docs/content/docs/development/operator-notifications.mdx @@ -43,7 +43,7 @@ X-Warmbly-Event: { "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 | diff --git a/internal/api/handler/admin.go b/internal/api/handler/admin.go index a122bd055..8c84f4918 100644 --- a/internal/api/handler/admin.go +++ b/internal/api/handler/admin.go @@ -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 diff --git a/internal/api/handler/fleet_nodes.go b/internal/api/handler/fleet_nodes.go index c58255363..92e908e46 100644 --- a/internal/api/handler/fleet_nodes.go +++ b/internal/api/handler/fleet_nodes.go @@ -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{ diff --git a/internal/api/handler/fleet_nodes_test.go b/internal/api/handler/fleet_nodes_test.go index 83382343e..12f732560 100644 --- a/internal/api/handler/fleet_nodes_test.go +++ b/internal/api/handler/fleet_nodes_test.go @@ -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"}, } diff --git a/internal/app/admin/service.go b/internal/app/admin/service.go index 4871717c6..fd33c633e 100644 --- a/internal/app/admin/service.go +++ b/internal/app/admin/service.go @@ -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") diff --git a/internal/app/consumer/dead_worker.go b/internal/app/consumer/dead_worker.go index 0972b6aec..7049102b0 100644 --- a/internal/app/consumer/dead_worker.go +++ b/internal/app/consumer/dead_worker.go @@ -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 { diff --git a/internal/app/consumer/dead_worker_test.go b/internal/app/consumer/dead_worker_test.go new file mode 100644 index 000000000..cfa88f55b --- /dev/null +++ b/internal/app/consumer/dead_worker_test.go @@ -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) + } +} diff --git a/internal/infrastructure/db/migrations/000169_worker_capacity_target.down.sql b/internal/infrastructure/db/migrations/000169_worker_capacity_target.down.sql new file mode 100644 index 000000000..97150ba6d --- /dev/null +++ b/internal/infrastructure/db/migrations/000169_worker_capacity_target.down.sql @@ -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; diff --git a/internal/infrastructure/db/migrations/000169_worker_capacity_target.up.sql b/internal/infrastructure/db/migrations/000169_worker_capacity_target.up.sql new file mode 100644 index 000000000..635aa8fe8 --- /dev/null +++ b/internal/infrastructure/db/migrations/000169_worker_capacity_target.up.sql @@ -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; diff --git a/internal/infrastructure/db/migrations/000169_worker_operational_capacity.down.sql b/internal/infrastructure/db/migrations/000170_worker_operational_capacity.down.sql similarity index 94% rename from internal/infrastructure/db/migrations/000169_worker_operational_capacity.down.sql rename to internal/infrastructure/db/migrations/000170_worker_operational_capacity.down.sql index 95bc26b63..8939e2609 100644 --- a/internal/infrastructure/db/migrations/000169_worker_operational_capacity.down.sql +++ b/internal/infrastructure/db/migrations/000170_worker_operational_capacity.down.sql @@ -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; diff --git a/internal/infrastructure/db/migrations/000169_worker_operational_capacity.up.sql b/internal/infrastructure/db/migrations/000170_worker_operational_capacity.up.sql similarity index 88% rename from internal/infrastructure/db/migrations/000169_worker_operational_capacity.up.sql rename to internal/infrastructure/db/migrations/000170_worker_operational_capacity.up.sql index 145ff3728..91b620e4d 100644 --- a/internal/infrastructure/db/migrations/000169_worker_operational_capacity.up.sql +++ b/internal/infrastructure/db/migrations/000170_worker_operational_capacity.up.sql @@ -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 = ( diff --git a/internal/infrastructure/db/migrations/000171_validate_worker_capacity_target.down.sql b/internal/infrastructure/db/migrations/000171_validate_worker_capacity_target.down.sql new file mode 100644 index 000000000..dd21fe742 --- /dev/null +++ b/internal/infrastructure/db/migrations/000171_validate_worker_capacity_target.down.sql @@ -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; diff --git a/internal/infrastructure/db/migrations/000171_validate_worker_capacity_target.up.sql b/internal/infrastructure/db/migrations/000171_validate_worker_capacity_target.up.sql new file mode 100644 index 000000000..ec17a3c96 --- /dev/null +++ b/internal/infrastructure/db/migrations/000171_validate_worker_capacity_target.up.sql @@ -0,0 +1,2 @@ +ALTER TABLE fleet_nodes + VALIDATE CONSTRAINT fleet_nodes_capacity_target_positive; diff --git a/internal/models/admin.go b/internal/models/admin.go index 5bcc06670..90c01bd3c 100644 --- a/internal/models/admin.go +++ b/internal/models/admin.go @@ -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 diff --git a/internal/repository/pg_admin.go b/internal/repository/pg_admin.go index 701b2bef0..b8b677d79 100644 --- a/internal/repository/pg_admin.go +++ b/internal/repository/pg_admin.go @@ -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 diff --git a/internal/repository/pg_worker_placement.go b/internal/repository/pg_worker_placement.go index 30866ad11..386d43eca 100644 --- a/internal/repository/pg_worker_placement.go +++ b/internal/repository/pg_worker_placement.go @@ -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, diff --git a/internal/repository/placement_queries_live_test.go b/internal/repository/placement_queries_live_test.go index 89499be62..d121d484a 100644 --- a/internal/repository/placement_queries_live_test.go +++ b/internal/repository/placement_queries_live_test.go @@ -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) + } +}