From 7d99424cd8c0306b020f2370a19eab9c30093063 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Tue, 15 Sep 2026 20:58:54 +0000 Subject: [PATCH] feat(stovepipe): cool down admissions after failures Summary: Intent: - Avoid immediately consuming more build resources after a runner-reported failure. Changes: - Persist the first terminal observation time on Build. - Advance the admission deadline while releasing the failed build slot in one queue CAS. - Add per-queue failure cooldown configuration to the YAML policy store. - Keep success, cancellation, and DLQ-forced failure outside the cooldown policy. This change builds on the per-queue admission configuration introduced by the parent PR. --- doc/rfc/stovepipe/steps/build.md | 4 +- doc/rfc/stovepipe/steps/buildsignal.md | 13 +- doc/rfc/stovepipe/steps/process.md | 1 + service/stovepipe/README.md | 2 +- service/stovepipe/server/main.go | 4 +- service/stovepipe/server/main_test.go | 3 + service/stovepipe/server/queues.yaml | 3 + stovepipe/controller/buildsignal/BUILD.bazel | 2 + .../controller/buildsignal/buildsignal.go | 67 +++++++-- .../buildsignal/buildsignal_test.go | 130 ++++++++++++++++-- stovepipe/entity/build.go | 9 +- stovepipe/entity/queue_config.go | 3 + stovepipe/extension/queueconfig/README.md | 6 +- .../extension/queueconfig/default/default.go | 2 + .../queueconfig/default/default_test.go | 1 + .../extension/queueconfig/yaml/yaml_test.go | 22 +-- .../extension/storage/mysql/build_store.go | 11 +- .../storage/mysql/build_store_test.go | 44 +++--- .../extension/storage/mysql/schema/build.sql | 1 + .../stovepipe/extension/storage/suite.go | 2 + 20 files changed, 256 insertions(+), 74 deletions(-) diff --git a/doc/rfc/stovepipe/steps/build.md b/doc/rfc/stovepipe/steps/build.md index f630870b2..f8eac5400 100644 --- a/doc/rfc/stovepipe/steps/build.md +++ b/doc/rfc/stovepipe/steps/build.md @@ -10,7 +10,7 @@ It handles only the trigger: it does not poll for completion, record greenness, `build` consumes a request id, published by `process` in Phase 1 (see [workflow.md](doc/rfc/stovepipe/workflow.md#workflow)). Both phases drive the same `build` → `buildsignal` machinery against the same `Request` row. The `process`/`analyze` → `build` topic is partitioned by **request id**; see [Partitioning](#partitioning) for the full rationale, including why `build` → `buildsignal` partitions finer (by build id) than SubmitQueue's equivalent topic. -**`build` is not the sole writer of the rows it touches, and its own writes are narrow.** On `Request`, `build` never writes at all — it only reads the `BuildStrategy`/`BaseURI`/`URI` fields `process` set at admit, and it leaves `Request.State` untouched at `processing` throughout (`process`, `buildsignal`, and the DLQ reconciler are `Request.State`'s only writers). On `Build`, `build` is the sole creator — it calls `BuildStore.Create` exactly once, at step 6 — and never mutates the row again; `buildsignal` is the sole writer of `Build.Status`/`Build.Version` afterward (see [buildsignal.md](doc/rfc/stovepipe/steps/buildsignal.md#input-and-re-entrancy)). This division is why `build` never needs a CAS/version write of its own: `Create` is the only storage mutation in its algorithm. +**`build` is not the sole writer of the rows it touches, and its own writes are narrow.** On `Request`, `build` never writes at all — it only reads the `BuildStrategy`/`BaseURI`/`URI` fields `process` set at admit, and it leaves `Request.State` untouched at `processing` throughout (`process`, `buildsignal`, and the DLQ reconciler are `Request.State`'s only writers). On `Build`, `build` is the sole creator — it calls `BuildStore.Create` exactly once, at step 6 — and never mutates the row again; `buildsignal` is the sole writer of `Build.Status`/`Build.TerminalAtMs`/`Build.Version` afterward (see [buildsignal.md](doc/rfc/stovepipe/steps/buildsignal.md#input-and-re-entrancy)). This division is why `build` never needs a CAS/version write of its own: `Create` is the only storage mutation in its algorithm. `build` is phase-agnostic: it never asks "which phase is this?" It reads whatever scope is already persisted and immutable on the `Request` and acts on it. Phase 2's project-scoped invocation is expected to read project-scoped equivalents of the scope fields off the same `Request`; the exact shape of a project-scoped trigger is left to the `analyze` design, consistent with `workflow.md`'s "project mapping contract" open question — see [Project-scoped `Trigger`: reserved, not yet designed](#project-scoped-trigger-reserved-not-yet-designed) for the resulting gap in the `BuildRunner` contract itself. @@ -52,7 +52,7 @@ For a delivery carrying request id `R`: either domain — the shape is deferred until then, not decided here. - failure -> return raw; classifier decides (transient runner blip retryable, bad URI not). -6. Persist Build{ID: buildID.ID, RequestID: R.ID, Status: accepted, Version: 1} +6. Persist Build{ID: buildID.ID, RequestID: R.ID, Status: accepted, TerminalAtMs: 0, Version: 1} via BuildStore.Create. - the row carries no scope; it is recoverable from the Request's immutable fields (see the entity table). diff --git a/doc/rfc/stovepipe/steps/buildsignal.md b/doc/rfc/stovepipe/steps/buildsignal.md index 7ce0dfcaa..a457d1cb8 100644 --- a/doc/rfc/stovepipe/steps/buildsignal.md +++ b/doc/rfc/stovepipe/steps/buildsignal.md @@ -12,7 +12,7 @@ It handles only the poll loop: it does not decide build strategy, write greennes Its logic does not branch on phase: it loads the `Build`, polls it toward terminal, persists the result, and publishes the request id onward to `record`. What differs between phases is what `record` does with that publish (whole-repo vs. per-project greenness) — not anything `buildsignal` decides. -`buildsignal` is the sole writer of `Build.Status`/`Build.Version` after `build` creates the row (see [build.md](doc/rfc/stovepipe/steps/build.md#input-partitioning-and-the-single-writer-property)). It reads `Request` via `RequestStore.Get` (for `R.Queue`, to resolve the build-runner) and writes it exactly once, at the terminal transition, to record the build's outcome — the one `Request.State` write outside `process` and the DLQ reconciler. +`buildsignal` is the sole writer of `Build.Status`/`Build.TerminalAtMs`/`Build.Version` after `build` creates the row (see [build.md](doc/rfc/stovepipe/steps/build.md#input-partitioning-and-the-single-writer-property)). It reads `Request` via `RequestStore.Get` (for `R.Queue`, to resolve the build-runner) and writes it exactly once, at the terminal transition, to record the build's outcome — the one `Request.State` write outside `process` and the DLQ reconciler. Its early-exit guard is deliberately narrower than `State.IsTerminal()`: it proceeds when the request is `processing` **or** already carries a build outcome. The second case matters because a redelivery after the outcome was stamped but before the `record` publish landed must re-publish rather than drop the signal; everything it re-runs is a no-op (the status is unchanged, the outcome is already recorded, the slot is not released twice) and the `record` publish is idempotent. @@ -55,15 +55,20 @@ For a delivery carrying build id `B`: - Build.Status already terminal, status differs -> terminal is WRITE-ONCE: do not overwrite; continue with the STORED status as authoritative (see Edge cases). - otherwise persist via BuildStore.Update(ctx, Build{...Status: status}, oldVersion, newVersion): + - when status is terminal, stamp TerminalAtMs with the observation time in the same update; - newVersion = Build.Version + 1; assign Build.Version = newVersion only on success. - ErrVersionMismatch -> return with its declaration-level retryable classification (a concurrent writer moved the row; reload and re-check). - with the write-once rule, accepted -> running -> {succeeded|failed|cancelled} is monotonic by mechanism, not by assumption about the backend. 7. If the stored status is terminal, and R does not already carry an outcome: - a. Release the queue's build slot: CAS-decrement Queue.in_flight_count, clamped at zero. + a. Load the queue config only when the runner-reported status is failed. Apply + failure_cooldown_ms only when it is positive; non-positive values disable the policy. + b. Release the queue's build slot in one Queue CAS: decrement in_flight_count, clamped at zero, + and for failed status advance build_admission_not_before_ms to at least + Build.TerminalAtMs + failure_cooldown_ms. - failure here aborts the step: R must not go terminal while still holding a slot. - b. CAS R from processing to the outcome the stored status projects onto it: + c. CAS R from processing to the outcome the stored status projects onto it: succeeded -> succeeded, failed -> failed, cancelled -> cancelled. First writer wins. Then publish R.ID to the record topic, partitioned by request id; ack, return. No re-publish to buildsignal. @@ -86,6 +91,8 @@ For a delivery carrying build id `B`: **Why `buildsignal` releases the slot rather than `record`**: the gate `process` claims is a *build* slot — it exists to bound concurrent builds per Queue — and once the build is terminal the build is over. Releasing here also keeps the invariant *a terminal `Request` has already released its slot*, which is what makes the DLQ reconciler's early-return on terminal requests safe. +**Why failure cooldown is part of the slot-release CAS**: both fields govern the next logical admission, so updating them from the same freshly loaded Queue snapshot prevents a process admission from slipping between capacity release and deadline advancement. A version conflict reloads the row and recomputes the maximum, preserving a later deadline written by the minimum-interval policy or another terminal result. The deadline is based on the persisted `Build.TerminalAtMs`, not redelivery time, so retries cannot slide the cooldown forward. Only an actual runner-reported `failed` status applies this policy; `cancelled` and DLQ-forced request failure do not. + **Why `record` hears only terminal signals**: `record` has no non-terminal work — by its own contract a non-terminal signal would be a pure no-op — and step 7 already branches on terminality to decide whether to keep polling, so gating the publish costs nothing and spares `record` a no-op delivery on every poll tick of every running build. Crash-safety is unaffected: a crash between the terminal `Update` and the publish redelivers the message; step 5 re-polls (the runner reports the same terminal status), step 6 no-ops, step 7 publishes. This is a deliberate divergence from SubmitQueue, whose buildsignal republishes to `speculate` on every tick — sound there because speculate is a state machine that may act on any signal; stovepipe has no such consumer. **Why step 6 guards on status and makes terminal write-once**: an unchanged status skips the CAS write entirely, so a long build being polled every couple of seconds doesn't churn `Build.Version` on every tick — the version only advances on a real state transition. The write-once rule exists because CAS alone cannot provide it: optimistic locking defends against *concurrent* writers, but a later delivery that polls a flaky backend and sees a different terminal status would CAS cleanly against the current version and overwrite (see Edge cases). A given `Build` has a single poll partition (see [Partitioning](doc/rfc/stovepipe/steps/build.md#partitioning)), so the only writer racing the CAS is a redelivery of the same message (e.g. after a lapsed visibility lease); `ErrVersionMismatch` carries a retryable classification and converges on redelivery. diff --git a/doc/rfc/stovepipe/steps/process.md b/doc/rfc/stovepipe/steps/process.md index 80aa719c2..adb2e0587 100644 --- a/doc/rfc/stovepipe/steps/process.md +++ b/doc/rfc/stovepipe/steps/process.md @@ -64,6 +64,7 @@ Validation is expensive and shares a baseline, so heads arriving while an earlie | Queue row | `build_admission_not_before_ms` | Durable earliest time for the next logical admission; zero until a time policy advances it. | | Queue config | `max_concurrent` | Cap on concurrent in-flight validations. **Default 1**; the YAML store supports per-queue overrides. | | Queue config | `minimum_build_admission_interval_ms` | Minimum start-to-start spacing between logical admissions. Positive values enable general throttling; non-positive values disable it. **Default 0**. | +| Queue config | `failure_cooldown_ms` | Delay after a runner-reported failed build before the next logical admission. Positive values enable the cooldown; non-positive values disable it. **Default 0**. | A slot is held from admit until the build goes terminal (`process → build → buildsignal`), not just while `process` runs. It is released when the Request reaches **any** terminal state and `in_flight_count` is decremented — `buildsignal` recording the build's outcome, success *or* failure, or the DLQ reconciler forcing a terminal `failed` (see [integrity](#in_flight_count-integrity)). A build *failure* frees the slot just like a success; only a Request that never terminates keeps its slot. diff --git a/service/stovepipe/README.md b/service/stovepipe/README.md index 9c67f31a9..70f273ac1 100644 --- a/service/stovepipe/README.md +++ b/service/stovepipe/README.md @@ -24,7 +24,7 @@ Stovepipe therefore needs two MySQL databases: a **storage** database (the `queu - **`inMemoryCounter`** — a process-local `counter.Counter` for sequence numbers; not durable. A real deployment uses a persistent implementation (e.g. `platform/extension/counter/mysql`). - **`fakeSourceControlFactory`** — seeds each queue with a deterministic single-commit history so ingest resolves a stable head URI (and re-ingesting the same queue exercises the dedup path). A real deployment supplies a VCS-backed `sourcecontrol.Factory`, which is also where a queue's promotion ref is resolved. The fake has no ref to move, so a promotion locally shows up only in the record consumer's logs. -`QUEUE_CONFIG_PATH` selects the YAML-backed queue policy store. When it is unset, the server uses the built-in defaults (`max_concurrent: 1`, `gate_wait_delay_ms: 5000`, and time-based admission throttling disabled). When it is set, the configured queue names must exactly match `MQ_TENANTS`. Each configured queue can set `minimum_build_admission_interval_ms` for general start-to-start throttling. For example, `minimum_build_admission_interval_ms: 3600000` permits at most one logical admission per hour. +`QUEUE_CONFIG_PATH` selects the YAML-backed queue policy store. When it is unset, the server uses the built-in defaults (`max_concurrent: 1`, `gate_wait_delay_ms: 5000`, and both time-based admission policies disabled). When it is set, the configured queue names must exactly match `MQ_TENANTS`. Each configured queue can set `minimum_build_admission_interval_ms` for general start-to-start throttling and `failure_cooldown_ms` for a delay after a runner-reported failed build. For example, `minimum_build_admission_interval_ms: 3600000` permits at most one logical admission per hour. ## Layout diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index f4a85f348..c794152e5 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -472,7 +472,7 @@ func registerPrimaryControllers( scope, store, materializer, - queueconfigdefault.NewStore(), + queueConfigs, sourceControl, registry, stovepipemq.TopicKeyProcess, @@ -489,7 +489,7 @@ func registerPrimaryControllers( } count++ - buildSignalController := buildsignal.NewController(logger, scope, store, materializer, brf, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal") + buildSignalController := buildsignal.NewController(logger, scope, store, materializer, queueConfigs, brf, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal") if err := c.Register(buildSignalController); err != nil { return count, fmt.Errorf("failed to register buildsignal controller: %w", err) } diff --git a/service/stovepipe/server/main_test.go b/service/stovepipe/server/main_test.go index d7c15c37c..bbe78bbf9 100644 --- a/service/stovepipe/server/main_test.go +++ b/service/stovepipe/server/main_test.go @@ -240,6 +240,7 @@ func TestLoadQueueConfigs(t *testing.T) { require.NoError(t, err) assert.EqualValues(t, 1, cfg.MaxConcurrent) assert.Zero(t, cfg.MinimumBuildAdmissionIntervalMs) + assert.Zero(t, cfg.FailureCooldownMs) }) t.Run("path loads per-queue policies", func(t *testing.T) { @@ -249,6 +250,7 @@ func TestLoadQueueConfigs(t *testing.T) { max_concurrent: 2 gate_wait_delay_ms: 5000 minimum_build_admission_interval_ms: 3600000 + failure_cooldown_ms: 900000 `), 0o600)) store, err := loadQueueConfigs(path) @@ -257,6 +259,7 @@ func TestLoadQueueConfigs(t *testing.T) { require.NoError(t, err) assert.EqualValues(t, 2, cfg.MaxConcurrent) assert.EqualValues(t, 3_600_000, cfg.MinimumBuildAdmissionIntervalMs) + assert.EqualValues(t, 900_000, cfg.FailureCooldownMs) require.NoError(t, validateQueueConfigTenants(context.Background(), path, store, []string{"monorepo/main"})) require.Error(t, validateQueueConfigTenants(context.Background(), path, store, []string{"other"})) }) diff --git a/service/stovepipe/server/queues.yaml b/service/stovepipe/server/queues.yaml index 07d70d20c..a68e816cb 100644 --- a/service/stovepipe/server/queues.yaml +++ b/service/stovepipe/server/queues.yaml @@ -5,11 +5,14 @@ queues: gate_wait_delay_ms: 5000 # Set to 3600000 to admit at most one build per hour. minimum_build_admission_interval_ms: 0 + failure_cooldown_ms: 0 - name: monorepo/release max_concurrent: 1 gate_wait_delay_ms: 5000 minimum_build_admission_interval_ms: 0 + failure_cooldown_ms: 0 - name: "monorepo/slow?buildrunner-fake=build-slow" max_concurrent: 1 gate_wait_delay_ms: 5000 minimum_build_admission_interval_ms: 0 + failure_cooldown_ms: 0 diff --git a/stovepipe/controller/buildsignal/BUILD.bazel b/stovepipe/controller/buildsignal/BUILD.bazel index 61088f534..8c978c4f0 100644 --- a/stovepipe/controller/buildsignal/BUILD.bazel +++ b/stovepipe/controller/buildsignal/BUILD.bazel @@ -15,6 +15,7 @@ go_library( "//stovepipe/core/requestlog:go_default_library", "//stovepipe/entity:go_default_library", "//stovepipe/extension/buildrunner:go_default_library", + "//stovepipe/extension/queueconfig:go_default_library", "//stovepipe/extension/storage:go_default_library", "@com_github_uber_go_tally//:go_default_library", "@org_uber_go_zap//:go_default_library", @@ -38,6 +39,7 @@ go_test( "//stovepipe/entity:go_default_library", "//stovepipe/extension/buildrunner:go_default_library", "//stovepipe/extension/buildrunner/mock:go_default_library", + "//stovepipe/extension/queueconfig/mock:go_default_library", "//stovepipe/extension/storage:go_default_library", "//stovepipe/extension/storage/mock:go_default_library", "@com_github_stretchr_testify//assert:go_default_library", diff --git a/stovepipe/controller/buildsignal/buildsignal.go b/stovepipe/controller/buildsignal/buildsignal.go index 2c56a04b8..793e0f1dc 100644 --- a/stovepipe/controller/buildsignal/buildsignal.go +++ b/stovepipe/controller/buildsignal/buildsignal.go @@ -24,6 +24,7 @@ import ( "context" "errors" "fmt" + "time" "github.com/uber-go/tally" entityqueue "github.com/uber/submitqueue/platform/base/messagequeue" @@ -35,6 +36,7 @@ import ( "github.com/uber/submitqueue/stovepipe/core/requestlog" "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/buildrunner" + "github.com/uber/submitqueue/stovepipe/extension/queueconfig" "github.com/uber/submitqueue/stovepipe/extension/storage" "go.uber.org/zap" ) @@ -63,10 +65,12 @@ type Controller struct { metricsScope tally.Scope stores storage.Factory materializer requestlog.Materializer + queueConfigs queueconfig.Store buildRunners buildrunner.Factory registry consumer.TopicRegistry topicKey consumer.TopicKey consumerGroup string + now func() time.Time } // Verify Controller implements consumer.Controller interface at compile time. @@ -81,6 +85,7 @@ func NewController( scope tally.Scope, stores storage.Factory, materializer requestlog.Materializer, + queueConfigs queueconfig.Store, buildRunners buildrunner.Factory, registry consumer.TopicRegistry, topicKey consumer.TopicKey, @@ -91,10 +96,12 @@ func NewController( metricsScope: scope.SubScope("buildsignal_controller"), stores: stores, materializer: materializer, + queueConfigs: queueConfigs, buildRunners: buildRunners, registry: registry, topicKey: topicKey, consumerGroup: consumerGroup, + now: time.Now, } } @@ -168,16 +175,16 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er return fmt.Errorf("failed to poll status for build %s: %w", build.ID, err) } - effective, err := c.reconcile(ctx, store, build, status) + build, err = c.reconcile(ctx, store, build, status) if err != nil { return err } - if effective.IsTerminal() { + if build.Status.IsTerminal() { if err := c.persistBuildFinishedLog(ctx, store, request, build.ID); err != nil { return err } - if err := c.finishRequest(ctx, store, &request, effective); err != nil { + if err := c.finishRequest(ctx, store, &request, build); err != nil { return err } if err := c.persistOutcomeLog(ctx, store, request); err != nil { @@ -189,7 +196,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er c.logger.Infow("build reached terminal status", "build_id", build.ID, "request_id", request.ID, - "status", string(effective), + "status", string(build.Status), "request_state", string(request.State), ) return nil @@ -198,11 +205,11 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er // Not terminal yet: hold the delivery so this same message redelivers // after the poll delay — the partition (keyed by build id) sleeps with it, // and the redelivery does not count toward the retry limit. - delayMs := pollDelay(effective) + delayMs := pollDelay(build.Status) delivery.Hold(delayMs) c.logger.Debugw("holding for next build status poll", "build_id", build.ID, - "status", string(effective), + "status", string(build.Status), "delay_ms", delayMs, ) return nil @@ -234,23 +241,41 @@ func (c *Controller) persistBuildFinishedLog(ctx context.Context, store storage. // the request non-terminal, so redelivery re-runs both steps and decrements again // — transiently over-admitting by one until releaseBuildSlot's zero clamp // reconverges, which is the failure mode this pipeline prefers. -func (c *Controller) finishRequest(ctx context.Context, store storage.Storage, request *entity.Request, status entity.BuildStatus) error { +func (c *Controller) finishRequest(ctx context.Context, store storage.Storage, request *entity.Request, build entity.Build) error { if request.State.HasBuildOutcome() { return nil } - if err := c.releaseBuildSlot(ctx, store, request.Queue); err != nil { + failureCooldownMs, err := c.failureCooldown(ctx, request.Queue, build) + if err != nil { + return err + } + if err := c.releaseBuildSlot(ctx, store, request.Queue, build.TerminalAtMs, failureCooldownMs); err != nil { metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1, metrics.TagsFromContext(ctx)...) return err } - if err := c.markOutcome(ctx, store, request, outcomeState(status)); err != nil { + if err := c.markOutcome(ctx, store, request, outcomeState(build.Status)); err != nil { metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1, metrics.TagsFromContext(ctx)...) return err } return nil } +func (c *Controller) failureCooldown(ctx context.Context, queueName string, build entity.Build) (int64, error) { + if build.Status != entity.BuildStatusFailed || build.TerminalAtMs == 0 { + return 0, nil + } + cfg, err := c.queueConfigs.Get(ctx, queueName) + if err != nil { + return 0, fmt.Errorf("failed to load queue config for %s: %w", queueName, err) + } + if cfg.FailureCooldownMs <= 0 { + return 0, nil + } + return cfg.FailureCooldownMs, nil +} + func (c *Controller) persistOutcomeLog(ctx context.Context, store storage.Storage, request entity.Request) error { var reason entity.RequestOutcomeReason // The durable request is authoritative when duplicate builds race to record different outcomes. @@ -326,7 +351,7 @@ func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, req // (preserving concurrent updates), clamps at zero, and retries on version conflicts. // Unlike process's unwind-path release this is not best-effort: the caller must not // mark the request terminal if the slot was not freed, so a hard failure is returned. -func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage, queueName string) error { +func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage, queueName string, terminalAtMs, failureCooldownMs int64) error { queueStore := store.GetQueueStore() for { @@ -340,6 +365,11 @@ func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage updated := queueRow updated.InFlightCount = queueRow.InFlightCount - 1 + cooldownNotBeforeMs := terminalAtMs + failureCooldownMs + cooldownAdvanced := failureCooldownMs > 0 && cooldownNotBeforeMs > updated.BuildAdmissionNotBeforeMs + if cooldownAdvanced { + updated.BuildAdmissionNotBeforeMs = cooldownNotBeforeMs + } newVersion := queueRow.Version + 1 if err := queueStore.Update(ctx, updated, queueRow.Version, newVersion); err != nil { if errors.Is(err, storage.ErrVersionMismatch) { @@ -348,6 +378,9 @@ func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage return fmt.Errorf("failed to release build slot for queue %s: %w", queueName, err) } metrics.NamedCounter(c.metricsScope, _opName, "slot_released", 1, metrics.TagsFromContext(ctx)...) + if cooldownAdvanced { + metrics.NamedCounter(c.metricsScope, _opName, "failure_cooldowns", 1, metrics.TagsFromContext(ctx)...) + } return nil } } @@ -356,24 +389,28 @@ func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage // should drive the rest of Process: the polled status when persisted (or // already unchanged), or the stored status when a stored terminal status is // write-once-protected against a differing poll. -func (c *Controller) reconcile(ctx context.Context, store storage.Storage, build entity.Build, status entity.BuildStatus) (entity.BuildStatus, error) { +func (c *Controller) reconcile(ctx context.Context, store storage.Storage, build entity.Build, status entity.BuildStatus) (entity.Build, error) { if status == build.Status { - return build.Status, nil + return build, nil } // Terminal is write-once: a later poll of a flaky backend must never // overwrite an already-committed terminal status. if build.Status.IsTerminal() { - return build.Status, nil + return build, nil } newVersion := build.Version + 1 updated := build updated.Status = status + if status.IsTerminal() { + updated.TerminalAtMs = c.now().UnixMilli() + } if err := store.GetBuildStore().Update(ctx, updated, build.Version, newVersion); err != nil { - return "", fmt.Errorf("failed to persist status for build %s: %w", build.ID, err) + return entity.Build{}, fmt.Errorf("failed to persist status for build %s: %w", build.ID, err) } - return status, nil + updated.Version = newVersion + return updated, nil } // loadBuild returns the build for id. diff --git a/stovepipe/controller/buildsignal/buildsignal_test.go b/stovepipe/controller/buildsignal/buildsignal_test.go index 0bfb76820..28f7e5a6e 100644 --- a/stovepipe/controller/buildsignal/buildsignal_test.go +++ b/stovepipe/controller/buildsignal/buildsignal_test.go @@ -18,6 +18,7 @@ import ( "context" "errors" "testing" + "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -34,6 +35,7 @@ import ( "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/buildrunner" buildrunnermock "github.com/uber/submitqueue/stovepipe/extension/buildrunner/mock" + queueconfigmock "github.com/uber/submitqueue/stovepipe/extension/queueconfig/mock" "github.com/uber/submitqueue/stovepipe/extension/storage" storagemock "github.com/uber/submitqueue/stovepipe/extension/storage/mock" "go.uber.org/mock/gomock" @@ -44,6 +46,7 @@ const ( testQueue = "monorepo/main" testID = "request/monorepo/main/7" testBuildID = "bk-1" + testNowMs = int64(1_750_000_000_000) ) func queueContext() context.Context { @@ -59,6 +62,7 @@ type buildsignalMocks struct { queueStore *storagemock.MockQueueStore store *storagemock.MockStorage materializer *requestlogmock.MockMaterializer + queueConfigs *queueconfigmock.MockStore runnerFactory *buildrunnermock.MockFactory runner *buildrunnermock.MockBuildRunner publisher *mqmock.MockPublisher @@ -81,6 +85,7 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildsig queueStore: storagemock.NewMockQueueStore(ctrl), store: storagemock.NewMockStorage(ctrl), materializer: requestlogmock.NewMockMaterializer(ctrl), + queueConfigs: queueconfigmock.NewMockStore(ctrl), runnerFactory: buildrunnermock.NewMockFactory(ctrl), runner: buildrunnermock.NewMockBuildRunner(ctrl), publisher: mqmock.NewMockPublisher(ctrl), @@ -100,7 +105,8 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildsig }) require.NoError(t, err) - c := NewController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: m.store}, m.materializer, m.runnerFactory, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal") + c := NewController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: m.store}, m.materializer, m.queueConfigs, m.runnerFactory, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal") + c.now = func() time.Time { return time.UnixMilli(testNowMs) } return c, m } @@ -165,6 +171,12 @@ func build(status entity.BuildStatus, version int32) entity.Build { } } +func terminalBuild(status entity.BuildStatus, version int32) entity.Build { + build := build(status, version) + build.TerminalAtMs = testNowMs + return build +} + // queueRow returns the testQueue's row holding inFlight admitted validations. func queueRow(inFlight, version int32) entity.Queue { return entity.Queue{ @@ -326,7 +338,7 @@ func TestProcess(t *testing.T) { m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) - m.buildStore.EXPECT().Update(gomock.Any(), build(entity.BuildStatusSucceeded, 2), int32(2), int32(3)).Return(nil) + m.buildStore.EXPECT().Update(gomock.Any(), terminalBuild(entity.BuildStatusSucceeded, 2), int32(2), int32(3)).Return(nil) logCall := expectFinish(m, entity.RequestStateSucceeded) m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil).After(logCall) }, @@ -338,7 +350,8 @@ func TestProcess(t *testing.T) { m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusFailed, nil, nil) - m.buildStore.EXPECT().Update(gomock.Any(), build(entity.BuildStatusFailed, 2), int32(2), int32(3)).Return(nil) + m.buildStore.EXPECT().Update(gomock.Any(), terminalBuild(entity.BuildStatusFailed, 2), int32(2), int32(3)).Return(nil) + m.queueConfigs.EXPECT().Get(gomock.Any(), testQueue).Return(entity.QueueConfig{Name: testQueue}, nil) logCall := expectFinish(m, entity.RequestStateFailed) m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil).After(logCall) }, @@ -350,7 +363,7 @@ func TestProcess(t *testing.T) { m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusCancelled, nil, nil) - m.buildStore.EXPECT().Update(gomock.Any(), build(entity.BuildStatusCancelled, 2), int32(2), int32(3)).Return(nil) + m.buildStore.EXPECT().Update(gomock.Any(), terminalBuild(entity.BuildStatusCancelled, 2), int32(2), int32(3)).Return(nil) logCall := expectFinish(m, entity.RequestStateCancelled) m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil).After(logCall) }, @@ -360,7 +373,7 @@ func TestProcess(t *testing.T) { setup: func(m buildsignalMocks) { request := requestWithState(entity.RequestStateSucceeded) request.Version = 2 - m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(entity.BuildStatusSucceeded, 3), nil) + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(terminalBuild(entity.BuildStatusSucceeded, 3), nil) m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(request, nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) @@ -375,7 +388,7 @@ func TestProcess(t *testing.T) { name: "build event failure stops before request outcome", wantErr: true, setup: func(m buildsignalMocks) { - m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(entity.BuildStatusSucceeded, 3), nil) + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(terminalBuild(entity.BuildStatusSucceeded, 3), nil) m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) @@ -398,7 +411,7 @@ func TestProcess(t *testing.T) { setup: func(m buildsignalMocks) { request := requestWithState(entity.RequestStateSucceeded) request.Version = 2 - m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(entity.BuildStatusSucceeded, 3), nil) + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(terminalBuild(entity.BuildStatusSucceeded, 3), nil) m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(request, nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) @@ -413,7 +426,7 @@ func TestProcess(t *testing.T) { { name: "write-once: stored terminal status is never overwritten", setup: func(m buildsignalMocks) { - m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(entity.BuildStatusSucceeded, 5), nil) + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(terminalBuild(entity.BuildStatusSucceeded, 5), nil) m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusFailed, nil, nil) @@ -452,7 +465,7 @@ func TestProcess(t *testing.T) { wantErr: true, wantRetry: false, setup: func(m buildsignalMocks) { - m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(entity.BuildStatusSucceeded, 3), nil) + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(terminalBuild(entity.BuildStatusSucceeded, 3), nil) m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) @@ -467,7 +480,7 @@ func TestProcess(t *testing.T) { wantErr: true, wantRetry: false, setup: func(m buildsignalMocks) { - m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(entity.BuildStatusSucceeded, 3), nil) + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(terminalBuild(entity.BuildStatusSucceeded, 3), nil) m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) @@ -482,7 +495,7 @@ func TestProcess(t *testing.T) { name: "outcome log failure retains terminal request and released slot", wantErr: true, setup: func(m buildsignalMocks) { - m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(entity.BuildStatusSucceeded, 3), nil) + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(terminalBuild(entity.BuildStatusSucceeded, 3), nil) m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) @@ -549,6 +562,101 @@ func TestProcess(t *testing.T) { } } +func TestProcessFailedBuildAppliesCooldown(t *testing.T) { + ctrl := gomock.NewController(t) + c, m := newController(t, ctrl) + const cooldownMs = int64(60_000) + + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(entity.BuildStatusRunning, 2), nil) + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) + m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) + m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusFailed, nil, nil) + m.buildStore.EXPECT().Update(gomock.Any(), terminalBuild(entity.BuildStatusFailed, 2), int32(2), int32(3)).Return(nil) + eventCall := expectBuildFinished(m) + configCall := m.queueConfigs.EXPECT().Get(gomock.Any(), testQueue).Return(entity.QueueConfig{ + Name: testQueue, + FailureCooldownMs: cooldownMs, + }, nil).After(eventCall) + m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(queueRow(1, 4), nil).After(configCall) + cooledDownQueue := queueRow(0, 4) + cooledDownQueue.BuildAdmissionNotBeforeMs = testNowMs + cooldownMs + queueCall := m.queueStore.EXPECT().Update(gomock.Any(), cooledDownQueue, int32(4), int32(5)).Return(nil) + requestCall := m.reqStore.EXPECT().Update(gomock.Any(), requestWithState(entity.RequestStateFailed), int32(1), int32(2)).Return(nil).After(queueCall) + logCall := expectOutcomeLog(m, entity.RequestStateFailed, 2).After(requestCall) + m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil).After(logCall) + + require.NoError(t, c.Process(queueContext(), delivery(t, ctrl, buildSignalPayload(t, testBuildID)))) +} + +func TestProcessStoredFailedBuildDoesNotSlideCooldown(t *testing.T) { + ctrl := gomock.NewController(t) + c, m := newController(t, ctrl) + c.now = func() time.Time { return time.UnixMilli(testNowMs + 10_000) } + const cooldownMs = int64(60_000) + + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(terminalBuild(entity.BuildStatusFailed, 3), nil) + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) + m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) + m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusFailed, nil, nil) + eventCall := expectBuildFinished(m) + configCall := m.queueConfigs.EXPECT().Get(gomock.Any(), testQueue).Return(entity.QueueConfig{ + Name: testQueue, + FailureCooldownMs: cooldownMs, + }, nil).After(eventCall) + m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(queueRow(1, 4), nil).After(configCall) + cooledDownQueue := queueRow(0, 4) + cooledDownQueue.BuildAdmissionNotBeforeMs = testNowMs + cooldownMs + queueCall := m.queueStore.EXPECT().Update(gomock.Any(), cooledDownQueue, int32(4), int32(5)).Return(nil) + requestCall := m.reqStore.EXPECT().Update(gomock.Any(), requestWithState(entity.RequestStateFailed), int32(1), int32(2)).Return(nil).After(queueCall) + logCall := expectOutcomeLog(m, entity.RequestStateFailed, 2).After(requestCall) + m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil).After(logCall) + + require.NoError(t, c.Process(queueContext(), delivery(t, ctrl, buildSignalPayload(t, testBuildID)))) +} + +func TestProcessNegativeFailureCooldownLeavesPolicyDisabled(t *testing.T) { + ctrl := gomock.NewController(t) + c, m := newController(t, ctrl) + + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(terminalBuild(entity.BuildStatusFailed, 3), nil) + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) + m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) + m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusFailed, nil, nil) + eventCall := expectBuildFinished(m) + m.queueConfigs.EXPECT().Get(gomock.Any(), testQueue).Return(entity.QueueConfig{ + Name: testQueue, + FailureCooldownMs: -1, + }, nil).After(eventCall) + m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(queueRow(1, 4), nil) + queueCall := m.queueStore.EXPECT().Update(gomock.Any(), queueRow(0, 4), int32(4), int32(5)).Return(nil) + requestCall := m.reqStore.EXPECT().Update(gomock.Any(), requestWithState(entity.RequestStateFailed), int32(1), int32(2)).Return(nil).After(queueCall) + logCall := expectOutcomeLog(m, entity.RequestStateFailed, 2).After(requestCall) + m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil).After(logCall) + + require.NoError(t, c.Process(queueContext(), delivery(t, ctrl, buildSignalPayload(t, testBuildID)))) +} + +func TestReleaseBuildSlotPreservesConcurrentLaterDeadline(t *testing.T) { + ctrl := gomock.NewController(t) + c, m := newController(t, ctrl) + const cooldownMs = int64(60_000) + + first := queueRow(1, 4) + m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(first, nil) + firstUpdate := queueRow(0, 4) + firstUpdate.BuildAdmissionNotBeforeMs = testNowMs + cooldownMs + m.queueStore.EXPECT().Update(gomock.Any(), firstUpdate, int32(4), int32(5)).Return(storage.ErrVersionMismatch) + + concurrent := queueRow(1, 5) + concurrent.BuildAdmissionNotBeforeMs = testNowMs + 2*cooldownMs + m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(concurrent, nil) + secondUpdate := concurrent + secondUpdate.InFlightCount = 0 + m.queueStore.EXPECT().Update(gomock.Any(), secondUpdate, int32(5), int32(6)).Return(nil) + + require.NoError(t, c.releaseBuildSlot(queueContext(), m.store, testQueue, testNowMs, cooldownMs)) +} + // TestPublishRecordCarriesRequestID pins the record contract: the payload and the // message id are the request id, not the build id, so record never has to reach a // Build. The message id is stable so a redelivery dedups instead of enqueuing twice. diff --git a/stovepipe/entity/build.go b/stovepipe/entity/build.go index e278bf904..110a24a5e 100644 --- a/stovepipe/entity/build.go +++ b/stovepipe/entity/build.go @@ -54,9 +54,9 @@ func (s BuildStatus) IsTerminal() bool { } // Build represents a single build triggered for a Request's commit. All -// fields except Status and Version are immutable after creation — build is -// the sole creator (via BuildStore.Create), and buildsignal is the sole -// writer of Status/Version afterward. +// fields except Status, TerminalAtMs, and Version are immutable after creation — +// build is the sole creator (via BuildStore.Create), and buildsignal is the sole +// writer of the lifecycle fields afterward. type Build struct { // ID is the build's own key: the runner-assigned id returned by // Trigger (e.g. a Buildkite build number). Opaque; never parsed or @@ -67,6 +67,9 @@ type Build struct { RequestID string `json:"request_id"` // Status is the build's lifecycle state. Status BuildStatus `json:"status"` + // TerminalAtMs is when the first terminal status was observed, in Unix milliseconds. + // Zero while the build is non-terminal and immutable once set. + TerminalAtMs int64 `json:"terminal_at_ms"` // Version is used for optimistic locking. Versioning starts at 1 and // is incremented for each change to the object. Version int32 `json:"version"` diff --git a/stovepipe/entity/queue_config.go b/stovepipe/entity/queue_config.go index 14ba5c159..59cb6301c 100644 --- a/stovepipe/entity/queue_config.go +++ b/stovepipe/entity/queue_config.go @@ -29,4 +29,7 @@ type QueueConfig struct { // MinimumBuildAdmissionIntervalMs is the minimum start-to-start spacing between logical // build admissions for this queue. Non-positive values disable time-based throttling. MinimumBuildAdmissionIntervalMs int64 `json:"minimum_build_admission_interval_ms" yaml:"minimum_build_admission_interval_ms"` + // FailureCooldownMs is the admission delay applied after the runner reports a failed + // build. Non-positive values disable failure-specific cooldown. + FailureCooldownMs int64 `json:"failure_cooldown_ms" yaml:"failure_cooldown_ms"` } diff --git a/stovepipe/extension/queueconfig/README.md b/stovepipe/extension/queueconfig/README.md index 9ac820909..42bca5916 100644 --- a/stovepipe/extension/queueconfig/README.md +++ b/stovepipe/extension/queueconfig/README.md @@ -10,9 +10,9 @@ Pipeline stages read mutable runtime state from storage and read knobs such as ` ## Entities -Queue configuration entity lives in `stovepipe/entity/queue_config.go` and carries deployment knobs (`max_concurrent`, `gate_wait_delay_ms`, `minimum_build_admission_interval_ms`) separate from the mutable `Queue` row. The minimum interval is start-to-start spacing between logical admissions. Positive values enable the policy; non-positive values disable it. +Queue configuration entity lives in `stovepipe/entity/queue_config.go` and carries deployment knobs (`max_concurrent`, `gate_wait_delay_ms`, `minimum_build_admission_interval_ms`, `failure_cooldown_ms`) separate from the mutable `Queue` row. The minimum interval is start-to-start spacing between logical admissions. The failure cooldown applies after a build runner reports `failed`; cancellation and DLQ-forced failure do not apply it. Positive values enable each time policy; non-positive values disable it. ## Implementations -- `default` returns the global wiring defaults for any non-empty queue name. Time-based admission throttling is disabled. -- `yaml` loads a validated immutable snapshot from a file. Every queue entry specifies its concurrency, gate delay, and general minimum admission interval. +- `default` returns the global wiring defaults for any non-empty queue name. Both admission intervals are disabled. +- `yaml` loads a validated immutable snapshot from a file. Every queue entry specifies its concurrency, gate delay, general minimum admission interval, and failure cooldown. diff --git a/stovepipe/extension/queueconfig/default/default.go b/stovepipe/extension/queueconfig/default/default.go index 29aa7d636..6dbde12ce 100644 --- a/stovepipe/extension/queueconfig/default/default.go +++ b/stovepipe/extension/queueconfig/default/default.go @@ -28,6 +28,7 @@ const ( _defaultMaxConcurrent = 1 _defaultGateWaitDelayMs = 5000 _defaultMinimumBuildAdmissionIntervalMs = 0 + _defaultFailureCooldownMs = 0 ) // Store is a queueconfig.Store that returns the same defaults for every queue. @@ -48,6 +49,7 @@ func (Store) Get(_ context.Context, name string) (entity.QueueConfig, error) { MaxConcurrent: _defaultMaxConcurrent, GateWaitDelayMs: _defaultGateWaitDelayMs, MinimumBuildAdmissionIntervalMs: _defaultMinimumBuildAdmissionIntervalMs, + FailureCooldownMs: _defaultFailureCooldownMs, }, nil } diff --git a/stovepipe/extension/queueconfig/default/default_test.go b/stovepipe/extension/queueconfig/default/default_test.go index 5df4e50fd..d60e181de 100644 --- a/stovepipe/extension/queueconfig/default/default_test.go +++ b/stovepipe/extension/queueconfig/default/default_test.go @@ -33,6 +33,7 @@ func TestStore_Get(t *testing.T) { assert.Equal(t, int32(1), cfg.MaxConcurrent) assert.Equal(t, int64(5000), cfg.GateWaitDelayMs) assert.Zero(t, cfg.MinimumBuildAdmissionIntervalMs) + assert.Zero(t, cfg.FailureCooldownMs) }) t.Run("empty name is not found", func(t *testing.T) { diff --git a/stovepipe/extension/queueconfig/yaml/yaml_test.go b/stovepipe/extension/queueconfig/yaml/yaml_test.go index f4dc27d0e..622689b93 100644 --- a/stovepipe/extension/queueconfig/yaml/yaml_test.go +++ b/stovepipe/extension/queueconfig/yaml/yaml_test.go @@ -48,19 +48,21 @@ func TestNewStore(t *testing.T) { max_concurrent: 2 gate_wait_delay_ms: 5000 minimum_build_admission_interval_ms: 3600000 + failure_cooldown_ms: 900000 `, }, {name: "empty list", content: "queues: []\n"}, {name: "empty document", content: "", wantErr: true}, {name: "malformed YAML", content: "queues: [", wantErr: true}, - {name: "unknown field", content: validQueueYAML("main", 1, 5000, 0) + "unexpected: true\n", wantErr: true}, - {name: "empty name", content: validQueueYAML("", 1, 5000, 0), wantErr: true}, - {name: "non-positive concurrency", content: validQueueYAML("main", 0, 5000, 0), wantErr: true}, - {name: "non-positive gate delay", content: validQueueYAML("main", 1, 0, 0), wantErr: true}, - {name: "negative minimum interval disables policy", content: validQueueYAML("main", 1, 5000, -1)}, + {name: "unknown field", content: validQueueYAML("main", 1, 5000, 0, 0) + "unexpected: true\n", wantErr: true}, + {name: "empty name", content: validQueueYAML("", 1, 5000, 0, 0), wantErr: true}, + {name: "non-positive concurrency", content: validQueueYAML("main", 0, 5000, 0, 0), wantErr: true}, + {name: "non-positive gate delay", content: validQueueYAML("main", 1, 0, 0, 0), wantErr: true}, + {name: "negative minimum interval disables policy", content: validQueueYAML("main", 1, 5000, -1, 0)}, + {name: "negative failure cooldown disables policy", content: validQueueYAML("main", 1, 5000, 0, -1)}, { name: "duplicate name", - content: validQueueYAML("main", 1, 5000, 0) + ` - name: main + content: validQueueYAML("main", 1, 5000, 0, 0) + ` - name: main max_concurrent: 1 gate_wait_delay_ms: 5000 `, @@ -82,17 +84,18 @@ func TestNewStore(t *testing.T) { } } -func validQueueYAML(name string, maxConcurrent int32, gateWaitDelayMs, minimumIntervalMs int64) string { +func validQueueYAML(name string, maxConcurrent int32, gateWaitDelayMs, minimumIntervalMs, failureCooldownMs int64) string { return fmt.Sprintf(`queues: - name: %s max_concurrent: %d gate_wait_delay_ms: %d minimum_build_admission_interval_ms: %d -`, name, maxConcurrent, gateWaitDelayMs, minimumIntervalMs) + failure_cooldown_ms: %d +`, name, maxConcurrent, gateWaitDelayMs, minimumIntervalMs, failureCooldownMs) } func TestStoreGetAndList(t *testing.T) { - store, err := NewStore(writeQueueConfig(t, validQueueYAML("monorepo/main", 2, 5000, 3_600_000))) + store, err := NewStore(writeQueueConfig(t, validQueueYAML("monorepo/main", 2, 5000, 3_600_000, 900_000))) require.NoError(t, err) got, err := store.Get(context.Background(), "monorepo/main") @@ -102,6 +105,7 @@ func TestStoreGetAndList(t *testing.T) { MaxConcurrent: 2, GateWaitDelayMs: 5000, MinimumBuildAdmissionIntervalMs: 3_600_000, + FailureCooldownMs: 900_000, }, got) _, err = store.Get(context.Background(), "missing") diff --git a/stovepipe/extension/storage/mysql/build_store.go b/stovepipe/extension/storage/mysql/build_store.go index ab079a04a..9bd0dda47 100644 --- a/stovepipe/extension/storage/mysql/build_store.go +++ b/stovepipe/extension/storage/mysql/build_store.go @@ -46,12 +46,13 @@ func (b *buildStore) Create(ctx context.Context, build entity.Build) (retErr err defer func() { op.Complete(retErr) }() _, err := b.db.ExecContext(ctx, - `INSERT INTO build (queue, id, request_id, status, version) - VALUES (?, ?, ?, ?, ?)`, + `INSERT INTO build (queue, id, request_id, status, terminal_at_ms, version) + VALUES (?, ?, ?, ?, ?, ?)`, b.queue, build.ID, build.RequestID, build.Status, + build.TerminalAtMs, build.Version, ) if err != nil { @@ -71,13 +72,14 @@ func (b *buildStore) Get(ctx context.Context, id string) (ret entity.Build, retE var build entity.Build err := b.db.QueryRowContext(ctx, - `SELECT id, request_id, status, version + `SELECT id, request_id, status, terminal_at_ms, version FROM build WHERE queue = ? AND id = ?`, b.queue, id, ).Scan( &build.ID, &build.RequestID, &build.Status, + &build.TerminalAtMs, &build.Version, ) @@ -101,9 +103,10 @@ func (b *buildStore) Update(ctx context.Context, build entity.Build, oldVersion, result, err := b.db.ExecContext(ctx, `UPDATE build - SET status = ?, version = ? + SET status = ?, terminal_at_ms = ?, version = ? WHERE queue = ? AND id = ? AND version = ?`, build.Status, + build.TerminalAtMs, newVersion, b.queue, build.ID, diff --git a/stovepipe/extension/storage/mysql/build_store_test.go b/stovepipe/extension/storage/mysql/build_store_test.go index 6e6e59118..3b7bb9469 100644 --- a/stovepipe/extension/storage/mysql/build_store_test.go +++ b/stovepipe/extension/storage/mysql/build_store_test.go @@ -42,10 +42,11 @@ func setupBuildStoreTest(t *testing.T) (*sql.DB, sqlmock.Sqlmock, storage.BuildS func TestBuildStore_Create(t *testing.T) { build := entity.Build{ - ID: "bk-1001", - RequestID: "request/monorepo/main/1", - Status: entity.BuildStatusAccepted, - Version: 1, + ID: "bk-1001", + RequestID: "request/monorepo/main/1", + Status: entity.BuildStatusAccepted, + TerminalAtMs: 1234, + Version: 1, } tests := []struct { @@ -58,7 +59,7 @@ func TestBuildStore_Create(t *testing.T) { name: "success", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("INSERT INTO build"). - WithArgs("monorepo/main", build.ID, build.RequestID, build.Status, build.Version). + WithArgs("monorepo/main", build.ID, build.RequestID, build.Status, build.TerminalAtMs, build.Version). WillReturnResult(sqlmock.NewResult(0, 1)) }, }, @@ -66,7 +67,7 @@ func TestBuildStore_Create(t *testing.T) { name: "duplicate id returns ErrAlreadyExists", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("INSERT INTO build"). - WithArgs("monorepo/main", build.ID, build.RequestID, build.Status, build.Version). + WithArgs("monorepo/main", build.ID, build.RequestID, build.Status, build.TerminalAtMs, build.Version). WillReturnError(&mysql.MySQLError{Number: mysqlErrDuplicateEntry}) }, wantErr: true, @@ -76,7 +77,7 @@ func TestBuildStore_Create(t *testing.T) { name: "other exec error", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("INSERT INTO build"). - WithArgs("monorepo/main", build.ID, build.RequestID, build.Status, build.Version). + WithArgs("monorepo/main", build.ID, build.RequestID, build.Status, build.TerminalAtMs, build.Version). WillReturnError(fmt.Errorf("connection reset")) }, wantErr: true, @@ -106,10 +107,11 @@ func TestBuildStore_Create(t *testing.T) { func TestBuildStore_Get(t *testing.T) { want := entity.Build{ - ID: "bk-1001", - RequestID: "request/monorepo/main/1", - Status: entity.BuildStatusRunning, - Version: 2, + ID: "bk-1001", + RequestID: "request/monorepo/main/1", + Status: entity.BuildStatusRunning, + TerminalAtMs: 1234, + Version: 2, } tests := []struct { @@ -124,9 +126,9 @@ func TestBuildStore_Get(t *testing.T) { name: "found", id: want.ID, setup: func(mock sqlmock.Sqlmock) { - rows := sqlmock.NewRows([]string{"id", "request_id", "status", "version"}). - AddRow(want.ID, want.RequestID, string(want.Status), want.Version) - mock.ExpectQuery("SELECT id, request_id, status, version"). + rows := sqlmock.NewRows([]string{"id", "request_id", "status", "terminal_at_ms", "version"}). + AddRow(want.ID, want.RequestID, string(want.Status), want.TerminalAtMs, want.Version) + mock.ExpectQuery("SELECT id, request_id, status, terminal_at_ms, version"). WithArgs("monorepo/main", want.ID). WillReturnRows(rows) }, @@ -136,7 +138,7 @@ func TestBuildStore_Get(t *testing.T) { name: "not found", id: "missing", setup: func(mock sqlmock.Sqlmock) { - mock.ExpectQuery("SELECT id, request_id, status, version"). + mock.ExpectQuery("SELECT id, request_id, status, terminal_at_ms, version"). WithArgs("monorepo/main", "missing"). WillReturnError(sql.ErrNoRows) }, @@ -147,7 +149,7 @@ func TestBuildStore_Get(t *testing.T) { name: "query error", id: "bad", setup: func(mock sqlmock.Sqlmock) { - mock.ExpectQuery("SELECT id, request_id, status, version"). + mock.ExpectQuery("SELECT id, request_id, status, terminal_at_ms, version"). WithArgs("monorepo/main", "bad"). WillReturnError(fmt.Errorf("connection reset")) }, @@ -178,7 +180,7 @@ func TestBuildStore_Get(t *testing.T) { } func TestBuildStore_Update(t *testing.T) { - build := entity.Build{ID: "bk-1001", Status: entity.BuildStatusRunning} + build := entity.Build{ID: "bk-1001", Status: entity.BuildStatusRunning, TerminalAtMs: 1234} const oldVersion, newVersion = int32(1), int32(2) tests := []struct { @@ -191,7 +193,7 @@ func TestBuildStore_Update(t *testing.T) { name: "success", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("UPDATE build"). - WithArgs(build.Status, newVersion, "monorepo/main", build.ID, oldVersion). + WithArgs(build.Status, build.TerminalAtMs, newVersion, "monorepo/main", build.ID, oldVersion). WillReturnResult(sqlmock.NewResult(0, 1)) }, }, @@ -199,7 +201,7 @@ func TestBuildStore_Update(t *testing.T) { name: "version mismatch", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("UPDATE build"). - WithArgs(build.Status, newVersion, "monorepo/main", build.ID, oldVersion). + WithArgs(build.Status, build.TerminalAtMs, newVersion, "monorepo/main", build.ID, oldVersion). WillReturnResult(sqlmock.NewResult(0, 0)) }, wantErr: true, @@ -209,7 +211,7 @@ func TestBuildStore_Update(t *testing.T) { name: "exec error", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("UPDATE build"). - WithArgs(build.Status, newVersion, "monorepo/main", build.ID, oldVersion). + WithArgs(build.Status, build.TerminalAtMs, newVersion, "monorepo/main", build.ID, oldVersion). WillReturnError(fmt.Errorf("connection reset")) }, wantErr: true, @@ -218,7 +220,7 @@ func TestBuildStore_Update(t *testing.T) { name: "rows affected error", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("UPDATE build"). - WithArgs(build.Status, newVersion, "monorepo/main", build.ID, oldVersion). + WithArgs(build.Status, build.TerminalAtMs, newVersion, "monorepo/main", build.ID, oldVersion). WillReturnResult(sqlmock.NewErrorResult(fmt.Errorf("driver error"))) }, wantErr: true, diff --git a/stovepipe/extension/storage/mysql/schema/build.sql b/stovepipe/extension/storage/mysql/schema/build.sql index 4c60f17c1..eddfaf01e 100644 --- a/stovepipe/extension/storage/mysql/schema/build.sql +++ b/stovepipe/extension/storage/mysql/schema/build.sql @@ -6,5 +6,6 @@ CREATE TABLE IF NOT EXISTS build ( request_id VARCHAR(255) NOT NULL, status VARCHAR(64) NOT NULL, version INT NOT NULL, + terminal_at_ms BIGINT NOT NULL DEFAULT 0, PRIMARY KEY (queue, id) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; diff --git a/test/integration/stovepipe/extension/storage/suite.go b/test/integration/stovepipe/extension/storage/suite.go index b5d5a3ff1..fae0f94d9 100644 --- a/test/integration/stovepipe/extension/storage/suite.go +++ b/test/integration/stovepipe/extension/storage/suite.go @@ -413,11 +413,13 @@ func (s *BuildStoreContractSuite) TestBuildStore_UpdateCAS() { updated := created updated.Status = entity.BuildStatusRunning + updated.TerminalAtMs = 1234 require.NoError(t, s.buildStore.Update(s.ctx, updated, 1, 2)) got, err := s.buildStore.Get(s.ctx, id) require.NoError(t, err) assert.Equal(t, entity.BuildStatusRunning, got.Status) + assert.EqualValues(t, 1234, got.TerminalAtMs) assert.Equal(t, int32(2), got.Version) err = s.buildStore.Update(s.ctx, updated, 1, 2)