From fd818de178268db3f22d13b9cf049cb5f52e4e80 Mon Sep 17 00:00:00 2001 From: Ahmustufa Date: Thu, 24 Sep 2026 15:43:51 +0500 Subject: [PATCH 1/4] feat(db): inbox poll backoff columns, ladder query and fan-out filter MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit An unreachable mail server costs ~60 dial+auth attempts and ~20 dead-letter rows per mailbox per hour today: inbox:sweep re-fans the poll out every 3 minutes forever, because nothing in the worker or coreapi writes mailboxes.status and the poll failure is simply retried and dropped. Adds mailboxes.inbox_poll_failures / inbox_poll_retry_after, read by exactly one query — ListActiveMailboxes, the poll fan-out. MailboxExists and ReserveMailboxSendSlot are untouched, so a backed-off mailbox still SENDS and status stays 'active'; writing status='error' would gate both and cost the user their mailbox, which is what "without deactivating the mailbox" forbids. RecordInboxPollFailure takes the schedule as a parameter array and indexes it by the new failure count, clamped to the last rung. The ladder therefore lives in Go and the cap is the clamp; the delay is added to the DATABASE clock, because ListActiveMailboxes compares it against the database clock. Clearing is folded into both cursor writes and into UpdateMailboxStatus, so a successful poll costs no extra round trip and pause→resume is the manual reset. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01VECELVAKe8Xcp7GH9wGR5t --- internal/platform/db/gen/mailbox.sql.go | 98 ++++- internal/platform/db/gen/models.go | 54 +-- .../db/inboxpollbackoff_integration_test.go | 376 ++++++++++++++++++ ...103858_mailbox_inbox_poll_backoff.down.sql | 6 + ...24103858_mailbox_inbox_poll_backoff.up.sql | 32 ++ internal/platform/db/queries/mailbox.sql | 63 ++- 6 files changed, 589 insertions(+), 40 deletions(-) create mode 100644 internal/platform/db/inboxpollbackoff_integration_test.go create mode 100644 internal/platform/db/migrations/20260924103858_mailbox_inbox_poll_backoff.down.sql create mode 100644 internal/platform/db/migrations/20260924103858_mailbox_inbox_poll_backoff.up.sql diff --git a/internal/platform/db/gen/mailbox.sql.go b/internal/platform/db/gen/mailbox.sql.go index 69b03432..1df46f3b 100644 --- a/internal/platform/db/gen/mailbox.sql.go +++ b/internal/platform/db/gen/mailbox.sql.go @@ -44,7 +44,7 @@ INSERT INTO mailboxes ( $13, $14, $15, $16, $17 ) -RETURNING id, workspace_id, provider, email, display_name, smtp_host, smtp_port, smtp_username, imap_host, imap_port, imap_username, secret_ciphertext, daily_cap, min_interval_seconds, ramp_enabled, ramp_start_cap, ramp_days, status, last_error, last_send_at, last_poll_at, created_at, inbox_last_seen_uid, inbox_uid_validity, inbox_cursor, allow_plaintext +RETURNING id, workspace_id, provider, email, display_name, smtp_host, smtp_port, smtp_username, imap_host, imap_port, imap_username, secret_ciphertext, daily_cap, min_interval_seconds, ramp_enabled, ramp_start_cap, ramp_days, status, last_error, last_send_at, last_poll_at, created_at, inbox_last_seen_uid, inbox_uid_validity, inbox_cursor, allow_plaintext, inbox_poll_failures, inbox_poll_retry_after ` type CreateMailboxParams struct { @@ -115,6 +115,8 @@ func (q *Queries) CreateMailbox(ctx context.Context, arg CreateMailboxParams) (M &i.InboxUidValidity, &i.InboxCursor, &i.AllowPlaintext, + &i.InboxPollFailures, + &i.InboxPollRetryAfter, ) return i, err } @@ -137,7 +139,7 @@ func (q *Queries) DeleteMailbox(ctx context.Context, arg DeleteMailboxParams) (i } const getMailbox = `-- name: GetMailbox :one -SELECT id, workspace_id, provider, email, display_name, smtp_host, smtp_port, smtp_username, imap_host, imap_port, imap_username, secret_ciphertext, daily_cap, min_interval_seconds, ramp_enabled, ramp_start_cap, ramp_days, status, last_error, last_send_at, last_poll_at, created_at, inbox_last_seen_uid, inbox_uid_validity, inbox_cursor, allow_plaintext FROM mailboxes WHERE id = $1 AND workspace_id = $2 +SELECT id, workspace_id, provider, email, display_name, smtp_host, smtp_port, smtp_username, imap_host, imap_port, imap_username, secret_ciphertext, daily_cap, min_interval_seconds, ramp_enabled, ramp_start_cap, ramp_days, status, last_error, last_send_at, last_poll_at, created_at, inbox_last_seen_uid, inbox_uid_validity, inbox_cursor, allow_plaintext, inbox_poll_failures, inbox_poll_retry_after FROM mailboxes WHERE id = $1 AND workspace_id = $2 ` type GetMailboxParams struct { @@ -175,12 +177,16 @@ func (q *Queries) GetMailbox(ctx context.Context, arg GetMailboxParams) (Mailbox &i.InboxUidValidity, &i.InboxCursor, &i.AllowPlaintext, + &i.InboxPollFailures, + &i.InboxPollRetryAfter, ) return i, err } const listActiveMailboxes = `-- name: ListActiveMailboxes :many -SELECT id, workspace_id FROM mailboxes WHERE status = 'active' +SELECT id, workspace_id FROM mailboxes +WHERE status = 'active' + AND (inbox_poll_retry_after IS NULL OR inbox_poll_retry_after <= now()) ` type ListActiveMailboxesRow struct { @@ -190,6 +196,14 @@ type ListActiveMailboxesRow struct { // Mailboxes eligible for inbox polling (reply/bounce detection). The poller // iterates these and calls GetMailbox per id to get IMAP config + cursor. +// +// inbox_poll_retry_after is the poll backoff, and this is the ONLY query that +// reads it. A mailbox whose server is unreachable is skipped for this tick +// rather than dialed again; it is still 'active', so MailboxExists still says +// yes and it still SENDS. The backoff is capped (see +// internal/worker/inbox.DefaultPollBackoff), so the worst case is a mailbox +// polled hourly instead of every three minutes — it recovers on its own the +// first time the server answers. func (q *Queries) ListActiveMailboxes(ctx context.Context) ([]ListActiveMailboxesRow, error) { rows, err := q.db.Query(ctx, listActiveMailboxes) if err != nil { @@ -211,7 +225,7 @@ func (q *Queries) ListActiveMailboxes(ctx context.Context) ([]ListActiveMailboxe } const listMailboxes = `-- name: ListMailboxes :many -SELECT id, workspace_id, provider, email, display_name, smtp_host, smtp_port, smtp_username, imap_host, imap_port, imap_username, secret_ciphertext, daily_cap, min_interval_seconds, ramp_enabled, ramp_start_cap, ramp_days, status, last_error, last_send_at, last_poll_at, created_at, inbox_last_seen_uid, inbox_uid_validity, inbox_cursor, allow_plaintext FROM mailboxes WHERE workspace_id = $1 ORDER BY created_at DESC +SELECT id, workspace_id, provider, email, display_name, smtp_host, smtp_port, smtp_username, imap_host, imap_port, imap_username, secret_ciphertext, daily_cap, min_interval_seconds, ramp_enabled, ramp_start_cap, ramp_days, status, last_error, last_send_at, last_poll_at, created_at, inbox_last_seen_uid, inbox_uid_validity, inbox_cursor, allow_plaintext, inbox_poll_failures, inbox_poll_retry_after FROM mailboxes WHERE workspace_id = $1 ORDER BY created_at DESC ` func (q *Queries) ListMailboxes(ctx context.Context, workspaceID uuid.UUID) ([]Mailbox, error) { @@ -250,6 +264,8 @@ func (q *Queries) ListMailboxes(ctx context.Context, workspaceID uuid.UUID) ([]M &i.InboxUidValidity, &i.InboxCursor, &i.AllowPlaintext, + &i.InboxPollFailures, + &i.InboxPollRetryAfter, ); err != nil { return nil, err } @@ -272,6 +288,55 @@ func (q *Queries) MailboxExists(ctx context.Context, id uuid.UUID) (bool, error) return exists, err } +const recordInboxPollFailure = `-- name: RecordInboxPollFailure :one +UPDATE mailboxes +SET inbox_poll_failures = inbox_poll_failures + 1, + inbox_poll_retry_after = now() + make_interval(secs => + ($1::double precision[])[ + LEAST(inbox_poll_failures + 1, array_length($1::double precision[], 1)) + ]) +WHERE id = $2 AND workspace_id = $3 +RETURNING inbox_poll_failures, inbox_poll_retry_after +` + +type RecordInboxPollFailureParams struct { + BackoffSeconds []float64 `json:"backoff_seconds"` + ID uuid.UUID `json:"id"` + WorkspaceID uuid.UUID `json:"workspace_id"` +} + +type RecordInboxPollFailureRow struct { + InboxPollFailures int32 `json:"inbox_poll_failures"` + InboxPollRetryAfter pgtype.Timestamptz `json:"inbox_poll_retry_after"` +} + +// Records one failed poll and schedules the next attempt, atomically. +// +// @backoff_seconds is the LADDER, passed in as data: rung N is the delay after +// the Nth consecutive failure, and the last rung is the cap (the clamp below is +// what makes it one). The schedule therefore lives in Go — see +// internal/worker/inbox.DefaultPollBackoff — and this statement knows only how +// to index it. A one-rung ladder is how the caller says "this failure cannot +// fix itself, go straight to the cap". +// +// The delay is added to the DATABASE clock, not the worker's, because +// ListActiveMailboxes compares it against the database clock; a few ms of +// host skew would otherwise decide whether a mailbox is due. +// +// Old-value semantics: inside SET, `inbox_poll_failures` is the value BEFORE +// this statement, so `inbox_poll_failures + 1` is the new count and the 1-based +// ladder index at the same time. +// +// Not gated on status: a mailbox paused mid-outage still records the failure it +// just had. It is not polled while paused (ListActiveMailboxes filters on +// status), and resuming clears the counter anyway (UpdateMailboxStatus). +func (q *Queries) RecordInboxPollFailure(ctx context.Context, arg RecordInboxPollFailureParams) (RecordInboxPollFailureRow, error) { + row := q.db.QueryRow(ctx, recordInboxPollFailure, arg.BackoffSeconds, arg.ID, arg.WorkspaceID) + var i RecordInboxPollFailureRow + err := row.Scan(&i.InboxPollFailures, &i.InboxPollRetryAfter) + return i, err +} + const reserveMailboxSendSlot = `-- name: ReserveMailboxSendSlot :one UPDATE mailboxes SET last_send_at = now() @@ -297,7 +362,8 @@ func (q *Queries) ReserveMailboxSendSlot(ctx context.Context, arg ReserveMailbox } const setInboxCursor = `-- name: SetInboxCursor :exec -UPDATE mailboxes SET inbox_last_seen_uid = $3, inbox_uid_validity = $4, last_poll_at = now() +UPDATE mailboxes SET inbox_last_seen_uid = $3, inbox_uid_validity = $4, last_poll_at = now(), + inbox_poll_failures = 0, inbox_poll_retry_after = NULL WHERE id = $1 AND workspace_id = $2 ` @@ -317,6 +383,11 @@ type SetInboxCursorParams struct { // many times it polled successfully, while Gmail and M365 mailboxes stamped it // correctly. That made "never polled" indistinguishable from "polling fine" for // exactly the transport with no provider dashboard to check instead. +// +// The poll backoff is cleared HERE rather than by a second call, for the same +// reason last_poll_at is stamped here: reaching this statement IS the proof the +// mailbox answered, and a separate "clear the failures" round trip could fail +// on its own and leave a healthy mailbox parked on the cap. func (q *Queries) SetInboxCursor(ctx context.Context, arg SetInboxCursorParams) error { _, err := q.db.Exec(ctx, setInboxCursor, arg.ID, @@ -328,7 +399,8 @@ func (q *Queries) SetInboxCursor(ctx context.Context, arg SetInboxCursorParams) } const setInboxCursorString = `-- name: SetInboxCursorString :exec -UPDATE mailboxes SET inbox_cursor = $3, last_poll_at = now() +UPDATE mailboxes SET inbox_cursor = $3, last_poll_at = now(), + inbox_poll_failures = 0, inbox_poll_retry_after = NULL WHERE id = $1 AND workspace_id = $2 ` @@ -338,7 +410,8 @@ type SetInboxCursorStringParams struct { InboxCursor string `json:"inbox_cursor"` } -// Persists an opaque provider cursor (Gmail historyId) after a poll pass. +// Persists an opaque provider cursor (Gmail historyId) after a poll pass, and +// clears the poll backoff for the reason SetInboxCursor gives above. func (q *Queries) SetInboxCursorString(ctx context.Context, arg SetInboxCursorStringParams) error { _, err := q.db.Exec(ctx, setInboxCursorString, arg.ID, arg.WorkspaceID, arg.InboxCursor) return err @@ -364,9 +437,9 @@ func (q *Queries) UpdateMailboxSecret(ctx context.Context, arg UpdateMailboxSecr const updateMailboxStatus = `-- name: UpdateMailboxStatus :one UPDATE mailboxes -SET status = $3, last_error = $4 +SET status = $3, last_error = $4, inbox_poll_failures = 0, inbox_poll_retry_after = NULL WHERE id = $1 AND workspace_id = $2 -RETURNING id, workspace_id, provider, email, display_name, smtp_host, smtp_port, smtp_username, imap_host, imap_port, imap_username, secret_ciphertext, daily_cap, min_interval_seconds, ramp_enabled, ramp_start_cap, ramp_days, status, last_error, last_send_at, last_poll_at, created_at, inbox_last_seen_uid, inbox_uid_validity, inbox_cursor, allow_plaintext +RETURNING id, workspace_id, provider, email, display_name, smtp_host, smtp_port, smtp_username, imap_host, imap_port, imap_username, secret_ciphertext, daily_cap, min_interval_seconds, ramp_enabled, ramp_start_cap, ramp_days, status, last_error, last_send_at, last_poll_at, created_at, inbox_last_seen_uid, inbox_uid_validity, inbox_cursor, allow_plaintext, inbox_poll_failures, inbox_poll_retry_after ` type UpdateMailboxStatusParams struct { @@ -376,6 +449,11 @@ type UpdateMailboxStatusParams struct { LastError string `json:"last_error"` } +// Pause/Resume. Clearing the inbox-poll backoff here makes pause→resume the +// operator's "try it again now" gesture: a mailbox parked on the hour-long cap +// because its password is wrong is otherwise up to an hour from noticing that +// the password was fixed. Pause clearing it too is harmless — a paused mailbox +// is not polled at all. func (q *Queries) UpdateMailboxStatus(ctx context.Context, arg UpdateMailboxStatusParams) (Mailbox, error) { row := q.db.QueryRow(ctx, updateMailboxStatus, arg.ID, @@ -411,6 +489,8 @@ func (q *Queries) UpdateMailboxStatus(ctx context.Context, arg UpdateMailboxStat &i.InboxUidValidity, &i.InboxCursor, &i.AllowPlaintext, + &i.InboxPollFailures, + &i.InboxPollRetryAfter, ) return i, err } diff --git a/internal/platform/db/gen/models.go b/internal/platform/db/gen/models.go index 1160f5c6..5bd251c0 100644 --- a/internal/platform/db/gen/models.go +++ b/internal/platform/db/gen/models.go @@ -674,32 +674,34 @@ type ListMember struct { } type Mailbox struct { - ID uuid.UUID `json:"id"` - WorkspaceID uuid.UUID `json:"workspace_id"` - Provider string `json:"provider"` - Email string `json:"email"` - DisplayName string `json:"display_name"` - SmtpHost string `json:"smtp_host"` - SmtpPort int32 `json:"smtp_port"` - SmtpUsername string `json:"smtp_username"` - ImapHost string `json:"imap_host"` - ImapPort int32 `json:"imap_port"` - ImapUsername string `json:"imap_username"` - SecretCiphertext string `json:"secret_ciphertext"` - DailyCap int32 `json:"daily_cap"` - MinIntervalSeconds int32 `json:"min_interval_seconds"` - RampEnabled bool `json:"ramp_enabled"` - RampStartCap int32 `json:"ramp_start_cap"` - RampDays int32 `json:"ramp_days"` - Status string `json:"status"` - LastError string `json:"last_error"` - LastSendAt pgtype.Timestamptz `json:"last_send_at"` - LastPollAt pgtype.Timestamptz `json:"last_poll_at"` - CreatedAt pgtype.Timestamptz `json:"created_at"` - InboxLastSeenUid int64 `json:"inbox_last_seen_uid"` - InboxUidValidity int64 `json:"inbox_uid_validity"` - InboxCursor string `json:"inbox_cursor"` - AllowPlaintext bool `json:"allow_plaintext"` + ID uuid.UUID `json:"id"` + WorkspaceID uuid.UUID `json:"workspace_id"` + Provider string `json:"provider"` + Email string `json:"email"` + DisplayName string `json:"display_name"` + SmtpHost string `json:"smtp_host"` + SmtpPort int32 `json:"smtp_port"` + SmtpUsername string `json:"smtp_username"` + ImapHost string `json:"imap_host"` + ImapPort int32 `json:"imap_port"` + ImapUsername string `json:"imap_username"` + SecretCiphertext string `json:"secret_ciphertext"` + DailyCap int32 `json:"daily_cap"` + MinIntervalSeconds int32 `json:"min_interval_seconds"` + RampEnabled bool `json:"ramp_enabled"` + RampStartCap int32 `json:"ramp_start_cap"` + RampDays int32 `json:"ramp_days"` + Status string `json:"status"` + LastError string `json:"last_error"` + LastSendAt pgtype.Timestamptz `json:"last_send_at"` + LastPollAt pgtype.Timestamptz `json:"last_poll_at"` + CreatedAt pgtype.Timestamptz `json:"created_at"` + InboxLastSeenUid int64 `json:"inbox_last_seen_uid"` + InboxUidValidity int64 `json:"inbox_uid_validity"` + InboxCursor string `json:"inbox_cursor"` + AllowPlaintext bool `json:"allow_plaintext"` + InboxPollFailures int32 `json:"inbox_poll_failures"` + InboxPollRetryAfter pgtype.Timestamptz `json:"inbox_poll_retry_after"` } type MailboxWorkerAssignment struct { diff --git a/internal/platform/db/inboxpollbackoff_integration_test.go b/internal/platform/db/inboxpollbackoff_integration_test.go new file mode 100644 index 00000000..24c04c27 --- /dev/null +++ b/internal/platform/db/inboxpollbackoff_integration_test.go @@ -0,0 +1,376 @@ +//go:build integration + +package db_test + +import ( + "context" + "testing" + "time" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/inroad/inroad/internal/platform/db" + "github.com/inroad/inroad/internal/platform/db/dbtest" + "github.com/inroad/inroad/internal/platform/db/gen" +) + +// The migration under test, named as a version rather than "one step down" so +// this keeps testing IT after later migrations land. +const ( + beforePollBackoff = 20260923144934 // the version preceding it + pollBackoffColumn = 20260924103858 +) + +// ladder is the production schedule (internal/worker/inbox.DefaultPollBackoff) +// expressed in the seconds the query takes. Duplicated here rather than +// imported because platform/* must never import worker/*; the unit test +// TestDefaultPollBackoffLadder pins the Go values these mirror. +var ladder = []float64{180, 360, 720, 1440, 2880, 3600} + +// capRung is a one-rung ladder — how a caller says "this failure cannot fix +// itself, go straight to the cap". +var capRung = []float64{3600} + +// backoffFixture is one workspace with one active mailbox, plus a second +// workspace used to prove the failure write is workspace-pinned. +type backoffFixture struct { + ws, other uuid.UUID + mailbox uuid.UUID +} + +func seedBackoffFixture(t *testing.T, ctx context.Context, pool *pgxpool.Pool) backoffFixture { + t.Helper() + scalar := func(sql string, args ...any) uuid.UUID { + t.Helper() + var id uuid.UUID + if err := pool.QueryRow(ctx, sql, args...).Scan(&id); err != nil { + t.Fatalf("seed (%s): %v", sql, err) + } + return id + } + ws := scalar(`INSERT INTO workspaces(name) VALUES($1) RETURNING id`, "Poll backoff "+uuid.NewString()) + other := scalar(`INSERT INTO workspaces(name) VALUES($1) RETURNING id`, "Poll backoff other "+uuid.NewString()) + return backoffFixture{ + ws: ws, + other: other, + mailbox: scalar(`INSERT INTO mailboxes(workspace_id,email,secret_ciphertext) + VALUES($1,$2,'sealed') RETURNING id`, ws, "poller-"+uuid.NewString()+"@backoff.test"), + } +} + +// listed reports whether the poll fan-out would pick this mailbox up now. +func listed(t *testing.T, ctx context.Context, q *gen.Queries, id uuid.UUID) bool { + t.Helper() + rows, err := q.ListActiveMailboxes(ctx) + if err != nil { + t.Fatalf("ListActiveMailboxes: %v", err) + } + for _, r := range rows { + if r.ID == id { + return true + } + } + return false +} + +// TestPollBackoffWidensAndCaps walks the ladder the way an hour-long outage +// does: every consecutive failure must schedule the next attempt FURTHER out +// than the last, and the last rung must be a ceiling rather than a step. +// +// The assertion is on the DELTA the database computed (retry_after - now()), +// not on a time this process constructed, because the value is compared against +// the database clock and a few ms of host skew must not be able to decide it. +func TestPollBackoffWidensAndCaps(t *testing.T) { + ctx := context.Background() + pool, q := connectBackoff(t, ctx) + fx := seedBackoffFixture(t, ctx, pool) + + var prev time.Duration + for attempt, want := range []time.Duration{ + 3 * time.Minute, 6 * time.Minute, 12 * time.Minute, + 24 * time.Minute, 48 * time.Minute, time.Hour, time.Hour, time.Hour, + } { + row, err := q.RecordInboxPollFailure(ctx, gen.RecordInboxPollFailureParams{ + BackoffSeconds: ladder, ID: fx.mailbox, WorkspaceID: fx.ws, + }) + if err != nil { + t.Fatalf("attempt %d: RecordInboxPollFailure: %v", attempt+1, err) + } + if int(row.InboxPollFailures) != attempt+1 { + t.Fatalf("attempt %d recorded failure count %d", attempt+1, row.InboxPollFailures) + } + got := scheduledDelay(t, ctx, pool, fx.mailbox) + if got < want-time.Minute || got > want+time.Minute { + t.Errorf("failure %d scheduled the retry in %v, want ~%v", attempt+1, got, want) + } + // The property that actually matters, independent of the exact rungs: + // it widens, and it never widens past the cap. + if attempt > 0 && got < prev && want != time.Hour { + t.Errorf("failure %d scheduled EARLIER (%v) than failure %d (%v)", attempt+1, got, attempt, prev) + } + if got > time.Hour+time.Minute { + t.Errorf("failure %d scheduled %v out — past the cap; an outage would cost the connection", attempt+1, got) + } + prev = got + } +} + +// TestPollBackoffSuppressesSchedulingButNeverEligibility is the one this whole +// change exists to guarantee. A backed-off mailbox must vanish from the POLL +// fan-out and from nothing else: it is still 'active', a send still finds it, +// and it is polled again the moment its retry time passes. +// +// A backoff that quietly became a deactivation is the exact failure the task +// exists to prevent, so both halves are asserted, not just the suppression. +func TestPollBackoffSuppressesSchedulingButNeverEligibility(t *testing.T) { + ctx := context.Background() + pool, q := connectBackoff(t, ctx) + fx := seedBackoffFixture(t, ctx, pool) + + if !listed(t, ctx, q, fx.mailbox) { + t.Fatal("a fresh mailbox is not in the poll fan-out") + } + if _, err := q.RecordInboxPollFailure(ctx, gen.RecordInboxPollFailureParams{ + BackoffSeconds: ladder, ID: fx.mailbox, WorkspaceID: fx.ws, + }); err != nil { + t.Fatalf("RecordInboxPollFailure: %v", err) + } + + if listed(t, ctx, q, fx.mailbox) { + t.Error("a backed-off mailbox is still fanned out; the hammering is unchanged") + } + + // Still ACTIVE, and therefore still sendable. MailboxExists is what the send + // path gates on; writing status='error' here would flip it and cost the user + // their mailbox. + if ok, err := q.MailboxExists(ctx, fx.mailbox); err != nil || !ok { + t.Errorf("MailboxExists = (%v, %v) for a backed-off mailbox; it must still SEND", ok, err) + } + if got := status(t, ctx, pool, fx.mailbox); got != "active" { + t.Errorf("status = %q after a failed poll, want %q — the poller must never write status", got, "active") + } + if _, err := q.ReserveMailboxSendSlot(ctx, gen.ReserveMailboxSendSlotParams{ + ID: fx.mailbox, WorkspaceID: fx.ws, + }); err != nil { + t.Errorf("ReserveMailboxSendSlot on a backed-off mailbox: %v; it must still send", err) + } + + // And it comes back on its own, with no operator action, once due. + if _, err := pool.Exec(ctx, + `UPDATE mailboxes SET inbox_poll_retry_after = now() - interval '1 second' WHERE id = $1`, + fx.mailbox); err != nil { + t.Fatalf("expire the backoff: %v", err) + } + if !listed(t, ctx, q, fx.mailbox) { + t.Error("a mailbox past its retry time is still suppressed; the backoff became a deactivation") + } +} + +// TestSuccessfulPollClearsTheBackoff covers both cursor writes — the IMAP UID +// cursor and the opaque provider cursor — because a mailbox left parked on the +// cap after the server came back is the same outage, just slower. +func TestSuccessfulPollClearsTheBackoff(t *testing.T) { + ctx := context.Background() + pool, q := connectBackoff(t, ctx) + + for _, tc := range []struct { + name string + clear func(t *testing.T, fx backoffFixture) + }{ + {"SetInboxCursor (IMAP)", func(t *testing.T, fx backoffFixture) { + if err := q.SetInboxCursor(ctx, gen.SetInboxCursorParams{ + ID: fx.mailbox, WorkspaceID: fx.ws, InboxLastSeenUid: 11, InboxUidValidity: 3, + }); err != nil { + t.Fatalf("SetInboxCursor: %v", err) + } + }}, + {"SetInboxCursorString (Gmail/Graph)", func(t *testing.T, fx backoffFixture) { + if err := q.SetInboxCursorString(ctx, gen.SetInboxCursorStringParams{ + ID: fx.mailbox, WorkspaceID: fx.ws, InboxCursor: "history-42", + }); err != nil { + t.Fatalf("SetInboxCursorString: %v", err) + } + }}, + {"UpdateMailboxStatus (the operator's resume)", func(t *testing.T, fx backoffFixture) { + if _, err := q.UpdateMailboxStatus(ctx, gen.UpdateMailboxStatusParams{ + ID: fx.mailbox, WorkspaceID: fx.ws, Status: "active", + }); err != nil { + t.Fatalf("UpdateMailboxStatus: %v", err) + } + }}, + } { + t.Run(tc.name, func(t *testing.T) { + fx := seedBackoffFixture(t, ctx, pool) + for range 4 { + if _, err := q.RecordInboxPollFailure(ctx, gen.RecordInboxPollFailureParams{ + BackoffSeconds: ladder, ID: fx.mailbox, WorkspaceID: fx.ws, + }); err != nil { + t.Fatalf("RecordInboxPollFailure: %v", err) + } + } + if listed(t, ctx, q, fx.mailbox) { + t.Fatal("four failures did not suppress the fan-out") + } + + tc.clear(t, fx) + + failures, retryAfter := backoffState(t, ctx, pool, fx.mailbox) + if failures != 0 || retryAfter != nil { + t.Errorf("after a success the backoff is (failures=%d, retry_after=%v), want (0, NULL)", failures, retryAfter) + } + if !listed(t, ctx, q, fx.mailbox) { + t.Error("a recovered mailbox is still suppressed") + } + }) + } +} + +// TestCapRungGoesStraightToTheCeiling is the AUTH shape: a one-rung ladder must +// schedule the cap on the very first failure, because a rejected sign-in cannot +// fix itself and repeating it is what gets an account locked. +func TestCapRungGoesStraightToTheCeiling(t *testing.T) { + ctx := context.Background() + pool, q := connectBackoff(t, ctx) + fx := seedBackoffFixture(t, ctx, pool) + + if _, err := q.RecordInboxPollFailure(ctx, gen.RecordInboxPollFailureParams{ + BackoffSeconds: capRung, ID: fx.mailbox, WorkspaceID: fx.ws, + }); err != nil { + t.Fatalf("RecordInboxPollFailure: %v", err) + } + if got := scheduledDelay(t, ctx, pool, fx.mailbox); got < 59*time.Minute { + t.Errorf("the first permanent failure scheduled a retry in %v, want the ~1h cap", got) + } +} + +// TestRecordInboxPollFailureIsWorkspacePinned: the mailbox UUID is unguessable, +// but the workspace pin is the invariant (docs/security.md 4), and a write that +// ignored it would let one tenant park another tenant's poller. +func TestRecordInboxPollFailureIsWorkspacePinned(t *testing.T) { + ctx := context.Background() + pool, q := connectBackoff(t, ctx) + fx := seedBackoffFixture(t, ctx, pool) + + _, err := q.RecordInboxPollFailure(ctx, gen.RecordInboxPollFailureParams{ + BackoffSeconds: ladder, ID: fx.mailbox, WorkspaceID: fx.other, + }) + if err == nil { + t.Fatal("a foreign workspace id recorded a failure; the query is not workspace-pinned") + } + if failures, retryAfter := backoffState(t, ctx, pool, fx.mailbox); failures != 0 || retryAfter != nil { + t.Errorf("a foreign write moved the row to (failures=%d, retry_after=%v)", failures, retryAfter) + } +} + +// TestPollBackoffMigrationRollsBack walks the migration forwards, back and +// forwards again — what a rollback plus redeploy does. Rolling back must +// restore the previous fan-out exactly: every active mailbox eligible, no +// residue. +// +// Runs on a scratch database because it moves the schema backwards, which would +// break every package sharing the test database. +func TestPollBackoffMigrationRollsBack(t *testing.T) { + ctx := context.Background() + dsn := dbtest.ScratchDSN(t, "poll_backoff_migration") + + if err := db.MigrateTo(dsn, beforePollBackoff); err != nil { + t.Fatalf("migrate to %d: %v", beforePollBackoff, err) + } + pool, err := db.Connect(ctx, dsn) + if err != nil { + t.Fatalf("connect: %v", err) + } + defer pool.Close() // before the scratch DROP, which dbtest's t.Cleanup runs after + + fx := seedBackoffFixture(t, ctx, pool) + q := gen.New(pool) + + if err := db.MigrateTo(dsn, pollBackoffColumn); err != nil { + t.Fatalf("migrate up: %v", err) + } + // An existing row must land eligible, not parked: the column defaults decide + // whether an upgrade silently stops polling every mailbox in the install. + if !listed(t, ctx, q, fx.mailbox) { + t.Fatal("a pre-existing mailbox is not eligible after the migration; the upgrade stopped polling") + } + if _, err := q.RecordInboxPollFailure(ctx, gen.RecordInboxPollFailureParams{ + BackoffSeconds: ladder, ID: fx.mailbox, WorkspaceID: fx.ws, + }); err != nil { + t.Fatalf("RecordInboxPollFailure: %v", err) + } + if listed(t, ctx, q, fx.mailbox) { + t.Fatal("the backoff does not suppress after the migration") + } + + if err := db.MigrateTo(dsn, beforePollBackoff); err != nil { + t.Fatalf("migrate down: %v", err) + } + var columns int + if err := pool.QueryRow(ctx, + `SELECT count(*) FROM information_schema.columns + WHERE table_name = 'mailboxes' + AND column_name IN ('inbox_poll_failures','inbox_poll_retry_after')`).Scan(&columns); err != nil { + t.Fatalf("column count: %v", err) + } + if columns != 0 { + t.Errorf("%d backoff columns survived the down migration", columns) + } + + if err := db.MigrateTo(dsn, pollBackoffColumn); err != nil { + t.Fatalf("migrate up again: %v", err) + } + // Forward again: the mailbox that WAS backed off is eligible, because the + // rollback dropped the state with the column. + if !listed(t, ctx, q, fx.mailbox) { + t.Error("after down+up the previously backed-off mailbox is still suppressed") + } +} + +func connectBackoff(t *testing.T, ctx context.Context) (*pgxpool.Pool, *gen.Queries) { + t.Helper() + dsn := dbtest.DSN(t) + if err := db.Migrate(dsn); err != nil { + t.Fatalf("migrate: %v", err) + } + pool, err := db.Connect(ctx, dsn) + if err != nil { + t.Fatalf("connect: %v", err) + } + t.Cleanup(pool.Close) + return pool, gen.New(pool) +} + +// scheduledDelay is how far out the DATABASE scheduled the next attempt, +// measured by the database, so host↔DB clock skew cannot decide the assertion. +func scheduledDelay(t *testing.T, ctx context.Context, pool *pgxpool.Pool, id uuid.UUID) time.Duration { + t.Helper() + var secs float64 + if err := pool.QueryRow(ctx, + `SELECT extract(epoch FROM inbox_poll_retry_after - now()) FROM mailboxes WHERE id = $1`, + id).Scan(&secs); err != nil { + t.Fatalf("scheduled delay: %v", err) + } + return time.Duration(secs * float64(time.Second)) +} + +func backoffState(t *testing.T, ctx context.Context, pool *pgxpool.Pool, id uuid.UUID) (int, *time.Time) { + t.Helper() + var failures int + var retryAfter *time.Time + if err := pool.QueryRow(ctx, + `SELECT inbox_poll_failures, inbox_poll_retry_after FROM mailboxes WHERE id = $1`, + id).Scan(&failures, &retryAfter); err != nil { + t.Fatalf("backoff state: %v", err) + } + return failures, retryAfter +} + +func status(t *testing.T, ctx context.Context, pool *pgxpool.Pool, id uuid.UUID) string { + t.Helper() + var s string + if err := pool.QueryRow(ctx, `SELECT status FROM mailboxes WHERE id = $1`, id).Scan(&s); err != nil { + t.Fatalf("status: %v", err) + } + return s +} diff --git a/internal/platform/db/migrations/20260924103858_mailbox_inbox_poll_backoff.down.sql b/internal/platform/db/migrations/20260924103858_mailbox_inbox_poll_backoff.down.sql new file mode 100644 index 00000000..6776e8b2 --- /dev/null +++ b/internal/platform/db/migrations/20260924103858_mailbox_inbox_poll_backoff.down.sql @@ -0,0 +1,6 @@ +-- Dropping the columns restores the previous behaviour exactly: the poll +-- fan-out's predicate goes back to `status = 'active'` alone, so every active +-- mailbox becomes eligible on the next sweep. Nothing else reads them. +ALTER TABLE mailboxes + DROP COLUMN IF EXISTS inbox_poll_failures, + DROP COLUMN IF EXISTS inbox_poll_retry_after; diff --git a/internal/platform/db/migrations/20260924103858_mailbox_inbox_poll_backoff.up.sql b/internal/platform/db/migrations/20260924103858_mailbox_inbox_poll_backoff.up.sql new file mode 100644 index 00000000..96e4617b --- /dev/null +++ b/internal/platform/db/migrations/20260924103858_mailbox_inbox_poll_backoff.up.sql @@ -0,0 +1,32 @@ +-- Widening backoff for an unreachable mail server, WITHOUT deactivating the +-- mailbox. +-- +-- Today a mailbox whose IMAP server is down is re-fanned-out by inbox:sweep +-- every 3 minutes forever: ~60 dial+auth attempts and ~20 dead-letter rows per +-- mailbox per hour. That reconnect hammering is itself what provokes the per-IP +-- sign-in throttles and rate limits the fleet architecture exists to avoid, and +-- it buries real failures in dead-letter noise. +-- +-- These two columns suppress SCHEDULING only. They are read by exactly one +-- query, ListActiveMailboxes (the inbox poll fan-out). They are deliberately +-- NOT read by MailboxExists or ReserveMailboxSendSlot, so a backed-off mailbox +-- still SENDS, and `status` is never written by the poller: writing +-- status='error' would gate both of those and cost the user their mailbox, +-- which is the opposite of what "without deactivating the mailbox" asks for. +-- +-- No index. ListActiveMailboxes is a full scan of the mailboxes table with or +-- without this predicate (there is no index on status either); the table holds +-- one row per connected mailbox, and an index that only ever narrows an already +-- sequential scan would be write amplification for nothing. +ALTER TABLE mailboxes + -- Consecutive failed poll attempts. Reset to 0 by every successful poll + -- (folded into SetInboxCursor / SetInboxCursorString, so a success costs no + -- extra round trip) and by any operator status change (pause/resume, which + -- is therefore the manual "try it again now" gesture). + ADD COLUMN IF NOT EXISTS inbox_poll_failures INTEGER NOT NULL DEFAULT 0, + -- The earliest time the poll fan-out may consider this mailbox again. + -- NULL means "eligible now", which is what every existing row gets and what + -- a successful poll restores. Computed from the DATABASE clock (now() plus + -- the rung), never the worker's, because it is compared against the + -- database clock in ListActiveMailboxes. + ADD COLUMN IF NOT EXISTS inbox_poll_retry_after TIMESTAMPTZ; diff --git a/internal/platform/db/queries/mailbox.sql b/internal/platform/db/queries/mailbox.sql index 184ba265..e940f20c 100644 --- a/internal/platform/db/queries/mailbox.sql +++ b/internal/platform/db/queries/mailbox.sql @@ -26,8 +26,13 @@ SELECT * FROM mailboxes WHERE workspace_id = $1 ORDER BY created_at DESC; SELECT count(*) FROM mailboxes WHERE workspace_id = $1 AND email = $2; -- name: UpdateMailboxStatus :one +-- Pause/Resume. Clearing the inbox-poll backoff here makes pause→resume the +-- operator's "try it again now" gesture: a mailbox parked on the hour-long cap +-- because its password is wrong is otherwise up to an hour from noticing that +-- the password was fixed. Pause clearing it too is harmless — a paused mailbox +-- is not polled at all. UPDATE mailboxes -SET status = $3, last_error = $4 +SET status = $3, last_error = $4, inbox_poll_failures = 0, inbox_poll_retry_after = NULL WHERE id = $1 AND workspace_id = $2 RETURNING *; @@ -51,7 +56,17 @@ SELECT EXISTS (SELECT 1 FROM mailboxes WHERE id = $1 AND status = 'active'); -- name: ListActiveMailboxes :many -- Mailboxes eligible for inbox polling (reply/bounce detection). The poller -- iterates these and calls GetMailbox per id to get IMAP config + cursor. -SELECT id, workspace_id FROM mailboxes WHERE status = 'active'; +-- +-- inbox_poll_retry_after is the poll backoff, and this is the ONLY query that +-- reads it. A mailbox whose server is unreachable is skipped for this tick +-- rather than dialed again; it is still 'active', so MailboxExists still says +-- yes and it still SENDS. The backoff is capped (see +-- internal/worker/inbox.DefaultPollBackoff), so the worst case is a mailbox +-- polled hourly instead of every three minutes — it recovers on its own the +-- first time the server answers. +SELECT id, workspace_id FROM mailboxes +WHERE status = 'active' + AND (inbox_poll_retry_after IS NULL OR inbox_poll_retry_after <= now()); -- name: SetInboxCursor :exec -- Persists the IMAP poll cursor after a poll pass, so the next pass resumes @@ -63,7 +78,13 @@ SELECT id, workspace_id FROM mailboxes WHERE status = 'active'; -- many times it polled successfully, while Gmail and M365 mailboxes stamped it -- correctly. That made "never polled" indistinguishable from "polling fine" for -- exactly the transport with no provider dashboard to check instead. -UPDATE mailboxes SET inbox_last_seen_uid = $3, inbox_uid_validity = $4, last_poll_at = now() +-- +-- The poll backoff is cleared HERE rather than by a second call, for the same +-- reason last_poll_at is stamped here: reaching this statement IS the proof the +-- mailbox answered, and a separate "clear the failures" round trip could fail +-- on its own and leave a healthy mailbox parked on the cap. +UPDATE mailboxes SET inbox_last_seen_uid = $3, inbox_uid_validity = $4, last_poll_at = now(), + inbox_poll_failures = 0, inbox_poll_retry_after = NULL WHERE id = $1 AND workspace_id = $2; -- name: UpdateMailboxSecret :exec @@ -73,6 +94,38 @@ UPDATE mailboxes SET secret_ciphertext = $3 WHERE id = $1 AND workspace_id = $2; -- name: SetInboxCursorString :exec --- Persists an opaque provider cursor (Gmail historyId) after a poll pass. -UPDATE mailboxes SET inbox_cursor = $3, last_poll_at = now() +-- Persists an opaque provider cursor (Gmail historyId) after a poll pass, and +-- clears the poll backoff for the reason SetInboxCursor gives above. +UPDATE mailboxes SET inbox_cursor = $3, last_poll_at = now(), + inbox_poll_failures = 0, inbox_poll_retry_after = NULL WHERE id = $1 AND workspace_id = $2; + +-- name: RecordInboxPollFailure :one +-- Records one failed poll and schedules the next attempt, atomically. +-- +-- @backoff_seconds is the LADDER, passed in as data: rung N is the delay after +-- the Nth consecutive failure, and the last rung is the cap (the clamp below is +-- what makes it one). The schedule therefore lives in Go — see +-- internal/worker/inbox.DefaultPollBackoff — and this statement knows only how +-- to index it. A one-rung ladder is how the caller says "this failure cannot +-- fix itself, go straight to the cap". +-- +-- The delay is added to the DATABASE clock, not the worker's, because +-- ListActiveMailboxes compares it against the database clock; a few ms of +-- host skew would otherwise decide whether a mailbox is due. +-- +-- Old-value semantics: inside SET, `inbox_poll_failures` is the value BEFORE +-- this statement, so `inbox_poll_failures + 1` is the new count and the 1-based +-- ladder index at the same time. +-- +-- Not gated on status: a mailbox paused mid-outage still records the failure it +-- just had. It is not polled while paused (ListActiveMailboxes filters on +-- status), and resuming clears the counter anyway (UpdateMailboxStatus). +UPDATE mailboxes +SET inbox_poll_failures = inbox_poll_failures + 1, + inbox_poll_retry_after = now() + make_interval(secs => + (@backoff_seconds::double precision[])[ + LEAST(inbox_poll_failures + 1, array_length(@backoff_seconds::double precision[], 1)) + ]) +WHERE id = @id AND workspace_id = @workspace_id +RETURNING inbox_poll_failures, inbox_poll_retry_after; From dab834807f87945679f552deb1d0e241f86a699e Mon Sep 17 00:00:00 2001 From: Ahmustufa Date: Thu, 24 Sep 2026 15:48:20 +0500 Subject: [PATCH 2/4] feat(mail): classify a provider connection failure into four shapes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A refused connection, a black-holed one, a wrong password and an SSRF refusal are four different signals, and the inbox poller needs to treat them differently: the first two fix themselves, the last two do not, and repeating a rejected sign-in is how a provider decides to lock an account. ClassifyConnectFailure reads sentinels, types and status codes only — never error text — and reads TRANSPORT EVIDENCE BEFORE the new ErrAuthRejected sentinel, so a connection that drops mid-LOGIN is not mistaken for a bad password. Cancellation is separated out so a shutting-down worker does not record failures against healthy mailboxes. Driven through the real IMAP client against PR 1's in-process server, which gains one capability: stopListening, a listener closed under the transport seam so the next dial gets a real ECONNREFUSED from the kernel rather than a stubbed reader. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01VECELVAKe8Xcp7GH9wGR5t --- internal/platform/mail/connectfailure.go | 201 +++++++++++++++ internal/platform/mail/connectfailure_test.go | 230 ++++++++++++++++++ internal/platform/mail/imapauth.go | 20 +- internal/platform/mail/imapfake_test.go | 15 ++ 4 files changed, 464 insertions(+), 2 deletions(-) create mode 100644 internal/platform/mail/connectfailure.go create mode 100644 internal/platform/mail/connectfailure_test.go diff --git a/internal/platform/mail/connectfailure.go b/internal/platform/mail/connectfailure.go new file mode 100644 index 00000000..53276d3f --- /dev/null +++ b/internal/platform/mail/connectfailure.go @@ -0,0 +1,201 @@ +package mail + +import ( + "context" + "crypto/tls" + "errors" + "io" + "net" + "syscall" + + "google.golang.org/api/googleapi" +) + +// ConnectFailure says WHY a provider connection failed, in the four shapes that +// call for different handling. It is the vocabulary the inbox poller's backoff +// is expressed in (internal/worker/inbox), and it lives here because the error +// shapes it reads — go-imap's replies, net's syscall errnos, tls's verification +// errors, googleapi's statuses — are this package's business and nobody else's. +// +// It is deliberately NOT Retryable() (retry.go). That answers a different and +// narrower question about a SEND: "can this be retried without risking a double +// delivery?". A poll cannot deliver anything twice, so nothing here is about +// safety; it is about how soon it is worth dialing again, and about not +// repeating a rejected sign-in. +type ConnectFailure uint8 + +const ( + // ConnectFailureNone is a nil error. + ConnectFailureNone ConnectFailure = iota + // ConnectFailureAborted is OUR OWN context being cancelled — a worker + // shutting down, or a caller that gave up. It says nothing about the server, + // and counting it would record a failure against a mailbox nothing is wrong + // with. + // + // It is recognised where cancellation is observable: the DNS resolve and the + // dial, the two legs that take a context. Once an IMAP session is up, + // go-imap reads under a socket deadline instead (dialIMAP's c.Timeout and + // newDeadlineConn), so a cancel arriving mid-session reads as a timeout and + // lands in ConnectFailureTransport — one recorded failure, one rung, which + // the next successful poll clears. + ConnectFailureAborted + // ConnectFailureTransport is the server being unreachable or unresponsive: + // refused, reset, black-holed, DNS-less, or a TLS handshake that did not + // complete. It is the one class that routinely fixes itself. + ConnectFailureTransport + // ConnectFailureAuth is the server understanding us perfectly and refusing + // the credential — a wrong password, a revoked token, or a server that will + // speak no mechanism we implement. It cannot fix itself, and every repeat is + // another rejected sign-in against the account. + ConnectFailureAuth + // ConnectFailurePolicy is OUR refusal, not the server's: the SSRF guard + // rejected the host or port. No socket was opened. Nothing but an operator + // editing the mailbox will change the answer. + ConnectFailurePolicy + // ConnectFailureUnknown is everything else — a provider status we do not + // recognise, a malformed reply, a library error. The caller should assume it + // will recur, because the alternative assumption is the one that produced the + // hammering in the first place. + ConnectFailureUnknown +) + +// String is the stable label this classification is logged under. +func (c ConnectFailure) String() string { + switch c { + case ConnectFailureNone: + return "none" + case ConnectFailureAborted: + return "aborted" + case ConnectFailureTransport: + return "transport" + case ConnectFailureAuth: + return "auth" + case ConnectFailurePolicy: + return "policy" + default: + return "unknown" + } +} + +// ClassifyConnectFailure sorts a provider connection error into one of the +// shapes above. +// +// THE ORDER IS THE DESIGN, not an implementation detail: +// +// - Cancellation is read first, because a cancelled dial arrives wrapped in +// the same *net.OpError a refused one does and would otherwise be +// indistinguishable from the server's fault. +// - Our own SSRF refusal next: it is a policy verdict, and it never touched +// the network, so no transport evidence exists to weigh against it. +// - TRANSPORT BEFORE AUTH, which is the subtle one. An authentication attempt +// that dies because the connection dropped mid-LOGIN is wrapped as an auth +// failure (it failed at the auth step) but is transport in substance. +// Reading transport first means a flaky network is never mistaken for a bad +// password — and mistaking it for one would park a healthy mailbox on the +// hour-long cap. +// +// It never inspects error TEXT. Every branch is a sentinel, a type or a status +// code, so a library changing its wording cannot silently reclassify a mailbox. +func ClassifyConnectFailure(err error) ConnectFailure { + if err == nil { + return ConnectFailureNone + } + if errors.Is(err, context.Canceled) { + return ConnectFailureAborted + } + if errors.Is(err, ErrHostNotPermitted) { + return ConnectFailurePolicy + } + if isTransportFailure(err) { + return ConnectFailureTransport + } + if errors.Is(err, ErrAuthRejected) || errors.Is(err, ErrNoIMAPAuthMechanism) { + return ConnectFailureAuth + } + if status, ok := apiStatus(err); ok { + return classifyAPIStatus(status) + } + return ConnectFailureUnknown +} + +// isTransportFailure reports whether the error is the network failing rather +// than anyone's decision. +// +// context.DeadlineExceeded counts, and is the case the fake IMAP server's +// silent-after-greeting mode exists for: a server that completes the TCP +// handshake and then says nothing is the most expensive failure there is, +// because each attempt holds a connection for the whole timeout. +func isTransportFailure(err error) bool { + var netErr net.Error + if errors.As(err, &netErr) && netErr.Timeout() { + return true + } + if errors.Is(err, context.DeadlineExceeded) || + errors.Is(err, io.EOF) || + errors.Is(err, io.ErrUnexpectedEOF) || + errors.Is(err, syscall.ECONNREFUSED) || + errors.Is(err, syscall.ECONNRESET) || + errors.Is(err, syscall.ECONNABORTED) || + errors.Is(err, syscall.EHOSTUNREACH) || + errors.Is(err, syscall.ENETUNREACH) || + errors.Is(err, syscall.EPIPE) || + errors.Is(err, syscall.ETIMEDOUT) { + return true + } + // A TLS handshake that did not complete: an expired or self-signed + // certificate, a protocol mismatch, an alert. It is not a credential and it + // is not our policy — the connection never came up — and it does resolve on + // its own when the operator's certificate is renewed. + var certErr *tls.CertificateVerificationError + if errors.As(err, &certErr) { + return true + } + var recordErr tls.RecordHeaderError + if errors.As(err, &recordErr) { + return true + } + var alertErr tls.AlertError + if errors.As(err, &alertErr) { + return true + } + var dnsErr *net.DNSError + if errors.As(err, &dnsErr) { + return true + } + // The catch-all for a dial/read/write that failed for a reason not named + // above. Last, so a more specific branch always wins. + var opErr *net.OpError + return errors.As(err, &opErr) +} + +// apiStatus extracts the HTTP status from a provider API failure, from either +// shape this package produces: its own *APIError, or the googleapi error the +// Gmail client library returns directly. +func apiStatus(err error) (int, bool) { + var apiErr *APIError + if errors.As(err, &apiErr) { + return apiErr.Status, true + } + var gerr *googleapi.Error + if errors.As(err, &gerr) { + return gerr.Code, true + } + return 0, false +} + +// classifyAPIStatus maps a provider's HTTP answer onto the same four shapes. +// +// 401/403 is the API transport's wrong password: a revoked or unconsented token. +// 429 and 5xx are the provider telling us to come back later, which is exactly +// what a widening backoff does. Anything else is a request we got wrong, which +// will recur until the code changes — unknown, not transport. +func classifyAPIStatus(status int) ConnectFailure { + switch { + case status == 401 || status == 403: + return ConnectFailureAuth + case status == 429 || status >= 500: + return ConnectFailureTransport + default: + return ConnectFailureUnknown + } +} diff --git a/internal/platform/mail/connectfailure_test.go b/internal/platform/mail/connectfailure_test.go new file mode 100644 index 00000000..d514949a --- /dev/null +++ b/internal/platform/mail/connectfailure_test.go @@ -0,0 +1,230 @@ +package mail + +import ( + "context" + "crypto/tls" + "errors" + "fmt" + "io" + "net" + "syscall" + "testing" + "time" + + "google.golang.org/api/googleapi" +) + +// The four signals this classification exists to tell apart, driven through the +// REAL IMAP client against the in-process server rather than against a +// hand-built error value. A refused connection, a black-holed one, a wrong +// password and a server that speaks no mechanism we implement all arrive here +// as whatever go-imap and net actually produce — which is the only version of +// them that matters, because the classifier's whole job is to read those. +// +// These must not call t.Parallel(): startFakeIMAP swaps a package-level +// transport (see its doc). + +// A dead listener is a real ECONNREFUSED from the kernel, not a stubbed reader. +func TestClassifyRefusedConnectionIsTransport(t *testing.T) { + srv := startFakeIMAP(t, imapScript{User: "u", Pass: "p"}) + srv.stopListening(t) + + r := &NetInboxReader{Timeout: 2 * time.Second} + _, _, err := r.CurrentState(t.Context(), fakeIMAPConfig("u", "p")) + if err == nil { + t.Fatal("a dead listener returned no error") + } + if got := ClassifyConnectFailure(err); got != ConnectFailureTransport { + t.Errorf("ClassifyConnectFailure(%v) = %v, want %v", err, got, ConnectFailureTransport) + } +} + +// Accept-then-stall: the most expensive failure there is, because each attempt +// holds a connection for the whole timeout. It must widen, not be mistaken for +// a credential problem. +func TestClassifySilentServerIsTransport(t *testing.T) { + startFakeIMAP(t, imapScript{SilentAfterGreeting: true}) + + r := &NetInboxReader{Timeout: 500 * time.Millisecond} + _, _, err := r.CurrentState(t.Context(), fakeIMAPConfig("u", "p")) + if err == nil { + t.Fatal("a silent server returned no error") + } + if got := ClassifyConnectFailure(err); got != ConnectFailureTransport { + t.Errorf("ClassifyConnectFailure(%v) = %v, want %v", err, got, ConnectFailureTransport) + } +} + +// A wrong password is not a network outage. It reaches the poller through the +// plain LOGIN command and through a SASL mechanism alike, so both are driven. +func TestClassifyRejectedCredentialIsAuth(t *testing.T) { + for _, tc := range []struct { + name string + script imapScript + }{ + {"LOGIN command", imapScript{User: "u", Pass: "right"}}, + {"AUTH=PLAIN, no LOGIN fallback", imapScript{ + User: "u", Pass: "right", Caps: []string{"AUTH=PLAIN", "LOGINDISABLED"}, + }}, + {"AUTH=CRAM-MD5, no LOGIN fallback", imapScript{ + User: "u", Pass: "right", Caps: []string{"AUTH=CRAM-MD5", "LOGINDISABLED"}, + }}, + } { + t.Run(tc.name, func(t *testing.T) { + startFakeIMAP(t, tc.script) + + r := &NetInboxReader{Timeout: 2 * time.Second} + _, _, err := r.CurrentState(t.Context(), fakeIMAPConfig("u", "wrong")) + if err == nil { + t.Fatal("a wrong password authenticated") + } + if got := ClassifyConnectFailure(err); got != ConnectFailureAuth { + t.Errorf("ClassifyConnectFailure(%v) = %v, want %v", err, got, ConnectFailureAuth) + } + // The credential must not be recoverable from the error a log line + // would print (docs/security.md, credential handling). + if msg := err.Error(); containsCommand([]string{msg}, "wrong") { + t.Errorf("the classified auth error carries the credential: %q", msg) + } + }) + } +} + +// A server that disables LOGIN and offers only mechanisms this package does not +// implement is the "this server wants GSSAPI" configuration fact: nothing but +// an operator changes the answer, so it is auth, not transport. +func TestClassifyNoSupportedMechanismIsAuth(t *testing.T) { + startFakeIMAP(t, imapScript{User: "u", Pass: "p", Caps: []string{"AUTH=GSSAPI", "LOGINDISABLED"}}) + + r := &NetInboxReader{Timeout: 2 * time.Second} + _, _, err := r.CurrentState(t.Context(), fakeIMAPConfig("u", "p")) + if !errors.Is(err, ErrNoIMAPAuthMechanism) { + t.Fatalf("got %v, want ErrNoIMAPAuthMechanism", err) + } + if got := ClassifyConnectFailure(err); got != ConnectFailureAuth { + t.Errorf("ClassifyConnectFailure(%v) = %v, want %v", err, got, ConnectFailureAuth) + } +} + +// Our own refusal, before any socket exists. +func TestClassifySSRFRefusalIsPolicy(t *testing.T) { + startFakeIMAP(t, imapScript{User: "u", Pass: "p"}) + + r := &NetInboxReader{Timeout: 2 * time.Second} // AllowPrivate false + _, _, err := r.CurrentState(t.Context(), IMAPConfig{ + Host: "169.254.169.254", Port: 993, Username: "u", Password: "p", + }) + if !errors.Is(err, ErrHostNotPermitted) { + t.Fatalf("got %v, want ErrHostNotPermitted", err) + } + if got := ClassifyConnectFailure(err); got != ConnectFailurePolicy { + t.Errorf("ClassifyConnectFailure(%v) = %v, want %v", err, got, ConnectFailurePolicy) + } +} + +// A cancelled context is not the server's fault, and counting it would record a +// failure against a mailbox nothing is wrong with. +// +// The cancellation is delivered where it is actually observable — the DIAL, +// which is the one leg of an IMAP poll that takes a context. Once the session is +// up, go-imap reads under a socket deadline rather than a context (see +// dialIMAP's c.Timeout and newDeadlineConn), so a cancel arriving mid-session +// surfaces as an i/o timeout and classifies as transport. That imprecision is +// bounded — one recorded failure, one rung — and narrowing it would mean +// rebuilding the reader's I/O around the context, which is not this change. +func TestClassifyCancellationIsNotTheServersFault(t *testing.T) { + startFakeIMAP(t, imapScript{User: "u", Pass: "p"}) + + ctx, cancel := context.WithCancel(t.Context()) + cancel() + + r := &NetInboxReader{Timeout: 2 * time.Second} + _, _, err := r.CurrentState(ctx, fakeIMAPConfig("u", "p")) + if err == nil { + t.Fatal("a cancelled poll returned no error") + } + if got := ClassifyConnectFailure(err); got != ConnectFailureAborted { + t.Errorf("ClassifyConnectFailure(%v) = %v, want %v", err, got, ConnectFailureAborted) + } +} + +// The ordering rule, stated as a test: an auth-step failure whose real cause is +// the connection dropping is TRANSPORT. Reading the sentinel first would park a +// mailbox on a flaky network at the hour-long cap as though its password were +// wrong. +func TestTransportEvidenceOutranksTheAuthSentinel(t *testing.T) { + err := fmt.Errorf("%w: imap login: %w", ErrAuthRejected, io.ErrUnexpectedEOF) + if got := ClassifyConnectFailure(err); got != ConnectFailureTransport { + t.Errorf("ClassifyConnectFailure(%v) = %v, want %v", err, got, ConnectFailureTransport) + } +} + +// The API transports' answers. 401/403 is their wrong password; 429 and 5xx are +// "come back later", which is what a widening backoff is. +func TestClassifyProviderAPIStatuses(t *testing.T) { + for _, tc := range []struct { + name string + err error + want ConnectFailure + }{ + {"graph 401", &APIError{Provider: "m365", Op: "inbox delta", Status: 401}, ConnectFailureAuth}, + {"graph 403", &APIError{Provider: "m365", Op: "inbox delta", Status: 403}, ConnectFailureAuth}, + {"graph 429", &APIError{Provider: "m365", Op: "inbox delta", Status: 429}, ConnectFailureTransport}, + {"graph 503", &APIError{Provider: "m365", Op: "inbox delta", Status: 503}, ConnectFailureTransport}, + {"graph 400", &APIError{Provider: "m365", Op: "inbox delta", Status: 400}, ConnectFailureUnknown}, + {"gmail 401", &googleapi.Error{Code: 401}, ConnectFailureAuth}, + {"gmail 429", &googleapi.Error{Code: 429}, ConnectFailureTransport}, + {"wrapped gmail 500", fmt.Errorf("gmail history: %w", &googleapi.Error{Code: 500}), ConnectFailureTransport}, + } { + t.Run(tc.name, func(t *testing.T) { + if got := ClassifyConnectFailure(tc.err); got != tc.want { + t.Errorf("ClassifyConnectFailure(%v) = %v, want %v", tc.err, got, tc.want) + } + }) + } +} + +// The remaining transport shapes, as the values net and crypto/tls hand us. No +// server can produce a DNS failure or a certificate error against an in-process +// listener, so these are constructed — but they are the library's own types, +// not strings, which is the property that makes the classifier robust. +func TestClassifyTransportShapes(t *testing.T) { + for _, tc := range []struct { + name string + err error + want ConnectFailure + }{ + {"nil", nil, ConnectFailureNone}, + {"refused", fmt.Errorf("imap dial: %w", syscall.ECONNREFUSED), ConnectFailureTransport}, + {"reset", fmt.Errorf("imap fetch: %w", syscall.ECONNRESET), ConnectFailureTransport}, + {"host unreachable", fmt.Errorf("imap dial: %w", syscall.EHOSTUNREACH), ConnectFailureTransport}, + {"eof", fmt.Errorf("imap dial: %w", io.EOF), ConnectFailureTransport}, + {"deadline", fmt.Errorf("imap dial: %w", context.DeadlineExceeded), ConnectFailureTransport}, + {"dns", &net.DNSError{Err: "no such host", Name: "mail.example.test", IsNotFound: true}, ConnectFailureTransport}, + {"tls certificate", &tls.CertificateVerificationError{Err: errors.New("expired")}, ConnectFailureTransport}, + {"tls alert", tls.AlertError(80), ConnectFailureTransport}, + {"not a network error at all", errors.New("imap select: NO mailbox does not exist"), ConnectFailureUnknown}, + } { + t.Run(tc.name, func(t *testing.T) { + if got := ClassifyConnectFailure(tc.err); got != tc.want { + t.Errorf("ClassifyConnectFailure(%v) = %v, want %v", tc.err, got, tc.want) + } + }) + } +} + +// The labels are logged, so they are part of the operator-facing contract. +func TestConnectFailureLabels(t *testing.T) { + for kind, want := range map[ConnectFailure]string{ + ConnectFailureNone: "none", + ConnectFailureAborted: "aborted", + ConnectFailureTransport: "transport", + ConnectFailureAuth: "auth", + ConnectFailurePolicy: "policy", + ConnectFailureUnknown: "unknown", + } { + if got := kind.String(); got != want { + t.Errorf("ConnectFailure(%d).String() = %q, want %q", kind, got, want) + } + } +} diff --git a/internal/platform/mail/imapauth.go b/internal/platform/mail/imapauth.go index 2d647e36..4e65fb3c 100644 --- a/internal/platform/mail/imapauth.go +++ b/internal/platform/mail/imapauth.go @@ -36,6 +36,22 @@ const ( mechPlain = "PLAIN" ) +// ErrAuthRejected marks a failure at the AUTHENTICATION step: the server +// understood the command and refused the credential. It wraps rather than +// replaces the server's own error, so the detail an operator needs is still +// there and the sentinel is still testable. +// +// It exists because "the password is wrong" and "the server is down" want +// opposite handling from the poller — see ClassifyConnectFailure, which reads +// transport evidence BEFORE this sentinel precisely because a connection that +// drops mid-LOGIN also fails here. +// +// No error wrapped by it carries the password: the digests, prompts and PLAIN +// payload stay inside the sasl.Client implementations below, and the underlying +// errors carry only what the server said (docs/security.md, credential +// handling). +var ErrAuthRejected = errors.New("mail: authentication rejected") + // ErrNoIMAPAuthMechanism is returned when a server disables the LOGIN command // and advertises no SASL mechanism this package implements. It is a // configuration fact an operator can act on ("this server wants GSSAPI"), which @@ -83,7 +99,7 @@ func authenticateIMAP(c *client.Client, cfg IMAPConfig) error { if mechErr == nil { return nil } - mechErr = fmt.Errorf("imap authenticate %s: %w", mech, mechErr) + mechErr = fmt.Errorf("%w: imap authenticate %s: %w", ErrAuthRejected, mech, mechErr) if loginDisabled { return mechErr } @@ -115,7 +131,7 @@ func pickIMAPMechanism(c *client.Client) (string, error) { // and the fallback for a mechanism that fails. func loginIMAP(c *client.Client, cfg IMAPConfig) error { if err := c.Login(cfg.Username, cfg.Password); err != nil { - return fmt.Errorf("imap login: %w", err) + return fmt.Errorf("%w: imap login: %w", ErrAuthRejected, err) } return nil } diff --git a/internal/platform/mail/imapfake_test.go b/internal/platform/mail/imapfake_test.go index 67a33710..a57c7264 100644 --- a/internal/platform/mail/imapfake_test.go +++ b/internal/platform/mail/imapfake_test.go @@ -145,6 +145,21 @@ func startFakeIMAP(t *testing.T, script imapScript) *fakeIMAP { return s } +// stopListening closes the listener while leaving the transport seam pointed at +// its (now dead) address, so the next dial gets a REAL ECONNREFUSED from the +// kernel. +// +// It is how "the mail server is down" is tested without mocking the reader: an +// unreachable-server test that substitutes a failing io.Reader is a test of the +// substitute. t.Cleanup closes the listener a second time, which is a no-op +// error this harness already ignores. +func (s *fakeIMAP) stopListening(t *testing.T) { + t.Helper() + if err := s.ln.Close(); err != nil { + t.Fatalf("closing the fake listener: %v", err) + } +} + // commandLog returns every command line the server received, in order, with the // client's tag stripped. Credentials are NOT stripped: asserting that a secret // did or did not cross the wire is one of the things this harness is for. From 10ea4da6781a694d63fa60998c43bbb8ed120595 Mon Sep 17 00:00:00 2001 From: Ahmustufa Date: Thu, 24 Sep 2026 16:07:49 +0500 Subject: [PATCH 3/4] feat(inbox): widening poll backoff on an unreachable server MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit An hour-long outage cost ~60 dial+auth attempts and ~20 dead-letter rows per mailbox, because a failed poll was simply retried and re-fanned-out every three minutes forever. That reconnect hammering is what provokes the per-IP sign-in throttles the fleet exists to avoid, and it buries real failures. The poller now classifies the failure and records it through a narrow, consumer-defined coreapi capability (PollBackoffCore) carried by BOTH clients, so a fleet host backs off exactly as a single-process install does: transport / unknown 3m -> 6 -> 12 -> 24 -> 48 -> 60m, capped auth / SSRF refusal straight to the 60m cap, and asynq.SkipRetry cancelled nothing recorded; our shutdown is not the server's fault our own failure nothing recorded; the dial legs are the only ones counted Auth is separated because a retry there is not free: it is another rejected sign-in, which is how a provider locks an account — the same reasoning authenticateIMAP already applies one level down. Transport evidence is read BEFORE the auth sentinel, so a connection that drops mid-LOGIN is not parked on the cap as though its password were wrong. The cap is the requirement, not the rungs: a backed-off mailbox is still polled, just hourly, and recovers on its own because the successful poll's cursor write clears the counter in the same statement. It is suppressed from the POLL fan-out and nothing else — status is never written, MailboxExists is untouched, and it still sends throughout. A failure that cannot be recorded fails OPEN, back to today's behaviour. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01VECELVAKe8Xcp7GH9wGR5t --- cmd/inroad/fleetlistener_test.go | 5 + internal/coreapi/inboxpollbackoff.go | 68 ++++ internal/coreapi/inprocess/inboxpoll.go | 50 +++ internal/coreapi/remote/handler.go | 2 + internal/coreapi/remote/inbound.go | 36 ++ internal/coreapi/remote/inbound_test.go | 45 ++- internal/coreapi/remote/inboundhandler.go | 36 ++ internal/coreapi/remote/wire.go | 40 ++- internal/worker/capabilities_test.go | 7 +- internal/worker/inbox/poll.go | 32 +- internal/worker/inbox/poll_test.go | 5 +- internal/worker/inbox/pollbackoff.go | 200 +++++++++++ .../inbox/pollbackoff_integration_test.go | 183 ++++++++++ internal/worker/inbox/pollbackoff_test.go | 336 ++++++++++++++++++ 14 files changed, 1033 insertions(+), 12 deletions(-) create mode 100644 internal/coreapi/inboxpollbackoff.go create mode 100644 internal/worker/inbox/pollbackoff.go create mode 100644 internal/worker/inbox/pollbackoff_integration_test.go create mode 100644 internal/worker/inbox/pollbackoff_test.go diff --git a/cmd/inroad/fleetlistener_test.go b/cmd/inroad/fleetlistener_test.go index 5e37d784..128d1235 100644 --- a/cmd/inroad/fleetlistener_test.go +++ b/cmd/inroad/fleetlistener_test.go @@ -101,6 +101,11 @@ func (f *fakeInboundWriter) SetInboxCursorString(context.Context, string, string return nil } +func (f *fakeInboundWriter) RecordInboxPollFailure(context.Context, string, string, []time.Duration) (coreapi.InboxPollBackoff, error) { + f.calls++ + return coreapi.InboxPollBackoff{}, nil +} + func (f *fakeInboundWriter) StoreInboundMessage(context.Context, coreapi.InboxMessageInput) error { f.calls++ return nil diff --git a/internal/coreapi/inboxpollbackoff.go b/internal/coreapi/inboxpollbackoff.go new file mode 100644 index 00000000..3bd2c736 --- /dev/null +++ b/internal/coreapi/inboxpollbackoff.go @@ -0,0 +1,68 @@ +package coreapi + +import ( + "errors" + "fmt" + "time" +) + +// The inbox poll backoff's vocabulary: what the execution plane tells the +// control plane when a mailbox's server would not answer, and what it gets back. +// +// It is NOT a coreapi.Client method. Like PendingSweepCore and ReplyCore, the +// capability is consumed through a narrow interface defined by its consumer +// (internal/worker/inbox.PollBackoffCore) and satisfied by both clients via type +// assertion at registration. Client already carries ~40 methods that a dozen +// test fakes implement in full, and every method on it is one a remote worker +// must be able to call — a bar this clears, which is why both implementations +// carry it, but not a reason to widen the seam itself. + +// ErrInvalidBackoffLadder rejects a schedule that could not produce a bounded +// retry time. It is a programming error rather than a runtime condition, but it +// crosses a process boundary on the remote transport, so it is validated there +// too: an empty ladder would leave inbox_poll_retry_after NULL and quietly +// disable the backoff, and a negative rung would schedule a retry in the past. +var ErrInvalidBackoffLadder = errors.New("coreapi: invalid inbox poll backoff ladder") + +// maxBackoffRungs bounds a ladder arriving over the wire. Nothing legitimate +// needs more than a handful of rungs, and an unbounded array from a compromised +// worker would be an unbounded array parameter in a SQL statement. +const maxBackoffRungs = 32 + +// InboxPollBackoff is the control plane's answer after recording one failed +// poll: how many consecutive failures this mailbox has now had, and the earliest +// time the poll fan-out will consider it again. +// +// RetryAfter is the DATABASE's clock, not the worker's — it is compared against +// the database clock in ListActiveMailboxes, and a worker's few-ms skew must not +// be what decides whether a mailbox is due. It is returned for the LOG LINE and +// for tests; nothing branches on it. +// +// snake_case json tags for the reason the package doc gives: the remote +// transport encodes THIS type rather than a parallel wire struct. +type InboxPollBackoff struct { + Failures int `json:"failures"` + RetryAfter time.Time `json:"retry_after"` +} + +// ValidateBackoffLadder checks a schedule before it reaches SQL. +// +// The ladder is indexed by consecutive-failure count and CLAMPED to its last +// rung, so the last rung is the cap — which is the whole reason a server down +// for an hour does not cost the user their connection. A one-rung ladder is +// therefore how a caller says "go straight to the cap"; that is a legitimate +// schedule, not a degenerate one. +func ValidateBackoffLadder(ladder []time.Duration) error { + if len(ladder) == 0 { + return fmt.Errorf("%w: empty", ErrInvalidBackoffLadder) + } + if len(ladder) > maxBackoffRungs { + return fmt.Errorf("%w: %d rungs exceeds the %d maximum", ErrInvalidBackoffLadder, len(ladder), maxBackoffRungs) + } + for i, d := range ladder { + if d <= 0 { + return fmt.Errorf("%w: rung %d is %v", ErrInvalidBackoffLadder, i+1, d) + } + } + return nil +} diff --git a/internal/coreapi/inprocess/inboxpoll.go b/internal/coreapi/inprocess/inboxpoll.go index 617112ff..e1e94422 100644 --- a/internal/coreapi/inprocess/inboxpoll.go +++ b/internal/coreapi/inprocess/inboxpoll.go @@ -128,6 +128,56 @@ func (c client) SetInboxCursorString(ctx context.Context, mailboxID, workspaceID }) } +// RecordInboxPollFailure records one failed poll against a mailbox and returns +// the schedule the database computed: the new consecutive-failure count and the +// earliest time the poll fan-out will consider it again. +// +// It writes ONLY the two backoff columns. It deliberately does not touch +// `status` — writing 'error' there would stop ListActiveMailboxes AND +// MailboxExists, so a mailbox whose server was down for an hour would stop +// SENDING too, which is exactly what "without deactivating the mailbox" +// forbids. Suppressing scheduling is the whole mechanism. +// +// ladder is the schedule, passed down as data and indexed by the new failure +// count with its last rung as the cap (see the query). It is validated here, at +// the seam, rather than trusted: an empty ladder would leave retry_after NULL +// and silently disable the backoff. +// +// workspaceID is pinned in the SQL WHERE, so a foreign id updates zero rows and +// surfaces as pgx.ErrNoRows rather than parking another tenant's poller. A +// mailbox deleted mid-poll takes the same path; the caller treats it as +// "nothing to record" rather than as a reason to fail the task. +func (c client) RecordInboxPollFailure(ctx context.Context, mailboxID, workspaceID string, ladder []time.Duration) (coreapi.InboxPollBackoff, error) { + if err := coreapi.ValidateBackoffLadder(ladder); err != nil { + return coreapi.InboxPollBackoff{}, err + } + id, err := uuid.Parse(mailboxID) + if err != nil { + return coreapi.InboxPollBackoff{}, err + } + ws, err := uuid.Parse(workspaceID) + if err != nil { + return coreapi.InboxPollBackoff{}, err + } + seconds := make([]float64, len(ladder)) + for i, d := range ladder { + seconds[i] = d.Seconds() + } + row, err := c.q.RecordInboxPollFailure(ctx, gen.RecordInboxPollFailureParams{ + BackoffSeconds: seconds, ID: id, WorkspaceID: ws, + }) + if err != nil { + if errors.Is(err, pgx.ErrNoRows) { + return coreapi.InboxPollBackoff{}, coreapi.ErrCrossTenant + } + return coreapi.InboxPollBackoff{}, err + } + return coreapi.InboxPollBackoff{ + Failures: int(row.InboxPollFailures), + RetryAfter: row.InboxPollRetryAfter.Time, + }, nil +} + // localFindSendByMessageID matches an inbound reply/bounce's Message-ID back to the // send that caused it, workspace-scoped. Returns ErrNoMatch when nothing // matches (unknown Message-ID — e.g. a reply to a message this workspace diff --git a/internal/coreapi/remote/handler.go b/internal/coreapi/remote/handler.go index 142258e3..e0c76d83 100644 --- a/internal/coreapi/remote/handler.go +++ b/internal/coreapi/remote/handler.go @@ -93,6 +93,7 @@ type JobReader interface { type InboundWriter interface { SetInboxCursor(ctx context.Context, mailboxID, workspaceID string, lastSeenUID, uidValidity uint32) error SetInboxCursorString(ctx context.Context, mailboxID, workspaceID, cursor string) error + RecordInboxPollFailure(ctx context.Context, mailboxID, workspaceID string, ladder []time.Duration) (coreapi.InboxPollBackoff, error) StoreInboundMessage(ctx context.Context, in coreapi.InboxMessageInput) error CaptureCRMReply(ctx context.Context, in coreapi.CRMReplyInput) error IngestComplaint(ctx context.Context, in coreapi.ComplaintInput) error @@ -311,6 +312,7 @@ func NewHandler(d Deps, token string, logger *slog.Logger) (http.Handler, error) // The inbound-mail routes and the last job read (slice 4). PathInboxCursorUID: h.setInboxCursorUID, PathInboxCursorString: h.setInboxCursorString, + PathInboxPollFailure: h.recordInboxPollFailure, PathInboxMessageStore: h.storeInboundMessage, PathCRMReplyCapture: h.captureCRMReply, PathComplaintIngest: h.ingestComplaint, diff --git a/internal/coreapi/remote/inbound.go b/internal/coreapi/remote/inbound.go index 9c06f872..5109e3ea 100644 --- a/internal/coreapi/remote/inbound.go +++ b/internal/coreapi/remote/inbound.go @@ -2,6 +2,7 @@ package remote import ( "context" + "time" "github.com/inroad/inroad/internal/coreapi" ) @@ -102,6 +103,41 @@ func (c *Client) SetInboxCursorString(ctx context.Context, mailboxID, workspaceI }, &ackResponse{}) } +// RecordInboxPollFailure records one failed poll and returns the schedule the +// control plane computed: the new consecutive-failure count and the earliest +// time the poll fan-out will consider this mailbox again. +// +// FAIL OPEN, which is the opposite direction from the cursor writes above, and +// deliberately so. If this call fails the mailbox simply keeps its old +// eligibility and is polled again on the next sweep — today's behaviour. The +// caller therefore treats an error here as something to log, never as a reason +// to fail the poll task: a backoff that could turn a control-plane hiccup into a +// mailbox that stops being polled would be the deactivation this whole change +// exists to avoid. +// +// The ladder is validated before it leaves, so a programming error surfaces at +// the caller rather than as a 400 from the control plane (which validates it +// again — see the handler). +func (c *Client) RecordInboxPollFailure(ctx context.Context, mailboxID, workspaceID string, ladder []time.Duration) (coreapi.InboxPollBackoff, error) { + if err := coreapi.ValidateBackoffLadder(ladder); err != nil { + return coreapi.InboxPollBackoff{}, err + } + if err := parseIDs(workspaceID, mailboxID); err != nil { + return coreapi.InboxPollBackoff{}, err + } + seconds := make([]float64, len(ladder)) + for i, d := range ladder { + seconds[i] = d.Seconds() + } + var out inboxPollFailureResponse + if err := c.post(ctx, c.outcomes, PathInboxPollFailure, inboxPollFailureRequest{ + WorkspaceID: workspaceID, MailboxID: mailboxID, BackoffSeconds: seconds, + }, &out); err != nil { + return coreapi.InboxPollBackoff{}, err + } + return coreapi.InboxPollBackoff{Failures: out.Failures, RetryAfter: out.RetryAfter}, nil +} + // StoreInboundMessage writes one matched inbound message onto its thread. // // The job travels WHOLE (coreapi.InboxMessageInput, not a subset) for the reason diff --git a/internal/coreapi/remote/inbound_test.go b/internal/coreapi/remote/inbound_test.go index 208ede87..7974fc8c 100644 --- a/internal/coreapi/remote/inbound_test.go +++ b/internal/coreapi/remote/inbound_test.go @@ -38,6 +38,7 @@ type fakeInbound struct { found bool send coreapi.WarmupSendRef matched bool + backoff coreapi.InboxPollBackoff } func (f *fakeInbound) record(args ...string) { f.calls++; f.gotArgs = args } @@ -53,6 +54,12 @@ func (f *fakeInbound) SetInboxCursorString(_ context.Context, mailboxID, workspa return f.err } +func (f *fakeInbound) RecordInboxPollFailure(_ context.Context, mailboxID, workspaceID string, ladder []time.Duration) (coreapi.InboxPollBackoff, error) { + f.record(mailboxID, workspaceID) + f.got = ladder + return f.backoff, f.err +} + func (f *fakeInbound) StoreInboundMessage(_ context.Context, in coreapi.InboxMessageInput) error { f.record(in.WorkspaceID, in.MailboxID) f.got = in @@ -164,6 +171,9 @@ func TestEverySlice4MethodRoundTrips(t *testing.T) { plan: coreapi.WarmupEngagePlan{ReceiptID: receipt, DoMarkRead: true, EngageAfter: 90 * time.Second}, label: coreapi.ReplyLabel{Key: "positive", StopsEnrollment: true}, send: coreapi.WarmupSendRef{WarmupSendID: send}, + backoff: coreapi.InboxPollBackoff{ + Failures: 3, RetryAfter: time.Now().Add(12 * time.Minute).UTC().Truncate(time.Millisecond), + }, } fl := &fakeFleet{queue: "w:box-1"} jobs := &fakeJobs{due: time.Now().Add(time.Hour).UTC().Truncate(time.Millisecond), sendNow: true} @@ -188,6 +198,39 @@ func TestEverySlice4MethodRoundTrips(t *testing.T) { } }) + t.Run("RecordInboxPollFailure", func(t *testing.T) { + ladder := []time.Duration{3 * time.Minute, 6 * time.Minute, time.Hour} + out, err := c.RecordInboxPollFailure(ctx, mailbox, ws, ladder) + if err != nil { + t.Fatalf("RecordInboxPollFailure: %v", err) + } + // The ladder crosses as SECONDS, so a transport that lost the unit would + // hand the control plane 180 nanoseconds and back a failing mailbox off + // by nothing at all. + got, ok := inbound.got.([]time.Duration) + if !ok { + t.Fatalf("handler got %T, want []time.Duration", inbound.got) + } + if len(got) != len(ladder) || got[0] != 3*time.Minute || got[2] != time.Hour { + t.Errorf("ladder = %v, want %v", got, ladder) + } + if out.Failures != inbound.backoff.Failures || !out.RetryAfter.Equal(inbound.backoff.RetryAfter) { + t.Errorf("backoff = %+v, want %+v", out, inbound.backoff) + } + }) + + t.Run("RecordInboxPollFailure refuses a ladder that could not bound a retry", func(t *testing.T) { + before := inbound.calls + for _, bad := range [][]time.Duration{nil, {}, {3 * time.Minute, 0}, {-time.Minute}} { + if _, err := c.RecordInboxPollFailure(ctx, mailbox, ws, bad); !errors.Is(err, coreapi.ErrInvalidBackoffLadder) { + t.Errorf("ladder %v: err = %v, want ErrInvalidBackoffLadder", bad, err) + } + } + if inbound.calls != before { + t.Error("an invalid ladder reached the control plane") + } + }) + t.Run("StoreInboundMessage", func(t *testing.T) { campaign := uuid.NewString() in := coreapi.InboxMessageInput{ @@ -735,7 +778,7 @@ func TestSlice4RoutesRequireTheFleetToken(t *testing.T) { func slice4Paths() []string { return []string{ - PathInboxCursorUID, PathInboxCursorString, PathInboxMessageStore, + PathInboxCursorUID, PathInboxCursorString, PathInboxPollFailure, PathInboxMessageStore, PathCRMReplyCapture, PathComplaintIngest, PathReplyLabelResolve, PathWarmupReceipt, PathWarmupSendByMessageID, PathWarmupTokenFailure, PathWarmupHardBounce, PathWarmupNextDue, diff --git a/internal/coreapi/remote/inboundhandler.go b/internal/coreapi/remote/inboundhandler.go index 06c0f682..87725d87 100644 --- a/internal/coreapi/remote/inboundhandler.go +++ b/internal/coreapi/remote/inboundhandler.go @@ -2,8 +2,11 @@ package remote import ( "net/http" + "time" "github.com/google/uuid" + + "github.com/inroad/inroad/internal/coreapi" ) // The control plane's side of the INBOUND MAIL path (slice 4). @@ -66,6 +69,39 @@ func (h *handler) setInboxCursorString(w http.ResponseWriter, r *http.Request) { respond(w, http.StatusOK, ackResponse{}) } +// recordInboxPollFailure records one failed poll and schedules the next attempt. +// +// The ladder is validated HERE as well as at the client, because this is a +// process boundary: the client's check protects a developer from a mistake, and +// this one protects the database from an unbounded or degenerate array arriving +// from a host the control plane does not compile. +func (h *handler) recordInboxPollFailure(w http.ResponseWriter, r *http.Request) { + var in inboxPollFailureRequest + if !decode(w, r, &in) { + return + } + ws, mailbox, ok := parsePair(w, in.WorkspaceID, "mailbox_id", in.MailboxID) + if !ok { + return + } + ladder := make([]time.Duration, len(in.BackoffSeconds)) + for i, s := range in.BackoffSeconds { + ladder[i] = time.Duration(s * float64(time.Second)) + } + if coreapi.ValidateBackoffLadder(ladder) != nil { + // The reason is not echoed: it is derived from caller input, and this + // route's replies stay as fixed as every other one here. + respond(w, http.StatusBadRequest, errorResponse{Error: "invalid backoff_seconds"}) + return + } + out, err := h.inbound.RecordInboxPollFailure(r.Context(), mailbox.String(), ws.String(), ladder) + if err != nil { + h.fail(w, failWrite, "inbox poll failure", err, "workspace_id", ws, "mailbox_id", mailbox) + return + } + respond(w, http.StatusOK, inboxPollFailureResponse{Failures: out.Failures, RetryAfter: out.RetryAfter}) +} + // storeInboundMessage writes one matched inbound message onto its thread. // // decodeMessage, not decode: the body is a customer's reply, capped by the same diff --git a/internal/coreapi/remote/wire.go b/internal/coreapi/remote/wire.go index b8c71ee0..e0caecdf 100644 --- a/internal/coreapi/remote/wire.go +++ b/internal/coreapi/remote/wire.go @@ -117,8 +117,13 @@ const ( // can only write the thread's history from what the poller hands it. The // same rule follows: no route here logs a body, a subject, a recipient or an // address, on either side. - PathInboxCursorUID = PathPrefix + "inbox-poll/cursor-uid" - PathInboxCursorString = PathPrefix + "inbox-poll/cursor-string" + PathInboxCursorUID = PathPrefix + "inbox-poll/cursor-uid" + PathInboxCursorString = PathPrefix + "inbox-poll/cursor-string" + // PathInboxPollFailure is the poll's OTHER outcome: the mailbox's server + // would not answer. It belongs with the cursor routes because it is the same + // decision — "this poll pass is finished, here is what the fan-out should do + // next" — and because the cursor writes are what CLEAR what it records. + PathInboxPollFailure = PathPrefix + "inbox-poll/failure" PathInboxMessageStore = PathPrefix + "inbox-message/store" PathCRMReplyCapture = PathPrefix + "crm-reply/capture" PathComplaintIngest = PathPrefix + "complaint/ingest" @@ -714,6 +719,37 @@ type inboxCursorStringRequest struct { Cursor string `json:"cursor"` } +// inboxPollFailureRequest records one failed poll and asks for the next attempt +// to be scheduled. +// +// BackoffSeconds carries the LADDER as data — rung N is the delay after the Nth +// consecutive failure, the last rung is the cap. The schedule is the worker's +// (internal/worker/inbox.DefaultPollBackoff) on both transports, so a fleet host +// and a single-process install back off identically; the control plane validates +// it (coreapi.ValidateBackoffLadder) rather than trusting it, because an empty +// ladder would leave retry_after NULL and silently disable the backoff. +// +// Seconds as float64 rather than a JSON-encoded time.Duration because this is a +// wire format read by whatever a future worker is written in, and "180" is a +// number every language agrees about while 180000000000 is a Go implementation +// detail. +// +// It carries NO error text. The failure's classification stays on the worker, +// which logs it; shipping a provider's error string here would put server +// output into a control-plane log line for nothing. +type inboxPollFailureRequest struct { + WorkspaceID string `json:"workspace_id"` + MailboxID string `json:"mailbox_id"` + BackoffSeconds []float64 `json:"backoff_seconds"` +} + +// inboxPollFailureResponse is what the control plane decided: the new +// consecutive-failure count and the DATABASE's retry time. +type inboxPollFailureResponse struct { + Failures int `json:"failures"` + RetryAfter time.Time `json:"retry_after"` +} + // inboxMessageRequest carries one inbound message onto its thread. // // It wraps coreapi.InboxMessageInput itself rather than mirroring its thirteen diff --git a/internal/worker/capabilities_test.go b/internal/worker/capabilities_test.go index e2ffcbdd..1755fedf 100644 --- a/internal/worker/capabilities_test.go +++ b/internal/worker/capabilities_test.go @@ -21,10 +21,11 @@ import ( // // # What this exists to catch // -// Eighteen capabilities are reached by COMMA-OK TYPE ASSERTION on whatever +// Nineteen capabilities are reached by COMMA-OK TYPE ASSERTION on whatever // value the composition root put in Deps.Core — six in handlers.go, three in // inbox.RegisterPerMessage, six feature-detected per message inside -// PollHandler, and three in cmd/worker. A type assertion is invisible to the +// PollHandler, one resolved once at PollHandler's wiring (the poll backoff), +// and three in cmd/worker. A type assertion is invisible to the // compiler. When one stops matching, nothing fails to build: the registrar // logs (or does not) and skips, the process starts, reports healthy, consumes // its queues, and silently does less than it did yesterday. @@ -124,6 +125,8 @@ func TestBothCoreAPIImplementationsCarryEveryOptionalCapability(t *testing.T) { assertsAs[coreapi.WarmupSendLookupClient](impl.core)}, {"coreapi.DeliverabilityComplaintClient", "inbox/poll.go, per message", "a mail-borne abuse report is never ingested and never suppresses", assertsAs[coreapi.DeliverabilityComplaintClient](impl.core)}, + {"inbox.PollBackoffCore", "inbox/poll.go, resolved once at wiring", "an unreachable mail server is re-dialed every sweep interval forever (inbox_poll_backoff_unavailable)", + assertsAs[inbox.PollBackoffCore](impl.core)}, {"coreapi.DeadLetterClient", "cmd/worker/main.go", "an exhausted task vanishes instead of reaching task_dead_letters", assertsAs[coreapi.DeadLetterClient](impl.core)}, {"coreapi.ProviderSignalClient", "cmd/worker/main.go", "per-worker provider verdicts are collected and never reported", diff --git a/internal/worker/inbox/poll.go b/internal/worker/inbox/poll.go index 8fb5f72a..9f040edf 100644 --- a/internal/worker/inbox/poll.go +++ b/internal/worker/inbox/poll.go @@ -222,6 +222,17 @@ func PollHandler(core coreapi.Client, reader mail.InboxReader, gmail GmailFetche "impact", "warmup token failures and warmup DSNs will not be recorded, "+ "and warmup DSNs may be misclassified as campaign bounces") } + // Resolved ONCE, for the same reason and with the same loud complaint: a + // core without it polls exactly as it did before the backoff existed, which + // means an unreachable server is re-dialed every three minutes forever. + backoff := pollBackoff{policy: DefaultPollBackoff} + if pbc, ok := core.(PollBackoffCore); ok { + backoff.core = pbc + } else { + slog.Error("inbox_poll_backoff_unavailable", + "impact", "a mailbox whose server is unreachable will be re-dialed every "+ + "sweep interval instead of backing off") + } return func(ctx context.Context, t *asynq.Task) error { var p queue.InboxPollPayload if err := json.Unmarshal(t.Payload(), &p); err != nil { @@ -249,18 +260,23 @@ func PollHandler(core coreapi.Client, reader mail.InboxReader, gmail GmailFetche } if job.Provider == "gmail" { - return pollAPI(ctx, core, gmail, classifier, hook, p, job, "gmail", gmailJunkScan(gmail)) + return pollAPI(ctx, core, backoff, gmail, classifier, hook, p, job, "gmail", gmailJunkScan(gmail)) } if job.Provider == "m365" { - return pollAPI(ctx, core, graph, classifier, hook, p, job, "m365", graphJunkScan(graph)) + return pollAPI(ctx, core, backoff, graph, classifier, hook, p, job, "m365", graphJunkScan(graph)) } defer zeroize(job.Password) cfg := mail.IMAPConfig{Host: job.Host, Port: job.Port, Username: job.Username, Password: string(job.Password)} + // The two calls below are the only ones in this branch that dial the + // mailbox's server, so they are the only ones the backoff counts. A + // failure anywhere AFTER them (a message that will not process, a cursor + // that will not persist) is ours, not the provider's, and backing off + // polling for it would delay reply detection over a bug of our own. uidValidity, uidNext, err := reader.CurrentState(ctx, cfg) if err != nil { - return err + return backoff.note(ctx, p, job.Provider, err) } // Re-baseline on a first poll (never-polled mailbox, UIDValidity==0) @@ -281,7 +297,7 @@ func PollHandler(core coreapi.Client, reader mail.InboxReader, gmail GmailFetche msgs, _, err := reader.Fetch(ctx, cfg, job.LastSeenUID, fetchBatchSize) if err != nil { - return err + return backoff.note(ctx, p, job.Provider, err) } var replies, bounces, skipped int @@ -354,12 +370,16 @@ type apiFetcher interface { // untouched). The short-lived access token is zeroized after the pass, like the // IMAP password. Only the transport (reader), the "provider" log value, and the // junk scanner differ between gmail and m365, so both providers share this body. -func pollAPI(ctx context.Context, core coreapi.Client, reader apiFetcher, classifier *replyclassify.Classifier, hook warmupHook, p queue.InboxPollPayload, job coreapi.InboxPollJob, provider string, junkScan apiJunkScan) error { +func pollAPI(ctx context.Context, core coreapi.Client, backoff pollBackoff, reader apiFetcher, classifier *replyclassify.Classifier, hook warmupHook, p queue.InboxPollPayload, job coreapi.InboxPollJob, provider string, junkScan apiJunkScan) error { defer zeroize(job.AccessToken) + // The one call here that reaches the provider, and so the one the backoff + // counts — see the IMAP branch's note. A 401 from Graph or Gmail is the API + // transport's refused credential and goes straight to the cap; a 429 or 5xx + // is "come back later", which is what the widening ladder is. msgs, newCursor, err := reader.Fetch(ctx, string(job.AccessToken), job.Cursor, fetchBatchSize) if err != nil { - return err + return backoff.note(ctx, p, provider, err) } var replies, bounces, skipped int diff --git a/internal/worker/inbox/poll_test.go b/internal/worker/inbox/poll_test.go index 7c8a1a99..eadf1171 100644 --- a/internal/worker/inbox/poll_test.go +++ b/internal/worker/inbox/poll_test.go @@ -85,6 +85,9 @@ type stubCore struct { cursorSet bool cursorUID uint32 cursorValidity uint32 + // cursorErr fails the cursor write — a CONTROL-PLANE failure, as opposed to + // a provider one, which the poll backoff must not count against the mailbox. + cursorErr error cursorStringSet bool cursorString string @@ -119,7 +122,7 @@ func (s *stubCore) GetInboxPollJob(context.Context, string, string) (coreapi.Inb func (s *stubCore) SetInboxCursor(_ context.Context, _, _ string, uid, validity uint32) error { s.cursorSet = true s.cursorUID, s.cursorValidity = uid, validity - return nil + return s.cursorErr } func (s *stubCore) SetInboxCursorString(_ context.Context, _, _, cursor string) error { diff --git a/internal/worker/inbox/pollbackoff.go b/internal/worker/inbox/pollbackoff.go new file mode 100644 index 00000000..584d5169 --- /dev/null +++ b/internal/worker/inbox/pollbackoff.go @@ -0,0 +1,200 @@ +package inbox + +import ( + "context" + "errors" + "fmt" + "log/slog" + "time" + + "github.com/hibiken/asynq" + + "github.com/inroad/inroad/internal/coreapi" + "github.com/inroad/inroad/internal/platform/mail" + "github.com/inroad/inroad/internal/platform/queue" +) + +// PollBackoffCore is the narrow coreapi capability the poll backoff needs: +// record one failed poll against one mailbox and be told when it may be tried +// again. Consumer-defined here rather than added to coreapi.Client — the same +// trade PendingSweepCore, ReplyCore and maintenance.Cleaner make, and for the +// same reason (that interface carries ~40 methods and a dozen test fakes +// implement it in full). +// +// ONE METHOD, AND IT TOUCHES TWO COLUMNS. That is the interface-segregation +// point and also the safety point: this seam is structurally incapable of +// writing mailboxes.status, so a backoff cannot become a deactivation however +// the calling code evolves. Writing status='error' would stop the mailbox +// SENDING as well as polling (MailboxExists gates on 'active'), which is the +// exact failure this whole change exists to prevent. +type PollBackoffCore interface { + RecordInboxPollFailure(ctx context.Context, mailboxID, workspaceID string, ladder []time.Duration) (coreapi.InboxPollBackoff, error) +} + +// PollBackoff is the schedule a failing mailbox is re-polled on. +// +// A struct rather than two loose slices for the reason StrandedPendingWindow is +// one: they are adjacent, interchangeable to the compiler, and swapping them +// changes behaviour silently instead of failing. +type PollBackoff struct { + // Ladder is indexed by consecutive-failure count and CLAMPED to its last + // rung, so the last rung is the cap. + Ladder []time.Duration + // Permanent is the single delay applied to a failure that cannot resolve on + // its own — a refused credential, an SSRF-blocked host. It is sent as a + // one-rung ladder, which is how the query is told "go straight to the cap". + Permanent time.Duration +} + +// DefaultPollBackoff is the schedule the poller ships with. +// +// # The rungs +// +// 3m → 6 → 12 → 24 → 48 → 60m, capped. The first rung is exactly +// inboxSweepInterval, so ONE failed poll changes nothing at all: a blip costs +// the mailbox no latency, and only a mailbox that keeps failing widens. From +// there each rung doubles, which takes an hour-long outage from ~20 fan-outs +// (60 dial+auth attempts, ~20 dead-letter rows) down to 5. +// +// # Why the cap matters more than the rungs +// +// "A server down for an hour must not cost the user their connection" is the +// requirement, and the cap is what delivers it: a backed-off mailbox is still +// polled, just hourly, so it RECOVERS ON ITS OWN the first time the server +// answers — the successful poll's cursor write clears the counter in the same +// statement. An uncapped exponential would reach a day and then a week, and a +// mailbox nobody thought to look at would stop detecting replies for good. That +// is a deactivation with extra steps. +// +// # Permanent +// +// One hour, the same ceiling, reached on the FIRST failure rather than the +// sixth. It is not longer than the cap: an operator who fixes a password should +// not wait a day, and pause→resume clears the counter for anyone who will not +// wait the hour (UpdateMailboxStatus). +var DefaultPollBackoff = PollBackoff{ + Ladder: []time.Duration{ + 3 * time.Minute, 6 * time.Minute, 12 * time.Minute, + 24 * time.Minute, 48 * time.Minute, 60 * time.Minute, + }, + Permanent: 60 * time.Minute, +} + +// pollFailureRecordTimeout bounds the failure write. It is deliberately short: +// the write is best-effort, and a control plane that is not answering promptly +// is not a reason to hold a poll slot open. +const pollFailureRecordTimeout = 5 * time.Second + +// rungsFor maps a classified failure onto the ladder to apply, or nil for +// "record nothing". +// +// FOUR SIGNALS, THREE ANSWERS: +// +// - A refused connection, a reset, a black-holed server, a failed TLS +// handshake and a provider 429/5xx are all TRANSPORT: the thing that +// routinely fixes itself, so it gets the widening ladder and recovers on its +// own. +// - A refused credential and an SSRF-blocked host go STRAIGHT TO THE CAP. They +// cannot fix themselves — no amount of retrying changes a wrong password — +// and an auth retry is not free the way a dial retry is: it is another +// rejected sign-in against the account, which is what makes a provider issue +// a challenge or lock the mailbox. This package already applies that +// reasoning one level down, where authenticateIMAP caps itself at two +// attempts per connection for the same reason; backing off to hourly is the +// same argument applied to the attempts themselves. +// - UNKNOWN takes the widening ladder too. The conservative choice is not +// "leave it alone": leaving it alone is precisely the behaviour that produced +// 60 attempts an hour, and an unrecognised failure that recurs is +// indistinguishable from an outage in every way that matters here. +// - A CANCELLED poll records nothing. It is our own shutdown, not the +// server's fault, and counting it would back off a mailbox nothing is wrong +// with every time a worker restarts. +func (b PollBackoff) rungsFor(kind mail.ConnectFailure) []time.Duration { + switch kind { + case mail.ConnectFailureTransport, mail.ConnectFailureUnknown: + return b.Ladder + case mail.ConnectFailureAuth, mail.ConnectFailurePolicy: + return []time.Duration{b.Permanent} + default: // ConnectFailureNone, ConnectFailureAborted + return nil + } +} + +// pollBackoff is the poller's failure recorder, resolved ONCE at wiring time. +// +// Resolved once, not per poll, for the reason PollHandler resolves +// WarmupEvidenceClient once: a comma-ok assertion evaluated per message means a +// core that does not implement it degrades in silence, and the whole point of +// this change is that silence is expensive. A nil core here is a poller that +// behaves exactly as it did before this existed. +type pollBackoff struct { + core PollBackoffCore + policy PollBackoff +} + +// note is called with the error from a provider dial or fetch. It classifies +// the failure, records it when it is the provider's fault, and returns the error +// the handler should return. +// +// The poll error is ALWAYS returned, never swallowed: asynq's retry and +// dead-letter capture are what make a genuine failure visible, and a backoff +// that also hid the failure would trade one invisible problem for another. What +// changes is the FAN-OUT rate, which is where the 60-attempts-an-hour came from. +// +// A permanent failure is additionally marked asynq.SkipRetry, so the two +// in-task retries do not turn one wrong password into three rejected sign-ins +// per rung. It is still captured as a dead letter — SkipRetry is terminal, and +// queue.IsTerminalFailure treats it as such on the first attempt. +func (b pollBackoff) note(ctx context.Context, p queue.InboxPollPayload, provider string, err error) error { + kind := mail.ClassifyConnectFailure(err) + ladder := b.policy.rungsFor(kind) + if ladder == nil { + return err + } + b.record(ctx, p, provider, kind, ladder) + if kind == mail.ConnectFailureAuth || kind == mail.ConnectFailurePolicy { + return fmt.Errorf("inbox poll: %w (%w)", err, asynq.SkipRetry) + } + return err +} + +// record writes the failure and logs what it decided. +// +// THE ERROR TEXT IS NOT LOGGED HERE and is not sent over the seam — only the +// classification. The handler's own error return is what carries the detail to +// asynq's log and to task_dead_letters; duplicating it here would put provider +// output, and on the auth path a server's response to a sign-in attempt, into a +// second place for no gain (docs/security.md, credential handling). +func (b pollBackoff) record(ctx context.Context, p queue.InboxPollPayload, provider string, kind mail.ConnectFailure, ladder []time.Duration) { + if b.core == nil { + return + } + // The record must OUTLIVE the failure that caused it. The most expensive + // failure this exists to stop — a server that accepts the connection and + // then says nothing — is exactly the one that can exhaust the poll's own + // deadline, and a record cancelled along with it would leave the mailbox on + // its old, tighter schedule and keep the hammering going. WithoutCancel + // keeps the trace values and drops only the cancellation; the deadline is + // then ours, chosen, not inherited. + recordCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), pollFailureRecordTimeout) + defer cancel() + + out, err := b.core.RecordInboxPollFailure(recordCtx, p.MailboxID, p.WorkspaceID, ladder) + if err != nil { + // FAIL OPEN. The mailbox keeps its old eligibility and is polled again on + // the next sweep — today's behaviour. A recording failure must never be + // able to widen a backoff or, worse, stop one being cleared. + if errors.Is(err, coreapi.ErrCrossTenant) { + // The mailbox was deleted (or never belonged to this workspace) + // between the fan-out and now. There is nothing to back off. + slog.Info("inbox_poll_backoff_no_mailbox", "mailbox_id", p.MailboxID, "workspace_id", p.WorkspaceID) + return + } + slog.Warn("inbox_poll_backoff_not_recorded", "mailbox_id", p.MailboxID, + "workspace_id", p.WorkspaceID, "reason", kind.String(), "error", err) + return + } + slog.Warn("inbox_poll_failed", "mailbox_id", p.MailboxID, "workspace_id", p.WorkspaceID, + "provider", provider, "reason", kind.String(), + "consecutive_failures", out.Failures, "retry_after", out.RetryAfter) +} diff --git a/internal/worker/inbox/pollbackoff_integration_test.go b/internal/worker/inbox/pollbackoff_integration_test.go new file mode 100644 index 00000000..5b412dbb --- /dev/null +++ b/internal/worker/inbox/pollbackoff_integration_test.go @@ -0,0 +1,183 @@ +//go:build integration + +package inbox + +import ( + "context" + "errors" + "fmt" + "syscall" + "testing" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/inroad/inroad/internal/platform/db/gen" + "github.com/inroad/inroad/internal/platform/mail" + "github.com/inroad/inroad/internal/platform/replyclassify" +) + +// The end-to-end proof, over the REAL inprocess coreapi and a real Postgres: +// an unreachable server backs the mailbox off the poll fan-out, the mailbox +// stays sendable throughout, and it comes back on its own. +// +// The unit tests above prove the classification and the ladder; this proves the +// three pieces are actually wired to each other — the handler to the coreapi +// capability, the capability to the columns, the columns to the fan-out query. +// Each of those seams is a type assertion or a generated query, and the unit +// tests are blind to all three. +func TestAnUnreachableServerBacksOffWithoutDeactivatingTheMailbox(t *testing.T) { + ctx := context.Background() + pool, q, closeFn := connect(t) + defer closeFn() + fx := seedActiveEnrollment(t, ctx, pool, q, newSealer(t), "") + + fannedOut := func() bool { + t.Helper() + rows, err := fx.core.ListActiveMailboxes(ctx) + if err != nil { + t.Fatalf("ListActiveMailboxes: %v", err) + } + for _, r := range rows { + if r.ID == fx.mailboxID.String() { + return true + } + } + return false + } + poll := func(reader *fakeReader) error { + t.Helper() + return PollHandler(fx.core, reader, nil, nil, replyclassify.New(nil), nil, noopEngageEnqueuer{})( + ctx, pollTaskFor(t, fx.mailboxID.String(), fx.ws.String())) + } + + if !fannedOut() { + t.Fatal("a healthy mailbox is not in the poll fan-out") + } + + // The server goes away. Three sweeps' worth of polls; the ladder widens under + // each one. + unreachable := &fakeReader{stateErr: fmt.Errorf("imap dial: %w", syscall.ECONNREFUSED)} + var previous int + for attempt := 1; attempt <= 3; attempt++ { + if err := poll(unreachable); !errors.Is(err, syscall.ECONNREFUSED) { + t.Fatalf("attempt %d: poll error = %v, want the dial failure", attempt, err) + } + failures, retryAfter := mailboxBackoff(t, ctx, pool, fx.mailboxID) + if failures != attempt { + t.Fatalf("attempt %d recorded failure count %d", attempt, failures) + } + if retryAfter <= previous { + t.Errorf("attempt %d scheduled the retry %ds out, no further than attempt %d's %ds", + attempt, retryAfter, attempt-1, previous) + } + previous = retryAfter + if fannedOut() { + t.Errorf("attempt %d left the mailbox in the fan-out; the hammering is unchanged", attempt) + } + } + + // THE ASSERTION THIS TASK EXISTS FOR. The mailbox is suppressed from POLLING + // and from nothing else — it is still 'active', so it still SENDS. + if got := mailboxStatus(t, ctx, pool, fx.mailboxID); got != "active" { + t.Fatalf("status = %q after three failed polls; the backoff became a deactivation", got) + } + if ok, err := q.MailboxExists(ctx, fx.mailboxID); err != nil || !ok { + t.Fatalf("MailboxExists = (%v, %v); a backed-off mailbox must still send", ok, err) + } + + // And it is polled again once due, with no operator action — the cap is what + // guarantees this happens at all. + if _, err := pool.Exec(ctx, + `UPDATE mailboxes SET inbox_poll_retry_after = now() - interval '1 second' WHERE id = $1`, + fx.mailboxID); err != nil { + t.Fatalf("expire the backoff: %v", err) + } + if !fannedOut() { + t.Fatal("a mailbox past its retry time is still suppressed") + } + + // The server comes back. One successful poll clears everything, in the same + // statement that advances the cursor. + if err := poll(&fakeReader{uidValidity: 7, uidNext: 30}); err != nil { + t.Fatalf("recovery poll: %v", err) + } + if failures, retryAfter := mailboxBackoff(t, ctx, pool, fx.mailboxID); failures != 0 || retryAfter != 0 { + t.Errorf("after a successful poll the backoff is (failures=%d, retry_after=+%ds), want cleared", failures, retryAfter) + } + if !fannedOut() { + t.Error("a recovered mailbox is still suppressed") + } +} + +// A refused credential is a different signal: it goes straight to the cap on the +// first failure, because retrying it is another rejected sign-in — and it still +// must not touch the mailbox's status. +func TestARejectedCredentialCapsImmediatelyAndStillSends(t *testing.T) { + ctx := context.Background() + pool, q, closeFn := connect(t) + defer closeFn() + fx := seedActiveEnrollment(t, ctx, pool, q, newSealer(t), "") + + reader := &fakeReader{stateErr: fmt.Errorf("%w: imap login: bad credentials", mail.ErrAuthRejected)} + if err := PollHandler(fx.core, reader, nil, nil, replyclassify.New(nil), nil, noopEngageEnqueuer{})( + ctx, pollTaskFor(t, fx.mailboxID.String(), fx.ws.String())); err == nil { + t.Fatal("a rejected credential did not fail the poll") + } + + failures, retryAfter := mailboxBackoff(t, ctx, pool, fx.mailboxID) + if failures != 1 { + t.Fatalf("failure count = %d, want 1", failures) + } + // The cap, on the FIRST failure — not the 3-minute first rung. + if want := int(DefaultPollBackoff.Permanent.Seconds()); retryAfter < want-60 || retryAfter > want+60 { + t.Errorf("retry scheduled +%ds out, want ~%ds (the cap, immediately)", retryAfter, want) + } + if got := mailboxStatus(t, ctx, pool, fx.mailboxID); got != "active" { + t.Errorf("status = %q; a wrong password must not deactivate a mailbox", got) + } + if ok, err := q.MailboxExists(ctx, fx.mailboxID); err != nil || !ok { + t.Errorf("MailboxExists = (%v, %v); a mailbox with a stale IMAP password still SENDS over SMTP", ok, err) + } + + // The operator's reset: pause→resume clears the counter, so nobody has to + // wait out the hour after fixing the password. + for _, status := range []string{"paused", "active"} { + if _, err := q.UpdateMailboxStatus(ctx, gen.UpdateMailboxStatusParams{ + ID: fx.mailboxID, WorkspaceID: fx.ws, Status: status, + }); err != nil { + t.Fatalf("%s: %v", status, err) + } + } + if failures, retryAfter := mailboxBackoff(t, ctx, pool, fx.mailboxID); failures != 0 || retryAfter != 0 { + t.Errorf("after resume the backoff is (failures=%d, retry_after=+%ds), want cleared", failures, retryAfter) + } +} + +// mailboxBackoff reads the two columns, measuring retry_after as SECONDS FROM +// THE DATABASE'S OWN now(). Comparing against a time this process built would +// make a few ms of host↔DB clock skew decide the assertion, and this test is all +// time arithmetic. 0 means NULL (no backoff). +func mailboxBackoff(t *testing.T, ctx context.Context, pool *pgxpool.Pool, id uuid.UUID) (int, int) { + t.Helper() + var failures int + var secs *float64 + if err := pool.QueryRow(ctx, + `SELECT inbox_poll_failures, extract(epoch FROM inbox_poll_retry_after - now()) + FROM mailboxes WHERE id = $1`, id).Scan(&failures, &secs); err != nil { + t.Fatalf("backoff state: %v", err) + } + if secs == nil { + return failures, 0 + } + return failures, int(*secs) +} + +func mailboxStatus(t *testing.T, ctx context.Context, pool *pgxpool.Pool, id uuid.UUID) string { + t.Helper() + var s string + if err := pool.QueryRow(ctx, `SELECT status FROM mailboxes WHERE id = $1`, id).Scan(&s); err != nil { + t.Fatalf("status: %v", err) + } + return s +} diff --git a/internal/worker/inbox/pollbackoff_test.go b/internal/worker/inbox/pollbackoff_test.go new file mode 100644 index 00000000..22544cd4 --- /dev/null +++ b/internal/worker/inbox/pollbackoff_test.go @@ -0,0 +1,336 @@ +package inbox + +import ( + "context" + "errors" + "fmt" + "io" + "syscall" + "testing" + "time" + + "github.com/hibiken/asynq" + + "github.com/inroad/inroad/internal/coreapi" + "github.com/inroad/inroad/internal/platform/mail" + "github.com/inroad/inroad/internal/platform/queue" + "github.com/inroad/inroad/internal/platform/replyclassify" +) + +// recordedFailure is one RecordInboxPollFailure the poller made. +type recordedFailure struct { + mailboxID string + workspaceID string + ladder []time.Duration + // ctxLive is whether the context the call arrived on was still usable. It is + // captured rather than asserted here because the whole point of one test + // below is that a poll killed by its own deadline can still record. + ctxLive bool +} + +// backoffCore is a stubCore that ALSO carries the optional backoff capability, +// which plain stubCore deliberately does not — so every pre-existing poll test +// keeps exercising the no-capability path. +type backoffCore struct { + *stubCore + recorded []recordedFailure + out coreapi.InboxPollBackoff + err error +} + +func (b *backoffCore) RecordInboxPollFailure(ctx context.Context, mailboxID, workspaceID string, ladder []time.Duration) (coreapi.InboxPollBackoff, error) { + b.recorded = append(b.recorded, recordedFailure{ + mailboxID: mailboxID, workspaceID: workspaceID, + ladder: append([]time.Duration(nil), ladder...), ctxLive: ctx.Err() == nil, + }) + return b.out, b.err +} + +func newBackoffCore(job coreapi.InboxPollJob) *backoffCore { + return &backoffCore{stubCore: &stubCore{job: job}} +} + +// a mailbox that has been polled before, so the poll takes the fetch path rather +// than the first-poll baseline. +func polledBefore() coreapi.InboxPollJob { + return coreapi.InboxPollJob{UIDValidity: 7, LastSeenUID: 10} +} + +// A transport failure is the case this whole change exists for: the server is +// unreachable, and the mailbox must be scheduled further out instead of re-dialed +// every sweep. +func TestTransportFailureRecordsTheWideningLadder(t *testing.T) { + core := newBackoffCore(polledBefore()) + reader := &fakeReader{stateErr: fmt.Errorf("imap dial: %w", syscall.ECONNREFUSED)} + + err := runPoll(t, core, reader) + if !errors.Is(err, syscall.ECONNREFUSED) { + t.Fatalf("poll error = %v, want the dial failure propagated", err) + } + // Still retried in-task: a refused connection may well succeed on the next + // attempt, and only the FAN-OUT rate is what was hammering. + if errors.Is(err, asynq.SkipRetry) { + t.Error("a transport failure was marked SkipRetry; a refused dial is worth retrying") + } + if len(core.recorded) != 1 { + t.Fatalf("recorded %d failures, want 1", len(core.recorded)) + } + got := core.recorded[0] + if !sameLadder(got.ladder, DefaultPollBackoff.Ladder) { + t.Errorf("ladder = %v, want the widening ladder %v", got.ladder, DefaultPollBackoff.Ladder) + } +} + +// The Fetch leg dials too, so a failure there counts the same as one at +// CurrentState. +func TestTransportFailureAtFetchAlsoRecords(t *testing.T) { + core := newBackoffCore(polledBefore()) + reader := &fakeReader{uidValidity: 7, uidNext: 20, fetchErr: fmt.Errorf("imap fetch: %w", io.EOF)} + + if err := runPoll(t, core, reader); !errors.Is(err, io.EOF) { + t.Fatalf("poll error = %v, want the fetch failure propagated", err) + } + if len(core.recorded) != 1 { + t.Fatalf("recorded %d failures, want 1", len(core.recorded)) + } +} + +// A refused credential goes STRAIGHT to the cap and is not retried in-task: +// every repeat is another rejected sign-in, which is what makes a provider lock +// an account. +func TestAuthFailureGoesStraightToTheCapAndSkipsRetry(t *testing.T) { + core := newBackoffCore(polledBefore()) + reader := &fakeReader{stateErr: fmt.Errorf("%w: imap login: bad credentials", mail.ErrAuthRejected)} + + err := runPoll(t, core, reader) + if !errors.Is(err, mail.ErrAuthRejected) { + t.Fatalf("poll error = %v, want the auth failure propagated (never swallowed)", err) + } + if !errors.Is(err, asynq.SkipRetry) { + t.Error("an auth failure was left retryable; two more rejected sign-ins per rung is the cost") + } + if len(core.recorded) != 1 { + t.Fatalf("recorded %d failures, want 1", len(core.recorded)) + } + ladder := core.recorded[0].ladder + if len(ladder) != 1 || ladder[0] != DefaultPollBackoff.Permanent { + t.Errorf("ladder = %v, want the one-rung cap [%v]", ladder, DefaultPollBackoff.Permanent) + } +} + +// Our own SSRF refusal never opened a socket and will never succeed until an +// operator edits the mailbox — same one-rung treatment, for a different reason. +func TestPolicyRefusalGoesStraightToTheCap(t *testing.T) { + core := newBackoffCore(polledBefore()) + reader := &fakeReader{stateErr: fmt.Errorf("imap: %w", mail.ErrHostNotPermitted)} + + if err := runPoll(t, core, reader); !errors.Is(err, mail.ErrHostNotPermitted) { + t.Fatalf("poll error = %v, want the guard's refusal propagated", err) + } + if len(core.recorded) != 1 || len(core.recorded[0].ladder) != 1 { + t.Fatalf("recorded = %+v, want one failure on a one-rung ladder", core.recorded) + } +} + +// A poll that reaches the server records nothing — the cursor write is what +// clears the counter, and it must not be preceded by a failure the poll did not +// have. +func TestSuccessfulPollRecordsNoFailure(t *testing.T) { + core := newBackoffCore(polledBefore()) + reader := &fakeReader{uidValidity: 7, uidNext: 20} + + if err := runPoll(t, core, reader); err != nil { + t.Fatalf("poll: %v", err) + } + if len(core.recorded) != 0 { + t.Errorf("a successful poll recorded %d failures", len(core.recorded)) + } + if !core.cursorSet { + t.Error("a successful poll did not advance the cursor, which is what clears the backoff") + } +} + +// A failure that is OURS, not the provider's, must not back off polling: the +// mailbox is fine, and delaying reply detection over our own bug helps nobody. +func TestNonProviderFailureRecordsNoBackoff(t *testing.T) { + core := newBackoffCore(polledBefore()) + core.cursorErr = errors.New("control plane unavailable") + reader := &fakeReader{uidValidity: 7, uidNext: 20} + + if err := runPoll(t, core, reader); err == nil { + t.Fatal("a failed cursor write did not fail the poll") + } + if len(core.recorded) != 0 { + t.Errorf("a control-plane failure recorded %d poll failures against the mailbox", len(core.recorded)) + } +} + +// A core without the capability polls exactly as it did before this existed: no +// panic, no swallowed error, no backoff. stubCore is that core. +func TestPollWithoutTheBackoffCapabilityIsUnchanged(t *testing.T) { + core := &stubCore{job: polledBefore()} + reader := &fakeReader{stateErr: fmt.Errorf("imap dial: %w", syscall.ECONNREFUSED)} + + if err := runPoll(t, core, reader); !errors.Is(err, syscall.ECONNREFUSED) { + t.Fatalf("poll error = %v, want the dial failure propagated", err) + } +} + +// A control plane that cannot record the failure must not change the poll's +// outcome. Failing open here means the mailbox keeps its old eligibility — the +// behaviour that shipped before the backoff — rather than a hiccup deciding a +// mailbox stops being polled. +func TestABackoffThatCannotBeRecordedStillFailsThePollAndNothingElse(t *testing.T) { + core := newBackoffCore(polledBefore()) + core.err = errors.New("control plane unavailable") + reader := &fakeReader{stateErr: fmt.Errorf("imap dial: %w", syscall.ECONNREFUSED)} + + if err := runPoll(t, core, reader); !errors.Is(err, syscall.ECONNREFUSED) { + t.Fatalf("poll error = %v, want the dial failure, not the recording failure", err) + } +} + +// A deleted mailbox (or a workspace that never owned it) surfaces as +// ErrCrossTenant from the workspace-pinned UPDATE. There is nothing to back off, +// and it must not change the poll's error either. +func TestABackoffForAMailboxThatIsGoneIsNotAnError(t *testing.T) { + core := newBackoffCore(polledBefore()) + core.err = coreapi.ErrCrossTenant + reader := &fakeReader{stateErr: fmt.Errorf("imap dial: %w", syscall.ECONNREFUSED)} + + if err := runPoll(t, core, reader); !errors.Is(err, syscall.ECONNREFUSED) { + t.Fatalf("poll error = %v, want the dial failure", err) + } +} + +// The record must OUTLIVE the failure that caused it. A server that accepts the +// connection and then stalls is exactly the failure that can burn the poll's own +// deadline, and a record cancelled along with it would leave the mailbox on its +// old, tighter schedule — the hammering this exists to stop, at its most +// expensive. +func TestTheFailureRecordSurvivesTheCancelledPoll(t *testing.T) { + core := newBackoffCore(polledBefore()) + b := pollBackoff{core: core, policy: DefaultPollBackoff} + + ctx, cancel := context.WithCancel(t.Context()) + cancel() + + p := queue.InboxPollPayload{MailboxID: "m1", WorkspaceID: "ws1"} + stalled := fmt.Errorf("imap dial: %w", context.DeadlineExceeded) + if err := b.note(ctx, p, "smtp", stalled); !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("note returned %v, want the stall propagated", err) + } + if len(core.recorded) != 1 { + t.Fatalf("recorded %d failures, want 1", len(core.recorded)) + } + if !core.recorded[0].ctxLive { + t.Error("the failure record arrived on a cancelled context; it would never have been written") + } + if core.recorded[0].mailboxID != "m1" || core.recorded[0].workspaceID != "ws1" { + t.Errorf("recorded against %+v, want (m1, ws1)", core.recorded[0]) + } +} + +// The classification-to-schedule map, stated once as a table. Four signals, +// three answers. +func TestRungsForEveryClassification(t *testing.T) { + for _, tc := range []struct { + kind mail.ConnectFailure + want []time.Duration + why string + }{ + {mail.ConnectFailureTransport, DefaultPollBackoff.Ladder, "an unreachable server routinely fixes itself"}, + {mail.ConnectFailureUnknown, DefaultPollBackoff.Ladder, "an unrecognised failure that recurs is an outage in every way that matters"}, + {mail.ConnectFailureAuth, []time.Duration{DefaultPollBackoff.Permanent}, "a rejected sign-in cannot fix itself and is not free to repeat"}, + {mail.ConnectFailurePolicy, []time.Duration{DefaultPollBackoff.Permanent}, "our own refusal will not change without an operator"}, + {mail.ConnectFailureAborted, nil, "our shutdown is not the server's fault"}, + {mail.ConnectFailureNone, nil, "there is no failure"}, + } { + t.Run(tc.kind.String(), func(t *testing.T) { + if got := DefaultPollBackoff.rungsFor(tc.kind); !sameLadder(got, tc.want) { + t.Errorf("rungsFor(%v) = %v, want %v — %s", tc.kind, got, tc.want, tc.why) + } + }) + } +} + +// The schedule's own properties, which the SQL and the integration test both +// depend on: it widens, it is capped, and the first rung costs a healthy mailbox +// nothing. +func TestDefaultPollBackoffLadder(t *testing.T) { + l := DefaultPollBackoff.Ladder + if len(l) == 0 { + t.Fatal("an empty ladder disables the backoff entirely") + } + // inboxSweepInterval, which is unexported in platform/queue. One failed poll + // must therefore change nothing at all: a blip costs no latency. + if l[0] != 3*time.Minute { + t.Errorf("first rung = %v, want the 3m sweep interval so one blip costs nothing", l[0]) + } + for i := 1; i < len(l); i++ { + if l[i] <= l[i-1] { + t.Errorf("rung %d (%v) does not widen on rung %d (%v)", i+1, l[i], i, l[i-1]) + } + } + ceiling := l[len(l)-1] + if ceiling > time.Hour { + t.Errorf("cap = %v; a server down for an hour must not cost the user their connection", ceiling) + } + if DefaultPollBackoff.Permanent != ceiling { + t.Errorf("permanent rung = %v, cap = %v; a fixable password must not wait longer than an outage does", + DefaultPollBackoff.Permanent, ceiling) + } + if err := coreapi.ValidateBackoffLadder(l); err != nil { + t.Errorf("the shipped ladder does not pass the seam's own validation: %v", err) + } +} + +// The API transports reach the same machinery through a different error shape. +func TestAPIPollFailureBacksOff(t *testing.T) { + for _, tc := range []struct { + name string + err error + wantCap bool + skipsNow bool + }{ + {"graph 401 is the API transport's wrong password", &mail.APIError{Provider: "m365", Op: "inbox delta", Status: 401}, true, true}, + {"graph 503 is come back later", &mail.APIError{Provider: "m365", Op: "inbox delta", Status: 503}, false, false}, + } { + t.Run(tc.name, func(t *testing.T) { + core := &backoffCore{stubCore: &stubCore{job: coreapi.InboxPollJob{Provider: "m365"}}} + err := PollHandler(core, nil, nil, failingGraph{tc.err}, replyclassify.New(nil), nil, noopEngageEnqueuer{})(t.Context(), pollTask(t)) + if !errors.Is(err, tc.err) { + t.Fatalf("poll error = %v, want the provider failure propagated", err) + } + if errors.Is(err, asynq.SkipRetry) != tc.skipsNow { + t.Errorf("SkipRetry = %v, want %v", errors.Is(err, asynq.SkipRetry), tc.skipsNow) + } + if len(core.recorded) != 1 { + t.Fatalf("recorded %d failures, want 1", len(core.recorded)) + } + atCap := len(core.recorded[0].ladder) == 1 + if atCap != tc.wantCap { + t.Errorf("ladder = %v, wantCap = %v", core.recorded[0].ladder, tc.wantCap) + } + }) + } +} + +// failingGraph is a GraphFetcher whose every fetch is the provider's answer. +type failingGraph struct{ err error } + +func (f failingGraph) Fetch(context.Context, string, string, int) ([]mail.InboundMessage, string, error) { + return nil, "", f.err +} + +func sameLadder(a, b []time.Duration) bool { + if len(a) != len(b) { + return false + } + for i := range a { + if a[i] != b[i] { + return false + } + } + return true +} From 74d2efebb6840164ae4e2d5aafe22baa63bf32a3 Mon Sep 17 00:00:00 2001 From: Ahmustufa Date: Thu, 24 Sep 2026 16:18:57 +0500 Subject: [PATCH 4/4] fix(mail): carry Graph inbox statuses as data so the poller can act on them MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The delta, junk-list and $value legs reported every non-2xx as fmt.Errorf("graph: ...: unexpected status %d"). That is the exact failure *APIError was introduced for on the SEND path: a status recoverable only by parsing English. It now matters here too. The poll backoff wants opposite handling for a Graph 401 (a revoked token — straight to the hour cap, because retrying is another rejected sign-in) and a 429/503 ("come back later" — the widening ladder), and through a sentence both classified as unknown and took the gentle ladder. The test drives graphDelta against a real httptest server rather than through Fetch, because Fetch host-pins the cursor to graph.microsoft.com before dialing it (security invariant 13) and must keep being unable to reach a test host. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01VECELVAKe8Xcp7GH9wGR5t --- internal/platform/mail/graphinbox.go | 12 +++++-- internal/platform/mail/graphinbox_test.go | 38 +++++++++++++++++++++++ 2 files changed, 47 insertions(+), 3 deletions(-) diff --git a/internal/platform/mail/graphinbox.go b/internal/platform/mail/graphinbox.go index 6c4e7d9f..9437e5ee 100644 --- a/internal/platform/mail/graphinbox.go +++ b/internal/platform/mail/graphinbox.go @@ -284,7 +284,7 @@ func graphJunkList(ctx context.Context, hc *http.Client, accessToken string, max } defer resp.Body.Close() if resp.StatusCode < 200 || resp.StatusCode >= 300 { - return nil, fmt.Errorf("graph: junk list: unexpected status %d", resp.StatusCode) + return nil, &APIError{Provider: "m365", Op: "junk list", Status: resp.StatusCode} } var body struct { Value []struct { @@ -345,8 +345,14 @@ func graphDelta(ctx context.Context, hc *http.Client, accessToken, u string) ([] return nil, "", "", errGraphDeltaExpired case resp.StatusCode == http.StatusBadRequest && graphResyncRequired(resp.Body): return nil, "", "", errGraphDeltaExpired + // *APIError rather than fmt.Errorf("unexpected status %d"), for the reason + // APIError's own doc gives about the SEND path: a status recoverable only by + // parsing English is a status nothing downstream can act on. The poller does + // act on it now — 401 is a revoked token and backs off straight to the cap, + // 429/5xx is "come back later" and takes the widening ladder — and it could + // not tell them apart through a sentence. case resp.StatusCode < 200 || resp.StatusCode >= 300: - return nil, "", "", fmt.Errorf("graph: delta: unexpected status %d", resp.StatusCode) + return nil, "", "", &APIError{Provider: "m365", Op: "inbox delta", Status: resp.StatusCode} } var body struct { Value []deltaMsg `json:"value"` @@ -406,7 +412,7 @@ func graphGetRaw(ctx context.Context, hc *http.Client, accessToken, id string) ( } defer resp.Body.Close() if resp.StatusCode < 200 || resp.StatusCode >= 300 { - return nil, fmt.Errorf("graph: raw: unexpected status %d", resp.StatusCode) + return nil, &APIError{Provider: "m365", Op: "message raw", Status: resp.StatusCode} } raw, err := io.ReadAll(resp.Body) if err != nil { diff --git a/internal/platform/mail/graphinbox_test.go b/internal/platform/mail/graphinbox_test.go index a4511c88..6b51cc66 100644 --- a/internal/platform/mail/graphinbox_test.go +++ b/internal/platform/mail/graphinbox_test.go @@ -3,6 +3,9 @@ package mail import ( "context" "encoding/json" + "errors" + "net/http" + "net/http/httptest" "strings" "testing" ) @@ -337,3 +340,38 @@ func TestGraphResyncRequired(t *testing.T) { }) } } + +// A Graph status the poller must ACT on has to arrive as DATA, not as a +// sentence. 401 (a revoked token) and 503 (come back later) want opposite +// handling from the inbox poll backoff — straight to the cap versus the +// widening ladder — and the only thing that can tell them apart is the status +// carried on *APIError. Until this, the delta path reported both as +// fmt.Errorf("unexpected status %d"), which classified as unknown. +// +// graphDelta is called directly rather than through Fetch because Fetch +// host-pins the cursor to graph.microsoft.com before dialing it (docs/security.md +// invariant 13), which an httptest server cannot satisfy — and must not be able to. +func TestGraphDeltaStatusIsCarriedAsDataNotProse(t *testing.T) { + for _, status := range []int{401, 403, 429, 500, 503} { + t.Run(http.StatusText(status), func(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(status) + })) + defer srv.Close() + + _, _, _, err := graphDelta(t.Context(), srv.Client(), "s3cret-token", srv.URL) + var apiErr *APIError + if !errors.As(err, &apiErr) { + t.Fatalf("graphDelta returned %v (%T), want an *APIError", err, err) + } + if apiErr.Status != status { + t.Errorf("status = %d, want %d", apiErr.Status, status) + } + // The error text must not echo the bearer or the provider's body: + // Graph puts request content in its error bodies (see the APIError doc). + if strings.Contains(apiErr.Error(), "s3cret-token") { + t.Errorf("the error carries the bearer token: %q", apiErr.Error()) + } + }) + } +}