diff --git a/internal/app/consumer/deleted_mailbox_test.go b/internal/app/consumer/deleted_mailbox_test.go new file mode 100644 index 000000000..412f173c7 --- /dev/null +++ b/internal/app/consumer/deleted_mailbox_test.go @@ -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)) + } + }) + } +} diff --git a/internal/app/consumer/event_email_error.go b/internal/app/consumer/event_email_error.go index bb5d3e285..1e9ab5831 100644 --- a/internal/app/consumer/event_email_error.go +++ b/internal/app/consumer/event_email_error.go @@ -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 +} diff --git a/internal/app/consumer/event_graph_delta_update.go b/internal/app/consumer/event_graph_delta_update.go index 066995a2f..317966199 100644 --- a/internal/app/consumer/event_graph_delta_update.go +++ b/internal/app/consumer/event_graph_delta_update.go @@ -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 } diff --git a/internal/app/consumer/event_history_id_update.go b/internal/app/consumer/event_history_id_update.go index 68bf00632..5fa657a78 100644 --- a/internal/app/consumer/event_history_id_update.go +++ b/internal/app/consumer/event_history_id_update.go @@ -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 } diff --git a/internal/app/consumer/event_mailbox_update.go b/internal/app/consumer/event_mailbox_update.go index 33ced8688..422574b5c 100644 --- a/internal/app/consumer/event_mailbox_update.go +++ b/internal/app/consumer/event_mailbox_update.go @@ -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 } diff --git a/internal/app/consumer/event_new_email.go b/internal/app/consumer/event_new_email.go index 11e2f3a16..f7de0585b 100644 --- a/internal/app/consumer/event_new_email.go +++ b/internal/app/consumer/event_new_email.go @@ -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 } diff --git a/internal/app/consumer/event_sync_state.go b/internal/app/consumer/event_sync_state.go index fe29960dc..dd6565523 100644 --- a/internal/app/consumer/event_sync_state.go +++ b/internal/app/consumer/event_sync_state.go @@ -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 } diff --git a/internal/repository/pg_danger_zone.go b/internal/repository/pg_danger_zone.go index 136e2d0bb..e8b00da89 100644 --- a/internal/repository/pg_danger_zone.go +++ b/internal/repository/pg_danger_zone.go @@ -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" diff --git a/internal/repository/pg_email_error.go b/internal/repository/pg_email_error.go index 07969b4a3..886068dd6 100644 --- a/internal/repository/pg_email_error.go +++ b/internal/repository/pg_email_error.go @@ -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 { diff --git a/internal/repository/pg_email_sync_state.go b/internal/repository/pg_email_sync_state.go index f919de7f4..76e4ce2f4 100644 --- a/internal/repository/pg_email_sync_state.go +++ b/internal/repository/pg_email_sync_state.go @@ -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 { diff --git a/internal/repository/pg_mailbox.go b/internal/repository/pg_mailbox.go index c7a58f6fd..46d5ae276 100644 --- a/internal/repository/pg_mailbox.go +++ b/internal/repository/pg_mailbox.go @@ -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 } diff --git a/internal/repository/pgerror_test.go b/internal/repository/pgerror_test.go index 4aeb7284d..2b61a2c47 100644 --- a/internal/repository/pgerror_test.go +++ b/internal/repository/pgerror_test.go @@ -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") } }