From 4a6f0c17c7f16b60246838905deae6db61050e9b Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Tue, 22 Sep 2026 04:39:18 -0700 Subject: [PATCH 1/3] feat: advance each IMAP folder's sync cursor to the SELECT view its search ran against instead of the earlier LIST-STATUS, so Sent copies appended between the two are no longer skipped, and drop the message-map entry when NEW_EMAIL fails to publish so an unpublished message is re-offered instead of read as known (#645) --- internal/app/worker/wmail/admit.go | 10 +- internal/app/worker/wmail/imap_conn.go | 3 + internal/app/worker/wmail/sync_imap.go | 51 ++++++--- internal/app/worker/wmail/sync_imap_test.go | 114 ++++++++++++++++++++ internal/client/smtpimap/imap/client.go | 28 +++++ 5 files changed, 192 insertions(+), 14 deletions(-) diff --git a/internal/app/worker/wmail/admit.go b/internal/app/worker/wmail/admit.go index 375fefabc..339da35b1 100644 --- a/internal/app/worker/wmail/admit.go +++ b/internal/app/worker/wmail/admit.go @@ -181,11 +181,19 @@ func (w *WMail) storeNew(ctx context.Context, msg *models.EmailMessageData, data } // The consumer decodes NEW_EMAIL as JobEventNewEmail{user_id, message}. - return w.onEvent(models.JobEventTypeNewEmail, &models.JobEventNewEmail{ + err := w.onEvent(models.JobEventTypeNewEmail, &models.JobEventNewEmail{ UserID: w.UserID, Message: data, ReportOriginalMessageID: reportAbout, }) + if err != nil { + // The entry would mark a message that never reached the unibox as + // known, and every later pass would skip it; drop it so it is re-offered. + if derr := w.EmailMessageMapRepository.Del(ctx, w.UserID, w.ID, mapKey, data.ID); derr != nil { + log.Warn().Err(derr).Str("email_id", w.ID.String()).Msg("sync: map entry for an unpublished message not removed") + } + } + return err } // capBody bounds a stored body part at MaxEmailBodySize. IMAP already reads diff --git a/internal/app/worker/wmail/imap_conn.go b/internal/app/worker/wmail/imap_conn.go index e5fd8d4fa..49c00ce2c 100644 --- a/internal/app/worker/wmail/imap_conn.go +++ b/internal/app/worker/wmail/imap_conn.go @@ -26,6 +26,9 @@ type ImapConn interface { HasCondStore() bool ReleaseMailbox() SelectForSync(mailbox string) (uint32, *errx.MailError) + // SelectForSyncState selects like SelectForSync and reports the selected + // view's cursors, which an incremental pass advances the folder to. + SelectForSyncState(mailbox string) (imap.Selected, *errx.MailError) SearchChangedSince(modSeq uint64) ([]goimap.UID, *errx.MailError) SearchNewSince(uidNext uint32) ([]goimap.UID, *errx.MailError) // SearchAll is the folder's complete UID set, the presence side of the diff --git a/internal/app/worker/wmail/sync_imap.go b/internal/app/worker/wmail/sync_imap.go index b8e943f0c..4eae926f6 100644 --- a/internal/app/worker/wmail/sync_imap.go +++ b/internal/app/worker/wmail/sync_imap.go @@ -95,13 +95,15 @@ func (w *WMail) Sync(ctx context.Context) *errx.MailError { changed := imapFolderChanged(befBox, box, condStore) fullyProcessed := true var touched map[string]struct{} + var view imap.Selected if changed && !stats.aborted { w.setWalking(box) - done, ids, err := w.imapIncremental(ctx, box, befBox, condStore, stats) + done, sel, ids, err := w.imapIncremental(ctx, box, befBox, condStore, stats) if err != nil { return err } fullyProcessed = done + view = sel touched = ids } else if changed { // The pass was aborted before this folder; hold its cursor too. @@ -115,6 +117,8 @@ func (w *WMail) Sync(ctx context.Context) *errx.MailError { if !fullyProcessed { next.HighestModSeq = befBox.HighestModSeq next.UIDNext = befBox.UIDNext + } else if changed { + advanceToView(&next, view) } if err := w.mboxEvent(&next); err != nil { return nil @@ -213,16 +217,22 @@ func imapFolderChanged(before, now *models.Mailbox, condStore bool) bool { // imapIncremental stores what changed in one folder since the held cursor. // Known messages relay their flags unbudgeted; new ones are admitted newest // first. It reports whether every change was stored, which is what lets the -// folder's cursor advance, plus the Message-IDs it fetched so the drafts +// folder's cursor advance, the selected view the search ran against, which is +// where it advances to, and the Message-IDs it fetched so the drafts // reconciliation can tell a re-appended draft from an expunged one. -func (w *WMail) imapIncremental(ctx context.Context, box, before *models.Mailbox, condStore bool, stats *tickStats) (bool, map[string]struct{}, *errx.MailError) { +func (w *WMail) imapIncremental(ctx context.Context, box, before *models.Mailbox, condStore bool, stats *tickStats) (bool, imap.Selected, map[string]struct{}, *errx.MailError) { client := w.SmtpImapData.ImapClient - count, err := client.SelectForSync(box.Name) + view, err := client.SelectForSyncState(box.Name) if err != nil { - return false, nil, err + return false, view, nil, err } - if count == 0 { - return true, nil, nil + // The listing and this view name different generations, so the search + // would answer about UIDs the cursor does not; the next pass re-baselines. + if view.UIDValidity != 0 && view.UIDValidity != box.UIDValidity { + return false, view, nil, nil + } + if view.Count == 0 { + return true, view, nil, nil } var uids []goimap.UID if condStore { @@ -231,10 +241,10 @@ func (w *WMail) imapIncremental(ctx context.Context, box, before *models.Mailbox uids, err = client.SearchNewSince(before.UIDNext) } if err != nil { - return false, nil, err + return false, view, nil, err } if len(uids) == 0 { - return true, nil, nil + return true, view, nil, nil } // Newest first: when budget is short, the freshest mail lands first. sort.Slice(uids, func(i, j int) bool { return uids[i] > uids[j] }) @@ -250,11 +260,11 @@ func (w *WMail) imapIncremental(ctx context.Context, box, before *models.Mailbox hi := min(lo+config.ImapFetchBatchSize, len(uids)) fetched, err := client.FetchEnvelopes(ctx, uids[lo:hi]) if err != nil { - return false, touched, err + return false, view, touched, err } done, err := w.imapApply(ctx, fetched, false, stats) if err != nil { - return false, touched, err + return false, view, touched, err } for _, f := range fetched { if touched != nil { @@ -266,10 +276,25 @@ func (w *WMail) imapIncremental(ctx context.Context, box, before *models.Mailbox // unfrozen on a long backlog would deactivate itself walking mail it // cannot keep. Stop here; the held mod-sequence re-offers the rest. if !done || stats.aborted || ctx.Err() != nil { - return false, touched, nil + return false, view, touched, nil } } - return true, touched, nil + return true, view, touched, nil +} + +// advanceToView moves a fully walked folder's cursor to the view its search +// ran against. The listing's STATUS is taken before the SELECT, and a server +// whose selected view lags it (a session snapshot, an APPEND from the send +// path in between) would otherwise record a cursor past mail the search never +// returned, and that mail would never be synced. A view that reports no +// cursor keeps the listing's. +func advanceToView(next *models.Mailbox, view imap.Selected) { + if view.UIDNext != 0 { + next.UIDNext = view.UIDNext + } + if view.HighestModSeq != 0 { + next.HighestModSeq = view.HighestModSeq + } } // imapReconcileDrafts removes the platform's rows for drafts the server no diff --git a/internal/app/worker/wmail/sync_imap_test.go b/internal/app/worker/wmail/sync_imap_test.go index 6409ae46c..ffc3e6341 100644 --- a/internal/app/worker/wmail/sync_imap_test.go +++ b/internal/app/worker/wmail/sync_imap_test.go @@ -36,6 +36,9 @@ type fakeImapConn struct { // selectGen, when non-zero, is the UIDVALIDITY SELECT reports, which the // reconciliation compares against the one the listing gave it. selectGen uint32 + // view, when set, is the cursors SELECT reports in place of the listing's: + // a server whose selected view lags or leads its STATUS. + view *imap.Selected } func (c *fakeImapConn) Folders() ([]models.Mailbox, *errx.MailError) { return c.folders, nil } @@ -75,6 +78,19 @@ func (c *fakeImapConn) SelectForSync(string) (uint32, *errx.MailError) { return uint32(len(c.changed)), nil } +func (c *fakeImapConn) SelectForSyncState(name string) (imap.Selected, *errx.MailError) { + sel := imap.Selected{Count: uint32(len(c.changed))} + for _, f := range c.folders { + if f.Name == name { + sel.UIDValidity, sel.UIDNext, sel.HighestModSeq = f.UIDValidity, f.UIDNext, f.HighestModSeq + } + } + if c.view != nil { + sel.UIDNext, sel.HighestModSeq = c.view.UIDNext, c.view.HighestModSeq + } + return sel, nil +} + func (c *fakeImapConn) SelectForSyncGen(name string) (uint32, uint32, *errx.MailError) { gen := c.selectGen if gen == 0 { @@ -790,3 +806,101 @@ func TestImapSyncDoesNotReconcileExpungesOutsideDrafts(t *testing.T) { t.Errorf("asked the backend about INBOX %d times; only drafts are reconciled", ctx.calls) } } + +// The cursor is the selected view's, not the listing's. A server whose STATUS +// is ahead of the view the search ran on (an APPEND to Sent landing between +// the two) reported mail the search never returned, and advancing to the +// listing skipped that mail for good (#645). +func TestImapSyncAdvancesToTheSelectedViewNotTheListing(t *testing.T) { + conn := &fakeImapConn{ + folders: []models.Mailbox{{Name: "Sent", Attrs: []string{"\\Sent"}, UIDValidity: 7, HighestModSeq: 300}}, + changed: []goimap.UID{1, 2}, + view: &imap.Selected{HighestModSeq: 200}, + } + w, _ := newIMAPTestMail(conn, &fixedBudget{allow: 10}, + &models.Mailbox{Name: "Sent", Attrs: []string{"\\Sent"}, UIDValidity: 7, HighestModSeq: 100}) + + if err := w.Sync(t.Context()); err != nil { + t.Fatalf("Sync: %v", err) + } + if got := w.SmtpImapData.Mailboxes[0].HighestModSeq; got != 200 { + t.Fatalf("mod-sequence = %d, want the selected view's 200", got) + } + + // The view catches up; the folder still differs from the cursor, so the + // next pass walks it again and reaches the message the first one missed. + conn.view = nil + conn.changed = []goimap.UID{1, 2, 3} + if err := w.Sync(t.Context()); err != nil { + t.Fatalf("second Sync: %v", err) + } + if conn.fetches != 2 { + t.Errorf("fetched %d batches over two passes, want 2", conn.fetches) + } + if got := w.SmtpImapData.Mailboxes[0].HighestModSeq; got != 300 { + t.Errorf("mod-sequence = %d, want 300 once the view caught up", got) + } +} + +func TestImapSyncAdvancesUIDNextToTheSelectedView(t *testing.T) { + conn := &fakeImapConn{ + noCondStore: true, + folders: []models.Mailbox{{Name: "Sent", Attrs: []string{"\\Sent"}, UIDValidity: 7, UIDNext: 104}}, + changed: []goimap.UID{101, 102}, + view: &imap.Selected{UIDNext: 103}, + } + w, _ := newIMAPTestMail(conn, &fixedBudget{allow: 10}, + &models.Mailbox{Name: "Sent", Attrs: []string{"\\Sent"}, UIDValidity: 7, UIDNext: 101}) + + if err := w.Sync(t.Context()); err != nil { + t.Fatalf("Sync: %v", err) + } + if got := w.SmtpImapData.Mailboxes[0].UIDNext; got != 103 { + t.Errorf("UIDNEXT cursor = %d, want the selected view's 103", got) + } +} + +// recordingMessageMap remembers what was added and removed. +type recordingMessageMap struct { + fakeMessageMap + added, removed []string +} + +func (m *recordingMessageMap) Add(_ context.Context, d repository.EmailMessageData) error { + m.added = append(m.added, d.MessageID) + return nil +} + +func (m *recordingMessageMap) Del(_ context.Context, _, _ uuid.UUID, messageID string, _ uuid.UUID) error { + m.removed = append(m.removed, messageID) + return nil +} + +// A message whose NEW_EMAIL never reached the bus must not stay in the map, +// or every later pass reads it as known and it never reaches the unibox. +func TestStoreNewDropsTheMapEntryWhenThePublishFails(t *testing.T) { + conn := &fakeImapConn{ + folders: []models.Mailbox{{Name: "Sent", Attrs: []string{"\\Sent"}, UIDValidity: 7, HighestModSeq: 300}}, + changed: []goimap.UID{1}, + } + w, _ := newIMAPTestMail(conn, &fixedBudget{allow: 10}, + &models.Mailbox{Name: "Sent", Attrs: []string{"\\Sent"}, UIDValidity: 7, HighestModSeq: 100}) + m := &recordingMessageMap{} + w.EmailMessageMapRepository = m + w.onEvent = func(kind models.JobEventType, _ any) error { + if kind == models.JobEventTypeNewEmail { + return fmt.Errorf("bus unavailable") + } + return nil + } + + if err := w.Sync(t.Context()); err != nil { + t.Fatalf("Sync: %v", err) + } + if len(m.added) != 1 || len(m.removed) != 1 || m.removed[0] != m.added[0] { + t.Fatalf("added %v, removed %v: the unpublished message must leave the map", m.added, m.removed) + } + if got := w.SmtpImapData.Mailboxes[0].HighestModSeq; got != 100 { + t.Errorf("mod-sequence = %d, want the held 100", got) + } +} diff --git a/internal/client/smtpimap/imap/client.go b/internal/client/smtpimap/imap/client.go index 7dae930f0..59eac37df 100644 --- a/internal/client/smtpimap/imap/client.go +++ b/internal/client/smtpimap/imap/client.go @@ -334,6 +334,34 @@ func (c *Client) SelectForSyncGen(mailbox string) (uint32, uint32, *errx.MailErr return data.NumMessages, data.UIDValidity, nil } +// Selected is the view a SELECT opened: the count and the cursors every +// SEARCH on it answers against. +type Selected struct { + Count uint32 + UIDValidity uint32 + UIDNext uint32 + HighestModSeq uint64 +} + +// SelectForSyncState selects exactly as SelectForSync does and reports the +// cursors of the selected view, which is what a folder's stored cursor must +// advance to: a STATUS taken earlier can be ahead of the view the SEARCH ran on. +func (c *Client) SelectForSyncState(mailbox string) (Selected, *errx.MailError) { + c.lifecycle.RLock() + defer c.lifecycle.RUnlock() + defer c.begin()() + data, err := c.selectMailbox(mailbox, &imap.SelectOptions{ReadOnly: true, CondStore: c.condStore.Load()}) + if err != nil { + return Selected{}, c.handleError(err) + } + return Selected{ + Count: data.NumMessages, + UIDValidity: data.UIDValidity, + UIDNext: uint32(data.UIDNext), + HighestModSeq: data.HighestModSeq, + }, nil +} + // ReleaseMailbox drops the selected mailbox. Dovecot answers LIST-STATUS for // the selected mailbox with the values it held at SELECT, so a loop that keeps // INBOX selected never sees another change land. Servers without UNSELECT keep From 9009c6ec246751bb75a7046b7051e6ac03def2c8 Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Tue, 22 Sep 2026 04:49:26 -0700 Subject: [PATCH 2/3] feat: build inbound bounce and complaint reports before NEW_EMAIL but publish them only after it succeeds, so a message re-offered after a failed arrival publish does not apply its deliverability report twice --- internal/app/worker/wmail/admit.go | 25 +++++++++--- internal/app/worker/wmail/bounce.go | 42 ++++++++++----------- internal/app/worker/wmail/sync_imap_test.go | 33 ++++++++++++++++ 3 files changed, 71 insertions(+), 29 deletions(-) diff --git a/internal/app/worker/wmail/admit.go b/internal/app/worker/wmail/admit.go index 339da35b1..2f9765d88 100644 --- a/internal/app/worker/wmail/admit.go +++ b/internal/app/worker/wmail/admit.go @@ -173,11 +173,14 @@ func (w *WMail) storeNew(ctx context.Context, msg *models.EmailMessageData, data // Both still run on every message: a report is one or the other in // practice, and deciding that here by skipping the second would be this // function guessing at a MIME question the parsers already answer. - bounceAbout := w.maybeEmitBounce(msg) - complaintAbout := w.maybeEmitComplaint(msg) - reportAbout := bounceAbout - if reportAbout == "" { - reportAbout = complaintAbout + bounce := w.bounceReport(msg) + complaint := w.complaintReport(msg) + var reportAbout string + switch { + case bounce != nil: + reportAbout = bounce.OriginalMessageID + case complaint != nil: + reportAbout = complaint.OriginalMessageID } // The consumer decodes NEW_EMAIL as JobEventNewEmail{user_id, message}. @@ -192,8 +195,18 @@ func (w *WMail) storeNew(ctx context.Context, msg *models.EmailMessageData, data if derr := w.EmailMessageMapRepository.Del(ctx, w.UserID, w.ID, mapKey, data.ID); derr != nil { log.Warn().Err(derr).Str("email_id", w.ID.String()).Msg("sync: map entry for an unpublished message not removed") } + return err } - return err + + // Reports go out only once the arrival is published, so a message + // re-offered after a failed publish does not apply its report twice. + if bounce != nil { + _ = w.onEvent(models.JobEventTypeInboundBounce, bounce) + } + if complaint != nil { + _ = w.onEvent(models.JobEventTypeInboundComplaint, complaint) + } + return nil } // capBody bounds a stored body part at MaxEmailBodySize. IMAP already reads diff --git a/internal/app/worker/wmail/bounce.go b/internal/app/worker/wmail/bounce.go index 1df4a31a4..995bc7381 100644 --- a/internal/app/worker/wmail/bounce.go +++ b/internal/app/worker/wmail/bounce.go @@ -8,8 +8,8 @@ import ( "github.com/warmbly/warmbly/internal/pkg/dsn" ) -// maybeEmitBounce inspects a freshly synced inbound message and, when it is a -// permanent delivery-status notification for one of our sends, emits an +// bounceReport inspects a freshly synced inbound message and, when it is a +// permanent delivery-status notification for one of our sends, builds the // INBOUND_BOUNCE event so the consumer can suppress the recipient and record the // bounce against the campaign. This is where API-sent (Gmail/Graph) mail finally // gets bounce tracking: those sends succeed synchronously, so the only bounce @@ -20,19 +20,19 @@ import ( // Best-effort and permanent-only: a message that doesn't parse to a permanent // failure with a resolvable original id is silently ignored, so a transient // (4.x.x) bounce never suppresses a valid recipient. -// It returns the send the NDR is about, so the arrival event can carry it: the +// The event names the send the NDR is about, so the arrival event can carry it: the // report's own body is the only place that id appears and the consumer cannot // read a body, but it is what decides whether the report is about a campaign // send the customer should see or a warmup send they never made. -func (w *WMail) maybeEmitBounce(msg *models.EmailMessageData) string { +func (w *WMail) bounceReport(msg *models.EmailMessageData) *models.JobEventInboundBounce { from := strings.Join(msg.From, " ") if !dsn.Detect(from, msg.Subject, headerFlagValue(msg.Flags, "Content-Type")) { - return "" + return nil } report := dsn.Parse(msg.BodyPlain + "\n" + msg.BodyHTML) if !report.Permanent { - return "" + return nil } // Resolve the original outbound Message-ID: the DSN body's returned headers @@ -42,18 +42,16 @@ func (w *WMail) maybeEmitBounce(msg *models.EmailMessageData) string { originalID = strings.Trim(msg.InReplyTo[len(msg.InReplyTo)-1], "<>") } if originalID == "" { - return "" // nothing to resolve the campaign send against + return nil // nothing to resolve the campaign send against } - originalID = strings.Trim(originalID, "<>") - _ = w.onEvent(models.JobEventTypeInboundBounce, &models.JobEventInboundBounce{ + return &models.JobEventInboundBounce{ UserID: w.UserID, EmailID: w.ID, - OriginalMessageID: originalID, + OriginalMessageID: strings.Trim(originalID, "<>"), FailedRecipient: report.FailedRecipient, Reason: msg.Subject, - }) - return originalID + } } // headerFlagValue reads a "Header:value" pseudo-flag out of the flag slice (the @@ -68,7 +66,7 @@ func headerFlagValue(flags []string, name string) string { return "" } -// maybeEmitComplaint emits INBOUND_COMPLAINT when a synced message is an abuse +// complaintReport builds INBOUND_COMPLAINT when a synced message is an abuse // feedback report for one of our sends. A complaint never arrives // synchronously; it comes back as mail, long after the send succeeded. // @@ -76,26 +74,24 @@ func headerFlagValue(flags []string, name string) string { // Microsoft Graph returns one rendered body and no parts, so a report synced // through Graph carries only its human notice and is not detected; reading // those needs a separate MIME fetch, which is not built. -// Like maybeEmitBounce, it returns the send the report is about so the arrival -// event can carry it. -func (w *WMail) maybeEmitComplaint(msg *models.EmailMessageData) string { +// Like bounceReport, the event names the send the report is about so the +// arrival event can carry it. +func (w *WMail) complaintReport(msg *models.EmailMessageData) *models.JobEventInboundComplaint { from := strings.Join(msg.From, " ") if !arf.Detect(from, msg.Subject, headerFlagValue(msg.Flags, "Content-Type")) { - return "" + return nil } report := arf.Parse(msg.BodyPlain + "\n" + msg.BodyHTML) if !report.IsComplaint || report.OriginalMessageID == "" { - return "" + return nil } - originalID := strings.Trim(report.OriginalMessageID, "<>") - _ = w.onEvent(models.JobEventTypeInboundComplaint, &models.JobEventInboundComplaint{ + return &models.JobEventInboundComplaint{ UserID: w.UserID, EmailID: w.ID, - OriginalMessageID: originalID, + OriginalMessageID: strings.Trim(report.OriginalMessageID, "<>"), ComplainedRecipient: report.ComplainedRecipient, Provider: report.UserAgent, - }) - return originalID + } } diff --git a/internal/app/worker/wmail/sync_imap_test.go b/internal/app/worker/wmail/sync_imap_test.go index ffc3e6341..6babed3a1 100644 --- a/internal/app/worker/wmail/sync_imap_test.go +++ b/internal/app/worker/wmail/sync_imap_test.go @@ -904,3 +904,36 @@ func TestStoreNewDropsTheMapEntryWhenThePublishFails(t *testing.T) { t.Errorf("mod-sequence = %d, want the held 100", got) } } + +// A bounce is applied once: its report goes out only after the arrival is +// published, so a message re-offered after a failed publish does not +// suppress and count against its campaign twice. +func TestStoreNewEmitsReportsOnlyAfterTheArrivalIsPublished(t *testing.T) { + msg := &models.EmailMessageData{ + MessageID: "", + From: []string{"Mail Delivery System (MAILER-DAEMON@mail.example.com)"}, + Subject: "Undelivered Mail Returned to Sender", + BodyPlain: "Final-Recipient: rfc822; nobody@invalid.example.com\nAction: failed\nStatus: 5.1.1\n\nMessage-ID: \n", + } + for _, publishFails := range []bool{true, false} { + w, _ := newIMAPTestMail(&fakeImapConn{}, &fixedBudget{}, &models.Mailbox{}) + var kinds []models.JobEventType + w.onEvent = func(kind models.JobEventType, _ any) error { + kinds = append(kinds, kind) + if kind == models.JobEventTypeNewEmail && publishFails { + return fmt.Errorf("bus unavailable") + } + return nil + } + data := &models.EmailMessageStoreData{ID: uuid.New(), MessageID: msg.MessageID} + _ = w.storeNew(t.Context(), msg, data, msg.MessageID) + + want := []models.JobEventType{models.JobEventTypeNewEmail} + if !publishFails { + want = append(want, models.JobEventTypeInboundBounce) + } + if fmt.Sprint(kinds) != fmt.Sprint(want) { + t.Errorf("publishFails=%v: events %v, want %v", publishFails, kinds, want) + } + } +} From fb453921db1b6e6625ae1c45756e1b3cad04306b Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Tue, 22 Sep 2026 04:59:27 -0700 Subject: [PATCH 3/3] feat: keep a failed message-map rollback pending on the mailbox and retry it at the start of every IMAP, Gmail and Graph sync pass, skipping the pass until it succeeds so an unpublished arrival is never read as known --- internal/app/worker/wmail/admit.go | 21 ++++++++- internal/app/worker/wmail/sync_google.go | 3 ++ internal/app/worker/wmail/sync_graph.go | 3 ++ internal/app/worker/wmail/sync_imap.go | 3 ++ internal/app/worker/wmail/sync_imap_test.go | 52 +++++++++++++++++++++ internal/app/worker/wmail/wmail.go | 3 ++ 6 files changed, 84 insertions(+), 1 deletion(-) diff --git a/internal/app/worker/wmail/admit.go b/internal/app/worker/wmail/admit.go index 2f9765d88..b1cf25062 100644 --- a/internal/app/worker/wmail/admit.go +++ b/internal/app/worker/wmail/admit.go @@ -4,6 +4,7 @@ import ( "context" "time" + "github.com/google/uuid" "github.com/rs/zerolog/log" "github.com/warmbly/warmbly/internal/config" "github.com/warmbly/warmbly/internal/errx" @@ -193,7 +194,11 @@ func (w *WMail) storeNew(ctx context.Context, msg *models.EmailMessageData, data // The entry would mark a message that never reached the unibox as // known, and every later pass would skip it; drop it so it is re-offered. if derr := w.EmailMessageMapRepository.Del(ctx, w.UserID, w.ID, mapKey, data.ID); derr != nil { - log.Warn().Err(derr).Str("email_id", w.ID.String()).Msg("sync: map entry for an unpublished message not removed") + log.Warn().Err(derr).Str("email_id", w.ID.String()).Msg("sync: map entry for an unpublished message not removed; retried next pass") + if w.unmapPending == nil { + w.unmapPending = map[string]uuid.UUID{} + } + w.unmapPending[mapKey] = data.ID } return err } @@ -209,6 +214,20 @@ func (w *WMail) storeNew(ctx context.Context, msg *models.EmailMessageData, data return nil } +// retryUnmap removes the map entries a failed publish could not, and reports +// whether none is left. A pass must not run while one is: it would read that +// message as known and move its cursor past it. +func (w *WMail) retryUnmap(ctx context.Context) bool { + for key, id := range w.unmapPending { + if err := w.EmailMessageMapRepository.Del(ctx, w.UserID, w.ID, key, id); err != nil { + log.Warn().Err(err).Str("email_id", w.ID.String()).Msg("sync: map entry for an unpublished message still not removed; pass skipped") + return false + } + delete(w.unmapPending, key) + } + return true +} + // capBody bounds a stored body part at MaxEmailBodySize. IMAP already reads // at most that much off the wire; Gmail and Graph hand over whole bodies, so // without this an oversized message stored unbounded on those providers. diff --git a/internal/app/worker/wmail/sync_google.go b/internal/app/worker/wmail/sync_google.go index 075cb9bd1..08f3ef7d8 100644 --- a/internal/app/worker/wmail/sync_google.go +++ b/internal/app/worker/wmail/sync_google.go @@ -22,6 +22,9 @@ func (w *WMail) SyncGoogle(ctx context.Context) *errx.MailError { w.beginTick() stats := &tickStats{} w.googleTick = stats + if !w.retryUnmap(ctx) { + return nil + } newHistoryID, err := w.GoogleData.Client.FetchHistory(ctx, w.GoogleData.LastHistoryID) if newHistoryID != 0 && newHistoryID != w.GoogleData.LastHistoryID { diff --git a/internal/app/worker/wmail/sync_graph.go b/internal/app/worker/wmail/sync_graph.go index 5514232a8..d34558db7 100644 --- a/internal/app/worker/wmail/sync_graph.go +++ b/internal/app/worker/wmail/sync_graph.go @@ -25,6 +25,9 @@ func (w *WMail) SyncGraph(ctx context.Context) *errx.MailError { w.beginTick() stats := &tickStats{} w.graphTick = stats + if !w.retryUnmap(ctx) { + return nil + } if err := w.GraphData.Client.Sync(ctx); err != nil { var mailErr *errx.MailError diff --git a/internal/app/worker/wmail/sync_imap.go b/internal/app/worker/wmail/sync_imap.go index 4eae926f6..20b90bd6e 100644 --- a/internal/app/worker/wmail/sync_imap.go +++ b/internal/app/worker/wmail/sync_imap.go @@ -31,6 +31,9 @@ func (w *WMail) Sync(ctx context.Context) *errx.MailError { } w.beginTick() stats := &tickStats{} + if !w.retryUnmap(ctx) { + return nil + } client := w.SmtpImapData.ImapClient // A mailbox left selected by the previous pass freezes LIST-STATUS on this diff --git a/internal/app/worker/wmail/sync_imap_test.go b/internal/app/worker/wmail/sync_imap_test.go index 6babed3a1..89a1c76ae 100644 --- a/internal/app/worker/wmail/sync_imap_test.go +++ b/internal/app/worker/wmail/sync_imap_test.go @@ -864,6 +864,8 @@ func TestImapSyncAdvancesUIDNextToTheSelectedView(t *testing.T) { type recordingMessageMap struct { fakeMessageMap added, removed []string + // failDel is how many removals fail before they start succeeding. + failDel int } func (m *recordingMessageMap) Add(_ context.Context, d repository.EmailMessageData) error { @@ -872,6 +874,10 @@ func (m *recordingMessageMap) Add(_ context.Context, d repository.EmailMessageDa } func (m *recordingMessageMap) Del(_ context.Context, _, _ uuid.UUID, messageID string, _ uuid.UUID) error { + if m.failDel > 0 { + m.failDel-- + return fmt.Errorf("backend unavailable") + } m.removed = append(m.removed, messageID) return nil } @@ -937,3 +943,49 @@ func TestStoreNewEmitsReportsOnlyAfterTheArrivalIsPublished(t *testing.T) { } } } + +// A removal that fails too is retried before the next pass looks at anything, +// and the pass waits for it, so the message is never read as known. +func TestAFailedUnmapIsRetriedBeforeThePassRuns(t *testing.T) { + conn := &fakeImapConn{ + folders: []models.Mailbox{{Name: "Sent", Attrs: []string{"\\Sent"}, UIDValidity: 7, HighestModSeq: 300}}, + changed: []goimap.UID{1}, + } + w, _ := newIMAPTestMail(conn, &fixedBudget{allow: 10}, + &models.Mailbox{Name: "Sent", Attrs: []string{"\\Sent"}, UIDValidity: 7, HighestModSeq: 100}) + m := &recordingMessageMap{failDel: 2} + w.EmailMessageMapRepository = m + publishFails := true + arrivals := 0 + w.onEvent = func(kind models.JobEventType, _ any) error { + if kind == models.JobEventTypeNewEmail { + arrivals++ + if publishFails { + return fmt.Errorf("bus unavailable") + } + } + return nil + } + + // Publish fails and so does the rollback; the next pass cannot remove it + // either, so it runs nothing. + for range 2 { + if err := w.Sync(t.Context()); err != nil { + t.Fatalf("Sync: %v", err) + } + } + if arrivals != 1 || conn.fetches != 1 { + t.Fatalf("arrivals %d, fetches %d: a pass ran with the entry still mapped", arrivals, conn.fetches) + } + + publishFails = false + if err := w.Sync(t.Context()); err != nil { + t.Fatalf("Sync: %v", err) + } + if len(m.removed) != 1 || arrivals != 2 { + t.Errorf("removed %v, arrivals %d: the message was not re-offered once the entry was gone", m.removed, arrivals) + } + if got := w.SmtpImapData.Mailboxes[0].HighestModSeq; got != 300 { + t.Errorf("mod-sequence = %d, want 300", got) + } +} diff --git a/internal/app/worker/wmail/wmail.go b/internal/app/worker/wmail/wmail.go index c1d3c9212..85ea745c5 100644 --- a/internal/app/worker/wmail/wmail.go +++ b/internal/app/worker/wmail/wmail.go @@ -98,6 +98,9 @@ type WMail struct { // flagScan is the previous flag snapshot per folder name, used only on // IMAP servers without CONDSTORE, which cannot say what changed. flagScan map[string]*folderFlagScan + // 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 // transportFailures counts consecutive passes that could not reach the // mail server, which paces the retry and keeps one outage to one warning. transportFailures int