Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions cmd/inroad/fleetlistener_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
68 changes: 68 additions & 0 deletions internal/coreapi/inboxpollbackoff.go
Original file line number Diff line number Diff line change
@@ -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
}
50 changes: 50 additions & 0 deletions internal/coreapi/inprocess/inboxpoll.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions internal/coreapi/remote/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down
36 changes: 36 additions & 0 deletions internal/coreapi/remote/inbound.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package remote

import (
"context"
"time"

"github.com/inroad/inroad/internal/coreapi"
)
Expand Down Expand Up @@ -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
Expand Down
45 changes: 44 additions & 1 deletion internal/coreapi/remote/inbound_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand All @@ -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
Expand Down Expand Up @@ -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}
Expand All @@ -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{
Expand Down Expand Up @@ -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,
Expand Down
36 changes: 36 additions & 0 deletions internal/coreapi/remote/inboundhandler.go
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down Expand Up @@ -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
Expand Down
40 changes: 38 additions & 2 deletions internal/coreapi/remote/wire.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading