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/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; 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/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()) + } + }) + } +} 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. 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 +}