Skip to content
Merged
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
17 changes: 11 additions & 6 deletions plugins/retention/consent_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,12 @@ import (
)

// policy builds a fixed consentPolicy for tests.
func policy(require bool, purpose string) func(context.Context, id.AppID) (bool, string) {
// policy builds a consent policy for the tests. The purpose is fixed at
// "marketing" because that is the only purpose these tests exercise; take a
// parameter again if a test needs a second one.
func policy(require bool) func(context.Context, id.AppID) (bool, string) {
const purpose = "marketing"

return func(context.Context, id.AppID) (bool, string) { return require, purpose }
}

Expand All @@ -28,7 +33,7 @@ func (s *stubConsent) HasConsent(context.Context, id.UserID, id.AppID, string) (

func TestAllowSendPassesWhenGateDisabled(t *testing.T) {
p := New()
p.consentPolicy = policy(false, "marketing")
p.consentPolicy = policy(false)
p.consent = &stubConsent{granted: false}

ok, _, err := p.allowSend(context.Background(), &Job{})
Expand All @@ -39,7 +44,7 @@ func TestAllowSendPassesWhenGateDisabled(t *testing.T) {

func TestAllowSendBlocksWithoutGrant(t *testing.T) {
p := New()
p.consentPolicy = policy(true, "marketing")
p.consentPolicy = policy(true)
p.consent = &stubConsent{granted: false}

ok, reason, err := p.allowSend(context.Background(), &Job{})
Expand All @@ -50,7 +55,7 @@ func TestAllowSendBlocksWithoutGrant(t *testing.T) {

func TestAllowSendPassesWithGrant(t *testing.T) {
p := New()
p.consentPolicy = policy(true, "marketing")
p.consentPolicy = policy(true)
p.consent = &stubConsent{granted: true}

ok, _, err := p.allowSend(context.Background(), &Job{})
Expand All @@ -60,7 +65,7 @@ func TestAllowSendPassesWithGrant(t *testing.T) {

func TestAllowSendBlocksWhenGateOnButConsentUnavailable(t *testing.T) {
p := New()
p.consentPolicy = policy(true, "marketing")
p.consentPolicy = policy(true)
p.consent = nil // consent plugin not registered

ok, reason, err := p.allowSend(context.Background(), &Job{})
Expand All @@ -72,7 +77,7 @@ func TestAllowSendBlocksWhenGateOnButConsentUnavailable(t *testing.T) {

func TestAllowSendReportsLookupFailureAsLocalError(t *testing.T) {
p := New()
p.consentPolicy = policy(true, "marketing")
p.consentPolicy = policy(true)
p.consent = &stubConsent{err: assert.AnError}

ok, reason, err := p.allowSend(context.Background(), &Job{})
Expand Down
16 changes: 6 additions & 10 deletions plugins/retention/provider_generic.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,16 +27,12 @@ type GenericProvider struct {
client *http.Client
}

// classifierPolicyDecided marks that the retry-classification policy in
// classifyHTTPError has been written and recorded in the spec (see
// "Retry classification" in
// docs/superpowers/specs/2026-09-03-crm-retention-delivery-design.md). It
// used to gate construction while the policy was still an open decision; now
// that it is decided, the constant just documents that the guard was here
// and why: an unwritten policy that only surfaced at the first CRM error
// would decide, by accident, whether a transient 503 permanently drops a
// customer's sync.
const classifierPolicyDecided = true
// The retry-classification policy in classifyHTTPError is written down in the
// spec, under "Retry classification" in
// docs/superpowers/specs/2026-09-03-crm-retention-delivery-design.md. It
// matters that it stays written down: an unwritten policy that only surfaced
// at the first CRM error would decide, by accident, whether a transient 503
// permanently drops a customer's sync.

// NewGenericProvider builds a Provider from a ProviderConfig. ContactURL is
// required: Capabilities always advertises CapContacts, so a provider that
Expand Down
53 changes: 29 additions & 24 deletions plugins/retention/provider_generic_classify_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,11 @@ import (

// fakeResp builds a *http.Response carrying just what classifyHTTPError and
// retryAfter look at: status, status code, and headers. No network involved.
// fakeResp builds a response for classifyHTTPError to inspect. It has no
// Body: classification reads only the status and headers, and nothing here
// came from a real transport, so there is no reader to drain. bodyclose
// still flags the call sites because the function returns *http.Response,
// hence the directives there. Closing a nil Body would panic.
func fakeResp(status int, header http.Header) *http.Response {
if header == nil {
header = http.Header{}
Expand Down Expand Up @@ -57,105 +62,105 @@ func TestClassifyHTTPError_Table(t *testing.T) {
},
{
name: "429 with Retry-After delta-seconds",
resp: fakeResp(http.StatusTooManyRequests, headerWithRetryAfter("120")),
resp: fakeResp(http.StatusTooManyRequests, headerWithRetryAfter("120")), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
retryAfter: 2 * time.Minute,
},
{
name: "429 with Retry-After clamped to 30m ceiling",
resp: fakeResp(http.StatusTooManyRequests, headerWithRetryAfter("99999")),
resp: fakeResp(http.StatusTooManyRequests, headerWithRetryAfter("99999")), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
retryAfter: 30 * time.Minute,
},
{
name: "429 with Retry-After floored to 1s",
resp: fakeResp(http.StatusTooManyRequests, headerWithRetryAfter("0")),
resp: fakeResp(http.StatusTooManyRequests, headerWithRetryAfter("0")), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
retryAfter: time.Second,
},
{
name: "429 with garbage Retry-After falls back to backoff",
resp: fakeResp(http.StatusTooManyRequests, headerWithRetryAfter("banana")),
resp: fakeResp(http.StatusTooManyRequests, headerWithRetryAfter("banana")), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
retryAfter: 0,
},
{
name: "429 with no Retry-After header",
resp: fakeResp(http.StatusTooManyRequests, nil),
resp: fakeResp(http.StatusTooManyRequests, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
retryAfter: 0,
},
{
name: "401 unauthorized",
resp: fakeResp(http.StatusUnauthorized, nil),
resp: fakeResp(http.StatusUnauthorized, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
retryAfter: 2 * time.Minute,
},
{
name: "403 forbidden",
resp: fakeResp(http.StatusForbidden, nil),
resp: fakeResp(http.StatusForbidden, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
retryAfter: 2 * time.Minute,
},
{
name: "404 drops the ref",
resp: fakeResp(http.StatusNotFound, nil),
resp: fakeResp(http.StatusNotFound, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
dropRef: true,
},
{
name: "500 internal server error",
resp: fakeResp(http.StatusInternalServerError, nil),
resp: fakeResp(http.StatusInternalServerError, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
},
{
name: "502 bad gateway",
resp: fakeResp(http.StatusBadGateway, nil),
resp: fakeResp(http.StatusBadGateway, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
},
{
name: "503 service unavailable",
resp: fakeResp(http.StatusServiceUnavailable, nil),
resp: fakeResp(http.StatusServiceUnavailable, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
},
{
name: "504 gateway timeout",
resp: fakeResp(http.StatusGatewayTimeout, nil),
resp: fakeResp(http.StatusGatewayTimeout, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
},
{
name: "501 not implemented is terminal",
resp: fakeResp(http.StatusNotImplemented, nil),
resp: fakeResp(http.StatusNotImplemented, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: false,
},
{
name: "400 bad request is terminal",
resp: fakeResp(http.StatusBadRequest, nil),
resp: fakeResp(http.StatusBadRequest, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: false,
},
{
name: "422 unprocessable entity is terminal",
resp: fakeResp(http.StatusUnprocessableEntity, nil),
resp: fakeResp(http.StatusUnprocessableEntity, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: false,
},
{
name: "413 payload too large is terminal",
resp: fakeResp(http.StatusRequestEntityTooLarge, nil),
resp: fakeResp(http.StatusRequestEntityTooLarge, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: false,
},
{
name: "408 request timeout",
resp: fakeResp(http.StatusRequestTimeout, nil),
resp: fakeResp(http.StatusRequestTimeout, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
},
{
name: "409 conflict",
resp: fakeResp(http.StatusConflict, nil),
resp: fakeResp(http.StatusConflict, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: true,
},
{
name: "418 unrecognised status is terminal",
resp: fakeResp(http.StatusTeapot, nil),
resp: fakeResp(http.StatusTeapot, nil), //nolint:bodyclose // synthetic response, Body is nil
retryable: false,
},
}
Expand All @@ -180,7 +185,7 @@ func TestClassifyHTTPError_NilResponsePreservesTransportError(t *testing.T) {

func TestRetryAfter_HTTPDateForm(t *testing.T) {
when := time.Now().Add(90 * time.Second).UTC()
resp := fakeResp(http.StatusTooManyRequests, headerWithRetryAfter(when.Format(http.TimeFormat)))
resp := fakeResp(http.StatusTooManyRequests, headerWithRetryAfter(when.Format(http.TimeFormat))) //nolint:bodyclose // synthetic response, Body is nil

got := retryAfter(resp)
assert.InDelta(t, 90*time.Second, got, float64(5*time.Second),
Expand All @@ -189,7 +194,7 @@ func TestRetryAfter_HTTPDateForm(t *testing.T) {

func TestClassifyHTTPError_429HTTPDateRetryAfter(t *testing.T) {
when := time.Now().Add(90 * time.Second).UTC()
resp := fakeResp(http.StatusTooManyRequests, headerWithRetryAfter(when.Format(http.TimeFormat)))
resp := fakeResp(http.StatusTooManyRequests, headerWithRetryAfter(when.Format(http.TimeFormat))) //nolint:bodyclose // synthetic response, Body is nil

pe := classifyHTTPError(resp, nil, nil)
require.NotNil(t, pe)
Expand All @@ -214,7 +219,7 @@ func TestTruncate(t *testing.T) {
// ──────────────────────────────────────────────────

func TestGenericProvider_UpsertContact_404EndToEnd(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusNotFound)
_, _ = w.Write([]byte(`{"error":"contact not found"}`))
}))
Expand All @@ -234,7 +239,7 @@ func TestGenericProvider_UpsertContact_404EndToEnd(t *testing.T) {
}

func TestGenericProvider_UpsertContact_429WithRetryAfterEndToEnd(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Retry-After", "45")
w.WriteHeader(http.StatusTooManyRequests)
_, _ = w.Write([]byte(`slow down`))
Expand All @@ -255,7 +260,7 @@ func TestGenericProvider_UpsertContact_429WithRetryAfterEndToEnd(t *testing.T) {
}

func TestGenericProvider_UpsertContact_400IsTerminalEndToEnd(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusBadRequest)
_, _ = w.Write([]byte(`{"error":"invalid email"}`))
}))
Expand Down
6 changes: 3 additions & 3 deletions plugins/retention/provider_generic_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ func TestGenericProviderName(t *testing.T) {
}

func TestGenericProviderUpsertContact_Success(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte(`{"id":"remote-501"}`))
Expand Down Expand Up @@ -170,7 +170,7 @@ func TestGenericProviderUpsertContact_FieldMapRenamesOutgoingFields(t *testing.T
}

func TestGenericProviderUpsertContact_RemoteIDDefaultsToID(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte(`{"id":"default-id-1","other":"ignored"}`))
}))
Expand All @@ -185,7 +185,7 @@ func TestGenericProviderUpsertContact_RemoteIDDefaultsToID(t *testing.T) {
}

func TestGenericProviderUpsertContact_RemoteIDFromConfiguredPath(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte(`{"result":{"contact_id":"nested-777"}}`))
}))
Expand Down
4 changes: 2 additions & 2 deletions plugins/retention/provider_hubspot.go
Original file line number Diff line number Diff line change
Expand Up @@ -152,8 +152,8 @@ func (h *HubSpotProvider) UpsertContact(ctx context.Context, c *Contact) (Remote

if remoteID != "" {
path := hubspotContactsPath + "/" + remoteID
if _, err := h.request(ctx, http.MethodPatch, path, map[string]interface{}{"properties": props}); err != nil {
return RemoteRef{}, err
if _, patchErr := h.request(ctx, http.MethodPatch, path, map[string]interface{}{"properties": props}); patchErr != nil {
return RemoteRef{}, patchErr
}
return RemoteRef{Provider: h.name, ObjectType: "contact", ID: remoteID}, nil
}
Expand Down
5 changes: 4 additions & 1 deletion plugins/retention/store_models.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,10 @@ type jobModel struct {
}

func fromJob(j *Job) *jobModel {
payload, _ := json.Marshal(j.Payload)
// Payload is a map[string]string, which has no channel, func or cyclic
// value that could make Marshal fail, so there is no error path here to
// handle. Marshalling anything richer would need this revisited.
payload, _ := json.Marshal(j.Payload) //nolint:errcheck // cannot fail for map[string]string
m := &jobModel{
ID: j.ID.String(), AppID: j.AppID.String(), EnvID: j.EnvID.String(),
UserID: j.UserID.String(), Provider: j.Provider, Kind: j.Kind,
Expand Down
Loading