Merge pull request #638 from warmbly/feat/warmup-mail-retention

feat: platform-owned retention of warmup mail in mailboxes, per-message record pruning, and a deletion strike limited to the first day
This commit is contained in:
Matthew Meszaros
2026-09-21 08:59:26 +00:00
committed by GitHub
45 changed files with 1601 additions and 97 deletions
+2
View File
@@ -757,6 +757,8 @@ 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
- 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
- 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
Current auto-block thresholds in code:
@@ -85,20 +85,37 @@ const RETENTION_FIELDS = [
key: "engagementDays",
setting: "engagement_event_days",
label: "Opens and clicks (days)",
min: RETENTION_MIN_DAYS,
help: "Per-event open and click logs, with the client, device and approximate location of each. Campaign counts and routing read a separate summary that is never pruned, so shortening this changes what a contact's timeline can show, not what a campaign does.",
},
{
key: "formDays",
setting: "form_event_days",
label: "Form funnel events (days)",
min: RETENTION_MIN_DAYS,
help: "Views, starts, field-level drop-off and submissions for hosted forms. Funnel reports range up to 90 days, so anything below that shortens the report too. Submitted contacts are unaffected.",
},
{
key: "auditDays",
setting: "audit_log_days",
label: "Audit log (days)",
min: RETENTION_MIN_DAYS,
help: "Who did what, from which IP address and user agent, with the change payload. This window is how long that record is held, and it is the one most likely to be set by a retention policy.",
},
{
key: "warmupMailDays",
setting: "warmup_mail_days",
label: "Warmup mail in mailboxes (days)",
min: 3,
help: "How long warmup mail stays in each mailbox before Warmbly deletes it, wherever the mailbox's filing setting keeps it (Trash on Gmail), with the stored copy of its body. A mailbox may set its own window in its drawer; this is the one every other mailbox follows. The floor leaves room for the engagement and a reply in the thread to finish.",
},
{
key: "warmupEventDays",
setting: "warmup_event_days",
label: "Warmup records (days)",
min: 30,
help: "Per-message warmup records: tokens, receipts, tampering events and spam reports. The pool health bands read the last 30 days, which is the floor. The daily sent and received counts behind the analytics are separate and never pruned.",
},
] as const;
type RetentionFieldKey = (typeof RETENTION_FIELDS)[number]["key"];
@@ -108,14 +125,26 @@ const RETENTION_PRESETS = [
{
id: "default",
label: "Defaults",
description: "365 / 180 / 90 days",
values: { engagementDays: "365", formDays: "180", auditDays: "90" },
description: "365 / 180 / 90 days, warmup 30 / 365",
values: {
engagementDays: "365",
formDays: "180",
auditDays: "90",
warmupMailDays: "30",
warmupEventDays: "365",
},
},
{
id: "minimal",
label: "Minimal retention",
description: "30 / 30 / 30 days",
values: { engagementDays: "30", formDays: "30", auditDays: "30" },
description: "30 / 30 / 30 days, warmup 7 / 30",
values: {
engagementDays: "30",
formDays: "30",
auditDays: "30",
warmupMailDays: "7",
warmupEventDays: "30",
},
},
] as const;
@@ -185,6 +214,8 @@ function toForm(s: InstanceSettings): FormState {
engagementDays: String(s.retention.engagement_event_days),
formDays: String(s.retention.form_event_days),
auditDays: String(s.retention.audit_log_days),
warmupMailDays: String(s.retention.warmup_mail_days),
warmupEventDays: String(s.retention.warmup_event_days),
},
tracking: {
machineWindowOpen: String(s.tracking.machine_window_open_seconds),
@@ -269,9 +300,7 @@ export function SettingsTab({ onDirtyChange, onSwitchTab }: SettingsTabProps) {
form !== null && SYNC_FIELDS.every((f) => syncFieldValid(form.sync[f.key], f.min, f.max));
const retentionValid =
form !== null &&
RETENTION_FIELDS.every((f) =>
syncFieldValid(form.retention[f.key], RETENTION_MIN_DAYS, RETENTION_MAX_DAYS),
);
RETENTION_FIELDS.every((f) => syncFieldValid(form.retention[f.key], f.min, RETENTION_MAX_DAYS));
const trackingValid =
form !== null &&
@@ -307,7 +336,7 @@ export function SettingsTab({ onDirtyChange, onSwitchTab }: SettingsTabProps) {
}
if (!retentionValid) {
toast.error(
`Every retention window must be a whole number of days between ${RETENTION_MIN_DAYS} and ${RETENTION_MAX_DAYS.toLocaleString()}`,
`Every retention window must be a whole number of days up to ${RETENTION_MAX_DAYS.toLocaleString()}, and not below the floor shown under it`,
);
return;
}
@@ -336,6 +365,8 @@ export function SettingsTab({ onDirtyChange, onSwitchTab }: SettingsTabProps) {
engagement_event_days: Number(form.retention.engagementDays),
form_event_days: Number(form.retention.formDays),
audit_log_days: Number(form.retention.auditDays),
warmup_mail_days: Number(form.retention.warmupMailDays),
warmup_event_days: Number(form.retention.warmupEventDays),
},
tracking: {
machine_window_open_seconds: Number(form.tracking.machineWindowOpen),
@@ -565,7 +596,7 @@ export function SettingsTab({ onDirtyChange, onSwitchTab }: SettingsTabProps) {
{RETENTION_FIELDS.map((f) => {
const valid = syncFieldValid(
form.retention[f.key],
RETENTION_MIN_DAYS,
f.min,
RETENTION_MAX_DAYS,
);
return (
@@ -590,13 +621,13 @@ export function SettingsTab({ onDirtyChange, onSwitchTab }: SettingsTabProps) {
className="mt-1"
/>
<p className="mt-1 text-xs text-muted-foreground">
{f.help} Between {RETENTION_MIN_DAYS} and{" "}
{f.help} Between {f.min} and{" "}
{RETENTION_MAX_DAYS.toLocaleString()} days.
</p>
{!valid && (
<p className="mt-1 text-xs text-red-600">
Enter a whole number of days between{" "}
{RETENTION_MIN_DAYS} and{" "}
{f.min} and{" "}
{RETENTION_MAX_DAYS.toLocaleString()}.
</p>
)}
@@ -120,6 +120,11 @@ export interface InstanceSettings {
engagement_event_days: number;
form_event_days: number;
audit_log_days: number;
// Warmup mail is deleted from each mailbox after this many days (a
// mailbox may set its own); the per-message warmup records after the
// second. The daily warmup statistics are never pruned.
warmup_mail_days: number;
warmup_event_days: number;
};
// How soon after a send an open or click is recorded as automated. The
// clock starts at dispatch to the worker, so the window also covers the
+4
View File
@@ -473,6 +473,7 @@ func main() {
AdvancedService: advancedService,
InboxTagger: inboxTagger,
Cache: redisCache,
Retention: instancesettings.NewService(instancesettings.NewStore(primaryDB.Pool)),
AdminRepo: repository.NewAdminRepository(primaryDB.Pool),
AssignmentService: workerAssignmentSvc,
Notifier: notificationService,
@@ -512,6 +513,9 @@ func main() {
// effective dwell close to the requested value.
go jobsService.StartWarmupEngagementPoller(ctx, 30*time.Second)
go jobsService.StartWarmupInboxCleanup(ctx)
// Deletes warmup mail past its retention window from the mailbox itself
// and prunes the per-message warmup records after theirs.
go jobsService.StartWarmupMailRetention(ctx)
go jobsService.StartPendingWarmupVerification(ctx)
// Re-offers inbound mail that reply processing never claimed, so a
// reply refused by a since-fixed check is still attributed to its lead.
@@ -161,6 +161,9 @@ Auth: **Scope** `WRITE_EMAILS` · **Org permission** `manage_emails`
| `warmup_start_time` | string | no | Daily warmup window start, `HH:MM`. |
| `warmup_end_time` | string | no | Daily warmup window end, `HH:MM`. |
| `warmup_days` | integer | no | Number of active warmup days per week. |
| `warmup_placement` | string | no | Where warmup mail is filed in the mailbox itself: `folder` (default), `inbox`, or `archive`. See [Warmup](/guides/warmup/#in-your-mail-client). |
| `warmup_folder` | string | no | Folder (Gmail label) warmup mail is filed into when `warmup_placement` is `folder`. Empty string means the instance default, `Warmbly`. No folder separators or control characters; 64 characters max. |
| `warmup_retention_days` | integer | no | How many days warmup mail stays in this mailbox before Warmbly deletes it, wherever `warmup_placement` keeps it: `3` to `3650`, or `0` to follow the instance setting (30 days unless the operator changed it). |
| `tags` | string[] | no | Tag ids assigned to the mailbox. |
```json
@@ -746,6 +746,8 @@ These are the only settings a browser can change, and no environment variable ow
| `retention.engagement_event_days` | integer, 1 to 3650 | `365` | How long the per-event open and click logs (client, device, approximate location) are kept. Campaign progress keeps its own summary that outlives them, so counts, filters and branching never change |
| `retention.form_event_days` | integer, 1 to 3650 | `180` | How long form funnel events (views, starts, field-level drop-off) are kept. Funnel reports range up to 90 days, so anything shorter shortens the report too |
| `retention.audit_log_days` | integer, 1 to 3650 | `90` | How long the audit trail is kept. It carries IP addresses, user agents and change payloads, so this is also how long that data is held |
| `retention.warmup_mail_days` | integer, 3 to 3650 | `30` | How long warmup mail stays in each mailbox before Warmbly deletes it, wherever the mailbox's filing setting keeps it (Trash on Gmail), with the stored copy of its body. A mailbox can set its own window on the Warmup tab of its drawer or through `warmup_retention_days` on the API; this is what every other mailbox follows. The floor leaves room for the engagement and for a reply in the thread to finish |
| `retention.warmup_event_days` | integer, 30 to 3650 | `365` | How long the per-message warmup records (tokens, receipts, tampering events, spam reports) are kept. The pool health bands read the last 30 days, which is the floor. The daily sent and received counts behind the warmup analytics are separate and never pruned |
| `tracking.machine_window_open_seconds` | integer, 1 to 900 | `60` | How soon after a send was dispatched an open is recorded as automated rather than a person's. The clock starts when the send is handed to a worker, so this window also covers the sending provider's queue and the transit to the recipient, not just reading time. Raise it when delivery-time scanners are being counted as opens; lower it when recipients who read immediately are being missed. A change applies within a minute and only to events recorded after it |
| `tracking.machine_window_click_seconds` | integer, 1 to 900 | `30` | The same window for click tickets. Kept separate because the two mistakes cost different things: a misjudged open loses a metric, a misjudged click loses the automation behind an interested lead |
| `tracking.machine_window_probable_seconds` | integer, 1 to 86400 | `600` | The window used instead of the two above when the tracking service recognised the source as a mail-security network that also renders clicked pages for people, which Proofpoint, Mimecast and Cisco do through browser isolation. Such a match cannot settle whether a person is behind the request, so it widens the window rather than deciding: inside it the event is the delivery-time scan, past it the recipient who got to the mail later. Never applied shorter than the window for the kind of event in hand, so naming a network can only ever catch more scans. Its ceiling reaches a day because how long a vendor takes to detonate a link is the vendor's property, not the instance's |
@@ -754,7 +756,7 @@ These are the only settings a browser can change, and no environment variable ow
The four `sync.*` values are read by the backend when a mailbox is loaded onto a worker (on connect, on reassignment, and by the reconciler's periodic republish), so a change reaches every mailbox within a few minutes without a restart. The fixed pacing numbers around them (burst per five minutes, hourly, backfill pace, the flood threshold and the chronic-overage rule) are compiled constants listed under **Instance > Configuration > Effective limits**; see [Mailboxes](/guides/mailboxes/#what-gets-synced) for how the budgets behave. The four `sync.*` values also have a companion read view on **Operations > Sync**, which shows each mailbox's backfill progress and fair-use throttle against them and can clear a throttle or restart a backfill.
The three `retention.*` values are read by the pruning sweeps on every pass, so shortening one takes effect on the next sweep rather than at the next restart. Deletion is permanent and there is no grace period: what already sits outside a shortened window goes on that sweep. See [data control](/development/data-control/#what-is-kept-and-for-how-long) for what each log holds and what a shorter window costs.
The five `retention.*` values are read by the pruning sweeps on every pass, so shortening one takes effect on the next sweep rather than at the next restart. Deletion is permanent and there is no grace period: what already sits outside a shortened window goes on that sweep. The warmup mail window is the one that reaches into customers' mailboxes: the consumer retires each message and the worker holding the mailbox deletes it there, so shortening it empties the warmup folders of the whole instance down to the new window within a few hours. See [data control](/development/data-control/#what-is-kept-and-for-how-long) for what each log holds and what a shorter window costs.
An unattended install can seed the whole document before anyone signs in, with `WARMBLY_SETTINGS_BOOTSTRAP` holding the same partial JSON the admin API takes:
@@ -94,8 +94,10 @@ See [Automatic inbox tagging](/guides/inbox-tagging/), [Advisor](/guides/advisor
| `retention.engagement_event_days` | 365 | Per-event open and click logs: client, device, approximate location |
| `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 |
Each is between 1 and 3,650 days. These are the three settings a retention or privacy policy applies to, because each window is also how long the personal data in that log is held.
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.
None of them change a number anyone reads. Campaign progress keeps its own summary of opens and clicks that outlives the per-event log, so counts, filters and branching are unaffected by shortening any of these. What gets shorter is what a contact's timeline can show, how far a funnel report reaches, and how far back an admin can audit.
@@ -103,7 +105,11 @@ None of them change a number anyone reads. Campaign progress keeps its own summa
There is no grace period and no copy. Take a backup first if you are not sure.
</Callout>
The **minimal retention** preset in the admin panel and in the installer sets all three to 30 days.
The **minimal retention** preset in the admin panel and in the installer sets the three event logs to 30 days, warmup mail to 7 and warmup records to 30.
### Warmup mail in mailboxes
Warmup mail is real mail in a real mailbox, filed into the warmup folder in both directions, and it is the platform's job to clear it. Once a message is older than the mailbox's window it is deleted where it sits: moved to Trash on Gmail, which Gmail empties after 30 days, and removed outright on Outlook and IMAP. The platform's stored copy of the body goes with it. Nothing about a deletion the platform made is ever held against the mailbox, and a message deleted by its owner after the first day is treated the same way. See [Warmup](/guides/warmup/#retention) for what the mailbox owner sees.
### Warmup standing by address
+3 -1
View File
@@ -176,8 +176,10 @@ Two groups, both editable afterwards in **Instance > Configuration > Settings**:
| Opens and clicks | 365 days | Per-event logs, with the client, device and approximate location of each |
| Form funnel events | 180 days | Views, starts, field-level drop-off |
| Audit log | 90 days | Who did what, from which IP address and user agent |
| Warmup mail in mailboxes | 30 days | Warmup mail is deleted from each mailbox once older than this, wherever it is filed |
| Warmup records | 365 days | Per-message warmup tokens, receipts, tampering events and spam reports |
The last three are also how long that personal data is held. A "minimal retention" preset sets every event window to 30 days.
The opens, forms and audit windows are also how long that personal data is held. A "minimal retention" preset sets those three to 30 days, warmup mail to 7 and warmup records to 30.
</Step>
+11 -1
View File
@@ -142,6 +142,14 @@ One case reaches back. If a warmup message is ever missed on arrival and shows u
The sweep repairs what it can still see, so it does not clear a backlog that built up before the setting existed. To do that on Gmail, search `label:Warmbly` and archive the results: warmup mail has always carried the label, so that selects exactly it. On Outlook and IMAP the received mail is already in the folder, and only the Sent copies from before this change stay in Sent.
### Retention
Warmup mail is real mail taking up real space, and once its engagement has been recorded there is nothing left to keep it for. So Warmbly clears the folder itself: every warmup message, received or sent, is deleted from the mailbox once it is older than the mailbox's retention window, 30 days unless you or your instance operator chose otherwise. On Gmail that means it is moved to Trash, which Gmail empties after another 30 days; on Outlook, Microsoft 365 and IMAP it is deleted outright. The copy of the message Warmbly itself stored goes with it.
The window is per mailbox, next to the filing setting on the **Warmup** tab of its drawer, and on the API as `warmup_retention_days`. Leave it empty to follow the instance setting; the shortest is 3 days, so a thread always has time to finish and the delayed engagement (read, starred, marked important) has always run before its message goes. A self-hosted instance sets its default under **Instance > Configuration > Data retention**, see [data control](/development/data-control/#warmup-mail-in-mailboxes).
At the volumes warmup runs, a few plaintext messages of a couple of kilobytes each, a mailbox in balance takes years to fill even on a small hosting quota, so the retention is there so that you never have to think about it rather than because the folder is about to overflow. You do not need to clear the folder by hand, and if you do, see [pool safety](#pool-safety) for what counts.
<Callout type="info" title="Gmail and the Sent copy">
Gmail exposes folders as labels and does not document whether `SENT` may be removed. Warmbly asks once per mailbox: where Gmail allows it, the sent copy leaves Sent entirely; where it refuses, the copy keeps its place in Sent and carries the warmup label, and nothing asks again. Outlook, Microsoft 365 and IMAP mailboxes move the copy out of Sent normally. SMTP mailboxes never get one in the first place, because Warmbly does not file a Sent copy for warmup.
</Callout>
@@ -179,7 +187,9 @@ Shared pools only work if participants behave, so every mailbox is judged on rat
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 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 tidying the Warmbly folder by hand once is a warning, not a ban, while flagging pool mail as spam twice is a block. A mailbox's own filing is never counted: Warmbly moving a warmup email into its folder, marking it read or rescuing it from spam is the platform acting, not the owner.
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.
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.
<Callout type="warn" title="Acting early protects everyone">
Warmbly intervenes well before providers would penalize a mailbox. One landing in spam a large share of the time is already dangerous to the pool, so it is throttled or removed rather than left collecting negative signals.
+26
View File
@@ -21425,6 +21425,19 @@
"type": "string",
"description": "Folder (Gmail label) warmup mail is filed into when warmup_placement is \"folder\". Empty means the instance default, \"Warmbly\"."
},
"warmup_retention_days": {
"type": "integer",
"anyOf": [
{
"const": 0
},
{
"minimum": 3,
"maximum": 3650
}
],
"description": "How many days warmup mail stays in this mailbox before Warmbly deletes it, wherever warmup_placement keeps it. 0 means the instance setting (30 days by default)."
},
"timezone": {
"type": "string"
},
@@ -21568,6 +21581,19 @@
"type": "string",
"description": "Folder (Gmail label) warmup mail is filed into when warmup_placement is \"folder\". Send an empty string to fall back to the instance default. No folder separators or control characters; 64 characters max."
},
"warmup_retention_days": {
"type": "integer",
"anyOf": [
{
"const": 0
},
{
"minimum": 3,
"maximum": 3650
}
],
"description": "How many days warmup mail stays in this mailbox before Warmbly deletes it, wherever warmup_placement keeps it: 3 to 3650, or 0 to follow the instance setting."
},
"tags": {
"type": "array",
"items": {
+32 -23
View File
@@ -33,18 +33,21 @@ func (d Deps) registerMailboxTools(r *Registry) {
Name: "update_mailbox",
Description: "Update a mailbox's settings: display name, reply-to, cold-send cap (campaign_limit), minimum gap between sends, status, and warmup parameters. Only provided fields change.",
InputSchema: objectSchema(map[string]any{
"email_account_id": strProp("The mailbox UUID."),
"name": strProp("Display name."),
"reply_to": strProp("Reply-to address."),
"status": enumProp("Mailbox status.", "active", "inactive"),
"campaign_limit": intProp("Max cold-campaign emails per day for this mailbox, 0 to 5000. Default 50; 30-50/day is the safe cold-outreach band."),
"min_wait_time": intProp("Minimum seconds between sends."),
"warmup": boolProp("Enable or disable warmup."),
"warmup_base": intProp("Warmup starting emails/day."),
"warmup_max": intProp("Warmup ceiling emails/day."),
"warmup_increase": intProp("Warmup daily ramp increment."),
"warmup_reply_rate": intProp("Warmup reply rate percent."),
"warmup_days": intProp("Warmup active days bitmask."),
"email_account_id": strProp("The mailbox UUID."),
"name": strProp("Display name."),
"reply_to": strProp("Reply-to address."),
"status": enumProp("Mailbox status.", "active", "inactive"),
"campaign_limit": intProp("Max cold-campaign emails per day for this mailbox, 0 to 5000. Default 50; 30-50/day is the safe cold-outreach band."),
"min_wait_time": intProp("Minimum seconds between sends."),
"warmup": boolProp("Enable or disable warmup."),
"warmup_base": intProp("Warmup starting emails/day."),
"warmup_max": intProp("Warmup ceiling emails/day."),
"warmup_increase": intProp("Warmup daily ramp increment."),
"warmup_reply_rate": intProp("Warmup reply rate percent."),
"warmup_days": intProp("Warmup active days bitmask."),
"warmup_placement": enumProp("Where warmup mail is filed in the mailbox itself.", "folder", "inbox", "archive"),
"warmup_folder": strProp("Folder (Gmail label) for warmup mail when warmup_placement is folder. Empty string means the instance default, Warmbly."),
"warmup_retention_days": intProp("Days warmup mail stays in the mailbox before Warmbly deletes it, wherever warmup_placement keeps it: 3 to 3650, or 0 to follow the instance setting."),
}, "email_account_id"),
Risk: generation.RiskWrite,
RequiredOrgPerm: models.PermManageEmails,
@@ -147,6 +150,9 @@ func (d Deps) updateMailbox(ctx context.Context, inv Invocation, args json.RawMe
WarmupIncrease *int `json:"warmup_increase"`
WarmupReplyRate *int `json:"warmup_reply_rate"`
WarmupDays *int `json:"warmup_days"`
WarmupPlacement *string `json:"warmup_placement"`
WarmupFolder *string `json:"warmup_folder"`
WarmupRetention *int `json:"warmup_retention_days"`
}](args)
if err != nil {
return "", err
@@ -156,17 +162,20 @@ func (d Deps) updateMailbox(ctx context.Context, inv Invocation, args json.RawMe
return "", err
}
upd := &models.UpdateEmail{
Name: in.Name,
ReplyTo: in.ReplyTo,
Status: in.Status,
CampaignLimit: in.CampaignLimit,
MinWaitTime: in.MinWaitTime,
Warmup: in.Warmup,
WarmupBase: in.WarmupBase,
WarmupMax: in.WarmupMax,
WarmupIncrease: in.WarmupIncrease,
WarmupReplyRate: in.WarmupReplyRate,
WarmupDays: in.WarmupDays,
Name: in.Name,
ReplyTo: in.ReplyTo,
Status: in.Status,
CampaignLimit: in.CampaignLimit,
MinWaitTime: in.MinWaitTime,
Warmup: in.Warmup,
WarmupBase: in.WarmupBase,
WarmupMax: in.WarmupMax,
WarmupIncrease: in.WarmupIncrease,
WarmupReplyRate: in.WarmupReplyRate,
WarmupDays: in.WarmupDays,
WarmupPlacement: in.WarmupPlacement,
WarmupFolder: in.WarmupFolder,
WarmupRetentionDays: in.WarmupRetention,
}
mb, xerr := d.Emails.Update(ctx, inv.OrgID.String(), inv.UserID.String(), in.EmailAccountID, upd)
if xerr != nil {
+23 -1
View File
@@ -5,6 +5,7 @@ import (
"fmt"
"slices"
"strings"
"time"
"github.com/google/uuid"
@@ -21,11 +22,21 @@ func (s *JobsService) HandleFlagsAdd(ctx context.Context, e *models.JobEventFlag
// tracked in the unibox, so there's nothing else to do for it.
if s.WarmupRepo != nil {
if rec, _ := s.WarmupRepo.GetWarmupReceived(ctx, e.EmailID, e.ID); rec != nil {
if containsSpamFlag(e.Flags) && s.WarmupService != nil {
switch {
case s.WarmupService == nil:
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)
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.
// The label is the owner's act, so it is judged here, on the
// same freshness rule as a removal; the later purge is then
// outside the window and reads as housekeeping.
hHarmer, _ := s.WarmupService.RecordTampering(ctx, e.EmailID, rec.MessageID, "deletion")
s.markRiskBandFromWarmupHealth(ctx, e.EmailID, hHarmer)
}
return nil
}
@@ -216,3 +227,14 @@ func (s *JobsService) HandleFlagsRemove(ctx context.Context, e *models.JobEventF
s.publishEmailUpdated(ctx, e.UserID, email)
return nil
}
// containsTrashFlag reports the transition Gmail emits for Delete: the TRASH
// label, passed through untranslated by the worker.
func containsTrashFlag(flags []string) bool {
for _, f := range flags {
if f == "TRASH" || f == "\\Trash" {
return true
}
}
return false
}
+28 -3
View File
@@ -2,19 +2,24 @@ package jobs
import (
"context"
"time"
"github.com/rs/zerolog/log"
"github.com/warmbly/warmbly/internal/config"
"github.com/warmbly/warmbly/internal/infrastructure/pubsub"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/repository"
)
// HandleRemoveEmail processes a message removal observed during mailbox sync.
//
// Tampering protection: if the removed message was a warmup email (tracked in
// warmup_received), the recipient deleted pool warmup mail. That is recorded
// as a strike and the health bands decide: one deletion warns, more pauses or
// blocks. The owner can appeal a block.
// warmup_received) and it went soon after it arrived, the recipient deleted
// pool warmup mail before its engagement was earned. That is recorded as a
// strike and the health bands decide: one deletion warns, more pauses or
// blocks. The owner can appeal a block. A removal later than that is
// housekeeping (see warmupDeletionCounts) and is not held against anyone.
//
// It also drops the local unibox entry for the removed message (best-effort).
func (s *JobsService) HandleRemoveEmail(ctx context.Context, e *models.JobEventRemoveEmail) error {
@@ -29,6 +34,11 @@ func (s *JobsService) HandleRemoveEmail(ctx context.Context, e *models.JobEventR
Str("email_id", e.EmailID.String()).
Str("message_id", rec.MessageID).
Msg("Warmup message left its folder because we moved it; not tampering")
case !warmupDeletionCounts(rec, time.Now()):
log.Debug().
Str("email_id", e.EmailID.String()).
Str("message_id", rec.MessageID).
Msg("Warmup message removed after its engagement window; housekeeping, not tampering")
case s.WarmupService != nil:
health, _ := s.WarmupService.RecordTampering(ctx, e.EmailID, rec.MessageID, "deletion")
s.markRiskBandFromWarmupHealth(ctx, e.EmailID, health)
@@ -56,3 +66,18 @@ func (s *JobsService) HandleRemoveEmail(ctx context.Context, e *models.JobEventR
}
return nil
}
// warmupDeletionCounts decides whether a deletion of a received warmup message
// is tampering. It is when the message is still fresh: the engagement legs
// run inside the first hours, and removing the mail before then costs the
// pool the signal it was sent for. Past config.WarmupDeletionStrikeHours the
// platform's own retention is going to delete it anyway, so an owner tidying
// the folder, Gmail purging its Trash or a server retention rule is doing the
// platform's job early, not harm. A message the retention sweep has already
// retired is the platform's own deletion whenever it is observed.
func warmupDeletionCounts(rec *repository.WarmupReceived, now time.Time) bool {
if rec == nil || rec.RetiredAt != nil {
return false
}
return now.Sub(rec.CreatedAt) < time.Duration(config.WarmupDeletionStrikeHours)*time.Hour
}
+5
View File
@@ -71,6 +71,11 @@ type JobsService struct {
// Cache for dead worker detection
Cache *cache.Cache
// Retention is the operator-editable retention section, read by the
// warmup mail retention sweep on every pass. Nil keeps the compiled
// defaults.
Retention RetentionSource
// AdminRepo for writing audit-log rows when the dead-worker job
// auto-reassigns email accounts (optional — heartbeat sync also writes
// here so admins can see why their fleet moved). Nil disables logging.
@@ -0,0 +1,166 @@
package jobs
import (
"context"
"errors"
"time"
"github.com/rs/zerolog/log"
"github.com/warmbly/warmbly/internal/config"
"github.com/warmbly/warmbly/internal/jobrun"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/repository"
)
// Warmup mail retention.
//
// Warmup mail is real mail in the customer's mailbox, and once its engagement
// has been recorded there is nothing left to keep it for. A mailbox on a fixed
// quota, which is every mailbox outside Google Workspace, would otherwise fill
// with it, and a full mailbox receives nothing at all. So the platform deletes
// it: every received copy and every sender's own copy past the mailbox's
// window (its own, else retention.warmup_mail_days) is retired here and a
// delete action goes to the worker holding the mailbox, which removes it from
// the provider and drops the stored body.
//
// The row is retired only after the action is on the bus, so a publish that
// fails leaves the message to be offered again next pass. The removal the
// sync then observes is outside config.WarmupDeletionStrikeHours by
// construction (the window floor is days), and the retired stamp says whose
// deletion it was, so it can never be a strike.
//
// The same loop prunes the per-message warmup records past
// retention.warmup_event_days once a full pass has completed.
const (
warmupRetentionBatch = 100
// warmupRetentionPassEvery is how often a completed pass is repeated. The
// windows are whole days, so nothing is gained by walking the tables
// more often than this.
warmupRetentionPassEvery = 6 * time.Hour
)
// StartWarmupMailRetention runs the retention sweep in bounded batches, and
// the record prune after each complete pass.
func (s *JobsService) StartWarmupMailRetention(ctx context.Context) {
if s.WarmupRepo == nil || s.Publisher == nil {
return
}
var nextPass time.Time
jobrun.Loop(ctx, "warmup_mail_retention", time.Minute, true, func(ctx context.Context) error {
if time.Now().Before(nextPass) {
return nil
}
batchCtx, cancel := context.WithTimeout(ctx, 45*time.Second)
defer cancel()
done, err := s.retireWarmupMailBatch(batchCtx)
if err != nil || !done {
return err
}
if err := s.pruneWarmupEvents(batchCtx); err != nil {
return err
}
nextPass = time.Now().Add(warmupRetentionPassEvery)
return nil
})
}
// warmupRetentionDays is the instance window every mailbox without one of its
// own follows, read on every batch so an admin edit applies to the next one.
func (s *JobsService) warmupRetentionDays(ctx context.Context) (mail, events int) {
mail, events = config.WarmupMailRetentionDaysDefault, config.WarmupEventRetentionDaysDefault
if s.Retention != nil {
r := s.Retention.RetentionWindows(ctx)
if r.WarmupMailDays >= config.WarmupMailRetentionDaysMin {
mail = r.WarmupMailDays
}
if r.WarmupEventDays >= config.WarmupEventRetentionDaysMin {
events = r.WarmupEventDays
}
}
return mail, events
}
// retireWarmupMailBatch offers one batch of received copies and one of sent
// copies to the workers. It reports done when both listings came back short,
// which means nothing older is waiting.
func (s *JobsService) retireWarmupMailBatch(ctx context.Context) (bool, error) {
days, _ := s.warmupRetentionDays(ctx)
received, err := s.WarmupRepo.ListWarmupMailToRetire(ctx, days, warmupRetentionBatch)
if err != nil {
return false, err
}
var failures []error
for i := range received {
m := &received[i]
if err := s.publishWarmupDelete(ctx, m); err != nil {
failures = append(failures, err)
continue
}
if err := s.WarmupRepo.RetireWarmupReceived(ctx, m.EmailAccountID, m.InternalID); err != nil {
failures = append(failures, err)
}
}
sent, err := s.WarmupRepo.ListWarmupSentCopiesToRetire(ctx, days, warmupRetentionBatch)
if err != nil {
return false, errors.Join(append(failures, err)...)
}
for i := range sent {
m := &sent[i]
if err := s.publishWarmupDelete(ctx, m); err != nil {
failures = append(failures, err)
continue
}
if err := s.WarmupRepo.RetireWarmupSentCopy(ctx, m.Token); err != nil {
failures = append(failures, err)
}
}
if len(failures) > 0 {
// A row whose publish failed was not retired and is offered again;
// the pass is not complete until every row of the batch went.
return false, errors.Join(failures...)
}
return len(received) < warmupRetentionBatch && len(sent) < warmupRetentionBatch, nil
}
// publishWarmupDelete sends the delete for one message to the worker holding
// its mailbox. The action carries every key the worker can find the message
// by, and where the mailbox files warmup, so the search starts in the folder
// the message is most likely in.
func (s *JobsService) publishWarmupDelete(ctx context.Context, m *repository.WarmupMailToRetire) error {
filing := models.Email{WarmupPlacement: m.Placement, WarmupFolder: m.Folder}
placement, folder := filing.WarmupFiling()
action := &models.WarmupEmailAction{
UserID: m.UserID,
EmailID: m.EmailAccountID,
GmailID: m.ProviderKey,
RFCMessageID: m.MessageID,
Actions: []string{models.WarmupActionDelete},
Placement: placement,
TargetFolder: folder,
}
if m.InternalID != [16]byte{} {
action.InternalID = m.InternalID.String()
}
return s.Publisher.PublishWarmupAction(ctx, m.WorkerID, action)
}
// pruneWarmupEvents drops the per-message warmup records past the instance
// window. Runs after a complete retention pass, so a record is only ever
// pruned once the mail it describes has been retired.
func (s *JobsService) pruneWarmupEvents(ctx context.Context) error {
_, days := s.warmupRetentionDays(ctx)
before := time.Now().AddDate(0, 0, -days)
pruned, err := s.WarmupRepo.PruneWarmupEventsBefore(ctx, before)
if err != nil {
return err
}
if pruned > 0 {
log.Info().Int64("pruned", pruned).Int("days", days).Msg("warmup: per-message records past the retention window pruned")
}
return nil
}
@@ -0,0 +1,255 @@
package jobs
import (
"context"
"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"
)
// retentionWarmupRepo answers the one receipt lookup the removal and flag
// handlers make.
type retentionWarmupRepo struct {
repository.WarmupRepository
rec *repository.WarmupReceived
}
func (r retentionWarmupRepo) GetWarmupReceived(context.Context, uuid.UUID, uuid.UUID) (*repository.WarmupReceived, error) {
return r.rec, nil
}
// retentionWarmupService records which strikes the handlers asked for.
type retentionWarmupService struct {
warmupapp.Service
strikes []string
}
func (s *retentionWarmupService) RecordTampering(_ context.Context, _ uuid.UUID, _, kind string) (*models.WarmupParticipantHealth, *errx.Error) {
s.strikes = append(s.strikes, kind)
return nil, nil
}
func (s *retentionWarmupService) ApplySpamReport(context.Context, uuid.UUID, uuid.UUID, string, string) (*models.WarmupParticipantHealth, *errx.Error) {
s.strikes = append(s.strikes, "spam_report")
return nil, nil
}
func retentionService(rec *repository.WarmupReceived) (*JobsService, *retentionWarmupService) {
svc := &retentionWarmupService{}
return &JobsService{
WarmupRepo: retentionWarmupRepo{rec: rec},
WarmupService: svc,
EmailRepository: warmupInboxEmailRepo{},
}, svc
}
func receivedAgo(age time.Duration) *repository.WarmupReceived {
return &repository.WarmupReceived{
EmailAccountID: uuid.New(), InternalID: uuid.New(),
MessageID: "<warmup@example.test>", SenderAccountID: uuid.New(),
CreatedAt: time.Now().Add(-age),
}
}
func TestWarmupDeletionCounts(t *testing.T) {
now := time.Now()
window := time.Duration(config.WarmupDeletionStrikeHours) * time.Hour
retired := now.Add(-time.Hour)
cases := []struct {
name string
rec *repository.WarmupReceived
want bool
}{
{"no receipt is not warmup", nil, false},
{"fresh is a strike", &repository.WarmupReceived{CreatedAt: now.Add(-time.Hour)}, true},
{"just inside the window is a strike", &repository.WarmupReceived{CreatedAt: now.Add(-window + time.Minute)}, true},
{"past the window is housekeeping", &repository.WarmupReceived{CreatedAt: now.Add(-window - time.Minute)}, false},
{"weeks later is housekeeping", &repository.WarmupReceived{CreatedAt: now.AddDate(0, 0, -31)}, false},
{"retired by the platform is never a strike, even fresh", &repository.WarmupReceived{CreatedAt: now.Add(-time.Hour), RetiredAt: &retired}, false},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
if got := warmupDeletionCounts(tc.rec, now); got != tc.want {
t.Fatalf("warmupDeletionCounts() = %v, want %v", got, tc.want)
}
})
}
}
// A removal of warmup mail is a strike only while the message is fresh. Later
// it is the mailbox owner tidying, Gmail purging its Trash, a server retention
// rule, or the platform's own retention, none of which is harm.
func TestRemoveEmailStrikesOnlyFreshWarmupMail(t *testing.T) {
cases := []struct {
name string
rec *repository.WarmupReceived
strikes int
}{
{"deleted an hour after arrival", receivedAgo(time.Hour), 1},
{"deleted a week after arrival", receivedAgo(7 * 24 * time.Hour), 0},
{"not warmup at all", nil, 0},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
s, svc := retentionService(tc.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) != tc.strikes {
t.Fatalf("strikes = %v, want %d", svc.strikes, tc.strikes)
}
})
}
}
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)
}
}
// 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.
func TestFlagsAddJudgesGmailTrashOnFreshness(t *testing.T) {
cases := []struct {
name string
rec *repository.WarmupReceived
flags []string
want []string
}{
{"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},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
s, svc := retentionService(tc.rec)
if err := s.HandleFlagsAdd(context.Background(), &models.JobEventFlags{
UserID: uuid.New(), EmailID: uuid.New(), ID: uuid.New(), Flags: tc.flags,
}); err != nil {
t.Fatal(err)
}
if len(svc.strikes) != len(tc.want) {
t.Fatalf("strikes = %v, want %v", svc.strikes, tc.want)
}
for i := range tc.want {
if svc.strikes[i] != tc.want[i] {
t.Fatalf("strikes = %v, want %v", svc.strikes, tc.want)
}
}
})
}
}
// retentionPublisher captures the delete actions the sweep publishes.
type retentionPublisher struct {
backfillPublisher
}
// retentionRepo serves one listing of each kind and records what was retired.
type retentionRepo struct {
repository.WarmupRepository
received []repository.WarmupMailToRetire
sent []repository.WarmupMailToRetire
retired []uuid.UUID
pruned *time.Time
}
func (r *retentionRepo) ListWarmupMailToRetire(context.Context, int, int) ([]repository.WarmupMailToRetire, error) {
return r.received, nil
}
func (r *retentionRepo) ListWarmupSentCopiesToRetire(context.Context, int, int) ([]repository.WarmupMailToRetire, error) {
return r.sent, nil
}
func (r *retentionRepo) RetireWarmupReceived(_ context.Context, _, internalID uuid.UUID) error {
r.retired = append(r.retired, internalID)
return nil
}
func (r *retentionRepo) RetireWarmupSentCopy(_ context.Context, token uuid.UUID) error {
r.retired = append(r.retired, token)
return nil
}
func (r *retentionRepo) PruneWarmupEventsBefore(_ context.Context, before time.Time) (int64, error) {
r.pruned = &before
return 0, nil
}
// The sweep publishes one delete per row, with every key the worker can find
// the message by, and retires the row only once the action is on the bus.
func TestRetireWarmupMailBatchPublishesAndRetires(t *testing.T) {
worker := uuid.New()
received := repository.WarmupMailToRetire{
UserID: uuid.New(), EmailAccountID: uuid.New(), WorkerID: worker,
InternalID: uuid.New(), MessageID: "<in@example.test>", ProviderKey: "gmail-1",
Placement: models.WarmupPlacementFolder, Folder: "Reputation",
}
sent := repository.WarmupMailToRetire{
UserID: uuid.New(), EmailAccountID: uuid.New(), WorkerID: worker,
Token: uuid.New(), MessageID: "<out@example.test>",
}
repo := &retentionRepo{received: []repository.WarmupMailToRetire{received}, sent: []repository.WarmupMailToRetire{sent}}
pub := &retentionPublisher{}
s := &JobsService{WarmupRepo: repo, Publisher: pub}
done, err := s.retireWarmupMailBatch(context.Background())
if err != nil {
t.Fatal(err)
}
if !done {
t.Fatal("a short listing should complete the pass")
}
if len(pub.actions) != 2 {
t.Fatalf("published %d actions, want 2", len(pub.actions))
}
in := pub.actions[0]
if in.Actions[0] != models.WarmupActionDelete || in.GmailID != "gmail-1" || in.RFCMessageID != received.MessageID ||
in.InternalID != received.InternalID.String() || in.TargetFolder != "Reputation" {
t.Fatalf("received delete carried %+v", in)
}
out := pub.actions[1]
if out.InternalID != "" || out.RFCMessageID != sent.MessageID || out.TargetFolder != config.WarmupFolderDefault {
t.Fatalf("sent-copy delete carried %+v", out)
}
if len(repo.retired) != 2 || repo.retired[0] != received.InternalID || repo.retired[1] != sent.Token {
t.Fatalf("retired %v, want the receipt then the token", repo.retired)
}
}
func TestPruneWarmupEventsUsesTheInstanceWindow(t *testing.T) {
repo := &retentionRepo{}
s := &JobsService{WarmupRepo: repo, Publisher: &retentionPublisher{}}
if err := s.pruneWarmupEvents(context.Background()); err != nil {
t.Fatal(err)
}
if repo.pruned == nil {
t.Fatal("nothing was pruned")
}
want := time.Now().AddDate(0, 0, -config.WarmupEventRetentionDaysDefault)
if d := repo.pruned.Sub(want); d < -time.Minute || d > time.Minute {
t.Fatalf("pruned before %v, want about %v", repo.pruned, want)
}
}
+29 -6
View File
@@ -67,6 +67,14 @@ type Retention struct {
FormEventDays int `json:"form_event_days"`
// AuditLogDays is how long the audit trail is kept.
AuditLogDays int `json:"audit_log_days"`
// WarmupMailDays is how long warmup mail stays in a mailbox before the
// platform deletes it from the warmup folder. A mailbox may set its own
// window; this is the one every other mailbox follows.
WarmupMailDays int `json:"warmup_mail_days"`
// WarmupEventDays is how long the per-message warmup records (tokens,
// receipts, tampering and spam reports) are kept. The daily warmup
// statistics behind the analytics are separate and never pruned.
WarmupEventDays int `json:"warmup_event_days"`
}
// Tracking holds the engagement-classification windows. Zero means "compiled
@@ -235,6 +243,8 @@ func DefaultRetention() Retention {
EngagementEventDays: config.EngagementEventRetentionDaysDefault,
FormEventDays: config.FormEventsRetentionDaysDefault,
AuditLogDays: config.AuditLogRetentionDaysDefault,
WarmupMailDays: config.WarmupMailRetentionDaysDefault,
WarmupEventDays: config.WarmupEventRetentionDaysDefault,
}
}
@@ -243,21 +253,26 @@ func DefaultRetention() Retention {
// written before this section existed must not silently start deleting
// everything on the next sweep.
func (r *Retention) Normalize() {
clamp := func(v, def int) int {
clamp := func(v, def, floor int) int {
if v <= 0 {
return def
}
if v < config.RetentionDaysMin {
return config.RetentionDaysMin
if v < floor {
return floor
}
if v > config.RetentionDaysMax {
return config.RetentionDaysMax
}
return v
}
r.EngagementEventDays = clamp(r.EngagementEventDays, config.EngagementEventRetentionDaysDefault)
r.FormEventDays = clamp(r.FormEventDays, config.FormEventsRetentionDaysDefault)
r.AuditLogDays = clamp(r.AuditLogDays, config.AuditLogRetentionDaysDefault)
r.EngagementEventDays = clamp(r.EngagementEventDays, config.EngagementEventRetentionDaysDefault, config.RetentionDaysMin)
r.FormEventDays = clamp(r.FormEventDays, config.FormEventsRetentionDaysDefault, config.RetentionDaysMin)
r.AuditLogDays = clamp(r.AuditLogDays, config.AuditLogRetentionDaysDefault, config.RetentionDaysMin)
// The two warmup windows have floors of their own: mail has to outlive
// the engagement legs and a reply-back, and the records have to outlive
// the thirty-day health bands that read them.
r.WarmupMailDays = clamp(r.WarmupMailDays, config.WarmupMailRetentionDaysDefault, config.WarmupMailRetentionDaysMin)
r.WarmupEventDays = clamp(r.WarmupEventDays, config.WarmupEventRetentionDaysDefault, config.WarmupEventRetentionDaysMin)
}
// Normalize clamps a document into its accepted range. It is applied on read
@@ -344,6 +359,8 @@ type Patch struct {
EngagementEventDays *int `json:"engagement_event_days"`
FormEventDays *int `json:"form_event_days"`
AuditLogDays *int `json:"audit_log_days"`
WarmupMailDays *int `json:"warmup_mail_days"`
WarmupEventDays *int `json:"warmup_event_days"`
} `json:"retention"`
Tracking *struct {
MachineWindowOpenSeconds *int `json:"machine_window_open_seconds"`
@@ -404,6 +421,12 @@ func (p Patch) Apply(doc Document) Document {
if p.Retention.AuditLogDays != nil {
doc.Retention.AuditLogDays = *p.Retention.AuditLogDays
}
if p.Retention.WarmupMailDays != nil {
doc.Retention.WarmupMailDays = *p.Retention.WarmupMailDays
}
if p.Retention.WarmupEventDays != nil {
doc.Retention.WarmupEventDays = *p.Retention.WarmupEventDays
}
}
if p.Tracking != nil {
if p.Tracking.MachineWindowOpenSeconds != nil {
@@ -245,3 +245,35 @@ func TestPatchTracking(t *testing.T) {
t.Errorf("absent tracking section changed the document: %+v", kept.Tracking)
}
}
// The warmup windows carry floors of their own. Zero still resolves to the
// default, so a document stored before they existed keeps every mailbox's
// warmup mail for the compiled window rather than deleting it on the next
// sweep, and a value under the floor is raised to it rather than refused.
func TestRetentionNormalizeWarmupWindows(t *testing.T) {
tests := []struct {
name string
mail int
events int
wantMail int
wantEvent int
}{
{"zero resolves to the defaults", 0, 0, config.WarmupMailRetentionDaysDefault, config.WarmupEventRetentionDaysDefault},
{"below the floors clamps up", 1, 7, config.WarmupMailRetentionDaysMin, config.WarmupEventRetentionDaysMin},
{"in range is kept", 14, 120, 14, 120},
{"above the ceiling clamps down", config.RetentionDaysMax + 1, config.RetentionDaysMax + 1, config.RetentionDaysMax, config.RetentionDaysMax},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
r := Retention{WarmupMailDays: tt.mail, WarmupEventDays: tt.events}
r.Normalize()
if r.WarmupMailDays != tt.wantMail || r.WarmupEventDays != tt.wantEvent {
t.Errorf("Normalize() = (%d, %d), want (%d, %d)", r.WarmupMailDays, r.WarmupEventDays, tt.wantMail, tt.wantEvent)
}
// The three older windows are untouched by the new floors.
if r.EngagementEventDays != config.EngagementEventRetentionDaysDefault {
t.Errorf("EngagementEventDays = %d, want the default", r.EngagementEventDays)
}
})
}
}
+4 -2
View File
@@ -672,7 +672,9 @@ func moreSevere(a, b evaluationDecision) evaluationDecision {
// 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.
// most likely cause is someone tidying the folder by hand. A deletion is only
// recorded at all inside config.WarmupDeletionStrikeHours of arrival; past
// that the platform's own retention would have removed the message anyway.
func evaluateTampering(metrics *models.WarmupHealthMetrics, now time.Time) evaluationDecision {
strikes := metrics.TamperingStrikes()
score := maxFloat(float64(strikes)*10, metrics.SpamPlacementRate)
@@ -696,7 +698,7 @@ func evaluateTampering(metrics *models.WarmupHealthMetrics, now time.Time) evalu
case strikes >= tamperingWatchStrikes:
return evaluationDecision{
State: models.WarmupHealthWatch,
Reason: "A warmup email was " + tamperingVerb(tamperingKind(metrics)) + ". Leave warmup mail where it is filed; a second one within 7 days pauses warmup.",
Reason: "A warmup email was " + tamperingVerb(tamperingKind(metrics)) + " soon after it arrived. Leave warmup mail where it is filed; Warmbly clears it on its own once its retention window passes. A second one within 7 days pauses warmup.",
Score: score,
}
}
+174 -19
View File
@@ -2,13 +2,18 @@ package worker
import (
"context"
"errors"
"fmt"
"strings"
"github.com/google/uuid"
"github.com/rs/zerolog/log"
"github.com/warmbly/warmbly/internal/app/worker/wmail"
"github.com/warmbly/warmbly/internal/client/smtpimap/imap"
"github.com/warmbly/warmbly/internal/config"
"github.com/warmbly/warmbly/internal/infrastructure/storage"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/repository"
)
// HandleWarmupAction executes recipient-side warmup actions on the mailbox.
@@ -36,27 +41,51 @@ func (w *WorkerService) HandleWarmupAction(ctx context.Context, action models.Wa
w.mailManager.RLock()
mail, exists := w.mailManager.Emails[action.EmailID]
w.mailManager.RUnlock()
if !exists {
log.Warn().Str("email_id", action.EmailID.String()).Msg("Email account not found for warmup action")
return nil
}
var err error
switch {
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.
log.Warn().Str("email_id", action.EmailID.String()).Msg("Email account not found for warmup action")
if hasWarmupAction(action.Actions, models.WarmupActionDelete) {
err = errors.New("mailbox not loaded on this worker")
}
case mail.GoogleData != nil && mail.GoogleData.Client != nil:
w.runGoogleWarmupActions(ctx, mail, action)
err = w.runGoogleWarmupActions(ctx, mail, action)
case mail.GraphData != nil && mail.GraphData.Client != nil:
w.runGraphWarmupActions(ctx, mail, action)
err = w.runGraphWarmupActions(ctx, mail, action)
case mail.SmtpImapData != nil && mail.SmtpImapData.ImapClient != nil:
w.runImapWarmupActions(ctx, mail, action)
err = w.runImapWarmupActions(ctx, mail, action)
default:
log.Warn().
Str("email_id", action.EmailID.String()).
Msg("No mail client available for warmup actions; skipping")
}
if err == nil {
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. The bus redelivers
// on an error, bounded so a message the provider will never give up is
// not retried forever.
if d := deliveryOf(ctx); d.redelivers && d.attempt < warmupDeleteRedeliveries {
return err
}
log.Error().Err(err).
Str("email_id", action.EmailID.String()).
Str("rfc_message_id", action.RFCMessageID).
Msg("Warmup delete gave up; the message stays in the mailbox")
return nil
}
func (w *WorkerService) runGoogleWarmupActions(ctx context.Context, mail *wmail.WMail, action models.WarmupEmailAction) {
// warmupDeleteRedeliveries bounds how many times a failed retention delete is
// redelivered before the message is left in place.
const warmupDeleteRedeliveries = 5
func (w *WorkerService) runGoogleWarmupActions(ctx context.Context, mail *wmail.WMail, action models.WarmupEmailAction) error {
placement, folder := warmupFiling(action)
// The archive placement wants the message out of every view without a label
// of its own, which on Gmail is "remove INBOX, add nothing".
@@ -65,12 +94,33 @@ func (w *WorkerService) runGoogleWarmupActions(ctx context.Context, mail *wmail.
label = ""
}
// The delete for the sender's own copy carries no Gmail id: the map is
// keyed by the provider's id and the control plane only knows the
// Message-ID, so it is resolved here.
gmailID := action.GmailID
if gmailID == "" && action.RFCMessageID != "" && hasWarmupAction(action.Actions, models.WarmupActionDelete) {
id, err := mail.GoogleData.Client.FindByRFCMessageID(ctx, action.RFCMessageID)
if err != nil {
return fmt.Errorf("search Gmail for the warmup message: %w", err)
}
gmailID = id
}
for _, act := range action.Actions {
switch act {
case models.WarmupActionFile:
if err := mail.GoogleData.Client.FileWarmup(ctx, action.GmailID, label); err != nil {
log.Error().Err(err).Str("gmail_id", action.GmailID).Str("folder", label).Msg("Failed to file warmup message (Gmail)")
}
case models.WarmupActionDelete:
if gmailID != "" {
if err := mail.GoogleData.Client.Trash(ctx, gmailID); err != nil {
return fmt.Errorf("delete warmup message (Gmail): %w", err)
}
}
if err := w.dropWarmupBody(ctx, mail, action, gmailID); err != nil {
return err
}
case models.WarmupActionMarkRead:
if err := mail.GoogleData.Client.MarkAsRead(ctx, action.GmailID); err != nil {
log.Error().Err(err).Str("gmail_id", action.GmailID).Msg("Failed to mark as read")
@@ -92,13 +142,14 @@ func (w *WorkerService) runGoogleWarmupActions(ctx context.Context, mail *wmail.
log.Warn().Str("action", act).Msg("Unknown warmup action")
}
}
return nil
}
// runGraphWarmupActions applies recipient-side warmup engagement to a Microsoft
// Graph mailbox. action.GmailID carries the Graph message id (the provider id
// field is provider-agnostic). Star maps to the follow-up flag, the closest
// Outlook equivalent of a Gmail star.
func (w *WorkerService) runGraphWarmupActions(ctx context.Context, mail *wmail.WMail, action models.WarmupEmailAction) {
func (w *WorkerService) runGraphWarmupActions(ctx context.Context, mail *wmail.WMail, action models.WarmupEmailAction) error {
client := mail.GraphData.Client
// A Graph message id changes whenever the message is moved (copy+delete), so
@@ -131,6 +182,7 @@ func (w *WorkerService) runGraphWarmupActions(ctx context.Context, mail *wmail.W
continue
}
if newID != "" {
w.remapProviderID(ctx, mail, msgID, newID)
msgID = newID // subsequent actions target the moved copy
}
case models.WarmupActionMarkRead:
@@ -144,6 +196,7 @@ func (w *WorkerService) runGraphWarmupActions(ctx context.Context, mail *wmail.W
continue
}
if newID != "" {
w.remapProviderID(ctx, mail, msgID, newID)
msgID = newID
}
case models.WarmupActionMarkImportant:
@@ -154,13 +207,46 @@ func (w *WorkerService) runGraphWarmupActions(ctx context.Context, mail *wmail.W
if err := client.AddFlag(ctx, msgID); err != nil {
log.Error().Err(err).Str("graph_id", msgID).Msg("Failed to flag warmup message (Graph)")
}
case models.WarmupActionDelete:
if msgID != "" {
if err := client.Delete(ctx, msgID); err != nil {
return fmt.Errorf("delete warmup message (Graph): %w", err)
}
}
// The live id resolved above is in the map because every move
// this worker made re-keyed it (remapProviderID), which is what
// lets the sender's own copy, whose internal id the control
// plane does not know, drop its body too.
if err := w.dropWarmupBody(ctx, mail, action, msgID); err != nil {
return err
}
default:
log.Warn().Str("action", act).Msg("Unknown warmup action")
}
}
return nil
}
func (w *WorkerService) runImapWarmupActions(ctx context.Context, mail *wmail.WMail, action models.WarmupEmailAction) {
// remapProviderID keeps the message map current across a Graph move, which
// is a copy plus delete and hands the message a new id. The map is keyed by
// the id the sync stored, so without this every message this worker filed is
// unreachable by its live id and the body it stored could never be dropped.
func (w *WorkerService) remapProviderID(ctx context.Context, mail *wmail.WMail, oldID, newID string) {
if mail.EmailMessageMapRepository == nil || oldID == "" || newID == "" || oldID == newID {
return
}
m, err := mail.EmailMessageMapRepository.Get(ctx, mail.UserID, mail.ID, oldID)
if err != nil || m == nil {
return
}
if err := mail.EmailMessageMapRepository.Add(ctx, repository.EmailMessageData{
UserID: m.UserID, EmailID: m.EmailID, MessageID: newID, ID: m.ID, ThreadID: m.ThreadID,
}); err != nil {
log.Debug().Err(err).Str("email_id", mail.ID.String()).Msg("Could not re-key the moved message in the map")
}
}
func (w *WorkerService) runImapWarmupActions(ctx context.Context, mail *wmail.WMail, action models.WarmupEmailAction) error {
boxes := mail.SmtpImapData.Mailboxes
imapClient := mail.SmtpImapData.ImapClient
sourceBox := lookupWarmupSourceFolder(boxes, action)
@@ -192,16 +278,33 @@ func (w *WorkerService) runImapWarmupActions(ctx context.Context, mail *wmail.WM
if sentBox := lookupSent(boxes); sentBox != nil {
sentName = sentBox.Name
}
trashName := ""
if trashBox := lookupTrash(boxes); trashBox != nil {
trashName = trashBox.Name
}
files := hasWarmupAction(action.Actions, models.WarmupActionFile)
boxName, uid := w.locateWarmupMessage(ctx, imapClient, action, sourceBox, files, dst, inboxName, sentName)
deletes := hasWarmupAction(action.Actions, models.WarmupActionDelete)
boxName, uid, searchErr := w.locateWarmupMessage(ctx, imapClient, action, sourceBox, files, dst, inboxName, sentName)
if boxName == "" {
if deletes && searchErr != nil {
// A folder could not be searched, so absence is not established;
// deleting the body now would leave the message with nothing to
// key a retry by.
return fmt.Errorf("locate warmup message for deletion: %w", searchErr)
}
log.Warn().
Str("folder", action.MailboxFolder).
Uint32("uid_validity", action.MailboxUIDValidity).
Str("email_id", action.EmailID.String()).
Msg("Warmup action skipped: the message could not be located in any folder")
return
if deletes {
// Already gone from the mailbox: an earlier delivery removed it
// before the body drop failed, or the owner did. Either way the
// body is the only thing left to do.
return w.dropWarmupBody(ctx, mail, action, action.RFCMessageID)
}
return nil
}
// moved is set once the message leaves boxName. Every later action would
@@ -252,10 +355,58 @@ func (w *WorkerService) runImapWarmupActions(ctx context.Context, mail *wmail.WM
// starring here would just re-flag the same message. Star is a
// Gmail-only distinct signal.
continue
case models.WarmupActionDelete:
if err := imapClient.DeleteUID(ctx, boxName, trashName, uid); err != nil {
return fmt.Errorf("delete warmup message (IMAP) uid %d in %q: %w", uid, boxName, err)
}
moved = true
// The map is keyed by the RFC Message-ID on IMAP.
if err := w.dropWarmupBody(ctx, mail, action, action.RFCMessageID); err != nil {
return err
}
default:
log.Warn().Str("action", act).Msg("Unknown warmup action")
}
}
return nil
}
// dropWarmupBody removes the platform's stored copy of a warmup message's
// body once the message is gone from the mailbox. The body was written when
// the message was synced, before anything knew it was warmup, and nothing
// else ever comes back for it. The internal id keying the blob travels with
// the action when the control plane has it; otherwise it is looked up from
// the provider key the worker acted on. Finding neither leaves the blob to
// the mailbox's erasure, which sweeps the whole prefix. A store that refuses
// the delete is an error, so the bus offers the action again.
func (w *WorkerService) dropWarmupBody(ctx context.Context, mail *wmail.WMail, action models.WarmupEmailAction, providerKey string) error {
if mail.Storage == nil {
return nil
}
internalID, err := uuid.Parse(action.InternalID)
if err != nil && providerKey != "" && mail.EmailMessageMapRepository != nil {
if m, lerr := mail.EmailMessageMapRepository.Get(ctx, mail.UserID, mail.ID, providerKey); lerr == nil && m != nil {
internalID, err = uuid.Parse(m.ID)
}
}
if err != nil || internalID == uuid.Nil {
log.Debug().Str("email_id", action.EmailID.String()).Msg("Warmup body left in place: no internal id to key it by")
return nil
}
key := config.StorageEndpointEmailBody(mail.UserID, mail.ID, internalID)
if err := mail.Storage.Delete(ctx, key); err != nil && !errors.Is(err, storage.ErrNotFound) {
return fmt.Errorf("drop the stored warmup body: %w", err)
}
return nil
}
func lookupTrash(boxes []*models.Mailbox) *models.Mailbox {
for _, b := range boxes {
if b != nil && imap.IsTrashMailbox(b.Name, b.Attrs) {
return b
}
}
return nil
}
// locateWarmupMessage resolves the folder and UID an action should act on.
@@ -270,7 +421,9 @@ func (w *WorkerService) runImapWarmupActions(ctx context.Context, mail *wmail.WM
//
// The Message-ID is the one identifier a move does not change, so the likely
// destinations are searched for it, most likely first. An empty folder name
// means the message is nowhere we know to look.
// means the message is nowhere we know to look; the error, when set with it,
// is the last search that failed, so the caller can tell "absent" from
// "could not look".
func (w *WorkerService) locateWarmupMessage(
ctx context.Context,
client warmupIMAPClient,
@@ -278,10 +431,11 @@ func (w *WorkerService) locateWarmupMessage(
sourceBox *models.Mailbox,
files bool,
dst, inboxName, sentName string,
) (string, uint32) {
) (string, uint32, error) {
if files && sourceBox != nil {
return sourceBox.Name, action.UID
return sourceBox.Name, action.UID, nil
}
var searchErr error
if action.RFCMessageID != "" {
// Ordered by likelihood: the destination a previous leg filed it into,
// then the inbox a rescue put it back in, then Sent for our own copy of
@@ -299,17 +453,18 @@ func (w *WorkerService) locateWarmupMessage(
uid, err := client.FindUIDByMessageID(ctx, name, action.RFCMessageID)
if err != nil {
log.Debug().Err(err).Str("folder", name).Str("email_id", action.EmailID.String()).Msg("Could not search a folder for the warmup message")
searchErr = err
continue
}
if uid != 0 {
return name, uid
return name, uid, nil
}
}
}
if sourceBox != nil {
return sourceBox.Name, action.UID
return sourceBox.Name, action.UID, nil
}
return "", 0
return "", 0, searchErr
}
// warmupIMAPClient is the slice of the IMAP client the warmup actions use, so
@@ -61,7 +61,7 @@ func TestLocateWarmupMessage(t *testing.T) {
t.Run("the filing leg acts where the mail arrived", func(t *testing.T) {
stub := &searchStub{}
box, uid := w.locateWarmupMessage(context.Background(), stub, warmupAction(), junk, true, "Warmbly", inboxName, "Sent")
box, uid, _ := w.locateWarmupMessage(context.Background(), stub, warmupAction(), junk, true, "Warmbly", inboxName, "Sent")
if box != "Junk" || uid != 41 {
t.Fatalf("got (%q, %d), want (Junk, 41)", box, uid)
}
@@ -74,7 +74,7 @@ func TestLocateWarmupMessage(t *testing.T) {
// This is the case the delayed leg used to get wrong: the UID it carries
// belongs to the arrival folder, which no longer has the message.
stub := &searchStub{found: map[string]uint32{"Warmbly": 7}}
box, uid := w.locateWarmupMessage(context.Background(), stub, warmupAction(), junk, false, "Warmbly", inboxName, "Sent")
box, uid, _ := w.locateWarmupMessage(context.Background(), stub, warmupAction(), junk, false, "Warmbly", inboxName, "Sent")
if box != "Warmbly" || uid != 7 {
t.Fatalf("got (%q, %d), want (Warmbly, 7)", box, uid)
}
@@ -87,7 +87,7 @@ func TestLocateWarmupMessage(t *testing.T) {
// Nothing is filed for the inbox placement, so the only move that
// happened was the rescue out of Junk.
stub := &searchStub{found: map[string]uint32{"INBOX": 12}}
box, uid := w.locateWarmupMessage(context.Background(), stub, warmupAction(), junk, false, "", inboxName, "Sent")
box, uid, _ := w.locateWarmupMessage(context.Background(), stub, warmupAction(), junk, false, "", inboxName, "Sent")
if box != "INBOX" || uid != 12 {
t.Fatalf("got (%q, %d), want (INBOX, 12)", box, uid)
}
@@ -95,7 +95,7 @@ func TestLocateWarmupMessage(t *testing.T) {
t.Run("a message nobody moved keeps its arrival folder and UID", func(t *testing.T) {
stub := &searchStub{}
box, uid := w.locateWarmupMessage(context.Background(), stub, warmupAction(), junk, false, "Warmbly", inboxName, "Sent")
box, uid, _ := w.locateWarmupMessage(context.Background(), stub, warmupAction(), junk, false, "Warmbly", inboxName, "Sent")
if box != "Junk" || uid != 41 {
t.Fatalf("got (%q, %d), want the arrival folder and UID", box, uid)
}
@@ -103,7 +103,7 @@ func TestLocateWarmupMessage(t *testing.T) {
t.Run("a folder that refuses a search does not end the hunt", func(t *testing.T) {
stub := &searchStub{fail: map[string]bool{"Warmbly": true}, found: map[string]uint32{"INBOX": 3}}
box, uid := w.locateWarmupMessage(context.Background(), stub, warmupAction(), junk, false, "Warmbly", inboxName, "Sent")
box, uid, _ := w.locateWarmupMessage(context.Background(), stub, warmupAction(), junk, false, "Warmbly", inboxName, "Sent")
if box != "INBOX" || uid != 3 {
t.Fatalf("got (%q, %d), want (INBOX, 3)", box, uid)
}
@@ -111,7 +111,7 @@ func TestLocateWarmupMessage(t *testing.T) {
t.Run("no source folder and nothing found is nothing to act on", func(t *testing.T) {
stub := &searchStub{}
box, uid := w.locateWarmupMessage(context.Background(), stub, warmupAction(), nil, false, "Warmbly", inboxName, "Sent")
box, uid, _ := w.locateWarmupMessage(context.Background(), stub, warmupAction(), nil, false, "Warmbly", inboxName, "Sent")
if box != "" || uid != 0 {
t.Fatalf("got (%q, %d), want an empty answer", box, uid)
}
@@ -124,7 +124,7 @@ func TestLocateWarmupMessage(t *testing.T) {
// The filing leg for a sent copy, on a worker whose folder list does not
// yet carry the UIDVALIDITY the event was published with.
stub := &searchStub{found: map[string]uint32{"Sent": 19}}
box, uid := w.locateWarmupMessage(context.Background(), stub, warmupAction(), nil, true, "Warmbly", inboxName, "Sent")
box, uid, _ := w.locateWarmupMessage(context.Background(), stub, warmupAction(), nil, true, "Warmbly", inboxName, "Sent")
if box != "Sent" || uid != 19 {
t.Fatalf("got (%q, %d), want (Sent, 19)", box, uid)
}
@@ -144,7 +144,7 @@ func TestLocateWarmupMessage(t *testing.T) {
stub := &searchStub{found: map[string]uint32{"Warmbly": 7}}
action := warmupAction()
action.RFCMessageID = ""
box, uid := w.locateWarmupMessage(context.Background(), stub, action, junk, false, "Warmbly", inboxName, "Sent")
box, uid, _ := w.locateWarmupMessage(context.Background(), stub, action, junk, false, "Warmbly", inboxName, "Sent")
if box != "Junk" || uid != 41 {
t.Fatalf("got (%q, %d), want the arrival folder and UID", box, uid)
}
+4
View File
@@ -55,6 +55,10 @@ type ImapConn interface {
// FindUIDByMessageID relocates a warmup message whose UID went void when an
// earlier engagement leg moved it.
FindUIDByMessageID(ctx context.Context, mailboxName, rfcMessageID string) (uint32, error)
// DeleteUID removes one message, the retention window's deletion: an
// 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
}
var _ ImapConn = (*imap.Client)(nil)
+39
View File
@@ -241,3 +241,42 @@ func (c *Client) SetSeen(ctx context.Context, messageIDs []string, seen bool) er
}
return nil
}
// Trash moves a message to Gmail's Trash, which Gmail empties on its own
// after thirty days. This is the deletion the retention window asks for: the
// modify scope the mailbox was connected with allows trashing and not the
// permanent delete, and Trash is also what a person pressing Delete gets.
func (c *Client) Trash(ctx context.Context, messageID string) error {
if c.srv == nil {
return fmt.Errorf("gmail service not initialized")
}
_, err := c.srv.Users.Messages.Trash("me", messageID).Context(ctx).Do()
if err != nil {
var gerr *googleapi.Error
if errors.As(err, &gerr) && gerr.Code == 404 {
// Already gone, which is the state being asked for.
return nil
}
return fmt.Errorf("failed to trash message: %w", err)
}
return nil
}
// FindByRFCMessageID returns the Gmail id of the message carrying the RFC
// 5322 Message-ID, or "" when the mailbox does not have it. Trash and Spam
// are excluded by Gmail's search by default, which is right here: a message
// already in either needs nothing more from the retention sweep.
func (c *Client) FindByRFCMessageID(ctx context.Context, rfcMessageID string) (string, error) {
rfcMessageID = strings.Trim(strings.TrimSpace(rfcMessageID), "<>")
if rfcMessageID == "" || c.srv == nil {
return "", nil
}
ids, _, err := c.ListMessages(ctx, "rfc822msgid:"+rfcMessageID, "", 1)
if err != nil {
return "", err
}
if len(ids) == 0 {
return "", nil
}
return ids[0], nil
}
+23
View File
@@ -2,6 +2,8 @@ package msgraph
import (
"context"
"io"
"net/http"
"net/url"
"strings"
)
@@ -196,3 +198,24 @@ func (c *Client) cacheFolder(name, id string) {
func (c *Client) SetSeen(ctx context.Context, messageID string, seen bool) error {
return c.doJSON(ctx, "PATCH", c.messageURL(messageID), map[string]any{"isRead": seen}, nil)
}
// Delete removes a message the way Outlook's Delete key does: into Deleted
// Items, where the mailbox's own retention policy takes it from. Exchange
// answers a message that is already gone with 404, which is the state being
// asked for, so that is not an error here.
func (c *Client) Delete(ctx context.Context, messageID string) error {
resp, err := c.do(ctx, http.MethodDelete, c.messageURL(messageID), "", nil)
if err != nil {
return transportError(err)
}
defer resp.Body.Close()
if resp.StatusCode == http.StatusNotFound {
_, _ = io.Copy(io.Discard, resp.Body)
return nil
}
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return HandleError(resp)
}
_, _ = io.Copy(io.Discard, resp.Body)
return nil
}
@@ -390,3 +390,66 @@ func (c *Client) SetSeen(ctx context.Context, mailboxName string, uids []uint32,
}
return nil
}
// ErrNoSingleMessageDelete is returned when the server offers neither UIDPLUS
// nor MOVE, so one message cannot be removed without expunging every message
// another client has flagged \Deleted in the same folder.
var ErrNoSingleMessageDelete = errors.New("imap: server offers neither UIDPLUS nor MOVE; not expunging a whole folder for one message")
// DeleteUID removes one message from mailboxName: the retention window's
// deletion once a warmup message has served its purpose.
//
// Only one message may go. With UIDPLUS that is \Deleted plus an expunge
// scoped to the UID. Without it, a plain EXPUNGE would also take every
// message some other client has flagged in the folder and not yet expunged,
// so the message is moved to trashName instead, where the server's own
// retention takes it from, and only with a real MOVE (the COPY fallback ends
// in that same folder-wide EXPUNGE). A server with neither gets
// ErrNoSingleMessageDelete and the message stays.
func (c *Client) DeleteUID(ctx context.Context, mailboxName, trashName string, uid uint32) error {
c.mu.Lock()
defer c.mu.Unlock()
if merr := c.ensureConnected(); merr != nil {
return merr
}
c.lifecycle.RLock()
defer c.lifecycle.RUnlock()
defer c.begin()()
name := c.qualifyMailboxLocked(mailboxName)
caps := c.client.Caps()
if caps.Has(imap.CapUIDPlus) {
if _, err := c.selectMailbox(name, nil); err != nil {
return fmt.Errorf("select %q: %w", name, err)
}
set := imap.UIDSetNum(imap.UID(uid))
storeCmd := c.client.Store(set, &imap.StoreFlags{
Op: imap.StoreFlagsAdd,
Silent: true,
Flags: []imap.Flag{imap.FlagDeleted},
}, nil)
if err := storeCmd.Close(); err != nil {
return fmt.Errorf("store \\Deleted on uid %d: %w", uid, err)
}
if err := c.client.UIDExpunge(set).Close(); err != nil {
return fmt.Errorf("uid expunge %d in %q: %w", uid, name, err)
}
return nil
}
if trashName != "" && caps.Has(imap.CapMove) && !strings.EqualFold(c.qualifyMailboxLocked(trashName), name) {
return c.moveUIDLocked(mailboxName, trashName, uid)
}
return ErrNoSingleMessageDelete
}
// IsTrashMailbox returns true if the mailbox's attributes or name identify it
// as the Trash folder, under RFC 6154 SPECIAL-USE or by name match.
func IsTrashMailbox(name string, attrs []string) bool {
for _, a := range attrs {
if strings.EqualFold(a, string(imap.MailboxAttrTrash)) {
return true
}
}
return matchesFolderName(strings.ToLower(leaf(strings.TrimSpace(name))), ImapTrash)
}
+24
View File
@@ -252,6 +252,30 @@ const (
RetentionDaysMin = 1
RetentionDaysMax = 3650
// Warmup mail is real mail in the customer's mailbox, and nothing about
// it is worth keeping once its engagement has been recorded. The platform
// deletes it from the warmup folder after this many days (per mailbox
// override in email_accounts.warmup_retention_days), so a mailbox on a
// fixed quota never fills up with it and the deletion is the platform's
// own. The floor leaves room for the delayed engagement legs and for a
// reply-back in the thread to finish before its opener goes.
WarmupMailRetentionDaysDefault = 30
WarmupMailRetentionDaysMin = 3
// WarmupEventRetentionDaysDefault is how long the per-message warmup
// records (tokens, receipts, tampering and spam reports) are kept. The
// health bands read at most thirty days, which is the floor; the daily
// warmup_statistics rows carry the analytics and are never pruned.
WarmupEventRetentionDaysDefault = 365
WarmupEventRetentionDaysMin = 30
// WarmupDeletionStrikeHours is how soon after arrival a deletion of a
// warmup email still costs the pool its engagement and so counts as
// tampering. Later on it is housekeeping: the platform was going to delete
// it anyway, and a mailbox owner tidying a folder, a provider purging its
// Trash or a server retention rule must not read as harm.
WarmupDeletionStrikeHours = 24
// 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,
+1
View File
@@ -170,6 +170,7 @@ var (
ErrEmailReplyRate = New(BadRequest, "Warmup reply rate must be between 0 and 100.")
ErrEmailWarmupPlacement = New(BadRequest, "Warmup filing must be one of: folder, inbox, archive.")
ErrEmailWarmupFolder = New(BadRequest, fmt.Sprintf("Warmup folder must be at most %d characters and contain no folder separators or control characters.", config.WarmupFolderMaxLen))
ErrEmailWarmupRetention = New(BadRequest, fmt.Sprintf("Warmup retention must be between %d and %d days, or 0 to follow the instance setting.", config.WarmupMailRetentionDaysMin, config.RetentionDaysMax))
// Disconnecting a mailbox has to reach the machine syncing it before the
// row goes: afterwards there is no assignment left to read and nothing that
@@ -0,0 +1,9 @@
ALTER TABLE public.warmup_tokens
DROP COLUMN IF EXISTS sent_retired_at;
ALTER TABLE public.warmup_received
DROP COLUMN IF EXISTS retired_at;
ALTER TABLE public.email_accounts
DROP CONSTRAINT IF EXISTS email_accounts_warmup_retention_days_check,
DROP COLUMN IF EXISTS warmup_retention_days;
@@ -0,0 +1,25 @@
-- Warmup mail is retained for a bounded time, by the platform.
--
-- Warmup mail is filed out of the way in the customer's mailbox, but it was
-- never removed: a mailbox on a fixed quota accumulated it indefinitely, and
-- the owner clearing the folder by hand was read as tampering. The platform
-- now deletes warmup mail from the mailbox after a retention window, and a
-- message it has retired can never be a strike against the mailbox.
--
-- email_accounts.warmup_retention_days: per-mailbox window, NULL meaning the
-- instance setting (retention.warmup_mail_days).
-- warmup_received.retired_at: when the platform sent the deletion for the
-- copy this mailbox received. Set before the worker acts, so a removal the
-- sync then observes is known to be ours.
-- warmup_tokens.sent_retired_at: the same for the sender's own copy of the
-- message, which is filed into the same folder.
ALTER TABLE public.email_accounts
ADD COLUMN warmup_retention_days integer,
ADD CONSTRAINT email_accounts_warmup_retention_days_check
CHECK (warmup_retention_days IS NULL OR (warmup_retention_days >= 3 AND warmup_retention_days <= 3650));
ALTER TABLE public.warmup_received
ADD COLUMN retired_at timestamp with time zone;
ALTER TABLE public.warmup_tokens
ADD COLUMN sent_retired_at timestamp with time zone;
@@ -0,0 +1 @@
DROP INDEX CONCURRENTLY IF EXISTS public.idx_warmup_received_live;
@@ -0,0 +1,6 @@
-- The retention sweep walks the live receipts oldest first; the retired ones
-- are the bulk and never need reading again. Built concurrently, on its own,
-- because warmup_received is written on every warmup arrival.
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_warmup_received_live
ON public.warmup_received (created_at)
WHERE retired_at IS NULL;
@@ -0,0 +1 @@
DROP INDEX CONCURRENTLY IF EXISTS public.idx_warmup_tokens_sent_live;
@@ -0,0 +1,5 @@
-- The same for the sender copies the retention sweep walks. Built
-- concurrently, on its own, because warmup_tokens is written on every send.
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_warmup_tokens_sent_live
ON public.warmup_tokens (created_at)
WHERE sent_retired_at IS NULL AND sent_message_id <> '';
+27
View File
@@ -105,6 +105,10 @@ type Email struct {
// means the instance default rather than "no folder".
WarmupPlacement string `json:"warmup_placement"`
WarmupFolder string `json:"warmup_folder"`
// WarmupRetentionDays is how long warmup mail stays in this mailbox before
// the platform deletes it. Zero means the instance setting: see
// WarmupMailRetentionDays.
WarmupRetentionDays int `json:"warmup_retention_days"`
Timezone string `json:"timezone"`
@@ -185,6 +189,26 @@ func (e *Email) WarmupFiling() (placement, folder string) {
return placement, folder
}
// ValidWarmupRetentionDays reports whether d is an accepted per-mailbox
// window: zero for the instance setting, or a number of days inside the band.
func ValidWarmupRetentionDays(d int) bool {
return d == 0 || (d >= config.WarmupMailRetentionDaysMin && d <= config.RetentionDaysMax)
}
// WarmupMailRetentionDays resolves how long this mailbox keeps warmup mail:
// its own window when it set one, otherwise the instance's. A row written
// before the column existed carries zero and follows the instance, which is
// what every mailbox did when the window was not a choice.
func (e *Email) WarmupMailRetentionDays(instanceDays int) int {
if ValidWarmupRetentionDays(e.WarmupRetentionDays) && e.WarmupRetentionDays != 0 {
return e.WarmupRetentionDays
}
if instanceDays < config.WarmupMailRetentionDaysMin {
return config.WarmupMailRetentionDaysDefault
}
return instanceDays
}
// SendFrom is the address this mailbox's mail is actually From. A verified
// alias when one was chosen, the mailbox address otherwise.
func (e *Email) SendFrom() string {
@@ -570,6 +594,9 @@ type UpdateEmail struct {
// an empty string.
WarmupPlacement *string `json:"warmup_placement"`
WarmupFolder *string `json:"warmup_folder"`
// WarmupRetentionDays is how long warmup mail is kept in the mailbox; 0
// goes back to the instance setting.
WarmupRetentionDays *int `json:"warmup_retention_days"`
// Timezone is the mailbox's own IANA zone, which its sending behaviour and
// business-hours window are evaluated in. Empty means not configured, so
+12
View File
@@ -56,6 +56,12 @@ const (
WarmupActionMarkRead = "mark_read"
WarmupActionMarkImportant = "mark_important"
WarmupActionStar = "star"
// WarmupActionDelete removes a warmup message the retention window has
// passed on from the mailbox (Trash on Gmail, Deleted Items on Outlook,
// expunged on IMAP) and drops the platform's copy of its body. Published
// 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"
)
// WarmupEmailAction represents actions to perform on a detected warmup email.
@@ -90,6 +96,12 @@ type WarmupEmailAction struct {
Placement string `json:"placement,omitempty" avro:"placement"`
TargetFolder string `json:"target_folder,omitempty" avro:"target_folder"`
// InternalID is the platform's id for the message, which keys the stored
// body the delete action drops. Empty when the control plane does not
// know it (the sender's own copy of a send), in which case the worker
// resolves it from the provider id it acted on.
InternalID string `json:"internal_id,omitempty" avro:"internal_id"`
// DelaySeconds is retained for wire compatibility but is now always 0: the
// recipient-side "dwell" is owned by the consumer's durable schedule
// (warmup_pending_engagements + the engagement poller), which publishes the
+43
View File
@@ -0,0 +1,43 @@
package models
import (
"testing"
"github.com/warmbly/warmbly/internal/config"
)
func TestWarmupMailRetentionDays(t *testing.T) {
cases := []struct {
name string
mailbox int
instance int
want int
}{
{"unset row follows the instance", 0, 45, 45},
{"own window wins", 14, 45, 14},
{"instance below the floor resolves to the default", 0, 0, config.WarmupMailRetentionDaysDefault},
{"garbage on the row follows the instance", 1, 45, 45},
{"garbage on the row and no instance value resolves to the default", -3, 0, config.WarmupMailRetentionDaysDefault},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
e := &Email{WarmupRetentionDays: tc.mailbox}
if got := e.WarmupMailRetentionDays(tc.instance); got != tc.want {
t.Fatalf("WarmupMailRetentionDays(%d) = %d, want %d", tc.instance, got, tc.want)
}
})
}
}
func TestValidWarmupRetentionDays(t *testing.T) {
for _, ok := range []int{0, config.WarmupMailRetentionDaysMin, 30, config.RetentionDaysMax} {
if !ValidWarmupRetentionDays(ok) {
t.Errorf("%d should be an accepted window", ok)
}
}
for _, bad := range []int{-1, 1, config.WarmupMailRetentionDaysMin - 1, config.RetentionDaysMax + 1} {
if ValidWarmupRetentionDays(bad) {
t.Errorf("%d should be refused", bad)
}
}
}
+18 -8
View File
@@ -665,7 +665,7 @@ func (r *emailRepository) Search(ctx context.Context, orgID, search string, curs
ea.min_wait_time, ea.reply_to, ea.tracking_domain, ea.tracking_domain_verified, ea.tracking_domain_verified_at, ea.track_direct_mail,
ea.auth_state, ea.auth_spf, ea.auth_dkim, ea.auth_dmarc, ea.auth_dmarc_policy, ea.auth_reason, ea.auth_checked_at, ea.auth_failing_since,
ea.warmup, ea.warmup_paused_at, ea.warmup_base,
ea.warmup_max, ea.warmup_increase, ea.warmup_reply_rate, ea.warmup_tag, COALESCE(ea.warmup_pool_type, 'free') AS warmup_pool_type, ea.warmup_start_time, ea.warmup_end_time, ea.warmup_days, ea.warmup_placement, ea.warmup_folder, ea.timezone, ea.save_to_sent,
ea.warmup_max, ea.warmup_increase, ea.warmup_reply_rate, ea.warmup_tag, COALESCE(ea.warmup_pool_type, 'free') AS warmup_pool_type, ea.warmup_start_time, ea.warmup_end_time, ea.warmup_days, ea.warmup_placement, ea.warmup_folder, COALESCE(ea.warmup_retention_days, 0) AS warmup_retention_days, ea.timezone, ea.save_to_sent,
ea.created_at, ea.updated_at,
COALESCE(
array_agg(eat.tag_id) FILTER (WHERE eat.tag_id IS NOT NULL), '{}'
@@ -716,7 +716,7 @@ func (r *emailRepository) Search(ctx context.Context, orgID, search string, curs
&i.LastSyncedAt, &i.LastID, &i.CampaignLimit, &i.MinWaitTime, &i.ReplyTo, &i.TrackingDomain, &i.TrackingDomainVerified, &i.TrackingDomainVerifiedAt, &i.TrackDirectMail,
&i.AuthState, &i.AuthSPF, &i.AuthDKIM, &i.AuthDMARC, &i.AuthDMARCPolicy, &i.AuthReason, &i.AuthCheckedAt, &i.AuthFailingSince,
&i.Warmup, &i.WarmupPausedAt, &i.WarmupBase, &i.WarmupMax, &i.WarmupIncrease, &i.WarmupReplyRate, &i.WarmupTag, &i.WarmupPoolType,
&i.WarmupStartTime, &i.WarmupEndTime, &i.WarmupDays, &i.WarmupPlacement, &i.WarmupFolder, &i.Timezone, &i.SaveToSent,
&i.WarmupStartTime, &i.WarmupEndTime, &i.WarmupDays, &i.WarmupPlacement, &i.WarmupFolder, &i.WarmupRetentionDays, &i.Timezone, &i.SaveToSent,
&i.CreatedAt, &i.UpdatedAt, &i.Tags,
)
if err != nil {
@@ -787,7 +787,7 @@ func (r *emailRepository) Get(ctx context.Context, orgID, emailAccountID string)
ea.min_wait_time, ea.reply_to, ea.tracking_domain, ea.tracking_domain_verified, ea.tracking_domain_verified_at, ea.track_direct_mail,
ea.auth_state, ea.auth_spf, ea.auth_dkim, ea.auth_dmarc, ea.auth_dmarc_policy, ea.auth_reason, ea.auth_checked_at, ea.auth_failing_since,
ea.warmup, ea.warmup_paused_at, ea.warmup_base,
ea.warmup_max, ea.warmup_increase, ea.warmup_reply_rate, ea.warmup_tag, COALESCE(ea.warmup_pool_type, 'free') AS warmup_pool_type, ea.warmup_start_time, ea.warmup_end_time, ea.warmup_days, ea.warmup_placement, ea.warmup_folder, ea.timezone, ea.save_to_sent,
ea.warmup_max, ea.warmup_increase, ea.warmup_reply_rate, ea.warmup_tag, COALESCE(ea.warmup_pool_type, 'free') AS warmup_pool_type, ea.warmup_start_time, ea.warmup_end_time, ea.warmup_days, ea.warmup_placement, ea.warmup_folder, COALESCE(ea.warmup_retention_days, 0) AS warmup_retention_days, ea.timezone, ea.save_to_sent,
ea.created_at, ea.updated_at,
COALESCE(array_agg(eat.tag_id) FILTER (WHERE eat.tag_id IS NOT NULL), '{}') AS tags
FROM email_accounts ea
@@ -811,7 +811,7 @@ func (r *emailRepository) Get(ctx context.Context, orgID, emailAccountID string)
&i.LastSyncedAt, &i.LastID, &i.CampaignLimit, &i.MinWaitTime, &i.ReplyTo, &i.TrackingDomain, &i.TrackingDomainVerified, &i.TrackingDomainVerifiedAt, &i.TrackDirectMail,
&i.AuthState, &i.AuthSPF, &i.AuthDKIM, &i.AuthDMARC, &i.AuthDMARCPolicy, &i.AuthReason, &i.AuthCheckedAt, &i.AuthFailingSince,
&i.Warmup, &i.WarmupPausedAt, &i.WarmupBase, &i.WarmupMax, &i.WarmupIncrease, &i.WarmupReplyRate, &i.WarmupTag, &i.WarmupPoolType,
&i.WarmupStartTime, &i.WarmupEndTime, &i.WarmupDays, &i.WarmupPlacement, &i.WarmupFolder, &i.Timezone, &i.SaveToSent,
&i.WarmupStartTime, &i.WarmupEndTime, &i.WarmupDays, &i.WarmupPlacement, &i.WarmupFolder, &i.WarmupRetentionDays, &i.Timezone, &i.SaveToSent,
&i.CreatedAt, &i.UpdatedAt, &i.Tags,
)
if err != nil {
@@ -1067,6 +1067,16 @@ func (r *emailRepository) Update(ctx context.Context, orgID, emailAccountID stri
args = append(args, folder)
argPos++
}
if udata.WarmupRetentionDays != nil {
if !models.ValidWarmupRetentionDays(*udata.WarmupRetentionDays) {
return nil, errx.ErrEmailWarmupRetention
}
// Zero is "follow the instance", stored as NULL so a later change to
// the instance setting reaches every mailbox that never chose.
setClauses = append(setClauses, fmt.Sprintf("%s = NULLIF($%d, 0)", "warmup_retention_days", argPos))
args = append(args, *udata.WarmupRetentionDays)
argPos++
}
// Tags are not a column on the row, so a patch that only moves them still
// leaves setClauses empty. Refusing it made the mailbox drawer's tag
@@ -1093,7 +1103,7 @@ func (r *emailRepository) Update(ctx context.Context, orgID, emailAccountID stri
COALESCE(last_synced_at, created_at) AS last_synced_at, last_id, campaign_limit, min_wait_time, reply_to, tracking_domain, tracking_domain_verified, tracking_domain_verified_at, track_direct_mail,
auth_state, auth_spf, auth_dkim, auth_dmarc, auth_dmarc_policy, auth_reason, auth_checked_at, auth_failing_since,
warmup, warmup_paused_at, warmup_base, warmup_max, warmup_increase, warmup_reply_rate, warmup_tag, warmup_pool_type,
warmup_start_time, warmup_end_time, warmup_days, warmup_placement, warmup_folder, save_to_sent, created_at, updated_at
warmup_start_time, warmup_end_time, warmup_days, warmup_placement, warmup_folder, COALESCE(warmup_retention_days, 0) AS warmup_retention_days, save_to_sent, created_at, updated_at
`, strings.Join(setClauses, ", "))
var i models.Email
@@ -1105,7 +1115,7 @@ func (r *emailRepository) Update(ctx context.Context, orgID, emailAccountID stri
// dashboard on every unrelated edit.
&i.AuthState, &i.AuthSPF, &i.AuthDKIM, &i.AuthDMARC, &i.AuthDMARCPolicy, &i.AuthReason, &i.AuthCheckedAt, &i.AuthFailingSince,
&i.Warmup, &i.WarmupPausedAt, &i.WarmupBase, &i.WarmupMax, &i.WarmupIncrease, &i.WarmupReplyRate, &i.WarmupTag, &i.WarmupPoolType,
&i.WarmupStartTime, &i.WarmupEndTime, &i.WarmupDays, &i.WarmupPlacement, &i.WarmupFolder, &i.SaveToSent,
&i.WarmupStartTime, &i.WarmupEndTime, &i.WarmupDays, &i.WarmupPlacement, &i.WarmupFolder, &i.WarmupRetentionDays, &i.SaveToSent,
&i.CreatedAt, &i.UpdatedAt,
)
if err != nil {
@@ -1575,7 +1585,7 @@ func (r *emailRepository) GetByID(ctx context.Context, emailAccountID uuid.UUID)
ea.provider, ea.status, COALESCE(ea.last_synced_at, ea.created_at) AS last_synced_at, ea.last_id, ea.campaign_limit,
ea.min_wait_time, ea.reply_to, ea.tracking_domain, ea.tracking_domain_verified, ea.tracking_domain_verified_at, ea.track_direct_mail, ea.warmup, ea.warmup_paused_at, ea.warmup_base,
ea.warmup_max, ea.warmup_increase, ea.warmup_reply_rate, ea.warmup_tag, ea.warmup_pool_type,
ea.warmup_start_time, ea.warmup_end_time, ea.warmup_days, ea.warmup_placement, ea.warmup_folder, ea.timezone, ea.save_to_sent,
ea.warmup_start_time, ea.warmup_end_time, ea.warmup_days, ea.warmup_placement, ea.warmup_folder, COALESCE(ea.warmup_retention_days, 0) AS warmup_retention_days, ea.timezone, ea.save_to_sent,
ea.auth_state, ea.auth_failing_since,
ea.created_at, ea.updated_at,
COALESCE(array_agg(eat.tag_id) FILTER (WHERE eat.tag_id IS NOT NULL), '{}') AS tags
@@ -1591,7 +1601,7 @@ func (r *emailRepository) GetByID(ctx context.Context, emailAccountID uuid.UUID)
&i.Provider, &i.Status, &i.LastSyncedAt, &i.LastID, &i.CampaignLimit,
&i.MinWaitTime, &i.ReplyTo, &i.TrackingDomain, &i.TrackingDomainVerified, &i.TrackingDomainVerifiedAt, &i.TrackDirectMail, &i.Warmup, &i.WarmupPausedAt, &i.WarmupBase,
&i.WarmupMax, &i.WarmupIncrease, &i.WarmupReplyRate, &i.WarmupTag, &i.WarmupPoolType,
&i.WarmupStartTime, &i.WarmupEndTime, &i.WarmupDays, &i.WarmupPlacement, &i.WarmupFolder, &i.Timezone, &i.SaveToSent,
&i.WarmupStartTime, &i.WarmupEndTime, &i.WarmupDays, &i.WarmupPlacement, &i.WarmupFolder, &i.WarmupRetentionDays, &i.Timezone, &i.SaveToSent,
&i.AuthState, &i.AuthFailingSince,
&i.CreatedAt, &i.UpdatedAt, &i.Tags,
)
+159 -2
View File
@@ -89,6 +89,27 @@ type WarmupReceived struct {
MessageID string
SenderAccountID uuid.UUID
CreatedAt time.Time
// 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
}
// WarmupMailToRetire is one warmup message whose retention window has passed,
// with what the worker needs to find it: the provider's key from the message
// map (a Gmail id, a Graph id, or the RFC Message-ID on IMAP), the immutable
// Message-ID, and where the mailbox files warmup so the search starts there.
// Token is set for the sender's own copy of a send, InternalID for the copy a
// recipient received.
type WarmupMailToRetire struct {
UserID uuid.UUID
EmailAccountID uuid.UUID
WorkerID uuid.UUID
InternalID uuid.UUID
Token uuid.UUID
MessageID string
ProviderKey string
Placement string
Folder string
}
// WarmupRepository defines methods for warmup data access
@@ -215,6 +236,25 @@ type WarmupRepository interface {
RecordWarmupTampering(ctx context.Context, accountID uuid.UUID, messageID, kind string) (bool, error)
CountWarmupTamperingSince(ctx context.Context, accountID uuid.UUID, since time.Time) (int, 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
// per-message records are pruned after theirs.
//
// ListWarmupMailToRetire returns received copies whose window (the
// mailbox's own, else defaultDays) has passed, oldest first, on active
// mailboxes that have a worker to act. ListWarmupSentCopiesToRetire is
// the same for the sender's own copy of each send. RetireWarmupReceived
// and RetireWarmupSentCopy stamp the row once the deletion is on the bus.
ListWarmupMailToRetire(ctx context.Context, defaultDays, limit int) ([]WarmupMailToRetire, error)
RetireWarmupReceived(ctx context.Context, accountID, internalID uuid.UUID) error
ListWarmupSentCopiesToRetire(ctx context.Context, defaultDays, limit int) ([]WarmupMailToRetire, error)
RetireWarmupSentCopy(ctx context.Context, token uuid.UUID) error
// PruneWarmupEventsBefore drops per-message records older than before:
// tampering events, spam reports, retired receipts, and tokens whose sent
// copy is gone. A receipt or token whose mail is still in the mailbox is
// kept, so a removal seen later can still be told apart from tampering.
PruneWarmupEventsBefore(ctx context.Context, before time.Time) (int64, error)
// Appeals (user-facing submission; admin review lives in the admin repo).
CreateWarmupAppeal(ctx context.Context, accountID, userID uuid.UUID, reason string) (uuid.UUID, error)
HasPendingWarmupAppeal(ctx context.Context, accountID uuid.UUID) (bool, error)
@@ -1663,13 +1703,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
SELECT email_account_id, internal_id, message_id, sender_account_id, created_at, retired_at
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.EmailAccountID, &w.InternalID, &w.MessageID, &w.SenderAccountID, &w.CreatedAt, &w.RetiredAt,
)
if errors.Is(err, sql.ErrNoRows) {
return nil, nil
@@ -1680,6 +1720,123 @@ func (r *warmupRepository) GetWarmupReceived(ctx context.Context, accountID, int
return &w, nil
}
// ListWarmupMailToRetire lists received warmup copies past their window:
// the mailbox's own when it set one, else the instance default handed in as
// $1, on active mailboxes that have a worker to act. The provider key comes
// from the message map by the internal id, which is what the worker's remove
// and flag events are keyed on.
func (r *warmupRepository) ListWarmupMailToRetire(ctx context.Context, defaultDays, limit int) ([]WarmupMailToRetire, error) {
query := `
SELECT ea.user_id, ea.id, ea.worker_id, wr.internal_id, wr.message_id,
COALESCE(m.message_id, ''), ea.warmup_placement, ea.warmup_folder
FROM warmup_received wr
JOIN email_accounts ea ON ea.id = wr.email_account_id
LEFT JOIN LATERAL (
SELECT em.message_id FROM email_message_map em
WHERE em.email_id = ea.id AND em.id = wr.internal_id
LIMIT 1
) m ON true
WHERE ea.status = 'active'
AND ea.worker_id IS NOT NULL
AND wr.retired_at IS NULL
AND wr.created_at < NOW() - make_interval(days => COALESCE(ea.warmup_retention_days, $1))
ORDER BY wr.created_at
LIMIT $2`
return r.scanMailToRetire(ctx, query, defaultDays, limit, false)
}
// ListWarmupSentCopiesToRetire lists the sender's own copies past their
// window. Gmail and Outlook file that copy themselves; an SMTP mailbox never
// gets one, because the worker does not append warmup to Sent, so those
// senders are left out rather than searched for a message that does not
// exist. The map is keyed by the provider's id on Gmail and Graph, so the key
// is usually empty there and the worker searches by Message-ID instead.
func (r *warmupRepository) ListWarmupSentCopiesToRetire(ctx context.Context, defaultDays, limit int) ([]WarmupMailToRetire, error) {
query := `
SELECT ea.user_id, ea.id, ea.worker_id, wt.token, wt.sent_message_id,
COALESCE(m.message_id, ''), ea.warmup_placement, ea.warmup_folder
FROM warmup_tokens wt
JOIN email_accounts ea ON ea.id = wt.sender_account_id
LEFT JOIN LATERAL (
SELECT em.message_id FROM email_message_map em
WHERE em.email_id = ea.id AND em.message_id = wt.sent_message_id
LIMIT 1
) m ON true
WHERE ea.status = 'active'
AND ea.worker_id IS NOT NULL
AND ea.provider <> 'smtp_imap'
AND wt.sent_message_id <> ''
AND wt.sent_retired_at IS NULL
AND wt.created_at < NOW() - make_interval(days => COALESCE(ea.warmup_retention_days, $1))
ORDER BY wt.created_at
LIMIT $2`
return r.scanMailToRetire(ctx, query, defaultDays, limit, true)
}
func (r *warmupRepository) scanMailToRetire(ctx context.Context, query string, defaultDays, limit int, sent bool) ([]WarmupMailToRetire, error) {
rows, err := r.db.Query(ctx, query, defaultDays, limit)
if err != nil {
return nil, err
}
defer rows.Close()
var out []WarmupMailToRetire
for rows.Next() {
var m WarmupMailToRetire
var id uuid.UUID
if err := rows.Scan(&m.UserID, &m.EmailAccountID, &m.WorkerID, &id, &m.MessageID, &m.ProviderKey, &m.Placement, &m.Folder); err != nil {
return nil, err
}
if sent {
m.Token = id
} else {
m.InternalID = id
}
out = append(out, m)
}
return out, rows.Err()
}
// RetireWarmupReceived stamps the receipt once its deletion is on the bus.
func (r *warmupRepository) RetireWarmupReceived(ctx context.Context, accountID, internalID uuid.UUID) error {
_, err := r.db.Exec(ctx, `
UPDATE warmup_received SET retired_at = NOW()
WHERE email_account_id = $1 AND internal_id = $2 AND retired_at IS NULL`, accountID, internalID)
return err
}
// RetireWarmupSentCopy stamps the token once the deletion of the sender's own
// copy is on the bus.
func (r *warmupRepository) RetireWarmupSentCopy(ctx context.Context, token uuid.UUID) error {
_, err := r.db.Exec(ctx, `
UPDATE warmup_tokens SET sent_retired_at = NOW()
WHERE token = $1 AND sent_retired_at IS NULL`, token)
return err
}
// PruneWarmupEventsBefore drops the per-message warmup records older than
// before. Receipts and tokens are kept while their mail may still be in the
// mailbox (not yet retired), so the retention sweep can still reach it and a
// removal seen later is still recognised as warmup rather than tampering.
func (r *warmupRepository) PruneWarmupEventsBefore(ctx context.Context, before time.Time) (int64, error) {
var total int64
for _, q := range []string{
`DELETE FROM warmup_tampering_events WHERE created_at < $1`,
`DELETE FROM warmup_spam_reports WHERE created_at < $1`,
`DELETE FROM warmup_received WHERE created_at < $1 AND retired_at IS NOT NULL`,
`DELETE FROM warmup_tokens
WHERE created_at < $1
AND (sent_message_id = '' OR sent_retired_at IS NOT NULL
OR sender_account_id IN (SELECT id FROM email_accounts WHERE provider = 'smtp_imap'))`,
} {
cmd, err := r.db.Exec(ctx, q, before)
if err != nil {
return total, err
}
total += cmd.RowsAffected()
}
return total, nil
}
// RecordWarmupTampering records one "harm" a participant did to a warmup email.
// Returns whether a new row was inserted (deduped per account+message+kind).
func (r *warmupRepository) RecordWarmupTampering(ctx context.Context, accountID uuid.UUID, messageID, kind string) (bool, error) {
@@ -0,0 +1,191 @@
package repository
import (
"context"
"testing"
"time"
"github.com/google/uuid"
)
// The retention sweep reaches only warmup mail past its window, on a mailbox
// that can act, and each mailbox's own window wins over the instance's.
func TestLiveWarmupMailRetention(t *testing.T) {
_, pool := liveContactDB(t)
ctx := context.Background()
userID, orgID := uuid.New(), uuid.New()
workerID := uuid.New()
// following: no window of its own, so the instance's 30 days apply.
// strict: a 5-day window of its own. idle: no worker, so nothing can act.
following, strict, idle, partner := uuid.New(), uuid.New(), uuid.New(), uuid.New()
exec := func(query string, args ...any) {
t.Helper()
if _, err := pool.Exec(ctx, query, args...); err != nil {
t.Fatalf("fixture: %v", err)
}
}
exec(`INSERT INTO users (id, first_name, last_name, email) VALUES ($1, 'Warmup', 'Retention', $2)`,
userID, "warmup-retention-"+uuid.NewString()+"@example.test")
exec(`INSERT INTO organizations (id, name, owner_user_id) VALUES ($1, 'Warmup retention', $2)`, orgID, userID)
exec(`INSERT INTO fleet_nodes (id, role, name) VALUES ($1, 'worker', 'retention-test')
ON CONFLICT DO NOTHING`, workerID)
exec(`INSERT INTO workers (id) VALUES ($1) ON CONFLICT DO NOTHING`, workerID)
mailbox := func(id uuid.UUID, name, provider string, worker *uuid.UUID, retention *int) {
exec(`INSERT INTO email_accounts
(id, user_id, organization_id, email, name, signature_plain, signature_html, provider, status, worker_id, warmup_retention_days, warmup_folder)
VALUES ($1, $2, $3, $4, $5, '', '', $6, 'active', $7, $8, 'Reputation')`,
id, userID, orgID, "warmup-retention-"+uuid.NewString()+"@example.test", name, provider, worker, retention)
}
five := 5
// following is on Gmail, which files a sent copy of its own; strict is an
// SMTP mailbox, which never has one.
mailbox(following, "following", "gmail", &workerID, nil)
mailbox(strict, "strict", "smtp_imap", &workerID, &five)
mailbox(idle, "idle", "smtp_imap", nil, nil)
mailbox(partner, "partner", "smtp_imap", &workerID, nil)
received := func(account uuid.UUID, age time.Duration, msgID string) uuid.UUID {
internal := uuid.New()
exec(`INSERT INTO warmup_received (email_account_id, internal_id, message_id, sender_account_id, created_at)
VALUES ($1, $2, $3, $4, NOW() - $5::interval)`, account, internal, msgID, partner, age)
exec(`INSERT INTO email_message_map (user_id, email_id, message_id, id) VALUES ($1, $2, $3, $4)`,
userID, account, msgID, internal)
return internal
}
oldFollowing := received(following, 40*24*time.Hour, "<old-following@example.test>")
freshFollowing := received(following, 10*24*time.Hour, "<fresh-following@example.test>")
oldStrict := received(strict, 10*24*time.Hour, "<old-strict@example.test>")
freshStrict := received(strict, 2*24*time.Hour, "<fresh-strict@example.test>")
oldIdle := received(idle, 40*24*time.Hour, "<old-idle@example.test>")
sentCopy := func(sender uuid.UUID, age time.Duration, msgID string) uuid.UUID {
token, task := uuid.New(), uuid.New()
exec(`INSERT INTO tasks (id, task_type, email_account_id, status, message_id)
VALUES ($1, 'warmup', $2, 'completed', $3)`, task, sender, msgID)
exec(`INSERT INTO warmup_tokens (token, task_id, sender_account_id, recipient_account_id, sent_message_id, created_at)
VALUES ($1, $2, $3, $4, $5, NOW() - $6::interval)`, token, task, sender, partner, msgID, age)
return token
}
oldSent := sentCopy(following, 40*24*time.Hour, "<old-sent@example.test>")
freshSent := sentCopy(following, 10*24*time.Hour, "<fresh-sent@example.test>")
smtpSent := sentCopy(strict, 40*24*time.Hour, "<smtp-sent@example.test>")
exec(`INSERT INTO email_message_map (user_id, email_id, message_id, id) VALUES ($1, $2, '<old-sent@example.test>', $3)`,
userID, following, uuid.New())
t.Cleanup(func() {
for _, step := range []struct {
query string
arg uuid.UUID
}{
{`DELETE FROM warmup_received WHERE email_account_id IN (SELECT id FROM email_accounts WHERE organization_id = $1)`, orgID},
{`DELETE FROM warmup_tokens WHERE sender_account_id IN (SELECT id FROM email_accounts WHERE organization_id = $1)`, orgID},
{`DELETE FROM tasks WHERE email_account_id IN (SELECT id FROM email_accounts WHERE organization_id = $1)`, orgID},
{`DELETE FROM email_message_map WHERE user_id = $1`, userID},
{`DELETE FROM email_accounts WHERE organization_id = $1`, orgID},
{`DELETE FROM organizations WHERE id = $1`, orgID},
{`DELETE FROM users WHERE id = $1`, userID},
{`DELETE FROM workers WHERE id = $1`, workerID},
{`DELETE FROM fleet_nodes WHERE id = $1`, workerID},
} {
if _, err := pool.Exec(context.Background(), step.query, step.arg); err != nil {
t.Errorf("cleanup: %v", err)
}
}
})
repo := &warmupRepository{db: pool}
rows, err := repo.ListWarmupMailToRetire(ctx, 30, 100)
if err != nil {
t.Fatalf("ListWarmupMailToRetire: %v", err)
}
got := map[uuid.UUID]WarmupMailToRetire{}
for _, r := range rows {
if r.EmailAccountID == following || r.EmailAccountID == strict || r.EmailAccountID == idle {
got[r.InternalID] = r
}
}
if len(got) != 2 {
t.Fatalf("listed %d of this org's receipts, want the two past their windows: %+v", len(got), got)
}
for _, want := range []uuid.UUID{oldFollowing, oldStrict} {
if _, ok := got[want]; !ok {
t.Errorf("receipt %s past its window was not listed", want)
}
}
for _, kept := range []uuid.UUID{freshFollowing, freshStrict, oldIdle} {
if _, ok := got[kept]; ok {
t.Errorf("receipt %s was listed: inside its window, or on a mailbox with no worker", kept)
}
}
// Oldest first, with everything the worker needs to find the message.
r := got[oldFollowing]
if r.WorkerID != workerID || r.UserID != userID || r.MessageID != "<old-following@example.test>" ||
r.ProviderKey != "<old-following@example.test>" || r.Folder != "Reputation" {
t.Fatalf("listed row = %+v, want the worker, user, Message-ID, map key and folder", r)
}
// Retiring takes the row out of the next listing.
if err := repo.RetireWarmupReceived(ctx, following, oldFollowing); err != nil {
t.Fatalf("RetireWarmupReceived: %v", err)
}
rec, err := repo.GetWarmupReceived(ctx, following, oldFollowing)
if err != nil || rec == nil || rec.RetiredAt == nil {
t.Fatalf("GetWarmupReceived after retiring = %+v, %v; want retired_at set", rec, err)
}
rows, err = repo.ListWarmupMailToRetire(ctx, 30, 100)
if err != nil {
t.Fatalf("ListWarmupMailToRetire again: %v", err)
}
for _, r := range rows {
if r.InternalID == oldFollowing {
t.Fatal("a retired receipt was offered again")
}
}
// The sender's own copies follow the same window, and the map resolves
// the provider key where it is keyed by Message-ID.
sent, err := repo.ListWarmupSentCopiesToRetire(ctx, 30, 100)
if err != nil {
t.Fatalf("ListWarmupSentCopiesToRetire: %v", err)
}
var listedOld, listedFresh, listedSMTP bool
for _, r := range sent {
switch r.Token {
case oldSent:
listedOld = true
if r.ProviderKey != "<old-sent@example.test>" || r.MessageID != "<old-sent@example.test>" || r.EmailAccountID != following {
t.Fatalf("sent copy row = %+v", r)
}
case freshSent:
listedFresh = true
case smtpSent:
listedSMTP = true
}
}
if !listedOld || listedFresh || listedSMTP {
t.Fatalf("sent copies listed: old=%v fresh=%v smtp=%v, want only the old Gmail one", listedOld, listedFresh, listedSMTP)
}
if err := repo.RetireWarmupSentCopy(ctx, oldSent); err != nil {
t.Fatalf("RetireWarmupSentCopy: %v", err)
}
// The prune drops retired receipts and retired tokens past the window,
// and leaves a receipt whose mail may still be in the mailbox.
exec(`UPDATE warmup_received SET created_at = NOW() - interval '400 days' WHERE internal_id IN ($1, $2)`, oldFollowing, oldIdle)
exec(`UPDATE warmup_tokens SET created_at = NOW() - interval '400 days' WHERE token IN ($1, $2)`, oldSent, smtpSent)
if _, err := repo.PruneWarmupEventsBefore(ctx, time.Now().AddDate(0, 0, -365)); err != nil {
t.Fatalf("PruneWarmupEventsBefore: %v", err)
}
if rec, _ := repo.GetWarmupReceived(ctx, following, oldFollowing); rec != nil {
t.Fatal("a retired receipt past the window survived the prune")
}
if rec, _ := repo.GetWarmupReceived(ctx, idle, oldIdle); rec == nil {
t.Fatal("a receipt whose mail was never retired was pruned; its removal could no longer be told from tampering")
}
var tokens int
if err := pool.QueryRow(ctx, `SELECT COUNT(*) FROM warmup_tokens WHERE token IN ($1, $2)`, oldSent, smtpSent).Scan(&tokens); err != nil || tokens != 0 {
t.Fatalf("retired and SMTP tokens past the window: count=%d err=%v, want both pruned", tokens, err)
}
}
+15 -7
View File
@@ -100,6 +100,8 @@ SYNC_DAILY_ORG="${WARMBLY_SYNC_DAILY_ORG:-25000}"
RET_ENGAGEMENT="${WARMBLY_RETENTION_ENGAGEMENT_DAYS:-365}"
RET_FORMS="${WARMBLY_RETENTION_FORM_DAYS:-180}"
RET_AUDIT="${WARMBLY_RETENTION_AUDIT_DAYS:-90}"
RET_WARMUP_MAIL="${WARMBLY_RETENTION_WARMUP_MAIL_DAYS:-30}"
RET_WARMUP_EVENTS="${WARMBLY_RETENTION_WARMUP_EVENT_DAYS:-365}"
# S3 blob answers, only read when BLOBS=s3.
S3_BUCKET="${WARMBLY_S3_BUCKET:-}"
@@ -1236,9 +1238,9 @@ FSEOF
}
render_settings_bootstrap() {
printf '{"sync":{"backfill_days":%s,"backfill_messages":%s,"daily_messages_per_mailbox":%s,"daily_messages_per_org":%s},"retention":{"engagement_event_days":%s,"form_event_days":%s,"audit_log_days":%s}}' \
printf '{"sync":{"backfill_days":%s,"backfill_messages":%s,"daily_messages_per_mailbox":%s,"daily_messages_per_org":%s},"retention":{"engagement_event_days":%s,"form_event_days":%s,"audit_log_days":%s,"warmup_mail_days":%s,"warmup_event_days":%s}}' \
"$SYNC_BACKFILL_DAYS" "$SYNC_BACKFILL_MESSAGES" "$SYNC_DAILY_MAILBOX" "$SYNC_DAILY_ORG" \
"$RET_ENGAGEMENT" "$RET_FORMS" "$RET_AUDIT"
"$RET_ENGAGEMENT" "$RET_FORMS" "$RET_AUDIT" "$RET_WARMUP_MAIL" "$RET_WARMUP_EVENTS"
}
render_env_bootstrap_owner() {
@@ -2098,13 +2100,14 @@ wiz_retention() {
out " ${WHITE}Start from which posture?${R}"
say ""
choose 1 \
"Defaults|90 days imported; opens/clicks 365 days, forms 180, audit 90." \
"Minimal retention|30 days imported; every event log 30 days. Reports get shorter." \
"Let me set each one|Seven questions, each with its default already filled in."
"Defaults|90 days imported; opens/clicks 365 days, forms 180, audit 90, warmup mail 30." \
"Minimal retention|30 days imported; every event log 30 days, warmup mail 7. Reports get shorter." \
"Let me set each one|Nine questions, each with its default already filled in."
case "$CHOICE" in
2)
SYNC_BACKFILL_DAYS=30; SYNC_BACKFILL_MESSAGES=2000
RET_ENGAGEMENT=30; RET_FORMS=30; RET_AUDIT=30
RET_WARMUP_MAIL=7; RET_WARMUP_EVENTS=30
ok "Minimal retention"
;;
3)
@@ -2121,6 +2124,10 @@ wiz_retention() {
ask_number "Opens and clicks (days)" "$RET_ENGAGEMENT" 1 3650; RET_ENGAGEMENT=$ANSWER
ask_number "Form funnel events (days)" "$RET_FORMS" 1 3650; RET_FORMS=$ANSWER
ask_number "Audit log, incl. IP addresses (days)" "$RET_AUDIT" 1 3650; RET_AUDIT=$ANSWER
say ""
out " ${DIM}Warmup. Mail is deleted from each mailbox after the first window, records after the second.${R}"
ask_number "Warmup mail in mailboxes (days)" "$RET_WARMUP_MAIL" 3 3650; RET_WARMUP_MAIL=$ANSWER
ask_number "Warmup records (days)" "$RET_WARMUP_EVENTS" 30 3650; RET_WARMUP_EVENTS=$ANSWER
;;
*) ok "Defaults" ;;
esac
@@ -2270,7 +2277,7 @@ summary_rows() {
"${DIM}Database${R} ${WHITE}$(db_label)${R}" \
"${DIM}Blobs${R} ${WHITE}$(blob_label)${R}" \
"${DIM}Import${R} ${WHITE}${SYNC_BACKFILL_DAYS} days, up to ${SYNC_BACKFILL_MESSAGES} messages per mailbox${R}" \
"${DIM}Kept${R} ${WHITE}opens/clicks ${RET_ENGAGEMENT}d, forms ${RET_FORMS}d, audit ${RET_AUDIT}d${R}" \
"${DIM}Kept${R} ${WHITE}opens/clicks ${RET_ENGAGEMENT}d, forms ${RET_FORMS}d, audit ${RET_AUDIT}d, warmup mail ${RET_WARMUP_MAIL}d, records ${RET_WARMUP_EVENTS}d${R}" \
"${DIM}Backups${R} ${WHITE}$(backup_label)${R}" \
"${DIM}First owner${R} ${WHITE}$(owner_label)${R}" \
"${DIM}Sign-ups${R} ${WHITE}$(registration_label)${R}" \
@@ -3053,7 +3060,7 @@ Answers environment variable
--registry PREFIX Image namespace WARMBLY_IMAGE_PREFIX [$DEFAULT_REGISTRY]
--backup-dir PATH Schedule backups into it WARMBLY_BACKUP_DIR
--backup-keep N How many to keep WARMBLY_BACKUP_KEEP [14]
--retention-preset P default|minimal (sets the three windows)
--retention-preset P default|minimal (sets the five windows)
--no-update-check No outbound call at all WARMBLY_UPDATE_CHECK
Other
@@ -3106,6 +3113,7 @@ parse_args() {
minimal)
SYNC_BACKFILL_DAYS=30; SYNC_BACKFILL_MESSAGES=2000
RET_ENGAGEMENT=30; RET_FORMS=30; RET_AUDIT=30
RET_WARMUP_MAIL=7; RET_WARMUP_EVENTS=30
;;
esac
shift
+1 -1
View File
@@ -1 +1 @@
98accfcb1d3a989ffb3bf10281567e9cd5a34c3f10644f29819bed8a1774ca7b install.sh
5857f585c4fff90f0e1fac743013fbce9016e401a14b2071da6ab576f01209d9 install.sh
+37 -1
View File
@@ -87,6 +87,7 @@ import { DitherBarChart } from "@/components/ui/dither";
import WeekdayBitmask from "../campaigns/schedule/WeekdayBitmask";
import { Loading } from "@/components/loader";
import { NumberInput, TextInput } from "@/components/ui/field";
import { clampWarmupRetentionDays } from "@/lib/warmupRetention";
import { useConfirm } from "@/hooks/context/confirm";
import { usePresenceResource } from "@/hooks/PresenceProvider";
import ResourceViewers from "@/components/app/presence/ResourceViewers";
@@ -377,7 +378,7 @@ const EDITABLE: (keyof Inbox)[] = [
"tags", "campaign_limit", "min_wait_time", "reply_to", "save_to_sent",
"warmup_base", "warmup_max", "warmup_increase", "warmup_reply_rate",
"warmup_tag", "warmup_start_time", "warmup_end_time", "warmup_days",
"warmup_placement", "warmup_folder",
"warmup_placement", "warmup_folder", "warmup_retention_days",
];
function Detail({ mailbox, onClose, initialTab = "overview", canWarmup = true }: { mailbox: Inbox; onClose: () => void; initialTab?: string; canWarmup?: boolean }) {
@@ -1365,10 +1366,45 @@ function WarmupPlacementFields({
/>
</FieldShell>
)}
<WarmupRetentionField form={form} update={update} gmail={gmail} />
</>
);
}
/* ── How long warmup mail stays in the mailbox ─────────────────────────── */
// Warmup mail is deleted by Warmbly once it has served its purpose, so a
// mailbox on a fixed quota never fills with it and the owner never has to
// clear the folder by hand. The window is per mailbox; empty follows the
// instance setting (30 days unless the operator changed it).
function WarmupRetentionField({
form,
update,
gmail,
}: {
form: Inbox;
update: (patch: Partial<Inbox>) => void;
gmail: boolean;
}) {
const days = form.warmup_retention_days ?? 0;
const where = gmail ? "moved to Trash, which Gmail empties after 30 days" : "deleted";
return (
<FieldShell
label="Keep warmup mail for (days)"
hint={`Warmup mail older than this is ${where} by Warmbly wherever the setting above keeps it, and the copy Warmbly stores goes with it. 0 follows the instance setting (30 days unless changed); otherwise 3 to 3650. A message is never deleted before its engagement is recorded.`}
>
<NumberInput
value={days}
min={0}
max={3650}
suffix="days"
onChange={(n) => update({ warmup_retention_days: clampWarmupRetentionDays(n) })}
className="w-full h-9"
/>
</FieldShell>
);
}
/* ── Direct-mail open/click tracking (per mailbox, off by default) ────── */
// Campaign mail has always carried a pixel and wrapped links. A reply written
@@ -61,6 +61,11 @@ export default interface Inbox {
warmup_placement?: "folder" | "inbox" | "archive";
/** Folder (Gmail label) for warmup mail. Empty = the default, "Warmbly". */
warmup_folder?: string;
/**
* Days warmup mail stays in this mailbox before Warmbly deletes it from
* the warmup folder. 0 = the instance setting (30 days by default).
*/
warmup_retention_days?: number;
created_at: Date;
updated_at: Date;
}
+19
View File
@@ -0,0 +1,19 @@
import { describe, expect, it } from "vitest";
import { clampWarmupRetentionDays } from "./warmupRetention";
describe("clampWarmupRetentionDays", () => {
it("keeps 0 as the instance setting", () => {
expect(clampWarmupRetentionDays(0)).toBe(0);
expect(clampWarmupRetentionDays(-4)).toBe(0);
expect(clampWarmupRetentionDays(Number.NaN)).toBe(0);
});
it("never yields 1 or 2, which the API refuses", () => {
expect(clampWarmupRetentionDays(1)).toBe(0); // stepping down from 3
expect(clampWarmupRetentionDays(2)).toBe(3); // stepping up from 0, or typed
});
it("keeps the band and drops fractions", () => {
expect(clampWarmupRetentionDays(3)).toBe(3);
expect(clampWarmupRetentionDays(14.9)).toBe(14);
expect(clampWarmupRetentionDays(9999)).toBe(3650);
});
});
+10
View File
@@ -0,0 +1,10 @@
// The API accepts 0 (follow the instance setting) or 3 to 3650 days for how
// long warmup mail stays in a mailbox. The stepper moves by one, so stepping
// down from 3 lands on 0 and stepping up from 0 lands on 3: 1 and 2, which the
// API refuses, can never be sent.
export function clampWarmupRetentionDays(n: number): number {
if (!Number.isFinite(n) || n <= 0) return 0;
const whole = Math.floor(n);
if (whole < 3) return whole === 1 ? 0 : 3;
return Math.min(3650, whole);
}