diff --git a/internal/api/handler/email_sync.go b/internal/api/handler/email_sync.go index d42e1f0b6..0accf2f60 100644 --- a/internal/api/handler/email_sync.go +++ b/internal/api/handler/email_sync.go @@ -34,12 +34,7 @@ func (h *Handler) GetEmailSync(c *gin.Context) { errx.JSON(c, errx.ErrUnauthorized) return } - state, policy, xerr := h.EmailService.GetSyncState(c.Request.Context(), orgID.String(), c.Param("id")) - if xerr != nil { - errx.JSON(c, xerr) - return - } - folders, xerr := h.EmailService.GetSyncFolders(c.Request.Context(), orgID.String(), c.Param("id")) + state, policy, folders, xerr := h.EmailService.GetSyncState(c.Request.Context(), orgID.String(), c.Param("id")) if xerr != nil { errx.JSON(c, xerr) return diff --git a/internal/app/consumer/event_mailbox_delete.go b/internal/app/consumer/event_mailbox_delete.go index 5f97f797b..8554b76d8 100644 --- a/internal/app/consumer/event_mailbox_delete.go +++ b/internal/app/consumer/event_mailbox_delete.go @@ -3,10 +3,31 @@ package jobs import ( "context" + "github.com/google/uuid" "github.com/rs/zerolog/log" + "github.com/warmbly/warmbly/internal/infrastructure/pubsub" "github.com/warmbly/warmbly/internal/models" ) +// publishFolderPurged tells open dashboards that rows left a mailbox in bulk. +// One org-scoped EMAIL_DELETED with no message id: the client drops every +// inbox list on that event whatever it names, and one event is what a purge +// of a whole folder deserves rather than one per row. +func (s *JobsService) publishFolderPurged(ctx context.Context, userID, emailID uuid.UUID) { + if s.StreamingPublisher == nil { + return + } + var orgID string + if account, err := s.EmailRepository.GetByID(ctx, emailID); err == nil && account != nil && account.OrganizationID != nil { + orgID = account.OrganizationID.String() + } + s.StreamingPublisher.PublishEmailDeleted(ctx, &pubsub.EmailInboxEvent{ + BaseEvent: pubsub.BaseEvent{UserID: userID.String()}, + OrgID: orgID, + EmailAccountID: emailID.String(), + }) +} + // HandleMailboxDelete retires a folder the last listing no longer had. // // A folder is identified by its name. UIDValidity is the fallback for an @@ -35,6 +56,9 @@ func (s *JobsService) HandleMailboxDelete(ctx context.Context, e *models.JobEven Str("folder", e.Mailbox). Int64("messages", n). Msg("folder excluded from sync: stored mail dropped") + if n > 0 { + s.publishFolderPurged(ctx, e.UserID, e.EmailID) + } } return nil } diff --git a/internal/app/email/loader.go b/internal/app/email/loader.go index b7a8b958e..763fa0ab1 100644 --- a/internal/app/email/loader.go +++ b/internal/app/email/loader.go @@ -2,6 +2,7 @@ package email import ( "context" + "fmt" "strings" "time" @@ -260,6 +261,10 @@ func (s *emailService) buildAddWorkerEmail(ctx context.Context, acc *models.Emai provider := models.InboxProvider(acc.Provider) saveToSent := acc.SaveToSent + sync, err := s.syncDataFor(ctx, acc.ID) + if err != nil { + return nil, err + } out := &models.AddWorkerEmail{ ID: acc.ID, UserID: userID, @@ -268,7 +273,7 @@ func (s *emailService) buildAddWorkerEmail(ctx context.Context, acc *models.Emai FirstName: first, LastName: last, Type: provider, - Sync: s.syncDataFor(ctx, acc.ID), + Sync: sync, // Only SMTP/IMAP acts on this; Gmail and Graph file their own copy. SaveToSent: &saveToSent, } @@ -355,7 +360,7 @@ func (s *emailService) lastHistoryFor(ctx context.Context, userID, emailID uuid. // state a previous worker left behind. Policy comes from instance settings // (compiled defaults when none are wired), so an operator's change applies at // the next load: onboarding, reassignment, or the reconciler's republish. -func (s *emailService) syncDataFor(ctx context.Context, emailID uuid.UUID) *models.AddWorkerEmailSyncData { +func (s *emailService) syncDataFor(ctx context.Context, emailID uuid.UUID) (*models.AddWorkerEmailSyncData, error) { budget := instancesettings.DefaultSync() if s.syncBudget != nil { budget = s.syncBudget.SyncBudget(ctx) @@ -368,11 +373,15 @@ func (s *emailService) syncDataFor(ctx context.Context, emailID uuid.UUID) *mode OrgDailyMessages: budget.DailyMessagesPerOrg, }, } - if skip, xerr := s.emailRepository.GetSyncSkipFolders(ctx, emailID); xerr == nil { - data.Policy.SkipFolders = skip - } else { - log.Warn().Str("email_id", emailID.String()).Msg("sync skip folders lookup failed; worker syncs every folder until the next republish") + // The skip list is part of the policy a republish replaces on the loaded + // mailbox, so a failed read cannot fall back to "skip nothing": that + // would have the worker baseline and import the excluded folders until + // the next republish. The load fails instead and the reconciler retries. + skip, xerr := s.emailRepository.GetSyncSkipFolders(ctx, emailID) + if xerr != nil { + return nil, fmt.Errorf("sync skip folders lookup: %w", xerr) } + data.Policy.SkipFolders = skip // A pool-linked mailbox is a warmup-only mirror: no history import. if s.poolLink != nil { if linked, err := s.poolLink.GetMailboxByAccount(ctx, emailID); err == nil && linked != nil { @@ -387,7 +396,7 @@ func (s *emailService) syncDataFor(ctx context.Context, emailID uuid.UUID) *mode log.Warn().Err(err).Str("email_id", emailID.String()).Msg("sync state lookup failed; worker starts fresh") } } - return data + return data, nil } // mailboxesFor is the IMAP folder state (name, UIDVALIDITY, HIGHESTMODSEQ) diff --git a/internal/app/email/service.go b/internal/app/email/service.go index 440ed0e36..14fe2a41b 100644 --- a/internal/app/email/service.go +++ b/internal/app/email/service.go @@ -2,6 +2,7 @@ package email import ( "context" + "slices" "sort" "strings" "time" @@ -133,11 +134,10 @@ type EmailService interface { // for a change outside the mailbox row (Warmbly Cloud enrollment). SyncWarmupPool(ctx context.Context, accountID uuid.UUID) // GetSyncState is the dashboard's view of a mailbox's sync: nil state when - // the worker has not reported yet. - GetSyncState(ctx context.Context, userID, emailID string) (*models.SyncState, models.SyncPolicy, *errx.Error) - // GetSyncFolders is the folders the sync has seen on the server, so a - // client can name one to skip. Empty for providers without folders. - GetSyncFolders(ctx context.Context, orgID, emailID string) ([]models.SyncFolder, *errx.Error) + // the worker has not reported yet, the policy in force, and the folders + // the sync has seen on the server so a client can name one to skip + // (empty for providers without folders). + GetSyncState(ctx context.Context, userID, emailID string) (*models.SyncState, models.SyncPolicy, []models.SyncFolder, *errx.Error) // UpdateSyncSettings replaces the mailbox's skip list, drops the mail // already stored from those folders, and re-ships the mailbox so the // worker applies it on its next pass. Returns the list as stored. @@ -366,29 +366,25 @@ func (s *emailService) publishAccountEvent(ctx context.Context, eventType pubsub // GetSyncState returns the persisted sync state and the policy currently in // force. It goes through Get so ownership is checked the same way as every // other per-mailbox read. -func (s *emailService) GetSyncState(ctx context.Context, orgID, emailID string) (*models.SyncState, models.SyncPolicy, *errx.Error) { +func (s *emailService) GetSyncState(ctx context.Context, orgID, emailID string) (*models.SyncState, models.SyncPolicy, []models.SyncFolder, *errx.Error) { acc, xerr := s.Get(ctx, orgID, emailID) if xerr != nil { - return nil, models.SyncPolicy{}, xerr + return nil, models.SyncPolicy{}, nil, xerr } - data := s.syncDataFor(ctx, acc.ID) - return data.State, data.Policy, nil + data, err := s.syncDataFor(ctx, acc.ID) + if err != nil { + log.Error().Err(err).Str("email_id", acc.ID.String()).Msg("sync state: policy lookup failed") + return nil, models.SyncPolicy{}, nil, errx.InternalError() + } + return data.State, data.Policy, s.syncFoldersFor(ctx, acc), nil } -// GetSyncFolders lists the IMAP folders the worker has reported for this +// syncFoldersFor lists the IMAP folders the worker has reported for this // mailbox, INBOX first and then by name, each with the canonical folder it // files under. Gmail and Outlook mailboxes have no folder list here. -func (s *emailService) GetSyncFolders(ctx context.Context, orgID, emailID string) ([]models.SyncFolder, *errx.Error) { - acc, xerr := s.Get(ctx, orgID, emailID) - if xerr != nil { - return nil, xerr - } +func (s *emailService) syncFoldersFor(ctx context.Context, acc *models.Email) []models.SyncFolder { out := []models.SyncFolder{} - userID, err := uuid.Parse(acc.UserID) - if acc.Provider != string(models.InboxProviderSMTPIMAP) || err != nil { - return out, nil - } - for _, box := range s.mailboxesFor(ctx, userID, acc.ID) { + for _, box := range s.imapFoldersFor(ctx, acc) { out = append(out, models.SyncFolder{Name: box.Name, Folder: imap.CanonicalFolder(box)}) } sort.SliceStable(out, func(i, j int) bool { @@ -398,14 +394,26 @@ func (s *emailService) GetSyncFolders(ctx context.Context, orgID, emailID string } return strings.ToLower(out[i].Name) < strings.ToLower(out[j].Name) }) - return out, nil + return out } -// UpdateSyncSettings stores a normalized skip list for an IMAP mailbox. The -// purge of already-stored mail is exact-name: the worker retires each -// subfolder it stops following by name on its next pass, and the re-ship -// makes that pass apply the new list within a minute rather than at the -// reconciler's next republish. +// imapFoldersFor is the saved folder listing of an IMAP mailbox, and nothing +// for any other provider. +func (s *emailService) imapFoldersFor(ctx context.Context, acc *models.Email) []models.Mailbox { + userID, err := uuid.Parse(acc.UserID) + if acc.Provider != string(models.InboxProviderSMTPIMAP) || err != nil { + return nil + } + return s.mailboxesFor(ctx, userID, acc.ID) +} + +// UpdateSyncSettings stores a normalized skip list for an IMAP mailbox and +// drops the mail already stored from every saved folder the list covers, +// decided by the same matcher the worker applies (case, subfolders), plus +// the names themselves for a folder not listed yet. The re-ship makes the +// worker's next pass apply the list within a minute rather than at the +// reconciler's next republish; a folder it retires then is purged again by +// name, which is idempotent. func (s *emailService) UpdateSyncSettings(ctx context.Context, orgID, emailID string, body *models.UpdateSyncSettings) ([]string, *errx.Error) { acc, xerr := s.Get(ctx, orgID, emailID) if xerr != nil { @@ -425,7 +433,13 @@ func (s *emailService) UpdateSyncSettings(ctx context.Context, orgID, emailID st return nil, xerr } if s.unibox != nil && len(folders) > 0 { - if n, err := s.unibox.DeleteByFolderPaths(ctx, acc.ID, folders); err != nil { + purge := append([]string(nil), folders...) + for _, box := range s.imapFoldersFor(ctx, acc) { + if imap.SkipsFolder(box, folders) && !slices.Contains(purge, box.Name) { + purge = append(purge, box.Name) + } + } + if n, err := s.unibox.DeleteByFolderPaths(ctx, acc.ID, purge); err != nil { log.Warn().Err(err).Str("email_id", acc.ID.String()).Msg("sync skip folders: purge of stored mail failed; the worker retires the folders on its next pass") } else if n > 0 { log.Info().Str("email_id", acc.ID.String()).Int64("messages", n).Msg("sync skip folders: stored mail from skipped folders dropped") diff --git a/internal/app/worker/wmail/sync_imap.go b/internal/app/worker/wmail/sync_imap.go index bcf453405..5c2a11e89 100644 --- a/internal/app/worker/wmail/sync_imap.go +++ b/internal/app/worker/wmail/sync_imap.go @@ -116,7 +116,9 @@ func (w *WMail) Sync(ctx context.Context) *errx.MailError { } changed := imapFolderChanged(befBox, box, condStore) - movedOut := len(skipped) > 0 && listedBefore && imapMovedOut(prevListing, box) + // A pass cut short by the search cap left rows unexamined, so the + // next pass looks again whether or not the count moved. + movedOut := len(skipped) > 0 && (w.skipPending[box.Name] || listedBefore && imapMovedOut(prevListing, box)) fullyProcessed := true var touched map[string]struct{} var view imap.Selected @@ -162,7 +164,7 @@ func (w *WMail) Sync(ctx context.Context) *errx.MailError { return err } } else if movedOut && !stats.aborted { - if err := w.imapReconcileSkipped(ctx, box, skipped, stats); err != nil { + if err := w.imapReconcileSkipped(ctx, box, skipped, touched, stats); err != nil { return err } } @@ -185,6 +187,7 @@ func (w *WMail) Sync(ctx context.Context) *errx.MailError { // Renames were already followed above, so a name missing from the listing // at this point really is a folder that is gone. var deleted []string + var gone []*models.Mailbox outer: for _, box := range w.SmtpImapData.Mailboxes { for _, f := range folders { @@ -192,7 +195,9 @@ outer: continue outer } } - + gone = append(gone, box) + } + for _, box := range gone { if err := w.onEvent(models.JobEventTypeMailboxDelete, &models.JobEventMailboxDelete{ UserID: w.UserID, EmailID: w.ID, @@ -200,7 +205,7 @@ outer: UIDValidity: box.UIDValidity, // A folder that is still on the server but now excluded takes // the mail already stored from it along. - Skipped: slices.ContainsFunc(skipped, func(s models.Mailbox) bool { return s.Name == box.Name }), + Skipped: imapRetiredIntoSkipped(box, gone, skipped), }); err != nil { return nil } @@ -211,6 +216,7 @@ outer: for _, name := range deleted { delete(w.flagScan, name) delete(w.listed, name) + delete(w.skipPending, name) // The backfill floor goes with the folder. A name is reusable, // and a floor left behind would be inherited by whatever is // created under it next. @@ -235,6 +241,32 @@ outer: return nil } +// imapRetiredIntoSkipped reports whether a folder leaving the listing is one +// the owner excluded: listed under the same name in the skipped set, or +// renamed into the skipped subtree, which the rename matcher could not see +// because skipped folders leave the listing before it runs. The rename is +// claimed on the matcher's own terms: exactly one folder gone and exactly +// one skipped folder carrying its UIDVALIDITY, so a server that stamps a +// whole tree from one creation time cannot make an unrelated deletion look +// like a move. +func imapRetiredIntoSkipped(box *models.Mailbox, gone []*models.Mailbox, skipped []models.Mailbox) bool { + if slices.ContainsFunc(skipped, func(s models.Mailbox) bool { return s.Name == box.Name }) { + return true + } + sameGone, sameSkipped := 0, 0 + for _, g := range gone { + if g.UIDValidity == box.UIDValidity { + sameGone++ + } + } + for i := range skipped { + if skipped[i].UIDValidity == box.UIDValidity { + sameSkipped++ + } + } + return sameGone == 1 && sameSkipped == 1 +} + // skipFolders is the owner's exclusion list as the policy in force carries // it; a republished ADD_EMAIL changes it between passes. func (w *WMail) skipFolders() []string { @@ -279,10 +311,19 @@ func imapMovedOut(before imapListed, now *models.Mailbox) bool { // Message-ID is found in a skipped folder. Rows checked once and found // nowhere are remembered for the session, so a folder the owner emptied by // hand does not cost a search per row on every later pass. -func (w *WMail) imapReconcileSkipped(ctx context.Context, box *models.Mailbox, skipped []models.Mailbox, stats *tickStats) *errx.MailError { +// +// The work is bounded per pass: at most imapSkipSearchesPerPass rows are +// looked for, and a folder with more left over is marked pending so the next +// pass continues where this one stopped. A row fetched this pass (touched) +// is live under a new UID, whatever an older copy elsewhere says. +func (w *WMail) imapReconcileSkipped(ctx context.Context, box *models.Mailbox, skipped []models.Mailbox, touched map[string]struct{}, stats *tickStats) *errx.MailError { if w.SyncContext == nil || len(skipped) == 0 { return nil } + if w.skipPending == nil { + w.skipPending = make(map[string]bool) + } + delete(w.skipPending, box.Name) stored, err := w.SyncContext.ListFolderMessages(ctx, w.UserID, w.ID, box.Name, box.UIDValidity) if err != nil { return w.controlPlaneError(err, stats) @@ -311,10 +352,14 @@ func (w *WMail) imapReconcileSkipped(ctx context.Context, box *models.Mailbox, s if w.skipChecked == nil || len(w.skipChecked) > imapSkipCheckedMax { w.skipChecked = make(map[string]struct{}) } + searches := 0 for _, m := range stored { if _, ok := live[m.UID]; ok { continue } + if _, refiled := touched[m.MessageID]; refiled { + continue + } if ctx.Err() != nil { return nil } @@ -322,6 +367,11 @@ func (w *WMail) imapReconcileSkipped(ctx context.Context, box *models.Mailbox, s if _, done := w.skipChecked[key]; done { continue } + if searches >= imapSkipSearchesPerPass { + w.skipPending[box.Name] = true + return nil + } + searches++ folder := w.imapFindInSkipped(ctx, skipped, m.MessageID) if folder == "" { w.skipChecked[key] = struct{}{} @@ -335,6 +385,12 @@ func (w *WMail) imapReconcileSkipped(ctx context.Context, box *models.Mailbox, s }); err != nil { return w.controlPlaneError(err, stats) } + // The map entry goes with the row: the sync reads a mapped + // Message-ID as already stored, and would never import the message + // again if it moved back into a synced folder. + if err := w.EmailMessageMapRepository.Del(ctx, w.UserID, w.ID, m.MessageID, m.ID); err != nil { + return w.controlPlaneError(err, stats) + } w.skipChecked[key] = struct{}{} } return nil @@ -342,7 +398,12 @@ func (w *WMail) imapReconcileSkipped(ctx context.Context, box *models.Mailbox, s // imapSkipCheckedMax bounds the per-session memory of rows already looked // for in the skipped folders; past it the memory starts over. -const imapSkipCheckedMax = 20_000 +const imapSkipCheckedMax = 250_000 + +// imapSkipSearchesPerPass caps the searches one folder's reconciliation +// spends in one pass, so an owner emptying a large folder by hand costs a +// bounded slice of every tick rather than one long one. +const imapSkipSearchesPerPass = 50 // imapFindInSkipped names the skipped folder holding the message, or "". // A key the sync made up for a message without a Message-ID was never on @@ -410,10 +471,10 @@ func (w *WMail) imapIncremental(ctx context.Context, box, before *models.Mailbox // Newest first: when budget is short, the freshest mail lands first. sort.Slice(uids, func(i, j int) bool { return uids[i] > uids[j] }) - // Only the drafts folder needs the fetched ids back, so nothing else pays - // for the set. + // Only the two reconciliations need the fetched ids back (drafts, and + // mail that left for a skipped folder), so nothing else pays for the set. var touched map[string]struct{} - if imapCanonicalFolder(box) == models.FolderDrafts { + if imapCanonicalFolder(box) == models.FolderDrafts || len(w.skipFolders()) > 0 { touched = make(map[string]struct{}, len(uids)) } diff --git a/internal/app/worker/wmail/sync_imap_skip_test.go b/internal/app/worker/wmail/sync_imap_skip_test.go index f18547696..3dfe5229a 100644 --- a/internal/app/worker/wmail/sync_imap_skip_test.go +++ b/internal/app/worker/wmail/sync_imap_skip_test.go @@ -2,6 +2,7 @@ package wmail import ( "context" + "fmt" "testing" goimap "github.com/emersion/go-imap/v2" @@ -243,3 +244,122 @@ func TestImapSyncBaselinesCountsOnFirstListing(t *testing.T) { t.Fatalf("removed %d rows on the second listing, want 1 (lookups=%d finds=%d events=%v listed=%+v)", len(removeIDs(*events)), ctx.calls, conn.finds, kinds, w.listed) } } + +// A synced folder renamed into the skipped subtree keeps its UIDVALIDITY, +// but the rename matcher never sees the new name because skipped folders +// leave the listing first. The retirement still carries the marker, so the +// mail stored under the old name is purged like any skipped folder's. +func TestImapSyncMarksFolderRenamedIntoSkippedSubtree(t *testing.T) { + conn := &fakeImapConn{folders: []models.Mailbox{ + {Name: "INBOX", UIDValidity: 7, HighestModSeq: 100, Delim: "/"}, + {Name: "Warmer", UIDValidity: 9, HighestModSeq: 100, Delim: "/"}, + {Name: "Warmer/Leads", UIDValidity: 21, HighestModSeq: 50, Delim: "/"}, + {Name: "Receipts", UIDValidity: 33, HighestModSeq: 10, Delim: "/"}, + }} + budget := &skipBudget{fixedBudget: &fixedBudget{allow: 10}, skip: []string{"Warmer"}} + w, events := newIMAPTestMail(conn, budget, &models.Mailbox{Name: "INBOX", UIDValidity: 7, HighestModSeq: 100}) + w.SmtpImapData.Mailboxes = append(w.SmtpImapData.Mailboxes, + &models.Mailbox{Name: "Leads", UIDValidity: 21, HighestModSeq: 50}, + // Gone for good, and its UIDVALIDITY matches nothing skipped. + &models.Mailbox{Name: "Old", UIDValidity: 44, HighestModSeq: 5}, + ) + + if err := w.Sync(t.Context()); err != nil { + t.Fatalf("Sync: %v", err) + } + got := map[string]bool{} + for _, e := range mailboxEvents(*events, models.JobEventTypeMailboxDelete) { + del := e.body.(*models.JobEventMailboxDelete) + got[del.Mailbox] = del.Skipped + } + if skipped, ok := got["Leads"]; !ok || !skipped { + t.Fatalf("Leads retired as %v (present %v), want skipped", skipped, ok) + } + if skipped, ok := got["Old"]; !ok || skipped { + t.Fatalf("Old retired as %v (present %v), want a plain deletion", skipped, ok) + } + if len(mailboxEvents(*events, models.JobEventTypeMailboxRename)) != 0 { + t.Fatal("a rename into a skipped folder must not be followed") + } +} + +// A row whose Message-ID was fetched this pass is live under a new UID, +// whatever a copy in a skipped folder says; it is never removed. +func TestImapSyncKeepsRowRefetchedThisPass(t *testing.T) { + conn := &fakeImapConn{ + folders: []models.Mailbox{ + {Name: "INBOX", UIDValidity: 7, HighestModSeq: 200, UIDNext: 12, Messages: 1, Delim: "/"}, + {Name: "Warmer", UIDValidity: 9, HighestModSeq: 100, Delim: "/"}, + }, + changed: []goimap.UID{11}, + all: []goimap.UID{11}, + inSkipped: map[string]map[string]uint32{"Warmer": {"<11@fake.test>": 3}}, + } + budget := &skipBudget{fixedBudget: &fixedBudget{allow: 10}, skip: []string{"Warmer"}} + w, events := newIMAPTestMail(conn, budget, &models.Mailbox{Name: "INBOX", UIDValidity: 7, HighestModSeq: 100, UIDNext: 10}) + w.rememberListing(&models.Mailbox{Name: "INBOX", Messages: 1, UIDNext: 10}) + w.EmailMessageMapRepository = knownMessageMap{id: uuid.New().String()} + // The stored row is the same message under its old UID: re-appended + // this pass as UID 11 (a warmup engagement leg moved it out and back). + w.SyncContext = &fakeSyncContext{stored: map[string][]repository.StoredFolderMessage{ + "INBOX": {{UID: 5, MessageID: "<11@fake.test>", ID: uuid.New()}}, + }} + + if err := w.Sync(t.Context()); err != nil { + t.Fatalf("Sync: %v", err) + } + if n := len(removeIDs(*events)); n != 0 { + t.Fatalf("removed %d rows, want none: the message is live under a new UID", n) + } + if conn.finds != 0 { + t.Fatalf("searched %d times for a row fetched this pass", conn.finds) + } +} + +// The reconciliation spends at most imapSkipSearchesPerPass searches on one +// folder per pass and continues on the next, so a folder emptied by hand is +// examined in slices rather than in one long tick. +func TestImapSyncCapsSkippedSearchesPerPass(t *testing.T) { + conn := &fakeImapConn{ + folders: []models.Mailbox{ + {Name: "INBOX", UIDValidity: 7, HighestModSeq: 100, UIDNext: 500, Messages: 0, Delim: "/"}, + {Name: "Warmer", UIDValidity: 9, HighestModSeq: 100, Delim: "/"}, + }, + inSkipped: map[string]map[string]uint32{"Warmer": {}}, + } + budget := &skipBudget{fixedBudget: &fixedBudget{allow: 10}, skip: []string{"Warmer"}} + w, events := newIMAPTestMail(conn, budget, &models.Mailbox{Name: "INBOX", UIDValidity: 7, HighestModSeq: 100, UIDNext: 500}) + w.rememberListing(&models.Mailbox{Name: "INBOX", Messages: 120, UIDNext: 500}) + w.EmailMessageMapRepository = knownMessageMap{id: uuid.New().String()} + stored := make([]repository.StoredFolderMessage, 0, 120) + for i := 1; i <= 120; i++ { + id := fmt.Sprintf("<%d@fake.test>", i) + stored = append(stored, repository.StoredFolderMessage{UID: uint32(i), MessageID: id, ID: uuid.New()}) + } + // The last one is in the skipped folder; it is reached on the third pass. + conn.inSkipped["Warmer"]["<120@fake.test>"] = 9 + ctx := &fakeSyncContext{stored: map[string][]repository.StoredFolderMessage{"INBOX": stored}} + w.SyncContext = ctx + + for pass, wantFinds := range []int{imapSkipSearchesPerPass, 2 * imapSkipSearchesPerPass, 120} { + if err := w.Sync(t.Context()); err != nil { + t.Fatalf("pass %d: %v", pass+1, err) + } + if conn.finds != wantFinds { + t.Fatalf("pass %d: %d searches so far, want %d", pass+1, conn.finds, wantFinds) + } + } + if n := len(removeIDs(*events)); n != 1 { + t.Fatalf("removed %d rows over three passes, want the one found in Warmer", n) + } + if w.skipPending["INBOX"] { + t.Fatal("the folder is still marked pending after every row was examined") + } + // A fourth pass with nothing new does not look again. + if err := w.Sync(t.Context()); err != nil { + t.Fatalf("fourth pass: %v", err) + } + if conn.finds != 120 { + t.Fatalf("a settled folder was searched again (finds = %d)", conn.finds) + } +} diff --git a/internal/app/worker/wmail/wmail.go b/internal/app/worker/wmail/wmail.go index 19da6a9b1..3e6e594c1 100644 --- a/internal/app/worker/wmail/wmail.go +++ b/internal/app/worker/wmail/wmail.go @@ -107,7 +107,12 @@ type WMail struct { // folder is baselined on first sight and the stored cursor is not used, // because a walked folder's cursor advances to the SELECT view and a // server whose view lags its STATUS would read as departures every pass. + // A move that lands while the worker is down is therefore not seen; the + // row it leaves is retired by the next departure from that folder. listed map[string]imapListed + // skipPending marks folders whose skipped-folder reconciliation hit the + // per-pass search cap, so the next pass continues it. + skipPending map[string]bool // unmapPending holds map entries for unpublished arrivals whose removal // failed, keyed by map key; every pass retries them before it looks. unmapPending map[string]uuid.UUID diff --git a/internal/infrastructure/db/migrations/000197_unibox_mailboxes_delim.down.sql b/internal/infrastructure/db/migrations/000197_unibox_mailboxes_delim.down.sql new file mode 100644 index 000000000..e8e53ed4f --- /dev/null +++ b/internal/infrastructure/db/migrations/000197_unibox_mailboxes_delim.down.sql @@ -0,0 +1,2 @@ +ALTER TABLE public.unibox_mailboxes + DROP COLUMN IF EXISTS delim; diff --git a/internal/infrastructure/db/migrations/000197_unibox_mailboxes_delim.up.sql b/internal/infrastructure/db/migrations/000197_unibox_mailboxes_delim.up.sql new file mode 100644 index 000000000..bd9da6ea6 --- /dev/null +++ b/internal/infrastructure/db/migrations/000197_unibox_mailboxes_delim.up.sql @@ -0,0 +1,8 @@ +-- The hierarchy delimiter a server reported for each folder is kept with the +-- folder, so the control plane can tell a subfolder from a sibling whose +-- name merely starts the same way. The worker matches a skipped folder's +-- subfolders on it; the purge that follows a folder being excluded from sync +-- has to reach the same set of stored rows. Empty when the server reported +-- none, where a folder has no subfolders to speak of. +ALTER TABLE public.unibox_mailboxes + ADD COLUMN delim text NOT NULL DEFAULT ''; diff --git a/internal/repository/pg_mailbox.go b/internal/repository/pg_mailbox.go index fd16bd4b8..c7a58f6fd 100644 --- a/internal/repository/pg_mailbox.go +++ b/internal/repository/pg_mailbox.go @@ -41,19 +41,20 @@ func (r *mailboxRepository) CreateEntry(ctx context.Context, userId, emailId uui mb.UpdatedAt = time.Now() query := ` - INSERT INTO unibox_mailboxes (email_id, uid_validity, mailbox, attributes, highestmodseq, uid_next, updated_at) - VALUES ($1, $2, $3, $4, $5, $6, $7) + INSERT INTO unibox_mailboxes (email_id, uid_validity, mailbox, attributes, highestmodseq, uid_next, updated_at, delim) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8) ON CONFLICT (email_id, mailbox) DO UPDATE SET uid_validity = EXCLUDED.uid_validity, attributes = EXCLUDED.attributes, highestmodseq = EXCLUDED.highestmodseq, uid_next = EXCLUDED.uid_next, - updated_at = EXCLUDED.updated_at + updated_at = EXCLUDED.updated_at, + delim = EXCLUDED.delim ` // attributes is NOT NULL; a nil slice binds as SQL NULL. See textArray. _, err := r.db.Exec(ctx, query, - emailId, mb.UIDValidity, mb.Name, textArray(mb.Attrs), mb.HighestModSeq, mb.UIDNext, mb.UpdatedAt, + 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 @@ -68,14 +69,14 @@ func (r *mailboxRepository) CreateEntry(ctx context.Context, userId, emailId uui func (r *mailboxRepository) GetMailbox(ctx context.Context, userId, emailId uuid.UUID, name string) (*models.Mailbox, error) { query := ` - SELECT mailbox, attributes, uid_validity, highestmodseq, uid_next, updated_at + SELECT mailbox, attributes, uid_validity, highestmodseq, uid_next, updated_at, delim FROM unibox_mailboxes WHERE email_id = $1 AND mailbox = $2 ` var mb models.Mailbox err := r.db.QueryRow(ctx, query, emailId, name).Scan( - &mb.Name, &mb.Attrs, &mb.UIDValidity, &mb.HighestModSeq, &mb.UIDNext, &mb.UpdatedAt, + &mb.Name, &mb.Attrs, &mb.UIDValidity, &mb.HighestModSeq, &mb.UIDNext, &mb.UpdatedAt, &mb.Delim, ) if err != nil { if err == pgx.ErrNoRows { @@ -89,7 +90,7 @@ func (r *mailboxRepository) GetMailbox(ctx context.Context, userId, emailId uuid func (r *mailboxRepository) ListMailboxes(ctx context.Context, userId, emailId uuid.UUID) ([]models.Mailbox, error) { query := ` - SELECT mailbox, attributes, uid_validity, highestmodseq, uid_next, updated_at + SELECT mailbox, attributes, uid_validity, highestmodseq, uid_next, updated_at, delim FROM unibox_mailboxes WHERE email_id = $1 ` @@ -103,7 +104,7 @@ func (r *mailboxRepository) ListMailboxes(ctx context.Context, userId, emailId u var mailboxes []models.Mailbox for rows.Next() { var mb models.Mailbox - if err := rows.Scan(&mb.Name, &mb.Attrs, &mb.UIDValidity, &mb.HighestModSeq, &mb.UIDNext, &mb.UpdatedAt); err != nil { + if err := rows.Scan(&mb.Name, &mb.Attrs, &mb.UIDValidity, &mb.HighestModSeq, &mb.UIDNext, &mb.UpdatedAt, &mb.Delim); err != nil { return nil, err } mailboxes = append(mailboxes, mb) diff --git a/internal/repository/pg_unibox.go b/internal/repository/pg_unibox.go index a42c2043c..34a45fbf2 100644 --- a/internal/repository/pg_unibox.go +++ b/internal/repository/pg_unibox.go @@ -924,14 +924,39 @@ func (r *uniboxRepository) Delete(ctx context.Context, userID, id uuid.UUID) err // DeleteByFolderPaths removes the mirror rows for whole source folders. The // mail stays where it is at the provider; only the platform's copy goes. +// The message map entries go with the rows: the sync reads a mapped +// Message-ID as already stored, so an entry left behind would keep the +// message from ever being imported again if it moved back into a synced +// folder. An arrival still parked on warmup verification is dropped too, or +// it would surface into a folder nobody follows. func (r *uniboxRepository) DeleteByFolderPaths(ctx context.Context, emailID uuid.UUID, folderPaths []string) (int64, error) { if len(folderPaths) == 0 { return 0, nil } - tag, err := r.db.Exec(ctx, `DELETE FROM unibox_emails WHERE email_id = $1 AND folder_path = ANY($2)`, emailID, folderPaths) + tx, err := r.db.Begin(ctx) if err != nil { return 0, err } + defer func() { _ = tx.Rollback(ctx) }() + if _, err := tx.Exec(ctx, ` + DELETE FROM email_message_map m + USING unibox_emails u + WHERE u.email_id = $1 AND u.folder_path = ANY($2) + AND m.email_id = u.email_id AND m.message_id = u.message_id`, emailID, folderPaths); err != nil { + return 0, err + } + if _, err := tx.Exec(ctx, ` + DELETE FROM unibox_pending_emails + WHERE email_account_id = $1 AND payload->'message'->>'folder_path' = ANY($2)`, emailID, folderPaths); err != nil { + return 0, err + } + tag, err := tx.Exec(ctx, `DELETE FROM unibox_emails WHERE email_id = $1 AND folder_path = ANY($2)`, emailID, folderPaths) + if err != nil { + return 0, err + } + if err := tx.Commit(ctx); err != nil { + return 0, err + } return tag.RowsAffected(), nil }