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
4 changes: 4 additions & 0 deletions internal/jobs/billing_coverage_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -405,6 +405,10 @@ func (g *stubGraceLocal) HasTerminatedGracePeriod(_ context.Context, _ uuid.UUID
return g.hasTerminated, g.termErr
}

func (g *stubGraceLocal) TerminateActiveGracePeriod(_ context.Context, _ uuid.UUID) error {
return nil
}

// TestBillingReconciler_Work_NotConfigured_AbortsBatch covers the
// errSubFetcherNotConfigured branch inside Work.
func TestBillingReconciler_Work_NotConfigured_AbortsBatch(t *testing.T) {
Expand Down
47 changes: 43 additions & 4 deletions internal/jobs/billing_reconciler.go
Original file line number Diff line number Diff line change
Expand Up @@ -330,6 +330,14 @@ type gracePeriodOpener interface {
// would see "no ACTIVE grace" and open a FRESH 7-day grace period,
// restarting the dunning-email cycle indefinitely.
HasTerminatedGracePeriod(ctx context.Context, teamID uuid.UUID, subscriptionID string) (bool, error)
// TerminateActiveGracePeriod closes any 'active' grace row for the team
// (status→'terminated', terminated_at=now()). Called when the subscription
// reaches a TERMINAL Razorpay status: without it the reconciler downgrades
// the team but leaves the grace row 'active', so payment_grace_reminder
// keeps emitting dunning emails forever and payment_grace_terminator later
// re-acts on an already-cancelled subscription (bug bash 2026-06-02 #5).
// Mirrors the terminate UPDATE in api models/payment_grace_periods.go.
TerminateActiveGracePeriod(ctx context.Context, teamID uuid.UUID) error
}

// gracePeriodTerminalStatuses are the payment_grace_periods.status values that
Expand Down Expand Up @@ -430,6 +438,23 @@ func (d *dbGracePeriodOpener) OpenGracePeriod(ctx context.Context, teamID uuid.U
return nil
}

// TerminateActiveGracePeriod closes any 'active' grace row for the team —
// status→'terminated', terminated_at=now(). Idempotent: a team with no active
// grace row updates zero rows and returns nil. Mirrors the terminate UPDATE in
// api/internal/models/payment_grace_periods.go so the worker and api converge
// on the same terminal state.
func (d *dbGracePeriodOpener) TerminateActiveGracePeriod(ctx context.Context, teamID uuid.UUID) error {
_, err := d.db.ExecContext(ctx, `
UPDATE payment_grace_periods
SET status = 'terminated', terminated_at = now()
WHERE team_id = $1 AND status = 'active'
`, teamID)
if err != nil {
return fmt.Errorf("dbGracePeriodOpener.TerminateActiveGracePeriod: %w", err)
}
return nil
}

// HasTerminatedGracePeriod reports whether the team already has a grace
// period for subscriptionID in a terminal status (see
// gracePeriodTerminalStatuses). When true the reconciler must NOT open a
Expand Down Expand Up @@ -1001,6 +1026,14 @@ func (w *BillingReconcilerWorker) Work(ctx context.Context, job *river.Job[Billi
}
correctedDowngrade++
metrics.BillingReconcilerGapCorrected.WithLabelValues("downgrade").Inc()
// Close any active grace period so the dunning reminder stops and
// the terminator doesn't re-act on this now-cancelled
// subscription (bug bash #5). Fail-open: the downgrade is already
// committed; a stuck grace row only costs extra dunning emails.
if gErr := w.grace.TerminateActiveGracePeriod(ctx, team.id); gErr != nil {
slog.Warn("billing.reconciler.grace_terminate_failed",
"team_id", team.id, "subscription_id", team.subscriptionID, "error", gErr)
}
// Emit audit for the event-email forwarder. Fail-open.
w.emitCancelAudit(ctx, team.id, team.planTier, targetTier, team.subscriptionID)

Expand Down Expand Up @@ -1107,11 +1140,17 @@ func (w *BillingReconcilerWorker) scanChargeUndeliverable(ctx context.Context) i
)
}

