Merge pull request #667 from warmbly/fix/consumer-evict-deleted-mailbox

feat: evict a deleted mailbox from the fleet when a worker event's write is refused for it
This commit is contained in:
Matthew Meszaros
2026-09-23 17:35:34 +00:00
committed by GitHub
12 changed files with 177 additions and 22 deletions
@@ -0,0 +1,134 @@
package jobs
import (
"context"
"errors"
"fmt"
"testing"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgconn"
"github.com/warmbly/warmbly/internal/errx"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/repository"
)
// mailboxFKError is what Postgres returns when a write names a mailbox whose
// email_accounts row has been deleted.
func mailboxFKError(table string) error {
return fmt.Errorf("insert: %w", &pgconn.PgError{Code: "23503", ConstraintName: table + "_email_id_fkey"})
}
// newDeletedMailboxFixture is a consumer with two live workers and a mailbox
// lookup that answers as the given row state.
func newDeletedMailboxFixture(workerID *uuid.UUID, lookupErr *errx.Error) (*JobsService, *stubEmailRepo, *stubPublisher, []models.Worker) {
workers := []models.Worker{{ID: uuid.New()}, {ID: uuid.New()}}
repo := &stubEmailRepo{workerID: workerID, workerErr: lookupErr}
pub := &stubPublisher{}
return &JobsService{EmailRepository: repo, Publisher: pub, WorkerRepo: &stubWorkerRepo{workers: workers}}, repo, pub, workers
}
func TestDropForDeletedMailboxEvictsFromEveryWorkerWhenTheRowIsGone(t *testing.T) {
s, _, pub, workers := newDeletedMailboxFixture(nil, errx.ErrNotFound)
emailID := uuid.New()
if !s.dropForDeletedMailbox(context.Background(), uuid.New(), emailID, mailboxFKError("unibox_emails")) {
t.Fatal("a write refused for a deleted mailbox was kept as an error; it is reported and the worker keeps syncing")
}
if len(pub.removed) != len(workers) {
t.Fatalf("published %d removals, want one per live worker (%d)", len(pub.removed), len(workers))
}
for i, r := range pub.removed {
if r.workerID != workers[i].ID || r.emailID != emailID.String() {
t.Errorf("removal %d = worker %s email %s, want worker %s email %s", i, r.workerID, r.emailID, workers[i].ID, emailID)
}
}
}
// A foreign-key refusal alone does not prove the mailbox is gone: the write
// may have tripped the owner's user_id. Evicting a live mailbox would stop it
// syncing, so the row decides.
func TestDropForDeletedMailboxKeepsTheErrorWhileTheMailboxExists(t *testing.T) {
assigned := uuid.New()
for name, worker := range map[string]*uuid.UUID{"assigned": &assigned, "unassigned": nil} {
t.Run(name, func(t *testing.T) {
s, _, pub, _ := newDeletedMailboxFixture(worker, nil)
if s.dropForDeletedMailbox(context.Background(), uuid.New(), uuid.New(), mailboxFKError("unibox_emails")) {
t.Fatal("dropped a refused write for a mailbox that still exists")
}
if len(pub.removed) != 0 {
t.Errorf("published %d removals for a mailbox that still exists", len(pub.removed))
}
})
}
}
func TestDropForDeletedMailboxKeepsTheErrorWhenTheLookupFails(t *testing.T) {
s, _, pub, _ := newDeletedMailboxFixture(nil, errx.InternalError())
if s.dropForDeletedMailbox(context.Background(), uuid.New(), uuid.New(), mailboxFKError("unibox_emails")) {
t.Fatal("dropped the event on a failed lookup; an outage is not a deleted mailbox")
}
if len(pub.removed) != 0 {
t.Errorf("published %d removals on a failed lookup", len(pub.removed))
}
}
func TestDropForDeletedMailboxLeavesOtherErrorsAlone(t *testing.T) {
s, repo, _, _ := newDeletedMailboxFixture(nil, errx.ErrNotFound)
for _, err := range []error{errors.New("connection reset"), &pgconn.PgError{Code: "23505"}, nil} {
if s.dropForDeletedMailbox(context.Background(), uuid.New(), uuid.New(), err) {
t.Errorf("dropped %v, which is not a deleted mailbox", err)
}
}
if repo.workerCalls != 0 {
t.Errorf("looked the mailbox up %d times for errors that are not foreign-key refusals", repo.workerCalls)
}
}
type refusingMailboxRepo struct {
repository.MailboxRepository
err error
}
func (r *refusingMailboxRepo) CreateEntry(context.Context, uuid.UUID, uuid.UUID, *models.Mailbox) error {
return r.err
}
type refusingSyncStateRepo struct {
repository.EmailSyncStateRepository
err error
}
func (r *refusingSyncStateRepo) Get(context.Context, uuid.UUID) (*models.SyncState, error) {
return nil, nil
}
func (r *refusingSyncStateRepo) Put(context.Context, uuid.UUID, uuid.UUID, *models.SyncState) error {
return r.err
}
// The handlers are where the refusal used to be discarded or reported; each
// must now end the event and evict.
func TestWorkerEventsForADeletedMailboxEvictInsteadOfFailing(t *testing.T) {
cases := map[string]func(s *JobsService, userID, emailID uuid.UUID) error{
"MAILBOX_UPDATE": func(s *JobsService, userID, emailID uuid.UUID) error {
s.MailboxRepository = &refusingMailboxRepo{err: mailboxFKError("unibox_mailboxes")}
return s.HandleMailboxUpdate(context.Background(), &models.JobEventMailboxUpdate{UserID: userID, EmailID: emailID, Data: &models.Mailbox{}})
},
"SYNC_STATE": func(s *JobsService, userID, emailID uuid.UUID) error {
s.EmailSyncStateRepository = &refusingSyncStateRepo{err: mailboxFKError("email_sync_state")}
return s.HandleSyncState(context.Background(), &models.JobEventSyncState{UserID: userID, EmailID: emailID})
},
}
for name, run := range cases {
t.Run(name, func(t *testing.T) {
s, _, pub, workers := newDeletedMailboxFixture(nil, errx.ErrNotFound)
if err := run(s, uuid.New(), uuid.New()); err != nil {
t.Fatalf("returned %v; the event is redelivered and reported for a mailbox that cannot come back", err)
}
if len(pub.removed) != len(workers) {
t.Errorf("published %d removals, want %d", len(pub.removed), len(workers))
}
})
}
}
@@ -423,3 +423,19 @@ func (s *JobsService) evictDeletedMailbox(ctx context.Context, userID, emailAcco
}
}
}
// dropForDeletedMailbox reports whether a worker event's write was refused
// because its mailbox row is gone, and if so evicts the mailbox from every
// worker. The event is then done: no retry can bring the row back.
func (s *JobsService) dropForDeletedMailbox(ctx context.Context, userID, emailAccountID uuid.UUID, err error) bool {
if !repository.IsForeignKeyViolation(err) || s.EmailRepository == nil {
return false
}
// Confirmed against the row, not inferred from the refusal: the write may
// have tripped a different foreign key, such as the owner's user_id.
if _, xerr := s.EmailRepository.GetWorkerID(ctx, emailAccountID); xerr == nil || !errors.Is(xerr, errx.ErrNotFound) {
return false
}
s.evictDeletedMailbox(ctx, userID, emailAccountID)
return true
}
@@ -14,6 +14,9 @@ func (s *JobsService) HandleGraphDeltaUpdate(ctx context.Context, e *models.JobE
return nil
}
if err := s.EmailGraphDeltaRepository.Put(ctx, e.UserID, e.EmailID, e.Folder, e.DeltaLink); err != nil {
if s.dropForDeletedMailbox(ctx, e.UserID, e.EmailID, err) {
return nil
}
CaptureError(e.UserID, e.EmailID, err)
return err
}
@@ -8,6 +8,9 @@ import (
func (s *JobsService) HandleHistoryIDUpdate(ctx context.Context, e *models.JobEventHistoryIDUpdate) error {
if err := s.EmailHistoryIDRepository.Put(ctx, e.UserID, e.EmailID, e.HistoryID); err != nil {
if s.dropForDeletedMailbox(ctx, e.UserID, e.EmailID, err) {
return nil
}
CaptureError(e.UserID, e.EmailID, err)
return err
}
@@ -13,6 +13,9 @@ func (s *JobsService) HandleMailboxUpdate(ctx context.Context, e *models.JobEven
e.EmailID,
e.Data,
); err != nil { // replaces the data if it already exists
if s.dropForDeletedMailbox(ctx, e.UserID, e.EmailID, err) {
return nil
}
CaptureError(e.UserID, e.EmailID, err)
return err
}
+6
View File
@@ -23,6 +23,9 @@ func (s *JobsService) HandleNewEmail(ctx context.Context, e *models.JobEventNewE
err := s.ingestNewEmail(ctx, e)
if errors.Is(err, errWarmupVerification) && s.UniboxRepository != nil {
if storeErr := s.UniboxRepository.DeferWarmupVerification(ctx, e); storeErr != nil {
if s.dropForDeletedMailbox(ctx, e.UserID, e.Message.EmailID, storeErr) {
return nil
}
return fmt.Errorf("defer inbox arrival: %w", storeErr)
}
log.Warn().Err(err).Str("email_id", e.Message.EmailID.String()).Msg("inbox arrival queued for warmup verification")
@@ -90,6 +93,9 @@ func (s *JobsService) ingestNewEmail(ctx context.Context, e *models.JobEventNewE
// Normal email processing
if err := s.UniboxRepository.CreateEntry(ctx, e.UserID, e.Message); err != nil {
if s.dropForDeletedMailbox(ctx, e.UserID, e.Message.EmailID, err) {
return nil
}
CaptureError(e.UserID, e.Message.EmailID, err)
return err
}
@@ -27,6 +27,9 @@ func (s *JobsService) HandleSyncState(ctx context.Context, e *models.JobEventSyn
}
if err := s.EmailSyncStateRepository.Put(ctx, e.UserID, e.EmailID, &e.State); err != nil {
if s.dropForDeletedMailbox(ctx, e.UserID, e.EmailID, err) {
return nil
}
CaptureError(e.UserID, e.EmailID, err)
return err
}
+2 -2
View File
@@ -463,11 +463,11 @@ func (r *dangerZoneRepository) hardDelete(ctx context.Context, mailboxScope, del
return placements, nil
}
// isForeignKeyViolation detects Postgres SQLSTATE 23503 (foreign_key_violation),
// IsForeignKeyViolation detects Postgres SQLSTATE 23503 (foreign_key_violation),
// which for a write keyed on a mailbox or a user means the parent row was
// deleted before the write landed. Matched through errors.As so a wrapped error
// still reports.
func isForeignKeyViolation(err error) bool {
func IsForeignKeyViolation(err error) bool {
var pgErr *pgconn.PgError
if errors.As(err, &pgErr) {
return pgErr.Code == "23503"
+1 -1
View File
@@ -141,7 +141,7 @@ func (r *emailAccountErrorRepository) CreateOnce(ctx context.Context, data *Crea
// incident: every one of these events names an id a worker held minutes
// ago, so a deleted mailbox with a busy sync loop reported one of these a
// minute.
if isForeignKeyViolation(err) {
if IsForeignKeyViolation(err) {
return nil, nil
}
if err != nil {
+1 -7
View File
@@ -78,13 +78,7 @@ func (r *pgEmailSyncStateRepository) Put(ctx context.Context, userID, emailID uu
state.BackfillSince, state.BackfillStartedAt, state.BackfillCompletedAt,
state.ThrottledUntil, state.ThrottleReason, state.Deferred, state.LastSyncedAt,
); err != nil {
// The mailbox was deleted while its sync pass was in flight. Nothing
// can own this state and no retry changes that, so the relay is done
// rather than failed; returned as an error it was reported and
// redelivered for as long as the worker kept relaying.
if isForeignKeyViolation(err) {
return nil
}
// A deleted mailbox refuses this as a foreign-key violation; the consumer evicts it.
return fmt.Errorf("email_sync_state: put: %w", err)
}
if state.LastSyncedAt != nil {
+1 -8
View File
@@ -56,14 +56,7 @@ func (r *mailboxRepository) CreateEntry(ctx context.Context, userId, emailId uui
_, err := r.db.Exec(ctx, query,
emailId, mb.UIDValidity, mb.Name, textArray(mb.Attrs), mb.HighestModSeq, mb.UIDNext, mb.UpdatedAt, mb.Delim,
)
// The mailbox was deleted between the worker listing its folders and this
// write landing. There is no parent to hang a folder off and never will be
// again, so the event is done rather than failed: returned as an error it
// was reported and redelivered forever, because no retry can bring the
// mailbox back. Same call as pg_email_error.go makes on the same race.
if isForeignKeyViolation(err) {
return nil
}
// A deleted mailbox refuses this as a foreign-key violation; the consumer evicts it.
return err
}
+4 -4
View File
@@ -14,16 +14,16 @@ import (
func TestIsForeignKeyViolation(t *testing.T) {
fk := &pgconn.PgError{Code: "23503", ConstraintName: "email_account_errors_email_account_id_fkey"}
if !isForeignKeyViolation(fk) {
if !IsForeignKeyViolation(fk) {
t.Fatal("bare 23503 not detected")
}
if !isForeignKeyViolation(fmt.Errorf("queryrow failed: %w", fk)) {
if !IsForeignKeyViolation(fmt.Errorf("queryrow failed: %w", fk)) {
t.Fatal("wrapped 23503 not detected")
}
if isForeignKeyViolation(&pgconn.PgError{Code: "23505"}) {
if IsForeignKeyViolation(&pgconn.PgError{Code: "23505"}) {
t.Fatal("unique violation reported as a foreign key violation")
}
if isForeignKeyViolation(errors.New("boom")) || isForeignKeyViolation(nil) {
if IsForeignKeyViolation(errors.New("boom")) || IsForeignKeyViolation(nil) {
t.Fatal("non-pg error reported as a foreign key violation")
}
}