mirror of
https://github.com/warmbly/warmbly.git
synced 2026-10-03 08:02:04 +00:00
feat: delete warmup mail from each mailbox once past a per-mailbox retention window (email_accounts.warmup_retention_days, else retention.warmup_mail_days, default 30) via a consumer sweep that retires the receipt and sender copy and a worker delete action that trashes on Gmail, deletes on Graph, expunges on IMAP and drops the stored body, prune per-message warmup records after retention.warmup_event_days, and count a warmup deletion as tampering only within 24 hours of arrival and never for a retired message, judging Gmail's Trash label on the same rule
This commit is contained in:
@@ -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 from the warmup folder (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,7 +621,7 @@ 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 && (
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 from the warmup folder: `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 from the warmup folder (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 from the warmup folder by the platform once older than this; 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
|
||||
|
||||
|
||||
@@ -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's warmup folder once older than this |
|
||||
| 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>
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -21425,6 +21425,10 @@
|
||||
"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",
|
||||
"description": "How many days warmup mail stays in this mailbox before Warmbly deletes it from the warmup folder. 0 means the instance setting (30 days by default)."
|
||||
},
|
||||
"timezone": {
|
||||
"type": "string"
|
||||
},
|
||||
@@ -21568,6 +21572,10 @@
|
||||
"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",
|
||||
"description": "How many days warmup mail stays in this mailbox before Warmbly deletes it from the warmup folder: 3 to 3650, or 0 to follow the instance setting."
|
||||
},
|
||||
"tags": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
|
||||
@@ -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 from the warmup folder: 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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,12 +2,15 @@ package worker
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"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"
|
||||
)
|
||||
|
||||
@@ -65,12 +68,32 @@ 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) {
|
||||
if id, err := mail.GoogleData.Client.FindByRFCMessageID(ctx, action.RFCMessageID); err != nil {
|
||||
log.Warn().Err(err).Str("email_id", action.EmailID.String()).Msg("Could not search Gmail for the warmup message")
|
||||
} else {
|
||||
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 {
|
||||
log.Error().Err(err).Str("gmail_id", gmailID).Msg("Failed to delete warmup message (Gmail)")
|
||||
continue
|
||||
}
|
||||
}
|
||||
w.dropWarmupBody(ctx, mail, action, gmailID)
|
||||
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")
|
||||
@@ -154,6 +177,18 @@ 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 {
|
||||
log.Error().Err(err).Str("graph_id", msgID).Msg("Failed to delete warmup message (Graph)")
|
||||
continue
|
||||
}
|
||||
}
|
||||
// The map is keyed by the id the message had when it was
|
||||
// synced, before our own filing moved it, so the live id
|
||||
// resolved above rarely finds it; the control plane's internal
|
||||
// id is what drops the body here.
|
||||
w.dropWarmupBody(ctx, mail, action, action.GmailID)
|
||||
default:
|
||||
log.Warn().Str("action", act).Msg("Unknown warmup action")
|
||||
}
|
||||
@@ -252,12 +287,47 @@ 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, uid); err != nil {
|
||||
log.Error().Err(err).Uint32("uid", uid).Str("folder", boxName).Msg("Failed to delete warmup message (IMAP)")
|
||||
continue
|
||||
}
|
||||
moved = true
|
||||
// The map is keyed by the RFC Message-ID on IMAP.
|
||||
w.dropWarmupBody(ctx, mail, action, action.RFCMessageID)
|
||||
default:
|
||||
log.Warn().Str("action", act).Msg("Unknown warmup action")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 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.
|
||||
func (w *WorkerService) dropWarmupBody(ctx context.Context, mail *wmail.WMail, action models.WarmupEmailAction, providerKey string) {
|
||||
if mail.Storage == nil {
|
||||
return
|
||||
}
|
||||
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
|
||||
}
|
||||
key := config.StorageEndpointEmailBody(mail.UserID, mail.ID, internalID)
|
||||
if err := mail.Storage.Delete(ctx, key); err != nil && !errors.Is(err, storage.ErrNotFound) {
|
||||
log.Warn().Err(err).Str("email_id", action.EmailID.String()).Msg("Failed to drop the stored warmup body")
|
||||
}
|
||||
}
|
||||
|
||||
// locateWarmupMessage resolves the folder and UID an action should act on.
|
||||
//
|
||||
// The arrival folder and UID travel with the action, and for the leg that files
|
||||
|
||||
@@ -55,6 +55,8 @@ 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 expunges one message: the retention window's deletion.
|
||||
DeleteUID(ctx context.Context, mailboxName string, uid uint32) error
|
||||
}
|
||||
|
||||
var _ ImapConn = (*imap.Client)(nil)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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,44 @@ func (c *Client) SetSeen(ctx context.Context, mailboxName string, uids []uint32,
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteUID removes one message from mailboxName for good: \Deleted, then an
|
||||
// expunge limited to that UID where the server offers UIDPLUS, so a message
|
||||
// somebody else flagged in the same folder is not taken along with it. This
|
||||
// is what the retention window asks for once a warmup message has served its
|
||||
// purpose, and IMAP has no Trash of its own to move it to instead.
|
||||
func (c *Client) DeleteUID(ctx context.Context, mailboxName 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)
|
||||
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 c.client.Caps().Has(imap.CapUIDPlus) {
|
||||
if err := c.client.UIDExpunge(set).Close(); err != nil {
|
||||
return fmt.Errorf("uid expunge %d in %q: %w", uid, name, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if err := c.client.Expunge().Close(); err != nil {
|
||||
return fmt.Errorf("expunge %q: %w", name, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,12 @@
|
||||
DROP INDEX IF EXISTS public.idx_warmup_tokens_sent_live;
|
||||
DROP INDEX IF EXISTS public.idx_warmup_received_live;
|
||||
|
||||
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,35 @@
|
||||
-- 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;
|
||||
|
||||
-- The retention sweep walks the live rows oldest first; the retired ones are
|
||||
-- the bulk and never need reading again.
|
||||
CREATE INDEX idx_warmup_received_live
|
||||
ON public.warmup_received (created_at)
|
||||
WHERE retired_at IS NULL;
|
||||
|
||||
CREATE INDEX idx_warmup_tokens_sent_live
|
||||
ON public.warmup_tokens (created_at)
|
||||
WHERE sent_retired_at IS NULL AND sent_message_id <> '';
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
@@ -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
@@ -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's warmup folder after the first window.${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 ${RET_WARMUP_MAIL}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 @@
|
||||
98accfcb1d3a989ffb3bf10281567e9cd5a34c3f10644f29819bed8a1774ca7b install.sh
|
||||
c7a12ec4b7d236358f63e3326c6bf7a125f288d41d46c8ab9268fb21c6256da3 install.sh
|
||||
|
||||
@@ -377,7 +377,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 +1365,46 @@ 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";
|
||||
const invalid = days !== 0 && days < 3;
|
||||
return (
|
||||
<FieldShell
|
||||
label="Keep warmup mail for (days)"
|
||||
hint={`Warmup mail older than this is ${where} by Warmbly, 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={invalid ? <span className="text-red-600">min 3</span> : "days"}
|
||||
onChange={(n) => update({ warmup_retention_days: Math.max(0, Math.min(3650, Math.floor(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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user