// Advance the cursor to the latest seen row. If count==0 we still
// advance to now() — saves re-scanning the same empty window next
// tick, and there's nothing in the window to lose.
// Advance the cursor ONLY when we actually saw rows — to the latest seen
// created_at. On an EMPTY window we must NOT jump the cursor to now():
// the strict `>` predicate combined with a now() that is ahead of a row's
// created_at — clock skew between this worker and the platform DB, or a
// transaction that committed late but stamped created_at with an earlier
// DB now() — would push the watermark past a row that becomes visible a
// moment later, permanently skipping it. Re-scanning the same small,
// (kind, created_at)-indexed window next tick is cheap, so leave the
// cursor unchanged when count==0 (bug bash 2026-06-02 #18).
if maxCreated.IsZero() {
maxCreated = time.Now().UTC()
return count // count == 0 — nothing seen, cursor stays put
}
w.chargeUndeliverableMu.Lock()
w.chargeUndeliverableCursor = maxCreated
Expand Down
48 changes: 39 additions & 9 deletions internal/jobs/billing_reconciler_charge_undeliverable_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ package jobs

import (
"context"
"database/sql/driver"
"errors"
"testing"
"time"
Expand Down Expand Up @@ -89,8 +90,12 @@ func TestScanChargeUndeliverable_NoNewRows(t *testing.T) {
}
w.chargeUndeliverableMu.Lock()
defer w.chargeUndeliverableMu.Unlock()
if !w.chargeUndeliverableCursor.After(prev) {
t.Fatalf("cursor should advance to now() even on zero rows: prev=%v cur=%v", prev, w.chargeUndeliverableCursor)
// bug bash #18: the cursor must NOT advance on an empty window. Jumping it
// to now() would let the strict `>` predicate skip a row that becomes
// visible a moment later (clock skew / late-committing INSERT). Re-scanning
// the same small window next tick is cheap, so the cursor stays put.
if !w.chargeUndeliverableCursor.Equal(prev) {
t.Fatalf("cursor must NOT advance on zero-row scan (bug #18): prev=%v cur=%v", prev, w.chargeUndeliverableCursor)
}
}

Expand Down Expand Up @@ -123,9 +128,12 @@ func TestScanChargeUndeliverable_DBErrorFailsOpen(t *testing.T) {
}
}

