feat: refresh self-hosted Cloud pool standing from validated mailbox listings and separate availability holds from reputation bans

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
Matthew Meszaros
2026-10-08 17:19:02 +00:00
co-authored by Devin AI
parent acbd79e3bf
commit b16dadd264
9 changed files with 303 additions and 7 deletions
+10
View File
@@ -1295,6 +1295,16 @@ func main() {
cloudLinkRepository := repository.NewCloudLinkRepository(primaryDB.Pool, credEncrypter)
emailService.WireCloudLink(cloudLinkRepository)
cloudLinkService = cloudlink.NewService(cloudLinkRepository, emailRepostory, emailService)
cloudLinkService.OnStandingChange(func(ctx context.Context, change models.CloudLinkStandingChange) {
state, _, err := warmupRepository.GetHealthState(ctx, change.EmailAccountID)
if err == nil {
err = workerRepository.SetEmailAccountRiskBand(ctx, change.EmailAccountID, models.RiskBandFromHealth(state))
}
if err != nil {
log.Printf("cloud standing refresh: risk band update failed: %v", err)
}
warmupService.PublishHealthTransition(ctx, change.EmailAccountID, change.Previous, change.Current, change.Reason)
})
if reports, ok := analyticsService.(analytics.CloudWarmupReportsAware); ok {
reports.WireCloudWarmupReports(cloudLinkService)
}
+1 -1
View File
@@ -54,7 +54,7 @@ Only the first two badge a nav tab.
**Mailbox configuration**: daily cap above the `50`/day safe band (gently with clean numbers, firmly without); a mailbox under `30` days old sending above `20`/day; a send gap under `5` minutes; unresolved errors on a mailbox campaigns still use; too few mailboxes carrying total volume; an inactive mailbox attached to a running campaign.
**Warmup**: warmup off while sending cold (critical on a new mailbox); warmup left paused; a ceiling too low to balance cold volume (roughly half the cold cap); reply rate under `20%`, making warmup one-way; quarantined or blocked from the pool.
**Warmup**: warmup off while sending cold (critical on a new mailbox); warmup left paused; a ceiling too low to balance cold volume (roughly half the cold cap); reply rate under `20%`, making warmup one-way; quarantined or blocked from the pool. Missing or stale Cloud status is reported separately from a real reputation ban, with steps to refresh it without disconnecting mailboxes.
**Campaigns**: no follow-ups; reply rate under `1%` at volume, escalating under `0.3%`; one step replying below a quarter of the campaign's best (compared only against siblings, which share an audience); one-click unsubscribe off; a daily limit above what its mailboxes can send; no active sender; follow-up spacing outside `2` to `14` days; under `7` days of leads left; no A/B test at high volume; a sending window too narrow for the volume; contacts enrolled in several running campaigns.
+3 -1
View File
@@ -122,7 +122,7 @@ The Warmup plan is the same plan a hosted workspace buys when it only warms and
Enrolled mailboxes are ordinary members of the hosted pool and are held to the same rules: verification tokens on every warmup message, complaint and bounce tracking with automatic quarantine or blocking when a mailbox starts hurting partners, and spam placement that slows a mailbox down. A mailbox that gets quarantined on the cloud stops being offered to partners and re-enters on the same probation terms as any hosted mailbox.
Your instance holds the mailbox to the cloud's verdict as well, exactly as it would one from its own pool. It reads every enrolled mailbox's standing from the cloud every five minutes, and on enrolling, pausing or resuming, and from then on:
Your instance holds the mailbox to the cloud's verdict as well, exactly as it would one from its own pool. It reads every enrolled mailbox's standing from the cloud every five minutes, on enrolling, pausing or resuming, and when loading the mailbox list in **Settings > Warmbly Cloud**. A recognized status from that list refreshes the local observation instead of being replaced by an older cached status. From then on:
| Cloud standing | On your instance |
|----------------|------------------|
@@ -134,6 +134,8 @@ The same standing feeds the mailbox's health on the **Mailboxes** page, the **Wa
Cloud standing is observed evidence, not a permanent permission. Healthy, watch and throttled observations have a 15-minute freshness limit. Missing, unrecognized or stale observations stop new control-plane admission with `cloud_evidence_unavailable` until a recognized cloud observation arrives. This availability hold does not claim the provider blocked the mailbox. Existing enrollments start with unknown freshness after upgrading, then the next successful standing sync establishes it.
The Advisor reports an unavailable Cloud status separately from a real pool quarantine or block. A mailbox may still be warming on Cloud while its self-hosted instance lacks a current status. Open **Settings > Warmbly Cloud** and refresh the mailbox list. If the status remains unavailable, check Cloud connectivity and the consumer's `cloud_standing_sync` job, and keep the backend and consumer on the same release. Do not disconnect or re-enroll the mailbox to refresh its status. The **Warming** count indicates warmup is enabled, not that a provider has confirmed a send; the warmup count and send warnings show progress and failures.
An active quarantine or block remains in force even when the cloud cannot be reached. A bounded hold still ends on its own date; an unbounded block needs review. Removing a mailbox from the cloud does not lift an active hold. Restrictive cloud evidence is retained by workspace and address in the local reputation ledger, so leaving or changing pools, reconnecting, or [exporting and importing the workspace](/guides/workspace-export-import/) does not reset it. The cloud link itself is instance-local and does not travel with the archive. Fresh positive observations from a new link do not erase an unexpired restriction from the previous link. Appeals go through the mailbox on Warmbly Cloud, where the verdict was reached.
Local and cloud warmup use the same active-mailbox, participant-role, workspace-standing and health checks. A recipient-only participant cannot originate diagnostic messages. Restricted or suspended workspaces cannot contribute candidates, and paying for Premium does not waive these rules. Spam placement slows participation; it is not by itself a reason to quarantine a mailbox.
+23
View File
@@ -765,6 +765,29 @@ func TestWarmupPoolFindingExplainsItselfWithTheBandsReason(t *testing.T) {
}
}
func TestUnavailableCloudStandingDoesNotInventAReputationBan(t *testing.T) {
m := healthyMailbox()
m.PoolHealth, m.PoolHealthReason = "blocked", "cloud_evidence_unavailable"
byKey := findingsByKey(Detect(snapshotOf(m), defaults()))
if _, ok := byKey["warmup_pool_blocked"]; ok {
t.Fatal("unavailable Cloud evidence was described as lost pool standing")
}
f, ok := byKey["cloud_standing_unavailable"]
if !ok || !strings.Contains(f.Remedy, "cloud_standing_sync") || len(f.Steps) == 0 || f.Action != nil {
t.Fatalf("availability finding does not offer a safe status refresh: %+v", f)
}
settings := defaults()
settings.MutedDetectors = []string{"cloud_standing_unavailable"}
if _, ok := findingsByKey(Detect(snapshotOf(m), settings))["cloud_standing_unavailable"]; ok {
t.Fatal("Cloud availability detector ignored its own mute setting")
}
m.PoolHealthReason = "complaints"
byKey = findingsByKey(Detect(snapshotOf(m), defaults()))
if _, ok := byKey["warmup_pool_blocked"]; !ok {
t.Fatal("a real Cloud block stopped being reported")
}
}
// judgedSnapshot is one active campaign with one email step, plus the verdict
// the copy judge would have attached to it.
func judgedSnapshot(v *copyjudge.Verdict) (*repository.AdvisorSnapshot, uuid.UUID) {
+28
View File
@@ -36,6 +36,12 @@ func warmupDetectors() []Detector {
About: "The share of warmup mail that gets replied to. Warmup works by looking like real correspondence; a low reply rate makes it one-way traffic, which is exactly the pattern it is supposed to counteract.",
Run: detectWarmupReplyRateLow,
},
{
Key: "cloud_standing_unavailable",
Category: models.AdvisorCategoryWarmup,
About: "A Cloud-linked mailbox with missing or stale status evidence. New sends are held until a recognized Cloud observation arrives; this is an availability hold, not a reputation ban or a provider authentication failure.",
Run: detectCloudStandingUnavailable,
},
{
Key: "warmup_pool_blocked",
Category: models.AdvisorCategoryWarmup,
@@ -214,9 +220,31 @@ func detectWarmupReplyRateLow(s *repository.AdvisorSnapshot) []Finding {
return out
}
func detectCloudStandingUnavailable(s *repository.AdvisorSnapshot) []Finding {
out := []Finding{}
for _, m := range s.Mailboxes {
if m.PoolHealthReason == "cloud_evidence_unavailable" {
out = append(out, Finding{
Key: "cloud_standing_unavailable", GroupTitle: "{count} mailboxes need a fresh Cloud pool status",
Category: models.AdvisorCategoryWarmup, Severity: models.AdvisorHigh, Surface: models.AdvisorSurfaceMailboxes,
EntityType: "email_account", EntityID: ref(m.ID), EntityLabel: m.Email, Impact: 70,
Title: fmt.Sprintf("Cloud pool status for %s is unavailable", m.Email),
Detail: "This instance has no current, recognized Cloud pool status for this mailbox. New sending is held until the status is refreshed. This is not evidence of spam, rejected credentials, or a Cloud pool ban; Cloud warmup may still be running independently.",
Remedy: "Open Settings > Warmbly Cloud to refresh the mailbox status. If Cloud cannot be reached, check the connection and the consumer's cloud_standing_sync job. Keep the backend and consumer on the same release. Do not disconnect or re-enroll the mailbox.",
Steps: []string{"Open Settings > Warmbly Cloud and reload the mailbox list to obtain a current status.", "If the status stays unavailable, check Cloud connectivity and the consumer's cloud_standing_sync job. Do not disconnect or re-enroll the mailbox."},
Evidence: map[string]any{"mailbox": m.Email, "pool_state": m.PoolHealth, "pool_health_reason": m.PoolHealthReason},
})
}
}
return out
}
func detectWarmupPoolBlocked(s *repository.AdvisorSnapshot) []Finding {
out := []Finding{}
for _, m := range s.Mailboxes {
if m.PoolHealthReason == "cloud_evidence_unavailable" {
continue
}
bad := m.PoolBlocked ||
m.PoolHealth == "quarantined" || m.PoolHealth == "blocked" || m.PoolHealth == "throttled"
if !bad {
@@ -0,0 +1,176 @@
package cloudlink
import (
"context"
"encoding/json"
"errors"
"net/http"
"net/http/httptest"
"testing"
"time"
"github.com/google/uuid"
"github.com/warmbly/warmbly/internal/errx"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/repository"
)
type listedStandingRepo struct {
*workspaceLinkRepo
locked bool
writeErr error
writes int
}
type listedStandingEmails struct {
repository.EmailRepository
accounts []models.Email
}
func (e listedStandingEmails) GetAllActiveInScope(context.Context, repository.AccountScope) ([]models.Email, *errx.Error) {
return e.accounts, nil
}
func TestMailboxListingCannotRefreshStandingFromAnotherLink(t *testing.T) {
org, first, second := uuid.New(), uuid.New(), uuid.New()
firstLink, secondLink := uuid.New(), uuid.New()
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) {
h := &models.WarmupHealthInfo{State: "healthy"}
if req.Header.Get("Authorization") == "Bearer other" {
h = &models.WarmupHealthInfo{State: "blocked", Reason: "unrelated_mailbox"}
}
_ = json.NewEncoder(w).Encode([]models.PoolLinkMailboxState{{RemoteID: first, Health: h}})
}))
defer srv.Close()
r := &listedStandingRepo{workspaceLinkRepo: &workspaceLinkRepo{
links: map[uuid.UUID]*models.CloudLink{
firstLink: {InstanceID: firstLink, CloudURL: srv.URL, Token: "owning"},
secondLink: {InstanceID: secondLink, CloudURL: srv.URL, Token: "other"},
},
mailboxes: map[uuid.UUID]*models.CloudLinkMailbox{
first: {EmailAccountID: first, RemoteID: first, InstanceID: firstLink, EnrollmentState: "active"},
second: {EmailAccountID: second, RemoteID: second, InstanceID: secondLink, EnrollmentState: "active"},
},
}}
s := &service{repo: r, emails: listedStandingEmails{accounts: []models.Email{{ID: first, OrganizationID: &org}, {ID: second, OrganizationID: &org}}}}
rows, xerr := s.ListMailboxes(context.Background(), org)
if xerr != nil || len(rows) != 2 || rows[0].Cloud.Health.State != "healthy" || rows[1].Cloud != nil || r.writes != 1 {
t.Fatalf("standing crossed link ownership: %+v, %v", rows, xerr)
}
}
func (r *listedStandingRepo) WithReconciliationLock(_ context.Context, fn func() error) error {
r.locked = true
defer func() { r.locked = false }()
return fn()
}
func (r *listedStandingRepo) SetStanding(_ context.Context, id uuid.UUID, h *models.WarmupHealthInfo, _ bool) (models.WarmupHealthState, error) {
if !r.locked {
return "", errors.New("observation not serialized with enrollment")
}
if r.writeErr != nil {
return "", r.writeErr
}
m := r.mailboxes[id]
var prev models.WarmupHealthState
if m.Standing != nil {
prev = models.WarmupHealthState(m.Standing.State)
}
now := time.Now()
m.Standing, m.StandingObservedAt = h, &now
r.writes++
return prev, nil
}
func (r *listedStandingRepo) InvalidateStanding(_ context.Context, id uuid.UUID) error {
r.mailboxes[id].StandingObservedAt = nil
return nil
}
func TestMailboxListRefreshesOnlyValidatedStandingOnItsOwningLink(t *testing.T) {
for _, scenario := range []string{"stale", "upgrade", "block", "unknown", "absent", "unreachable", "write_failure", "wrong_remote", "disconnect", "pending_remove"} {
t.Run(scenario, func(t *testing.T) {
org, account, instance := uuid.New(), uuid.New(), uuid.New()
old := time.Now().Add(-time.Hour)
m := &models.CloudLinkMailbox{EmailAccountID: account, RemoteID: account, InstanceID: instance, EnrollmentState: "active", Standing: &models.WarmupHealthInfo{State: "healthy"}, StandingObservedAt: &old}
if scenario == "upgrade" {
m.StandingObservedAt = nil
}
if scenario == "absent" {
m.Standing = &models.WarmupHealthInfo{State: "blocked", Reason: "complaints"}
}
if scenario == "pending_remove" {
m.EnrollmentState = "pending_remove"
}
state := &models.WarmupHealthInfo{State: "healthy"}
if scenario == "block" {
state = &models.WarmupHealthInfo{State: "quarantined", Reason: "complaints"}
}
if scenario == "unknown" {
state.State = "future_state"
}
if scenario == "absent" {
state = nil
}
remote := account
if scenario == "wrong_remote" {
remote = uuid.New()
}
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) {
if req.Header.Get("Authorization") != "Bearer owning-link" {
t.Error("wrong source link")
}
if scenario == "unreachable" {
w.WriteHeader(http.StatusServiceUnavailable)
return
}
_ = json.NewEncoder(w).Encode([]models.PoolLinkMailboxState{{RemoteID: remote, Health: state}})
}))
t.Cleanup(srv.Close)
r := &listedStandingRepo{workspaceLinkRepo: &workspaceLinkRepo{links: map[uuid.UUID]*models.CloudLink{instance: {InstanceID: instance, CloudURL: srv.URL, Token: "owning-link", DisconnectPending: scenario == "disconnect"}}, mailboxes: map[uuid.UUID]*models.CloudLinkMailbox{account: m}}}
if scenario == "write_failure" {
r.writeErr = errors.New("database unavailable")
}
s := &service{repo: r, emails: enrollmentEmails{stubEmails{account: &models.Email{ID: account, OrganizationID: &org, Status: "active"}}}}
var changes []models.CloudLinkStandingChange
s.OnStandingChange(func(_ context.Context, c models.CloudLinkStandingChange) { changes = append(changes, c) })
rows, xerr := s.ListMailboxes(context.Background(), org)
if scenario == "write_failure" {
if xerr == nil || len(changes) != 0 || m.StandingObservedAt != &old {
t.Fatalf("failed persistence was accepted: %+v, %v", rows, xerr)
}
return
}
if xerr != nil || len(rows) != 1 {
t.Fatalf("list: %+v, %v", rows, xerr)
}
row := rows[0]
switch scenario {
case "stale", "upgrade":
if row.Cloud == nil || row.Cloud.Health.State != "healthy" || row.StandingObservedAt == nil || !row.StandingObservedAt.After(old) || r.writes != 1 {
t.Fatalf("fresh recognized status lost: %+v", row)
}
case "block":
if row.Cloud.Health.State != "quarantined" || len(changes) != 1 || changes[0].Current != models.WarmupHealthQuarantined {
t.Fatalf("real hold or transition lost: %+v, %+v", row, changes)
}
if _, xerr := s.ListMailboxes(context.Background(), org); xerr != nil || len(changes) != 1 {
t.Fatal("unchanged observation repeated transition")
}
case "unknown", "pending_remove":
if row.Cloud.Health.State != "blocked" || row.Cloud.Health.Reason != "cloud_evidence_unavailable" {
t.Fatalf("invalid admission became healthy: %+v", row.Cloud.Health)
}
case "absent":
if row.Cloud.Health.State != "blocked" || row.Cloud.Health.Reason != "complaints" || r.writes != 0 {
t.Fatal("absent status lifted actual block")
}
case "unreachable", "wrong_remote", "disconnect":
if row.Cloud != nil || r.writes != 0 {
t.Fatalf("unvalidated observation admitted: %+v", row)
}
}
})
}
}
+41 -5
View File
@@ -148,6 +148,7 @@ type Service interface {
// enrolled mailbox, so this instance's send gates hold the same verdict,
// and returns the transitions it saw.
SyncStanding(ctx context.Context) ([]models.CloudLinkStandingChange, *errx.Error)
OnStandingChange(func(context.Context, models.CloudLinkStandingChange))
// IsCloudWarmupThreadReply asks by ancestry: whether what a tokenless
// message answers is a turn of one of the cloud's warmup conversations.
IsCloudWarmupThreadReply(ctx context.Context, accountID uuid.UUID, messageID string, inReplyTo []string) (bool, error)
@@ -174,7 +175,12 @@ type service struct {
// offerFetch lets one caller ask Cloud for the offer while the others wait for its answer.
offerFetch sync.Mutex
disconnected []func(context.Context, uuid.UUID)
disconnected []func(context.Context, uuid.UUID)
standingChanged func(context.Context, models.CloudLinkStandingChange)
}
func (s *service) OnStandingChange(fn func(context.Context, models.CloudLinkStandingChange)) {
s.standingChanged = fn
}
func NewService(repo repository.CloudLinkRepository, emails repository.EmailRepository, emailSvc email.EmailService) Service {
@@ -470,6 +476,15 @@ func (s *service) CheckEnrollment(ctx context.Context, accountID uuid.UUID) (boo
}
func (s *service) ListMailboxes(ctx context.Context, orgID uuid.UUID) ([]models.CloudLinkMailboxRow, *errx.Error) {
if !reconciliationLocked(ctx) {
var rows []models.CloudLinkMailboxRow
xerr := s.reconcileLocked(ctx, func(ctx context.Context) *errx.Error {
var xerr *errx.Error
rows, xerr = s.ListMailboxes(ctx, orgID)
return xerr
})
return rows, xerr
}
accounts, xerr := s.emails.GetAllActiveInScope(ctx, repository.NewAccountScope(&orgID))
if xerr != nil {
return nil, xerr
@@ -484,15 +499,36 @@ func (s *service) ListMailboxes(ctx context.Context, orgID uuid.UUID) ([]models.
}
// One round trip for every enrolled mailbox's cloud state.
cloudByRemote := map[uuid.UUID]*models.PoolLinkMailboxState{}
cloudByAccount := map[uuid.UUID]*models.PoolLinkMailboxState{}
legacyInstances := map[uuid.UUID]bool{}
for instanceID := range mailboxGroups(enrolled) {
for instanceID, group := range mailboxGroups(enrolled) {
if l, err := s.repo.GetByInstance(ctx, instanceID); err == nil && l != nil {
legacyInstances[instanceID] = l.OrganizationID == nil
if l.DisconnectPending {
continue
}
var states []models.PoolLinkMailboxState
if xerr := s.clientFor(l).do(ctx, http.MethodGet, "/instance/mailboxes", nil, &states); xerr == nil {
byRemote := map[uuid.UUID]*models.PoolLinkMailboxState{}
for i := range states {
cloudByRemote[states[i].RemoteID] = &states[i]
byRemote[states[i].RemoteID] = &states[i]
}
for _, e := range group {
state := byRemote[e.RemoteID]
cloudByAccount[e.EmailAccountID] = state
if state == nil || state.Health == nil || !knownHealthState(state.Health.State) {
if err := s.repo.InvalidateStanding(ctx, e.EmailAccountID); err != nil {
return nil, errx.InternalError()
}
e.StandingObservedAt = nil
} else {
if _, ok := s.recordStanding(ctx, e.EmailAccountID, state.Health, false); !ok {
return nil, errx.InternalError()
}
observed := time.Now()
e.Standing, e.StandingObservedAt = state.Health, &observed
}
byAccount[e.EmailAccountID] = e
}
}
}
@@ -511,7 +547,7 @@ func (s *service) ListMailboxes(ctx context.Context, orgID uuid.UUID) ([]models.
row.Managed = e.Managed
row.EnrollmentState = e.EnrollmentState
row.StandingObservedAt = e.StandingObservedAt
row.Cloud = cloudByRemote[e.RemoteID]
row.Cloud = cloudByAccount[a.ID]
// The recorded standing covers a mailbox the cloud holds out of its pool.
if row.Cloud != nil {
row.Cloud.Health = e.EffectiveStanding(time.Now())
+11
View File
@@ -34,6 +34,17 @@ func (s *service) recordStanding(ctx context.Context, accountID uuid.UUID, h *mo
log.Warn().Err(err).Str("account_id", accountID.String()).Msg("cloud link: warmup standing could not be recorded")
return "", false
}
if !initial && s.standingChanged != nil {
previous := prev
if previous == "" {
previous = models.WarmupHealthHealthy
}
if previous != models.WarmupHealthState(h.State) {
s.standingChanged(ctx, models.CloudLinkStandingChange{
EmailAccountID: accountID, Previous: previous, Current: models.WarmupHealthState(h.State), Reason: h.Reason,
})
}
}
return prev, true
}
@@ -8,6 +8,7 @@ import (
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/warmbly/warmbly/internal/infrastructure/db"
"github.com/warmbly/warmbly/internal/models"
)
@@ -89,6 +90,15 @@ func TestLiveCloudStandingFreshnessIsBoundedWithoutInventingAProviderBlock(t *te
if reason != "" && (err != nil || h == nil || h.Reason != reason) {
t.Fatalf("cloud diagnostics = %+v, %v", h, err)
}
snap, err := NewAdvisorRepository(&db.DB{Pool: f.pool}).LoadSnapshot(ctx, f.org, time.Now())
if err != nil {
t.Fatal(err)
}
for _, m := range snap.Mailboxes {
if m.ID == f.sender && (m.PoolHealth != string(want) || m.PoolHealthReason != reason) {
t.Fatalf("advisor standing = %s (%s), want %s (%s)", m.PoolHealth, m.PoolHealthReason, want, reason)
}
}
}
assertState(models.WarmupHealthBlocked, "cloud_evidence_unavailable")
if _, err := links.SetStanding(ctx, f.sender, &models.WarmupHealthInfo{State: "healthy"}, false); err != nil {