diff --git a/.github/workflows/agent-ci.yml b/.github/workflows/agent-ci.yml index c9a8d9d9..6f843117 100644 --- a/.github/workflows/agent-ci.yml +++ b/.github/workflows/agent-ci.yml @@ -43,7 +43,7 @@ jobs: run: go vet ./... - name: Install staticcheck - run: GOBIN=$PWD/.bin go install honnef.co/go/tools/cmd/staticcheck@latest + run: GOBIN=$PWD/.bin go install honnef.co/go/tools/cmd/staticcheck@v0.7.0 - name: Staticcheck run: ./.bin/staticcheck ./... diff --git a/.github/workflows/cli-ci.yml b/.github/workflows/cli-ci.yml index 3bd81025..1caf8caa 100644 --- a/.github/workflows/cli-ci.yml +++ b/.github/workflows/cli-ci.yml @@ -43,7 +43,7 @@ jobs: run: go vet ./... - name: Install staticcheck - run: GOBIN=$PWD/.bin go install honnef.co/go/tools/cmd/staticcheck@latest + run: GOBIN=$PWD/.bin go install honnef.co/go/tools/cmd/staticcheck@v0.7.0 - name: Staticcheck run: ./.bin/staticcheck ./... diff --git a/.github/workflows/updater-ci.yml b/.github/workflows/updater-ci.yml index 44d3fe51..6146f048 100644 --- a/.github/workflows/updater-ci.yml +++ b/.github/workflows/updater-ci.yml @@ -45,7 +45,7 @@ jobs: run: go vet ./... - name: Install staticcheck - run: GOBIN=$PWD/.bin go install honnef.co/go/tools/cmd/staticcheck@latest + run: GOBIN=$PWD/.bin go install honnef.co/go/tools/cmd/staticcheck@v0.7.0 - name: Staticcheck run: ./.bin/staticcheck ./... diff --git a/agent/internal/traefik/reload.go b/agent/internal/traefik/reload.go index fdee5bf0..21a1438a 100644 --- a/agent/internal/traefik/reload.go +++ b/agent/internal/traefik/reload.go @@ -15,6 +15,7 @@ import ( const ( lastReloadSuccessMetric = "traefik_config_last_reload_success" pendingReloadMarkerName = ".routing-reload-pending" + metricsReadyTimeout = 15 * time.Second ) var ( @@ -30,6 +31,22 @@ func LastSuccessfulReload() (time.Time, error) { return readLastSuccessfulReload() } +func waitForMetricsReady(timeout time.Duration) error { + deadline := time.Now().Add(timeout) + var lastErr error + for { + if _, err := LastSuccessfulReload(); err == nil { + return nil + } else { + lastErr = err + } + if time.Now().After(deadline) { + return fmt.Errorf("traefik metrics did not become ready within %s: %w", timeout, lastErr) + } + time.Sleep(reloadPollInterval) + } +} + func fetchLastSuccessfulReload() (time.Time, error) { response, err := metricsHTTPClient.Get(traefikMetricsURL) if err != nil { diff --git a/agent/internal/traefik/reload_test.go b/agent/internal/traefik/reload_test.go index 5a9bb34e..f3519ac2 100644 --- a/agent/internal/traefik/reload_test.go +++ b/agent/internal/traefik/reload_test.go @@ -20,6 +20,54 @@ traefik_config_last_reload_success 1.725e+09 } } +func TestWaitForMetricsReadyRetriesTemporaryErrors(t *testing.T) { + originalReader := readLastSuccessfulReload + originalPollInterval := reloadPollInterval + t.Cleanup(func() { + readLastSuccessfulReload = originalReader + reloadPollInterval = originalPollInterval + }) + + attempts := 0 + readLastSuccessfulReload = func() (time.Time, error) { + attempts++ + if attempts < 3 { + return time.Time{}, os.ErrNotExist + } + return time.Now(), nil + } + reloadPollInterval = time.Millisecond + + if err := waitForMetricsReady(50 * time.Millisecond); err != nil { + t.Fatal(err) + } + if attempts != 3 { + t.Fatalf("metrics read attempted %d times, want 3", attempts) + } +} + +func TestWaitForMetricsReadyTimesOut(t *testing.T) { + originalReader := readLastSuccessfulReload + originalPollInterval := reloadPollInterval + t.Cleanup(func() { + readLastSuccessfulReload = originalReader + reloadPollInterval = originalPollInterval + }) + + readLastSuccessfulReload = func() (time.Time, error) { + return time.Time{}, os.ErrDeadlineExceeded + } + reloadPollInterval = time.Millisecond + + err := waitForMetricsReady(5 * time.Millisecond) + if err == nil { + t.Fatal("metrics readiness wait unexpectedly succeeded") + } + if !strings.Contains(err.Error(), "traefik metrics did not become ready within 5ms") { + t.Fatalf("unexpected timeout error: %v", err) + } +} + func TestDynamicConfigReloadedRequiresReloadAtOrAfterNewestFile(t *testing.T) { originalDir := dynamicConfigDir originalReader := readLastSuccessfulReload diff --git a/agent/internal/traefik/static.go b/agent/internal/traefik/static.go index 6c9cf2bd..d022a84c 100644 --- a/agent/internal/traefik/static.go +++ b/agent/internal/traefik/static.go @@ -6,7 +6,6 @@ import ( "os" "os/exec" "reflect" - "time" "gopkg.in/yaml.v3" ) @@ -16,6 +15,7 @@ const ( metricsEntryPointAddr = "127.0.0.1:9100" ) +// Whole-number buckets must be ints to match yaml.v3's decoded types. var prometheusLatencyBuckets = []interface{}{ 0.005, 0.01, @@ -26,12 +26,12 @@ var prometheusLatencyBuckets = []interface{}{ 0.25, 0.5, 0.75, - 1.0, + 1, 2.5, - 5.0, - 10.0, - 30.0, - 60.0, + 5, + 10, + 30, + 60, } func validateStaticConfig(data []byte) error { @@ -218,6 +218,5 @@ func ReloadTraefik() error { return fmt.Errorf("failed to restart traefik: %w", err) } log.Printf("[traefik] restarted traefik to apply static config changes") - time.Sleep(2 * time.Second) - return nil + return waitForMetricsReady(metricsReadyTimeout) } diff --git a/agent/internal/traefik/static_test.go b/agent/internal/traefik/static_test.go index 7f2b47af..71d99e9e 100644 --- a/agent/internal/traefik/static_test.go +++ b/agent/internal/traefik/static_test.go @@ -1,6 +1,10 @@ package traefik -import "testing" +import ( + "testing" + + "gopkg.in/yaml.v3" +) func TestEnsurePrometheusMetricsConfigAddsPrivateMetricsEndpoint(t *testing.T) { config := map[string]interface{}{ @@ -44,7 +48,17 @@ func TestEnsurePrometheusMetricsConfigIsStable(t *testing.T) { if !ensurePrometheusMetricsConfig(config) { t.Fatal("expected first call to modify config") } - if ensurePrometheusMetricsConfig(config) { + + data, err := yaml.Marshal(config) + if err != nil { + t.Fatalf("failed to marshal config: %v", err) + } + var roundTripped map[string]interface{} + if err := yaml.Unmarshal(data, &roundTripped); err != nil { + t.Fatalf("failed to unmarshal config: %v", err) + } + + if ensurePrometheusMetricsConfig(roundTripped) { t.Fatal("expected second call to be stable") } } diff --git a/web/lib/inngest/functions/rollout-workflow.ts b/web/lib/inngest/functions/rollout-workflow.ts index 0c54f5c4..03c25530 100644 --- a/web/lib/inngest/functions/rollout-workflow.ts +++ b/web/lib/inngest/functions/rollout-workflow.ts @@ -418,33 +418,57 @@ export const rolloutWorkflow = inngest.createFunction( }); } - const certResult = await step.run("issue-certificates", async () => { - await db - .update(rollouts) - .set({ currentStage: "certificates" }) - .where(eq(rollouts.id, rolloutId)); - try { - const result = await issueCertificatesForRevision(specification); - if (result.issuedDomains.length > 0) { - await ingestRolloutLog( - rolloutId, - serviceId, - "certificates", - `Certificates issued for ${result.issuedDomains.length} domain(s)`, - ); - } - return { success: true as const }; - } catch (error) { - const message = - error instanceof Error - ? error.message - : "Certificate provisioning failed"; - await ingestRolloutLog(rolloutId, serviceId, "certificates", message); - return { success: false as const, reason: message }; + let certificatesIssued = false; + let certificateFailureReason = "Certificate provisioning failed"; + for (let attempt = 1; attempt <= 3; attempt++) { + const certResult = await step.run( + `issue-certificates-${attempt}`, + async () => { + await db + .update(rollouts) + .set({ currentStage: "certificates" }) + .where(eq(rollouts.id, rolloutId)); + try { + const result = await issueCertificatesForRevision(specification); + if (result.issuedDomains.length > 0) { + await ingestRolloutLog( + rolloutId, + serviceId, + "certificates", + `Certificates issued for ${result.issuedDomains.length} domain(s)`, + ); + } + return { success: true as const }; + } catch (error) { + const message = + error instanceof Error + ? error.message + : "Certificate provisioning failed"; + await ingestRolloutLog( + rolloutId, + serviceId, + "certificates", + message, + ); + return { success: false as const, reason: message }; + } + }, + ); + if (certResult.success) { + certificatesIssued = true; + break; } - }); - if (!certResult.success) { + certificateFailureReason = certResult.reason; + if (attempt < 3) { + await step.sleep( + `wait-for-certificate-retry-${attempt}`, + attempt === 1 ? "10s" : "20s", + ); + } + } + + if (!certificatesIssued) { await step.run("handle-certificate-failure", async () => { await handleRolloutFailure( rolloutId, @@ -455,7 +479,7 @@ export const rolloutWorkflow = inngest.createFunction( }); return { status: "failed", - reason: certResult.reason, + reason: certificateFailureReason, }; } diff --git a/web/tests/rollout-certificate-retry.test.ts b/web/tests/rollout-certificate-retry.test.ts new file mode 100644 index 00000000..5651f3fb --- /dev/null +++ b/web/tests/rollout-certificate-retry.test.ts @@ -0,0 +1,189 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => { + function updateQuery() { + const query = { + set: vi.fn(() => query), + where: vi.fn(async () => undefined), + }; + return query; + } + + return { + update: vi.fn(updateQuery), + issueCertificatesForRevision: vi.fn(), + handleRolloutFailure: vi.fn(), + ingestRolloutLog: vi.fn(), + }; +}); + +vi.mock("@/db", () => ({ db: { update: mocks.update } })); +vi.mock("@/db/queries", () => ({ getService: vi.fn() })); +vi.mock("@/lib/deployment-status", () => ({ + isObservedReady: vi.fn(), + observedReadyPhases: [], +})); +vi.mock("@/lib/preview-deployments", () => ({ + canDeployServiceRevision: vi.fn(), + updatePreviewGitHubStatus: vi.fn(), +})); +vi.mock("@/lib/routing-sync", () => ({ buildRoutingTargets: vi.fn() })); +vi.mock("@/lib/service-revisions", () => ({ + getRolloutServiceRevision: vi.fn(), +})); +vi.mock("@/lib/victoria-logs", () => ({ + ingestRolloutLog: mocks.ingestRolloutLog, +})); +vi.mock("@/lib/work-queue", () => ({ + enqueueReconcileForAllOnlineServers: vi.fn(), +})); +vi.mock("@/lib/inngest/client", () => ({ + inngest: { + createFunction: vi.fn( + (_options: unknown, handler: (input: unknown) => unknown) => handler, + ), + }, +})); +vi.mock("@/lib/inngest/events", () => ({ + inngestEvents: { + rolloutCreated: { name: "rollout/created" }, + rolloutCancelled: { name: "rollout/cancelled" }, + resourceStatusChanged: { name: "resource/status.changed" }, + serverDnsSynced: { name: "server/dns.synced" }, + }, +})); +vi.mock("@/lib/inngest/functions/rollout-helpers", () => ({ + checkForRollingUpdate: vi.fn(), + cleanupExistingDeployments: vi.fn(), + cleanupTerminalDeployments: vi.fn(), + completeRollout: vi.fn(), + createDeploymentRecords: vi.fn(), + issueCertificatesForRevision: mocks.issueCertificatesForRevision, + resolveRevisionPlacements: vi.fn(), + validateServers: vi.fn(), +})); +vi.mock("@/lib/inngest/functions/rollout-utils", () => ({ + handleRolloutFailure: mocks.handleRolloutFailure, +})); + +import { rolloutWorkflow } from "@/lib/inngest/functions/rollout-workflow"; + +function invokeRollout() { + const step = { + run: vi.fn(async (name: string, operation: () => unknown) => { + if ( + name.startsWith("issue-certificates-") || + name === "handle-certificate-failure" + ) { + return operation(); + } + + switch (name) { + case "validate-service": + return false; + case "acquire-rollout-turn-0": + return "acquired"; + case "load-service-revision": + return { + id: "revision-1", + specification: { + ports: [], + serverless: { enabled: false }, + }, + }; + case "log-rollout-started": + return undefined; + case "load-placements": + return { + success: true, + placements: [{ serverId: "server-1", replicas: 1 }], + totalReplicas: 1, + }; + case "validate-servers": + return { success: true, serverIds: ["server-1"] }; + case "cleanup-terminal-deployments": + return undefined; + case "check-rolling-update": + return false; + case "cleanup-existing": + return undefined; + case "create-deployments": + throw new Error("continued after certificates"); + default: + throw new Error(`unexpected step: ${name}`); + } + }), + sleep: vi.fn(async () => undefined), + waitForEvent: vi.fn(), + }; + const handler = rolloutWorkflow as unknown as (input: { + event: { data: { rolloutId: string; serviceId: string } }; + step: typeof step; + }) => Promise; + + return { + result: handler({ + event: { + data: { rolloutId: "rollout-1", serviceId: "service-1" }, + }, + step, + }), + step, + }; +} + +describe("rollout certificate retries", () => { + beforeEach(() => { + vi.clearAllMocks(); + }); + + it("uses durable backoff and continues after a later attempt succeeds", async () => { + mocks.issueCertificatesForRevision + .mockRejectedValueOnce(new Error("first outage")) + .mockRejectedValueOnce(new Error("second outage")) + .mockResolvedValueOnce({ issuedDomains: [] }); + + const { result, step } = invokeRollout(); + + await expect(result).rejects.toThrow("continued after certificates"); + expect(mocks.issueCertificatesForRevision).toHaveBeenCalledTimes(3); + expect(step.sleep).toHaveBeenNthCalledWith( + 1, + "wait-for-certificate-retry-1", + "10s", + ); + expect(step.sleep).toHaveBeenNthCalledWith( + 2, + "wait-for-certificate-retry-2", + "20s", + ); + expect(step.run).toHaveBeenCalledWith( + "create-deployments", + expect.any(Function), + ); + expect(mocks.handleRolloutFailure).not.toHaveBeenCalled(); + }); + + it("uses the existing failure path once after three failed attempts", async () => { + mocks.issueCertificatesForRevision + .mockRejectedValueOnce(new Error("first outage")) + .mockRejectedValueOnce(new Error("second outage")) + .mockRejectedValueOnce(new Error("final outage")); + + const { result, step } = invokeRollout(); + + await expect(result).resolves.toEqual({ + status: "failed", + reason: "final outage", + }); + expect(mocks.issueCertificatesForRevision).toHaveBeenCalledTimes(3); + expect(step.sleep).toHaveBeenCalledTimes(2); + expect(mocks.handleRolloutFailure).toHaveBeenCalledOnce(); + expect(mocks.handleRolloutFailure).toHaveBeenCalledWith( + "rollout-1", + "service-1", + "certificate_provisioning_failed", + false, + ); + }); +});