feat: make every skipped-folder purge drop the message map entries and parked arrivals with the rows so a message can be imported again if it moves back, persist each folder's delimiter (migration 000197) so the backend purge matches the same subfolders and case the worker skips, fail the mailbox load closed when the skip list cannot be read, mark a folder renamed into the skipped subtree as skipped, never remove a row re-fetched this pass, cap the reconciliation at 50 searches per pass and continue next pass, publish one EMAIL_DELETED after a folder purge, and resolve the account once for GET /emails/:id/sync

This commit is contained in:
Matthew Meszaros
2026-09-22 05:33:56 -07:00
parent ad353b385a
commit ee9f5bf84e
11 changed files with 322 additions and 58 deletions
+1 -6
View File
@@ -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
@@ -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
}
+16 -7
View File
@@ -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)
+41 -27
View File
@@ -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")
+70 -9
View File
@@ -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))
}
@@ -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)
}
}
+5
View File
@@ -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
@@ -0,0 +1,2 @@
ALTER TABLE public.unibox_mailboxes
DROP COLUMN IF EXISTS delim;
@@ -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 '';
+9 -8
View File
@@ -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)
+26 -1
View File
@@ -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
}