// TestScanChargeUndeliverable_FirstTickUsesLookback — zero-value cursor
// causes the scanner to seed at now()-1h on the first tick after pod
// boot.
// TestScanChargeUndeliverable_FirstTickUsesLookback — a zero-value cursor
// makes the scanner QUERY from now()-1h on every tick after pod boot, until a
// row is actually seen. bug bash #18: on an empty result the persisted cursor
// is NOT advanced (it stays zero), so the 1h look-back keeps re-applying — a
// row landing in that window is always caught, never skipped. The query arg is
// asserted to be ~now()-1h to prove the look-back is applied.
func TestScanChargeUndeliverable_FirstTickUsesLookback(t *testing.T) {
db, mock, err := sqlmock.New()
if err != nil {
Expand All @@ -134,17 +142,39 @@ func TestScanChargeUndeliverable_FirstTickUsesLookback(t *testing.T) {
defer db.Close()

mock.ExpectQuery(`SELECT created_at FROM audit_log`).
WithArgs(chargeUndeliverableAuditKind, sqlmock.AnyArg()).
WithArgs(chargeUndeliverableAuditKind, lookbackArg{around: time.Now().UTC().Add(-1 * time.Hour), tol: 2 * time.Minute}).
WillReturnRows(sqlmock.NewRows([]string{"created_at"}))

w := &BillingReconcilerWorker{db: db}
before := time.Now().UTC()
_ = w.scanChargeUndeliverable(context.Background())

w.chargeUndeliverableMu.Lock()
cursor := w.chargeUndeliverableCursor
w.chargeUndeliverableMu.Unlock()
if cursor.Before(before) {
t.Fatalf("first-tick cursor should advance to ~now: cursor=%v before=%v", cursor, before)
// Empty result → cursor stays zero so the look-back re-applies next tick.
if !cursor.IsZero() {
t.Fatalf("empty first-tick must leave the cursor unadvanced (zero) so the look-back re-applies (bug #18); got %v", cursor)
}
if err := mock.ExpectationsWereMet(); err != nil {
t.Fatalf("query did not use the now()-1h look-back arg: %v", err)
}
}

// lookbackArg is a sqlmock matcher asserting a time.Time arg is within tol of
// the expected look-back instant.
type lookbackArg struct {
around time.Time
tol time.Duration
}

func (m lookbackArg) Match(v driver.Value) bool {
t, ok := v.(time.Time)
if !ok {
return false
}
d := t.Sub(m.around)
if d < 0 {
d = -d
}
return d <= m.tol
}
50 changes: 46 additions & 4 deletions internal/jobs/billing_reconciler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,10 +60,12 @@ func (s *stubFetcher) FetchSubscriptionForReconciler(_ context.Context, _ string

// stubGrace implements gracePeriodOpener for tests.
type stubGrace struct {
hasActive bool
hasTerminated bool // P1-F(b): a prior grace period reached a terminal status
openCalls int
openErr error
hasActive bool
hasTerminated bool // P1-F(b): a prior grace period reached a terminal status
openCalls int
openErr error
terminateCalls int // #5: grace closed on terminal downgrade
terminateErr error
}

func (g *stubGrace) GetActiveGracePeriod(_ context.Context, _ uuid.UUID) (bool, error) {
Expand All @@ -79,6 +81,11 @@ func (g *stubGrace) HasTerminatedGracePeriod(_ context.Context, _ uuid.UUID, _ s
return g.hasTerminated, nil
}

func (g *stubGrace) TerminateActiveGracePeriod(_ context.Context, _ uuid.UUID) error {
g.terminateCalls++
return g.terminateErr
}

// teamRowCols are the columns the billing reconciler SELECT returns.
var teamRowCols = []string{"id", "stripe_customer_id", "plan_tier"}

Expand Down Expand Up @@ -1076,3 +1083,38 @@ func TestBillingReconciler_OrphanSweep_QueryFailure_FailOpen(t *testing.T) {
t.Errorf("unmet expectations: %v", err)
}
}

// bug bash #5: a terminal downgrade closes the active grace period; if the
// close itself errors the downgrade still succeeds (fail-open) — exercises the
// grace_terminate_failed warn branch.
func TestBillingReconciler_CancelledSubscription_GraceTerminateError_StillDowngrades(t *testing.T) {
db, mock, err := sqlmock.New(sqlmock.QueryMatcherOption(sqlmock.QueryMatcherRegexp))
if err != nil {
t.Fatalf("sqlmock.New: %v", err)
}
defer db.Close()

teamID := uuid.New()
mock.ExpectQuery(`SELECT id, stripe_customer_id, plan_tier`).
WillReturnRows(sqlmock.NewRows(teamRowCols).AddRow(teamID, "sub_grace_err", "pro"))
mock.ExpectExec(`UPDATE teams SET plan_tier`).
WithArgs("hobby", teamID).
WillReturnResult(sqlmock.NewResult(1, 1))
mock.ExpectExec(`INSERT INTO audit_log`).
WillReturnResult(sqlmock.NewResult(1, 1))
expectEmptyOrphanSweep(mock)

fetcher := &stubFetcher{details: &jobs.ReconcilerSubDetails{Status: "cancelled", PlanID: "", PaidCount: 3}}
grace := &stubGrace{terminateErr: errors.New("grace close failed")}

w := jobs.NewBillingReconcilerWorker(db, fetcher, grace)
if err := w.Work(context.Background(), fakeJob[jobs.BillingReconcilerArgs]()); err != nil {
t.Fatalf("downgrade must succeed despite grace-close error (fail-open): %v", err)
}
if grace.terminateCalls != 1 {
t.Errorf("TerminateActiveGracePeriod calls = %d; want 1", grace.terminateCalls)
}
if err := mock.ExpectationsWereMet(); err != nil {
t.Errorf("unmet: %v", err)
}
}
73 changes: 73 additions & 0 deletions internal/jobs/bugbash2_coverage_internal_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
package jobs

// bugbash2_coverage_internal_test.go — internal (package jobs) coverage for
// the bug-bash batch-2 changes whose lines aren't reachable from the external
// test package: the autopsy deploy.failed dedup-hit branch and the real
// dbGracePeriodOpener.TerminateActiveGracePeriod UPDATE. All hermetic (sqlmock).

import (
"context"
"errors"
"testing"

sqlmock "github.com/DATA-DOG/go-sqlmock"
"github.com/google/uuid"
)

// #15: when a deploy.failed audit row already exists for the deployment,
// emitDeployFailedAudit must SKIP the INSERT (idempotent — no duplicate email).
func TestEmitDeployFailedAudit_SkipsWhenAlreadyEmitted(t *testing.T) {
db, mock, err := sqlmock.New(sqlmock.QueryMatcherOption(sqlmock.QueryMatcherRegexp))
if err != nil {
t.Fatalf("sqlmock.New: %v", err)
}
defer db.Close()
depID := uuid.New()

mock.ExpectQuery(`SELECT team_id FROM deployments`).
WithArgs(depID).
WillReturnRows(sqlmock.NewRows([]string{"team_id"}).AddRow(uuid.New()))
// Dedup probe finds an existing deploy.failed row → no INSERT must follow.
mock.ExpectQuery(`SELECT EXISTS`).
WillReturnRows(sqlmock.NewRows([]string{"exists"}).AddRow(true))

if err := emitDeployFailedAudit(context.Background(), db, depID, "BuildFailed", "boom"); err != nil {
t.Fatalf("emitDeployFailedAudit: %v", err)
}
if err := mock.ExpectationsWereMet(); err != nil {
t.Errorf("an INSERT must NOT run when a deploy.failed row already exists: %v", err)
}
}

// #5: the real dbGracePeriodOpener.TerminateActiveGracePeriod issues the
// status→'terminated' UPDATE; an exec error is wrapped and returned.
func TestDBGracePeriodOpener_TerminateActiveGracePeriod(t *testing.T) {
teamID := uuid.New()

t.Run("success", func(t *testing.T) {
db, mock, _ := sqlmock.New(sqlmock.QueryMatcherOption(sqlmock.QueryMatcherRegexp))
defer db.Close()
mock.ExpectExec(`UPDATE payment_grace_periods\s+SET status = 'terminated'`).
WithArgs(teamID).
WillReturnResult(sqlmock.NewResult(0, 1))
d := &dbGracePeriodOpener{db: db}
if err := d.TerminateActiveGracePeriod(context.Background(), teamID); err != nil {
t.Fatalf("TerminateActiveGracePeriod: %v", err)
}
if err := mock.ExpectationsWereMet(); err != nil {
t.Errorf("unmet: %v", err)
}
})

t.Run("db error wrapped", func(t *testing.T) {
db, mock, _ := sqlmock.New(sqlmock.QueryMatcherOption(sqlmock.QueryMatcherRegexp))
defer db.Close()
mock.ExpectExec(`UPDATE payment_grace_periods`).
WithArgs(teamID).
WillReturnError(errors.New("boom"))
d := &dbGracePeriodOpener{db: db}
if err := d.TerminateActiveGracePeriod(context.Background(), teamID); err == nil {
t.Fatal("expected error to propagate")
}
})
}
26 changes: 26 additions & 0 deletions internal/jobs/deploy_failure_autopsy.go
Original file line number Diff line number Diff line change
Expand Up @@ -633,6 +633,32 @@ func emitDeployFailedAudit(ctx context.Context, db *sql.DB, deploymentID uuid.UU
"error_summary": summary,
"source": "worker_autopsy",
}
// Idempotency guard (bug bash 2026-06-02 #15): this runs every time the
// status reconciler observes the deployment in a failed state, and the
// reconciler re-lists the row whenever the subsequent status UPDATE fails
// (it only excludes terminal rows once the flip succeeds). Without this
// guard each retry — and the api's own deploy.failed emit — inserts a
// fresh audit_log row with a NEW id, and the email forwarder (which dedups
// by audit_id, not by deployment) sends a duplicate failure email per
// retry. Skip the INSERT when a deploy.failed row already exists for this
// deployment.
var alreadyEmitted bool
if err := db.QueryRowContext(ctx, `
SELECT EXISTS (
SELECT 1 FROM audit_log
WHERE kind = $1 AND metadata->>'deploy_id' = $2
)
`, auditKindDeployFailed, deploymentID.String()).Scan(&alreadyEmitted); err != nil {
// Fail-open: if the dedup probe errors, fall through and insert — a
// possible duplicate email is better than dropping the failure
// notification entirely.
slog.Warn("jobs.deploy_failure_autopsy.dedup_probe_failed",
"deploy_id", deploymentID, "error", err,
"note", "inserting deploy.failed anyway (fail-open)")
} else if alreadyEmitted {
return nil // a deploy.failed audit row already exists for this deployment
}

// json.Marshal of a map[string]any with string keys + string values is
// total — unreachable error path. The orphan-sweep audit emit follows
// the same _-ignore pattern (orphan_sweep_reconciler.go:emitOrphanAudit).
Expand Down
Loading
Loading