diff --git a/.github/workflows/golangci-lint.yml b/.github/workflows/golangci-lint.yml new file mode 100644 index 0000000..c097451 --- /dev/null +++ b/.github/workflows/golangci-lint.yml @@ -0,0 +1,48 @@ +name: golangci-lint + +on: + push: + branches: [master, main] + pull_request: + branches: [master, main] + schedule: + - cron: "23 6 * * 1" + +permissions: + contents: read + pull-requests: read + +jobs: + lint: + name: lint + runs-on: ubuntu-latest + timeout-minutes: 10 + steps: + - uses: actions/checkout@v4 + with: + path: worker + # Sibling checkouts (proto/common) for repos with replace directives. + # No-op for repos that do not need them. + - uses: actions/checkout@v4 + if: ${{ hashFiles('worker/go.mod') != '' }} + with: + repository: InstaNode-dev/common + path: common + continue-on-error: true + - uses: actions/checkout@v4 + with: + repository: InstaNode-dev/proto + path: proto + continue-on-error: true + - uses: actions/setup-go@v5 + with: + go-version-file: worker/go.mod + - uses: golangci/golangci-lint-action@v8 + with: + # Pinned (was `latest`) so a new upstream release can't silently + # introduce analyzers/findings that fail a previously-green PR — + # `latest` floated to v2.12.2 and turned on a govet `inline` check + # mid-PR. Bump deliberately + re-run `make gate` when upgrading. + version: v2.12.2 + working-directory: worker + args: --timeout=5m diff --git a/.golangci.yml b/.golangci.yml new file mode 100644 index 0000000..a0a3149 --- /dev/null +++ b/.golangci.yml @@ -0,0 +1,43 @@ +# golangci-lint v2 config — start conservative, expand once baseline is clean +version: "2" + +run: + timeout: 5m + tests: true + +linters: + # Default linter set is govet+errcheck+ineffassign+staticcheck+unused. + # We explicitly add misspell + gocyclo on top. gosimple folded into staticcheck in v2. + enable: + - errcheck # checks unchecked errors + - govet # standard vet + - ineffassign # ineffective assignments + - staticcheck # bug detection (subsumes gosimple in v2) + - unused # unused code + - misspell # spelling + - gocyclo # cyclomatic complexity + settings: + gocyclo: + # Set just above the current top offender (33: entitlement_reconciler + # / event_email_forwarder Work funcs). These are inherently branchy + # job dispatch/reconcile loops; a config bump is preferred over a + # behaviour-changing refactor for this mechanical lint-gate pass. + min-complexity: 34 + errcheck: + # The fmt.Fprint* family writes to stdout/stderr/buffers in CLI and + # diagnostic code where a write error is neither actionable nor + # meaningful. Matches the convention adopted by the sibling repos. + exclude-functions: + - fmt.Fprint + - fmt.Fprintf + - fmt.Fprintln + exclusions: + rules: + - path: _test\.go + linters: + - errcheck + - gocyclo + +issues: + max-issues-per-linter: 0 + max-same-issues: 0 diff --git a/internal/email/brevo_provider.go b/internal/email/brevo_provider.go index 6ee3913..6d8f472 100644 --- a/internal/email/brevo_provider.go +++ b/internal/email/brevo_provider.go @@ -388,7 +388,7 @@ func (p *BrevoProvider) doRequest(ctx context.Context, evt EventEmail, body []by ) return "", &SendError{Class: SendClassTransient, Cause: err, Message: "brevo: http"} } - defer resp.Body.Close() + defer func() { _ = resp.Body.Close() }() respBody, _ := io.ReadAll(io.LimitReader(resp.Body, brevoBodyReadCap)) statusLabel := strconv.Itoa(resp.StatusCode) diff --git a/internal/jobs/backup_extra_test.go b/internal/jobs/backup_extra_test.go index 5bb6c5e..6952406 100644 --- a/internal/jobs/backup_extra_test.go +++ b/internal/jobs/backup_extra_test.go @@ -1592,7 +1592,7 @@ func TestDefaultPgDumpExec_StderrTruncates(t *testing.T) { dir := t.TempDir() bin := filepath.Join(dir, "noisy_pg_dump") // Write 1000 bytes to stderr then exit 1. - script := fmt.Sprintf("#!/bin/sh\nfor i in $(seq 1 100); do printf 'XXXXXXXXXX' >&2; done\nexit 1\n") + script := "#!/bin/sh\nfor i in $(seq 1 100); do printf 'XXXXXXXXXX' >&2; done\nexit 1\n" if err := os.WriteFile(bin, []byte(script), 0o755); err != nil { t.Fatalf("write: %v", err) } diff --git a/internal/jobs/billing_reconciler.go b/internal/jobs/billing_reconciler.go index 5f0d35e..a9fd020 100644 --- a/internal/jobs/billing_reconciler.go +++ b/internal/jobs/billing_reconciler.go @@ -803,7 +803,7 @@ func (w *BillingReconcilerWorker) Work(ctx context.Context, job *river.Job[Billi if err != nil { return fmt.Errorf("BillingReconcilerWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var teams []billingReconcilerTeamRow for rows.Next() { @@ -817,7 +817,7 @@ func (w *BillingReconcilerWorker) Work(ctx context.Context, job *river.Job[Billi if rowsErr := rows.Err(); rowsErr != nil { return fmt.Errorf("BillingReconcilerWorker: rows error: %w", rowsErr) } - rows.Close() + _ = rows.Close() metrics.BillingReconcilerTeamsScanned.Add(float64(len(teams))) @@ -1071,7 +1071,7 @@ func (w *BillingReconcilerWorker) scanChargeUndeliverable(ctx context.Context) i "error", err, "note", "fail-open — retry next tick") return 0 } - defer rows.Close() + defer func() { _ = rows.Close() }() var count int var maxCreated time.Time @@ -1176,7 +1176,7 @@ func (w *BillingReconcilerWorker) runOrphanSweep(ctx context.Context) (scanned, ) return 0, 0 } - defer rows.Close() + defer func() { _ = rows.Close() }() var candidates []billingReconcilerOrphanRow for rows.Next() { @@ -1191,7 +1191,7 @@ func (w *BillingReconcilerWorker) runOrphanSweep(ctx context.Context) (scanned, slog.Warn("billing.reconciler.orphan_rows_error", "error", rowsErr) return 0, 0 } - rows.Close() + _ = rows.Close() for i, c := range candidates { // 100ms stagger between Razorpay calls, same as the primary sweep. diff --git a/internal/jobs/checkout_reconcile.go b/internal/jobs/checkout_reconcile.go index 9104497..798039a 100644 --- a/internal/jobs/checkout_reconcile.go +++ b/internal/jobs/checkout_reconcile.go @@ -158,7 +158,7 @@ func (w *CheckoutReconcileWorker) Work(ctx context.Context, job *river.Job[Check if err != nil { return fmt.Errorf("CheckoutReconcileWorker: candidate query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var candidates []checkoutRow for rows.Next() { @@ -172,7 +172,7 @@ func (w *CheckoutReconcileWorker) Work(ctx context.Context, job *river.Job[Check if err := rows.Err(); err != nil { return fmt.Errorf("CheckoutReconcileWorker: rows error: %w", err) } - rows.Close() + _ = rows.Close() if len(candidates) == 0 { // P1-1 (BugBash 2026-05-19): idle tick — zero candidates carries diff --git a/internal/jobs/churn_predictor.go b/internal/jobs/churn_predictor.go index 65741a3..040e723 100644 --- a/internal/jobs/churn_predictor.go +++ b/internal/jobs/churn_predictor.go @@ -76,12 +76,6 @@ type ChurnPredictorArgs struct{} // Kind is the River worker key. func (ChurnPredictorArgs) Kind() string { return "churn_predictor" } -// churnPredictorInterval is the periodic dispatch cadence — daily. -// The 30-day dedupe window means cadence below ~daily produces no -// additional flags; daily gives at most one chance per day to catch a -// team that crossed the 7-day threshold yesterday. -const churnPredictorInterval = 24 * time.Hour - // churnInactivityWindow is the silence period that defines a churn // candidate. 7 days is the brief's threshold — long enough that a // vacation or busy week doesn't trigger a "we miss you" email @@ -231,7 +225,7 @@ func (w *ChurnPredictorWorker) Work(ctx context.Context, job *river.Job[ChurnPre if err != nil { return fmt.Errorf("ChurnPredictorWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var candidates []churnCandidateRow for rows.Next() { @@ -249,7 +243,7 @@ func (w *ChurnPredictorWorker) Work(ctx context.Context, job *river.Job[ChurnPre if err := rows.Err(); err != nil { return fmt.Errorf("ChurnPredictorWorker: rows error: %w", err) } - rows.Close() + _ = rows.Close() if len(candidates) == 0 { slog.Info("jobs.churn_predictor.completed", diff --git a/internal/jobs/coverage_misc_test.go b/internal/jobs/coverage_misc_test.go index 3b759f3..a0431cf 100644 --- a/internal/jobs/coverage_misc_test.go +++ b/internal/jobs/coverage_misc_test.go @@ -21,7 +21,6 @@ import ( "fmt" "net/http" "net/http/httptest" - "net/url" "os" "path/filepath" "strings" @@ -1277,10 +1276,6 @@ func TestProber_UptimeRetention_DBError_Propagates(t *testing.T) { // ─── real_prober.go ──────────────────────────────────────────────────────────── func TestProber_NormalizeStorageURL_AllArms(t *testing.T) { - if _, err := url.Parse("://bad"); err == nil { - // just here to silence the unused-import linter if it ever fires - } - // http/https → returned untouched. if got, err := normalizeStorageURL("https://example.com/bucket"); err != nil || got != "https://example.com/bucket" { t.Errorf("https returned (%q, %v)", got, err) diff --git a/internal/jobs/custom_domain_reconcile.go b/internal/jobs/custom_domain_reconcile.go index 7b3c96b..9d39e0e 100644 --- a/internal/jobs/custom_domain_reconcile.go +++ b/internal/jobs/custom_domain_reconcile.go @@ -294,7 +294,7 @@ func (w *CustomDomainReconciler) reconcileCertReady(ctx context.Context, d activ _ = w.updateLastCheck(ctx, d.id, fmt.Sprintf("HTTPS HEAD probe failed: %v", err)) return reconcileNoop } - defer resp.Body.Close() + defer func() { _ = resp.Body.Close() }() if resp.StatusCode >= 200 && resp.StatusCode < 400 { if err := w.updateStatus(ctx, d.id, statusLive, ""); err != nil { @@ -366,7 +366,7 @@ func (w *CustomDomainReconciler) listActiveDomains(ctx context.Context) ([]activ if err != nil { return nil, fmt.Errorf("listActiveDomains: query: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var out []activeCustomDomain for rows.Next() { diff --git a/internal/jobs/customer_backup_runner.go b/internal/jobs/customer_backup_runner.go index b467a40..26b1bce 100644 --- a/internal/jobs/customer_backup_runner.go +++ b/internal/jobs/customer_backup_runner.go @@ -231,7 +231,7 @@ func (w *CustomerBackupRunnerWorker) Work(ctx context.Context, job *river.Job[Cu if err != nil { return fmt.Errorf("CustomerBackupRunnerWorker: select pending failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() type pending struct { backupID string @@ -258,7 +258,7 @@ func (w *CustomerBackupRunnerWorker) Work(ctx context.Context, job *river.Job[Cu if err := rows.Err(); err != nil { return fmt.Errorf("CustomerBackupRunnerWorker: rows error: %w", err) } - rows.Close() + _ = rows.Close() processed := 0 succeeded := 0 @@ -695,7 +695,7 @@ func (w *CustomerBackupRunnerWorker) runRetentionSweep(ctx context.Context) { } victims = append(victims, v) } - rows.Close() + _ = rows.Close() for _, v := range victims { if delErr := w.store.DeleteObject(ctx, w.bucket, v.s3Key); delErr != nil { @@ -796,7 +796,7 @@ func (w *CustomerBackupRunnerWorker) refundManualBackupQuota(teamID uuid.UUID, b } return fmt.Errorf("api request: %w", doErr) } - defer resp.Body.Close() + defer func() { _ = resp.Body.Close() }() if resp.StatusCode < 200 || resp.StatusCode >= 300 { body, _ := io.ReadAll(io.LimitReader(resp.Body, 1024)) return fmt.Errorf("api status %d: %s", resp.StatusCode, strings.TrimSpace(string(body))) diff --git a/internal/jobs/customer_backup_scheduler.go b/internal/jobs/customer_backup_scheduler.go index d9c65e9..92ae94e 100644 --- a/internal/jobs/customer_backup_scheduler.go +++ b/internal/jobs/customer_backup_scheduler.go @@ -128,7 +128,7 @@ func (w *CustomerBackupSchedulerWorker) Work(ctx context.Context, job *river.Job if err != nil { return fmt.Errorf("CustomerBackupSchedulerWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() type cand struct { id string diff --git a/internal/jobs/customer_restore_runner.go b/internal/jobs/customer_restore_runner.go index 0068ccb..e12606e 100644 --- a/internal/jobs/customer_restore_runner.go +++ b/internal/jobs/customer_restore_runner.go @@ -161,7 +161,7 @@ func (w *CustomerRestoreRunnerWorker) Work(ctx context.Context, job *river.Job[C if err != nil { return fmt.Errorf("CustomerRestoreRunnerWorker: select pending failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() type pending struct { restoreID string @@ -189,7 +189,7 @@ func (w *CustomerRestoreRunnerWorker) Work(ctx context.Context, job *river.Job[C if err := rows.Err(); err != nil { return fmt.Errorf("CustomerRestoreRunnerWorker: rows error: %w", err) } - rows.Close() + _ = rows.Close() processed := 0 succeeded := 0 @@ -352,7 +352,7 @@ func (w *CustomerRestoreRunnerWorker) processRestore(parentCtx context.Context, tmpFile, tmpErr := os.CreateTemp("", "instant-restore-*.dump.gz") if tmpErr != nil { - obj.Close() + _ = obj.Close() w.markRestoreFailed(ctx, p.restoreID, fmt.Sprintf("create temp file: %v", tmpErr), start, p) return false } @@ -366,11 +366,11 @@ func (w *CustomerRestoreRunnerWorker) processRestore(parentCtx context.Context, hasher := sha256.New() tee := io.TeeReader(obj, hasher) if _, copyErr := io.Copy(tmpFile, tee); copyErr != nil { - obj.Close() + _ = obj.Close() w.markRestoreFailed(ctx, p.restoreID, fmt.Sprintf("S3 read failed: %v", copyErr), start, p) return false } - obj.Close() + _ = obj.Close() // Integrity gate (same semantics as before — only the hashing path // changed). The backup runner (customer_backup_runner.go, FIX-H #59) @@ -412,7 +412,7 @@ func (w *CustomerRestoreRunnerWorker) processRestore(parentCtx context.Context, w.markRestoreFailed(ctx, p.restoreID, fmt.Sprintf("gunzip header: %v", gzErr), start, p) return false } - defer gzReader.Close() + defer func() { _ = gzReader.Close() }() if runErr := w.pgRestore.Run(ctx, plainConn, gzReader); runErr != nil { w.markRestoreFailed(ctx, p.restoreID, fmt.Sprintf("pg_restore failed: %v", runErr), start, p) diff --git a/internal/jobs/deploy_failure_autopsy.go b/internal/jobs/deploy_failure_autopsy.go index e6027a8..2201acb 100644 --- a/internal/jobs/deploy_failure_autopsy.go +++ b/internal/jobs/deploy_failure_autopsy.go @@ -192,7 +192,7 @@ func (c *k8sAutopsyClient) GetPodLogs(ctx context.Context, namespace, podName st return nil, nil //nolint:nilerr } } - defer stream.Close() + defer func() { _ = stream.Close() }() return readLogLines(stream) } diff --git a/internal/jobs/deploy_notify_webhook.go b/internal/jobs/deploy_notify_webhook.go index da7db8f..3f33537 100644 --- a/internal/jobs/deploy_notify_webhook.go +++ b/internal/jobs/deploy_notify_webhook.go @@ -461,7 +461,7 @@ func (w *DeployNotifyWebhookWorker) fetchBatch(ctx context.Context, c deployNoti if err != nil { return nil, fmt.Errorf("fetchBatch query: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var out []deployNotifyAuditRow for rows.Next() { @@ -569,7 +569,7 @@ func (w *DeployNotifyWebhookWorker) dispatch(ctx context.Context, urlStr string, // Drain + close body unconditionally so the connection // returns to the keep-alive pool. _, _ = io.Copy(io.Discard, resp.Body) - resp.Body.Close() + _ = resp.Body.Close() if resp.StatusCode >= 200 && resp.StatusCode < 300 { return nil } diff --git a/internal/jobs/deploy_status_reconcile.go b/internal/jobs/deploy_status_reconcile.go index bd98aad..e56bd5a 100644 --- a/internal/jobs/deploy_status_reconcile.go +++ b/internal/jobs/deploy_status_reconcile.go @@ -517,7 +517,7 @@ func (w *DeployStatusReconciler) listActiveDeployments(ctx context.Context) ([]a if err != nil { return nil, fmt.Errorf("listActiveDeployments: query: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var out []activeDeployment for rows.Next() { diff --git a/internal/jobs/deployment_expirer.go b/internal/jobs/deployment_expirer.go index 7462e76..922c638 100644 --- a/internal/jobs/deployment_expirer.go +++ b/internal/jobs/deployment_expirer.go @@ -104,7 +104,7 @@ func (w *DeploymentExpirerWorker) Work(ctx context.Context, job *river.Job[Deplo if err != nil { return fmt.Errorf("DeploymentExpirerWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var candidates []deployExpirerRow for rows.Next() { @@ -119,7 +119,7 @@ func (w *DeploymentExpirerWorker) Work(ctx context.Context, job *river.Job[Deplo if err := rows.Err(); err != nil { return fmt.Errorf("DeploymentExpirerWorker: rows error: %w", err) } - rows.Close() + _ = rows.Close() if len(candidates) == 0 { // T21 P1-1 (BugBash 2026-05-20): idle-tick demoted INFO→DEBUG. diff --git a/internal/jobs/deployment_reminder.go b/internal/jobs/deployment_reminder.go index 008bab4..07ef666 100644 --- a/internal/jobs/deployment_reminder.go +++ b/internal/jobs/deployment_reminder.go @@ -227,7 +227,7 @@ func (w *DeploymentReminderWorker) Work(ctx context.Context, job *river.Job[Depl if err != nil { return fmt.Errorf("DeploymentReminderWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var candidates []deployReminderRow for rows.Next() { @@ -242,7 +242,7 @@ func (w *DeploymentReminderWorker) Work(ctx context.Context, job *river.Job[Depl if err := rows.Err(); err != nil { return fmt.Errorf("DeploymentReminderWorker: rows error: %w", err) } - rows.Close() + _ = rows.Close() // Sample TTL state for the gauge — counts apply to BOTH the // candidates-in-window set and the broader population. The gauge @@ -346,7 +346,7 @@ func (w *DeploymentReminderWorker) sampleTTLGauge(ctx context.Context) { slog.Warn("jobs.deployment_reminder.gauge_sample_failed", "error", err) return } - defer rows.Close() + defer func() { _ = rows.Close() }() seen := map[string]bool{} for rows.Next() { var policy string diff --git a/internal/jobs/email.go b/internal/jobs/email.go index 8500138..780fe6b 100644 --- a/internal/jobs/email.go +++ b/internal/jobs/email.go @@ -106,7 +106,7 @@ func (w *WeeklyDigestWorker) Work(ctx context.Context, job *river.Job[WeeklyDige if err != nil { return fmt.Errorf("weekly_digest.Work query users: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() type userTeam struct { email string @@ -180,7 +180,7 @@ func (w *WeeklyDigestWorker) buildResourceDigestCounts(ctx context.Context, team if err != nil { return nil, fmt.Errorf("buildResourceDigestCounts query: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var stats []DigestResourceCount for rows.Next() { diff --git a/internal/jobs/entitlement_reconciler.go b/internal/jobs/entitlement_reconciler.go index 21e2fcf..7c3d133 100644 --- a/internal/jobs/entitlement_reconciler.go +++ b/internal/jobs/entitlement_reconciler.go @@ -341,7 +341,7 @@ func (w *EntitlementReconcilerWorker) Work(ctx context.Context, job *river.Job[E if err != nil { return fmt.Errorf("EntitlementReconcilerWorker: postgres query failed: %w", err) } - defer pgRows.Close() + defer func() { _ = pgRows.Close() }() var pgCandidates []entitlementCandidate for pgRows.Next() { @@ -358,7 +358,7 @@ func (w *EntitlementReconcilerWorker) Work(ctx context.Context, job *river.Job[E if rowsErr := pgRows.Err(); rowsErr != nil { return fmt.Errorf("EntitlementReconcilerWorker: postgres rows error: %w", rowsErr) } - pgRows.Close() + _ = pgRows.Close() var pgScanned, pgDrifted, pgRegraded, pgFailed, pgSkippedTier int for _, c := range pgCandidates { @@ -513,7 +513,7 @@ func (w *EntitlementReconcilerWorker) Work(ctx context.Context, job *river.Job[E // River retries the whole tick, which will re-run both sweeps. return fmt.Errorf("EntitlementReconcilerWorker: redis query failed: %w", redisErr) } - defer redisRows.Close() + defer func() { _ = redisRows.Close() }() var redisCandidates []redisEntitlementCandidate for redisRows.Next() { @@ -538,7 +538,7 @@ func (w *EntitlementReconcilerWorker) Work(ctx context.Context, job *river.Job[E if rowsErr := redisRows.Err(); rowsErr != nil { return fmt.Errorf("EntitlementReconcilerWorker: redis rows error: %w", rowsErr) } - redisRows.Close() + _ = redisRows.Close() var redisChecked, redisApplied, redisSkipped, redisFailed, redisSkippedTier int for _, c := range redisCandidates { @@ -711,7 +711,7 @@ func (w *EntitlementReconcilerWorker) sweepMongoEntitlements(ctx context.Context ) return 0, 0, 0, 0, 0 } - defer mongoRows.Close() + defer func() { _ = mongoRows.Close() }() var mongoCandidates []mongoEntitlementCandidate for mongoRows.Next() { @@ -732,7 +732,7 @@ func (w *EntitlementReconcilerWorker) sweepMongoEntitlements(ctx context.Context slog.Warn("jobs.entitlement_reconciler.mongo.rows_failed", "error", rowsErr) return checked, applied, skipped, failed, skippedTier } - mongoRows.Close() + _ = mongoRows.Close() for _, c := range mongoCandidates { checked++ diff --git a/internal/jobs/event_email_forwarder.go b/internal/jobs/event_email_forwarder.go index 7697031..c689725 100644 --- a/internal/jobs/event_email_forwarder.go +++ b/internal/jobs/event_email_forwarder.go @@ -898,7 +898,10 @@ batchLoop: // on a cursor reset, which is harmless. classification is // permanent_drop so a support grep against the 059 columns // finds these alongside the F4 missing_renderer drops. - if claimed, ledgerErr := w.ledger.markSent(ctx, ledgerClaim{ + // A non-nil error is logged; a successful claim and an + // already-claimed (claimed==false) outcome are both benign, so + // the boolean return is intentionally discarded. + if _, ledgerErr := w.ledger.markSent(ctx, ledgerClaim{ AuditID: row.ID, Provider: w.provider.Name(), ProviderID: evt.IdempotencyKey, @@ -913,8 +916,6 @@ batchLoop: "error", ledgerErr, "note", "terminal class but ledger claim failed — cursor still advances; cursor-reset may re-attempt this row", ) - } else if !claimed { - // Already claimed by a prior attempt; benign. } slog.Info("jobs.event_email_forwarder.row_skipped", "kind", row.Kind, @@ -1037,7 +1038,7 @@ func (w *EventEmailForwarderWorker) fetchBatch(ctx context.Context, c eventCurso if err != nil { return nil, fmt.Errorf("fetchBatch query: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var out []auditRow for rows.Next() { diff --git a/internal/jobs/event_email_forwarder_coverage_test.go b/internal/jobs/event_email_forwarder_coverage_test.go index f19a6a0..fa451df 100644 --- a/internal/jobs/event_email_forwarder_coverage_test.go +++ b/internal/jobs/event_email_forwarder_coverage_test.go @@ -435,9 +435,6 @@ func TestForwarder_NewEventEmailForwarderWorker_ProductionConstructorWiresProvid // ── Work() error branches not yet exercised ─────────────────────────────── -// flakyProvider returns a configurable messageId+err per call. -type flakyProvider struct{ fakeProvider } - // failingLedger captures call counts and lets the test fail any of the // three operations. Used to drive Work() branches where the ledger probe // errors or markSent fails post-2xx. diff --git a/internal/jobs/expire.go b/internal/jobs/expire.go index 79819ac..b378426 100644 --- a/internal/jobs/expire.go +++ b/internal/jobs/expire.go @@ -195,7 +195,7 @@ func (w *ExpireAnonymousWorker) Work(ctx context.Context, job *river.Job[ExpireA if err != nil { return fmt.Errorf("ExpireAnonymousWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var candidates []toExpire for rows.Next() { @@ -209,7 +209,7 @@ func (w *ExpireAnonymousWorker) Work(ctx context.Context, job *river.Job[ExpireA if err := rows.Err(); err != nil { return fmt.Errorf("ExpireAnonymousWorker: rows error: %w", err) } - rows.Close() + _ = rows.Close() if len(candidates) == 0 { return nil diff --git a/internal/jobs/expire_imminent.go b/internal/jobs/expire_imminent.go index 6e7bb53..53de29d 100644 --- a/internal/jobs/expire_imminent.go +++ b/internal/jobs/expire_imminent.go @@ -181,7 +181,7 @@ func (w *ExpireImminentWorker) Work(ctx context.Context, job *river.Job[ExpireIm if err != nil { return fmt.Errorf("ExpireImminentWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var candidates []expireImminentRow for rows.Next() { @@ -199,7 +199,7 @@ func (w *ExpireImminentWorker) Work(ctx context.Context, job *river.Job[ExpireIm if err := rows.Err(); err != nil { return fmt.Errorf("ExpireImminentWorker: rows error: %w", err) } - rows.Close() + _ = rows.Close() if len(candidates) == 0 { // #146 (BugBash 2026-05-20 idle-tick noise pass): 10-min tick = diff --git a/internal/jobs/expire_stacks.go b/internal/jobs/expire_stacks.go index 3a2d88f..9368039 100644 --- a/internal/jobs/expire_stacks.go +++ b/internal/jobs/expire_stacks.go @@ -94,7 +94,7 @@ func deleteK8sNamespace(ctx context.Context, client *http.Client, namespace, nsP if err != nil { return fmt.Errorf("deleteK8sNamespace: DELETE %s: %w", namespace, err) } - defer resp.Body.Close() + defer func() { _ = resp.Body.Close() }() switch resp.StatusCode { case http.StatusOK, http.StatusAccepted, http.StatusNotFound: @@ -139,7 +139,7 @@ func (w *ExpireStacksWorker) Work(ctx context.Context, job *river.Job[ExpireStac if err != nil { return fmt.Errorf("ExpireStacksWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() type expiredStack struct { id string @@ -157,7 +157,7 @@ func (w *ExpireStacksWorker) Work(ctx context.Context, job *river.Job[ExpireStac if err := rows.Err(); err != nil { return fmt.Errorf("ExpireStacksWorker: rows error: %w", err) } - rows.Close() + _ = rows.Close() var deleted int for _, s := range expired { diff --git a/internal/jobs/expire_unexported_coverage_test.go b/internal/jobs/expire_unexported_coverage_test.go index 3cf7184..84e29cf 100644 --- a/internal/jobs/expire_unexported_coverage_test.go +++ b/internal/jobs/expire_unexported_coverage_test.go @@ -30,13 +30,20 @@ import ( "github.com/google/uuid" madmin "github.com/minio/madmin-go/v3" minio "github.com/minio/minio-go/v7" + "github.com/minio/minio-go/v7/pkg/credentials" "instant.dev/worker/internal/provisioner" ) -// madminNew is a thin alias so the test reads naturally; madmin.New -// builds an admin client that signs requests against the given endpoint. -var madminNew = madmin.New +// madminNew is a thin alias so the test reads naturally; it builds an admin +// client (via madmin.NewWithOptions) that signs requests against the given +// endpoint. Keeps the historical 4-arg call shape at the test call site. +func madminNew(endpoint, accessKey, secretKey string, secure bool) (*madmin.AdminClient, error) { + return madmin.NewWithOptions(endpoint, &madmin.Options{ + Creds: credentials.NewStaticV4(accessKey, secretKey, ""), + Secure: secure, + }) +} // TestInClusterK8sClient_NotInCluster: the function checks for the // /var/run/secrets/kubernetes.io/serviceaccount/token file. Outside a diff --git a/internal/jobs/expiry_reminder.go b/internal/jobs/expiry_reminder.go index e8046e2..9fc0f8f 100644 --- a/internal/jobs/expiry_reminder.go +++ b/internal/jobs/expiry_reminder.go @@ -192,7 +192,7 @@ func (w *ExpiryReminderWorker) Work(ctx context.Context, job *river.Job[ExpiryRe if err != nil { return fmt.Errorf("ExpiryReminderWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var candidates []expiryReminderRow for rows.Next() { @@ -207,7 +207,7 @@ func (w *ExpiryReminderWorker) Work(ctx context.Context, job *river.Job[ExpiryRe if err := rows.Err(); err != nil { return fmt.Errorf("ExpiryReminderWorker: rows error: %w", err) } - rows.Close() + _ = rows.Close() if len(candidates) == 0 { // P1-1 (BugBash 2026-05-19): idle tick — demoted INFO → DEBUG. diff --git a/internal/jobs/expiry_reminder_email.go b/internal/jobs/expiry_reminder_email.go index 81a18aa..b5716f3 100644 --- a/internal/jobs/expiry_reminder_email.go +++ b/internal/jobs/expiry_reminder_email.go @@ -196,20 +196,20 @@ func renderAnonExpiryEmail(params map[string]string) (subject, html, text string // forwarder. The bug will show up in slog (downstream) — the // forwarder logs every send. htmlBuf.Reset() - htmlBuf.WriteString(fmt.Sprintf( + fmt.Fprintf(&htmlBuf, "

