diff --git a/AGENTS.md b/AGENTS.md
index 3ad68442a..1066a2c7d 100644
--- a/AGENTS.md
+++ b/AGENTS.md
@@ -768,7 +768,7 @@ Signals used:
- every warmup email carries a verification token, minted by the platform, single-use, bound to its recipient
- 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
+- 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` counts a deletion and a spam move as one strike each over the seven-day window, and warns at one, quarantines at two and blocks at four. No provider names who moved a message into spam (Microsoft's ZAP, Workspace post-delivery scanning and client junk filters look exactly like a user's report), so a spam move is never charged on sight. A spam label on mail that arrived in spam is the filter's (`warmup_received.landed_spam`; Gmail can report it as a later label change) and is not even held. Any other move is held in `warmup_spam_moves` and `attributeSpamMove` (`internal/app/consumer/warmup_spam_attribution.go`) decides it after `config.WarmupSpamMoveSettleMinutes`: provider when the same sender was junked in another workspace within a day (which also withdraws owner verdicts it explains) or the move came straight after arrival with nobody there; owner when `mailbox_owner_activity` shows the owner at the mailbox around it, or a mailbox in use keeps junking many senders nobody else does; nobody otherwise. Owner activity is only a read, unread or star change the provider reported that our store did not already hold, on mail past its arrival grace, so Warmbly's own echoes and a filter finishing delivery never count. Only an owner verdict strikes the recipient and files a complaint against the sender; the rest are placement (`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`), 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 address's ledger row when no mailbox with the address is in a 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 (`Recheck`), stamps `verify_requested_at` and asks again after six hours while unanswered. A recheck never adds a strike: present or retired-by-retention withdraws it, anything else stamps `verified_at` and it stands. Tampering events are kept at least `config.WarmupTamperingKeepDays` so the strikes behind a live hold are there to re-decide it. 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
diff --git a/cmd/consumer/main.go b/cmd/consumer/main.go
index e78e19928..0911f2d18 100644
--- a/cmd/consumer/main.go
+++ b/cmd/consumer/main.go
@@ -525,6 +525,8 @@ func main() {
// Searches the mailbox for deletion strikes recorded before removals were
// checked, withdrawing any whose message is still there.
go jobsService.StartWarmupTamperingRecheck(ctx)
+ // Attributes each warmup email moved to spam once the activity around it settles.
+ go jobsService.StartWarmupSpamMoveAttribution(ctx)
go jobsService.StartWarmupPlacementSweep(ctx)
go jobsService.StartPendingWarmupVerification(ctx)
// Re-offers inbound mail that reply processing never claimed, so a
diff --git a/docs/content/docs/development/data-control.mdx b/docs/content/docs/development/data-control.mdx
index f1cbef9ed..90677ead2 100644
--- a/docs/content/docs/development/data-control.mdx
+++ b/docs/content/docs/development/data-control.mdx
@@ -98,7 +98,7 @@ See [Automatic inbox tagging](/guides/inbox-tagging/), [Advisor](/guides/advisor
| `retention.form_event_days` | 180 | Form funnel events: views, starts, field-level drop-off |
| `retention.audit_log_days` | 90 | The audit trail: actor, IP address, user agent, change payload |
| `retention.warmup_mail_days` | 30 | Warmup mail in the mailboxes themselves, and the stored copy of each message's body. Deleted by the platform once older than this, wherever the mailbox files it; a mailbox can set its own window |
-| `retention.warmup_event_days` | 365 | Per-message warmup records: tokens, receipts, tampering events, spam reports. The daily sent and received counts behind the analytics are kept |
+| `retention.warmup_event_days` | 365 | Per-message warmup records: tokens, receipts, tampering events, spam moves and their attribution, spam reports. The daily sent and received counts behind the analytics are kept. When each mailbox's owner last acted in it (five-minute marks, no message content) is kept 30 days whatever this is set to |
The first three are between 1 and 3,650 days, and they are the settings a retention or privacy policy applies to, because each window is also how long the personal data in that log is held. The warmup mail window starts at 3 days and the warmup records window at 30, the least the engagement legs and the pool health bands need.
diff --git a/docs/content/docs/guides/deliverability.mdx b/docs/content/docs/guides/deliverability.mdx
index a76130892..a57196b62 100644
--- a/docs/content/docs/guides/deliverability.mdx
+++ b/docs/content/docs/guides/deliverability.mdx
@@ -42,7 +42,7 @@ Each mailbox sits in a band, shown as a colored chip. The band controls how Warm
| Watch | `>= 10%` | `>= 0.03%` | n/a | n/a |
| Throttled | `>= 20%` | n/a | n/a | n/a |
| Quarantine | n/a | `>= 0.10%` | `>= 5%` | Repeated tampering with received warmup mail |
-| Blocked | n/a | `>= 0.30%` | `>= 10%` | Clear abuse signals (repeated spam flags on received warmup mail) |
+| Blocked | n/a | `>= 0.30%` | `>= 10%` | Clear abuse signals (four or more deletions or spam moves of received warmup mail in 7 days) |
A mailbox enters a band by crossing any single threshold. What each band does:
diff --git a/docs/content/docs/guides/warmup.mdx b/docs/content/docs/guides/warmup.mdx
index 8f2c32a1a..ca63b41ec 100644
--- a/docs/content/docs/guides/warmup.mdx
+++ b/docs/content/docs/guides/warmup.mdx
@@ -233,19 +233,29 @@ Partner selection favours the recipients where warmup earns something. A small-h
Quarantined and blocked mailboxes are selected as neither sender nor recipient. Mail arriving in a mailbox is never held against it: every warmup message carries a single-use token bound to its recipient, so a token cannot be replayed or redirected, and one that lands where it does not belong is simply filed as ordinary mail.
-What a mailbox does to warmup mail it received is held against it, on a ladder rather than at once. Deleting a warmup email within a day of its arrival counts as one strike and marking one as spam counts as two, over the last seven days. One strike puts the mailbox on watch with the reason shown in its drawer, two pause it from the pool for seven days, and four block it for thirty. So deleting one fresh warmup message is a warning, not a ban, while flagging pool mail as spam twice is a block.
+What a mailbox does to warmup mail it received is held against it, on a ladder rather than at once. Deleting a warmup email within a day of its arrival counts as one strike, and so does moving one to spam, over the last seven days. One strike puts the mailbox on watch with the reason shown in its drawer, two pause it from the pool for seven days, and four block it for thirty. So one deleted or junked warmup message is a warning, not a ban.
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 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.
+No provider says who moved a message into spam. Gmail, Microsoft 365 and IMAP report it the same way whether you pressed "Report spam", the provider re-filed it after delivery (Microsoft's zero-hour auto purge, Google Workspace's post-delivery scanning), a desktop mail client's junk filter moved it, or a security tool did. So Warmbly charges a move to spam only when the evidence points at a person, and waits 30 minutes after the move to see it:
+
+- **Arrived in spam.** A warmup email that was in spam when it arrived is never a strike. The label on it is the provider's filter, and it counts as spam placement against the sender.
+- **The provider at work.** When the same sender's warmup mail was moved to spam in another workspace within a day, the provider is re-judging that sender, and nobody is charged. If a mailbox was already charged for one of those, the strike is withdrawn and any pause it caused is decided again. A move within 15 minutes of arrival while nobody is using the mailbox is the filter catching up, and is not charged either.
+- **You at the mailbox.** A move while someone was reading, marking unread or starring mail in the mailbox, within 30 minutes either side, is taken as yours. Only changes made at the provider count: reading a message in Warmbly's inbox is never mistaken for it.
+- **A pattern.** In a mailbox someone has used in the last 14 days, unexplained moves of warmup mail from three or more different senders in a week, which no other workspace sees, are the mailbox's own filtering and are charged.
+- **Anything else charges nobody**, including every move in a mailbox nobody has used for 14 days. The sender still has it counted as spam placement, which only slows sending down.
+
+Only a move charged to you files a complaint against the sender. A move attributed to the provider, or to nobody, is spam placement for the sender instead. On Outlook, Microsoft 365 and IMAP mailboxes, a move to Junk is found by the search above and is not charged at all.
+
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.
Getting back in requires requalifying, not just waiting: healthy authentication, no recent complaints or hard-bounce spikes, and spam placement back to a low level. Return is gradual, not a jump back to the old ceiling.
-A quarantine or a block also holds for its full term. The signals behind it age out of their windows long before it ends, and that does not release it early; a more serious finding can still replace it. A throttle is different and lifts as soon as the mailbox recovers, as the table says.
+A quarantine or a block also holds for its full term. The signals behind it age out of their windows long before it ends, and that does not release it early; a more serious finding can still replace it. A throttle is different and lifts as soon as the mailbox recovers, as the table says. When support lifts a pause or block, or approves an appeal against one, the strikes behind it are cleared with it, so the next evaluation does not reimpose it.
On a self-hosted instance linked to [Warmbly Cloud](/guides/warmbly-cloud/#safety-and-enforcement), a mailbox the cloud warms is judged there, and the instance applies the cloud's standing to its own campaigns with the same effects as this table.
diff --git a/internal/app/consumer/event_flags_update.go b/internal/app/consumer/event_flags_update.go
index c010621ed..4322a3150 100644
--- a/internal/app/consumer/event_flags_update.go
+++ b/internal/app/consumer/event_flags_update.go
@@ -8,6 +8,7 @@ import (
"time"
"github.com/google/uuid"
+ "github.com/rs/zerolog/log"
"github.com/warmbly/warmbly/internal/config"
"github.com/warmbly/warmbly/internal/models"
@@ -16,19 +17,22 @@ import (
func (s *JobsService) HandleFlagsAdd(ctx context.Context, e *models.JobEventFlags) error {
// Tampering check first: verified warmup mail is NOT in the unibox, so we
- // detect it via the warmup_received record. If the recipient marked a
- // warmup email as spam, that harms the pool — penalise the sender for the
- // spam signal AND ban the harmer (they can appeal). Warmup mail isn't
- // tracked in the unibox, so there's nothing else to do for it.
+ // detect it via the warmup_received record. A spam label names no actor on
+ // any provider: on mail that arrived in spam it is the filter's own, and
+ // any other move is held and attributed on the evidence around it
+ // (attributeSpamMove) before anyone is charged.
if s.WarmupRepo != nil {
if rec, _ := s.WarmupRepo.GetWarmupReceived(ctx, e.EmailID, e.ID); rec != nil {
switch {
case s.WarmupService == nil:
+ case containsSpamFlag(e.Flags) && (rec.LandedSpam || rec.MessageID == ""):
case containsSpamFlag(e.Flags):
- hSender, _ := s.WarmupService.ApplySpamReport(ctx, e.EmailID, rec.SenderAccountID, rec.MessageID, "user_complaint")
- s.markRiskBandFromWarmupHealth(ctx, rec.SenderAccountID, hSender)
- hHarmer, _ := s.WarmupService.RecordTampering(ctx, e.EmailID, rec.MessageID, "spam_flag")
- s.markRiskBandFromWarmupHealth(ctx, e.EmailID, hHarmer)
+ if _, err := s.WarmupRepo.RecordWarmupSpamMove(ctx, repository.WarmupSpamMove{
+ EmailAccountID: e.EmailID, MessageID: rec.MessageID,
+ SenderAccountID: rec.SenderAccountID, ReceivedAt: rec.CreatedAt,
+ }); err != nil {
+ return fmt.Errorf("hold warmup spam move: %w", err)
+ }
case containsTrashFlag(e.Flags) && warmupDeletionCounts(rec, time.Now()):
// Gmail reports Delete as gaining the TRASH label and only
// reports the message gone when Trash is emptied, weeks later.
@@ -98,6 +102,9 @@ func (s *JobsService) HandleFlagsAdd(ctx context.Context, e *models.JobEventFlag
if !slices.Contains(email.Flags, e.Flags[i]) {
email.Flags = append(email.Flags, e.Flags[i])
updated = true
+ if e.Flags[i] == models.FlagFlagged {
+ s.noteOwnerActivity(ctx, e.EmailID, email.InternalDate)
+ }
}
}
@@ -111,6 +118,7 @@ func (s *JobsService) HandleFlagsAdd(ctx context.Context, e *models.JobEventFlag
update.Seen = &seen
email.Seen = true
updated = true
+ s.noteOwnerActivity(ctx, e.EmailID, email.InternalDate)
}
if !updated {
@@ -175,6 +183,9 @@ func (s *JobsService) HandleFlagsRemove(ctx context.Context, e *models.JobEventF
// Losing \Seen is the provider reporting the message back to unread, and
// that is the column the inbox reads, not the flag array.
unread := models.SeenFromFlags(e.Flags) && email.Seen
+ if unread || (slices.Contains(e.Flags, models.FlagFlagged) && slices.Contains(email.Flags, models.FlagFlagged)) {
+ s.noteOwnerActivity(ctx, e.EmailID, email.InternalDate)
+ }
if len(email.Flags) == 0 && !unread {
return nil
@@ -228,6 +239,22 @@ func (s *JobsService) HandleFlagsRemove(ctx context.Context, e *models.JobEventF
return nil
}
+// ownerActivityArrivalGrace is how long after arrival a filter may still be labelling a message.
+const ownerActivityArrivalGrace = 2 * time.Minute
+
+// noteOwnerActivity records the owner acting on their own mail at the
+// provider. Callers pass only changes our store did not already hold, so a
+// change made in Warmbly and echoed back by the sync never counts, and a
+// change to mail that only just arrived may be a filter finishing delivery.
+func (s *JobsService) noteOwnerActivity(ctx context.Context, accountID uuid.UUID, arrived time.Time) {
+ if s.WarmupRepo == nil || arrived.IsZero() || time.Since(arrived) < ownerActivityArrivalGrace {
+ return
+ }
+ if err := s.WarmupRepo.RecordOwnerActivity(ctx, accountID, time.Now()); err != nil {
+ log.Warn().Err(err).Str("email_id", accountID.String()).Msg("owner activity not recorded")
+ }
+}
+
// containsTrashFlag reports the transition Gmail emits for Delete: the TRASH
// label, passed through untranslated by the worker.
func containsTrashFlag(flags []string) bool {
diff --git a/internal/app/consumer/event_new_email.go b/internal/app/consumer/event_new_email.go
index d5b30bb5d..fb02a9db6 100644
--- a/internal/app/consumer/event_new_email.go
+++ b/internal/app/consumer/event_new_email.go
@@ -447,18 +447,19 @@ func firstSenderAddress(from []string) string {
func (s *JobsService) acceptWarmupEmail(ctx context.Context, e *models.JobEventNewEmail, token *models.WarmupToken) {
s.WarmupRepo.ConsumeWarmupToken(ctx, token.Token)
+ landed := models.ClassifyWarmupLanding(e.Message.Folder, e.Message.Flags)
+
// Record the receipt so a later deletion or spam-flag of THIS message can be
// attributed back to warmup and to the sender. Verified warmup mail is not
// stored in the unibox, so this is the only record that the message was a
// warmup email.
if e.Message != nil {
- if err := s.WarmupRepo.RecordWarmupReceived(ctx, e.Message.EmailID, e.Message.ID, e.Message.MessageID, token.SenderAccountID); err != nil {
+ if err := s.WarmupRepo.RecordWarmupReceived(ctx, e.Message.EmailID, e.Message.ID, e.Message.MessageID, token.SenderAccountID, landed == models.WarmupLandedSpam); err != nil {
log.Warn().Err(err).Str("email_id", e.Message.EmailID.String()).Msg("Failed to record warmup receipt")
}
}
recipient := s.recipientAccount(ctx, e.Message.EmailID)
- landed := models.ClassifyWarmupLanding(e.Message.Folder, e.Message.Flags)
// If the warmup mail arrived in a Junk/Spam state, record a
// spam_placement event against the sender. This is distinct from a
diff --git a/internal/app/consumer/event_update_email.go b/internal/app/consumer/event_update_email.go
index e09868d78..66a364001 100644
--- a/internal/app/consumer/event_update_email.go
+++ b/internal/app/consumer/event_update_email.go
@@ -45,6 +45,7 @@ func (s *JobsService) HandleUpdateEmail(ctx context.Context, e *models.JobEventE
// provider: mail read in the customer's own client is read here too.
if seen := models.SeenFromFlags(e.Flags); seen != email.Seen {
updateData.Seen = &seen
+ s.noteOwnerActivity(ctx, e.EmailID, email.InternalDate)
}
if email.UID != e.UID {
updateData.UID = &e.UID
diff --git a/internal/app/consumer/warmup_retention_test.go b/internal/app/consumer/warmup_retention_test.go
index 6c7973fe2..0f4ce575b 100644
--- a/internal/app/consumer/warmup_retention_test.go
+++ b/internal/app/consumer/warmup_retention_test.go
@@ -19,17 +19,24 @@ import (
// handlers make.
type retentionWarmupRepo struct {
repository.WarmupRepository
- rec *repository.WarmupReceived
+ rec *repository.WarmupReceived
+ held *[]repository.WarmupSpamMove
}
func (r retentionWarmupRepo) GetWarmupReceived(context.Context, uuid.UUID, uuid.UUID) (*repository.WarmupReceived, error) {
return r.rec, nil
}
+func (r retentionWarmupRepo) RecordWarmupSpamMove(_ context.Context, m repository.WarmupSpamMove) (bool, error) {
+ *r.held = append(*r.held, m)
+ return true, nil
+}
+
// retentionWarmupService records which strikes the handlers asked for.
type retentionWarmupService struct {
warmupapp.Service
strikes []string
+ held []repository.WarmupSpamMove
fail bool
}
@@ -54,7 +61,7 @@ func (s *retentionWarmupService) ApplySpamReport(context.Context, uuid.UUID, uui
func retentionService(rec *repository.WarmupReceived) (*JobsService, *retentionWarmupService) {
svc := &retentionWarmupService{}
return &JobsService{
- WarmupRepo: retentionWarmupRepo{rec: rec},
+ WarmupRepo: retentionWarmupRepo{rec: rec, held: &svc.held},
WarmupService: svc,
EmailRepository: warmupInboxEmailRepo{},
}, svc
@@ -275,19 +282,24 @@ func TestRecheckTamperingSearchesEachOldStrike(t *testing.T) {
}
// Gmail reports Delete as gaining the TRASH label. That is the owner's act
-// and is judged on the same freshness rule; a spam flag is still the graver
-// strike and is never subject to the window.
+// and is judged on the same freshness rule. A spam label charges nobody on
+// sight: a move after arrival is held for attribution, and the label on mail
+// that arrived in spam is the filter's own and is not even held.
func TestFlagsAddJudgesGmailTrashOnFreshness(t *testing.T) {
+ landedSpam := receivedAgo(time.Second)
+ landedSpam.LandedSpam = true
cases := []struct {
name string
rec *repository.WarmupReceived
flags []string
want []string
+ held int
}{
- {"trashed an hour after arrival", receivedAgo(time.Hour), []string{"TRASH"}, []string{"deletion"}},
- {"trashed a month after arrival", receivedAgo(30 * 24 * time.Hour), []string{"TRASH"}, nil},
- {"flagged as spam a month after arrival", receivedAgo(30 * 24 * time.Hour), []string{"SPAM"}, []string{"spam_report", "spam_flag"}},
- {"read is not a strike", receivedAgo(time.Hour), []string{models.FlagSeen}, nil},
+ {"trashed an hour after arrival", receivedAgo(time.Hour), []string{"TRASH"}, []string{"deletion"}, 0},
+ {"trashed a month after arrival", receivedAgo(30 * 24 * time.Hour), []string{"TRASH"}, nil, 0},
+ {"moved to spam a month after arrival is held", receivedAgo(30 * 24 * time.Hour), []string{"SPAM"}, nil, 1},
+ {"spam label on mail that arrived in spam", landedSpam, []string{"SPAM"}, nil, 0},
+ {"read is not a strike", receivedAgo(time.Hour), []string{models.FlagSeen}, nil, 0},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
@@ -300,6 +312,9 @@ func TestFlagsAddJudgesGmailTrashOnFreshness(t *testing.T) {
if len(svc.strikes) != len(tc.want) {
t.Fatalf("strikes = %v, want %v", svc.strikes, tc.want)
}
+ if len(svc.held) != tc.held {
+ t.Fatalf("held %d moves, want %d", len(svc.held), tc.held)
+ }
for i := range tc.want {
if svc.strikes[i] != tc.want[i] {
t.Fatalf("strikes = %v, want %v", svc.strikes, tc.want)
diff --git a/internal/app/consumer/warmup_spam_attribution.go b/internal/app/consumer/warmup_spam_attribution.go
new file mode 100644
index 000000000..61d19cbad
--- /dev/null
+++ b/internal/app/consumer/warmup_spam_attribution.go
@@ -0,0 +1,167 @@
+package jobs
+
+import (
+ "context"
+ "fmt"
+ "slices"
+ "time"
+
+ "github.com/rs/zerolog/log"
+
+ "github.com/warmbly/warmbly/internal/config"
+ "github.com/warmbly/warmbly/internal/jobrun"
+ "github.com/warmbly/warmbly/internal/repository"
+)
+
+const (
+ spamMoveAttributionInterval = 5 * time.Minute
+ spamMoveAttributionBatch = 200
+ // spamMoveClaimLease outlives one move's work; a consumer that dies holding it frees it on expiry.
+ spamMoveClaimLease = 5 * time.Minute
+)
+
+// Evidence names stored with each verdict, so an operator can read why.
+const (
+ spamMoveCorrelated = "correlated"
+ spamMoveOnArrival = "on_arrival"
+ spamMoveOwnerActive = "owner_active"
+ spamMoveRepeated = "repeated"
+ spamMoveDormant = "dormant"
+ spamMoveNoOwnerTrace = "no_owner_activity"
+)
+
+// attributeSpamMove decides who moved a received warmup email into spam. No
+// provider says, so the verdict rests on what the pool and the mailbox show:
+// the provider when other workspaces saw the same sender junked or the move
+// came straight after arrival with nobody there, the owner when they were
+// active in the mailbox around it or keep junking pool mail nobody else does,
+// and nobody otherwise. Only the owner is charged.
+func attributeSpamMove(m repository.WarmupSpamMove, ev repository.WarmupSpamMoveEvidence) (string, []string) {
+ quick := m.ObservedAt.Sub(m.ReceivedAt) < time.Duration(config.WarmupSpamMoveQuickMinutes)*time.Minute
+ switch {
+ case ev.CorrelatedElsewhere > 0:
+ return repository.SpamMoveProvider, []string{spamMoveCorrelated}
+ case quick && !ev.OwnerActiveNear:
+ return repository.SpamMoveProvider, []string{spamMoveOnArrival}
+ case ev.OwnerActiveNear:
+ return repository.SpamMoveOwner, []string{spamMoveOwnerActive}
+ case !ev.OwnerActiveRecently:
+ return repository.SpamMoveUnattributed, []string{spamMoveDormant}
+ case ev.PatternSenders >= config.WarmupSpamMovePatternSenders:
+ return repository.SpamMoveOwner, []string{spamMoveRepeated}
+ }
+ return repository.SpamMoveUnattributed, []string{spamMoveNoOwnerTrace}
+}
+
+// StartWarmupSpamMoveAttribution decides each held spam move once it has settled.
+func (s *JobsService) StartWarmupSpamMoveAttribution(ctx context.Context) {
+ if s.WarmupRepo == nil || s.WarmupService == nil {
+ return
+ }
+ jobrun.Loop(ctx, "warmup_spam_move_attribution", spamMoveAttributionInterval, true, func(ctx context.Context) error {
+ batchCtx, cancel := context.WithTimeout(ctx, 2*time.Minute)
+ defer cancel()
+ return s.attributeSpamMoves(batchCtx, time.Now())
+ })
+}
+
+func (s *JobsService) attributeSpamMoves(ctx context.Context, now time.Time) error {
+ settled := now.Add(-time.Duration(config.WarmupSpamMoveSettleMinutes) * time.Minute)
+ moves, err := s.WarmupRepo.ListSettledWarmupSpamMoves(ctx, settled, spamMoveAttributionBatch)
+ if err != nil {
+ return fmt.Errorf("list warmup spam moves: %w", err)
+ }
+ for _, m := range moves {
+ if err := s.attributeOneSpamMove(ctx, m); err != nil {
+ return err
+ }
+ }
+ return nil
+}
+
+// attributeOneSpamMove claims the move, fixes its verdict once and applies it.
+// Effects are idempotent and the move is completed only after all of them, so
+// a failure part way re-applies the same verdict on a later pass.
+func (s *JobsService) attributeOneSpamMove(ctx context.Context, m repository.WarmupSpamMove) error {
+ claimed, err := s.WarmupRepo.ClaimWarmupSpamMove(ctx, m.EmailAccountID, m.MessageID, spamMoveClaimLease)
+ if err != nil {
+ return fmt.Errorf("claim warmup spam move: %w", err)
+ }
+ if !claimed {
+ return nil
+ }
+
+ verdict, signals := m.Verdict, m.Signals
+ if verdict == repository.SpamMovePending {
+ ev, err := s.WarmupRepo.WarmupSpamMoveEvidence(ctx, m)
+ if err != nil {
+ return fmt.Errorf("warmup spam move evidence: %w", err)
+ }
+ verdict, signals = attributeSpamMove(m, ev)
+ fixed, err := s.WarmupRepo.FixWarmupSpamMoveVerdict(ctx, m.EmailAccountID, m.MessageID, verdict, signals)
+ if err != nil {
+ return fmt.Errorf("fix warmup spam move verdict: %w", err)
+ }
+ if !fixed {
+ return nil
+ }
+ }
+
+ if verdict == repository.SpamMoveOwner {
+ hSender, xerr := s.WarmupService.ApplySpamReport(ctx, m.EmailAccountID, m.SenderAccountID, m.MessageID, "user_complaint")
+ if xerr != nil {
+ return fmt.Errorf("record warmup spam complaint: %w", xerr)
+ }
+ s.markRiskBandFromWarmupHealth(ctx, m.SenderAccountID, hSender)
+ hOwner, xerr := s.WarmupService.RecordTampering(ctx, m.EmailAccountID, m.MessageID, "spam_flag")
+ if xerr != nil {
+ return fmt.Errorf("record warmup spam strike: %w", xerr)
+ }
+ s.markRiskBandFromWarmupHealth(ctx, m.EmailAccountID, hOwner)
+ } else {
+ // The provider junked the sender's mail, or may have: a placement
+ // reading against the sender, which only ever slows it down.
+ provider, domain := recipientProviderDomain(s.recipientAccount(ctx, m.EmailAccountID))
+ hSender, xerr := s.WarmupService.RecordSpamPlacement(ctx, m.EmailAccountID, m.SenderAccountID, m.MessageID, "", provider, domain)
+ if xerr != nil {
+ return fmt.Errorf("record warmup spam placement: %w", xerr)
+ }
+ s.markRiskBandFromWarmupHealth(ctx, m.SenderAccountID, hSender)
+ }
+
+ if slices.Contains(signals, spamMoveCorrelated) {
+ if err := s.withdrawCorrelatedOwnerMoves(ctx, m); err != nil {
+ return err
+ }
+ }
+
+ if err := s.WarmupRepo.CompleteWarmupSpamMove(ctx, m.EmailAccountID, m.MessageID); err != nil {
+ return fmt.Errorf("complete warmup spam move: %w", err)
+ }
+ log.Info().Str("email_id", m.EmailAccountID.String()).Str("verdict", verdict).Strs("signals", signals).
+ Msg("warmup spam move attributed")
+ return nil
+}
+
+// withdrawCorrelatedOwnerMoves takes back the strikes charged for the same
+// sender's mail in other workspaces before this move showed the provider at work.
+func (s *JobsService) withdrawCorrelatedOwnerMoves(ctx context.Context, m repository.WarmupSpamMove) error {
+ owners, err := s.WarmupRepo.CorrelatedOwnerSpamMoves(ctx, m.SenderAccountID, m.EmailAccountID, m.ObservedAt)
+ if err != nil {
+ return fmt.Errorf("correlated spam moves: %w", err)
+ }
+ for _, o := range owners {
+ health, xerr := s.WarmupService.WithdrawTampering(ctx, o.EmailAccountID, o.MessageID, "spam_flag")
+ if xerr != nil {
+ return fmt.Errorf("withdraw correlated spam strike: %w", xerr)
+ }
+ s.markRiskBandFromWarmupHealth(ctx, o.EmailAccountID, health)
+ }
+ if len(owners) == 0 {
+ return nil
+ }
+ if err := s.WarmupRepo.ReattributeOwnerSpamMoves(ctx, m.SenderAccountID, m.EmailAccountID, m.ObservedAt); err != nil {
+ return fmt.Errorf("reattribute correlated spam moves: %w", err)
+ }
+ return nil
+}
diff --git a/internal/app/consumer/warmup_spam_attribution_test.go b/internal/app/consumer/warmup_spam_attribution_test.go
new file mode 100644
index 000000000..ab5314e4c
--- /dev/null
+++ b/internal/app/consumer/warmup_spam_attribution_test.go
@@ -0,0 +1,238 @@
+package jobs
+
+import (
+ "context"
+ "slices"
+ "testing"
+ "time"
+
+ "github.com/google/uuid"
+
+ warmupapp "github.com/warmbly/warmbly/internal/app/warmup"
+ "github.com/warmbly/warmbly/internal/config"
+ "github.com/warmbly/warmbly/internal/errx"
+ "github.com/warmbly/warmbly/internal/models"
+ "github.com/warmbly/warmbly/internal/repository"
+)
+
+func spamMoveAfter(sinceArrival time.Duration) repository.WarmupSpamMove {
+ observed := time.Date(2026, 9, 28, 14, 18, 36, 0, time.UTC)
+ return repository.WarmupSpamMove{
+ EmailAccountID: uuid.New(), MessageID: "", SenderAccountID: uuid.New(),
+ ReceivedAt: observed.Add(-sinceArrival), ObservedAt: observed, Verdict: repository.SpamMovePending,
+ }
+}
+
+func TestAttributeSpamMove(t *testing.T) {
+ quick := time.Duration(config.WarmupSpamMoveQuickMinutes)*time.Minute - time.Second
+ later := 6 * time.Hour
+ pattern := config.WarmupSpamMovePatternSenders
+ cases := []struct {
+ name string
+ since time.Duration
+ ev repository.WarmupSpamMoveEvidence
+ want string
+ }{
+ {"the filter catching up with nobody there", quick, repository.WarmupSpamMoveEvidence{OwnerActiveRecently: true}, repository.SpamMoveProvider},
+ {"the owner at the mailbox right after arrival", quick, repository.WarmupSpamMoveEvidence{OwnerActiveNear: true, OwnerActiveRecently: true}, repository.SpamMoveOwner},
+ {"the owner at the mailbox later", later, repository.WarmupSpamMoveEvidence{OwnerActiveNear: true, OwnerActiveRecently: true}, repository.SpamMoveOwner},
+ {"other workspaces junked the same sender", later, repository.WarmupSpamMoveEvidence{OwnerActiveNear: true, OwnerActiveRecently: true, CorrelatedElsewhere: 1}, repository.SpamMoveProvider},
+ {"a mailbox nobody uses", later, repository.WarmupSpamMoveEvidence{PatternSenders: pattern + 5}, repository.SpamMoveUnattributed},
+ {"one unexplained move in a used mailbox", later, repository.WarmupSpamMoveEvidence{OwnerActiveRecently: true, PatternSenders: 1}, repository.SpamMoveUnattributed},
+ {"a used mailbox junking many senders nobody else does", later, repository.WarmupSpamMoveEvidence{OwnerActiveRecently: true, PatternSenders: pattern}, repository.SpamMoveOwner},
+ }
+ for _, tc := range cases {
+ t.Run(tc.name, func(t *testing.T) {
+ got, signals := attributeSpamMove(spamMoveAfter(tc.since), tc.ev)
+ if got != tc.want {
+ t.Fatalf("verdict = %s (%v), want %s", got, signals, tc.want)
+ }
+ if len(signals) == 0 {
+ t.Fatal("a verdict carries no evidence for an operator to read")
+ }
+ })
+ }
+}
+
+// attributionRepo serves settled moves and their evidence, and records what was decided.
+type attributionRepo struct {
+ repository.WarmupRepository
+ moves []repository.WarmupSpamMove
+ evidence repository.WarmupSpamMoveEvidence
+ correlated []repository.WarmupSpamMove
+ claimLost bool
+ log *[]string
+}
+
+func (r attributionRepo) ClaimWarmupSpamMove(context.Context, uuid.UUID, string, time.Duration) (bool, error) {
+ return !r.claimLost, nil
+}
+
+func (r attributionRepo) FixWarmupSpamMoveVerdict(_ context.Context, _ uuid.UUID, _, verdict string, _ []string) (bool, error) {
+ *r.log = append(*r.log, "fix:"+verdict)
+ return true, nil
+}
+
+func (r attributionRepo) CompleteWarmupSpamMove(context.Context, uuid.UUID, string) error {
+ *r.log = append(*r.log, "complete")
+ return nil
+}
+
+func (r attributionRepo) ListSettledWarmupSpamMoves(context.Context, time.Time, int) ([]repository.WarmupSpamMove, error) {
+ return r.moves, nil
+}
+
+func (r attributionRepo) WarmupSpamMoveEvidence(context.Context, repository.WarmupSpamMove) (repository.WarmupSpamMoveEvidence, error) {
+ return r.evidence, nil
+}
+
+func (r attributionRepo) CorrelatedOwnerSpamMoves(context.Context, uuid.UUID, uuid.UUID, time.Time) ([]repository.WarmupSpamMove, error) {
+ return r.correlated, nil
+}
+
+func (r attributionRepo) ReattributeOwnerSpamMoves(context.Context, uuid.UUID, uuid.UUID, time.Time) error {
+ *r.log = append(*r.log, "reattribute")
+ return nil
+}
+
+// attributionService records the effects a verdict has, in order.
+type attributionService struct {
+ warmupapp.Service
+ senderFails bool
+ log *[]string
+}
+
+func (s attributionService) ApplySpamReport(_ context.Context, _, _ uuid.UUID, _, reportType string) (*models.WarmupParticipantHealth, *errx.Error) {
+ if s.senderFails {
+ return nil, errx.InternalError()
+ }
+ *s.log = append(*s.log, "sender:"+reportType)
+ return nil, nil
+}
+
+func (s attributionService) RecordSpamPlacement(context.Context, uuid.UUID, uuid.UUID, string, string, string, string) (*models.WarmupParticipantHealth, *errx.Error) {
+ if s.senderFails {
+ return nil, errx.InternalError()
+ }
+ *s.log = append(*s.log, "sender:spam_placement")
+ return nil, nil
+}
+
+func (s attributionService) RecordTampering(_ context.Context, _ uuid.UUID, _, kind string) (*models.WarmupParticipantHealth, *errx.Error) {
+ *s.log = append(*s.log, "strike:"+kind)
+ return nil, nil
+}
+
+func (s attributionService) WithdrawTampering(_ context.Context, _ uuid.UUID, _, kind string) (*models.WarmupParticipantHealth, *errx.Error) {
+ *s.log = append(*s.log, "withdraw:"+kind)
+ return nil, nil
+}
+
+// Only an owner verdict charges the recipient; the rest read as placement
+// against the sender. The verdict is fixed before any effect and the move is
+// completed after all of them, so a failure part way re-applies the same
+// verdict, and a correlation takes back the owner verdicts it explains first.
+func TestAttributeSpamMovesAppliesTheVerdict(t *testing.T) {
+ owner := repository.WarmupSpamMoveEvidence{OwnerActiveNear: true, OwnerActiveRecently: true}
+ fixedProvider := spamMoveAfter(6 * time.Hour)
+ fixedProvider.Verdict, fixedProvider.Signals = repository.SpamMoveProvider, []string{spamMoveOnArrival}
+ cases := []struct {
+ name string
+ move repository.WarmupSpamMove
+ ev repository.WarmupSpamMoveEvidence
+ correlated int
+ claimLost bool
+ senderFails bool
+ wantErr bool
+ want []string
+ }{
+ {name: "owner", move: spamMoveAfter(6 * time.Hour), ev: owner,
+ want: []string{"fix:owner", "sender:user_complaint", "strike:spam_flag", "complete"}},
+ {name: "provider on arrival", move: spamMoveAfter(time.Minute),
+ want: []string{"fix:provider", "sender:spam_placement", "complete"}},
+ {name: "unattributed", move: spamMoveAfter(6 * time.Hour), ev: repository.WarmupSpamMoveEvidence{OwnerActiveRecently: true, PatternSenders: 1},
+ want: []string{"fix:unattributed", "sender:spam_placement", "complete"}},
+ {name: "correlated withdraws earlier owner verdicts", move: spamMoveAfter(6 * time.Hour), ev: repository.WarmupSpamMoveEvidence{CorrelatedElsewhere: 2}, correlated: 2,
+ want: []string{"fix:provider", "sender:spam_placement", "withdraw:spam_flag", "withdraw:spam_flag", "reattribute", "complete"}},
+ {name: "a move another consumer holds is left alone", move: spamMoveAfter(6 * time.Hour), ev: owner, claimLost: true,
+ want: nil},
+ {name: "a fixed verdict is re-applied, not decided again", move: fixedProvider, ev: owner,
+ want: []string{"sender:spam_placement", "complete"}},
+ {name: "a failed sender write leaves the move to retry", move: spamMoveAfter(6 * time.Hour), ev: owner, senderFails: true, wantErr: true,
+ want: []string{"fix:owner"}},
+ }
+ for _, tc := range cases {
+ t.Run(tc.name, func(t *testing.T) {
+ var log []string
+ repo := attributionRepo{moves: []repository.WarmupSpamMove{tc.move}, evidence: tc.ev, claimLost: tc.claimLost, log: &log}
+ for range tc.correlated {
+ repo.correlated = append(repo.correlated, spamMoveAfter(6*time.Hour))
+ }
+ s := &JobsService{WarmupRepo: repo, WarmupService: attributionService{log: &log, senderFails: tc.senderFails}, EmailRepository: warmupInboxEmailRepo{}}
+ err := s.attributeSpamMoves(context.Background(), time.Now())
+ if (err != nil) != tc.wantErr {
+ t.Fatalf("err = %v, want error %v", err, tc.wantErr)
+ }
+ if !slices.Equal(log, tc.want) {
+ t.Fatalf("effects = %v, want %v", log, tc.want)
+ }
+ })
+ }
+}
+
+// activityRepo is a mailbox with no warmup receipt that records owner activity.
+type activityRepo struct {
+ repository.WarmupRepository
+ noted *int
+}
+
+func (activityRepo) GetWarmupReceived(context.Context, uuid.UUID, uuid.UUID) (*repository.WarmupReceived, error) {
+ return nil, nil
+}
+
+func (r activityRepo) RecordOwnerActivity(context.Context, uuid.UUID, time.Time) error {
+ *r.noted++
+ return nil
+}
+
+// Owner activity is a change the provider reports that our store did not
+// already hold: a read made in Warmbly and echoed back is not the owner at the
+// mailbox, and a label on mail that only just arrived may be a filter.
+func TestOwnerActivityIsOnlyTheProvidersOwnChange(t *testing.T) {
+ old := time.Now().Add(-time.Hour)
+ cases := []struct {
+ name string
+ stored models.EmailMessageStoreData
+ add bool
+ flags []string
+ want int
+ }{
+ {"read at the provider", models.EmailMessageStoreData{Flags: []string{}, InternalDate: old}, true, []string{models.FlagSeen}, 1},
+ {"a read made in Warmbly echoed back", models.EmailMessageStoreData{Flags: []string{models.FlagSeen}, Seen: true, InternalDate: old}, true, []string{models.FlagSeen}, 0},
+ {"read by a filter on arrival", models.EmailMessageStoreData{Flags: []string{}, InternalDate: time.Now()}, true, []string{models.FlagSeen}, 0},
+ {"starred at the provider", models.EmailMessageStoreData{Flags: []string{}, Seen: true, InternalDate: old}, true, []string{models.FlagFlagged}, 1},
+ {"marked unread at the provider", models.EmailMessageStoreData{Flags: []string{models.FlagSeen}, Seen: true, InternalDate: old}, false, []string{models.FlagSeen}, 1},
+ {"a label no person sets", models.EmailMessageStoreData{Flags: []string{}, Seen: true, InternalDate: old}, true, []string{`\Important`}, 0},
+ }
+ for _, tc := range cases {
+ t.Run(tc.name, func(t *testing.T) {
+ stored := tc.stored
+ s, _ := seenSyncService(&stored)
+ noted := 0
+ s.WarmupRepo = activityRepo{noted: ¬ed}
+ ev := &models.JobEventFlags{UserID: uuid.New(), EmailID: uuid.New(), ID: uuid.New(), Flags: tc.flags}
+ var err error
+ if tc.add {
+ err = s.HandleFlagsAdd(context.Background(), ev)
+ } else {
+ err = s.HandleFlagsRemove(context.Background(), ev)
+ }
+ if err != nil {
+ t.Fatal(err)
+ }
+ if noted != tc.want {
+ t.Fatalf("noted %d owner activities, want %d", noted, tc.want)
+ }
+ })
+ }
+}
diff --git a/internal/app/orgtransfer/spec.go b/internal/app/orgtransfer/spec.go
index a490848bf..7ffe4c992 100644
--- a/internal/app/orgtransfer/spec.go
+++ b/internal/app/orgtransfer/spec.go
@@ -905,6 +905,8 @@ var ExcludedTables = map[string]string{
"oauth_authorization_codes": "Single-use authorization codes, valid for seconds.",
"scheduled_deletions": "Instance lifecycle state. Importing a pending deletion would schedule the destination workspace for destruction.",
"dedicated_worker_assignments": "Worker topology, which is a property of the instance rather than the workspace.",
+ "warmup_spam_moves": "Per-message attribution evidence for warmup mail this instance synced, kept only to decide recent tampering; the destination judges its own.",
+ "mailbox_owner_activity": "Five-minute buckets of sync-observed owner activity on this instance, read only to attribute recent spam moves.",
"warmup_pools": "Instance-global pool definitions shared by every workspace on the instance.",
"pool_link_codes": "In-flight link handshakes between a self-hosted instance and this cloud, valid for minutes.",
"cli_auth_codes": "In-flight `warmbly auth login` handshakes, valid for minutes. The API key an approval mints does travel, with the api_keys rows.",
diff --git a/internal/app/warmup/service.go b/internal/app/warmup/service.go
index 46fa0987c..71611579d 100644
--- a/internal/app/warmup/service.go
+++ b/internal/app/warmup/service.go
@@ -59,10 +59,10 @@ const (
minComplaintSample = 100
- // Tampering: harm done to warmup mail the mailbox received, as weighted
- // strikes over the seven-day window (a deletion is one, a spam flag two).
- // One deletion is housekeeping until proven otherwise, so it only warns;
- // the ladder climbs from there and every step lapses on its own.
+ // Tampering: harm done to warmup mail the mailbox received, one strike per
+ // deletion or spam move over the seven-day window. One is housekeeping or a
+ // provider's filter until proven otherwise, so it only warns; the ladder
+ // climbs from there and every step lapses on its own.
tamperingWatchStrikes = 1
tamperingQuarantineStrikes = 2
tamperingBlockStrikes = 4
@@ -446,7 +446,7 @@ func tamperingVerb(kind string) string {
case "deletion":
return "deleted"
case "spam_flag":
- return "marked as spam"
+ return "moved to spam"
default:
return "tampered with"
}
@@ -464,7 +464,7 @@ func tamperingKind(m *models.WarmupHealthMetrics) string {
func tamperingSummary(m *models.WarmupHealthMetrics) string {
parts := []string{}
if m.SpamFlagsLast7d > 0 {
- parts = append(parts, fmt.Sprintf("%d warmup %s marked as spam", m.SpamFlagsLast7d, plural(m.SpamFlagsLast7d, "email", "emails")))
+ parts = append(parts, fmt.Sprintf("%d warmup %s moved to spam", m.SpamFlagsLast7d, plural(m.SpamFlagsLast7d, "email", "emails")))
}
if m.DeletionsLast7d > 0 {
parts = append(parts, fmt.Sprintf("%d warmup %s deleted", m.DeletionsLast7d, plural(m.DeletionsLast7d, "email", "emails")))
@@ -769,9 +769,9 @@ func moreSevere(a, b evaluationDecision) evaluationDecision {
return a
}
-// evaluateTampering needs no sample: each strike is one deliberate act on mail
-// the mailbox verifiably received. A single deletion only warns, because the
-// most likely cause is someone tidying the folder by hand. A deletion is only
+// evaluateTampering needs no sample: each strike is one act on mail the mailbox
+// verifiably received. A single one only warns, because the likeliest cause is
+// someone tidying the folder or a provider filing it as spam. A deletion is only
// recorded inside config.WarmupDeletionStrikeHours of arrival, and only once a
// search of the mailbox found the message in the trash or gone.
func evaluateTampering(metrics *models.WarmupHealthMetrics, now time.Time) evaluationDecision {
diff --git a/internal/app/warmup/service_test.go b/internal/app/warmup/service_test.go
index fe09dc446..91d17059c 100644
--- a/internal/app/warmup/service_test.go
+++ b/internal/app/warmup/service_test.go
@@ -187,10 +187,11 @@ func TestEvaluateMetricsTamperingLadder(t *testing.T) {
{"nothing", 0, 0, models.WarmupHealthHealthy, 0},
{"one deletion warns", 1, 0, models.WarmupHealthWatch, 0},
{"two deletions pause", 2, 0, models.WarmupHealthQuarantined, warmupQuarantineDuration},
- {"one spam flag pauses", 0, 1, models.WarmupHealthQuarantined, warmupQuarantineDuration},
+ {"one spam flag warns", 0, 1, models.WarmupHealthWatch, 0},
+ {"two spam flags pause", 0, 2, models.WarmupHealthQuarantined, warmupQuarantineDuration},
{"four deletions block", 4, 0, models.WarmupHealthBlocked, warmupBlockDuration},
- {"two spam flags block", 0, 2, models.WarmupHealthBlocked, warmupBlockDuration},
- {"a flag and two deletions block", 2, 1, models.WarmupHealthBlocked, warmupBlockDuration},
+ {"four spam flags block", 0, 4, models.WarmupHealthBlocked, warmupBlockDuration},
+ {"a flag and three deletions block", 3, 1, models.WarmupHealthBlocked, warmupBlockDuration},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
@@ -238,7 +239,7 @@ func TestEvaluateMetricsTamperingCombinesWithRates(t *testing.T) {
{"one deletion does not mask a placement throttle", models.WarmupHealthMetrics{DeletionsLast7d: 1, PlacementSample: 20, SpamPlacementRate: 50}, models.WarmupHealthThrottled, func() *time.Time { u := now.Add(warmupThrottleDuration); return &u }()},
{"a placement watch does not mask a tampering quarantine", models.WarmupHealthMetrics{DeletionsLast7d: 2, PlacementSample: 20, SpamPlacementRate: 10}, models.WarmupHealthQuarantined, &quarantine},
{"a complaint-rate quarantine does not mask a tampering block", models.WarmupHealthMetrics{DeletionsLast7d: 4, DeliveredLast30d: 100, ComplaintRate: complaintRateQuarantinePct}, models.WarmupHealthBlocked, &block},
- {"a bounce-rate quarantine does not mask a tampering block", models.WarmupHealthMetrics{SpamFlagsLast7d: 2, DeliveredLast30d: 100, BounceRate: bounceRateQuarantinePct}, models.WarmupHealthBlocked, &block},
+ {"a bounce-rate quarantine does not mask a tampering block", models.WarmupHealthMetrics{SpamFlagsLast7d: 4, DeliveredLast30d: 100, BounceRate: bounceRateQuarantinePct}, models.WarmupHealthBlocked, &block},
{"a warmup-complaint quarantine does not mask a tampering block", models.WarmupHealthMetrics{DeletionsLast7d: 4, SentLast7d: 20, WarmupComplaintRate: warmupComplaintQuarantinePct}, models.WarmupHealthBlocked, &block},
{"a placement throttle does not mask a tampering quarantine", models.WarmupHealthMetrics{DeletionsLast7d: 2, PlacementSample: 20, SpamPlacementRate: spamPlacementThrottlePct}, models.WarmupHealthQuarantined, &quarantine},
{"heavy placement does not soften a tampering block", models.WarmupHealthMetrics{DeletionsLast7d: 4, PlacementSample: 20, SpamPlacementRate: 90}, models.WarmupHealthBlocked, &block},
diff --git a/internal/config/constants.go b/internal/config/constants.go
index 77a72d398..edbb21db1 100644
--- a/internal/config/constants.go
+++ b/internal/config/constants.go
@@ -346,6 +346,22 @@ const (
// strikes behind a live hold are always there to re-decide it.
WarmupTamperingKeepDays = 37
+ // A warmup email moved to spam names no actor on any provider, so the move
+ // is held this long before it is attributed, to see the activity around it.
+ WarmupSpamMoveSettleMinutes = 30
+ // Owner activity this close to a move, either side, attributes it to the owner.
+ WarmupSpamMoveActivityMinutes = 30
+ // A move this soon after arrival, with nobody active, is the filter catching up.
+ WarmupSpamMoveQuickMinutes = 15
+ // The same sender's mail moved to spam in another workspace this close is the provider re-judging it.
+ WarmupSpamMoveCorrelationHours = 24
+ // Unexplained moves from this many distinct senders in seven days, in a mailbox someone uses, are its owner's.
+ WarmupSpamMovePatternSenders = 3
+ // A mailbox with no owner activity this long has nobody to have moved anything.
+ WarmupOwnerDormantDays = 14
+ // Owner activity is kept this long; it only answers the two windows above.
+ WarmupOwnerActivityKeepDays = 30
+
// CampaignSendStampAttempts is how many times the control plane retries the
// sent_at stamp after a send is already on the bus. The reservation is what
// keeps the step from being re-sent, so a lost stamp is a pacing problem,
diff --git a/internal/infrastructure/db/migrations/000233_warmup_spam_move_attribution.down.sql b/internal/infrastructure/db/migrations/000233_warmup_spam_move_attribution.down.sql
new file mode 100644
index 000000000..85058ed4d
--- /dev/null
+++ b/internal/infrastructure/db/migrations/000233_warmup_spam_move_attribution.down.sql
@@ -0,0 +1,3 @@
+DROP TABLE IF EXISTS mailbox_owner_activity;
+DROP TABLE IF EXISTS warmup_spam_moves;
+ALTER TABLE warmup_received DROP COLUMN IF EXISTS landed_spam;
diff --git a/internal/infrastructure/db/migrations/000233_warmup_spam_move_attribution.up.sql b/internal/infrastructure/db/migrations/000233_warmup_spam_move_attribution.up.sql
new file mode 100644
index 000000000..01d30531c
--- /dev/null
+++ b/internal/infrastructure/db/migrations/000233_warmup_spam_move_attribution.up.sql
@@ -0,0 +1,47 @@
+-- landed_spam: the warmup email was in the recipient's spam folder when it arrived, so a spam label on it is the filter's, not the owner's.
+ALTER TABLE warmup_received ADD COLUMN IF NOT EXISTS landed_spam boolean NOT NULL DEFAULT false;
+
+UPDATE warmup_received wr
+SET landed_spam = true
+WHERE wr.message_id <> ''
+ AND EXISTS (
+ SELECT 1 FROM warmup_spam_reports sr
+ WHERE sr.reporter_account_id = wr.email_account_id
+ AND sr.message_id = wr.message_id
+ AND sr.report_type = 'spam_placement'
+ );
+
+-- A spam strike on mail that arrived in spam was never the owner's act.
+DELETE FROM warmup_tampering_events t
+USING warmup_received wr
+WHERE t.kind = 'spam_flag'
+ AND wr.email_account_id = t.email_account_id
+ AND wr.message_id = t.message_id
+ AND wr.landed_spam;
+
+-- A received warmup email moved to spam after arrival, held until the activity around it attributes it.
+CREATE TABLE IF NOT EXISTS warmup_spam_moves (
+ email_account_id uuid NOT NULL REFERENCES email_accounts(id) ON DELETE CASCADE,
+ message_id text NOT NULL,
+ sender_account_id uuid NOT NULL REFERENCES email_accounts(id) ON DELETE CASCADE,
+ received_at timestamptz NOT NULL,
+ observed_at timestamptz NOT NULL DEFAULT now(),
+ verdict text NOT NULL DEFAULT 'pending'
+ CHECK (verdict IN ('pending', 'owner', 'provider', 'unattributed')),
+ signals text[] NOT NULL DEFAULT '{}',
+ -- claimed_until: one consumer holds the move while it applies the verdict; decided_at: the verdict's effects are applied.
+ claimed_until timestamptz,
+ decided_at timestamptz,
+ PRIMARY KEY (email_account_id, message_id)
+);
+CREATE INDEX IF NOT EXISTS idx_warmup_spam_moves_undecided ON warmup_spam_moves (observed_at) WHERE decided_at IS NULL;
+CREATE INDEX IF NOT EXISTS idx_warmup_spam_moves_sender ON warmup_spam_moves (sender_account_id, observed_at);
+
+-- Five-minute buckets in which the owner acted on their own mail at the provider (read, unread, star), never Warmbly's echo.
+CREATE TABLE IF NOT EXISTS mailbox_owner_activity (
+ email_account_id uuid NOT NULL REFERENCES email_accounts(id) ON DELETE CASCADE,
+ bucket timestamptz NOT NULL,
+ events integer NOT NULL DEFAULT 1,
+ PRIMARY KEY (email_account_id, bucket)
+);
+CREATE INDEX IF NOT EXISTS idx_mailbox_owner_activity_bucket ON mailbox_owner_activity (bucket);
diff --git a/internal/models/unibox.go b/internal/models/unibox.go
index bb1793f72..28bab199e 100644
--- a/internal/models/unibox.go
+++ b/internal/models/unibox.go
@@ -595,6 +595,9 @@ type SeenRelayTarget struct {
// isRead.
const FlagSeen = `\Seen`
+// FlagFlagged is the star: Gmail STARRED, IMAP \Flagged.
+const FlagFlagged = `\Flagged`
+
// SeenFromFlags reads a message's read state out of its flags. The stored
// `seen` column has to follow the provider: mail the customer already read in
// their own client is read in Warmbly, and a copy the worker files in Sent
diff --git a/internal/models/warmup.go b/internal/models/warmup.go
index a2db2a9fc..2b5c7c0f6 100644
--- a/internal/models/warmup.go
+++ b/internal/models/warmup.go
@@ -339,13 +339,13 @@ type WarmupHealthMetrics struct {
BounceRate float64 `json:"bounce_rate"`
// DeletionsLast7d and SpamFlagsLast7d are warmup messages this mailbox
- // received and then deleted or flagged as spam. TamperingStrikes weighs
- // them: a spam flag counts double, because nobody flags mail by accident.
+ // received and then deleted or moved to spam. TamperingStrikes weighs them
+ // equally: no provider says who moved a message into spam.
DeletionsLast7d int `json:"deletions_last_7d"`
SpamFlagsLast7d int `json:"spam_flags_last_7d"`
}
// TamperingStrikes is the weighted harm count the tampering band reads.
func (m *WarmupHealthMetrics) TamperingStrikes() int {
- return m.DeletionsLast7d + 2*m.SpamFlagsLast7d
+ return m.DeletionsLast7d + m.SpamFlagsLast7d
}
diff --git a/internal/repository/pg_admin.go b/internal/repository/pg_admin.go
index 37c382173..09b7ecac5 100644
--- a/internal/repository/pg_admin.go
+++ b/internal/repository/pg_admin.go
@@ -1146,7 +1146,13 @@ func (r *adminRepository) BlockAccount(ctx context.Context, accountID uuid.UUID,
// UnblockAccount unblocks an account from warmup pools
func (r *adminRepository) UnblockAccount(ctx context.Context, accountID uuid.UUID) error {
- _, err := r.db.Exec(ctx, `
+ tx, err := r.db.Begin(ctx)
+ if err != nil {
+ return err
+ }
+ defer tx.Rollback(ctx)
+
+ if _, err := tx.Exec(ctx, `
UPDATE warmup_pool_participants
SET blocked_at = NULL,
blocked_reason = NULL,
@@ -1156,7 +1162,19 @@ func (r *adminRepository) UnblockAccount(ctx context.Context, accountID uuid.UUI
last_health_evaluated_at = NOW(),
last_health_score = 0
WHERE email_account_id = $1
- `, accountID)
+ `, accountID); err != nil {
+ return err
+ }
+ if err := forgiveWarmupStrikes(ctx, tx, accountID); err != nil {
+ return err
+ }
+ return tx.Commit(ctx)
+}
+
+// forgiveWarmupStrikes clears the tampering strikes behind a hold an admin
+// lifted, or the next health evaluation reimposes it on the same strikes.
+func forgiveWarmupStrikes(ctx context.Context, tx pgx.Tx, accountID uuid.UUID) error {
+ _, err := tx.Exec(ctx, `DELETE FROM warmup_tampering_events WHERE email_account_id = $1`, accountID)
return err
}
@@ -1297,6 +1315,9 @@ func (r *adminRepository) ReviewAppeal(ctx context.Context, appealID uuid.UUID,
if err != nil {
return err
}
+ if err := forgiveWarmupStrikes(ctx, tx, accountID); err != nil {
+ return err
+ }
_, _ = tx.Exec(ctx, `
INSERT INTO warmup_admin_actions (admin_user_id, email_account_id, action, reason)
diff --git a/internal/repository/pg_warmup.go b/internal/repository/pg_warmup.go
index 17bfafad3..ca5c9c3ad 100644
--- a/internal/repository/pg_warmup.go
+++ b/internal/repository/pg_warmup.go
@@ -92,6 +92,9 @@ type WarmupReceived struct {
// RetiredAt is when the retention sweep sent the deletion for this
// message. A removal observed after that is the platform's own.
RetiredAt *time.Time
+ // LandedSpam is set when the message arrived in the spam folder, so a
+ // spam label on it is the provider's filter, not the owner.
+ LandedSpam bool
}
// WarmupMailToRetire is one warmup message whose retention window has passed,
@@ -236,8 +239,18 @@ type WarmupRepository interface {
// Tampering protection: track delivered warmup mail so a later deletion or
// spam-flag can be attributed, and count "harm" events per mailbox.
- RecordWarmupReceived(ctx context.Context, accountID, internalID uuid.UUID, messageID string, senderAccountID uuid.UUID) error
+ RecordWarmupReceived(ctx context.Context, accountID, internalID uuid.UUID, messageID string, senderAccountID uuid.UUID, landedSpam bool) error
GetWarmupReceived(ctx context.Context, accountID, internalID uuid.UUID) (*WarmupReceived, error)
+ // Spam moves after arrival are held and attributed on the activity around them.
+ RecordWarmupSpamMove(ctx context.Context, m WarmupSpamMove) (bool, error)
+ ListSettledWarmupSpamMoves(ctx context.Context, settledBefore time.Time, limit int) ([]WarmupSpamMove, error)
+ WarmupSpamMoveEvidence(ctx context.Context, m WarmupSpamMove) (WarmupSpamMoveEvidence, error)
+ ClaimWarmupSpamMove(ctx context.Context, accountID uuid.UUID, messageID string, lease time.Duration) (bool, error)
+ FixWarmupSpamMoveVerdict(ctx context.Context, accountID uuid.UUID, messageID, verdict string, signals []string) (bool, error)
+ CompleteWarmupSpamMove(ctx context.Context, accountID uuid.UUID, messageID string) error
+ CorrelatedOwnerSpamMoves(ctx context.Context, senderID, exceptAccountID uuid.UUID, at time.Time) ([]WarmupSpamMove, error)
+ ReattributeOwnerSpamMoves(ctx context.Context, senderID, exceptAccountID uuid.UUID, at time.Time) error
+ RecordOwnerActivity(ctx context.Context, accountID uuid.UUID, at time.Time) 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 did not earn,
@@ -1749,13 +1762,13 @@ func (r *warmupRepository) GetRecentPartnerCounts(ctx context.Context, accountID
// RecordWarmupReceived stores a delivered warmup email keyed by recipient +
// internal message id. Idempotent on re-delivery of the same message.
-func (r *warmupRepository) RecordWarmupReceived(ctx context.Context, accountID, internalID uuid.UUID, messageID string, senderAccountID uuid.UUID) error {
+func (r *warmupRepository) RecordWarmupReceived(ctx context.Context, accountID, internalID uuid.UUID, messageID string, senderAccountID uuid.UUID, landedSpam bool) error {
query := `
- INSERT INTO warmup_received (email_account_id, internal_id, message_id, sender_account_id)
- VALUES ($1, $2, $3, $4)
+ INSERT INTO warmup_received (email_account_id, internal_id, message_id, sender_account_id, landed_spam)
+ VALUES ($1, $2, $3, $4, $5)
ON CONFLICT (email_account_id, internal_id) DO NOTHING
`
- _, err := r.db.Exec(ctx, query, accountID, internalID, messageID, senderAccountID)
+ _, err := r.db.Exec(ctx, query, accountID, internalID, messageID, senderAccountID, landedSpam)
return err
}
@@ -1763,13 +1776,13 @@ func (r *warmupRepository) RecordWarmupReceived(ctx context.Context, accountID,
// message id. Returns nil when the message was not a warmup email.
func (r *warmupRepository) GetWarmupReceived(ctx context.Context, accountID, internalID uuid.UUID) (*WarmupReceived, error) {
query := `
- SELECT email_account_id, internal_id, message_id, sender_account_id, created_at, retired_at
+ SELECT email_account_id, internal_id, message_id, sender_account_id, created_at, retired_at, landed_spam
FROM warmup_received
WHERE email_account_id = $1 AND internal_id = $2
`
var w WarmupReceived
err := r.db.QueryRow(ctx, query, accountID, internalID).Scan(
- &w.EmailAccountID, &w.InternalID, &w.MessageID, &w.SenderAccountID, &w.CreatedAt, &w.RetiredAt,
+ &w.EmailAccountID, &w.InternalID, &w.MessageID, &w.SenderAccountID, &w.CreatedAt, &w.RetiredAt, &w.LandedSpam,
)
if errors.Is(err, sql.ErrNoRows) {
return nil, nil
@@ -1883,6 +1896,11 @@ func (r *warmupRepository) PruneWarmupEventsBefore(ctx context.Context, before t
`DELETE FROM warmup_tampering_events
WHERE created_at < LEAST($1, NOW() - make_interval(days => ` + strconv.Itoa(config.WarmupTamperingKeepDays) + `))`,
`DELETE FROM warmup_spam_reports WHERE created_at < $1`,
+ `DELETE FROM warmup_spam_moves
+ WHERE decided_at IS NOT NULL
+ AND observed_at < LEAST($1, NOW() - make_interval(days => ` + strconv.Itoa(config.WarmupTamperingKeepDays) + `))`,
+ `DELETE FROM mailbox_owner_activity
+ WHERE bucket < NOW() - make_interval(days => ` + strconv.Itoa(config.WarmupOwnerActivityKeepDays) + `)`,
`DELETE FROM warmup_received WHERE created_at < $1 AND retired_at IS NOT NULL`,
`DELETE FROM warmup_tokens
WHERE created_at < $1
diff --git a/internal/repository/pg_warmup_spam_moves.go b/internal/repository/pg_warmup_spam_moves.go
new file mode 100644
index 000000000..9b4aa5298
--- /dev/null
+++ b/internal/repository/pg_warmup_spam_moves.go
@@ -0,0 +1,217 @@
+package repository
+
+import (
+ "context"
+ "time"
+
+ "github.com/google/uuid"
+
+ "github.com/warmbly/warmbly/internal/config"
+)
+
+// Verdicts a warmup spam move is attributed to.
+const (
+ SpamMovePending = "pending"
+ SpamMoveOwner = "owner"
+ SpamMoveProvider = "provider"
+ SpamMoveUnattributed = "unattributed"
+)
+
+// WarmupSpamMove is a received warmup email seen moving into spam after it
+// arrived, which no provider attributes to anyone.
+type WarmupSpamMove struct {
+ EmailAccountID uuid.UUID
+ MessageID string
+ SenderAccountID uuid.UUID
+ ReceivedAt time.Time
+ ObservedAt time.Time
+ // Verdict is pending until fixed; a fixed verdict whose effects were not
+ // all applied is re-applied as it stands, never decided again.
+ Verdict string
+ Signals []string
+}
+
+// WarmupSpamMoveEvidence is what the pool and the mailbox show around one move.
+type WarmupSpamMoveEvidence struct {
+ // OwnerActiveNear: the owner acted on their own mail within
+ // config.WarmupSpamMoveActivityMinutes of the move.
+ OwnerActiveNear bool
+ // OwnerActiveRecently: any owner activity in config.WarmupOwnerDormantDays.
+ OwnerActiveRecently bool
+ // CorrelatedElsewhere counts the same sender's mail moved to spam in other
+ // workspaces within config.WarmupSpamMoveCorrelationHours.
+ CorrelatedElsewhere int
+ // PatternSenders is distinct senders across this mailbox's unexplained
+ // moves in the last seven days, this one included.
+ PatternSenders int
+}
+
+// RecordWarmupSpamMove holds a move for attribution; the first sighting wins.
+func (r *warmupRepository) RecordWarmupSpamMove(ctx context.Context, m WarmupSpamMove) (bool, error) {
+ cmd, err := r.db.Exec(ctx, `
+ INSERT INTO warmup_spam_moves (email_account_id, message_id, sender_account_id, received_at)
+ VALUES ($1, $2, $3, $4)
+ ON CONFLICT (email_account_id, message_id) DO NOTHING`,
+ m.EmailAccountID, m.MessageID, m.SenderAccountID, m.ReceivedAt)
+ if err != nil {
+ return false, err
+ }
+ return cmd.RowsAffected() > 0, nil
+}
+
+// ListSettledWarmupSpamMoves is unclaimed moves observed before settledBefore
+// whose verdict is not yet applied, oldest first.
+func (r *warmupRepository) ListSettledWarmupSpamMoves(ctx context.Context, settledBefore time.Time, limit int) ([]WarmupSpamMove, error) {
+ rows, err := r.db.Query(ctx, `
+ SELECT email_account_id, message_id, sender_account_id, received_at, observed_at, verdict, signals
+ FROM warmup_spam_moves
+ WHERE decided_at IS NULL AND observed_at < $1
+ AND (claimed_until IS NULL OR claimed_until < NOW())
+ ORDER BY observed_at
+ LIMIT $2`, settledBefore, limit)
+ if err != nil {
+ return nil, err
+ }
+ defer rows.Close()
+ var out []WarmupSpamMove
+ for rows.Next() {
+ var m WarmupSpamMove
+ if err := rows.Scan(&m.EmailAccountID, &m.MessageID, &m.SenderAccountID, &m.ReceivedAt, &m.ObservedAt, &m.Verdict, &m.Signals); err != nil {
+ return nil, err
+ }
+ out = append(out, m)
+ }
+ return out, rows.Err()
+}
+
+// ClaimWarmupSpamMove gives one consumer the move for lease; false when another holds it or it is done.
+func (r *warmupRepository) ClaimWarmupSpamMove(ctx context.Context, accountID uuid.UUID, messageID string, lease time.Duration) (bool, error) {
+ cmd, err := r.db.Exec(ctx, `
+ UPDATE warmup_spam_moves
+ SET claimed_until = NOW() + make_interval(secs => $3)
+ WHERE email_account_id = $1 AND message_id = $2 AND decided_at IS NULL
+ AND (claimed_until IS NULL OR claimed_until < NOW())`,
+ accountID, messageID, lease.Seconds())
+ if err != nil {
+ return false, err
+ }
+ return cmd.RowsAffected() > 0, nil
+}
+
+// WarmupSpamMoveEvidence gathers the evidence for one move in one round trip.
+func (r *warmupRepository) WarmupSpamMoveEvidence(ctx context.Context, m WarmupSpamMove) (WarmupSpamMoveEvidence, error) {
+ var ev WarmupSpamMoveEvidence
+ err := r.db.QueryRow(ctx, `
+ SELECT
+ EXISTS (SELECT 1 FROM mailbox_owner_activity a
+ WHERE a.email_account_id = $1
+ AND a.bucket >= $2::timestamptz - make_interval(mins => $4 + 5)
+ AND a.bucket <= $2::timestamptz + make_interval(mins => $4)),
+ EXISTS (SELECT 1 FROM mailbox_owner_activity a
+ WHERE a.email_account_id = $1
+ AND a.bucket >= $2::timestamptz - make_interval(days => $5)),
+ (SELECT COUNT(*) FROM warmup_spam_moves o
+ JOIN email_accounts oa ON oa.id = o.email_account_id
+ WHERE o.sender_account_id = $3
+ AND o.email_account_id <> $1
+ AND oa.organization_id IS DISTINCT FROM (SELECT organization_id FROM email_accounts WHERE id = $1)
+ AND o.observed_at BETWEEN $2::timestamptz - make_interval(hours => $6) AND $2::timestamptz + make_interval(hours => $6)),
+ (SELECT COUNT(DISTINCT s) FROM (
+ SELECT sender_account_id AS s FROM warmup_spam_moves
+ WHERE email_account_id = $1
+ AND verdict IN ('owner', 'unattributed')
+ AND observed_at >= $2::timestamptz - INTERVAL '7 days' AND observed_at <= $2::timestamptz
+ UNION SELECT $3::uuid) p)`,
+ m.EmailAccountID, m.ObservedAt, m.SenderAccountID,
+ config.WarmupSpamMoveActivityMinutes, config.WarmupOwnerDormantDays, config.WarmupSpamMoveCorrelationHours,
+ ).Scan(&ev.OwnerActiveNear, &ev.OwnerActiveRecently, &ev.CorrelatedElsewhere, &ev.PatternSenders)
+ return ev, err
+}
+
+// FixWarmupSpamMoveVerdict records the verdict once, before its effects are applied.
+func (r *warmupRepository) FixWarmupSpamMoveVerdict(ctx context.Context, accountID uuid.UUID, messageID, verdict string, signals []string) (bool, error) {
+ cmd, err := r.db.Exec(ctx, `
+ UPDATE warmup_spam_moves
+ SET verdict = $3, signals = $4
+ WHERE email_account_id = $1 AND message_id = $2 AND verdict = 'pending' AND decided_at IS NULL`,
+ accountID, messageID, verdict, signals)
+ if err != nil {
+ return false, err
+ }
+ return cmd.RowsAffected() > 0, nil
+}
+
+// CompleteWarmupSpamMove marks the verdict's effects applied and releases the claim.
+func (r *warmupRepository) CompleteWarmupSpamMove(ctx context.Context, accountID uuid.UUID, messageID string) error {
+ _, err := r.db.Exec(ctx, `
+ UPDATE warmup_spam_moves
+ SET decided_at = NOW(), claimed_until = NULL
+ WHERE email_account_id = $1 AND message_id = $2 AND decided_at IS NULL`,
+ accountID, messageID)
+ return err
+}
+
+// correlatedOwnerMovesWhere selects the sender's applied owner verdicts in
+// other workspaces near $3, for sender $1 and the correlating mailbox $2. One
+// still being applied keeps its verdict, so a strike always matches its row.
+const correlatedOwnerMovesWhere = `
+ o.sender_account_id = $1
+ AND o.email_account_id <> $2
+ AND o.verdict = 'owner'
+ AND o.decided_at IS NOT NULL
+ AND o.email_account_id IN (
+ SELECT oa.id FROM email_accounts oa
+ WHERE oa.organization_id IS DISTINCT FROM (SELECT organization_id FROM email_accounts WHERE id = $2))
+ AND o.observed_at BETWEEN $3::timestamptz - make_interval(hours => $4) AND $3::timestamptz + make_interval(hours => $4)`
+
+// CorrelatedOwnerSpamMoves lists the owner verdicts a correlation now explains.
+func (r *warmupRepository) CorrelatedOwnerSpamMoves(ctx context.Context, senderID, exceptAccountID uuid.UUID, at time.Time) ([]WarmupSpamMove, error) {
+ rows, err := r.db.Query(ctx, `
+ SELECT o.email_account_id, o.message_id, o.sender_account_id, o.received_at, o.observed_at
+ FROM warmup_spam_moves o
+ WHERE `+correlatedOwnerMovesWhere,
+ senderID, exceptAccountID, at, config.WarmupSpamMoveCorrelationHours)
+ if err != nil {
+ return nil, err
+ }
+ defer rows.Close()
+ var out []WarmupSpamMove
+ for rows.Next() {
+ var m WarmupSpamMove
+ if err := rows.Scan(&m.EmailAccountID, &m.MessageID, &m.SenderAccountID, &m.ReceivedAt, &m.ObservedAt); err != nil {
+ return nil, err
+ }
+ out = append(out, m)
+ }
+ return out, rows.Err()
+}
+
+// ReattributeOwnerSpamMoves turns those moves into provider moves, and the
+// complaint each filed against the sender into placement.
+func (r *warmupRepository) ReattributeOwnerSpamMoves(ctx context.Context, senderID, exceptAccountID uuid.UUID, at time.Time) error {
+ _, err := r.db.Exec(ctx, `
+ WITH moved AS (
+ UPDATE warmup_spam_moves o
+ SET verdict = 'provider', signals = array_append(o.signals, 'correlated_later'), decided_at = NOW()
+ WHERE `+correlatedOwnerMovesWhere+`
+ RETURNING o.email_account_id, o.message_id
+ )
+ UPDATE warmup_spam_reports sr
+ SET report_type = 'spam_placement'
+ FROM moved
+ WHERE sr.reporter_account_id = moved.email_account_id
+ AND sr.message_id = moved.message_id
+ AND sr.report_type = 'user_complaint'`,
+ senderID, exceptAccountID, at, config.WarmupSpamMoveCorrelationHours)
+ return err
+}
+
+// RecordOwnerActivity marks the five-minute bucket the owner acted in.
+func (r *warmupRepository) RecordOwnerActivity(ctx context.Context, accountID uuid.UUID, at time.Time) error {
+ _, err := r.db.Exec(ctx, `
+ INSERT INTO mailbox_owner_activity (email_account_id, bucket)
+ VALUES ($1, to_timestamp(floor(extract(epoch FROM $2::timestamptz) / 300) * 300))
+ ON CONFLICT (email_account_id, bucket) DO UPDATE SET events = mailbox_owner_activity.events + 1`,
+ accountID, at)
+ return err
+}
diff --git a/internal/repository/pool_link_live_test.go b/internal/repository/pool_link_live_test.go
index 729fabd2d..c7e8e662b 100644
--- a/internal/repository/pool_link_live_test.go
+++ b/internal/repository/pool_link_live_test.go
@@ -143,7 +143,7 @@ func TestLivePoolLinkWarmupDeliveryMatchesAndRefuses(t *testing.T) {
// A delivery already recorded, whose token has since been cleaned up.
const receivedID = ""
- if err := f.warmup.RecordWarmupReceived(ctx, f.recipient, uuid.New(), receivedID, f.sender); err != nil {
+ if err := f.warmup.RecordWarmupReceived(ctx, f.recipient, uuid.New(), receivedID, f.sender, false); err != nil {
t.Fatalf("RecordWarmupReceived: %v", err)
}
ok, err = f.warmup.IsWarmupDelivery(ctx, f.recipient, "someone@elsewhere.test", receivedID, "Anything")