From c354c8ec8146dddf173b906612c8d1dde793d2e2 Mon Sep 17 00:00:00 2001 From: "prath.shenoy" Date: Thu, 24 Sep 2026 23:17:30 +0000 Subject: [PATCH] feat(stovepipe): Persist terminal build identity Summary: Persist the winning terminal build identity in the terminal request-history entry, atomically with the request outcome, so record can later resolve final artifacts without a Build reverse index. Test Plan: ./tool/bazel test //stovepipe/controller/buildsignal:go_default_test //stovepipe/core/messagequeue:go_default_test //stovepipe/extension/storage/mysql:go_default_test --test_output=errors Revert Plan: Revert this commit to stop persisting terminal build identity. Jira Issues: None --- doc/rfc/stovepipe/request-log.md | 4 +- doc/rfc/stovepipe/steps/build.md | 16 ++-- doc/rfc/stovepipe/steps/buildsignal.md | 12 ++- doc/rfc/stovepipe/steps/record.md | 2 +- .../controller/buildsignal/buildsignal.go | 87 ++++++++++-------- .../buildsignal/buildsignal_test.go | 91 +++++++++++++------ .../storage/mock/request_store_mock.go | 14 +++ .../extension/storage/mysql/request_store.go | 64 +++++++++++++ .../storage/mysql/request_store_test.go | 52 +++++++++++ stovepipe/extension/storage/request_store.go | 6 ++ 10 files changed, 270 insertions(+), 78 deletions(-) diff --git a/doc/rfc/stovepipe/request-log.md b/doc/rfc/stovepipe/request-log.md index 761159c67..a2a7ba636 100644 --- a/doc/rfc/stovepipe/request-log.md +++ b/doc/rfc/stovepipe/request-log.md @@ -191,7 +191,7 @@ Metadata is serialized as a JSON object and normalized to an empty map on write ## Write and Repair Protocol -Request-log durability is part of completing a pipeline transition. The source write succeeds first, the required log record is retained second, and a dependent handoff is published only after log creation or identical-existing reconciliation succeeds. +Request-log durability is part of completing a pipeline transition. The source write succeeds first, the required log record is retained second, and a dependent handoff is published only after log creation or identical-existing reconciliation succeeds. The terminal build outcome is the exception: its Request transition and terminal state entry commit atomically, because that entry selects the winning build. For a Request transition, the controller: @@ -211,7 +211,7 @@ Request creation, Build changes, and fact creation use the same source-write, lo | Ingest | Create accepted Request, then retain accepted. | An existing Request ensures accepted before process publication. | | Process | CAS to superseded or processing, then retain that state. | An existing state is reconstructed from Request context before ack or build publication. | | Build | Create Build after runner acceptance, then retain `build_triggered`. | An identical existing Build ensures the event before buildsignal publication. | -| Buildsignal | Persist terminal Build and retain `build_finished`; CAS the Request outcome and retain its terminal state. | Existing terminal Build and Request outcome each ensure their own entry before record publication. | +| Buildsignal | Persist terminal Build and retain `build_finished`; atomically CAS the Request outcome and its terminal state entry, including the winning `build_id`. | Existing terminal Build and Request outcome each ensure their own entry before record publication. | | Record | Create or verify the whole-repository fact, then retain `validation_fact_recorded`; after recording project facts, retain one `project_facts_recorded` event. | An identical whole-repository fact and project-fact batch ensure their events before bookmark or promotion work. | | Record DLQ | Retain `record_abandoned`, then acknowledge the remaining record work. | The stable event ID makes history retention idempotent without replaying facts, bookmarks, promotion, or hooks. | | Reconciler | CAS an unrecoverable non-terminal Request to failed, then retain failed. | An existing terminal Request is repaired from its persisted outcome without relabeling it. | diff --git a/doc/rfc/stovepipe/steps/build.md b/doc/rfc/stovepipe/steps/build.md index f630870b2..c4d249a9c 100644 --- a/doc/rfc/stovepipe/steps/build.md +++ b/doc/rfc/stovepipe/steps/build.md @@ -54,8 +54,10 @@ For a delivery carrying request id `R`: 6. Persist Build{ID: buildID.ID, RequestID: R.ID, Status: accepted, Version: 1} via BuildStore.Create. - - the row carries no scope; it is recoverable from the Request's immutable fields - (see the entity table). + - the row carries no validation scope; that is recoverable from the Request's immutable + fields (see the entity table). `Build.RequestID` remains the only Build-to-Request link. + When a build wins the terminal-outcome race, buildsignal records its id in the immutable + terminal request-history entry. - a crash between step 5 and this write orphans the triggered build (see Idempotency). - ErrAlreadyExists -> benign (reachable only with a backend that returns deterministic ids for retried triggers); continue to step 7. @@ -69,7 +71,7 @@ For a delivery carrying request id `R`: 8. ack. ``` -`Build.ID` is **minted by the runner at `Trigger`**, exactly as in SubmitQueue's build controller: the runner returns its native id (a Buildkite build number, a CI-gateway job id), `build` adopts it as the `Build`'s key, and that same id travels on every hop that needs a build — `build` → `buildsignal` carries the build id in the message, so the poll loop reaches the `Build` by a direct get on identity it was handed. `buildsignal` → `record` carries the **request id** instead: `record`'s unit of work is a Request, and the build's terminal status is projected onto `Request.State` before the publish, so `record` never reaches a `Build` at all. No reader ever *derives* a build id or needs a reverse index; `Build.RequestID` covers the one navigation the pipeline needs in the other direction (`Build` → `Request`). Another approach is deriving the key from the Request (`buildKey(R)`) and/or passing a caller-supplied idempotency key to `Trigger`; see [Alternatives considered](#alternatives-considered-for-the-build-identity) for what each would buy and cost. +`Build.ID` is **minted by the runner at `Trigger`**, exactly as in SubmitQueue's build controller: the runner returns its native id (a Buildkite build number, a CI-gateway job id), `build` adopts it as the `Build`'s key, and that same id travels on every hop that needs a build — `build` → `buildsignal` carries the build id in the message, so the poll loop reaches the `Build` by a direct get on identity it was handed. `buildsignal` → `record` carries the **request id** instead: `record`'s unit of work is a Request, and buildsignal atomically retains the winning build id in that Request's terminal history entry before publishing. A later artifact reader can resolve that entry and load the selected Build; it neither derives a build id nor queries builds by RequestID. Another approach is deriving the key from the Request (`buildKey(R)`) and/or passing a caller-supplied idempotency key to `Trigger`; see [Alternatives considered](#alternatives-considered-for-the-build-identity) for what each would buy and cost. `build` writes only the `Build`, and only at creation; it never mutates `Request.State`. The Request stays `processing` (set by `process`) through `build` until `buildsignal` moves it terminal by recording the build's outcome. `Build.Status` is the fine-grained build lifecycle; `Request.State` is the coarse pipeline lifecycle. This is also what keeps `process.md` step 3 correct: because `build` leaves the Request at `processing`, a redelivered `process` message still matches its "if processing, re-publish to build" guard. @@ -202,7 +204,7 @@ Trigger(ctx context.Context, baseURI, headURI string, projectScope entity.Projec Both `Trigger` and `Status`/`Cancel` differ *in contract* between domains, even though `Status`/`Cancel` happen to be identical in shape: both domains poll and cancel by the same opaque, runner-minted id with the same async semantics. Rather than promoting that shape parity into a shared `platform/base`/`platform/extension/buildrunner` type and interface — which would force a one-time migration of SubmitQueue's already-shipped controllers, storage, and protobuf mappings onto the shared type — each domain keeps its own `BuildRunner` interface and its own local `BuildID`/`BuildStatus`/`BuildMetadata`, and real code reuse happens one layer down, in a shared backend implementation (e.g. a Buildkite client) that both domains' concrete runners wrap. [Alternatives considered for sharing the contract](#alternatives-considered-for-sharing-the-contract) below records the shapes weighed against this one, including the shared-interface alternative that was set aside. -There is exactly one build id: the runner mints it at `Trigger`, `build` adopts it as `Build.ID`, and every later call and message carries it verbatim — `Status`/`Cancel` take the same value `Trigger` returned, the queue payload is the same value, the store key is the same value. This is SubmitQueue's convention end to end. The id is opaque: no stovepipe reader parses it, derives it, or equates it with another entity's id — the trap SubmitQueue's speculate/cancel path falls into. And per the extension rules a runner keeps only transient local state, so the durable `Request` ↔ `Build` linkage lives in **our** store as `Build.RequestID`, never in the runner. +There is exactly one build id: the runner mints it at `Trigger`, `build` adopts it as `Build.ID`, and every later call and message carries it verbatim — `Status`/`Cancel` take the same value `Trigger` returned, the queue payload is the same value, the store key is the same value. This is SubmitQueue's convention end to end. The id is opaque: no stovepipe reader parses it, derives it, or equates it with another entity's id. The terminal request-history state entry retains the winning id alongside the terminal Request state; it is selected by buildsignal's outcome CAS, not a derived key or a reverse index. Per the extension rules a runner keeps only transient local state, so that entry and `Build.RequestID` live in **our** store, never in the runner. Supporting entity types: `BuildStatus`, `BuildMetadata`, and `BuildID` live in `stovepipe/entity`, shaped the same as SubmitQueue's `submitqueue/entity` equivalents but defined and duplicated locally rather than shared — `BuildStatus` is the narrow lowercase enum `"" (unknown) / accepted / running / succeeded / failed / cancelled` with an `IsTerminal()` predicate covering the last three, `BuildMetadata` is the free-form `map[string]string`, and `BuildID` is a `{ID string}` wire struct wrapping the one runner-assigned id everywhere it appears — `Trigger`'s return, `Status`/`Cancel`'s parameter, the queue payload. `stovepipe/entity/build.go` keeps what's stovepipe-specific: the `Build` entity itself (`RequestID` alongside `ID`/`Status`/`Version`). How a target graph reaches `analyze` is out of scope for this doc — left to the `analyze` design. @@ -266,12 +268,12 @@ Key the `Build` by identity derived from the Request — `buildKey(R) = R.ID` fo | Pros | Cons | |---|---| | Redelivery dedup by direct get: checking `BuildStore.Get(buildKey(R))` before triggering means at-least-once delivery never starts a second build | A second id concept (`Build.ID` beside `Build.RunnerBuildID`) carried by every entity, signature, and reader forever | -| `Request` → `Build` navigation with no reverse index, per the KV key-derivation rule in [AGENTS.md](AGENTS.md) | No current reader needs to *derive* a build id — the id travels in every message hop, so each consumer already holds the key it needs | +| `Request` → `Build` navigation with no reverse index, per the KV key-derivation rule in [AGENTS.md](AGENTS.md) | `record` needs the winning build's id, but the terminal state entry selects it atomically with the outcome; changing the Build key would still add a second identity and would not remove that winner-selection state | | Enforces (rather than assumes) the direct-navigation property SubmitQueue's speculate takes on faith | Diverges entity shape and controller flow from SubmitQueue, weakening the "structurally the same controller" claim and dual-implementing-backend symmetry | Trade-offs: the dedup guards a rare event at a permanent modeling cost. The duplicate it prevents arises only from a redelivery inside the trigger window — rare, and already harmless (identical scope; `buildsignal`'s superseded short-circuit and its first-writer-wins outcome CAS make the loser a no-op — see [Idempotency](#idempotency)). The prospective key-derivers — a future canceller, or `analyze` reaching back to the Phase-1 target graph — would need to be handed the id by their producing stage instead, if those designs land. -Note that moving `record`'s input from the build id to the request id does *not* trigger this alternative, even though it removes the last hop that carried a build id to a Request-scoped consumer. The trigger condition is a stage that must **derive a build's key from a Request**, and `record` does not: the build's terminal status is projected onto `Request.State` before the publish, so `record` reads the Request and never reaches a `Build`. +Moving `record`'s input from the build id to the request id now requires the winning identity to survive that handoff. The trigger condition for this alternative is still a stage that must **derive a build's key from a Request**. `record` does not: buildsignal retains the terminal state and its winning build id before publishing, so `record` can resolve that state entry without a reverse Build index. #### Alternative B: caller-supplied idempotency key on `Trigger` @@ -320,7 +322,7 @@ The row deliberately carries no scope: `R.URI`, `R.BaseURI`, and `R.BuildStrateg `IsTerminal()` on `entity.BuildStatus` covers exactly the three terminal rows. Once `buildsignal` persists one of them, that status is **write-once** — a later poll reporting a different terminal value never overwrites it (see [buildsignal.md](doc/rfc/stovepipe/steps/buildsignal.md#algorithm), step 6). -Plus the `BuildID{ID string}` wire type in `stovepipe/entity` (same "id only travels" convention as `RequestID`, shaped like SubmitQueue's own `entity.BuildID` but not the same Go type — see the [contract sketch](#stovepipe-buildrunner-contract-design-sketch)), wrapping the one runner-assigned id everywhere it appears — `Trigger`'s return, the queue payload, `Status`/`Cancel`'s parameter. `buildsignal` reaches a build by the id carried in its message, and `record` reads the `Request` (whose state carries the build's outcome) rather than a `Build`, so no reverse index from `Request` to its builds is ever needed. +Plus the `BuildID{ID string}` wire type in `stovepipe/entity` (same "id only travels" convention as `RequestID`, shaped like SubmitQueue's own `entity.BuildID` but not the same Go type — see the [contract sketch](#stovepipe-buildrunner-contract-design-sketch)), wrapping the one runner-assigned id everywhere it appears — `Trigger`'s return, the queue payload, `Status`/`Cancel`'s parameter. `buildsignal` reaches a build by the id carried in its message, then atomically retains the winning id in the terminal request-history state entry. An artifact reader resolves that entry instead of querying a Build by request. **`BuildStore`** (new, added to the `Storage` aggregator via `GetBuildStore()`), matching stovepipe's existing `RequestStore` conventions — **generic `Update` with caller-owned version arithmetic**: diff --git a/doc/rfc/stovepipe/steps/buildsignal.md b/doc/rfc/stovepipe/steps/buildsignal.md index 7ce0dfcaa..a84d4417d 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.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. That write atomically retains the matching terminal history entry, including the build identity — 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. @@ -63,12 +63,14 @@ For a delivery carrying build id `B`: 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. - 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: + b. Atomically CAS R from processing to the outcome the stored status projects onto it + and retain R's terminal request-history state entry with B as its `build_id` metadata: 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. - - record loads the Request directly by this key and derives greenness from its outcome, - so it never reaches a Build and no reverse lookup from Request to its builds is needed. + - record loads the Request directly by this key and derives greenness from its outcome. A + later artifact reader loads the terminal state entry by request version to obtain the + selected build id; no reverse lookup from Request to its builds is needed. - the message id is the request id, so a redelivery republishing the same terminal signal dedups into the original message; record is idempotent regardless. - publish failure -> return raw (non-retryable); the outcome is persisted, operational @@ -164,4 +166,4 @@ One boundary of that posture is worth stating: the `MaxAttempts` path fires only ## Entity, storage, and queue additions -No additions beyond [build.md](doc/rfc/stovepipe/steps/build.md#entity-and-storage-additions-needed): `buildsignal` calls `BuildStore.Get`/`Update` and `RequestStore.Get`/`Update` against the `Build`/`Request` shapes defined there — `Build.ID` being the runner-assigned id it hands straight back to `Status` — plus `QueueStore.Get`/`Update` to release the build slot, and consumes/re-produces the `BuildSignal` message on `TopicKeyBuildSignal` introduced there. `Request.State` gains the three build outcomes (`succeeded`, `failed`, `cancelled`), all terminal. The message it publishes to `record`, and the `record` topic key itself, are owned by the `record` stage and land with `record.md`; `buildsignal` only needs that the **request id** reaches the record topic once the build is terminal, partitioned by request id. +No additions beyond [build.md](doc/rfc/stovepipe/steps/build.md#entity-and-storage-additions-needed): `buildsignal` calls `BuildStore.Get`/`Update` and `RequestStore.Get`/`FinalizeOutcome` against the `Build`/`Request` shapes defined there — `Build.ID` being the runner-assigned id it hands straight back to `Status` — plus `QueueStore.Get`/`Update` to release the build slot, and consumes/re-produces the `BuildSignal` message on `TopicKeyBuildSignal` introduced there. `Request.State` gains the three build outcomes (`succeeded`, `failed`, `cancelled`), all terminal. The message it publishes to `record`, and the `record` topic key itself, are owned by the `record` stage and land with `record.md`; `buildsignal` only needs that the **request id** reaches the record topic once the build is terminal, partitioned by request id. diff --git a/doc/rfc/stovepipe/steps/record.md b/doc/rfc/stovepipe/steps/record.md index 435cebeb1..af9fd3812 100644 --- a/doc/rfc/stovepipe/steps/record.md +++ b/doc/rfc/stovepipe/steps/record.md @@ -206,7 +206,7 @@ Hooks here must be idempotent on `id`, as everywhere. "Fire-and-forget" describe ## Request lifecycle -Phase 1 uses the states in [stovepipe/entity/request.go](../../../../stovepipe/entity/request.go). `record` runs *after* the Request is terminal: `buildsignal` projects the build's terminal status onto it as `succeeded`, `failed`, or `cancelled` (`RequestState.HasBuildOutcome()`), and only then publishes. So `record` reads an outcome and writes no state. `superseded` is terminal without an outcome. +Phase 1 uses the states in [stovepipe/entity/request.go](../../../../stovepipe/entity/request.go). `record` runs *after* the Request is terminal: `buildsignal` projects the build's terminal status onto it as `succeeded`, `failed`, or `cancelled` (`RequestState.HasBuildOutcome()`), atomically retaining the matching terminal request-history state entry with the winning build id, and only then publishes. So `record` reads an outcome and writes no state. `superseded` is terminal without an outcome. Phase 2 broadens "complete" to "all planned facts recorded", which needs a marker this stage does not own; see [Completion marker: open](#completion-marker-open). diff --git a/stovepipe/controller/buildsignal/buildsignal.go b/stovepipe/controller/buildsignal/buildsignal.go index 2c56a04b8..8c517cfb6 100644 --- a/stovepipe/controller/buildsignal/buildsignal.go +++ b/stovepipe/controller/buildsignal/buildsignal.go @@ -177,10 +177,11 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er if err := c.persistBuildFinishedLog(ctx, store, request, build.ID); err != nil { return err } - if err := c.finishRequest(ctx, store, &request, effective); err != nil { + terminalStateLog, err := c.finishRequest(ctx, store, &request, effective, build.ID) + if err != nil { return err } - if err := c.persistOutcomeLog(ctx, store, request); err != nil { + if err := c.persistOutcomeLog(ctx, store, terminalStateLog); err != nil { return err } if err := c.publishRecord(ctx, request.ID, request.Queue); err != nil { @@ -234,40 +235,29 @@ 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, status entity.BuildStatus, buildID string) (entity.RequestLog, error) { if request.State.HasBuildOutcome() { - return nil + return c.existingOutcomeLog(ctx, store, *request) } if err := c.releaseBuildSlot(ctx, store, request.Queue); err != nil { metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1, metrics.TagsFromContext(ctx)...) - return err + return entity.RequestLog{}, err } - if err := c.markOutcome(ctx, store, request, outcomeState(status)); err != nil { + log, err := c.markOutcome(ctx, store, request, outcomeState(status), buildID) + if err != nil { metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1, metrics.TagsFromContext(ctx)...) - return err + return entity.RequestLog{}, err } - return nil + return log, 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. - switch request.State { - case entity.RequestStateSucceeded: - reason = entity.RequestOutcomeReasonBuildSucceeded - case entity.RequestStateFailed: - reason = entity.RequestOutcomeReasonBuildFailed - case entity.RequestStateCancelled: - reason = entity.RequestOutcomeReasonBuildCancelled - default: - return fmt.Errorf("request %s has no build outcome to record", request.ID) - } - - log := requestlog.NewRequestStateLog(request, reason) +// persistOutcomeLog reconciles the state entry the outcome transaction already +// retained and updates its derived read models. +func (c *Controller) persistOutcomeLog(ctx context.Context, store storage.Storage, log entity.RequestLog) error { if err := c.materializer.PersistLog(ctx, store, log); err != nil { - return fmt.Errorf("failed to record %s state for request %s: %w", request.State, request.ID, err) + return fmt.Errorf("failed to materialize terminal state for request %s: %w", log.RequestID, err) } return nil } @@ -286,38 +276,63 @@ func outcomeState(status entity.BuildStatus) entity.RequestState { } } -// markOutcome CAS-transitions request from processing to state, retrying on version -// conflicts. First writer wins: once any outcome is recorded a later caller leaves it -// alone, so duplicate builds for one request (which build.md accepts) cannot flip the -// verdict back and forth. -func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, request *entity.Request, state entity.RequestState) error { +// markOutcome CAS-transitions request from processing to state and atomically +// retains the corresponding state entry. First writer wins: duplicate builds cannot +// flip the verdict or replace the build identity retained by that entry. +func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, request *entity.Request, state entity.RequestState, buildID string) (entity.RequestLog, error) { reqStore := store.GetRequestStore() for { if request.State != entity.RequestStateProcessing { - return nil + return c.existingOutcomeLog(ctx, store, *request) } updated := *request updated.State = state - newVersion := request.Version + 1 - if err := reqStore.Update(ctx, updated, request.Version, newVersion); err != nil { + updated.Version = request.Version + 1 + log := requestlog.NewRequestStateLog(updated, outcomeReason(updated.State)) + log.Metadata[requestlog.MetadataKeyBuildID] = buildID + if err := reqStore.FinalizeOutcome(ctx, updated, request.Version, updated.Version, log); err != nil { if errors.Is(err, storage.ErrVersionMismatch) { got, getErr := reqStore.Get(ctx, request.ID) if getErr != nil { - return fmt.Errorf("failed to reload request %s after version mismatch: %w", request.ID, getErr) + return entity.RequestLog{}, fmt.Errorf("failed to reload request %s after version mismatch: %w", request.ID, getErr) } *request = got continue } - return fmt.Errorf("failed to mark request %s %s: %w", request.ID, state, err) + return entity.RequestLog{}, fmt.Errorf("failed to mark request %s %s: %w", request.ID, state, err) } - updated.Version = newVersion *request = updated metrics.NamedCounter(c.metricsScope, _opName, "outcomes", 1, metrics.TagsFromContext(ctx, metrics.NewTag("state", string(state)))..., ) - return nil + return log, nil + } +} + +func (c *Controller) existingOutcomeLog(ctx context.Context, store storage.Storage, request entity.Request) (entity.RequestLog, error) { + if !request.State.HasBuildOutcome() { + return entity.RequestLog{}, fmt.Errorf("request %s has no build outcome", request.ID) + } + logID := requestlog.NewRequestStateLog(request, outcomeReason(request.State)).ID + log, err := store.GetRequestLogStore().Get(ctx, request.ID, logID) + if err != nil { + return entity.RequestLog{}, fmt.Errorf("failed to load terminal state for request %s: %w", request.ID, err) + } + return log, nil +} + +func outcomeReason(state entity.RequestState) entity.RequestOutcomeReason { + switch state { + case entity.RequestStateSucceeded: + return entity.RequestOutcomeReasonBuildSucceeded + case entity.RequestStateFailed: + return entity.RequestOutcomeReasonBuildFailed + case entity.RequestStateCancelled: + return entity.RequestOutcomeReasonBuildCancelled + default: + return entity.RequestOutcomeReasonUnknown } } diff --git a/stovepipe/controller/buildsignal/buildsignal_test.go b/stovepipe/controller/buildsignal/buildsignal_test.go index 0bfb76820..60d9d266d 100644 --- a/stovepipe/controller/buildsignal/buildsignal_test.go +++ b/stovepipe/controller/buildsignal/buildsignal_test.go @@ -54,15 +54,16 @@ func queueContext() context.Context { // buildsignalMocks bundles the mocks a buildsignal controller test case wires // expectations on. type buildsignalMocks struct { - reqStore *storagemock.MockRequestStore - buildStore *storagemock.MockBuildStore - queueStore *storagemock.MockQueueStore - store *storagemock.MockStorage - materializer *requestlogmock.MockMaterializer - runnerFactory *buildrunnermock.MockFactory - runner *buildrunnermock.MockBuildRunner - publisher *mqmock.MockPublisher - metricsScope tally.TestScope + reqStore *storagemock.MockRequestStore + requestLogStore *storagemock.MockRequestLogStore + buildStore *storagemock.MockBuildStore + queueStore *storagemock.MockQueueStore + store *storagemock.MockStorage + materializer *requestlogmock.MockMaterializer + runnerFactory *buildrunnermock.MockFactory + runner *buildrunnermock.MockBuildRunner + publisher *mqmock.MockPublisher + metricsScope tally.TestScope } // staticStorageFactory resolves every queue to one fixed store aggregate. @@ -76,18 +77,20 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildsig scope := tally.NewTestScope("test", nil) m := buildsignalMocks{ - reqStore: storagemock.NewMockRequestStore(ctrl), - buildStore: storagemock.NewMockBuildStore(ctrl), - queueStore: storagemock.NewMockQueueStore(ctrl), - store: storagemock.NewMockStorage(ctrl), - materializer: requestlogmock.NewMockMaterializer(ctrl), - runnerFactory: buildrunnermock.NewMockFactory(ctrl), - runner: buildrunnermock.NewMockBuildRunner(ctrl), - publisher: mqmock.NewMockPublisher(ctrl), - metricsScope: scope, + reqStore: storagemock.NewMockRequestStore(ctrl), + requestLogStore: storagemock.NewMockRequestLogStore(ctrl), + buildStore: storagemock.NewMockBuildStore(ctrl), + queueStore: storagemock.NewMockQueueStore(ctrl), + store: storagemock.NewMockStorage(ctrl), + materializer: requestlogmock.NewMockMaterializer(ctrl), + runnerFactory: buildrunnermock.NewMockFactory(ctrl), + runner: buildrunnermock.NewMockBuildRunner(ctrl), + publisher: mqmock.NewMockPublisher(ctrl), + metricsScope: scope, } m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() + m.store.EXPECT().GetRequestLogStore().Return(m.requestLogStore).AnyTimes() m.store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes() m.store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() @@ -147,12 +150,13 @@ func buildSignalPayload(t *testing.T, id string) []byte { // requestWithState returns a Request past process's admit, in the given state. func requestWithState(state entity.RequestState) entity.Request { - return entity.Request{ + request := entity.Request{ ID: testID, Queue: testQueue, State: state, Version: 1, } + return request } // build returns a Build with the given status/version, tied to testID. @@ -174,11 +178,19 @@ func queueRow(inFlight, version int32) entity.Queue { } } +func newTerminalStateLog(request entity.Request, buildID string) entity.RequestLog { + log := requestlog.NewRequestStateLog(request, outcomeReason(request.State)) + log.Metadata[requestlog.MetadataKeyBuildID] = buildID + return log +} + func expectFinishWrites(m buildsignalMocks, state entity.RequestState) *gomock.Call { eventCall := expectBuildFinished(m) m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(queueRow(1, 4), nil).After(eventCall) m.queueStore.EXPECT().Update(gomock.Any(), queueRow(0, 4), int32(4), int32(5)).Return(nil) - return m.reqStore.EXPECT().Update(gomock.Any(), requestWithState(state), int32(1), int32(2)).Return(nil) + request := requestWithState(state) + request.Version = 2 + return m.reqStore.EXPECT().FinalizeOutcome(gomock.Any(), request, int32(1), int32(2), newTerminalStateLog(request, testBuildID)).Return(nil) } func expectBuildFinished(m buildsignalMocks) *gomock.Call { @@ -196,7 +208,9 @@ func expectBuildFinished(m buildsignalMocks) *gomock.Call { } func expectFinish(m buildsignalMocks, state entity.RequestState) *gomock.Call { - return expectOutcomeLog(m, state, 2).After(expectFinishWrites(m, state)) + request := requestWithState(state) + request.Version = 2 + return m.materializer.EXPECT().PersistLog(gomock.Any(), m.store, newTerminalStateLog(request, testBuildID)).Return(nil).After(expectFinishWrites(m, state)) } func expectOutcomeLog(m buildsignalMocks, state entity.RequestState, version int32) *gomock.Call { @@ -207,11 +221,32 @@ func expectOutcomeLog(m buildsignalMocks, state entity.RequestState, version int }[state] request := requestWithState(state) request.Version = version + log := requestlog.NewRequestStateLog(request, reason) + log.Metadata[requestlog.MetadataKeyBuildID] = testBuildID + getCall := m.requestLogStore.EXPECT().Get(gomock.Any(), request.ID, log.ID).Return(log, nil) return m.materializer.EXPECT().PersistLog( gomock.Any(), m.store, - requestlog.NewRequestStateLog(request, reason), - ).Return(nil) + log, + ).Return(nil).After(getCall) +} + +func TestMarkOutcomePreservesFirstBuildIdentity(t *testing.T) { + ctrl := gomock.NewController(t) + c, m := newController(t, ctrl) + request := requestWithState(entity.RequestStateProcessing) + winner := requestWithState(entity.RequestStateSucceeded) + winner.Version = 2 + winnerLog := newTerminalStateLog(winner, "winning-build") + + m.reqStore.EXPECT().FinalizeOutcome(gomock.Any(), gomock.Any(), int32(1), int32(2), gomock.Any()).Return(storage.ErrVersionMismatch) + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(winner, nil) + m.requestLogStore.EXPECT().Get(gomock.Any(), testID, winnerLog.ID).Return(winnerLog, nil) + + gotLog, err := c.markOutcome(context.Background(), m.store, &request, entity.RequestStateFailed, "losing-build") + require.NoError(t, err) + assert.Equal(t, winner, request) + assert.Equal(t, winnerLog, gotLog) } func TestProcess(t *testing.T) { @@ -403,11 +438,13 @@ func TestProcess(t *testing.T) { 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) eventCall := expectBuildFinished(m) + terminalStateLog := newTerminalStateLog(request, testBuildID) + m.requestLogStore.EXPECT().Get(gomock.Any(), testID, terminalStateLog.ID).Return(terminalStateLog, nil).After(eventCall) m.materializer.EXPECT().PersistLog( gomock.Any(), m.store, - requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonBuildSucceeded), - ).Return(errors.New("db down")).After(eventCall) + newTerminalStateLog(request, testBuildID), + ).Return(errors.New("db down")) }, }, { @@ -459,7 +496,7 @@ func TestProcess(t *testing.T) { eventCall := expectBuildFinished(m) m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(queueRow(1, 4), nil).After(eventCall) m.queueStore.EXPECT().Update(gomock.Any(), queueRow(0, 4), int32(4), int32(5)).Return(nil) - m.reqStore.EXPECT().Update(gomock.Any(), gomock.Any(), int32(1), int32(2)).Return(errors.New("db down")) + m.reqStore.EXPECT().FinalizeOutcome(gomock.Any(), gomock.Any(), int32(1), int32(2), gomock.Any()).Return(errors.New("db down")) }, }, { @@ -492,7 +529,7 @@ func TestProcess(t *testing.T) { m.materializer.EXPECT().PersistLog( gomock.Any(), m.store, - requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonBuildSucceeded), + newTerminalStateLog(request, testBuildID), ).Return(errors.New("db down")).After(updateCall) }, }, diff --git a/stovepipe/extension/storage/mock/request_store_mock.go b/stovepipe/extension/storage/mock/request_store_mock.go index 089ae8a11..b7f3f4ac0 100644 --- a/stovepipe/extension/storage/mock/request_store_mock.go +++ b/stovepipe/extension/storage/mock/request_store_mock.go @@ -55,6 +55,20 @@ func (mr *MockRequestStoreMockRecorder) Create(ctx, request any) *gomock.Call { return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Create", reflect.TypeOf((*MockRequestStore)(nil).Create), ctx, request) } +// FinalizeOutcome mocks base method. +func (m *MockRequestStore) FinalizeOutcome(ctx context.Context, request entity.Request, oldVersion, newVersion int32, log entity.RequestLog) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "FinalizeOutcome", ctx, request, oldVersion, newVersion, log) + ret0, _ := ret[0].(error) + return ret0 +} + +// FinalizeOutcome indicates an expected call of FinalizeOutcome. +func (mr *MockRequestStoreMockRecorder) FinalizeOutcome(ctx, request, oldVersion, newVersion, log any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "FinalizeOutcome", reflect.TypeOf((*MockRequestStore)(nil).FinalizeOutcome), ctx, request, oldVersion, newVersion, log) +} + // Get mocks base method. func (m *MockRequestStore) Get(ctx context.Context, id string) (entity.Request, error) { m.ctrl.T.Helper() diff --git a/stovepipe/extension/storage/mysql/request_store.go b/stovepipe/extension/storage/mysql/request_store.go index 0c2492b13..e2c49bbc2 100644 --- a/stovepipe/extension/storage/mysql/request_store.go +++ b/stovepipe/extension/storage/mysql/request_store.go @@ -17,8 +17,10 @@ package mysql import ( "context" "database/sql" + "encoding/json" "errors" "fmt" + "time" "github.com/go-sql-driver/mysql" "github.com/uber-go/tally" @@ -155,6 +157,68 @@ func (r *requestStore) Update(ctx context.Context, request entity.Request, oldVe return nil } +// FinalizeOutcome atomically advances a Request to its terminal state and retains +// the matching state entry. Keeping the winning build id in that immutable entry +// means a crash cannot leave a terminal request without its selected build. +func (r *requestStore) FinalizeOutcome(ctx context.Context, request entity.Request, oldVersion, newVersion int32, log entity.RequestLog) (retErr error) { + op := metrics.Begin(r.scope, "finalize_outcome", metrics.StorageLatencyBuckets) + defer func() { op.Complete(retErr) }() + + if request.Queue != r.queue || log.Queue != r.queue || log.RequestID != request.ID { + return fmt.Errorf("request outcome queue or request binding does not match store") + } + if log.TimestampMs == 0 { + log.TimestampMs = time.Now().UnixMilli() + } + if err := log.Validate(); err != nil { + return fmt.Errorf("invalid terminal request log: %w", err) + } + + tx, err := r.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("failed to begin request outcome transaction: %w", err) + } + defer func() { + if retErr != nil { + _ = tx.Rollback() + } + }() + + result, err := tx.ExecContext(ctx, + `UPDATE request SET uri = ?, state = ?, build_strategy = ?, base_uri = ?, version = ? + WHERE queue = ? AND id = ? AND version = ?`, + request.URI, request.State, request.BuildStrategy, request.BaseURI, newVersion, + request.Queue, request.ID, oldVersion, + ) + if err != nil { + return fmt.Errorf("failed to finalize request outcome: %w", err) + } + rows, err := result.RowsAffected() + if err != nil { + return fmt.Errorf("failed to inspect finalized request outcome: %w", err) + } + if rows != 1 { + return fmt.Errorf("version mismatch for request outcome: id=%q expected_version=%d: %w", request.ID, oldVersion, storage.ErrVersionMismatch) + } + + metadata, err := json.Marshal(log.Metadata) + if err != nil { + return fmt.Errorf("failed to marshal terminal request log metadata: %w", err) + } + _, err = tx.ExecContext(ctx, `INSERT INTO request_log ( + queue, request_id, log_id, timestamp_ms, state, event, request_version, outcome_reason, metadata + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`, + log.Queue, log.RequestID, log.ID, log.TimestampMs, log.State, log.Event, log.RequestVersion, log.OutcomeReason, metadata, + ) + if err != nil { + return fmt.Errorf("failed to retain terminal request log: %w", err) + } + if err := tx.Commit(); err != nil { + return fmt.Errorf("failed to commit request outcome: %w", err) + } + return nil +} + // isDuplicateEntry reports whether err is a MySQL duplicate-key (1062) error. func isDuplicateEntry(err error) bool { var mysqlErr *mysql.MySQLError diff --git a/stovepipe/extension/storage/mysql/request_store_test.go b/stovepipe/extension/storage/mysql/request_store_test.go index 0fc657cfa..ba2dc31ad 100644 --- a/stovepipe/extension/storage/mysql/request_store_test.go +++ b/stovepipe/extension/storage/mysql/request_store_test.go @@ -260,6 +260,58 @@ func TestRequestStore_Update(t *testing.T) { } } +func TestRequestStore_FinalizeOutcome(t *testing.T) { + request := entity.Request{ + ID: "request/monorepo/main/1", + Queue: "monorepo/main", + URI: "git://remote/monorepo/main/deadbeef", + State: entity.RequestStateFailed, + Version: 2, + } + log := entity.RequestLog{ + ID: "state/2", + Queue: request.Queue, + RequestID: request.ID, + TimestampMs: 1, + State: request.State, + RequestVersion: request.Version, + OutcomeReason: entity.RequestOutcomeReasonBuildFailed, + Metadata: map[string]string{"build_id": "bk-1"}, + } + + t.Run("commits request and terminal log together", func(t *testing.T) { + db, mock, store := setupRequestStoreTest(t) + defer db.Close() + + mock.ExpectBegin() + mock.ExpectExec("UPDATE request"). + WithArgs(request.URI, request.State, request.BuildStrategy, request.BaseURI, int32(2), request.Queue, request.ID, int32(1)). + WillReturnResult(sqlmock.NewResult(0, 1)) + mock.ExpectExec("INSERT INTO request_log"). + WithArgs(log.Queue, log.RequestID, log.ID, log.TimestampMs, log.State, log.Event, log.RequestVersion, log.OutcomeReason, sqlmock.AnyArg()). + WillReturnResult(sqlmock.NewResult(0, 1)) + mock.ExpectCommit() + + require.NoError(t, store.FinalizeOutcome(context.Background(), request, 1, 2, log)) + require.NoError(t, mock.ExpectationsWereMet()) + }) + + t.Run("rolls back when another build won", func(t *testing.T) { + db, mock, store := setupRequestStoreTest(t) + defer db.Close() + + mock.ExpectBegin() + mock.ExpectExec("UPDATE request"). + WithArgs(request.URI, request.State, request.BuildStrategy, request.BaseURI, int32(2), request.Queue, request.ID, int32(1)). + WillReturnResult(sqlmock.NewResult(0, 0)) + mock.ExpectRollback() + + err := store.FinalizeOutcome(context.Background(), request, 1, 2, log) + require.ErrorIs(t, err, storage.ErrVersionMismatch) + require.NoError(t, mock.ExpectationsWereMet()) + }) +} + func TestIsDuplicateEntry(t *testing.T) { tests := []struct { name string diff --git a/stovepipe/extension/storage/request_store.go b/stovepipe/extension/storage/request_store.go index 4001bf0fa..01f873fb6 100644 --- a/stovepipe/extension/storage/request_store.go +++ b/stovepipe/extension/storage/request_store.go @@ -41,4 +41,10 @@ type RequestStore interface { // and only assigns request.Version = newVersion after this call succeeds. The store performs // a pure conditional write and does not read request.Version. Update(ctx context.Context, request entity.Request, oldVersion, newVersion int32) error + + // FinalizeOutcome atomically writes a terminal request state and its matching + // request-history state entry. The log records the winning build id, so a + // later reader can identify the build that established the terminal outcome. + // Returns ErrVersionMismatch if the request was changed first. + FinalizeOutcome(ctx context.Context, request entity.Request, oldVersion, newVersion int32, log entity.RequestLog) error }