Your instanode %s resource expires in %s %s. Visit %s to keep it.

", view.ResourceType, view.HoursRemaining, hourWord(view.Plural), view.UpgradeURL, - )) + ) } html = htmlBuf.String() var textBuf bytes.Buffer if err := anonExpiryTextTmpl.Execute(&textBuf, view); err != nil { textBuf.Reset() - textBuf.WriteString(fmt.Sprintf( + fmt.Fprintf(&textBuf, "Your instanode %s resource expires in %s %s. Visit %s to keep it.\n", view.ResourceType, view.HoursRemaining, hourWord(view.Plural), view.UpgradeURL, - )) + ) } text = textBuf.String() diff --git a/internal/jobs/geodb.go b/internal/jobs/geodb.go index 2a40c7d..2e6dc16 100644 --- a/internal/jobs/geodb.go +++ b/internal/jobs/geodb.go @@ -160,7 +160,7 @@ func (w *RefreshGeoDBWorker) Work(ctx context.Context, job *river.Job[RefreshGeo if err != nil { return fmt.Errorf("RefreshGeoDBWorker: download failed: %w", err) } - defer resp.Body.Close() + defer func() { _ = resp.Body.Close() }() if resp.StatusCode != http.StatusOK { return fmt.Errorf("RefreshGeoDBWorker: unexpected status %d", resp.StatusCode) @@ -174,12 +174,12 @@ func (w *RefreshGeoDBWorker) Work(ctx context.Context, job *river.Job[RefreshGeo // extract only the *.mmdb member before the atomic rename. tmpPath := job.Args.DBPath + ".tmp" if err := extractGeoLite2MMDB(resp.Body, tmpPath); err != nil { - os.Remove(tmpPath) + _ = os.Remove(tmpPath) return fmt.Errorf("RefreshGeoDBWorker: extract failed: %w", err) } if err := os.Rename(tmpPath, job.Args.DBPath); err != nil { - os.Remove(tmpPath) + _ = os.Remove(tmpPath) return fmt.Errorf("RefreshGeoDBWorker: rename failed: %w", err) } @@ -203,7 +203,7 @@ func extractGeoLite2MMDB(r io.Reader, dstPath string) error { if err != nil { return fmt.Errorf("extractGeoLite2MMDB: gzip reader: %w", err) } - defer gz.Close() + defer func() { _ = gz.Close() }() tr := tar.NewReader(gz) for { @@ -228,7 +228,7 @@ func extractGeoLite2MMDB(r io.Reader, dstPath string) error { // Bound the copy to the declared tar-header size to avoid a // decompression-bomb writing unbounded bytes. if _, err := io.Copy(f, io.LimitReader(tr, hdr.Size)); err != nil { - f.Close() + _ = f.Close() return fmt.Errorf("extractGeoLite2MMDB: write temp file: %w", err) } if err := f.Close(); err != nil { diff --git a/internal/jobs/github_deploy_dispatcher.go b/internal/jobs/github_deploy_dispatcher.go index eef5356..f26e21b 100644 --- a/internal/jobs/github_deploy_dispatcher.go +++ b/internal/jobs/github_deploy_dispatcher.go @@ -185,7 +185,7 @@ func (w *GitHubDeployDispatcher) claimBatch(ctx context.Context) ([]pendingGitHu if err != nil { return nil, err } - defer tx.Rollback() + defer func() { _ = tx.Rollback() }() rs, err := tx.QueryContext(ctx, ` WITH claimed AS ( @@ -211,12 +211,12 @@ func (w *GitHubDeployDispatcher) claimBatch(ctx context.Context) ([]pendingGitHu for rs.Next() { var r pendingGitHubDeploy if err := rs.Scan(&r.id, &r.connectionID, &r.appUUID, &r.commitSHA, &r.attempts); err != nil { - rs.Close() + _ = rs.Close() return nil, err } out = append(out, r) } - rs.Close() + _ = rs.Close() // Enrich with parent metadata. Done one-by-one inside the same tx so // the claim and the read are atomic. @@ -272,7 +272,7 @@ func (w *GitHubDeployDispatcher) fetchTarball(ctx context.Context, url string) ( if err != nil { return nil, err } - defer resp.Body.Close() + defer func() { _ = resp.Body.Close() }() if resp.StatusCode >= 400 && resp.StatusCode < 500 { // Permanent — repo / ref doesn't exist (or is private + no auth). // Don't retry. @@ -306,7 +306,7 @@ func (w *GitHubDeployDispatcher) postRedeploy(ctx context.Context, appIDSlug str if _, err := part.Write(tarball); err != nil { return err } - mw.Close() + _ = mw.Close() url := w.apiBaseURL + "/deploy/" + appIDSlug + "/redeploy" req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, &body) @@ -338,8 +338,8 @@ func (w *GitHubDeployDispatcher) postRedeploy(ctx context.Context, appIDSlug str } return err } - defer resp.Body.Close() - io.Copy(io.Discard, resp.Body) // drain so keep-alive is happy + defer func() { _ = resp.Body.Close() }() + _, _ = io.Copy(io.Discard, resp.Body) // drain so keep-alive is happy if resp.StatusCode >= 400 && resp.StatusCode < 500 { return &permanentError{Code: resp.StatusCode, Msg: "api redeploy 4xx"} } diff --git a/internal/jobs/magic_link_reconciler.go b/internal/jobs/magic_link_reconciler.go index ec0120c..7161ade 100644 --- a/internal/jobs/magic_link_reconciler.go +++ b/internal/jobs/magic_link_reconciler.go @@ -217,7 +217,7 @@ func (w *MagicLinkReconcilerWorker) listReconcileCandidates(ctx context.Context) if err != nil { return nil, fmt.Errorf("query magic_links: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var out []magicLinkReconcileRow for rows.Next() { @@ -271,7 +271,7 @@ func (w *MagicLinkReconcilerWorker) driveResend(ctx context.Context, r magicLink ) return reconcileOutcomeSkipped } - defer resp.Body.Close() + defer func() { _ = resp.Body.Close() }() if resp.StatusCode < 200 || resp.StatusCode >= 300 { bodyBytes, _ := io.ReadAll(io.LimitReader(resp.Body, 1024)) diff --git a/internal/jobs/orphan_sweep_canceler_test.go b/internal/jobs/orphan_sweep_canceler_test.go index 3dfdd2f..4335cda 100644 --- a/internal/jobs/orphan_sweep_canceler_test.go +++ b/internal/jobs/orphan_sweep_canceler_test.go @@ -217,7 +217,7 @@ func reachRazorpayHTTPTimeout(t *testing.T, c jobs.OrphanSubscriptionCanceler) t t.Helper() v := reflect.ValueOf(c) // canceler is *razorpayOrphanCanceler — dereference. - for v.Kind() == reflect.Ptr || v.Kind() == reflect.Interface { + for v.Kind() == reflect.Pointer || v.Kind() == reflect.Interface { v = v.Elem() } clientFld := v.FieldByName("client") @@ -226,7 +226,7 @@ func reachRazorpayHTTPTimeout(t *testing.T, c jobs.OrphanSubscriptionCanceler) t } // clientFld is razorpaySubCancelClient (interface). adapter := clientFld.Elem() - for adapter.Kind() == reflect.Ptr { + for adapter.Kind() == reflect.Pointer { adapter = adapter.Elem() } cFld := adapter.FieldByName("c") @@ -235,7 +235,7 @@ func reachRazorpayHTTPTimeout(t *testing.T, c jobs.OrphanSubscriptionCanceler) t } // cFld is *razorpay.Client; deref then read Request.HTTPClient.Timeout. sdkClient := cFld - for sdkClient.Kind() == reflect.Ptr { + for sdkClient.Kind() == reflect.Pointer { sdkClient = sdkClient.Elem() } reqFld := sdkClient.FieldByName("Request") @@ -243,7 +243,7 @@ func reachRazorpayHTTPTimeout(t *testing.T, c jobs.OrphanSubscriptionCanceler) t t.Fatalf("razorpay.Client has no `Request` field. type=%s", sdkClient.Type()) } req := reqFld - for req.Kind() == reflect.Ptr { + for req.Kind() == reflect.Pointer { req = req.Elem() } httpFld := req.FieldByName("HTTPClient") @@ -251,7 +251,7 @@ func reachRazorpayHTTPTimeout(t *testing.T, c jobs.OrphanSubscriptionCanceler) t t.Fatalf("requests.Request has no `HTTPClient`. type=%s", req.Type()) } httpClient := httpFld - for httpClient.Kind() == reflect.Ptr { + for httpClient.Kind() == reflect.Pointer { httpClient = httpClient.Elem() } timeoutFld := httpClient.FieldByName("Timeout") diff --git a/internal/jobs/orphan_sweep_reconciler.go b/internal/jobs/orphan_sweep_reconciler.go index 3c21030..679e430 100644 --- a/internal/jobs/orphan_sweep_reconciler.go +++ b/internal/jobs/orphan_sweep_reconciler.go @@ -386,12 +386,12 @@ func (w *OrphanSweepReconciler) sweepStuckPendingTeams(ctx context.Context) (fin for rows.Next() { var c teamPendingDeletion if scanErr := rows.Scan(&c.teamID, &c.deletionRequestedAt); scanErr != nil { - rows.Close() + _ = rows.Close() return 0, 0, fmt.Errorf("scan stuck-pending team: %w", scanErr) } candidates = append(candidates, c) } - rows.Close() + _ = rows.Close() if err := rows.Err(); err != nil { return 0, 0, err } @@ -449,12 +449,12 @@ func (w *OrphanSweepReconciler) sweepOrphanedSubscriptions(ctx context.Context) for rows.Next() { var o orphanedSub if scanErr := rows.Scan(&o.teamID, &o.subscriptionID); scanErr != nil { - rows.Close() + _ = rows.Close() return 0, 0, fmt.Errorf("scan orphaned sub: %w", scanErr) } orphans = append(orphans, o) } - rows.Close() + _ = rows.Close() if err := rows.Err(); err != nil { return 0, 0, err } @@ -697,7 +697,7 @@ func (w *OrphanSweepReconciler) fetchDeployRowsByAppID(ctx context.Context) (map if err != nil { return nil, err } - defer rows.Close() + defer func() { _ = rows.Close() }() out := make(map[string]deployRowSnapshot) for rows.Next() { var ( @@ -816,7 +816,7 @@ func (w *OrphanSweepReconciler) fetchLiveResourceTokens(ctx context.Context) (ma if err != nil { return nil, err } - defer rows.Close() + defer func() { _ = rows.Close() }() out := make(map[string]bool) for rows.Next() { var token string @@ -923,7 +923,7 @@ func (w *OrphanSweepReconciler) fetchLiveStackIDs(ctx context.Context) (map[stri if err != nil { return nil, err } - defer rows.Close() + defer func() { _ = rows.Close() }() out := make(map[string]bool) for rows.Next() { var id string @@ -1083,7 +1083,7 @@ func (w *OrphanSweepReconciler) fetchStuckBuildCandidates(ctx context.Context) ( if err != nil { return nil, err } - defer rows.Close() + defer func() { _ = rows.Close() }() var out []stuckBuildCandidate for rows.Next() { var c stuckBuildCandidate diff --git a/internal/jobs/orphan_sweep_reconciler_test.go b/internal/jobs/orphan_sweep_reconciler_test.go index ddd97e6..1f41bd0 100644 --- a/internal/jobs/orphan_sweep_reconciler_test.go +++ b/internal/jobs/orphan_sweep_reconciler_test.go @@ -224,25 +224,6 @@ func (f *fakePodStateProvider) ListPodWaitingReasons(_ context.Context, namespac return out, nil } -// expectEmptyPass4 queues the sqlmock expectation for PASS 4's live-token -// query returning nothing — used when the fake reports customer namespaces -// but a test only cares about the earlier passes, OR when the fake reports -// no customer namespaces (PASS 4 short-circuits before the query, so callers -// with an empty customer-namespace set must NOT queue this). -func expectEmptyPass4(mock sqlmock.Sqlmock) { - mock.ExpectQuery(`SELECT DISTINCT token::text\s+FROM resources`). - WillReturnRows(sqlmock.NewRows([]string{"token"})) -} - -// expectEmptyPass5 queues the sqlmock expectation for PASS 5's live-stack-id -// query returning nothing — same shape as expectEmptyPass4 but for stacks. -// Only call this when the fake reports a non-empty stackNamespaces slice; -// PASS 5 short-circuits before the query on an empty namespaces list. -func expectEmptyPass5(mock sqlmock.Sqlmock) { - mock.ExpectQuery(`SELECT id::text FROM stacks`). - WillReturnRows(sqlmock.NewRows([]string{"id"})) -} - // captureBytesArgOSR captures a JSONB column's bytes for audit assertions. type captureBytesArgOSR struct{ out *[]byte } diff --git a/internal/jobs/payment_grace_reminder.go b/internal/jobs/payment_grace_reminder.go index 421e5d0..8729fa0 100644 --- a/internal/jobs/payment_grace_reminder.go +++ b/internal/jobs/payment_grace_reminder.go @@ -123,7 +123,7 @@ func (w *PaymentGraceReminderWorker) Work(ctx context.Context, job *river.Job[Pa if err != nil { return fmt.Errorf("PaymentGraceReminderWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var candidates []paymentGraceReminderRow for rows.Next() { @@ -137,7 +137,7 @@ func (w *PaymentGraceReminderWorker) Work(ctx context.Context, job *river.Job[Pa if err := rows.Err(); err != nil { return fmt.Errorf("PaymentGraceReminderWorker: rows error: %w", err) } - rows.Close() + _ = rows.Close() if len(candidates) == 0 { // P1-1 (BugBash 2026-05-19): idle tick — demoted INFO → DEBUG. diff --git a/internal/jobs/payment_grace_terminator.go b/internal/jobs/payment_grace_terminator.go index 20a0ad4..9819ed4 100644 --- a/internal/jobs/payment_grace_terminator.go +++ b/internal/jobs/payment_grace_terminator.go @@ -158,7 +158,7 @@ func (w *PaymentGraceTerminatorWorker) Work(ctx context.Context, job *river.Job[ if err != nil { return fmt.Errorf("PaymentGraceTerminatorWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var candidates []paymentGraceTerminatorRow for rows.Next() { @@ -172,7 +172,7 @@ func (w *PaymentGraceTerminatorWorker) Work(ctx context.Context, job *river.Job[ if err := rows.Err(); err != nil { return fmt.Errorf("PaymentGraceTerminatorWorker: rows error: %w", err) } - rows.Close() + _ = rows.Close() if len(candidates) == 0 { // P1-1 (BugBash 2026-05-19): idle tick — demoted INFO → DEBUG. @@ -291,7 +291,7 @@ func (w *PaymentGraceTerminatorWorker) terminate(ctx context.Context, teamID uui } return fmt.Errorf("api request: %w", doErr) } - defer resp.Body.Close() + defer func() { _ = resp.Body.Close() }() if resp.StatusCode < 200 || resp.StatusCode >= 300 { body, _ := io.ReadAll(io.LimitReader(resp.Body, 1024)) return fmt.Errorf("api status %d: %s", resp.StatusCode, strings.TrimSpace(string(body))) diff --git a/internal/jobs/pending_deletion_expirer.go b/internal/jobs/pending_deletion_expirer.go index 4e64122..ea984e8 100644 --- a/internal/jobs/pending_deletion_expirer.go +++ b/internal/jobs/pending_deletion_expirer.go @@ -116,7 +116,7 @@ func (w *PendingDeletionExpirerWorker) Work(ctx context.Context, job *river.Job[ if err != nil { return fmt.Errorf("pending_deletion_expirer: sweep failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var expired []expiredPendingDeletionRow for rows.Next() { diff --git a/internal/jobs/platform_db_backup.go b/internal/jobs/platform_db_backup.go index 865e80e..97df79a 100644 --- a/internal/jobs/platform_db_backup.go +++ b/internal/jobs/platform_db_backup.go @@ -323,7 +323,7 @@ func (w *PlatformDBBackupWorker) Work(ctx context.Context, job *river.Job[Platfo if err != nil { return fmt.Errorf("PlatformDBBackupWorker: acquire conn: %w", err) } - defer conn.Close() + defer func() { _ = conn.Close() }() var locked bool if err := conn.QueryRowContext(ctx, `SELECT pg_try_advisory_lock($1)`, platformDBBackupLockKey).Scan(&locked); err != nil { diff --git a/internal/jobs/propagation_runner.go b/internal/jobs/propagation_runner.go index 3b502bc..552e7ee 100644 --- a/internal/jobs/propagation_runner.go +++ b/internal/jobs/propagation_runner.go @@ -543,7 +543,7 @@ func (w *PropagationRunnerWorker) pickEligible(ctx context.Context) ([]propagati if err != nil { return nil, fmt.Errorf("select eligible: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var out []propagationRow for rows.Next() { @@ -557,7 +557,7 @@ func (w *PropagationRunnerWorker) pickEligible(ctx context.Context) ([]propagati if rowsErr := rows.Err(); rowsErr != nil { return nil, fmt.Errorf("rows iteration: %w", rowsErr) } - rows.Close() + _ = rows.Close() // D22-P3 lease bump (2026-05-21). Push next_attempt_at on the // just-picked, still-locked rows so a pod crash between COMMIT below @@ -1077,7 +1077,7 @@ func handleTierElevation(ctx context.Context, db *sql.DB, regrader propagationRe if err != nil { return fmt.Errorf("select team resources: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() type res struct { id uuid.UUID @@ -1099,7 +1099,7 @@ func handleTierElevation(ctx context.Context, db *sql.DB, regrader propagationRe if rowsErr := rows.Err(); rowsErr != nil { return fmt.Errorf("iterate resources: %w", rowsErr) } - rows.Close() + _ = rows.Close() // Empty team is success — there is nothing to regrade. The customer // upgraded but hasn't provisioned anything yet; future provisions diff --git a/internal/jobs/provisioner_reconciler.go b/internal/jobs/provisioner_reconciler.go index 33ddcc4..ef6d641 100644 --- a/internal/jobs/provisioner_reconciler.go +++ b/internal/jobs/provisioner_reconciler.go @@ -155,7 +155,7 @@ func (w *ProvisionerReconcilerWorker) Work(ctx context.Context, job *river.Job[P if err != nil { return fmt.Errorf("ProvisionerReconcilerWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var candidates []reconcilerCandidate for rows.Next() { @@ -176,7 +176,7 @@ func (w *ProvisionerReconcilerWorker) Work(ctx context.Context, job *river.Job[P if rowsErr := rows.Err(); rowsErr != nil { return fmt.Errorf("ProvisionerReconcilerWorker: rows error: %w", rowsErr) } - rows.Close() + _ = rows.Close() if len(candidates) == 0 { // T21 P1-1 (BugBash 2026-05-20): idle-tick demoted INFO→DEBUG. diff --git a/internal/jobs/quota.go b/internal/jobs/quota.go index ab036e8..bbe5618 100644 --- a/internal/jobs/quota.go +++ b/internal/jobs/quota.go @@ -222,7 +222,7 @@ func (w *EnforceStorageQuotaWorker) runRedisEvictionLoop(ctx context.Context) (i if err != nil { return 0, fmt.Errorf("EnforceStorageQuotaWorker.redisEvictionLoop: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() checked, enforced := 0, 0 @@ -325,7 +325,7 @@ func (w *EnforceStorageQuotaWorker) runSuspendLoop(ctx context.Context) ([]strin if err != nil { return nil, fmt.Errorf("EnforceStorageQuotaWorker.suspendLoop: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() checked := 0 suspendedIDs := make([]string, 0) @@ -475,7 +475,7 @@ func (w *EnforceStorageQuotaWorker) runUnsuspendLoop(ctx context.Context, skipID if err != nil { return 0, fmt.Errorf("EnforceStorageQuotaWorker.unsuspendLoop: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() unsuspended := 0 diff --git a/internal/jobs/quota_infra.go b/internal/jobs/quota_infra.go index 9756455..323c1fb 100644 --- a/internal/jobs/quota_infra.go +++ b/internal/jobs/quota_infra.go @@ -152,7 +152,7 @@ func (r *directResourceRevoker) revokePostgres(ctx context.Context, token string slog.Warn("quota_infra.revokePostgres: open failed (fail-open)", "token", logsafe.Token(token), "error", err) return nil } - defer conn.Close() + defer func() { _ = conn.Close() }() if _, err := conn.ExecContext(ctx, fmt.Sprintf(`REVOKE CONNECT ON DATABASE %q FROM %q`, dbName, username)); err != nil { @@ -191,7 +191,7 @@ func (r *directResourceRevoker) grantPostgres(ctx context.Context, token string) slog.Warn("quota_infra.grantPostgres: open failed (fail-open)", "token", logsafe.Token(token), "error", err) return nil } - defer conn.Close() + defer func() { _ = conn.Close() }() if _, err := conn.ExecContext(ctx, fmt.Sprintf(`GRANT CONNECT ON DATABASE %q TO %q`, dbName, username)); err != nil { @@ -316,7 +316,7 @@ func setCustomerRedisACL(ctx context.Context, adminURL, username string, enable return nil } client := goredis.NewClient(opts) - defer client.Close() + defer func() { _ = client.Close() }() state := "off" if enable { state = "on" @@ -408,7 +408,8 @@ func validateSuspendIdent(s string) error { return fmt.Errorf("empty identifier") } for _, ch := range s { - if !((ch >= 'a' && ch <= 'z') || (ch >= '0' && ch <= '9') || ch == '_' || ch == '-') { + allowed := (ch >= 'a' && ch <= 'z') || (ch >= '0' && ch <= '9') || ch == '_' || ch == '-' + if !allowed { return fmt.Errorf("unsafe identifier %q", s) } } diff --git a/internal/jobs/quota_redis_eviction.go b/internal/jobs/quota_redis_eviction.go index 740fa79..b157d71 100644 --- a/internal/jobs/quota_redis_eviction.go +++ b/internal/jobs/quota_redis_eviction.go @@ -183,7 +183,7 @@ func (e *directRedisEvictor) EvictTenantToCap(ctx context.Context, token string, return 0, 0, fmt.Errorf("EvictTenantToCap: parse admin URL: %w", err) } client := goredis.NewClient(opts) - defer client.Close() + defer func() { _ = client.Close() }() return evictTenantToCap(ctx, client, token, limitBytes) } diff --git a/internal/jobs/quota_wall_nudge.go b/internal/jobs/quota_wall_nudge.go index 3b4d14e..e307d90 100644 --- a/internal/jobs/quota_wall_nudge.go +++ b/internal/jobs/quota_wall_nudge.go @@ -121,7 +121,7 @@ func (w *QuotaWallNudgeWorker) Work(ctx context.Context, job *river.Job[QuotaWal if err != nil { return fmt.Errorf("QuotaWallNudgeWorker: list teams: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() scanned, nudged, skipped := 0, 0, 0 diff --git a/internal/jobs/real_prober.go b/internal/jobs/real_prober.go index 9e78b5e..e94fc37 100644 --- a/internal/jobs/real_prober.go +++ b/internal/jobs/real_prober.go @@ -209,7 +209,7 @@ func (p *realProber) probePostgres(ctx context.Context, connURL string) (ProbeOu if err != nil { return ProbeUnreachable, fmt.Errorf("postgres: sql.Open: %w", err) } - defer db.Close() + defer func() { _ = db.Close() }() // Tighten dial timeout so a black-holed host fails fast rather than // hanging the goroutine for the driver default. @@ -241,7 +241,7 @@ func (p *realProber) probeRedis(ctx context.Context, connURL string) (ProbeOutco opts.MinIdleConns = 0 client := redis.NewClient(opts) - defer client.Close() + defer func() { _ = client.Close() }() if err := client.Ping(ctx).Err(); err != nil { return ProbeUnreachable, fmt.Errorf("redis: PING: %w", err) @@ -296,7 +296,7 @@ func (p *realProber) probeStorage(ctx context.Context, connURL string) (ProbeOut if err != nil { return ProbeUnreachable, fmt.Errorf("storage: HEAD %s: %w", target, err) } - defer resp.Body.Close() + defer func() { _ = resp.Body.Close() }() // Any HTTP status is fine — we just needed a TCP+TLS handshake. return ProbeReachable, nil } @@ -352,7 +352,7 @@ func (p *realProber) probeQueue(ctx context.Context, connURL string) (ProbeOutco if err != nil { return ProbeUnreachable, fmt.Errorf("queue: GET %s: %w", target, err) } - defer resp.Body.Close() + defer func() { _ = resp.Body.Close() }() if resp.StatusCode != http.StatusOK { return ProbeUnreachable, fmt.Errorf("queue: NATS unhealthy (HTTP %d)", resp.StatusCode) } diff --git a/internal/jobs/resource_heartbeat.go b/internal/jobs/resource_heartbeat.go index 204fba1..873a2f6 100644 --- a/internal/jobs/resource_heartbeat.go +++ b/internal/jobs/resource_heartbeat.go @@ -157,7 +157,7 @@ func (w *ResourceHeartbeatWorker) Work(ctx context.Context, job *river.Job[Resou if err != nil { return fmt.Errorf("ResourceHeartbeatWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() var candidates []heartbeatCandidate for rows.Next() { @@ -178,7 +178,7 @@ func (w *ResourceHeartbeatWorker) Work(ctx context.Context, job *river.Job[Resou if rowsErr := rows.Err(); rowsErr != nil { return fmt.Errorf("ResourceHeartbeatWorker: rows error: %w", rowsErr) } - rows.Close() + _ = rows.Close() if len(candidates) == 0 { // Wave 3 / Worker T21 P1-1 follow-up (#146): demote idle-tick INFO → @@ -391,7 +391,7 @@ func (w *ResourceHeartbeatWorker) sampleDegradedGauge(ctx context.Context) { slog.Warn("jobs.resource_heartbeat.sample_gauge_failed", "error", err) return } - defer rows.Close() + defer func() { _ = rows.Close() }() seen := map[string]bool{} for rows.Next() { var rt string diff --git a/internal/jobs/storage.go b/internal/jobs/storage.go index 748eeb7..1bdf603 100644 --- a/internal/jobs/storage.go +++ b/internal/jobs/storage.go @@ -85,7 +85,7 @@ func (w *UpdateStorageBytesWorker) Work(ctx context.Context, job *river.Job[Upda if err != nil { return fmt.Errorf("UpdateStorageBytesWorker: query failed: %w", err) } - defer rows.Close() + defer func() { _ = rows.Close() }() updated := 0 for rows.Next() { diff --git a/internal/jobs/team_deletion_executor.go b/internal/jobs/team_deletion_executor.go index c78130f..99ecfb9 100644 --- a/internal/jobs/team_deletion_executor.go +++ b/internal/jobs/team_deletion_executor.go @@ -311,7 +311,7 @@ func (w *TeamDeletionExecutorWorker) fetchCandidates(ctx context.Context) ([]tea if err != nil { return nil, err } - defer rows.Close() + defer func() { _ = rows.Close() }() var out []teamPendingDeletion for rows.Next() { @@ -525,7 +525,7 @@ func (w *TeamDeletionExecutorWorker) fetchTeamDeployAppIDs(ctx context.Context, if err != nil { return nil, err } - defer rows.Close() + defer func() { _ = rows.Close() }() var out []string for rows.Next() { @@ -551,7 +551,7 @@ func (w *TeamDeletionExecutorWorker) fetchTeamResources(ctx context.Context, tea if err != nil { return nil, err } - defer rows.Close() + defer func() { _ = rows.Close() }() var out []pendingResource for rows.Next() { diff --git a/internal/jobs/uptime_prober.go b/internal/jobs/uptime_prober.go index 808a3e2..615194f 100644 --- a/internal/jobs/uptime_prober.go +++ b/internal/jobs/uptime_prober.go @@ -350,7 +350,7 @@ func (w *UptimeProberWorker) httpHEAD(ctx context.Context, url string, useHead b if err != nil { return false, nil } - defer resp.Body.Close() + defer func() { _ = resp.Body.Close() }() // 5xx → unhealthy; everything else → healthy (TLS+routing+upstream // returned a response). 4xx is healthy because the dispatcher / // ingress is functional even if the route doesn't exist on the diff --git a/internal/jobs/workers.go b/internal/jobs/workers.go index f9430e6..47031a7 100644 --- a/internal/jobs/workers.go +++ b/internal/jobs/workers.go @@ -8,6 +8,7 @@ import ( "time" madmin "github.com/minio/madmin-go/v3" + "github.com/minio/minio-go/v7/pkg/credentials" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" "github.com/newrelic/go-agent/v3/newrelic" @@ -293,6 +294,24 @@ func (dailyAt2UTCSchedule) Next(t time.Time) time.Time { // // Pass nil to fall back to the legacy hardcoded 7-day retention default // — retentionDaysForTier WARNs in that case. + +// newMinioAdminClient builds the MinIO admin client used for storage IAM +// cleanup. It returns (nil, nil) when no MINIO_* endpoint is configured — +// the common case with the OBJECT_STORE_* shared-key backend (DO Spaces / +// AWS / GCS / R2), where no per-customer IAM user was ever created. Only +// the legacy self-hosted MinIO backend sets MinioEndpoint. Extracted from +// StartWorkers so the credential-construction path is unit-testable without +// booting the full worker pool. +func newMinioAdminClient(cfg *config.Config) (*madmin.AdminClient, error) { + if cfg.MinioEndpoint == "" { + return nil, nil + } + return madmin.NewWithOptions(cfg.MinioEndpoint, &madmin.Options{ + Creds: credentials.NewStaticV4(cfg.MinioRootUser, cfg.MinioRootPassword, ""), + Secure: false, + }) +} + func StartWorkers(ctx context.Context, db *sql.DB, rdb *redis.Client, cfg *config.Config, provClient *provisioner.Client, planRegistry PlanRegistry, backupPlans BackupPlanRegistry, deployStatusK8s deployStatusK8sProvider, deployAutopsyK8s deployAutopsyK8sProvider, nrApp *newrelic.Application) *Workers { // rdb is used by LoopsEventForwarderWorker (cursor storage). Other // workers access redis indirectly via the platform DB. @@ -351,13 +370,9 @@ func StartWorkers(ctx context.Context, db *sql.DB, rdb *redis.Client, cfg *confi // release per-IAM-user resources (i.e. self-hosted MinIO backend). With // the OBJECT_STORE_* shared-key backend (DO Spaces / AWS / GCS / R2) this // stays nil because no per-customer IAM was created in the first place. - var minioClient *madmin.AdminClient - if cfg.MinioEndpoint != "" { - if mc, err := madmin.New(cfg.MinioEndpoint, cfg.MinioRootUser, cfg.MinioRootPassword, false); err != nil { - slog.Warn("jobs.workers.minio_client_init_failed", "error", err) - } else { - minioClient = mc - } + minioClient, err := newMinioAdminClient(cfg) + if err != nil { + slog.Warn("jobs.workers.minio_client_init_failed", "error", err) } // Build the storage_bytes scanner — provider-agnostic, uses plain S3 API diff --git a/internal/jobs/workers_lifecycle_test.go b/internal/jobs/workers_lifecycle_test.go index 441deaa..544a746 100644 --- a/internal/jobs/workers_lifecycle_test.go +++ b/internal/jobs/workers_lifecycle_test.go @@ -198,9 +198,15 @@ func TestWorkers_Stop_NoCancelBeforeStopReturns(t *testing.T) { // a seam (or we wire in a `riverclient.Stoppable` interface). // // Field-shape sanity-check (no nil deref): a real value retains - // the *river.Client[pgx.Tx] and context.CancelFunc field types. + // the *river.Client[pgx.Tx] and context.CancelFunc field types. The + // explicitly-typed parameters of these no-op sinks are the assertion — + // the call fails to compile if a field's type ever drifts. w := Workers{} - var _ *river.Client[pgx.Tx] = w.client - var _ context.CancelFunc = w.cancel + pinRiverClient(w.client) + pinCancelFunc(w.cancel) _ = w } + +// pinRiverClient / pinCancelFunc pin the Workers field types at compile time. +func pinRiverClient(*river.Client[pgx.Tx]) {} +func pinCancelFunc(context.CancelFunc) {} diff --git a/internal/jobs/workers_schedule_coverage_test.go b/internal/jobs/workers_schedule_coverage_test.go index d251bce..a606dc0 100644 --- a/internal/jobs/workers_schedule_coverage_test.go +++ b/internal/jobs/workers_schedule_coverage_test.go @@ -9,6 +9,8 @@ package jobs import ( "context" "database/sql" + "fmt" + "net/url" "os" "testing" "time" @@ -109,6 +111,33 @@ func TestEntitlementRegraderAdapter_ErrorArm(t *testing.T) { } } +// newMinioAdminClient: empty endpoint short-circuits to (nil, nil); a +// configured endpoint builds an admin client without contacting it +// (madmin.NewWithOptions only parses the endpoint URL, no network I/O). +func TestNewMinioAdminClient_EmptyEndpoint_ReturnsNil(t *testing.T) { + mc, err := newMinioAdminClient(&config.Config{MinioEndpoint: ""}) + if err != nil { + t.Fatalf("empty endpoint returned err: %v", err) + } + if mc != nil { + t.Errorf("empty endpoint returned non-nil client %v, want nil", mc) + } +} + +func TestNewMinioAdminClient_ConfiguredEndpoint_BuildsClient(t *testing.T) { + mc, err := newMinioAdminClient(&config.Config{ + MinioEndpoint: "minio.example.com:9000", + MinioRootUser: "minioadmin", + MinioRootPassword: "minioadmin", + }) + if err != nil { + t.Fatalf("configured endpoint returned err: %v", err) + } + if mc == nil { + t.Error("configured endpoint returned nil client, want non-nil") + } +} + // StartWorkers early-return guards. We don't start a real worker pool here — // only the synchronous early-exit branches. func TestStartWorkers_BadDatabaseURL_ReturnsEmpty(t *testing.T) { @@ -179,3 +208,63 @@ func TestStartWorkers_FullBoot(t *testing.T) { t.Skip("StartWorkers did not reach started=true (River migrate/start gated by DB perms) — body still executed for coverage") } } + +// TestStartWorkers_BootInIsolatedDB drives StartWorkers past the +// email-provider init into the MinIO-client + worker-registration body using +// a throwaway database created from the root test-pg connection. Unlike +// TestStartWorkers_FullBoot (which needs the dedicated TEST_WORKER_STARTUP_DSN +// and skips in CI), this self-provisions its own DB so River's schema +// migrations never touch the shared TEST_DATABASE_URL schema — giving CI +// coverage of the post-email-init body (incl. the newMinioAdminClient call +// site) without the cross-test perturbation risk that motivated the separate +// DSN. Skips cleanly when the root pg container is unreachable. +func TestStartWorkers_BootInIsolatedDB(t *testing.T) { + root, err := sql.Open("postgres", pgTestDSN()) + if err != nil { + t.Skipf("postgres open: %v", err) + } + defer root.Close() + root.SetConnMaxLifetime(5 * time.Second) + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := root.PingContext(ctx); err != nil { + t.Skipf("postgres ping failed (docker test-pg not reachable): %v", err) + } + + dbName := fmt.Sprintf("worker_boot_%d", time.Now().UnixNano()%1_000_000) + if _, err := root.ExecContext(ctx, `CREATE DATABASE `+quoteIdent(dbName)); err != nil { + t.Skipf("CREATE DATABASE: %v", err) + } + t.Cleanup(func() { _, _ = root.Exec(`DROP DATABASE IF EXISTS ` + quoteIdent(dbName)) }) + + // Rewrite the path of pgTestDSN to point at the throwaway DB. + bootDSN := dsnWithDB(pgTestDSN(), dbName) + + cfg := &config.Config{ + DatabaseURL: bootDSN, + Environment: "development", + // EmailProvider empty → NoopProvider; all optional deps unset → + // every worker wired fail-open. A deliberately-malformed MinioEndpoint + // makes newMinioAdminClient return an error so the call site's + // fail-open WARN branch is exercised (minioClient stays nil and boot + // continues — proving the cleanup-client init never blocks startup). + MinioEndpoint: "::::bad", + } + + w := StartWorkers(context.Background(), nil, nil, cfg, nil, nil, nil, nil, nil, nil) + if w == nil { + t.Fatal("StartWorkers returned nil") + } + defer w.Stop() +} + +// dsnWithDB replaces the database name (the URL path) in a postgres DSN. +func dsnWithDB(dsn, db string) string { + u, err := url.Parse(dsn) + if err != nil { + return dsn + } + u.Path = "/" + db + return u.String() +} diff --git a/main.go b/main.go index 6c2bf27..f132651 100644 --- a/main.go +++ b/main.go @@ -277,7 +277,7 @@ func run(ctx context.Context, d deps) int { cfg := d.loadConfig() // panics on missing required env vars in production database := d.connectPostgres(cfg.DatabaseURL) - defer database.Close() + defer func() { _ = database.Close() }() // Pool-saturation observability (Wave-3 chaos verify, 2026-05-21): // pushes *sql.DB.Stats onto instant_pg_pool_* gauges so operators can @@ -287,7 +287,7 @@ func run(ctx context.Context, d deps) int { d.startPoolStats(poolStatsCtx, database, "platform_db") rdb := d.connectRedis(cfg.RedisURL) - defer rdb.Close() + defer func() { _ = rdb.Close() }() workers := d.startWorkers(ctx, database, rdb, cfg) defer workers.Stop()