feat: re-decide a withdrawn warmup tampering hold on the strikes left in the seven days before it was imposed (on the pool row, or the reputation ledger for a mailbox out of every pool), stamp an old strike verified only when a worker answers its search and ask again after six hours, retry a removal check for a mailbox still loading, redeliver a failed strike, search IMAP folders in one session hold, and share the Message-ID search helpers

This commit is contained in:
Matthew Meszaros
2026-09-29 04:57:45 -07:00
parent c5a981d86c
commit e5eb3b0ff1
20 changed files with 581 additions and 317 deletions
+1 -1
View File
@@ -770,7 +770,7 @@ Signals used:
- no inbound token is evidence against the mailbox that received it. It did not present the token; its worker synced whatever landed in its inbox, and inbound mail is attacker-controlled: every pool member holds tokens naming itself and a partner, and forwarding three to another member used to block that member for 30 days. The recipient check already makes a token worthless anywhere but its own destination, so nothing is charged on that path (#468, #481). Do not reintroduce a charge there, whether gated by a window, a folder check, a clock or by which pair the token names; each of those was tried and each was a way to be wrong (#477, #480)
- tampering with warmup mail a mailbox verifiably received (deleting it, flagging it as spam) is attributed to that mailbox, because only its owner can do it. It is a ladder, not a first-strike ban: `evaluateMetrics` weighs a deletion as one strike and a spam flag as two over the seven-day window, and warns at one, quarantines at two and blocks at four (`tampering*Strikes` in `internal/app/warmup/service.go`). One deletion is someone tidying the folder by hand until proven otherwise (#635). `RecordTampering` only records the event and re-evaluates, so a sweep reaches the same answer; a tampering block carries a term like every other band and never requires review
- a deletion is a strike only inside `config.WarmupDeletionStrikeHours` of arrival (`warmupDeletionCounts` in `internal/app/consumer/event_remove_email.go`), and never for a receipt the retention sweep has retired. The engagement a message earns happens in its first hours; after that the platform deletes it itself (#637), so a later removal, whichever of the owner, Gmail's Trash purge, a server retention rule or our own sweep did it, is housekeeping. Gmail's Delete arrives as the `TRASH` label and is judged there on the same rule, because the `messagesDeleted` history record only comes when Trash is emptied, weeks later and in a burst. Do not widen the window or count a removal past it: every mailbox on a fixed quota has to be able to clear the folder
- a removal is never a strike on its own. Graph reports a move exactly like a delete, and a provider's filter, a mailbox rule, another Warmbly instance syncing the same mailbox (a self-hosted instance warming in Warmbly Cloud) or our own filing can all move warmup mail. A fresh removal publishes a `verify_removal` warmup action; the worker searches the whole mailbox by Message-ID (`internal/app/worker/event_warmup_verify.go`) and answers `WARMUP_REMOVAL_CHECKED`, and `HandleWarmupRemovalChecked` (`internal/app/consumer/warmup_removal_check.go`) records a strike only for a message in the trash or gone. Found anywhere else withdraws any strike for it (`WithdrawTampering`), lifting a pause or block that strike imposed. A search that fails or cannot tell (IMAP only sees synced folders) charges nothing. `warmup_tampering_events.verified_at` is NULL only on strikes from before the search, which `StartWarmupTamperingRecheck` searches once. Gmail's `TRASH` label needs no search, since it is the message entering the trash. Do not add another path that records a deletion without the search; the self-move marker is only a shortcut that skips it for our own filing
- a removal is never a strike on its own. Graph reports a move exactly like a delete, and a provider's filter, a mailbox rule, another Warmbly instance syncing the same mailbox (a self-hosted instance warming in Warmbly Cloud) or our own filing can all move warmup mail. A fresh removal publishes a `verify_removal` warmup action; the worker searches the whole mailbox by Message-ID (`internal/app/worker/event_warmup_verify.go`) and answers `WARMUP_REMOVAL_CHECKED`, and `HandleWarmupRemovalChecked` (`internal/app/consumer/warmup_removal_check.go`) records a strike only for a message in the trash or gone. Found anywhere else withdraws any strike for it (`WithdrawTampering`), and a tampering pause or block is re-decided on the strikes left in the seven days before it was imposed (never on today's window, which old strikes have aged out of), then lowered, shortened from its original decision time, or lifted. The revision lands on the pool row, or on the ledger when the mailbox is in no pool, the one write to `warmup_reputation_ledger` outside its trigger, so a withdrawn hold is not seeded back on rejoin. A search that fails or cannot tell (IMAP only sees synced folders) charges nothing. `warmup_tampering_events.verified_at` is NULL only on strikes from before the search; `StartWarmupTamperingRecheck` asks for those, stamps `verify_requested_at`, asks again after six hours while unanswered, and the answer stamps `verified_at`. Gmail's `TRASH` label needs no search, since it is the message entering the trash. Do not add another path that records a deletion without the search; the self-move marker is only a shortcut that skips it for our own filing
- warmup mail is retained by the platform, not the owner: `StartWarmupMailRetention` (`internal/app/consumer/warmup_mail_retention.go`) retires every received copy and every sender's copy past the mailbox's window (`email_accounts.warmup_retention_days`, else `retention.warmup_mail_days`) and publishes `WarmupActionDelete` to the worker, which trashes it on Gmail, deletes it on Graph, expunges it on IMAP and drops the stored body. The row is retired only after the action is on the bus, so a failed publish is re-offered. The same loop prunes tokens, receipts, tampering events and spam reports past `retention.warmup_event_days`; `warmup_statistics` carries the analytics and is never pruned
- accounts can be auto-blocked from warmup pools
+1 -1
View File
@@ -237,7 +237,7 @@ What a mailbox does to warmup mail it received is held against it, on a ladder r
Only a fresh deletion counts, because that is the one that costs the pool something: the engagement a warmup message earns happens in its first hours, and removing it before then takes that signal away. A warmup message deleted later is housekeeping, whether by you, by Gmail emptying its Trash, by a retention rule on your mail server or by Warmbly's own [retention](#retention), and it is never held against the mailbox. On Gmail, pressing Delete is what is judged, not the purge from Trash weeks later. A mailbox's own filing is never counted either: Warmbly moving a warmup email into its folder, marking it read, rescuing it from spam or deleting it once its window has passed is the platform acting, not the owner.
Moving a warmup email is not deleting it. Outlook and Microsoft 365 report a message moved to another folder exactly as they report one deleted, and a mailbox's own rules, its provider's filter or a second Warmbly instance syncing the same mailbox can all move mail. So before a removal counts, Warmbly searches the whole mailbox for the message: found in any folder other than Trash (Deleted Items on Outlook), it was filed, and nothing is held against the mailbox. Only a message in Trash or gone for good is a strike. A search that cannot run, or cannot tell, counts as nothing. Deletion strikes recorded before this search existed are searched for once in the same way, and a strike whose message is still in the mailbox is withdrawn, lifting the pause or block it caused.
Moving a warmup email is not deleting it. Outlook and Microsoft 365 report a message moved to another folder exactly as they report one deleted, and a mailbox's own rules, its provider's filter or a second Warmbly instance syncing the same mailbox can all move mail. So before a removal counts, Warmbly searches the whole mailbox for the message: found in any folder other than Trash (Deleted Items on Outlook), it was filed, and nothing is held against the mailbox. Only a message in Trash or gone for good is a strike. A search that cannot run, or cannot tell, counts as nothing. Deletion strikes recorded before this search existed are searched for in the same way. A strike whose message is still in the mailbox is withdrawn, and the pause or block it caused is decided again on the strikes that remain, so it is lifted or shortened only when those no longer earn it.
<Callout type="warn" title="Acting early protects everyone">
Warmbly intervenes well before providers would penalize a mailbox. Complaints, bounces and tampering take a mailbox out of the pool early. A mailbox landing in spam at the major providers is slowed down instead, so it keeps warming, which is how it recovers.
+1 -1
View File
@@ -42,7 +42,7 @@ func (s *JobsService) HandleRemoveEmail(ctx context.Context, e *models.JobEventR
Str("message_id", rec.MessageID).
Msg("Warmup message removed after its engagement window; housekeeping, not tampering")
default:
checkErr = s.checkWarmupRemoval(ctx, e.UserID, e.EmailID, rec.MessageID, &rec.InternalID)
checkErr = s.checkWarmupRemoval(ctx, e.UserID, e.EmailID, rec.MessageID)
}
}
}
+29 -52
View File
@@ -14,16 +14,9 @@ import (
"github.com/warmbly/warmbly/internal/models"
)
// A removal the sync reports is only a message leaving a folder. Graph reports
// every move that way, and a message filed elsewhere by a provider's filter, a
// mailbox rule, another Warmbly instance syncing the same mailbox or our own
// filing is still in the mailbox. So a fresh removal of warmup mail is never a
// strike on its own: the worker searches the mailbox for the message and a
// strike is recorded only when it is in the trash or gone.
// checkWarmupRemoval asks the worker holding the mailbox where a removed
// warmup message went. A mailbox with nowhere to ask is not charged.
func (s *JobsService) checkWarmupRemoval(ctx context.Context, userID, accountID uuid.UUID, rfcMessageID string, internalID *uuid.UUID) error {
// warmup message went; a removal is only a message leaving a folder.
func (s *JobsService) checkWarmupRemoval(ctx context.Context, userID, accountID uuid.UUID, rfcMessageID string) error {
if s.Publisher == nil || s.EmailRepository == nil || rfcMessageID == "" {
return nil
}
@@ -38,101 +31,85 @@ func (s *JobsService) checkWarmupRemoval(ctx context.Context, userID, accountID
log.Info().Str("email_id", accountID.String()).Msg("Warmup removal not checked: mailbox has no worker; nothing charged")
return nil
}
return s.publishRemovalCheck(ctx, *account.WorkerID, userID, accountID, rfcMessageID, internalID)
return s.publishRemovalCheck(ctx, *account.WorkerID, userID, accountID, rfcMessageID)
}
func (s *JobsService) publishRemovalCheck(ctx context.Context, workerID, userID, accountID uuid.UUID, rfcMessageID string, internalID *uuid.UUID) error {
action := &models.WarmupEmailAction{
func (s *JobsService) publishRemovalCheck(ctx context.Context, workerID, userID, accountID uuid.UUID, rfcMessageID string) error {
if err := s.Publisher.PublishWarmupAction(ctx, workerID, &models.WarmupEmailAction{
UserID: userID,
EmailID: accountID,
RFCMessageID: rfcMessageID,
Actions: []string{models.WarmupActionVerifyRemoval},
}
if internalID != nil && *internalID != uuid.Nil {
action.InternalID = internalID.String()
}
if err := s.Publisher.PublishWarmupAction(ctx, workerID, action); err != nil {
}); err != nil {
return fmt.Errorf("warmup removal check: publish: %w", err)
}
return nil
}
// HandleWarmupRemovalChecked judges a removal on where the worker found the
// message. Found anywhere outside the trash withdraws any strike for it,
// which is also how a strike recorded before this check is corrected.
// message: outside the trash withdraws any strike for it, in the trash or
// nowhere records one.
func (s *JobsService) HandleWarmupRemovalChecked(ctx context.Context, e *models.JobEventWarmupRemovalChecked) error {
if s.WarmupService == nil || e == nil || e.RFCMessageID == "" {
return nil
}
withdraw := e.Outcome == models.WarmupRemovalPresent
var health *models.WarmupParticipantHealth
var xerr *errx.Error
switch e.Outcome {
case models.WarmupRemovalPresent:
health, xerr = s.WarmupService.WithdrawTampering(ctx, e.EmailID, e.RFCMessageID, "deletion")
case models.WarmupRemovalTrashed, models.WarmupRemovalGone:
// The retention sweep deletes warmup mail itself; a message it has
// retired being gone says nothing about the owner.
withdraw = s.retiredByPlatform(ctx, e)
health, xerr = s.WarmupService.RecordTampering(ctx, e.EmailID, e.RFCMessageID, "deletion")
if xerr == nil && s.WarmupRepo != nil {
// Confirms a strike recorded before the search existed.
if err := s.WarmupRepo.MarkTamperingVerified(ctx, e.EmailID, e.RFCMessageID, "deletion"); err != nil {
return fmt.Errorf("confirm warmup strike: %w", err)
}
}
default:
return nil
}
if withdraw {
health, xerr := s.WarmupService.WithdrawTampering(ctx, e.EmailID, e.RFCMessageID, "deletion")
if xerr != nil {
return fmt.Errorf("withdraw warmup strike: %w", xerr)
}
s.markRiskBandFromWarmupHealth(ctx, e.EmailID, health)
return nil
if xerr != nil {
return fmt.Errorf("judge warmup removal: %w", xerr)
}
health, _ := s.WarmupService.RecordTampering(ctx, e.EmailID, e.RFCMessageID, "deletion")
s.markRiskBandFromWarmupHealth(ctx, e.EmailID, health)
return nil
}
func (s *JobsService) retiredByPlatform(ctx context.Context, e *models.JobEventWarmupRemovalChecked) bool {
internalID, err := uuid.Parse(e.InternalID)
if err != nil || s.WarmupRepo == nil {
return false
}
rec, _ := s.WarmupRepo.GetWarmupReceived(ctx, e.EmailID, internalID)
return rec != nil && rec.RetiredAt != nil
}
const (
warmupTamperingRecheckBatch = 100
// warmupTamperingRecheckWindow is the seven days a strike counts plus the
// thirty-day block it can lead to; an older strike decides nothing.
// The seven days a strike counts plus the thirty-day block it can impose.
warmupTamperingRecheckWindow = 37 * 24 * time.Hour
// A search nobody answered is asked for again after this.
warmupTamperingRecheckRetry = 6 * time.Hour
)
// StartWarmupTamperingRecheck searches the mailbox once for every deletion
// strike recorded before removals were checked, so a strike for a message
// that was only moved is withdrawn and the hold it caused is lifted.
// StartWarmupTamperingRecheck searches for every deletion strike recorded
// before removals were checked, until each is answered or ages out.
func (s *JobsService) StartWarmupTamperingRecheck(ctx context.Context) {
if s.WarmupRepo == nil || s.Publisher == nil {
return
}
jobrun.Loop(ctx, "warmup_tampering_recheck", 10*time.Minute, true, func(ctx context.Context) error {
jobrun.Loop(ctx, "warmup_tampering_recheck", time.Hour, true, func(ctx context.Context) error {
batchCtx, cancel := context.WithTimeout(ctx, 45*time.Second)
defer cancel()
return s.recheckTamperingBatch(batchCtx)
})
}
// recheckTamperingBatch stamps a strike once its search is on the bus; a
// search that never answers leaves the strike as it was.
func (s *JobsService) recheckTamperingBatch(ctx context.Context) error {
rows, err := s.WarmupRepo.ListUnverifiedDeletions(ctx, time.Now().Add(-warmupTamperingRecheckWindow), warmupTamperingRecheckBatch)
rows, err := s.WarmupRepo.ListUnverifiedDeletions(ctx, time.Now().Add(-warmupTamperingRecheckWindow), warmupTamperingRecheckRetry, warmupTamperingRecheckBatch)
if err != nil {
return err
}
var failures []error
for i := range rows {
r := &rows[i]
if err := s.publishRemovalCheck(ctx, r.WorkerID, r.UserID, r.EmailAccountID, r.MessageID, r.InternalID); err != nil {
if err := s.publishRemovalCheck(ctx, r.WorkerID, r.UserID, r.EmailAccountID, r.MessageID); err != nil {
failures = append(failures, err)
continue
}
if err := s.WarmupRepo.MarkTamperingVerified(ctx, r.EmailAccountID, r.MessageID, "deletion"); err != nil {
if err := s.WarmupRepo.MarkTamperingVerifyRequested(ctx, r.EmailAccountID, r.MessageID); err != nil {
failures = append(failures, err)
}
}
+64 -72
View File
@@ -30,9 +30,13 @@ func (r retentionWarmupRepo) GetWarmupReceived(context.Context, uuid.UUID, uuid.
type retentionWarmupService struct {
warmupapp.Service
strikes []string
fail bool
}
func (s *retentionWarmupService) RecordTampering(_ context.Context, _ uuid.UUID, _, kind string) (*models.WarmupParticipantHealth, *errx.Error) {
if s.fail {
return nil, errx.InternalError()
}
s.strikes = append(s.strikes, kind)
return nil, nil
}
@@ -124,7 +128,7 @@ func TestRemoveEmailChecksOnlyFreshWarmupMail(t *testing.T) {
if tc.checks == 1 {
a := pub.actions[0]
if len(a.Actions) != 1 || a.Actions[0] != models.WarmupActionVerifyRemoval ||
a.RFCMessageID != tc.rec.MessageID || a.InternalID != tc.rec.InternalID.String() || pub.workers[0] != worker {
a.RFCMessageID != tc.rec.MessageID || pub.workers[0] != worker {
t.Fatalf("check carried %+v to %v", a, pub.workers[0])
}
}
@@ -148,33 +152,26 @@ func TestRemoveEmailWithNowhereToCheckChargesNothing(t *testing.T) {
}
// The strike follows where the worker found the message: anywhere outside the
// trash withdraws it, the trash or nowhere records it, and a message the
// retention sweep has retired is the platform's deletion whatever the search
// says.
// trash withdraws it, the trash or nowhere records and confirms it.
func TestRemovalCheckedJudgesOnWhereTheMessageIs(t *testing.T) {
retired := time.Now().Add(-time.Minute)
cases := []struct {
name string
outcome string
retired bool
want []string
name string
outcome string
want []string
verified int
}{
{"moved to another folder", models.WarmupRemovalPresent, false, []string{"withdraw:deletion"}},
{"in the trash", models.WarmupRemovalTrashed, false, []string{"deletion"}},
{"gone for good", models.WarmupRemovalGone, false, []string{"deletion"}},
{"gone after the platform retired it", models.WarmupRemovalGone, true, []string{"withdraw:deletion"}},
{"an answer this consumer does not know", "sideways", false, nil},
{"moved to another folder", models.WarmupRemovalPresent, []string{"withdraw:deletion"}, 0},
{"in the trash", models.WarmupRemovalTrashed, []string{"deletion"}, 1},
{"gone for good", models.WarmupRemovalGone, []string{"deletion"}, 1},
{"an answer this consumer does not know", "sideways", nil, 0},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
rec := receivedAgo(time.Hour)
if tc.retired {
rec.RetiredAt = &retired
}
s, svc := retentionService(rec)
s, svc := retentionService(receivedAgo(time.Hour))
repo := &verifiedRepo{}
s.WarmupRepo = repo
if err := s.HandleWarmupRemovalChecked(context.Background(), &models.JobEventWarmupRemovalChecked{
UserID: uuid.New(), EmailID: uuid.New(), InternalID: rec.InternalID.String(),
RFCMessageID: rec.MessageID, Outcome: tc.outcome,
UserID: uuid.New(), EmailID: uuid.New(), RFCMessageID: "<m@example.test>", Outcome: tc.outcome,
}); err != nil {
t.Fatal(err)
}
@@ -186,33 +183,57 @@ func TestRemovalCheckedJudgesOnWhereTheMessageIs(t *testing.T) {
t.Fatalf("strikes = %v, want %v", svc.strikes, tc.want)
}
}
if repo.verified != tc.verified {
t.Fatalf("confirmed %d strikes, want %d", repo.verified, tc.verified)
}
})
}
}
// unverifiedRepo serves one listing of strikes recorded before removals were
// checked and records which were stamped.
type unverifiedRepo struct {
// A strike the service could not record is redelivered, not acked.
func TestRemovalCheckedRedeliversAFailedStrike(t *testing.T) {
s, svc := retentionService(nil)
svc.fail = true
if err := s.HandleWarmupRemovalChecked(context.Background(), &models.JobEventWarmupRemovalChecked{
UserID: uuid.New(), EmailID: uuid.New(), RFCMessageID: "<m@example.test>", Outcome: models.WarmupRemovalTrashed,
}); err == nil {
t.Fatal("a failed strike was acked")
}
}
// verifiedRepo counts the strikes a search confirmed.
type verifiedRepo struct {
repository.WarmupRepository
rows []repository.WarmupTamperingToVerify
verified []string
verified int
}
func (r *unverifiedRepo) ListUnverifiedDeletions(context.Context, time.Time, int) ([]repository.WarmupTamperingToVerify, error) {
return r.rows, nil
}
func (r *unverifiedRepo) MarkTamperingVerified(_ context.Context, _ uuid.UUID, messageID, _ string) error {
r.verified = append(r.verified, messageID)
func (r *verifiedRepo) MarkTamperingVerified(context.Context, uuid.UUID, string, string) error {
r.verified++
return nil
}
// Every old strike is searched for once, and stamped only once its search is
// on the bus, so a publish that fails is offered again.
func TestRecheckTamperingSearchesEachOldStrikeOnce(t *testing.T) {
internal := uuid.New()
// unverifiedRepo serves one listing of strikes recorded before removals were
// checked and records which searches were asked for.
type unverifiedRepo struct {
repository.WarmupRepository
rows []repository.WarmupTamperingToVerify
requested []string
}
func (r *unverifiedRepo) ListUnverifiedDeletions(context.Context, time.Time, time.Duration, int) ([]repository.WarmupTamperingToVerify, error) {
return r.rows, nil
}
func (r *unverifiedRepo) MarkTamperingVerifyRequested(_ context.Context, _ uuid.UUID, messageID string) error {
r.requested = append(r.requested, messageID)
return nil
}
// Every old strike is searched for, and marked asked only once its search is
// on the bus, so a publish that fails is offered again next pass.
func TestRecheckTamperingSearchesEachOldStrike(t *testing.T) {
rows := []repository.WarmupTamperingToVerify{
{EmailAccountID: uuid.New(), UserID: uuid.New(), WorkerID: uuid.New(), MessageID: "<a@example.test>", InternalID: &internal},
{EmailAccountID: uuid.New(), UserID: uuid.New(), WorkerID: uuid.New(), MessageID: "<a@example.test>"},
{EmailAccountID: uuid.New(), UserID: uuid.New(), WorkerID: uuid.New(), MessageID: "<b@example.test>"},
}
repo := &unverifiedRepo{rows: rows}
@@ -221,10 +242,11 @@ func TestRecheckTamperingSearchesEachOldStrikeOnce(t *testing.T) {
if err := s.recheckTamperingBatch(context.Background()); err != nil {
t.Fatal(err)
}
if len(pub.actions) != 2 || len(repo.verified) != 2 {
t.Fatalf("published %d, stamped %d; want 2 and 2", len(pub.actions), len(repo.verified))
if len(pub.actions) != 2 || len(repo.requested) != 2 {
t.Fatalf("published %d, marked %d; want 2 and 2", len(pub.actions), len(repo.requested))
}
if pub.actions[0].InternalID != internal.String() || pub.actions[1].InternalID != "" || pub.workers[1] != rows[1].WorkerID {
if pub.actions[1].RFCMessageID != rows[1].MessageID || pub.workers[1] != rows[1].WorkerID ||
pub.actions[1].Actions[0] != models.WarmupActionVerifyRemoval {
t.Fatalf("checks carried %+v to %v", pub.actions, pub.workers)
}
@@ -233,38 +255,8 @@ func TestRecheckTamperingSearchesEachOldStrikeOnce(t *testing.T) {
if err := s.recheckTamperingBatch(context.Background()); err == nil {
t.Fatal("a failed publish should be reported")
}
if len(failing.verified) != 0 {
t.Fatalf("stamped %v without a search on the bus", failing.verified)
}
}
func TestRemoveEmailNeverStrikesARetiredMessage(t *testing.T) {
rec := receivedAgo(time.Hour)
retired := time.Now().Add(-time.Minute)
rec.RetiredAt = &retired
s, svc := retentionService(rec)
if err := s.HandleRemoveEmail(context.Background(), &models.JobEventRemoveEmail{
UserID: uuid.New(), EmailID: uuid.New(), ID: uuid.New(),
}); err != nil {
t.Fatal(err)
}
if len(svc.strikes) != 0 {
t.Fatalf("the platform's own deletion was recorded as tampering: %v", svc.strikes)
}
}
// A fresh warmup message that the sync found in a folder the owner excluded
// from sync was filed, not deleted: it is still in the mailbox, so it is not
// a strike whatever its age.
func TestRemoveEmailNeverStrikesAMessageFiledIntoASkippedFolder(t *testing.T) {
s, svc := retentionService(receivedAgo(time.Hour))
if err := s.HandleRemoveEmail(context.Background(), &models.JobEventRemoveEmail{
UserID: uuid.New(), EmailID: uuid.New(), ID: uuid.New(), SkippedFolder: "Warmer",
}); err != nil {
t.Fatal(err)
}
if len(svc.strikes) != 0 {
t.Fatalf("a move into a skipped folder was recorded as tampering: %v", svc.strikes)
if len(failing.requested) != 0 {
t.Fatalf("marked %v asked without a search on the bus", failing.requested)
}
}
+59 -27
View File
@@ -371,41 +371,73 @@ func (s *service) WithdrawTampering(ctx context.Context, accountID uuid.UUID, me
if err != nil {
return nil, errx.InternalError()
}
participant, xerr := s.getParticipantForAnyPool(ctx, accountID)
if xerr != nil || participant == nil || !removed {
return participant, xerr
// Run even when the strike is already gone, so a retry finishes a revision
// an earlier attempt did not.
revised, xerr := s.reviseTamperingHold(ctx, accountID)
if xerr != nil {
return nil, xerr
}
// A hold is a sentence the bands never lower, so the one this strike
// imposed is lifted here, and only when what remains decides less.
if reason := tamperingHoldReason(participant); reason != "" {
metrics, err := s.loadMetrics(ctx, accountID, participant)
if err != nil {
return nil, errx.InternalError()
}
decision := evaluateMetrics(metrics, placementPrior(participant), s.now().UTC())
if healthSeverity(decision.State) < healthSeverity(participant.HealthState) {
if _, err := s.repo.LiftTamperingHold(ctx, accountID, reason); err != nil {
return nil, errx.InternalError()
}
}
participant, xerr := s.getParticipantForAnyPool(ctx, accountID)
if xerr != nil || participant == nil || (!removed && !revised) {
return participant, xerr
}
return s.evaluateAndPersist(ctx, participant)
}
// tamperingHoldReason is the reason on a live pause or block the tampering
// band imposed, or "" when the row holds nothing or something else holds it.
func tamperingHoldReason(p *models.WarmupParticipantHealth) string {
if p.BlockedReason == nil || p.BlockedUntil == nil {
return ""
// reviseTamperingHold re-decides a live tampering hold on the strikes left in
// the seven days before it was imposed, and lowers it when they decide less.
// The bands never lower a hold on their own, so this is the only way a
// withdrawn strike reaches one.
func (s *service) reviseTamperingHold(ctx context.Context, accountID uuid.UUID) (bool, *errx.Error) {
hold, err := s.repo.GetWarmupHold(ctx, accountID)
if err != nil {
return false, errx.InternalError()
}
if p.HealthState != models.WarmupHealthQuarantined && p.HealthState != models.WarmupHealthBlocked {
return ""
if !isTamperingHold(hold) {
return false, nil
}
reason := *p.BlockedReason
if strings.HasPrefix(reason, tamperingPausePrefix) || strings.HasPrefix(reason, tamperingBlockPrefix) {
return reason
decidedAt := tamperingHoldDecidedAt(hold)
deletions, spamFlags, err := s.repo.CountWarmupTamperingBetween(ctx, accountID, decidedAt.Add(-7*24*time.Hour), decidedAt)
if err != nil {
return false, errx.InternalError()
}
return ""
decision := evaluateTampering(&models.WarmupHealthMetrics{DeletionsLast7d: deletions, SpamFlagsLast7d: spamFlags}, decidedAt)
if healthSeverity(decision.State) >= healthSeverity(hold.State) {
return false, nil
}
state, until, reason := decision.State, decision.BlockedUntil, decision.Reason
if until == nil || !until.After(s.now()) {
state, until, reason = models.WarmupHealthHealthy, nil, ""
}
revised, err := s.repo.ReviseWarmupHold(ctx, accountID, hold.Reason, state, until, reason)
if err != nil {
return false, errx.InternalError()
}
return revised, nil
}
// isTamperingHold is a live pause or block the tampering band imposed.
func isTamperingHold(h *repository.WarmupHold) bool {
if h == nil || h.BlockedUntil == nil {
return false
}
if h.State != models.WarmupHealthQuarantined && h.State != models.WarmupHealthBlocked {
return false
}
return strings.HasPrefix(h.Reason, tamperingPausePrefix) || strings.HasPrefix(h.Reason, tamperingBlockPrefix)
}
// tamperingHoldDecidedAt is when the hold was imposed, from its term when the
// row does not record it.
func tamperingHoldDecidedAt(h *repository.WarmupHold) time.Time {
if h.BlockedAt != nil {
return *h.BlockedAt
}
term := warmupQuarantineDuration
if h.State == models.WarmupHealthBlocked {
term = warmupBlockDuration
}
return h.BlockedUntil.Add(-term)
}
func tamperingVerb(kind string) string {
+86 -34
View File
@@ -8,68 +8,120 @@ import (
"github.com/google/uuid"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/repository"
)
// withdrawRepo holds one participant row and whatever strikes remain after a
// withdrawal, and records whether the hold was lifted.
// withdrawRepo holds one standing and the strikes left in its deciding window,
// and records how the hold was revised.
type withdrawRepo struct {
ownPoolRepo
deletionsLeft int
lifted string
hold *repository.WarmupHold
removed bool
deletions int
spamFlags int
revised bool
revisedState models.WarmupHealthState
revisedUntil *time.Time
window [2]time.Time
}
func (r *withdrawRepo) WithdrawWarmupTampering(context.Context, uuid.UUID, string, string) (bool, error) {
return true, nil
return r.removed, nil
}
func (r *withdrawRepo) HealthMetricCounts(context.Context, uuid.UUID, time.Time, time.Time) (models.WarmupHealthCounts, error) {
return models.WarmupHealthCounts{DeletionsLast7d: r.deletionsLeft}, nil
func (r *withdrawRepo) GetWarmupHold(context.Context, uuid.UUID) (*repository.WarmupHold, error) {
return r.hold, nil
}
func (r *withdrawRepo) LiftTamperingHold(_ context.Context, _ uuid.UUID, reason string) (bool, error) {
r.lifted = reason
return true, nil
func (r *withdrawRepo) CountWarmupTamperingBetween(_ context.Context, _ uuid.UUID, from, to time.Time) (int, int, error) {
r.window = [2]time.Time{from, to}
return r.deletions, r.spamFlags, nil
}
func heldRow(reason string) *models.WarmupParticipantHealth {
until := time.Now().Add(6 * 24 * time.Hour)
return &models.WarmupParticipantHealth{
PoolType: "premium", HealthState: models.WarmupHealthQuarantined,
BlockedUntil: &until, BlockedReason: &reason,
func (r *withdrawRepo) ReviseWarmupHold(_ context.Context, _ uuid.UUID, reason string, state models.WarmupHealthState, until *time.Time, _ string) (bool, error) {
if reason != r.hold.Reason {
return false, nil
}
r.revised, r.revisedState, r.revisedUntil = true, state, until
return true, nil
}
// A pause that rested on a withdrawn strike is lifted; one the remaining
// strikes still earn, or one something else imposed, stays.
func TestWithdrawTamperingLiftsOnlyTheHoldItImposed(t *testing.T) {
func tamperingHold(state models.WarmupHealthState, decidedAgo time.Duration, reason string) *repository.WarmupHold {
at := time.Now().Add(-decidedAgo)
term := warmupQuarantineDuration
if state == models.WarmupHealthBlocked {
term = warmupBlockDuration
}
until := at.Add(term)
return &repository.WarmupHold{State: state, BlockedAt: &at, BlockedUntil: &until, Reason: reason, InPool: true}
}
// A hold is re-decided on the strikes left in the seven days before it was
// imposed, not on today's window, so strikes that aged out still count.
func TestWithdrawTamperingRevisesOnTheWindowThatDecidedTheHold(t *testing.T) {
pause := tamperingPausePrefix + "2 warmup emails deleted in the last 7 days."
block := tamperingBlockPrefix + "5 warmup emails deleted in the last 7 days."
cases := []struct {
name string
reason string
left int
wantLift bool
wantState models.WarmupHealthState
hold *repository.WarmupHold
deletions int
want models.WarmupHealthState // "" means left alone
wantUntil bool
}{
{"one strike left only warns", pause, 1, true, models.WarmupHealthWatch},
{"two strikes left still pause", pause, 2, false, models.WarmupHealthQuarantined},
{"a complaint quarantine is not ours to lift", "complaint rate 0.20% exceeded quarantine threshold", 0, false, models.WarmupHealthHealthy},
{"a pause left with one strike lifts", tamperingHold(models.WarmupHealthQuarantined, 24*time.Hour, pause), 1, models.WarmupHealthHealthy, false},
{"a pause left with two strikes stands", tamperingHold(models.WarmupHealthQuarantined, 24*time.Hour, pause), 2, "", false},
{"an old block left with four strikes stands", tamperingHold(models.WarmupHealthBlocked, 10*24*time.Hour, block), 4, "", false},
{"a block left with three strikes becomes a pause from when it was decided", tamperingHold(models.WarmupHealthBlocked, 2*24*time.Hour, block), 3, models.WarmupHealthQuarantined, true},
{"a block whose pause would already have ended lifts", tamperingHold(models.WarmupHealthBlocked, 10*24*time.Hour, block), 3, models.WarmupHealthHealthy, false},
{"a complaint quarantine is not ours to lift", tamperingHold(models.WarmupHealthQuarantined, 24*time.Hour, "complaint rate 0.20% exceeded quarantine threshold"), 0, "", false},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
repo := &withdrawRepo{ownPoolRepo: ownPoolRepo{health: heldRow(tc.reason)}, deletionsLeft: tc.left}
health, err := NewService(repo).WithdrawTampering(context.Background(), uuid.New(), "<m@example.test>", "deletion")
if err != nil {
repo := &withdrawRepo{
ownPoolRepo: ownPoolRepo{health: &models.WarmupParticipantHealth{PoolType: "premium", HealthState: tc.hold.State}},
hold: tc.hold, removed: true, deletions: tc.deletions,
}
if _, err := NewService(repo).WithdrawTampering(context.Background(), uuid.New(), "<m@example.test>", "deletion"); err != nil {
t.Fatalf("unexpected error: %v", err)
}
if (repo.lifted != "") != tc.wantLift {
t.Fatalf("lifted = %q, want lift %v", repo.lifted, tc.wantLift)
if tc.want == "" {
if repo.revised {
t.Fatalf("revised to %v, want the hold left alone", repo.revisedState)
}
return
}
if tc.wantLift && repo.lifted != tc.reason {
t.Fatalf("lifted the hold with reason %q, want %q", repo.lifted, tc.reason)
if !repo.revised || repo.revisedState != tc.want {
t.Fatalf("revised %v to %v, want %v", repo.revised, repo.revisedState, tc.want)
}
if health == nil || health.HealthState != tc.wantState {
t.Fatalf("decided %v, want %v", health, tc.wantState)
if (repo.revisedUntil != nil) != tc.wantUntil {
t.Fatalf("revised until %v, want a term %v", repo.revisedUntil, tc.wantUntil)
}
if tc.wantUntil && !repo.revisedUntil.Equal(tc.hold.BlockedAt.Add(warmupQuarantineDuration)) {
t.Fatalf("pause ends %v, want seven days from when the block was decided", repo.revisedUntil)
}
if !repo.window[1].Equal(*tc.hold.BlockedAt) || !repo.window[0].Equal(tc.hold.BlockedAt.Add(-7*24*time.Hour)) {
t.Fatalf("counted strikes over %v, want the seven days before the hold", repo.window)
}
})
}
}
// A retry after the strike is already gone still finishes the revision, and a
// mailbox out of every pool has its ledger standing revised.
func TestWithdrawTamperingRetriesAndLedger(t *testing.T) {
pause := tamperingPausePrefix + "2 warmup emails deleted in the last 7 days."
hold := tamperingHold(models.WarmupHealthQuarantined, 24*time.Hour, pause)
hold.InPool = false
repo := &withdrawRepo{hold: hold, removed: false, deletions: 1}
health, err := NewService(repo).WithdrawTampering(context.Background(), uuid.New(), "<m@example.test>", "deletion")
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if !repo.revised || repo.revisedState != models.WarmupHealthHealthy {
t.Fatalf("revised %v to %v, want the ledger hold lifted", repo.revised, repo.revisedState)
}
if health != nil {
t.Fatalf("a mailbox in no pool came back with a standing: %+v", health)
}
}
+6 -8
View File
@@ -44,10 +44,11 @@ func (w *WorkerService) HandleWarmupAction(ctx context.Context, action models.Wa
case !exists:
// Engagement on a mailbox this worker is not holding is dropped; it
// is best effort and a mailbox mid-move earns its signal elsewhere.
// A retention delete is redelivered instead, because the control
// plane has already retired the row and will not send it again.
// A retention delete (its row is already retired) and a removal check
// are redelivered instead, in case the mailbox is still loading.
log.Warn().Str("email_id", action.EmailID.String()).Msg("Email account not found for warmup action")
if hasWarmupAction(action.Actions, models.WarmupActionDelete) {
if hasWarmupAction(action.Actions, models.WarmupActionDelete) ||
hasWarmupAction(action.Actions, models.WarmupActionVerifyRemoval) {
err = errors.New("mailbox not loaded on this worker")
}
case hasWarmupAction(action.Actions, models.WarmupActionVerifyRemoval):
@@ -67,11 +68,8 @@ func (w *WorkerService) HandleWarmupAction(ctx context.Context, action models.Wa
return nil
}
// Engagement is best effort and never returns here. A retention delete
// is not: the control plane has already retired the row, so a failure
// here is the only chance the message has of going. A removal check is
// retried too, since giving up means the removal is never judged. The
// bus redelivers on an error, bounded so a message the provider will
// never give up is not retried forever.
// (the row is already retired) and a removal check are redelivered, a
// bounded number of times.
if d := deliveryOf(ctx); d.redelivers && d.attempt < warmupDeleteRedeliveries {
return err
}
+22 -30
View File
@@ -12,10 +12,8 @@ import (
"github.com/warmbly/warmbly/internal/models"
)
// verifyWarmupRemoval looks for a warmup message the sync reported removed and
// reports where it is. A failed lookup is an error so the bus offers it again;
// one that never succeeds, or cannot tell, reports nothing and nothing is
// charged.
// verifyWarmupRemoval reports where a warmup message the sync saw removed is.
// A search that cannot tell reports nothing, so nothing is charged.
func (w *WorkerService) verifyWarmupRemoval(ctx context.Context, mail *wmail.WMail, action models.WarmupEmailAction) error {
outcome, err := locateWarmupMessage(ctx, mail, action.RFCMessageID)
if err != nil {
@@ -32,7 +30,6 @@ func (w *WorkerService) verifyWarmupRemoval(ctx context.Context, mail *wmail.WMa
return w.Produce(models.JobEventTypeWarmupRemovalChecked, action.EmailID.String(), &models.JobEventWarmupRemovalChecked{
UserID: action.UserID,
EmailID: action.EmailID,
InternalID: action.InternalID,
RFCMessageID: action.RFCMessageID,
Outcome: outcome,
})
@@ -64,40 +61,35 @@ func removalOutcome(found, trashed bool, err error) (string, error) {
return models.WarmupRemovalPresent, nil
}
// imapMessageHolder is the slice of the IMAP client the search needs.
type imapMessageHolder interface {
HoldsMessageID(ctx context.Context, mailboxName, rfcMessageID string) (bool, error)
// imapMessageLocator is the slice of the IMAP client the search needs.
type imapMessageLocator interface {
LocateMessageID(ctx context.Context, mailboxes []string, rfcMessageID string) (held []string, unsearched int, err error)
}
// locateImapMessage asks every synced folder, trash last. Not finding it is
// inconclusive: the synced list leaves out folders the owner excluded and
// folders past the sync's cap.
func locateImapMessage(ctx context.Context, client imapMessageHolder, boxes []*models.Mailbox, rfcMessageID string) (string, error) {
var trash []string
// locateImapMessage searches every synced folder. Not found, or found only in
// the trash with a folder left unsearched, cannot be told apart from a move
// into a folder the sync does not list.
func locateImapMessage(ctx context.Context, client imapMessageLocator, boxes []*models.Mailbox, rfcMessageID string) (string, error) {
var names []string
trash := map[string]bool{}
for _, b := range boxes {
if b == nil || !imap.SelectableFolder(b.Attrs) {
continue
}
if imap.IsTrashMailbox(b.Name, b.Attrs) {
trash = append(trash, b.Name)
continue
}
held, err := client.HoldsMessageID(ctx, b.Name, rfcMessageID)
if err != nil {
return "", err
}
if held {
names = append(names, b.Name)
trash[b.Name] = imap.IsTrashMailbox(b.Name, b.Attrs)
}
held, unsearched, err := client.LocateMessageID(ctx, names, rfcMessageID)
if err != nil {
return "", err
}
for _, name := range held {
if !trash[name] {
return models.WarmupRemovalPresent, nil
}
}
for _, name := range trash {
held, err := client.HoldsMessageID(ctx, name, rfcMessageID)
if err != nil {
return "", err
}
if held {
return models.WarmupRemovalTrashed, nil
}
if len(held) > 0 && unsearched == 0 {
return models.WarmupRemovalTrashed, nil
}
return "", nil
}
+19 -11
View File
@@ -8,19 +8,27 @@ import (
"github.com/warmbly/warmbly/internal/models"
)
// folderHolder answers which folders hold the message, and fails on one.
// folderHolder answers which folders hold the message, and which it cannot open.
type folderHolder struct {
holds map[string]bool
fails string
opened []string
holds map[string]bool
fails string
searched []string
}
func (h *folderHolder) HoldsMessageID(_ context.Context, mailbox, _ string) (bool, error) {
h.opened = append(h.opened, mailbox)
if mailbox == h.fails {
return false, errors.New("cannot select")
func (h *folderHolder) LocateMessageID(_ context.Context, mailboxes []string, _ string) ([]string, int, error) {
h.searched = mailboxes
var held []string
unsearched := 0
for _, m := range mailboxes {
if m == h.fails {
unsearched++
continue
}
if h.holds[m] {
held = append(held, m)
}
}
return h.holds[mailbox], nil
return held, unsearched, nil
}
func TestLocateImapMessageReportsWhereTheMessageIs(t *testing.T) {
@@ -40,8 +48,8 @@ func TestLocateImapMessageReportsWhereTheMessageIs(t *testing.T) {
{"moved to another folder", map[string]bool{"Warmbly": true}, "", models.WarmupRemovalPresent, false},
{"a copy outside the trash wins", map[string]bool{"Trash": true, "INBOX": true}, "", models.WarmupRemovalPresent, false},
{"only in the trash", map[string]bool{"Trash": true}, "", models.WarmupRemovalTrashed, false},
{"only in the trash with a folder unsearched is inconclusive", map[string]bool{"Trash": true}, "Warmbly", "", false},
{"nowhere synced is inconclusive", nil, "", "", false},
{"a folder that cannot be opened is an error", nil, "Warmbly", "", true},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
@@ -53,7 +61,7 @@ func TestLocateImapMessageReportsWhereTheMessageIs(t *testing.T) {
if got != tc.want {
t.Fatalf("outcome = %q, want %q", got, tc.want)
}
for _, name := range h.opened {
for _, name := range h.searched {
if name == "[Parent]" {
t.Fatal("searched a folder listed only as hierarchy")
}
+3 -3
View File
@@ -68,9 +68,9 @@ type ImapConn interface {
// expunge scoped to the UID, or a move into trashName where the server
// cannot scope one.
DeleteUID(ctx context.Context, mailboxName, trashName string, uid uint32) error
// HoldsMessageID answers whether one folder holds a message, and fails
// on a folder it cannot open: the check behind a removal's verdict.
HoldsMessageID(ctx context.Context, mailboxName, rfcMessageID string) (bool, error)
// LocateMessageID lists which folders hold a message: the search behind
// a warmup removal's verdict.
LocateMessageID(ctx context.Context, mailboxes []string, rfcMessageID string) (held []string, unsearched int, err error)
}
var _ ImapConn = (*imap.Client)(nil)
+2 -3
View File
@@ -282,9 +282,8 @@ func (c *Client) FindByRFCMessageID(ctx context.Context, rfcMessageID string) (s
return ids[0], nil
}
// LocateRFCMessageID reports whether the mailbox still holds the message
// carrying this RFC 5322 Message-ID, and whether every copy it holds is in
// Trash. Spam is searched too: a message filed there was moved, not deleted.
// LocateRFCMessageID reports whether the mailbox still holds the message, and
// whether every copy is in Trash. Spam counts as held: it was moved there.
func (c *Client) LocateRFCMessageID(ctx context.Context, rfcMessageID string) (found, trashed bool, err error) {
rfcMessageID = strings.Trim(strings.TrimSpace(rfcMessageID), "<>")
if rfcMessageID == "" || c.srv == nil {
+10 -8
View File
@@ -128,8 +128,7 @@ func (c *Client) move(ctx context.Context, messageID, destinationID string) (str
// ids change on move, so warmup actions re-resolve against this stable key.
// Returns an empty string (no error) when the message can't be found.
func (c *Client) ResolveMessageID(ctx context.Context, internetMessageID string) (string, error) {
filter := "internetMessageId eq '" + strings.ReplaceAll(internetMessageID, "'", "''") + "'"
u := c.root() + "/messages?$select=id&$top=1&$filter=" + url.QueryEscape(filter)
u := c.messagesByInternetID(internetMessageID, "id", 1)
var resp struct {
Value []struct {
ID string `json:"id"`
@@ -144,10 +143,8 @@ func (c *Client) ResolveMessageID(ctx context.Context, internetMessageID string)
return resp.Value[0].ID, nil
}
// LocateRFCMessageID reports whether the mailbox still holds the message with
// this internetMessageId in any folder, and whether every copy it holds is in
// Deleted Items. A move is copy plus delete, so the removal the delta reports
// for it says nothing about where the message went; this does.
// LocateRFCMessageID reports whether any folder still holds the message, and
// whether every copy is in Deleted Items. Delta reports a move as a removal.
func (c *Client) LocateRFCMessageID(ctx context.Context, internetMessageID string) (found, trashed bool, err error) {
id := strings.TrimSpace(internetMessageID)
if id == "" {
@@ -163,8 +160,7 @@ func (c *Client) LocateRFCMessageID(ctx context.Context, internetMessageID strin
forms = append(forms, "<"+id+">")
}
for _, form := range forms {
filter := "internetMessageId eq '" + strings.ReplaceAll(form, "'", "''") + "'"
u := c.root() + "/messages?$select=id,parentFolderId&$top=10&$filter=" + url.QueryEscape(filter)
u := c.messagesByInternetID(form, "id,parentFolderId", 10)
var resp struct {
Value []struct {
ParentFolderID string `json:"parentFolderId"`
@@ -186,6 +182,12 @@ func (c *Client) LocateRFCMessageID(ctx context.Context, internetMessageID strin
return false, false, nil
}
// messagesByInternetID lists the mailbox's messages carrying one internetMessageId.
func (c *Client) messagesByInternetID(internetMessageID, fields string, top int) string {
filter := "internetMessageId eq '" + strings.ReplaceAll(internetMessageID, "'", "''") + "'"
return c.root() + "/messages?$select=" + fields + "&$top=" + itoa(top) + "&$filter=" + url.QueryEscape(filter)
}
func (c *Client) messageURL(messageID string) string {
return c.root() + "/messages/" + url.PathEscape(messageID)
}
+35 -23
View File
@@ -38,16 +38,10 @@ func (c *Client) FindUIDByMessageID(ctx context.Context, mailboxName, rfcMessage
return 0, nil
}
// SEARCH HEADER matches on a substring of the header value, so the angle
// brackets are kept: a bare id would also match any message whose
// References or In-Reply-To names it.
data, err := c.client.UIDSearch(&imap.SearchCriteria{
Header: []imap.SearchCriteriaHeaderField{{Key: "Message-Id", Value: "<" + strings.Trim(rfcMessageID, "<>") + ">"}},
}, nil).Wait()
uids, err := c.searchMessageIDLocked(rfcMessageID)
if err != nil {
return 0, fmt.Errorf("search %q for message id: %w", name, err)
}
uids := data.AllUIDs()
if len(uids) == 0 {
return 0, nil
}
@@ -56,35 +50,55 @@ func (c *Client) FindUIDByMessageID(ctx context.Context, mailboxName, rfcMessage
return uint32(uids[len(uids)-1]), nil
}
// HoldsMessageID reports whether mailboxName holds the message with the given
// RFC 5322 Message-ID. Unlike FindUIDByMessageID, a folder that cannot be
// selected is an error: the caller is establishing that the message is gone.
func (c *Client) HoldsMessageID(ctx context.Context, mailboxName, rfcMessageID string) (bool, error) {
// LocateMessageID lists which of mailboxes hold the message, in one hold of
// the session. A folder that cannot be selected is counted as unsearched, not
// as not holding it.
func (c *Client) LocateMessageID(ctx context.Context, mailboxes []string, rfcMessageID string) (held []string, unsearched int, err error) {
rfcMessageID = strings.TrimSpace(rfcMessageID)
if mailboxName == "" || rfcMessageID == "" {
return false, errors.New("imap: nothing to look up")
if rfcMessageID == "" {
return nil, 0, errors.New("imap: no message id to look up")
}
c.mu.Lock()
defer c.mu.Unlock()
if merr := c.ensureConnected(); merr != nil {
return false, merr
return nil, 0, merr
}
c.lifecycle.RLock()
defer c.lifecycle.RUnlock()
defer c.begin()()
name := c.qualifyMailboxLocked(mailboxName)
if _, err := c.selectMailbox(name, nil); err != nil {
return false, fmt.Errorf("select %q: %w", name, err)
for _, mailbox := range mailboxes {
if ctx.Err() != nil {
return nil, 0, ctx.Err()
}
name := c.qualifyMailboxLocked(mailbox)
if _, err := c.selectMailbox(name, nil); err != nil {
unsearched++
continue
}
uids, err := c.searchMessageIDLocked(rfcMessageID)
if err != nil {
return nil, 0, fmt.Errorf("search %q for message id: %w", name, err)
}
if len(uids) > 0 {
held = append(held, mailbox)
}
}
return held, unsearched, nil
}
// searchMessageIDLocked searches the selected folder for one Message-ID. The
// brackets are kept because SEARCH HEADER matches substrings, and a bare id
// would also match a message whose References names it.
func (c *Client) searchMessageIDLocked(rfcMessageID string) ([]imap.UID, error) {
data, err := c.client.UIDSearch(&imap.SearchCriteria{
Header: []imap.SearchCriteriaHeaderField{{Key: "Message-Id", Value: "<" + strings.Trim(rfcMessageID, "<>") + ">"}},
}, nil).Wait()
if err != nil {
return false, fmt.Errorf("search %q for message id: %w", name, err)
return nil, err
}
return len(data.AllUIDs()) > 0, nil
return data.AllUIDs(), nil
}
// SelectableFolder is false for a folder listed only as hierarchy.
@@ -120,13 +134,11 @@ func (c *Client) FindUIDsByMessageIDs(ctx context.Context, mailboxName string, r
if id == "" {
continue
}
data, err := c.client.UIDSearch(&imap.SearchCriteria{
Header: []imap.SearchCriteriaHeaderField{{Key: "Message-Id", Value: "<" + strings.Trim(id, "<>") + ">"}},
}, nil).Wait()
uids, err := c.searchMessageIDLocked(id)
if err != nil {
return nil, fmt.Errorf("search %q for message id: %w", name, err)
}
if uids := data.AllUIDs(); len(uids) > 0 {
if len(uids) > 0 {
found[raw] = uint32(uids[len(uids)-1])
}
}
@@ -1 +1,3 @@
DROP INDEX IF EXISTS idx_warmup_tampering_unverified;
ALTER TABLE warmup_tampering_events DROP COLUMN IF EXISTS verify_requested_at;
ALTER TABLE warmup_tampering_events DROP COLUMN IF EXISTS verified_at;
@@ -1,5 +1,7 @@
-- When a deletion strike was confirmed by searching the mailbox for the
-- message. Rows recorded before that search existed stay NULL and are searched
-- once by the consumer's recheck; every row written from now on is confirmed.
-- verified_at: when a search of the mailbox confirmed the strike; NULL on rows from before the search.
-- verify_requested_at: when the consumer last asked a worker to search for one of those.
ALTER TABLE warmup_tampering_events ADD COLUMN IF NOT EXISTS verified_at timestamptz;
ALTER TABLE warmup_tampering_events ALTER COLUMN verified_at SET DEFAULT now();
ALTER TABLE warmup_tampering_events ADD COLUMN IF NOT EXISTS verify_requested_at timestamptz;
CREATE INDEX IF NOT EXISTS idx_warmup_tampering_unverified
ON warmup_tampering_events (created_at) WHERE verified_at IS NULL;
+1 -3
View File
@@ -27,12 +27,10 @@ type JobEventRemoveEmail struct {
}
// JobEventWarmupRemovalChecked is where the worker found a warmup message
// after the sync reported it removed. InternalID is empty when the control
// plane did not know it.
// after the sync reported it removed.
type JobEventWarmupRemovalChecked struct {
UserID uuid.UUID `json:"user_id" avro:"user_id"`
EmailID uuid.UUID `json:"email_id" avro:"email_id"`
InternalID string `json:"internal_id,omitempty" avro:"internal_id"`
RFCMessageID string `json:"rfc_message_id" avro:"rfc_message_id"`
Outcome string `json:"outcome" avro:"outcome"`
}
+2 -4
View File
@@ -62,10 +62,8 @@ const (
// by the retention sweep alone, and only for a message it has retired
// first, so the removal the sync then observes is never a strike.
WarmupActionDelete = "delete"
// WarmupActionVerifyRemoval asks the worker where a warmup message the
// sync reported gone actually is. It changes nothing in the mailbox; the
// answer comes back as WARMUP_REMOVAL_CHECKED, and only a message in the
// trash or gone for good is held against the mailbox.
// WarmupActionVerifyRemoval asks where a removed warmup message went,
// answered by WARMUP_REMOVAL_CHECKED; it changes nothing in the mailbox.
WarmupActionVerifyRemoval = "verify_removal"
)
+112 -33
View File
@@ -240,18 +240,23 @@ type WarmupRepository interface {
GetWarmupReceived(ctx context.Context, accountID, internalID uuid.UUID) (*WarmupReceived, error)
RecordWarmupTampering(ctx context.Context, accountID uuid.UUID, messageID, kind string) (bool, error)
CountWarmupTamperingSince(ctx context.Context, accountID uuid.UUID, since time.Time) (int, error)
// WithdrawWarmupTampering deletes a strike the mailbox turned out not to
// have earned, reporting whether there was one.
// WithdrawWarmupTampering deletes a strike the mailbox did not earn,
// reporting whether there was one.
WithdrawWarmupTampering(ctx context.Context, accountID uuid.UUID, messageID, kind string) (bool, error)
// ListUnverifiedDeletions is deletion strikes since the given time that
// were recorded before removals were checked against the mailbox, on
// active mailboxes with a worker to search them. MarkTamperingVerified
// stamps one once its search is on the bus.
ListUnverifiedDeletions(ctx context.Context, since time.Time, limit int) ([]WarmupTamperingToVerify, error)
// CountWarmupTamperingBetween is the strikes a hold decided at `to` read.
CountWarmupTamperingBetween(ctx context.Context, accountID uuid.UUID, from, to time.Time) (deletions, spamFlags int, err error)
// ListUnverifiedDeletions is deletion strikes since `since` that no search
// has confirmed and none was asked for within `retryAfter`, on active
// mailboxes with a worker.
ListUnverifiedDeletions(ctx context.Context, since time.Time, retryAfter time.Duration, limit int) ([]WarmupTamperingToVerify, error)
MarkTamperingVerifyRequested(ctx context.Context, accountID uuid.UUID, messageID string) error
MarkTamperingVerified(ctx context.Context, accountID uuid.UUID, messageID, kind string) error
// LiftTamperingHold clears a live pause or block whose reason is still
// the given one, so a hold set since by something else is left alone.
LiftTamperingHold(ctx context.Context, accountID uuid.UUID, reason string) (bool, error)
// GetWarmupHold is the mailbox's pool row standing, else the ledger row
// that would seed one; nil when neither holds anything.
GetWarmupHold(ctx context.Context, accountID uuid.UUID) (*WarmupHold, error)
// ReviseWarmupHold replaces a hold still carrying reason, on the pool row
// or, with none, on the ledger. A healthy state clears it.
ReviseWarmupHold(ctx context.Context, accountID uuid.UUID, reason string, state models.WarmupHealthState, until *time.Time, newReason string) (bool, error)
// Retention: warmup mail is deleted from the mailbox once its window has
// passed, the platform's own copy of the body with it, and the
@@ -1915,35 +1920,38 @@ func (r *warmupRepository) WithdrawWarmupTampering(ctx context.Context, accountI
return cmd.RowsAffected() > 0, nil
}
// WarmupTamperingToVerify is one unverified deletion strike and what the
// worker needs to search the mailbox for its message.
func (r *warmupRepository) CountWarmupTamperingBetween(ctx context.Context, accountID uuid.UUID, from, to time.Time) (int, int, error) {
var deletions, spamFlags int
err := r.db.QueryRow(ctx, `
SELECT COUNT(*) FILTER (WHERE kind = 'deletion'), COUNT(*) FILTER (WHERE kind = 'spam_flag')
FROM warmup_tampering_events
WHERE email_account_id = $1 AND created_at >= $2 AND created_at <= $3`,
accountID, from, to).Scan(&deletions, &spamFlags)
return deletions, spamFlags, err
}
// WarmupTamperingToVerify is one unverified deletion strike and where to ask.
type WarmupTamperingToVerify struct {
EmailAccountID uuid.UUID
UserID uuid.UUID
WorkerID uuid.UUID
MessageID string
// InternalID is the receipt's message id; nil once the receipt is pruned.
InternalID *uuid.UUID
}
func (r *warmupRepository) ListUnverifiedDeletions(ctx context.Context, since time.Time, limit int) ([]WarmupTamperingToVerify, error) {
func (r *warmupRepository) ListUnverifiedDeletions(ctx context.Context, since time.Time, retryAfter time.Duration, limit int) ([]WarmupTamperingToVerify, error) {
rows, err := r.db.Query(ctx, `
SELECT t.email_account_id, ea.user_id, ea.worker_id, t.message_id, wr.internal_id
SELECT t.email_account_id, ea.user_id, ea.worker_id, t.message_id
FROM warmup_tampering_events t
JOIN email_accounts ea ON ea.id = t.email_account_id
LEFT JOIN LATERAL (
SELECT r.internal_id FROM warmup_received r
WHERE r.email_account_id = t.email_account_id AND r.message_id = t.message_id
LIMIT 1
) wr ON true
WHERE t.kind = 'deletion'
AND t.verified_at IS NULL
WHERE t.verified_at IS NULL
AND t.kind = 'deletion'
AND t.created_at >= $1
AND (t.verify_requested_at IS NULL OR t.verify_requested_at < NOW() - make_interval(secs => $2))
AND t.message_id <> ''
AND ea.status = 'active'
AND ea.worker_id IS NOT NULL
ORDER BY t.created_at
LIMIT $2`, since, limit)
LIMIT $3`, since, retryAfter.Seconds(), limit)
if err != nil {
return nil, err
}
@@ -1951,7 +1959,7 @@ func (r *warmupRepository) ListUnverifiedDeletions(ctx context.Context, since ti
var out []WarmupTamperingToVerify
for rows.Next() {
var v WarmupTamperingToVerify
if err := rows.Scan(&v.EmailAccountID, &v.UserID, &v.WorkerID, &v.MessageID, &v.InternalID); err != nil {
if err := rows.Scan(&v.EmailAccountID, &v.UserID, &v.WorkerID, &v.MessageID); err != nil {
return nil, err
}
out = append(out, v)
@@ -1959,6 +1967,14 @@ func (r *warmupRepository) ListUnverifiedDeletions(ctx context.Context, since ti
return out, rows.Err()
}
func (r *warmupRepository) MarkTamperingVerifyRequested(ctx context.Context, accountID uuid.UUID, messageID string) error {
_, err := r.db.Exec(ctx, `
UPDATE warmup_tampering_events SET verify_requested_at = NOW()
WHERE email_account_id = $1 AND message_id = $2 AND kind = 'deletion' AND verified_at IS NULL`,
accountID, messageID)
return err
}
func (r *warmupRepository) MarkTamperingVerified(ctx context.Context, accountID uuid.UUID, messageID, kind string) error {
_, err := r.db.Exec(ctx, `
UPDATE warmup_tampering_events SET verified_at = NOW()
@@ -1967,18 +1983,81 @@ func (r *warmupRepository) MarkTamperingVerified(ctx context.Context, accountID
return err
}
func (r *warmupRepository) LiftTamperingHold(ctx context.Context, accountID uuid.UUID, reason string) (bool, error) {
// WarmupHold is a mailbox's standing as a pool row or ledger row holds it.
type WarmupHold struct {
State models.WarmupHealthState
BlockedAt *time.Time
BlockedUntil *time.Time
Reason string
// InPool is false when the standing is only on the ledger.
InPool bool
}
func (r *warmupRepository) GetWarmupHold(ctx context.Context, accountID uuid.UUID) (*WarmupHold, error) {
var h WarmupHold
var state string
err := r.db.QueryRow(ctx, `
SELECT health_state, blocked_at, blocked_until, COALESCE(blocked_reason, ''), true
FROM warmup_pool_participants WHERE email_account_id = $1
UNION ALL
SELECT l.health_state, l.blocked_at, l.blocked_until, COALESCE(l.blocked_reason, ''), false
FROM warmup_reputation_ledger l
JOIN email_accounts a ON a.organization_id = l.organization_id AND lower(btrim(a.email)) = l.email
WHERE a.id = $1
AND NOT EXISTS (SELECT 1 FROM warmup_pool_participants WHERE email_account_id = $1)
LIMIT 1`, accountID).Scan(&state, &h.BlockedAt, &h.BlockedUntil, &h.Reason, &h.InPool)
if errors.Is(err, sql.ErrNoRows) {
return nil, nil
}
if err != nil {
return nil, err
}
h.State = models.WarmupHealthState(state)
return &h, nil
}
func (r *warmupRepository) ReviseWarmupHold(ctx context.Context, accountID uuid.UUID, reason string, state models.WarmupHealthState, until *time.Time, newReason string) (bool, error) {
held := state == models.WarmupHealthQuarantined || state == models.WarmupHealthBlocked
if !held {
until, newReason = nil, ""
}
cmd, err := r.db.Exec(ctx, `
UPDATE warmup_pool_participants
SET health_state = 'healthy',
blocked_at = NULL,
blocked_until = NULL,
blocked_reason = NULL
SET health_state = $3::text,
blocked_until = $4::timestamptz,
blocked_reason = NULLIF($5::text, ''),
blocked_at = CASE WHEN $4::timestamptz IS NULL THEN NULL ELSE blocked_at END
WHERE email_account_id = $1
AND blocked_reason = $2
AND health_state IN ('quarantined', 'blocked')
AND blocked_until IS NOT NULL`,
accountID, reason)
AND health_state IN ('quarantined', 'blocked')`,
accountID, reason, string(state), until, newReason)
if err != nil {
return false, err
}
if cmd.RowsAffected() > 0 {
return true, nil
}
// The pool row is the standing whenever there is one; the ledger only
// speaks for a mailbox out of every pool.
if !held {
cmd, err = r.db.Exec(ctx, `
DELETE FROM warmup_reputation_ledger l
USING email_accounts a
WHERE a.id = $1 AND l.organization_id = a.organization_id AND l.email = lower(btrim(a.email))
AND l.blocked_reason = $2
AND NOT EXISTS (SELECT 1 FROM warmup_pool_participants WHERE email_account_id = $1)`,
accountID, reason)
} else {
cmd, err = r.db.Exec(ctx, `
UPDATE warmup_reputation_ledger l
SET health_state = $3::text, blocked_until = $4::timestamptz, blocked_reason = $5::text,
standing_until = GREATEST($4::timestamptz, NOW())
FROM email_accounts a
WHERE a.id = $1 AND l.organization_id = a.organization_id AND l.email = lower(btrim(a.email))
AND l.blocked_reason = $2
AND NOT EXISTS (SELECT 1 FROM warmup_pool_participants WHERE email_account_id = $1)`,
accountID, reason, string(state), until, newReason)
}
if err != nil {
return false, err
}
@@ -0,0 +1,121 @@
package repository
import (
"context"
"testing"
"time"
"github.com/warmbly/warmbly/internal/models"
)
// A withdrawn tampering strike revises the hold it imposed wherever that hold
// lives: the pool row, or the ledger of a mailbox out of every pool, which
// would otherwise seed the hold back on rejoin.
//
// WARMBLY_TEST_DB=postgres://warmbly:warmbly@localhost:15432/warmbly_x?sslmode=disable \
// go test ./internal/repository/ -run LiveTamperingHold -v
func TestLiveTamperingHoldRevision(t *testing.T) {
f := newLedgerFixture(t)
ctx := context.Background()
id := f.addMailbox(t, f.user)
f.join(t, id)
until := time.Now().Add(6 * 24 * time.Hour)
f.exec(t, `UPDATE warmup_pool_participants SET health_state = 'quarantined', blocked_at = now(),
blocked_until = $2, blocked_reason = 'Paused from warmup: x' WHERE email_account_id = $1`, id, until)
hold, err := f.warmups.GetWarmupHold(ctx, id)
if err != nil || hold == nil || !hold.InPool || hold.State != models.WarmupHealthQuarantined || hold.Reason != "Paused from warmup: x" {
t.Fatalf("pool hold = %+v, %v", hold, err)
}
if ok, err := f.warmups.ReviseWarmupHold(ctx, id, "something else", models.WarmupHealthHealthy, nil, ""); err != nil || ok {
t.Fatalf("revised a hold whose reason changed: %v %v", ok, err)
}
// Out of every pool, the ledger speaks for the mailbox.
if err := f.warmups.LeaveAllPools(ctx, id); err != nil {
t.Fatal(err)
}
hold, err = f.warmups.GetWarmupHold(ctx, id)
if err != nil || hold == nil || hold.InPool || hold.State != models.WarmupHealthQuarantined {
t.Fatalf("ledger hold = %+v, %v", hold, err)
}
shorter := time.Now().Add(time.Hour)
if ok, err := f.warmups.ReviseWarmupHold(ctx, id, "Paused from warmup: x", models.WarmupHealthQuarantined, &shorter, "Paused from warmup: y"); err != nil || !ok {
t.Fatalf("ledger revision: %v %v", ok, err)
}
if ok, err := f.warmups.ReviseWarmupHold(ctx, id, "Paused from warmup: y", models.WarmupHealthHealthy, nil, ""); err != nil || !ok {
t.Fatalf("ledger lift: %v %v", ok, err)
}
if hold, err = f.warmups.GetWarmupHold(ctx, id); err != nil || hold != nil {
t.Fatalf("ledger still holds %+v, %v", hold, err)
}
// Rejoining seeds nothing, and a pool-row lift clears the ledger mirror too.
f.join(t, id)
f.exec(t, `UPDATE warmup_pool_participants SET health_state = 'blocked', blocked_at = now(),
blocked_until = $2, blocked_reason = 'Blocked from warmup: x' WHERE email_account_id = $1`, id, until)
if ok, err := f.warmups.ReviseWarmupHold(ctx, id, "Blocked from warmup: x", models.WarmupHealthHealthy, nil, ""); err != nil || !ok {
t.Fatalf("pool lift: %v %v", ok, err)
}
var ledger int
if err := f.pool.QueryRow(ctx, `SELECT COUNT(*) FROM warmup_reputation_ledger WHERE organization_id = $1`, f.org).Scan(&ledger); err != nil || ledger != 0 {
t.Fatalf("ledger rows after lift = %d, %v", ledger, err)
}
if hold, err = f.warmups.GetWarmupHold(ctx, id); err != nil || hold == nil || hold.State != models.WarmupHealthHealthy || hold.BlockedUntil != nil {
t.Fatalf("pool row after lift = %+v, %v", hold, err)
}
}
// Strikes from before the search are listed until one is asked for, then
// again after the retry window; confirmed and new strikes never are.
func TestLiveTamperingUnverifiedListing(t *testing.T) {
f := newLedgerFixture(t)
ctx := context.Background()
id := f.addMailbox(t, f.user)
worker := f.org // any uuid; the listing only needs one set
f.exec(t, `INSERT INTO fleet_nodes (id, role, name) VALUES ($1, 'worker', 'tamper-test') ON CONFLICT DO NOTHING`, worker)
f.exec(t, `INSERT INTO workers (id) VALUES ($1) ON CONFLICT DO NOTHING`, worker)
f.exec(t, `UPDATE email_accounts SET worker_id = $2 WHERE id = $1`, id, worker)
t.Cleanup(func() {
_, _ = f.pool.Exec(ctx, `DELETE FROM warmup_tampering_events WHERE email_account_id = $1`, id)
_, _ = f.pool.Exec(ctx, `UPDATE email_accounts SET worker_id = NULL WHERE id = $1`, id)
_, _ = f.pool.Exec(ctx, `DELETE FROM workers WHERE id = $1`, worker)
_, _ = f.pool.Exec(ctx, `DELETE FROM fleet_nodes WHERE id = $1`, worker)
})
f.exec(t, `INSERT INTO warmup_tampering_events (email_account_id, message_id, kind, verified_at) VALUES ($1, '<old@t>', 'deletion', NULL)`, id)
if _, err := f.warmups.RecordWarmupTampering(ctx, id, "<new@t>", "deletion"); err != nil {
t.Fatal(err)
}
list := func() []WarmupTamperingToVerify {
rows, err := f.warmups.ListUnverifiedDeletions(ctx, time.Now().Add(-time.Hour), 6*time.Hour, 10)
if err != nil {
t.Fatal(err)
}
return rows
}
rows := list()
if len(rows) != 1 || rows[0].MessageID != "<old@t>" || rows[0].WorkerID != worker || rows[0].UserID != f.user {
t.Fatalf("listed %+v, want only the pre-search strike", rows)
}
if err := f.warmups.MarkTamperingVerifyRequested(ctx, id, "<old@t>"); err != nil {
t.Fatal(err)
}
if rows = list(); len(rows) != 0 {
t.Fatalf("listed %+v right after asking", rows)
}
f.exec(t, `UPDATE warmup_tampering_events SET verify_requested_at = now() - interval '7 hours' WHERE email_account_id = $1`, id)
if rows = list(); len(rows) != 1 {
t.Fatalf("an unanswered search was not offered again: %+v", rows)
}
if err := f.warmups.MarkTamperingVerified(ctx, id, "<old@t>", "deletion"); err != nil {
t.Fatal(err)
}
if rows = list(); len(rows) != 0 {
t.Fatalf("listed a confirmed strike: %+v", rows)
}
d, s, err := f.warmups.CountWarmupTamperingBetween(ctx, id, time.Now().Add(-time.Hour), time.Now().Add(time.Minute))
if err != nil || d != 2 || s != 0 {
t.Fatalf("counted %d deletions, %d flags, %v", d, s, err)
}
}