From 52e64ceb3d76a9006925110e06941b0453190bfc Mon Sep 17 00:00:00 2001 From: "chenghan.ying" Date: Fri, 17 Jul 2026 18:38:18 +0000 Subject: [PATCH 1/7] feat(platform): add shared buildkite extension client --- .../buildrunner/buildkite/BUILD.bazel | 19 ++ .../extension/buildrunner/buildkite/README.md | 27 ++ .../extension/buildrunner/buildkite/client.go | 237 ++++++++++++++++++ .../buildrunner/buildkite/client_test.go | 218 ++++++++++++++++ 4 files changed, 501 insertions(+) create mode 100644 platform/extension/buildrunner/buildkite/BUILD.bazel create mode 100644 platform/extension/buildrunner/buildkite/README.md create mode 100644 platform/extension/buildrunner/buildkite/client.go create mode 100644 platform/extension/buildrunner/buildkite/client_test.go diff --git a/platform/extension/buildrunner/buildkite/BUILD.bazel b/platform/extension/buildrunner/buildkite/BUILD.bazel new file mode 100644 index 00000000..8f701d42 --- /dev/null +++ b/platform/extension/buildrunner/buildkite/BUILD.bazel @@ -0,0 +1,19 @@ +load("@rules_go//go:def.bzl", "go_library", "go_test") + +go_library( + name = "go_default_library", + srcs = ["client.go"], + importpath = "github.com/uber/submitqueue/platform/extension/buildrunner/buildkite", + visibility = ["//visibility:public"], +) + +go_test( + name = "go_default_test", + srcs = ["client_test.go"], + embed = [":go_default_library"], + deps = [ + "//platform/http:go_default_library", + "@com_github_stretchr_testify//assert:go_default_library", + "@com_github_stretchr_testify//require:go_default_library", + ], +) diff --git a/platform/extension/buildrunner/buildkite/README.md b/platform/extension/buildrunner/buildkite/README.md new file mode 100644 index 00000000..d946cdbc --- /dev/null +++ b/platform/extension/buildrunner/buildkite/README.md @@ -0,0 +1,27 @@ +# Buildkite client + +Shared HTTP client and Buildkite-specific facts for every domain's Buildkite-backed `BuildRunner`. There is no `BuildRunner` interface here by design — each domain (`submitqueue`, `stovepipe`, ...) defines its own `BuildRunner` and its own `BuildStatus`, and adapts this package's `State` to it. See [`doc/rfc/stovepipe/steps/build.md`](../../../../doc/rfc/stovepipe/steps/build.md#alternatives-considered-for-sharing-the-contract) for why the contract stays per-domain while the backend is shared. + +## What lives here + +- `Client`: `CreateBuild` / `GetBuild` / `CancelBuild` against the Buildkite REST API, plus `EncodeBuildNumber` / `ParseBuildNumber` for the build-number-as-id convention. Construct one with `NewClient(httpClient)`, where `httpClient` already has the pipeline's base URL (via `platform/http.BaseURLTransport`) and auth configured. +- `State` / `ParseState`: Buildkite's own build-state vocabulary (`creating`, `scheduled`, `running`, `blocked`, `passed`, `failed`, `canceling`, `canceled`, `skipped`, `not_run`), collapsed into the five states every domain's `BuildStatus` already distinguishes. Each domain still does its own trivial `State` → its own `BuildStatus` switch — this package does not know either domain's entity types. +- `DecodeMetadataEnv`: recovers a JSON-encoded `map[string]string` from a build's echoed env vars, given the env key the caller used at trigger time. + +## Who consumes it + +- `submitqueue/extension/buildrunner/buildkite` — batch-identity `BuildRunner`, resolves changes via `changeset.Resolver` before triggering. +- `stovepipe/extension/buildrunner/buildkite` — URI-identity `BuildRunner`, triggers directly from `headURI`/`baseURI`. + +Both wrap a `*Client` built at the wiring layer; this package never constructs one from raw config (no credentials, no pipeline slug) itself. + +## How the same `Client` stays safe to share across two different checkout strategies + +SubmitQueue's build materializes state that doesn't exist yet — the runner resolves a batch DAG into composite base/head commits by applying patches. Stovepipe's build checks out a commit that already exists on trunk and, for an incremental build, diffs it against a baseline. These are different problems at the CI-pipeline level (see [build.md's "Why separate contracts"](../../../../doc/rfc/stovepipe/steps/build.md#why-separate-contracts)), yet both go through the same `Client.CreateBuild` call — the `Client` never inspects `CreateBuildRequest.Env` or picks a strategy, so there is nothing here that needs to "know" which pattern applies. + +The split happens entirely outside this package, at two layers below it: + +1. **Each domain's own adapter shapes a different `Env` payload.** SubmitQueue's runner sends `SQ_BASE_URIS`/`SQ_HEAD_URIS` as JSON-encoded arrays of resolved change URIs; stovepipe's sends `STOVEPIPE_HEAD_URI`/`STOVEPIPE_BASE_URI` as single plain strings. Both just call `Client.CreateBuild` with their own `Env` map — this package treats it as an opaque `map[string]string`. +2. **Each domain's `Client` targets a different Buildkite pipeline**, bound once at wiring time via the `httpClient`'s `BaseURLTransport.BaseURL` (e.g. `.../pipelines/submitqueue-go-code` vs `.../pipelines/stovepipe-go-code`). Each pipeline has its own Buildkite pipeline script (owned by Buildkite config, outside this repo) written against its domain's env-var contract — one applies patches into composite commits, the other checks out `STOVEPIPE_HEAD_URI` and diffs against `STOVEPIPE_BASE_URI`. + +So the checkout strategy is a static, wiring-time binding — one `Client` instance ↔ one pipeline ↔ one pipeline script ↔ one env-var contract ↔ one domain's adapter — never a runtime decision made anywhere in Go code. diff --git a/platform/extension/buildrunner/buildkite/client.go b/platform/extension/buildrunner/buildkite/client.go new file mode 100644 index 00000000..723d8398 --- /dev/null +++ b/platform/extension/buildrunner/buildkite/client.go @@ -0,0 +1,237 @@ +// Copyright (c) 2025 Uber Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +// Package buildkite provides the HTTP client and Buildkite-specific facts +// (build state vocabulary, the env-var metadata round-trip convention, and +// build number id encoding) shared by every domain's Buildkite-backed +// BuildRunner. It intentionally holds no BuildRunner interface or domain +// entity types — each domain (submitqueue, stovepipe, ...) defines its own +// BuildRunner and its own BuildStatus, and adapts this package's State to it. +// See doc/rfc/stovepipe/steps/build.md's "Alternatives considered for +// sharing the contract" for the rationale. +package buildkite + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "strconv" +) + +// Client is a thin wrapper around the Buildkite REST endpoints a BuildRunner +// needs: create, get, and cancel a build. +type Client struct { + httpClient *http.Client +} + +// NewClient wraps a pre-configured *http.Client as a Buildkite Client. The +// caller is responsible for the base URL (via platform/http.BaseURLTransport) +// and auth (via an Authorization-header transport). +func NewClient(httpClient *http.Client) *Client { + return &Client{httpClient: httpClient} +} + +// CreateBuildRequest is the payload for POST /builds. +type CreateBuildRequest struct { + Message string `json:"message"` + Env map[string]string `json:"env"` +} + +// BuildResponse is the subset of fields callers need from a Buildkite build +// object. +type BuildResponse struct { + Number int `json:"number"` + State string `json:"state"` + WebURL string `json:"web_url"` + Env map[string]string `json:"env"` +} + +// CreateBuild creates a new Buildkite build. +func (c *Client) CreateBuild(ctx context.Context, req CreateBuildRequest) (BuildResponse, error) { + body, err := json.Marshal(req) + if err != nil { + return BuildResponse{}, fmt.Errorf("marshal request: %w", err) + } + var build BuildResponse + if err := c.do(ctx, http.MethodPost, "/builds", body, &build); err != nil { + return BuildResponse{}, err + } + return build, nil +} + +// GetBuild fetches a build by its Buildkite build number. +func (c *Client) GetBuild(ctx context.Context, number int) (BuildResponse, error) { + u := fmt.Sprintf("/builds/%d", number) + var build BuildResponse + if err := c.do(ctx, http.MethodGet, u, nil, &build); err != nil { + return BuildResponse{}, err + } + return build, nil +} + +// CancelBuild requests cancellation. Returns nil when the build is already +// terminal (HTTP 422) — the Buildkite API uses that status to indicate a +// non-cancellable build, which the BuildRunner contract treats as a no-op. +func (c *Client) CancelBuild(ctx context.Context, number int) error { + u := fmt.Sprintf("/builds/%d/cancel", number) + req, err := http.NewRequestWithContext(ctx, http.MethodPut, u, nil) + if err != nil { + return fmt.Errorf("create request: %w", err) + } + c.setHeaders(req) + + resp, err := c.httpClient.Do(req) + if err != nil { + return fmt.Errorf("send request: %w", err) + } + defer resp.Body.Close() + _, _ = io.Copy(io.Discard, resp.Body) + + switch resp.StatusCode { + case http.StatusOK: + return nil + case http.StatusUnprocessableEntity: + // Already terminal — no-op per BuildRunner.Cancel contract. + return nil + default: + return fmt.Errorf("unexpected status %d from cancel", resp.StatusCode) + } +} + +// do sends an HTTP request with the standard Buildkite headers and, on a 2xx +// response, decodes the body into out (when non-nil). A 404 is reported as a +// "build not found" error per the BuildRunner contract. +func (c *Client) do(ctx context.Context, method, rawURL string, body []byte, out any) error { + var bodyReader io.Reader + if body != nil { + bodyReader = bytes.NewReader(body) + } + + req, err := http.NewRequestWithContext(ctx, method, rawURL, bodyReader) + if err != nil { + return fmt.Errorf("create request: %w", err) + } + c.setHeaders(req) + + resp, err := c.httpClient.Do(req) + if err != nil { + return fmt.Errorf("send request: %w", err) + } + defer resp.Body.Close() + + respBody, err := io.ReadAll(resp.Body) + if err != nil { + return fmt.Errorf("read response: %w", err) + } + + if resp.StatusCode == http.StatusNotFound { + return fmt.Errorf("build not found") + } + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return fmt.Errorf("API returned status %d: %s", resp.StatusCode, respBody) + } + + if out != nil { + if err := json.Unmarshal(respBody, out); err != nil { + return fmt.Errorf("unmarshal response: %w", err) + } + } + return nil +} + +func (c *Client) setHeaders(req *http.Request) { + req.Header.Set("Content-Type", "application/json") +} + +// EncodeBuildNumber encodes a Buildkite build number as an opaque build id +// string. +func EncodeBuildNumber(number int) string { + return strconv.Itoa(number) +} + +// ParseBuildNumber is the inverse of EncodeBuildNumber. +func ParseBuildNumber(id string) (int, error) { + n, err := strconv.Atoi(id) + if err != nil { + return 0, fmt.Errorf("invalid build ID %q", id) + } + return n, nil +} + +// State is Buildkite's own build-state vocabulary, interpreted from the raw +// state string every domain's Buildkite adapter receives on a build object. +// Buildkite's raw states are: creating, scheduled, running, blocked, passed, +// failed, canceling, canceled, skipped, not_run. +type State string + +const ( + // StateUnknown is returned for a raw state this package does not + // recognize. Not terminal — callers should keep polling rather than + // treat it as a final outcome. + StateUnknown State = "" + // StateAccepted means the build has been accepted for execution but has + // not started running yet (Buildkite: creating, scheduled). + StateAccepted State = "accepted" + // StateRunning means the build is currently executing, including while + // blocked on a manual block step (Buildkite: running, blocked). + StateRunning State = "running" + // StateSucceeded means the build completed successfully (Buildkite: + // passed). + StateSucceeded State = "succeeded" + // StateFailed means the build did not produce a passing result + // (Buildkite: failed, not_run, skipped). + StateFailed State = "failed" + // StateCancelled means the build was cancelled, or is in the process of + // being cancelled (Buildkite: canceling, canceled). + StateCancelled State = "cancelled" +) + +// ParseState maps a raw Buildkite build-state string to a State. An +// unrecognized raw string maps to StateUnknown rather than being assumed +// terminal. +func ParseState(raw string) State { + switch raw { + case "creating", "scheduled": + return StateAccepted + case "running", "blocked": + return StateRunning + case "passed": + return StateSucceeded + case "failed", "not_run", "skipped": + return StateFailed + case "canceling", "canceled": + return StateCancelled + default: + return StateUnknown + } +} + +// DecodeMetadataEnv recovers a JSON-encoded map[string]string from env[key]. +// Buildkite echoes env vars back on the build object, so a caller that seeded +// a JSON-encoded map at Trigger time can recover it here at Status time +// without any local state. Returns an empty non-nil map when the key is +// absent or its value cannot be decoded — a corrupt env var must not fail the +// caller. +func DecodeMetadataEnv(env map[string]string, key string) map[string]string { + meta := make(map[string]string) + raw, ok := env[key] + if !ok || raw == "" { + return meta + } + _ = json.Unmarshal([]byte(raw), &meta) + return meta +} diff --git a/platform/extension/buildrunner/buildkite/client_test.go b/platform/extension/buildrunner/buildkite/client_test.go new file mode 100644 index 00000000..4984c329 --- /dev/null +++ b/platform/extension/buildrunner/buildkite/client_test.go @@ -0,0 +1,218 @@ +// Copyright (c) 2025 Uber Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package buildkite + +import ( + "context" + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + phttp "github.com/uber/submitqueue/platform/http" +) + +// newTestClient creates a Client backed by a test HTTP server. +func newTestClient(t *testing.T, handler http.Handler) *Client { + t.Helper() + srv := httptest.NewServer(handler) + t.Cleanup(srv.Close) + c, err := phttp.NewClient(srv.URL) + require.NoError(t, err) + return NewClient(c) +} + +func buildJSON(number int, state, webURL string) []byte { + return buildJSONWithEnv(number, state, webURL, nil) +} + +func buildJSONWithEnv(number int, state, webURL string, env map[string]string) []byte { + b, _ := json.Marshal(BuildResponse{Number: number, State: state, WebURL: webURL, Env: env}) + return b +} + +// --- CreateBuild --- + +func TestCreateBuild_SubmitsPayloadAndReturnsResponse(t *testing.T) { + var capturedMethod string + var capturedBody []byte + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + capturedMethod = req.Method + capturedBody, _ = io.ReadAll(req.Body) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(buildJSON(42, "scheduled", "https://buildkite.com/test-org/my-pipeline/builds/42")) + })) + + resp, err := c.CreateBuild(context.Background(), CreateBuildRequest{ + Message: "test build", + Env: map[string]string{"KEY": "value"}, + }) + require.NoError(t, err) + assert.Equal(t, http.MethodPost, capturedMethod) + assert.Equal(t, 42, resp.Number) + assert.Equal(t, "scheduled", resp.State) + assert.Equal(t, "https://buildkite.com/test-org/my-pipeline/builds/42", resp.WebURL) + + var req CreateBuildRequest + require.NoError(t, json.Unmarshal(capturedBody, &req)) + assert.Equal(t, "test build", req.Message) + assert.Equal(t, "value", req.Env["KEY"]) +} + +func TestCreateBuild_ErrorStatus_ReturnsError(t *testing.T) { + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + + _, err := c.CreateBuild(context.Background(), CreateBuildRequest{}) + require.Error(t, err) +} + +// --- GetBuild --- + +func TestGetBuild_ReturnsResponse(t *testing.T) { + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + assert.Equal(t, http.MethodGet, req.Method) + assert.Equal(t, "/builds/7", req.URL.Path) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(buildJSON(7, "running", "https://buildkite.com/test-org/my-pipeline/builds/7")) + })) + + resp, err := c.GetBuild(context.Background(), 7) + require.NoError(t, err) + assert.Equal(t, "running", resp.State) +} + +func TestGetBuild_NotFound_ReturnsError(t *testing.T) { + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusNotFound) + })) + + _, err := c.GetBuild(context.Background(), 99) + require.Error(t, err) +} + +func TestGetBuild_EchoesEnv(t *testing.T) { + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(buildJSONWithEnv(7, "passed", "", map[string]string{"FOO": "bar"})) + })) + + resp, err := c.GetBuild(context.Background(), 7) + require.NoError(t, err) + assert.Equal(t, "bar", resp.Env["FOO"]) +} + +// --- CancelBuild --- + +func TestCancelBuild_CallsBuildkite(t *testing.T) { + var cancelled bool + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + cancelled = true + assert.Equal(t, http.MethodPut, req.Method) + w.WriteHeader(http.StatusOK) + })) + + require.NoError(t, c.CancelBuild(context.Background(), 5)) + assert.True(t, cancelled) +} + +func TestCancelBuild_AlreadyTerminal_Noop(t *testing.T) { + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusUnprocessableEntity) + })) + + require.NoError(t, c.CancelBuild(context.Background(), 5)) +} + +func TestCancelBuild_ErrorStatus_ReturnsError(t *testing.T) { + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + + require.Error(t, c.CancelBuild(context.Background(), 5)) +} + +// --- EncodeBuildNumber / ParseBuildNumber --- + +func TestEncodeParseBuildNumber_RoundTrip(t *testing.T) { + for _, n := range []int{1, 9999, 0} { + id := EncodeBuildNumber(n) + got, err := ParseBuildNumber(id) + require.NoError(t, err) + assert.Equal(t, n, got) + } +} + +func TestParseBuildNumber_Invalid(t *testing.T) { + for _, id := range []string{"", "notanumber", "org/pipeline/1"} { + t.Run(id, func(t *testing.T) { + _, err := ParseBuildNumber(id) + require.Error(t, err) + }) + } +} + +// --- ParseState --- + +func TestParseState(t *testing.T) { + tests := []struct { + raw string + want State + }{ + {"creating", StateAccepted}, + {"scheduled", StateAccepted}, + {"running", StateRunning}, + {"blocked", StateRunning}, + {"passed", StateSucceeded}, + {"failed", StateFailed}, + {"not_run", StateFailed}, + {"skipped", StateFailed}, + {"canceling", StateCancelled}, + {"canceled", StateCancelled}, + {"some_future_state", StateUnknown}, + {"", StateUnknown}, + } + for _, tt := range tests { + t.Run(tt.raw, func(t *testing.T) { + assert.Equal(t, tt.want, ParseState(tt.raw)) + }) + } +} + +// --- DecodeMetadataEnv --- + +func TestDecodeMetadataEnv_PresentAndValid(t *testing.T) { + env := map[string]string{"SQ_METADATA": `{"requester":"alice","ticket":"SQ-42"}`} + got := DecodeMetadataEnv(env, "SQ_METADATA") + assert.Equal(t, map[string]string{"requester": "alice", "ticket": "SQ-42"}, got) +} + +func TestDecodeMetadataEnv_Absent_ReturnsEmptyMap(t *testing.T) { + got := DecodeMetadataEnv(map[string]string{}, "SQ_METADATA") + assert.Empty(t, got) + assert.NotNil(t, got) +} + +func TestDecodeMetadataEnv_Malformed_ReturnsEmptyMap(t *testing.T) { + env := map[string]string{"SQ_METADATA": "not json"} + got := DecodeMetadataEnv(env, "SQ_METADATA") + assert.Empty(t, got) + assert.NotNil(t, got) +} From 78267432bddb7c5fd1a41a5e92fd32ee9dbd7710 Mon Sep 17 00:00:00 2001 From: "chenghan.ying" Date: Fri, 17 Jul 2026 18:38:38 +0000 Subject: [PATCH 2/7] refactor(submitqueue): rebase buildkite adapter onto shared platform client --- submitqueue/extension/buildrunner/README.md | 2 +- .../buildrunner/buildkite/BUILD.bazel | 7 +- .../buildrunner/buildkite/buildkite.go | 89 ++++------- .../buildrunner/buildkite/buildkite_test.go | 51 ++----- .../extension/buildrunner/buildkite/client.go | 139 ------------------ 5 files changed, 47 insertions(+), 241 deletions(-) delete mode 100644 submitqueue/extension/buildrunner/buildkite/client.go diff --git a/submitqueue/extension/buildrunner/README.md b/submitqueue/extension/buildrunner/README.md index 29d8f512..7c3d610f 100644 --- a/submitqueue/extension/buildrunner/README.md +++ b/submitqueue/extension/buildrunner/README.md @@ -17,4 +17,4 @@ See [`doc/rfc/submitqueue/build-runner.md`](../../../doc/rfc/submitqueue/build-r - `githubactions`: proof-of-architecture backend that dispatches a GitHub Actions workflow. See [`githubactions/README.md`](githubactions/README.md) for the workflow inputs and example orchestrator environment variables. -- `buildkite`: Buildkite-backed backend. +- `buildkite`: Buildkite-backed backend. Its HTTP client and Buildkite-specific facts (state vocabulary, metadata env-var round-trip) live in [`platform/extension/buildrunner/buildkite`](../../../platform/extension/buildrunner/buildkite/README.md), shared with `stovepipe`'s own Buildkite backend. diff --git a/submitqueue/extension/buildrunner/buildkite/BUILD.bazel b/submitqueue/extension/buildrunner/buildkite/BUILD.bazel index 36bd4de8..d73b45a0 100644 --- a/submitqueue/extension/buildrunner/buildkite/BUILD.bazel +++ b/submitqueue/extension/buildrunner/buildkite/BUILD.bazel @@ -2,14 +2,12 @@ load("@rules_go//go:def.bzl", "go_library", "go_test") go_library( name = "go_default_library", - srcs = [ - "buildkite.go", - "client.go", - ], + srcs = ["buildkite.go"], importpath = "github.com/uber/submitqueue/submitqueue/extension/buildrunner/buildkite", visibility = ["//visibility:public"], deps = [ "//platform/base/change:go_default_library", + "//platform/extension/buildrunner/buildkite:go_default_library", "//submitqueue/core/changeset:go_default_library", "//submitqueue/entity:go_default_library", "//submitqueue/extension/buildrunner:go_default_library", @@ -23,6 +21,7 @@ go_test( embed = [":go_default_library"], deps = [ "//platform/base/change:go_default_library", + "//platform/extension/buildrunner/buildkite:go_default_library", "//platform/http:go_default_library", "//submitqueue/core/changeset:go_default_library", "//submitqueue/core/changeset/fake:go_default_library", diff --git a/submitqueue/extension/buildrunner/buildkite/buildkite.go b/submitqueue/extension/buildrunner/buildkite/buildkite.go index 05898561..f91e2c93 100644 --- a/submitqueue/extension/buildrunner/buildkite/buildkite.go +++ b/submitqueue/extension/buildrunner/buildkite/buildkite.go @@ -33,12 +33,11 @@ import ( "context" "encoding/json" "fmt" - "net/http" - "strconv" "go.uber.org/zap" "github.com/uber/submitqueue/platform/base/change" + platformbuildkite "github.com/uber/submitqueue/platform/extension/buildrunner/buildkite" "github.com/uber/submitqueue/submitqueue/core/changeset" "github.com/uber/submitqueue/submitqueue/entity" "github.com/uber/submitqueue/submitqueue/extension/buildrunner" @@ -69,23 +68,21 @@ const ( // runner implements buildrunner.BuildRunner. type runner struct { cfg buildrunner.Config - client *client + client *platformbuildkite.Client resolver changeset.Resolver logger *zap.SugaredLogger } var _ buildrunner.BuildRunner = (*runner)(nil) -// Params holds the dependencies for a Buildkite BuildRunner. The caller is -// responsible for configuring HTTPClient with the base URL (via -// platform/http.BaseURLTransport) and auth (via an Authorization-header transport). +// Params holds the dependencies for a Buildkite BuildRunner. type Params struct { // Config holds the per-queue identity for this BuildRunner. Config buildrunner.Config - // HTTPClient is a pre-configured HTTP client. The caller is responsible - // for the base URL (via platform/http.BaseURLTransport) and auth (via a - // transport layer). If nil, http.DefaultClient is used. - HTTPClient *http.Client + // Client is a pre-constructed Buildkite client. The wiring layer builds + // it once via platformbuildkite.NewClient, with the pipeline's base URL + // (via platform/http.BaseURLTransport) and auth already configured. + Client *platformbuildkite.Client // Resolver resolves a batch's changes (base and head batches). Resolver changeset.Resolver // Logger is the structured logger. @@ -94,22 +91,18 @@ type Params struct { // NewBuildRunner constructs a Buildkite-backed BuildRunner bound to a single // pipeline. -// -// The HTTPClient must have BaseURLTransport configured to the pipeline's API -// root (e.g. "https://api.buildkite.com/v2/organizations/{org}/pipelines/{slug}"), -// and an auth transport that injects the Authorization header. func NewBuildRunner(params Params) (buildrunner.BuildRunner, error) { - if params.HTTPClient == nil { - return nil, fmt.Errorf("http client is required") + if params.Client == nil { + return nil, fmt.Errorf("buildkite client is required") } if params.Logger == nil { return nil, fmt.Errorf("logger is required") } - return newRunner(params.Config, &client{httpClient: params.HTTPClient}, params.Resolver, params.Logger.Named("buildkite_buildrunner")), nil + return newRunner(params.Config, params.Client, params.Resolver, params.Logger.Named("buildkite_buildrunner")), nil } // newRunner constructs a runner. Used by NewBuildRunner and by tests. -func newRunner(cfg buildrunner.Config, c *client, resolver changeset.Resolver, logger *zap.SugaredLogger) *runner { +func newRunner(cfg buildrunner.Config, c *platformbuildkite.Client, resolver changeset.Resolver, logger *zap.SugaredLogger) *runner { return &runner{ cfg: cfg, client: c, @@ -144,12 +137,12 @@ func (r *runner) Trigger(ctx context.Context, base []entity.Batch, head entity.B env[EnvKeyMetadata] = string(metaJSON) } - req := createBuildRequest{ + req := platformbuildkite.CreateBuildRequest{ Message: "submitqueue speculative build", Env: env, } - resp, err := r.client.createBuild(ctx, req) + resp, err := r.client.CreateBuild(ctx, req) if err != nil { return entity.BuildID{}, fmt.Errorf("buildkite: create build: %w", err) } @@ -157,18 +150,18 @@ func (r *runner) Trigger(ctx context.Context, base []entity.Batch, head entity.B r.logger.Debugw("triggered Buildkite build", "buildkite_number", resp.Number, ) - return entity.BuildID{ID: encodeBuildNumber(resp.Number)}, nil + return entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(resp.Number)}, nil } // Status fetches the current state of the build from Buildkite and returns it // with the build URL and any caller-supplied metadata in BuildMetadata. func (r *runner) Status(ctx context.Context, buildID entity.BuildID) (entity.BuildStatus, entity.BuildMetadata, error) { - number, err := parseBuildNumber(buildID.ID) + number, err := platformbuildkite.ParseBuildNumber(buildID.ID) if err != nil { return entity.BuildStatusUnknown, nil, fmt.Errorf("buildkite: malformed build ID: %w", err) } - resp, err := r.client.getBuild(ctx, number) + resp, err := r.client.GetBuild(ctx, number) if err != nil { return entity.BuildStatusUnknown, nil, fmt.Errorf("buildkite: get build: %w", err) } @@ -181,12 +174,12 @@ func (r *runner) Status(ctx context.Context, buildID entity.BuildID) (entity.Bui // Cancel calls the Buildkite API to cancel the build. A no-op on already-terminal // builds (Buildkite returns 422 for those). func (r *runner) Cancel(ctx context.Context, buildID entity.BuildID) error { - number, err := parseBuildNumber(buildID.ID) + number, err := platformbuildkite.ParseBuildNumber(buildID.ID) if err != nil { return fmt.Errorf("buildkite: malformed build ID: %w", err) } - if err := r.client.cancelBuild(ctx, number); err != nil { + if err := r.client.CancelBuild(ctx, number); err != nil { return fmt.Errorf("buildkite: cancel build: %w", err) } r.logger.Debugw("cancelled Buildkite build", @@ -204,52 +197,24 @@ func flattenURIs(changes []change.Change) []string { return uris } -// encodeBuildNumber encodes a Buildkite build number as the SQ build ID. -func encodeBuildNumber(number int) string { - return strconv.Itoa(number) -} - -// parseBuildNumber is the inverse of encodeBuildNumber. -func parseBuildNumber(id string) (int, error) { - n, err := strconv.Atoi(id) - if err != nil { - return 0, fmt.Errorf("invalid build ID %q", id) - } - return n, nil -} - // decodeMetadata recovers the caller-supplied BuildMetadata from the env vars -// Buildkite echoes back on the build object. Returns an empty non-nil map when -// SQ_METADATA is absent or cannot be decoded — a corrupt env var must not fail -// a Status call. +// Buildkite echoes back on the build object. func decodeMetadata(env map[string]string) entity.BuildMetadata { - meta := make(entity.BuildMetadata) - raw, ok := env[EnvKeyMetadata] - if !ok || raw == "" { - return meta - } - _ = json.Unmarshal([]byte(raw), &meta) - return meta + return entity.BuildMetadata(platformbuildkite.DecodeMetadataEnv(env, EnvKeyMetadata)) } -// mapState maps a Buildkite build state string to a BuildStatus. -// -// Buildkite states: creating, scheduled, running, blocked, passed, failed, -// canceling, canceled, skipped, not_run. +// mapState maps a Buildkite build state to a BuildStatus. func mapState(state string) entity.BuildStatus { - switch state { - case "creating", "scheduled": + switch platformbuildkite.ParseState(state) { + case platformbuildkite.StateAccepted: return entity.BuildStatusAccepted - case "running", "blocked": - // blocked = waiting on a block step; still live, not yet terminal. + case platformbuildkite.StateRunning: return entity.BuildStatusRunning - case "passed": + case platformbuildkite.StateSucceeded: return entity.BuildStatusSucceeded - case "failed", "not_run", "skipped": - // not_run/skipped never produced a passing result; treat them as - // terminal failure so the batch is not merged on a non-success verdict. + case platformbuildkite.StateFailed: return entity.BuildStatusFailed - case "canceling", "canceled": + case platformbuildkite.StateCancelled: return entity.BuildStatusCancelled default: // Unrecognised Buildkite state. Do NOT assume terminal: Unknown is diff --git a/submitqueue/extension/buildrunner/buildkite/buildkite_test.go b/submitqueue/extension/buildrunner/buildkite/buildkite_test.go index 90dcfa86..4decca58 100644 --- a/submitqueue/extension/buildrunner/buildkite/buildkite_test.go +++ b/submitqueue/extension/buildrunner/buildkite/buildkite_test.go @@ -27,6 +27,7 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/uber/submitqueue/platform/base/change" + platformbuildkite "github.com/uber/submitqueue/platform/extension/buildrunner/buildkite" phttp "github.com/uber/submitqueue/platform/http" "github.com/uber/submitqueue/submitqueue/core/changeset" changesetfake "github.com/uber/submitqueue/submitqueue/core/changeset/fake" @@ -49,7 +50,7 @@ func newTestRunner(t *testing.T, handler http.Handler, resolver ...changeset.Res } return newRunner( buildrunner.Config{QueueName: "my-queue"}, - &client{httpClient: c}, + platformbuildkite.NewClient(c), r, zap.NewNop().Sugar(), ) @@ -62,7 +63,7 @@ func buildJSON(number int, state, webURL string) []byte { // buildJSONWithEnv encodes fields into a Buildkite build JSON response including env vars. func buildJSONWithEnv(number int, state, webURL string, env map[string]string) []byte { - b, _ := json.Marshal(buildResponse{Number: number, State: state, WebURL: webURL, Env: env}) + b, _ := json.Marshal(platformbuildkite.BuildResponse{Number: number, State: state, WebURL: webURL, Env: env}) return b } @@ -70,7 +71,7 @@ func buildJSONWithEnv(number int, state, webURL string, env map[string]string) [ func TestNew_ImplementsInterface(t *testing.T) { r, err := NewBuildRunner(Params{Logger: zap.NewNop().Sugar()}) - require.Error(t, err, "http client is required") + require.Error(t, err, "buildkite client is required") var _ buildrunner.BuildRunner = r } @@ -93,11 +94,11 @@ func TestTrigger_SubmitsCorrectPayloadAndReturnsBuildkiteNumber(t *testing.T) { id, err := r.Trigger(context.Background(), []entity.Batch{{ID: "base-batch"}}, entity.Batch{ID: "head-batch"}, nil) require.NoError(t, err) - assert.Equal(t, encodeBuildNumber(42), id.ID) + assert.Equal(t, platformbuildkite.EncodeBuildNumber(42), id.ID) assert.Equal(t, http.MethodPost, capturedMethod) - var req createBuildRequest + var req platformbuildkite.CreateBuildRequest require.NoError(t, json.Unmarshal(capturedBody, &req)) assert.Equal(t, `["github://github.example.com/org/repo/pull/1/aaa111"]`, req.Env[EnvKeyBaseURIs]) assert.Equal(t, `["github://github.example.com/org/repo/pull/2/bbb222"]`, req.Env[EnvKeyHeadURIs]) @@ -116,7 +117,7 @@ func TestTrigger_EmptyBase_ProducesJSONArray(t *testing.T) { _, err := r.Trigger(context.Background(), nil, entity.Batch{ID: "head-batch"}, nil) require.NoError(t, err) - var req createBuildRequest + var req platformbuildkite.CreateBuildRequest require.NoError(t, json.Unmarshal(capturedBody, &req)) // nil base must produce [] in JSON, not null. assert.Equal(t, "[]", req.Env[EnvKeyBaseURIs]) @@ -137,7 +138,7 @@ func TestTrigger_MultipleChangesFlattened(t *testing.T) { _, err := r.Trigger(context.Background(), nil, entity.Batch{ID: "head-batch"}, nil) require.NoError(t, err) - var req createBuildRequest + var req platformbuildkite.CreateBuildRequest require.NoError(t, json.Unmarshal(capturedBody, &req)) assert.Equal(t, `["github://github.example.com/org/repo/pull/1/aaa","github://github.example.com/org/repo/pull/2/bbb","github://github.example.com/org/repo/pull/3/ccc"]`, @@ -167,7 +168,7 @@ func TestTrigger_WithMetadata_SetsEnvVar(t *testing.T) { _, err := r.Trigger(context.Background(), nil, entity.Batch{ID: "head-batch"}, metadata) require.NoError(t, err) - var req createBuildRequest + var req platformbuildkite.CreateBuildRequest require.NoError(t, json.Unmarshal(capturedBody, &req)) require.Contains(t, req.Env, EnvKeyMetadata) @@ -188,7 +189,7 @@ func TestTrigger_NilMetadata_NoMetadataEnvVar(t *testing.T) { _, err := r.Trigger(context.Background(), nil, entity.Batch{ID: "head-batch"}, nil) require.NoError(t, err) - var req createBuildRequest + var req platformbuildkite.CreateBuildRequest require.NoError(t, json.Unmarshal(capturedBody, &req)) assert.NotContains(t, req.Env, EnvKeyMetadata) } @@ -228,7 +229,7 @@ func TestStatus_ReturnsLiveBuildkiteState(t *testing.T) { _, _ = w.Write(buildJSON(7, "running", "https://buildkite.com/test-org/my-pipeline/builds/7")) })) - status, meta, err := r.Status(context.Background(), entity.BuildID{ID: encodeBuildNumber(7)}) + status, meta, err := r.Status(context.Background(), entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(7)}) require.NoError(t, err) assert.Equal(t, entity.BuildStatusRunning, status) assert.Equal(t, "https://buildkite.com/test-org/my-pipeline/builds/7", meta["url"]) @@ -239,7 +240,7 @@ func TestStatus_BuildkiteNotFound(t *testing.T) { w.WriteHeader(http.StatusNotFound) })) - _, _, err := r.Status(context.Background(), entity.BuildID{ID: encodeBuildNumber(99)}) + _, _, err := r.Status(context.Background(), entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(99)}) require.Error(t, err) } @@ -254,7 +255,7 @@ func TestStatus_EchosCallerMetadata(t *testing.T) { )) })) - _, meta, err := r.Status(context.Background(), entity.BuildID{ID: encodeBuildNumber(7)}) + _, meta, err := r.Status(context.Background(), entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(7)}) require.NoError(t, err) assert.Equal(t, "alice", meta["requester"]) assert.Equal(t, "SQ-42", meta["ticket"]) @@ -267,7 +268,7 @@ func TestStatus_NoMetadata_ReturnsOnlyURL(t *testing.T) { _, _ = w.Write(buildJSON(8, "running", "https://buildkite.com/test-org/my-pipeline/builds/8")) })) - _, meta, err := r.Status(context.Background(), entity.BuildID{ID: encodeBuildNumber(8)}) + _, meta, err := r.Status(context.Background(), entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(8)}) require.NoError(t, err) assert.Equal(t, entity.BuildMetadata{"url": "https://buildkite.com/test-org/my-pipeline/builds/8"}, meta) } @@ -282,7 +283,7 @@ func TestCancel_CallsBuildkite(t *testing.T) { _, _ = w.Write(buildJSON(5, "canceled", "")) })) - require.NoError(t, r.Cancel(context.Background(), entity.BuildID{ID: encodeBuildNumber(5)})) + require.NoError(t, r.Cancel(context.Background(), entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(5)})) assert.True(t, cancelCalled) } @@ -292,25 +293,5 @@ func TestCancel_AlreadyTerminal_Noop(t *testing.T) { w.WriteHeader(http.StatusUnprocessableEntity) })) - require.NoError(t, r.Cancel(context.Background(), entity.BuildID{ID: encodeBuildNumber(5)})) -} - -// --- Internal helpers --- - -func TestEncodeParseBuildNumber_RoundTrip(t *testing.T) { - for _, n := range []int{1, 9999, 0} { - id := encodeBuildNumber(n) - got, err := parseBuildNumber(id) - require.NoError(t, err) - assert.Equal(t, n, got) - } -} - -func TestParseBuildNumber_Invalid(t *testing.T) { - for _, id := range []string{"", "notanumber", "org/pipeline/1"} { - t.Run(id, func(t *testing.T) { - _, err := parseBuildNumber(id) - require.Error(t, err) - }) - } + require.NoError(t, r.Cancel(context.Background(), entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(5)})) } diff --git a/submitqueue/extension/buildrunner/buildkite/client.go b/submitqueue/extension/buildrunner/buildkite/client.go deleted file mode 100644 index dcf09d99..00000000 --- a/submitqueue/extension/buildrunner/buildkite/client.go +++ /dev/null @@ -1,139 +0,0 @@ -// Copyright (c) 2025 Uber Technologies, Inc. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -package buildkite - -import ( - "bytes" - "context" - "encoding/json" - "fmt" - "io" - "net/http" -) - -// client is a thin wrapper around the Buildkite REST endpoints that BuildRunner -// needs: create, get, and cancel a build. -type client struct { - httpClient *http.Client -} - -type createBuildRequest struct { - Message string `json:"message"` - Env map[string]string `json:"env"` -} - -// buildResponse is the subset of fields the runner needs from a Buildkite -// build object. -type buildResponse struct { - Number int `json:"number"` - State string `json:"state"` - WebURL string `json:"web_url"` - Env map[string]string `json:"env"` -} - -func (c *client) createBuild(ctx context.Context, req createBuildRequest) (buildResponse, error) { - body, err := json.Marshal(req) - if err != nil { - return buildResponse{}, fmt.Errorf("marshal request: %w", err) - } - var build buildResponse - if err := c.do(ctx, http.MethodPost, "/builds", body, &build); err != nil { - return buildResponse{}, err - } - return build, nil -} - -func (c *client) getBuild(ctx context.Context, number int) (buildResponse, error) { - u := fmt.Sprintf("/builds/%d", number) - var build buildResponse - if err := c.do(ctx, http.MethodGet, u, nil, &build); err != nil { - return buildResponse{}, err - } - return build, nil -} - -// cancelBuild requests cancellation. Returns nil when the build is already -// terminal (HTTP 422) — the Buildkite API uses that status to indicate a -// non-cancellable build, which the BuildRunner contract treats as a no-op. -func (c *client) cancelBuild(ctx context.Context, number int) error { - u := fmt.Sprintf("/builds/%d/cancel", number) - req, err := http.NewRequestWithContext(ctx, http.MethodPut, u, nil) - if err != nil { - return fmt.Errorf("create request: %w", err) - } - c.setHeaders(req) - - resp, err := c.httpClient.Do(req) - if err != nil { - return fmt.Errorf("send request: %w", err) - } - defer resp.Body.Close() - _, _ = io.Copy(io.Discard, resp.Body) - - switch resp.StatusCode { - case http.StatusOK: - return nil - case http.StatusUnprocessableEntity: - // Already terminal — no-op per BuildRunner.Cancel contract. - return nil - default: - return fmt.Errorf("unexpected status %d from cancel", resp.StatusCode) - } -} - -// do sends an HTTP request with the standard Buildkite headers and, on a 2xx -// response, decodes the body into out (when non-nil). A 404 is reported as a -// "build not found" error per the BuildRunner contract. -func (c *client) do(ctx context.Context, method, rawURL string, body []byte, out any) error { - var bodyReader io.Reader - if body != nil { - bodyReader = bytes.NewReader(body) - } - - req, err := http.NewRequestWithContext(ctx, method, rawURL, bodyReader) - if err != nil { - return fmt.Errorf("create request: %w", err) - } - c.setHeaders(req) - - resp, err := c.httpClient.Do(req) - if err != nil { - return fmt.Errorf("send request: %w", err) - } - defer resp.Body.Close() - - respBody, err := io.ReadAll(resp.Body) - if err != nil { - return fmt.Errorf("read response: %w", err) - } - - if resp.StatusCode == http.StatusNotFound { - return fmt.Errorf("build not found") - } - if resp.StatusCode < 200 || resp.StatusCode >= 300 { - return fmt.Errorf("API returned status %d: %s", resp.StatusCode, respBody) - } - - if out != nil { - if err := json.Unmarshal(respBody, out); err != nil { - return fmt.Errorf("unmarshal response: %w", err) - } - } - return nil -} - -func (c *client) setHeaders(req *http.Request) { - req.Header.Set("Content-Type", "application/json") -} From 4f58d41711a421d550156c4b95ca153475c1729a Mon Sep 17 00:00:00 2001 From: "chenghan.ying" Date: Fri, 17 Jul 2026 18:39:00 +0000 Subject: [PATCH 3/7] feat(stovepipe): add buildkite-backed BuildRunner adapter --- stovepipe/extension/buildrunner/README.md | 2 + .../buildrunner/buildkite/BUILD.bazel | 29 ++ .../buildrunner/buildkite/buildkite.go | 196 +++++++++++++ .../buildrunner/buildkite/buildkite_test.go | 274 ++++++++++++++++++ 4 files changed, 501 insertions(+) create mode 100644 stovepipe/extension/buildrunner/buildkite/BUILD.bazel create mode 100644 stovepipe/extension/buildrunner/buildkite/buildkite.go create mode 100644 stovepipe/extension/buildrunner/buildkite/buildkite_test.go diff --git a/stovepipe/extension/buildrunner/README.md b/stovepipe/extension/buildrunner/README.md index 3abb05c0..b25d64fa 100644 --- a/stovepipe/extension/buildrunner/README.md +++ b/stovepipe/extension/buildrunner/README.md @@ -8,4 +8,6 @@ Vendor-agnostic interface through which Stovepipe triggers and polls builds agai Implementations return plain, unclassified errors — the calling controller decides retryable-vs-not and user-vs-infra, per `platform/errs`. +Real backends (`buildkite`, `githubactions`) are thin adapters over a shared platform client (`platform/extension/buildrunner/{backend}`) — see that package's README for the HTTP client and vendor-specific details. + See [doc/rfc/stovepipe/steps/build.md](../../../doc/rfc/stovepipe/steps/build.md#why-separate-contracts) for why this is a separate contract from SubmitQueue's own `buildrunner` rather than a shared one. To add a backend, create `buildrunner/{backend}/`, implement `BuildRunner`, and return it from a `New(...)` constructor. diff --git a/stovepipe/extension/buildrunner/buildkite/BUILD.bazel b/stovepipe/extension/buildrunner/buildkite/BUILD.bazel new file mode 100644 index 00000000..64c3dda6 --- /dev/null +++ b/stovepipe/extension/buildrunner/buildkite/BUILD.bazel @@ -0,0 +1,29 @@ +load("@rules_go//go:def.bzl", "go_library", "go_test") + +go_library( + name = "go_default_library", + srcs = ["buildkite.go"], + importpath = "github.com/uber/submitqueue/stovepipe/extension/buildrunner/buildkite", + visibility = ["//visibility:public"], + deps = [ + "//platform/extension/buildrunner/buildkite:go_default_library", + "//stovepipe/entity:go_default_library", + "//stovepipe/extension/buildrunner:go_default_library", + "@org_uber_go_zap//:go_default_library", + ], +) + +go_test( + name = "go_default_test", + srcs = ["buildkite_test.go"], + embed = [":go_default_library"], + deps = [ + "//platform/extension/buildrunner/buildkite:go_default_library", + "//platform/http:go_default_library", + "//stovepipe/entity:go_default_library", + "//stovepipe/extension/buildrunner:go_default_library", + "@com_github_stretchr_testify//assert:go_default_library", + "@com_github_stretchr_testify//require:go_default_library", + "@org_uber_go_zap//:go_default_library", + ], +) diff --git a/stovepipe/extension/buildrunner/buildkite/buildkite.go b/stovepipe/extension/buildrunner/buildkite/buildkite.go new file mode 100644 index 00000000..db890e08 --- /dev/null +++ b/stovepipe/extension/buildrunner/buildkite/buildkite.go @@ -0,0 +1,196 @@ +// Copyright (c) 2025 Uber Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +// Package buildkite implements buildrunner.BuildRunner backed by the Buildkite +// CI platform. +// +// Trigger calls the Buildkite API to create the build and returns the Buildkite +// build number as the build ID. Status and Cancel parse the number directly +// from the build ID — no local state is required. +// +// The Buildkite build receives the head and base URIs as environment +// variables (STOVEPIPE_HEAD_URI, STOVEPIPE_BASE_URI, STOVEPIPE_QUEUE). The +// pipeline script checks out headURI and, when baseURI is non-empty, diffs +// against it for an incremental build. +// +// Caller-supplied BuildMetadata is forwarded to the build as +// STOVEPIPE_METADATA (JSON-encoded). Buildkite echoes env vars back on the +// build object, so Status recovers and returns the original metadata without +// any local state. +package buildkite + +import ( + "context" + "encoding/json" + "fmt" + + "go.uber.org/zap" + + platformbuildkite "github.com/uber/submitqueue/platform/extension/buildrunner/buildkite" + "github.com/uber/submitqueue/stovepipe/entity" + "github.com/uber/submitqueue/stovepipe/extension/buildrunner" +) + +// Env var keys set on every triggered Buildkite build. +const ( + // EnvKeyHeadURI carries the URI of the commit under validation. + EnvKeyHeadURI = "STOVEPIPE_HEAD_URI" + + // EnvKeyBaseURI carries the incremental baseline URI, empty for a full + // build. + EnvKeyBaseURI = "STOVEPIPE_BASE_URI" + + // EnvKeyQueue carries the queue name so the pipeline script can select + // queue-specific test targets. + EnvKeyQueue = "STOVEPIPE_QUEUE" + + // EnvKeyMetadata carries the JSON-encoded BuildMetadata provided by the + // caller to Trigger. Buildkite echoes env vars on the build object, so + // Status can recover and return the original metadata without local state. + EnvKeyMetadata = "STOVEPIPE_METADATA" +) + +// runner implements buildrunner.BuildRunner. +type runner struct { + cfg buildrunner.Config + client *platformbuildkite.Client + logger *zap.SugaredLogger +} + +var _ buildrunner.BuildRunner = (*runner)(nil) + +// Params holds the dependencies for a Buildkite BuildRunner. +type Params struct { + // Config holds the per-queue identity for this BuildRunner. + Config buildrunner.Config + // Client is a pre-constructed Buildkite client. The wiring layer builds + // it once via platformbuildkite.NewClient, with the pipeline's base URL + // (via platform/http.BaseURLTransport) and auth already configured. + Client *platformbuildkite.Client + // Logger is the structured logger. + Logger *zap.SugaredLogger +} + +// NewBuildRunner constructs a Buildkite-backed BuildRunner bound to a single +// pipeline. +func NewBuildRunner(params Params) (buildrunner.BuildRunner, error) { + if params.Client == nil { + return nil, fmt.Errorf("buildkite client is required") + } + if params.Logger == nil { + return nil, fmt.Errorf("logger is required") + } + return newRunner(params.Config, params.Client, params.Logger.Named("buildkite_buildrunner")), nil +} + +// newRunner constructs a runner. Used by NewBuildRunner and by tests. +func newRunner(cfg buildrunner.Config, c *platformbuildkite.Client, logger *zap.SugaredLogger) *runner { + return &runner{ + cfg: cfg, + client: c, + logger: logger, + } +} + +// Trigger calls the Buildkite API to create the build and returns the +// Buildkite build number as the build ID. Errors are propagated to the +// caller so the queue consumer can nack and retry. +func (r *runner) Trigger(ctx context.Context, headURI, baseURI string, metadata entity.BuildMetadata) (entity.BuildID, error) { + env := map[string]string{ + EnvKeyHeadURI: headURI, + EnvKeyBaseURI: baseURI, + EnvKeyQueue: r.cfg.QueueName, + } + if len(metadata) > 0 { + metaJSON, _ := json.Marshal(metadata) + env[EnvKeyMetadata] = string(metaJSON) + } + + req := platformbuildkite.CreateBuildRequest{ + Message: "stovepipe build", + Env: env, + } + + resp, err := r.client.CreateBuild(ctx, req) + if err != nil { + return entity.BuildID{}, fmt.Errorf("buildkite: create build: %w", err) + } + + r.logger.Debugw("triggered Buildkite build", + "buildkite_number", resp.Number, + ) + return entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(resp.Number)}, nil +} + +// Status fetches the current state of the build from Buildkite and returns it +// with the build URL and any caller-supplied metadata in BuildMetadata. +func (r *runner) Status(ctx context.Context, buildID entity.BuildID) (entity.BuildStatus, entity.BuildMetadata, error) { + number, err := platformbuildkite.ParseBuildNumber(buildID.ID) + if err != nil { + return entity.BuildStatusUnknown, nil, fmt.Errorf("buildkite: malformed build ID: %w", err) + } + + resp, err := r.client.GetBuild(ctx, number) + if err != nil { + return entity.BuildStatusUnknown, nil, fmt.Errorf("buildkite: get build: %w", err) + } + + meta := decodeMetadata(resp.Env) + meta["url"] = resp.WebURL + return mapState(resp.State), meta, nil +} + +// Cancel calls the Buildkite API to cancel the build. A no-op on +// already-terminal builds (Buildkite returns 422 for those). +func (r *runner) Cancel(ctx context.Context, buildID entity.BuildID) error { + number, err := platformbuildkite.ParseBuildNumber(buildID.ID) + if err != nil { + return fmt.Errorf("buildkite: malformed build ID: %w", err) + } + + if err := r.client.CancelBuild(ctx, number); err != nil { + return fmt.Errorf("buildkite: cancel build: %w", err) + } + r.logger.Debugw("cancelled Buildkite build", + "buildkite_number", number, + ) + return nil +} + +// decodeMetadata recovers the caller-supplied BuildMetadata from the env vars +// Buildkite echoes back on the build object. +func decodeMetadata(env map[string]string) entity.BuildMetadata { + return entity.BuildMetadata(platformbuildkite.DecodeMetadataEnv(env, EnvKeyMetadata)) +} + +// mapState maps a Buildkite build state to a BuildStatus. +func mapState(state string) entity.BuildStatus { + switch platformbuildkite.ParseState(state) { + case platformbuildkite.StateAccepted: + return entity.BuildStatusAccepted + case platformbuildkite.StateRunning: + return entity.BuildStatusRunning + case platformbuildkite.StateSucceeded: + return entity.BuildStatusSucceeded + case platformbuildkite.StateFailed: + return entity.BuildStatusFailed + case platformbuildkite.StateCancelled: + return entity.BuildStatusCancelled + default: + // Unrecognised Buildkite state. Do NOT assume terminal: Unknown is + // non-terminal, so the buildsignal poll loop keeps waiting rather than + // failing the build on a state this code does not understand. + return entity.BuildStatusUnknown + } +} diff --git a/stovepipe/extension/buildrunner/buildkite/buildkite_test.go b/stovepipe/extension/buildrunner/buildkite/buildkite_test.go new file mode 100644 index 00000000..06093714 --- /dev/null +++ b/stovepipe/extension/buildrunner/buildkite/buildkite_test.go @@ -0,0 +1,274 @@ +// Copyright (c) 2025 Uber Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package buildkite + +import ( + "context" + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "testing" + + "go.uber.org/zap" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + platformbuildkite "github.com/uber/submitqueue/platform/extension/buildrunner/buildkite" + phttp "github.com/uber/submitqueue/platform/http" + "github.com/uber/submitqueue/stovepipe/entity" + "github.com/uber/submitqueue/stovepipe/extension/buildrunner" +) + +// newTestRunner creates a runner backed by a test HTTP server. +func newTestRunner(t *testing.T, handler http.Handler) *runner { + t.Helper() + srv := httptest.NewServer(handler) + t.Cleanup(srv.Close) + c, err := phttp.NewClient(srv.URL) + require.NoError(t, err) + return newRunner( + buildrunner.Config{QueueName: "my-queue"}, + platformbuildkite.NewClient(c), + zap.NewNop().Sugar(), + ) +} + +// buildJSON encodes fields into a minimal Buildkite build JSON response. +func buildJSON(number int, state, webURL string) []byte { + return buildJSONWithEnv(number, state, webURL, nil) +} + +// buildJSONWithEnv encodes fields into a Buildkite build JSON response including env vars. +func buildJSONWithEnv(number int, state, webURL string, env map[string]string) []byte { + b, _ := json.Marshal(platformbuildkite.BuildResponse{Number: number, State: state, WebURL: webURL, Env: env}) + return b +} + +// --- Interface / constructor --- + +func TestNew_ImplementsInterface(t *testing.T) { + r, err := NewBuildRunner(Params{Logger: zap.NewNop().Sugar()}) + require.Error(t, err, "buildkite client is required") + var _ buildrunner.BuildRunner = r +} + +// --- Trigger --- + +func TestTrigger_SubmitsCorrectPayloadAndReturnsBuildkiteNumber(t *testing.T) { + var capturedMethod string + var capturedBody []byte + + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + capturedMethod = req.Method + capturedBody, _ = io.ReadAll(req.Body) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(buildJSON(42, "scheduled", "https://buildkite.com/test-org/my-pipeline/builds/42")) + })) + + id, err := r.Trigger(context.Background(), "github://repo/head/aaa", "github://repo/base/bbb", nil) + require.NoError(t, err) + assert.Equal(t, platformbuildkite.EncodeBuildNumber(42), id.ID) + + assert.Equal(t, http.MethodPost, capturedMethod) + + var req platformbuildkite.CreateBuildRequest + require.NoError(t, json.Unmarshal(capturedBody, &req)) + assert.Equal(t, "github://repo/head/aaa", req.Env[EnvKeyHeadURI]) + assert.Equal(t, "github://repo/base/bbb", req.Env[EnvKeyBaseURI]) + assert.Equal(t, "my-queue", req.Env[EnvKeyQueue]) +} + +func TestTrigger_EmptyBaseURI_FullBuild(t *testing.T) { + var capturedBody []byte + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + capturedBody, _ = io.ReadAll(req.Body) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(buildJSON(1, "scheduled", "")) + })) + + _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + require.NoError(t, err) + + var req platformbuildkite.CreateBuildRequest + require.NoError(t, json.Unmarshal(capturedBody, &req)) + assert.Equal(t, "", req.Env[EnvKeyBaseURI]) +} + +func TestTrigger_BuildkiteError_ReturnsError(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + + _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + require.Error(t, err) +} + +func TestTrigger_WithMetadata_SetsEnvVar(t *testing.T) { + var capturedBody []byte + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + capturedBody, _ = io.ReadAll(req.Body) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(buildJSON(10, "scheduled", "")) + })) + + metadata := entity.BuildMetadata{"requester": "alice", "ticket": "SQ-42"} + _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", metadata) + require.NoError(t, err) + + var req platformbuildkite.CreateBuildRequest + require.NoError(t, json.Unmarshal(capturedBody, &req)) + require.Contains(t, req.Env, EnvKeyMetadata) + + var got entity.BuildMetadata + require.NoError(t, json.Unmarshal([]byte(req.Env[EnvKeyMetadata]), &got)) + assert.Equal(t, metadata, got) +} + +func TestTrigger_NilMetadata_NoMetadataEnvVar(t *testing.T) { + var capturedBody []byte + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + capturedBody, _ = io.ReadAll(req.Body) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(buildJSON(11, "scheduled", "")) + })) + + _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + require.NoError(t, err) + + var req platformbuildkite.CreateBuildRequest + require.NoError(t, json.Unmarshal(capturedBody, &req)) + assert.NotContains(t, req.Env, EnvKeyMetadata) +} + +// --- Status --- + +func TestStatus_StateMapping(t *testing.T) { + tests := []struct { + bkState string + want entity.BuildStatus + }{ + {"creating", entity.BuildStatusAccepted}, + {"scheduled", entity.BuildStatusAccepted}, + {"running", entity.BuildStatusRunning}, + {"blocked", entity.BuildStatusRunning}, + {"passed", entity.BuildStatusSucceeded}, + {"failed", entity.BuildStatusFailed}, + {"not_run", entity.BuildStatusFailed}, + {"skipped", entity.BuildStatusFailed}, + {"canceling", entity.BuildStatusCancelled}, + {"canceled", entity.BuildStatusCancelled}, + // Unrecognised states map to the non-terminal Unknown, not Failed, so a + // state this code doesn't know about doesn't terminally fail a build. + {"some_future_state", entity.BuildStatusUnknown}, + {"", entity.BuildStatusUnknown}, + } + for _, tt := range tests { + t.Run(tt.bkState, func(t *testing.T) { + assert.Equal(t, tt.want, mapState(tt.bkState)) + }) + } +} + +func TestStatus_ReturnsLiveBuildkiteState(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(buildJSON(7, "running", "https://buildkite.com/test-org/my-pipeline/builds/7")) + })) + + status, meta, err := r.Status(context.Background(), entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(7)}) + require.NoError(t, err) + assert.Equal(t, entity.BuildStatusRunning, status) + assert.Equal(t, "https://buildkite.com/test-org/my-pipeline/builds/7", meta["url"]) +} + +func TestStatus_MalformedBuildID(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + + _, _, err := r.Status(context.Background(), entity.BuildID{ID: "not-a-number"}) + require.Error(t, err) +} + +func TestStatus_BuildkiteNotFound(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusNotFound) + })) + + _, _, err := r.Status(context.Background(), entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(99)}) + require.Error(t, err) +} + +func TestStatus_EchosCallerMetadata(t *testing.T) { + metadata := entity.BuildMetadata{"requester": "alice", "ticket": "SQ-42"} + metaJSON, _ := json.Marshal(metadata) + + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(buildJSONWithEnv(7, "passed", "https://buildkite.com/test-org/my-pipeline/builds/7", + map[string]string{EnvKeyMetadata: string(metaJSON)}, + )) + })) + + _, meta, err := r.Status(context.Background(), entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(7)}) + require.NoError(t, err) + assert.Equal(t, "alice", meta["requester"]) + assert.Equal(t, "SQ-42", meta["ticket"]) + assert.Equal(t, "https://buildkite.com/test-org/my-pipeline/builds/7", meta["url"]) +} + +func TestStatus_NoMetadata_ReturnsOnlyURL(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(buildJSON(8, "running", "https://buildkite.com/test-org/my-pipeline/builds/8")) + })) + + _, meta, err := r.Status(context.Background(), entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(8)}) + require.NoError(t, err) + assert.Equal(t, entity.BuildMetadata{"url": "https://buildkite.com/test-org/my-pipeline/builds/8"}, meta) +} + +// --- Cancel --- + +func TestCancel_CallsBuildkite(t *testing.T) { + var cancelCalled bool + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + cancelCalled = true + w.WriteHeader(http.StatusOK) + _, _ = w.Write(buildJSON(5, "canceled", "")) + })) + + require.NoError(t, r.Cancel(context.Background(), entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(5)})) + assert.True(t, cancelCalled) +} + +func TestCancel_AlreadyTerminal_Noop(t *testing.T) { + // Buildkite returns 422 when the build cannot be cancelled (already terminal). + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusUnprocessableEntity) + })) + + require.NoError(t, r.Cancel(context.Background(), entity.BuildID{ID: platformbuildkite.EncodeBuildNumber(5)})) +} + +func TestCancel_MalformedBuildID(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + + require.Error(t, r.Cancel(context.Background(), entity.BuildID{ID: "not-a-number"})) +} From a31f44da2da9fe764032a2f4c7b6430e284a7a62 Mon Sep 17 00:00:00 2001 From: "chenghan.ying" Date: Mon, 20 Jul 2026 18:04:20 +0000 Subject: [PATCH 4/7] feat(platform): add shared githubactions extension client --- .../buildrunner/githubactions/BUILD.bazel | 19 ++ .../buildrunner/githubactions/README.md | 26 ++ .../buildrunner/githubactions/client.go | 308 ++++++++++++++++++ .../buildrunner/githubactions/client_test.go | 222 +++++++++++++ 4 files changed, 575 insertions(+) create mode 100644 platform/extension/buildrunner/githubactions/BUILD.bazel create mode 100644 platform/extension/buildrunner/githubactions/README.md create mode 100644 platform/extension/buildrunner/githubactions/client.go create mode 100644 platform/extension/buildrunner/githubactions/client_test.go diff --git a/platform/extension/buildrunner/githubactions/BUILD.bazel b/platform/extension/buildrunner/githubactions/BUILD.bazel new file mode 100644 index 00000000..7e8b938c --- /dev/null +++ b/platform/extension/buildrunner/githubactions/BUILD.bazel @@ -0,0 +1,19 @@ +load("@rules_go//go:def.bzl", "go_library", "go_test") + +go_library( + name = "go_default_library", + srcs = ["client.go"], + importpath = "github.com/uber/submitqueue/platform/extension/buildrunner/githubactions", + visibility = ["//visibility:public"], +) + +go_test( + name = "go_default_test", + srcs = ["client_test.go"], + embed = [":go_default_library"], + deps = [ + "//platform/http:go_default_library", + "@com_github_stretchr_testify//assert:go_default_library", + "@com_github_stretchr_testify//require:go_default_library", + ], +) diff --git a/platform/extension/buildrunner/githubactions/README.md b/platform/extension/buildrunner/githubactions/README.md new file mode 100644 index 00000000..c5532816 --- /dev/null +++ b/platform/extension/buildrunner/githubactions/README.md @@ -0,0 +1,26 @@ +# GitHub Actions client + +Shared HTTP client and GitHub Actions-specific facts for every domain's GitHub Actions-backed `BuildRunner`. There is no `BuildRunner` interface here by design — each domain (`submitqueue`, `stovepipe`, ...) defines its own `BuildRunner` and its own `BuildStatus`, and adapts this package's `RunStatus` to it. Mirrors [`platform/extension/buildrunner/buildkite`](../buildkite/README.md)'s split; see that package's README and [`doc/rfc/stovepipe/steps/build.md`](../../../../doc/rfc/stovepipe/steps/build.md#alternatives-considered-for-sharing-the-contract) for the shared rationale. + +## What lives here + +- `Client`: `DispatchWorkflow` / `GetRun` / `CancelRun` against the GitHub Actions REST API, bound to a single repository and workflow. Construct one with `NewClient(httpClient, owner, repo, workflowID)`, where `httpClient` already has the GitHub API root (via `platform/http.BaseURLTransport`, typically `https://api.github.com`) and a token with `actions:read`/`actions:write` configured. +- `RunStatus` / `ParseRunStatus`: GitHub's own run status/conclusion vocabulary (`queued`/`in_progress`/`completed` crossed with a conclusion), collapsed into the five states every domain's `BuildStatus` already distinguishes. Each domain still does its own trivial `RunStatus` → its own `BuildStatus` switch — this package does not know either domain's entity types. +- `EncodeRunID` / `ParseRunID`: the run-id-as-build-id convention, mirroring Buildkite's build-number encoding. +- `RunMetadata`: builds the caller-facing metadata map (run id, attempt, status, conclusion, title, URL, branch, created-at) from a `WorkflowRun`. + +## Who consumes it + +- `submitqueue/extension/buildrunner/githubactions` — batch-identity `BuildRunner`, resolves changes via `changeset.Resolver` before dispatching. +- `stovepipe/extension/buildrunner/githubactions` — URI-identity `BuildRunner`, dispatches directly from `headURI`/`baseURI`. + +Both wrap a `*Client` built at the wiring layer; this package never constructs one from raw config (no credentials, no repo/workflow) itself. + +## How the same `Client` stays safe to share across two different checkout strategies + +As with the Buildkite client, `Client.DispatchWorkflow` never inspects `DispatchWorkflowRequest.Inputs` or picks a strategy — the split happens entirely outside this package: + +1. **Each domain's own adapter shapes a different `Inputs` payload.** SubmitQueue's runner sends `sq_base_uris`/`sq_head_uris` as JSON-encoded arrays of resolved change URIs; stovepipe's sends `stovepipe_head_uri`/`stovepipe_base_uri` as single plain strings. Both just call `Client.DispatchWorkflow` with their own `Inputs` map — this package treats it as an opaque `map[string]string`. +2. **Each domain's `Client` targets a different repository/workflow**, bound once at wiring time via `NewClient`'s `owner`/`repo`/`workflowID` arguments. Each workflow file (owned by the target repository, outside this package) is written against its domain's input contract — one applies patches into composite commits, the other checks out `stovepipe_head_uri` and diffs against `stovepipe_base_uri`. + +So the checkout strategy is a static, wiring-time binding — one `Client` instance ↔ one repo/workflow ↔ one workflow definition ↔ one input contract ↔ one domain's adapter — never a runtime decision made anywhere in Go code. diff --git a/platform/extension/buildrunner/githubactions/client.go b/platform/extension/buildrunner/githubactions/client.go new file mode 100644 index 00000000..f66b2e60 --- /dev/null +++ b/platform/extension/buildrunner/githubactions/client.go @@ -0,0 +1,308 @@ +// Copyright (c) 2025 Uber Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +// Package githubactions provides the HTTP client and GitHub Actions-specific +// facts (run status/conclusion vocabulary and run id encoding) shared by +// every domain's GitHub Actions-backed BuildRunner. It intentionally holds no +// BuildRunner interface or domain entity types — each domain (submitqueue, +// stovepipe, ...) defines its own BuildRunner and its own BuildStatus, and +// adapts this package's RunStatus to it. See +// platform/extension/buildrunner/buildkite's README for the analogous +// rationale applied to the Buildkite backend. +package githubactions + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/url" + "strconv" +) + +const githubAPIVersion = "2026-03-10" + +// Client is a thin wrapper around the GitHub Actions REST endpoints a +// BuildRunner needs: workflow dispatch, get workflow run, and cancel run. It +// is bound to a single repository and workflow. +type Client struct { + httpClient *http.Client + owner string + repo string + workflowID string +} + +// NewClient wraps a pre-configured *http.Client as a GitHub Actions Client +// bound to one repository and workflow. The caller is responsible for base +// URL resolution (e.g. via platform/http.BaseURLTransport, typically +// "https://api.github.com") and auth (a token with actions:read/actions:write +// injected via a transport). +func NewClient(httpClient *http.Client, owner, repo, workflowID string) *Client { + return &Client{ + httpClient: httpClient, + owner: owner, + repo: repo, + workflowID: workflowID, + } +} + +// Owner is the repository owner or organization this Client is bound to. +func (c *Client) Owner() string { return c.owner } + +// Repo is the repository name this Client is bound to. +func (c *Client) Repo() string { return c.repo } + +// WorkflowID is the workflow file name or numeric workflow ID this Client is +// bound to. +func (c *Client) WorkflowID() string { return c.workflowID } + +// DispatchWorkflowRequest is the payload for POST .../dispatches. +type DispatchWorkflowRequest struct { + Ref string `json:"ref"` + ReturnRunDetails bool `json:"return_run_details,omitempty"` + Inputs map[string]string `json:"inputs,omitempty"` +} + +// DispatchWorkflowResponse is the subset of fields callers need from a +// dispatch-workflow response. +type DispatchWorkflowResponse struct { + WorkflowRunID int64 `json:"workflow_run_id"` + RunURL string `json:"run_url"` + HTMLURL string `json:"html_url"` +} + +// WorkflowRun is the subset of a GitHub Actions workflow run object callers +// need. +type WorkflowRun struct { + ID int64 `json:"id"` + Name string `json:"name"` + DisplayTitle string `json:"display_title"` + Status string `json:"status"` + Conclusion string `json:"conclusion"` + HTMLURL string `json:"html_url"` + RunAttempt int `json:"run_attempt"` + Event string `json:"event"` + HeadBranch string `json:"head_branch"` + CreatedAt string `json:"created_at"` +} + +// DispatchWorkflow dispatches the bound workflow. +func (c *Client) DispatchWorkflow(ctx context.Context, req DispatchWorkflowRequest) (DispatchWorkflowResponse, error) { + body, err := json.Marshal(req) + if err != nil { + return DispatchWorkflowResponse{}, fmt.Errorf("marshal request: %w", err) + } + var resp DispatchWorkflowResponse + if err := c.do(ctx, http.MethodPost, c.workflowPath("dispatches"), body, &resp); err != nil { + return DispatchWorkflowResponse{}, err + } + return resp, nil +} + +// GetRun fetches a workflow run by its GitHub run id. +func (c *Client) GetRun(ctx context.Context, runID int64) (WorkflowRun, error) { + var run WorkflowRun + if err := c.do(ctx, http.MethodGet, c.runPath(runID), nil, &run); err != nil { + return WorkflowRun{}, err + } + return run, nil +} + +// CancelRun requests cancellation of the workflow run. Returns nil when the +// run is already terminal or otherwise not cancellable (HTTP 409/422) — the +// BuildRunner contract treats that as a no-op. +func (c *Client) CancelRun(ctx context.Context, runID int64) error { + req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.runPath(runID)+"/cancel", nil) + if err != nil { + return fmt.Errorf("create request: %w", err) + } + c.setHeaders(req) + + resp, err := c.httpClient.Do(req) + if err != nil { + return fmt.Errorf("send request: %w", err) + } + defer resp.Body.Close() + _, _ = io.Copy(io.Discard, resp.Body) + + switch resp.StatusCode { + case http.StatusAccepted, http.StatusOK, http.StatusCreated, http.StatusNoContent: + return nil + case http.StatusConflict, http.StatusUnprocessableEntity: + // Already terminal or otherwise not cancellable: no-op per + // BuildRunner.Cancel contract. + return nil + case http.StatusNotFound: + return fmt.Errorf("workflow run not found") + default: + return fmt.Errorf("unexpected status %d from cancel", resp.StatusCode) + } +} + +func (c *Client) workflowPath(suffix string) string { + path := fmt.Sprintf( + "/repos/%s/%s/actions/workflows/%s", + url.PathEscape(c.owner), + url.PathEscape(c.repo), + url.PathEscape(c.workflowID), + ) + if suffix != "" { + path += "/" + suffix + } + return path +} + +func (c *Client) runPath(runID int64) string { + return fmt.Sprintf( + "/repos/%s/%s/actions/runs/%s", + url.PathEscape(c.owner), + url.PathEscape(c.repo), + url.PathEscape(strconv.FormatInt(runID, 10)), + ) +} + +func (c *Client) do(ctx context.Context, method, rawURL string, body []byte, out any) error { + var bodyReader io.Reader + if body != nil { + bodyReader = bytes.NewReader(body) + } + + req, err := http.NewRequestWithContext(ctx, method, rawURL, bodyReader) + if err != nil { + return fmt.Errorf("create request: %w", err) + } + c.setHeaders(req) + + resp, err := c.httpClient.Do(req) + if err != nil { + return fmt.Errorf("send request: %w", err) + } + defer resp.Body.Close() + + respBody, err := io.ReadAll(resp.Body) + if err != nil { + return fmt.Errorf("read response: %w", err) + } + + if resp.StatusCode == http.StatusNotFound { + return fmt.Errorf("GitHub API returned 404 for %s %s", method, rawURL) + } + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return fmt.Errorf("API returned status %d: %s", resp.StatusCode, respBody) + } + + if out != nil && len(respBody) > 0 { + if err := json.Unmarshal(respBody, out); err != nil { + return fmt.Errorf("unmarshal response: %w", err) + } + } + return nil +} + +func (c *Client) setHeaders(req *http.Request) { + req.Header.Set("Accept", "application/vnd.github+json") + req.Header.Set("Content-Type", "application/json") + req.Header.Set("X-GitHub-Api-Version", githubAPIVersion) +} + +// EncodeRunID encodes a GitHub Actions workflow run id as an opaque build id +// string. +func EncodeRunID(runID int64) string { + return strconv.FormatInt(runID, 10) +} + +// ParseRunID is the inverse of EncodeRunID. +func ParseRunID(id string) (int64, error) { + runID, err := strconv.ParseInt(id, 10, 64) + if err != nil || runID <= 0 { + return 0, fmt.Errorf("invalid build ID %q", id) + } + return runID, nil +} + +// RunStatus is GitHub Actions' own run status/conclusion vocabulary, +// collapsed into the five states every domain's BuildStatus already +// distinguishes. +type RunStatus string + +const ( + // RunStatusUnknown is returned for a raw status/conclusion pair this + // package does not recognize. Not terminal — callers should keep polling + // rather than treat it as a final outcome. + RunStatusUnknown RunStatus = "" + // RunStatusAccepted means the run has been accepted for execution but has + // not started yet (GitHub: queued, requested, waiting, pending). + RunStatusAccepted RunStatus = "accepted" + // RunStatusRunning means the run is currently executing (GitHub: + // in_progress). + RunStatusRunning RunStatus = "running" + // RunStatusSucceeded means the run completed successfully (GitHub: + // completed with conclusion success). + RunStatusSucceeded RunStatus = "succeeded" + // RunStatusFailed means the run completed without a passing result + // (GitHub: completed with any conclusion other than success/cancelled). + RunStatusFailed RunStatus = "failed" + // RunStatusCancelled means the run was cancelled (GitHub: completed with + // conclusion cancelled). + RunStatusCancelled RunStatus = "cancelled" +) + +// ParseRunStatus maps a raw GitHub Actions run status/conclusion pair to a +// RunStatus. conclusion is only consulted when status is "completed". An +// unrecognized status, or a completed run with an empty conclusion, maps to +// RunStatusUnknown rather than being assumed terminal. +func ParseRunStatus(status, conclusion string) RunStatus { + switch status { + case "queued", "requested", "waiting", "pending": + return RunStatusAccepted + case "in_progress": + return RunStatusRunning + case "completed": + switch conclusion { + case "success": + return RunStatusSucceeded + case "cancelled": + return RunStatusCancelled + case "": + return RunStatusUnknown + default: + return RunStatusFailed + } + default: + return RunStatusUnknown + } +} + +// RunMetadata builds the caller-facing metadata map for a workflow run: the +// run's own identity and outcome fields, keyed for direct use as (or merge +// into) a domain's BuildMetadata. +func RunMetadata(run WorkflowRun) map[string]string { + meta := map[string]string{ + "github_run_id": strconv.FormatInt(run.ID, 10), + "github_run_attempt": strconv.Itoa(run.RunAttempt), + "github_status": run.Status, + "github_conclusion": run.Conclusion, + "github_display_title": run.DisplayTitle, + "url": run.HTMLURL, + } + if run.HeadBranch != "" { + meta["github_head_branch"] = run.HeadBranch + } + if run.CreatedAt != "" { + meta["github_created_at"] = run.CreatedAt + } + return meta +} diff --git a/platform/extension/buildrunner/githubactions/client_test.go b/platform/extension/buildrunner/githubactions/client_test.go new file mode 100644 index 00000000..82f0ce52 --- /dev/null +++ b/platform/extension/buildrunner/githubactions/client_test.go @@ -0,0 +1,222 @@ +// Copyright (c) 2025 Uber Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package githubactions + +import ( + "context" + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "strconv" + "strings" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + phttp "github.com/uber/submitqueue/platform/http" +) + +// newTestClient creates a Client backed by a test HTTP server. +func newTestClient(t *testing.T, handler http.Handler) *Client { + t.Helper() + srv := httptest.NewServer(handler) + t.Cleanup(srv.Close) + c, err := phttp.NewClient(srv.URL) + require.NoError(t, err) + return NewClient(c, "uber", "submitqueue", "submitqueue-ci.yml") +} + +// --- NewClient / accessors --- + +func TestNewClient_ExposesIdentity(t *testing.T) { + c := NewClient(http.DefaultClient, "uber", "submitqueue", "submitqueue-ci.yml") + assert.Equal(t, "uber", c.Owner()) + assert.Equal(t, "submitqueue", c.Repo()) + assert.Equal(t, "submitqueue-ci.yml", c.WorkflowID()) +} + +// --- DispatchWorkflow --- + +func TestDispatchWorkflow_SubmitsPayloadAndReturnsResponse(t *testing.T) { + var capturedMethod, capturedPath string + var capturedBody []byte + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + capturedMethod = req.Method + capturedPath = req.URL.Path + capturedBody, _ = io.ReadAll(req.Body) + _ = json.NewEncoder(w).Encode(DispatchWorkflowResponse{ + WorkflowRunID: 42, + RunURL: "https://api.github.com/repos/uber/submitqueue/actions/runs/42", + HTMLURL: "https://github.com/uber/submitqueue/actions/runs/42", + }) + })) + + resp, err := c.DispatchWorkflow(context.Background(), DispatchWorkflowRequest{ + Ref: "main", + ReturnRunDetails: true, + Inputs: map[string]string{"key": "value"}, + }) + require.NoError(t, err) + assert.Equal(t, http.MethodPost, capturedMethod) + assert.Equal(t, "/repos/uber/submitqueue/actions/workflows/submitqueue-ci.yml/dispatches", capturedPath) + assert.Equal(t, int64(42), resp.WorkflowRunID) + + var req DispatchWorkflowRequest + require.NoError(t, json.Unmarshal(capturedBody, &req)) + assert.Equal(t, "main", req.Ref) + assert.True(t, req.ReturnRunDetails) + assert.Equal(t, "value", req.Inputs["key"]) +} + +func TestDispatchWorkflow_ErrorStatus_ReturnsError(t *testing.T) { + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + + _, err := c.DispatchWorkflow(context.Background(), DispatchWorkflowRequest{}) + require.Error(t, err) +} + +// --- GetRun --- + +func TestGetRun_ReturnsResponse(t *testing.T) { + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + assert.Equal(t, http.MethodGet, req.Method) + assert.Equal(t, "/repos/uber/submitqueue/actions/runs/42", req.URL.Path) + _ = json.NewEncoder(w).Encode(WorkflowRun{ID: 42, Status: "in_progress"}) + })) + + run, err := c.GetRun(context.Background(), 42) + require.NoError(t, err) + assert.Equal(t, "in_progress", run.Status) +} + +func TestGetRun_NotFound_ReturnsError(t *testing.T) { + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusNotFound) + })) + + _, err := c.GetRun(context.Background(), 99) + require.Error(t, err) +} + +// --- CancelRun --- + +func TestCancelRun_CallsGitHub(t *testing.T) { + var cancelledRunID int64 + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + require.True(t, req.Method == http.MethodPost && strings.HasSuffix(req.URL.Path, "/cancel")) + parts := strings.Split(req.URL.Path, "/") + id, err := strconv.ParseInt(parts[len(parts)-2], 10, 64) + require.NoError(t, err) + cancelledRunID = id + w.WriteHeader(http.StatusAccepted) + })) + + require.NoError(t, c.CancelRun(context.Background(), 99)) + assert.Equal(t, int64(99), cancelledRunID) +} + +func TestCancelRun_AlreadyTerminal_Noop(t *testing.T) { + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusUnprocessableEntity) + })) + + require.NoError(t, c.CancelRun(context.Background(), 5)) +} + +func TestCancelRun_NotFound_ReturnsError(t *testing.T) { + c := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusNotFound) + })) + + require.Error(t, c.CancelRun(context.Background(), 5)) +} + +// --- EncodeRunID / ParseRunID --- + +func TestEncodeParseRunID_RoundTrip(t *testing.T) { + for _, id := range []int64{1, 9999} { + got, err := ParseRunID(EncodeRunID(id)) + require.NoError(t, err) + assert.Equal(t, id, got) + } +} + +func TestParseRunID_Invalid(t *testing.T) { + for _, id := range []string{"", "notanumber", "0", "-1", "gha-old-id"} { + t.Run(id, func(t *testing.T) { + _, err := ParseRunID(id) + require.Error(t, err) + }) + } +} + +// --- ParseRunStatus --- + +func TestParseRunStatus(t *testing.T) { + tests := []struct { + name string + status string + conclusion string + want RunStatus + }{ + {name: "queued", status: "queued", want: RunStatusAccepted}, + {name: "pending", status: "pending", want: RunStatusAccepted}, + {name: "requested", status: "requested", want: RunStatusAccepted}, + {name: "waiting", status: "waiting", want: RunStatusAccepted}, + {name: "running", status: "in_progress", want: RunStatusRunning}, + {name: "success", status: "completed", conclusion: "success", want: RunStatusSucceeded}, + {name: "failure", status: "completed", conclusion: "failure", want: RunStatusFailed}, + {name: "timed out", status: "completed", conclusion: "timed_out", want: RunStatusFailed}, + {name: "cancelled", status: "completed", conclusion: "cancelled", want: RunStatusCancelled}, + {name: "completed missing conclusion", status: "completed", want: RunStatusUnknown}, + {name: "unknown", status: "some_future_state", want: RunStatusUnknown}, + {name: "empty", want: RunStatusUnknown}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, ParseRunStatus(tt.status, tt.conclusion)) + }) + } +} + +// --- RunMetadata --- + +func TestRunMetadata_IncludesIdentityAndOutcome(t *testing.T) { + meta := RunMetadata(WorkflowRun{ + ID: 42, + DisplayTitle: "SubmitQueue gha-trace", + Status: "in_progress", + HTMLURL: "https://github.com/uber/submitqueue/actions/runs/42", + RunAttempt: 2, + HeadBranch: "main", + CreatedAt: "2026-06-05T00:00:00Z", + }) + assert.Equal(t, "42", meta["github_run_id"]) + assert.Equal(t, "2", meta["github_run_attempt"]) + assert.Equal(t, "in_progress", meta["github_status"]) + assert.Equal(t, "https://github.com/uber/submitqueue/actions/runs/42", meta["url"]) + assert.Equal(t, "main", meta["github_head_branch"]) + assert.Equal(t, "2026-06-05T00:00:00Z", meta["github_created_at"]) +} + +func TestRunMetadata_OmitsEmptyOptionalFields(t *testing.T) { + meta := RunMetadata(WorkflowRun{ID: 7, Status: "queued"}) + assert.NotContains(t, meta, "github_head_branch") + assert.NotContains(t, meta, "github_created_at") +} From fa7875ff15da4124e24b79cd970ce58ddf17dc4f Mon Sep 17 00:00:00 2001 From: "chenghan.ying" Date: Mon, 20 Jul 2026 18:04:44 +0000 Subject: [PATCH 5/7] refactor(submitqueue): rebase githubactions adapter onto shared platform client --- .../buildrunner/githubactions/BUILD.bazel | 7 +- .../buildrunner/githubactions/README.md | 5 + .../buildrunner/githubactions/client.go | 178 ------------------ .../githubactions/githubactions.go | 140 ++++---------- .../githubactions/githubactions_test.go | 27 +-- 5 files changed, 57 insertions(+), 300 deletions(-) delete mode 100644 submitqueue/extension/buildrunner/githubactions/client.go diff --git a/submitqueue/extension/buildrunner/githubactions/BUILD.bazel b/submitqueue/extension/buildrunner/githubactions/BUILD.bazel index c63d3eca..61bda338 100644 --- a/submitqueue/extension/buildrunner/githubactions/BUILD.bazel +++ b/submitqueue/extension/buildrunner/githubactions/BUILD.bazel @@ -2,14 +2,12 @@ load("@rules_go//go:def.bzl", "go_library", "go_test") go_library( name = "go_default_library", - srcs = [ - "client.go", - "githubactions.go", - ], + srcs = ["githubactions.go"], importpath = "github.com/uber/submitqueue/submitqueue/extension/buildrunner/githubactions", visibility = ["//visibility:public"], deps = [ "//platform/base/change:go_default_library", + "//platform/extension/buildrunner/githubactions:go_default_library", "//submitqueue/core/changeset:go_default_library", "//submitqueue/entity:go_default_library", "//submitqueue/extension/buildrunner:go_default_library", @@ -23,6 +21,7 @@ go_test( embed = [":go_default_library"], deps = [ "//platform/base/change:go_default_library", + "//platform/extension/buildrunner/githubactions:go_default_library", "//platform/http:go_default_library", "//submitqueue/core/changeset:go_default_library", "//submitqueue/core/changeset/fake:go_default_library", diff --git a/submitqueue/extension/buildrunner/githubactions/README.md b/submitqueue/extension/buildrunner/githubactions/README.md index e69d1a2b..69802ac0 100644 --- a/submitqueue/extension/buildrunner/githubactions/README.md +++ b/submitqueue/extension/buildrunner/githubactions/README.md @@ -4,6 +4,11 @@ `workflow_dispatch`. It is intended to prove the SubmitQueue BuildRunner architecture against a common CI system without adding local state. +Its HTTP client and GitHub Actions-specific facts (run status/conclusion +vocabulary, run id encoding) live in +[`platform/extension/buildrunner/githubactions`](../../../../platform/extension/buildrunner/githubactions/README.md), +shared with `stovepipe`'s own GitHub Actions backend. + ## How it works 1. `Trigger` dispatches the configured workflow on a trusted ref, usually diff --git a/submitqueue/extension/buildrunner/githubactions/client.go b/submitqueue/extension/buildrunner/githubactions/client.go deleted file mode 100644 index 535cc105..00000000 --- a/submitqueue/extension/buildrunner/githubactions/client.go +++ /dev/null @@ -1,178 +0,0 @@ -// Copyright (c) 2025 Uber Technologies, Inc. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -package githubactions - -import ( - "bytes" - "context" - "encoding/json" - "fmt" - "io" - "net/http" - "net/url" - "strconv" -) - -const githubAPIVersion = "2026-03-10" - -// client is a thin wrapper around the GitHub Actions REST endpoints that -// BuildRunner needs: workflow dispatch, get workflow run, and cancel run. -type client struct { - httpClient *http.Client - owner string - repo string - workflowID string -} - -type dispatchWorkflowRequest struct { - Ref string `json:"ref"` - ReturnRunDetails bool `json:"return_run_details,omitempty"` - Inputs map[string]string `json:"inputs,omitempty"` -} - -type dispatchWorkflowResponse struct { - WorkflowRunID int64 `json:"workflow_run_id"` - RunURL string `json:"run_url"` - HTMLURL string `json:"html_url"` -} - -// workflowRun is the subset of a GitHub Actions workflow run object the -// runner needs. -type workflowRun struct { - ID int64 `json:"id"` - Name string `json:"name"` - DisplayTitle string `json:"display_title"` - Status string `json:"status"` - Conclusion string `json:"conclusion"` - HTMLURL string `json:"html_url"` - RunAttempt int `json:"run_attempt"` - Event string `json:"event"` - HeadBranch string `json:"head_branch"` - CreatedAt string `json:"created_at"` -} - -func (c *client) dispatchWorkflow(ctx context.Context, req dispatchWorkflowRequest) (dispatchWorkflowResponse, error) { - body, err := json.Marshal(req) - if err != nil { - return dispatchWorkflowResponse{}, fmt.Errorf("marshal request: %w", err) - } - var resp dispatchWorkflowResponse - if err := c.do(ctx, http.MethodPost, c.workflowPath("dispatches"), body, &resp); err != nil { - return dispatchWorkflowResponse{}, err - } - return resp, nil -} - -func (c *client) getRun(ctx context.Context, runID int64) (workflowRun, error) { - var run workflowRun - if err := c.do(ctx, http.MethodGet, c.runPath(runID), nil, &run); err != nil { - return workflowRun{}, err - } - return run, nil -} - -func (c *client) cancelRun(ctx context.Context, runID int64) error { - req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.runPath(runID)+"/cancel", nil) - if err != nil { - return fmt.Errorf("create request: %w", err) - } - c.setHeaders(req) - - resp, err := c.httpClient.Do(req) - if err != nil { - return fmt.Errorf("send request: %w", err) - } - defer resp.Body.Close() - _, _ = io.Copy(io.Discard, resp.Body) - - switch resp.StatusCode { - case http.StatusAccepted, http.StatusOK, http.StatusCreated, http.StatusNoContent: - return nil - case http.StatusConflict, http.StatusUnprocessableEntity: - // Already terminal or otherwise not cancellable: no-op per - // BuildRunner.Cancel contract. - return nil - case http.StatusNotFound: - return fmt.Errorf("workflow run not found") - default: - return fmt.Errorf("unexpected status %d from cancel", resp.StatusCode) - } -} - -func (c *client) workflowPath(suffix string) string { - path := fmt.Sprintf( - "/repos/%s/%s/actions/workflows/%s", - url.PathEscape(c.owner), - url.PathEscape(c.repo), - url.PathEscape(c.workflowID), - ) - if suffix != "" { - path += "/" + suffix - } - return path -} - -func (c *client) runPath(runID int64) string { - return fmt.Sprintf( - "/repos/%s/%s/actions/runs/%s", - url.PathEscape(c.owner), - url.PathEscape(c.repo), - url.PathEscape(strconv.FormatInt(runID, 10)), - ) -} - -func (c *client) do(ctx context.Context, method, rawURL string, body []byte, out any) error { - var bodyReader io.Reader - if body != nil { - bodyReader = bytes.NewReader(body) - } - - req, err := http.NewRequestWithContext(ctx, method, rawURL, bodyReader) - if err != nil { - return fmt.Errorf("create request: %w", err) - } - c.setHeaders(req) - - resp, err := c.httpClient.Do(req) - if err != nil { - return fmt.Errorf("send request: %w", err) - } - defer resp.Body.Close() - - respBody, err := io.ReadAll(resp.Body) - if err != nil { - return fmt.Errorf("read response: %w", err) - } - - if resp.StatusCode == http.StatusNotFound { - return fmt.Errorf("GitHub API returned 404 for %s %s", method, rawURL) - } - if resp.StatusCode < 200 || resp.StatusCode >= 300 { - return fmt.Errorf("API returned status %d: %s", resp.StatusCode, respBody) - } - - if out != nil && len(respBody) > 0 { - if err := json.Unmarshal(respBody, out); err != nil { - return fmt.Errorf("unmarshal response: %w", err) - } - } - return nil -} - -func (c *client) setHeaders(req *http.Request) { - req.Header.Set("Accept", "application/vnd.github+json") - req.Header.Set("Content-Type", "application/json") - req.Header.Set("X-GitHub-Api-Version", githubAPIVersion) -} diff --git a/submitqueue/extension/buildrunner/githubactions/githubactions.go b/submitqueue/extension/buildrunner/githubactions/githubactions.go index 4b1c70b3..b69a5bc4 100644 --- a/submitqueue/extension/buildrunner/githubactions/githubactions.go +++ b/submitqueue/extension/buildrunner/githubactions/githubactions.go @@ -24,12 +24,11 @@ import ( "context" "encoding/json" "fmt" - "net/http" - "strconv" "go.uber.org/zap" "github.com/uber/submitqueue/platform/base/change" + platformgithubactions "github.com/uber/submitqueue/platform/extension/buildrunner/githubactions" "github.com/uber/submitqueue/submitqueue/core/changeset" "github.com/uber/submitqueue/submitqueue/entity" "github.com/uber/submitqueue/submitqueue/extension/buildrunner" @@ -50,27 +49,19 @@ const ( defaultRef = "main" ) -// Params holds the dependencies and GitHub workflow identity for a GitHub -// Actions BuildRunner. +// Params holds the dependencies for a GitHub Actions BuildRunner. type Params struct { // Config holds the per-queue identity for this BuildRunner. Config buildrunner.Config - // HTTPClient is a pre-configured GitHub API client. The caller is responsible - // for base URL resolution (e.g. via platform/http.BaseURLTransport) and auth. - // The token needs actions:write to dispatch/cancel workflows and actions:read - // to poll status. - HTTPClient *http.Client + // Client is a pre-constructed GitHub Actions client, bound to one + // repository and workflow. The wiring layer builds it once via + // platformgithubactions.NewClient, with the GitHub API root (via + // platform/http.BaseURLTransport) and auth already configured. + Client *platformgithubactions.Client // Resolver resolves a batch's changes (base and head batches). Resolver changeset.Resolver // Logger is the structured logger. Logger *zap.SugaredLogger - // Owner is the repository owner or organization, for example "uber". - Owner string - // Repo is the repository name, for example "submitqueue". - Repo string - // WorkflowID is the workflow file name or numeric workflow ID accepted by - // the GitHub Actions API, for example "submitqueue-ci.yml". - WorkflowID string // Ref is the branch, tag, or SHA where the trusted workflow is read from. // Defaults to "main". Ref string @@ -83,7 +74,7 @@ type runner struct { cfg buildrunner.Config ref string extraInputs map[string]string - client *client + client *platformgithubactions.Client resolver changeset.Resolver logger *zap.SugaredLogger } @@ -93,29 +84,28 @@ var _ buildrunner.BuildRunner = (*runner)(nil) // NewBuildRunner constructs a GitHub Actions-backed BuildRunner bound to one // repository/workflow and one queue config. func NewBuildRunner(params Params) (buildrunner.BuildRunner, error) { - if err := validateConfig(params.HTTPClient, params.Logger, params.Owner, params.Repo, params.WorkflowID); err != nil { - return nil, err + if params.Client == nil { + return nil, fmt.Errorf("github actions client is required") } - if params.Ref == "" { - params.Ref = defaultRef + if params.Logger == nil { + return nil, fmt.Errorf("logger is required") + } + ref := params.Ref + if ref == "" { + ref = defaultRef } return newRunner( params.Config, - params.Ref, + ref, params.ExtraInputs, - &client{ - httpClient: params.HTTPClient, - owner: params.Owner, - repo: params.Repo, - workflowID: params.WorkflowID, - }, + params.Client, params.Resolver, params.Logger.Named("githubactions_buildrunner"), ), nil } -func newRunner(cfg buildrunner.Config, ref string, extraInputs map[string]string, c *client, resolver changeset.Resolver, logger *zap.SugaredLogger) *runner { +func newRunner(cfg buildrunner.Config, ref string, extraInputs map[string]string, c *platformgithubactions.Client, resolver changeset.Resolver, logger *zap.SugaredLogger) *runner { copied := make(map[string]string, len(extraInputs)) for k, v := range extraInputs { copied[k] = v @@ -130,25 +120,6 @@ func newRunner(cfg buildrunner.Config, ref string, extraInputs map[string]string } } -func validateConfig(httpClient *http.Client, logger *zap.SugaredLogger, owner, repo, workflowID string) error { - if httpClient == nil { - return fmt.Errorf("http client is required") - } - if logger == nil { - return fmt.Errorf("logger is required") - } - if owner == "" { - return fmt.Errorf("owner is required") - } - if repo == "" { - return fmt.Errorf("repo is required") - } - if workflowID == "" { - return fmt.Errorf("workflow ID is required") - } - return nil -} - // Trigger dispatches the configured GitHub Actions workflow and returns the // GitHub workflow run ID as the SubmitQueue build ID. Errors are propagated to // the caller so the queue consumer can nack and retry. @@ -167,7 +138,7 @@ func (r *runner) Trigger(ctx context.Context, base []entity.Batch, head entity.B return entity.BuildID{}, err } - resp, err := r.client.dispatchWorkflow(ctx, dispatchWorkflowRequest{ + resp, err := r.client.DispatchWorkflow(ctx, platformgithubactions.DispatchWorkflowRequest{ Ref: r.ref, ReturnRunDetails: true, Inputs: inputs, @@ -176,17 +147,17 @@ func (r *runner) Trigger(ctx context.Context, base []entity.Batch, head entity.B return entity.BuildID{}, fmt.Errorf("github actions: dispatch workflow: %w", err) } if resp.WorkflowRunID <= 0 { - return entity.BuildID{}, fmt.Errorf("github actions: dispatch workflow: response missing workflow_run_id (requires X-GitHub-Api-Version %s and return_run_details support)", githubAPIVersion) + return entity.BuildID{}, fmt.Errorf("github actions: dispatch workflow: response missing workflow_run_id (requires X-GitHub-Api-Version and return_run_details support)") } r.logger.Debugw("dispatched GitHub Actions workflow", - "owner", r.client.owner, - "repo", r.client.repo, - "workflow_id", r.client.workflowID, + "owner", r.client.Owner(), + "repo", r.client.Repo(), + "workflow_id", r.client.WorkflowID(), "ref", r.ref, "github_run_id", resp.WorkflowRunID, ) - return entity.BuildID{ID: strconv.FormatInt(resp.WorkflowRunID, 10)}, nil + return entity.BuildID{ID: platformgithubactions.EncodeRunID(resp.WorkflowRunID)}, nil } func (r *runner) dispatchInputs(base, head []change.Change, metadata entity.BuildMetadata) (map[string]string, error) { @@ -219,25 +190,25 @@ func (r *runner) dispatchInputs(base, head []change.Change, metadata entity.Buil // Status fetches the workflow run by GitHub run ID and maps the run's // status/conclusion onto BuildStatus. func (r *runner) Status(ctx context.Context, buildID entity.BuildID) (entity.BuildStatus, entity.BuildMetadata, error) { - runID, err := parseRunID(buildID.ID) + runID, err := platformgithubactions.ParseRunID(buildID.ID) if err != nil { return entity.BuildStatusUnknown, nil, fmt.Errorf("github actions: malformed build ID: %w", err) } - run, err := r.client.getRun(ctx, runID) + run, err := r.client.GetRun(ctx, runID) if err != nil { return entity.BuildStatusUnknown, nil, fmt.Errorf("github actions: get run: %w", err) } - return mapRunStatus(run.Status, run.Conclusion), metadataFromRun(run), nil + return mapRunStatus(run.Status, run.Conclusion), entity.BuildMetadata(platformgithubactions.RunMetadata(run)), nil } // Cancel requests cancellation of the workflow run identified by buildID. func (r *runner) Cancel(ctx context.Context, buildID entity.BuildID) error { - runID, err := parseRunID(buildID.ID) + runID, err := platformgithubactions.ParseRunID(buildID.ID) if err != nil { return fmt.Errorf("github actions: malformed build ID: %w", err) } - if err := r.client.cancelRun(ctx, runID); err != nil { + if err := r.client.CancelRun(ctx, runID); err != nil { return fmt.Errorf("github actions: cancel run: %w", err) } r.logger.Debugw("cancelled GitHub Actions workflow run", @@ -246,43 +217,18 @@ func (r *runner) Cancel(ctx context.Context, buildID entity.BuildID) error { return nil } -func metadataFromRun(run workflowRun) entity.BuildMetadata { - meta := entity.BuildMetadata{ - "github_run_id": strconv.FormatInt(run.ID, 10), - "github_run_attempt": strconv.Itoa(run.RunAttempt), - "github_status": run.Status, - "github_conclusion": run.Conclusion, - "github_display_title": run.DisplayTitle, - "url": run.HTMLURL, - } - if run.HeadBranch != "" { - meta["github_head_branch"] = run.HeadBranch - } - if run.CreatedAt != "" { - meta["github_created_at"] = run.CreatedAt - } - return meta -} - func mapRunStatus(status, conclusion string) entity.BuildStatus { - switch status { - case "queued", "requested", "waiting", "pending": + switch platformgithubactions.ParseRunStatus(status, conclusion) { + case platformgithubactions.RunStatusAccepted: return entity.BuildStatusAccepted - case "in_progress": + case platformgithubactions.RunStatusRunning: return entity.BuildStatusRunning - case "completed": - switch conclusion { - case "success": - return entity.BuildStatusSucceeded - case "cancelled": - return entity.BuildStatusCancelled - case "": - return entity.BuildStatusUnknown - default: - // Any completed non-success/non-cancelled conclusion is a terminal - // failure for SubmitQueue purposes. - return entity.BuildStatusFailed - } + case platformgithubactions.RunStatusSucceeded: + return entity.BuildStatusSucceeded + case platformgithubactions.RunStatusFailed: + return entity.BuildStatusFailed + case platformgithubactions.RunStatusCancelled: + return entity.BuildStatusCancelled default: return entity.BuildStatusUnknown } @@ -295,11 +241,3 @@ func flattenURIs(changes []change.Change) []string { } return uris } - -func parseRunID(id string) (int64, error) { - runID, err := strconv.ParseInt(id, 10, 64) - if err != nil || runID <= 0 { - return 0, fmt.Errorf("invalid build ID %q", id) - } - return runID, nil -} diff --git a/submitqueue/extension/buildrunner/githubactions/githubactions_test.go b/submitqueue/extension/buildrunner/githubactions/githubactions_test.go index d1067947..93a0b238 100644 --- a/submitqueue/extension/buildrunner/githubactions/githubactions_test.go +++ b/submitqueue/extension/buildrunner/githubactions/githubactions_test.go @@ -29,6 +29,7 @@ import ( "go.uber.org/zap" "github.com/uber/submitqueue/platform/base/change" + platformgithubactions "github.com/uber/submitqueue/platform/extension/buildrunner/githubactions" phttp "github.com/uber/submitqueue/platform/http" "github.com/uber/submitqueue/submitqueue/core/changeset" changesetfake "github.com/uber/submitqueue/submitqueue/core/changeset/fake" @@ -53,12 +54,7 @@ func newTestRunner(t *testing.T, handler http.Handler, resolver ...changeset.Res buildrunner.Config{QueueName: "my-queue"}, "main", map[string]string{"custom": "value"}, - &client{ - httpClient: c, - owner: "uber", - repo: "submitqueue", - workflowID: "submitqueue-ci.yml", - }, + platformgithubactions.NewClient(c, "uber", "submitqueue", "submitqueue-ci.yml"), r, zap.NewNop().Sugar(), ) @@ -67,17 +63,14 @@ func newTestRunner(t *testing.T, handler http.Handler, resolver ...changeset.Res func TestNewBuildRunner_ValidatesConfig(t *testing.T) { _, err := NewBuildRunner(Params{Logger: zap.NewNop().Sugar()}) require.Error(t, err) - assert.Contains(t, err.Error(), "http client is required") + assert.Contains(t, err.Error(), "github actions client is required") } func TestNewBuildRunner_BindsQueueConfigAndExtraInputs(t *testing.T) { br, err := NewBuildRunner(Params{ Config: buildrunner.Config{QueueName: "queue-a"}, - HTTPClient: http.DefaultClient, + Client: platformgithubactions.NewClient(http.DefaultClient, "uber", "submitqueue", "submitqueue-ci.yml"), Logger: zap.NewNop().Sugar(), - Owner: "uber", - Repo: "submitqueue", - WorkflowID: "submitqueue-ci.yml", Ref: "main", ExtraInputs: map[string]string{"runner": "ubuntu-latest"}, }) @@ -101,7 +94,7 @@ func TestTrigger_DispatchesWorkflowAndReturnsRunID(t *testing.T) { capturedMethod = req.Method capturedPath = req.URL.String() capturedBody, _ = io.ReadAll(req.Body) - _ = json.NewEncoder(w).Encode(dispatchWorkflowResponse{ + _ = json.NewEncoder(w).Encode(platformgithubactions.DispatchWorkflowResponse{ WorkflowRunID: 42, RunURL: "https://api.github.com/repos/uber/submitqueue/actions/runs/42", HTMLURL: "https://github.com/uber/submitqueue/actions/runs/42", @@ -117,7 +110,7 @@ func TestTrigger_DispatchesWorkflowAndReturnsRunID(t *testing.T) { assert.Equal(t, http.MethodPost, capturedMethod) assert.Equal(t, "/repos/uber/submitqueue/actions/workflows/submitqueue-ci.yml/dispatches", capturedPath) - var req dispatchWorkflowRequest + var req platformgithubactions.DispatchWorkflowRequest require.NoError(t, json.Unmarshal(capturedBody, &req)) assert.Equal(t, "main", req.Ref) assert.True(t, req.ReturnRunDetails) @@ -137,13 +130,13 @@ func TestTrigger_EmptyBaseProducesJSONArray(t *testing.T) { resolver := changesetfake.New().Set("head-batch", change.Change{URIs: []string{"u"}}) r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { capturedBody, _ = io.ReadAll(req.Body) - _ = json.NewEncoder(w).Encode(dispatchWorkflowResponse{WorkflowRunID: 7}) + _ = json.NewEncoder(w).Encode(platformgithubactions.DispatchWorkflowResponse{WorkflowRunID: 7}) }), resolver) _, err := r.Trigger(context.Background(), nil, entity.Batch{ID: "head-batch"}, nil) require.NoError(t, err) - var req dispatchWorkflowRequest + var req platformgithubactions.DispatchWorkflowRequest require.NoError(t, json.Unmarshal(capturedBody, &req)) assert.Equal(t, "[]", req.Inputs[InputKeyBaseURIs]) assert.NotContains(t, req.Inputs, InputKeyMetadata) @@ -151,7 +144,7 @@ func TestTrigger_EmptyBaseProducesJSONArray(t *testing.T) { func TestTrigger_ErrorsWhenDispatchResponseHasNoRunID(t *testing.T) { r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { - _ = json.NewEncoder(w).Encode(dispatchWorkflowResponse{}) + _ = json.NewEncoder(w).Encode(platformgithubactions.DispatchWorkflowResponse{}) })) _, err := r.Trigger(context.Background(), nil, entity.Batch{ID: "head-batch"}, nil) @@ -163,7 +156,7 @@ func TestStatus_GetsRunAndMapsState(t *testing.T) { r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { assert.Equal(t, http.MethodGet, req.Method) assert.Equal(t, "/repos/uber/submitqueue/actions/runs/42", req.URL.Path) - _ = json.NewEncoder(w).Encode(workflowRun{ + _ = json.NewEncoder(w).Encode(platformgithubactions.WorkflowRun{ ID: 42, DisplayTitle: "SubmitQueue gha-trace", Status: "in_progress", From 85fc006e13bcc25788640b03609d40373e49eb91 Mon Sep 17 00:00:00 2001 From: "chenghan.ying" Date: Mon, 20 Jul 2026 18:05:11 +0000 Subject: [PATCH 6/7] feat(stovepipe): add githubactions-backed BuildRunner adapter --- .../buildrunner/githubactions/BUILD.bazel | 29 ++ .../githubactions/githubactions.go | 219 +++++++++++++++ .../githubactions/githubactions_test.go | 254 ++++++++++++++++++ 3 files changed, 502 insertions(+) create mode 100644 stovepipe/extension/buildrunner/githubactions/BUILD.bazel create mode 100644 stovepipe/extension/buildrunner/githubactions/githubactions.go create mode 100644 stovepipe/extension/buildrunner/githubactions/githubactions_test.go diff --git a/stovepipe/extension/buildrunner/githubactions/BUILD.bazel b/stovepipe/extension/buildrunner/githubactions/BUILD.bazel new file mode 100644 index 00000000..77f60560 --- /dev/null +++ b/stovepipe/extension/buildrunner/githubactions/BUILD.bazel @@ -0,0 +1,29 @@ +load("@rules_go//go:def.bzl", "go_library", "go_test") + +go_library( + name = "go_default_library", + srcs = ["githubactions.go"], + importpath = "github.com/uber/submitqueue/stovepipe/extension/buildrunner/githubactions", + visibility = ["//visibility:public"], + deps = [ + "//platform/extension/buildrunner/githubactions:go_default_library", + "//stovepipe/entity:go_default_library", + "//stovepipe/extension/buildrunner:go_default_library", + "@org_uber_go_zap//:go_default_library", + ], +) + +go_test( + name = "go_default_test", + srcs = ["githubactions_test.go"], + embed = [":go_default_library"], + deps = [ + "//platform/extension/buildrunner/githubactions:go_default_library", + "//platform/http:go_default_library", + "//stovepipe/entity:go_default_library", + "//stovepipe/extension/buildrunner:go_default_library", + "@com_github_stretchr_testify//assert:go_default_library", + "@com_github_stretchr_testify//require:go_default_library", + "@org_uber_go_zap//:go_default_library", + ], +) diff --git a/stovepipe/extension/buildrunner/githubactions/githubactions.go b/stovepipe/extension/buildrunner/githubactions/githubactions.go new file mode 100644 index 00000000..548cc8ea --- /dev/null +++ b/stovepipe/extension/buildrunner/githubactions/githubactions.go @@ -0,0 +1,219 @@ +// Copyright (c) 2025 Uber Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +// Package githubactions implements buildrunner.BuildRunner backed by GitHub +// Actions workflow_dispatch. +// +// Trigger dispatches the configured workflow on a trusted ref and returns the +// GitHub workflow run ID as the build ID. Status and Cancel use that ID to +// call GitHub's workflow-run endpoints — no local state is required. +// +// The workflow receives the head and base URIs as workflow inputs +// (stovepipe_head_uri, stovepipe_base_uri, stovepipe_queue). The workflow +// definition checks out stovepipe_head_uri and, when stovepipe_base_uri is +// non-empty, diffs against it for an incremental build. +// +// Caller-supplied BuildMetadata is forwarded to the dispatch as +// stovepipe_metadata (JSON-encoded). Unlike Buildkite, GitHub does not echo +// dispatch inputs back on the run object, so Status returns the run's own +// identity/outcome metadata instead of the caller-supplied metadata. +package githubactions + +import ( + "context" + "encoding/json" + "fmt" + + "go.uber.org/zap" + + platformgithubactions "github.com/uber/submitqueue/platform/extension/buildrunner/githubactions" + "github.com/uber/submitqueue/stovepipe/entity" + "github.com/uber/submitqueue/stovepipe/extension/buildrunner" +) + +// Workflow input keys set on every dispatched run. +const ( + // InputKeyHeadURI carries the URI of the commit under validation. + InputKeyHeadURI = "stovepipe_head_uri" + + // InputKeyBaseURI carries the incremental baseline URI, empty for a full + // build. + InputKeyBaseURI = "stovepipe_base_uri" + + // InputKeyQueue carries the queue name so the workflow can select + // queue-specific test targets. + InputKeyQueue = "stovepipe_queue" + + // InputKeyMetadata carries the JSON-encoded BuildMetadata provided by the + // caller to Trigger. + InputKeyMetadata = "stovepipe_metadata" + + defaultRef = "main" +) + +// Params holds the dependencies for a GitHub Actions BuildRunner. +type Params struct { + // Config holds the per-queue identity for this BuildRunner. + Config buildrunner.Config + // Client is a pre-constructed GitHub Actions client, bound to one + // repository and workflow. The wiring layer builds it once via + // platformgithubactions.NewClient, with the GitHub API root (via + // platform/http.BaseURLTransport) and auth already configured. + Client *platformgithubactions.Client + // Logger is the structured logger. + Logger *zap.SugaredLogger + // Ref is the branch, tag, or SHA where the trusted workflow is read from. + // Defaults to "main". + Ref string + // ExtraInputs are copied into every workflow_dispatch request. Reserved + // stovepipe_* inputs managed by this package override conflicting keys. + ExtraInputs map[string]string +} + +// runner implements buildrunner.BuildRunner. +type runner struct { + cfg buildrunner.Config + ref string + extraInputs map[string]string + client *platformgithubactions.Client + logger *zap.SugaredLogger +} + +var _ buildrunner.BuildRunner = (*runner)(nil) + +// NewBuildRunner constructs a GitHub Actions-backed BuildRunner bound to one +// repository/workflow and one queue config. +func NewBuildRunner(params Params) (buildrunner.BuildRunner, error) { + if params.Client == nil { + return nil, fmt.Errorf("github actions client is required") + } + if params.Logger == nil { + return nil, fmt.Errorf("logger is required") + } + ref := params.Ref + if ref == "" { + ref = defaultRef + } + return newRunner(params.Config, ref, params.ExtraInputs, params.Client, params.Logger.Named("githubactions_buildrunner")), nil +} + +// newRunner constructs a runner. Used by NewBuildRunner and by tests. +func newRunner(cfg buildrunner.Config, ref string, extraInputs map[string]string, c *platformgithubactions.Client, logger *zap.SugaredLogger) *runner { + copied := make(map[string]string, len(extraInputs)) + for k, v := range extraInputs { + copied[k] = v + } + return &runner{ + cfg: cfg, + ref: ref, + extraInputs: copied, + client: c, + logger: logger, + } +} + +// Trigger dispatches the configured GitHub Actions workflow and returns the +// GitHub workflow run ID as the build ID. Errors are propagated to the caller +// so the queue consumer can nack and retry. +func (r *runner) Trigger(ctx context.Context, headURI, baseURI string, metadata entity.BuildMetadata) (entity.BuildID, error) { + inputs := make(map[string]string, len(r.extraInputs)+4) + for k, v := range r.extraInputs { + inputs[k] = v + } + inputs[InputKeyHeadURI] = headURI + inputs[InputKeyBaseURI] = baseURI + inputs[InputKeyQueue] = r.cfg.QueueName + if len(metadata) > 0 { + metaJSON, err := json.Marshal(metadata) + if err != nil { + return entity.BuildID{}, fmt.Errorf("github actions: marshal metadata: %w", err) + } + inputs[InputKeyMetadata] = string(metaJSON) + } + + resp, err := r.client.DispatchWorkflow(ctx, platformgithubactions.DispatchWorkflowRequest{ + Ref: r.ref, + ReturnRunDetails: true, + Inputs: inputs, + }) + if err != nil { + return entity.BuildID{}, fmt.Errorf("github actions: dispatch workflow: %w", err) + } + if resp.WorkflowRunID <= 0 { + return entity.BuildID{}, fmt.Errorf("github actions: dispatch workflow: response missing workflow_run_id (requires return_run_details support)") + } + + r.logger.Debugw("dispatched GitHub Actions workflow", + "owner", r.client.Owner(), + "repo", r.client.Repo(), + "workflow_id", r.client.WorkflowID(), + "ref", r.ref, + "github_run_id", resp.WorkflowRunID, + ) + return entity.BuildID{ID: platformgithubactions.EncodeRunID(resp.WorkflowRunID)}, nil +} + +// Status fetches the workflow run by GitHub run ID and maps the run's +// status/conclusion onto BuildStatus, returning the run's own identity and +// outcome fields as BuildMetadata. +func (r *runner) Status(ctx context.Context, buildID entity.BuildID) (entity.BuildStatus, entity.BuildMetadata, error) { + runID, err := platformgithubactions.ParseRunID(buildID.ID) + if err != nil { + return entity.BuildStatusUnknown, nil, fmt.Errorf("github actions: malformed build ID: %w", err) + } + + run, err := r.client.GetRun(ctx, runID) + if err != nil { + return entity.BuildStatusUnknown, nil, fmt.Errorf("github actions: get run: %w", err) + } + return mapRunStatus(run.Status, run.Conclusion), entity.BuildMetadata(platformgithubactions.RunMetadata(run)), nil +} + +// Cancel requests cancellation of the workflow run identified by buildID. +func (r *runner) Cancel(ctx context.Context, buildID entity.BuildID) error { + runID, err := platformgithubactions.ParseRunID(buildID.ID) + if err != nil { + return fmt.Errorf("github actions: malformed build ID: %w", err) + } + if err := r.client.CancelRun(ctx, runID); err != nil { + return fmt.Errorf("github actions: cancel run: %w", err) + } + r.logger.Debugw("cancelled GitHub Actions workflow run", + "github_run_id", runID, + ) + return nil +} + +// mapRunStatus maps a GitHub Actions run status/conclusion pair to a +// BuildStatus. +func mapRunStatus(status, conclusion string) entity.BuildStatus { + switch platformgithubactions.ParseRunStatus(status, conclusion) { + case platformgithubactions.RunStatusAccepted: + return entity.BuildStatusAccepted + case platformgithubactions.RunStatusRunning: + return entity.BuildStatusRunning + case platformgithubactions.RunStatusSucceeded: + return entity.BuildStatusSucceeded + case platformgithubactions.RunStatusFailed: + return entity.BuildStatusFailed + case platformgithubactions.RunStatusCancelled: + return entity.BuildStatusCancelled + default: + // Unrecognised GitHub status/conclusion. Do NOT assume terminal: + // Unknown is non-terminal, so the buildsignal poll loop keeps waiting + // rather than failing the build on a state this code does not + // understand. + return entity.BuildStatusUnknown + } +} diff --git a/stovepipe/extension/buildrunner/githubactions/githubactions_test.go b/stovepipe/extension/buildrunner/githubactions/githubactions_test.go new file mode 100644 index 00000000..47030f08 --- /dev/null +++ b/stovepipe/extension/buildrunner/githubactions/githubactions_test.go @@ -0,0 +1,254 @@ +// Copyright (c) 2025 Uber Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package githubactions + +import ( + "context" + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "strconv" + "strings" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.uber.org/zap" + + platformgithubactions "github.com/uber/submitqueue/platform/extension/buildrunner/githubactions" + phttp "github.com/uber/submitqueue/platform/http" + "github.com/uber/submitqueue/stovepipe/entity" + "github.com/uber/submitqueue/stovepipe/extension/buildrunner" +) + +// newTestRunner creates a runner backed by a test HTTP server. +func newTestRunner(t *testing.T, handler http.Handler) *runner { + t.Helper() + srv := httptest.NewServer(handler) + t.Cleanup(srv.Close) + c, err := phttp.NewClient(srv.URL) + require.NoError(t, err) + return newRunner( + buildrunner.Config{QueueName: "my-queue"}, + "main", + map[string]string{"custom": "value"}, + platformgithubactions.NewClient(c, "uber", "stovepipe", "stovepipe-ci.yml"), + zap.NewNop().Sugar(), + ) +} + +// --- Interface / constructor --- + +func TestNew_ImplementsInterface(t *testing.T) { + r, err := NewBuildRunner(Params{Logger: zap.NewNop().Sugar()}) + require.Error(t, err, "github actions client is required") + var _ buildrunner.BuildRunner = r +} + +// --- Trigger --- + +func TestTrigger_DispatchesWorkflowAndReturnsRunID(t *testing.T) { + var capturedMethod, capturedPath string + var capturedBody []byte + + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + capturedMethod = req.Method + capturedPath = req.URL.Path + capturedBody, _ = io.ReadAll(req.Body) + _ = json.NewEncoder(w).Encode(platformgithubactions.DispatchWorkflowResponse{WorkflowRunID: 42}) + })) + + id, err := r.Trigger(context.Background(), "github://repo/head/aaa", "github://repo/base/bbb", nil) + require.NoError(t, err) + assert.Equal(t, platformgithubactions.EncodeRunID(42), id.ID) + + assert.Equal(t, http.MethodPost, capturedMethod) + assert.Equal(t, "/repos/uber/stovepipe/actions/workflows/stovepipe-ci.yml/dispatches", capturedPath) + + var req platformgithubactions.DispatchWorkflowRequest + require.NoError(t, json.Unmarshal(capturedBody, &req)) + assert.Equal(t, "main", req.Ref) + assert.True(t, req.ReturnRunDetails) + assert.Equal(t, "github://repo/head/aaa", req.Inputs[InputKeyHeadURI]) + assert.Equal(t, "github://repo/base/bbb", req.Inputs[InputKeyBaseURI]) + assert.Equal(t, "my-queue", req.Inputs[InputKeyQueue]) + assert.Equal(t, "value", req.Inputs["custom"]) +} + +func TestTrigger_EmptyBaseURI_FullBuild(t *testing.T) { + var capturedBody []byte + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + capturedBody, _ = io.ReadAll(req.Body) + _ = json.NewEncoder(w).Encode(platformgithubactions.DispatchWorkflowResponse{WorkflowRunID: 1}) + })) + + _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + require.NoError(t, err) + + var req platformgithubactions.DispatchWorkflowRequest + require.NoError(t, json.Unmarshal(capturedBody, &req)) + assert.Equal(t, "", req.Inputs[InputKeyBaseURI]) +} + +func TestTrigger_DispatchError_ReturnsError(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + + _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + require.Error(t, err) +} + +func TestTrigger_ErrorsWhenDispatchResponseHasNoRunID(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + _ = json.NewEncoder(w).Encode(platformgithubactions.DispatchWorkflowResponse{}) + })) + + _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + require.Error(t, err) + assert.Contains(t, err.Error(), "response missing workflow_run_id") +} + +func TestTrigger_WithMetadata_SetsInput(t *testing.T) { + var capturedBody []byte + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + capturedBody, _ = io.ReadAll(req.Body) + _ = json.NewEncoder(w).Encode(platformgithubactions.DispatchWorkflowResponse{WorkflowRunID: 10}) + })) + + metadata := entity.BuildMetadata{"requester": "alice", "ticket": "SQ-42"} + _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", metadata) + require.NoError(t, err) + + var req platformgithubactions.DispatchWorkflowRequest + require.NoError(t, json.Unmarshal(capturedBody, &req)) + require.Contains(t, req.Inputs, InputKeyMetadata) + + var got entity.BuildMetadata + require.NoError(t, json.Unmarshal([]byte(req.Inputs[InputKeyMetadata]), &got)) + assert.Equal(t, metadata, got) +} + +func TestTrigger_NilMetadata_NoMetadataInput(t *testing.T) { + var capturedBody []byte + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + capturedBody, _ = io.ReadAll(req.Body) + _ = json.NewEncoder(w).Encode(platformgithubactions.DispatchWorkflowResponse{WorkflowRunID: 11}) + })) + + _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + require.NoError(t, err) + + var req platformgithubactions.DispatchWorkflowRequest + require.NoError(t, json.Unmarshal(capturedBody, &req)) + assert.NotContains(t, req.Inputs, InputKeyMetadata) +} + +// --- Status --- + +func TestStatus_StatusMapping(t *testing.T) { + tests := []struct { + name string + status string + conclusion string + want entity.BuildStatus + }{ + {name: "queued", status: "queued", want: entity.BuildStatusAccepted}, + {name: "in_progress", status: "in_progress", want: entity.BuildStatusRunning}, + {name: "success", status: "completed", conclusion: "success", want: entity.BuildStatusSucceeded}, + {name: "failure", status: "completed", conclusion: "failure", want: entity.BuildStatusFailed}, + {name: "cancelled", status: "completed", conclusion: "cancelled", want: entity.BuildStatusCancelled}, + // Unrecognised status/conclusion maps to the non-terminal Unknown, not + // Failed, so a state this code doesn't know about doesn't terminally + // fail a build. + {name: "some_future_state", status: "some_future_state", want: entity.BuildStatusUnknown}, + {name: "empty", want: entity.BuildStatusUnknown}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, mapRunStatus(tt.status, tt.conclusion)) + }) + } +} + +func TestStatus_ReturnsRunMetadata(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + assert.Equal(t, http.MethodGet, req.Method) + assert.Equal(t, "/repos/uber/stovepipe/actions/runs/7", req.URL.Path) + _ = json.NewEncoder(w).Encode(platformgithubactions.WorkflowRun{ + ID: 7, + Status: "in_progress", + HTMLURL: "https://github.com/uber/stovepipe/actions/runs/7", + }) + })) + + status, meta, err := r.Status(context.Background(), entity.BuildID{ID: platformgithubactions.EncodeRunID(7)}) + require.NoError(t, err) + assert.Equal(t, entity.BuildStatusRunning, status) + assert.Equal(t, "https://github.com/uber/stovepipe/actions/runs/7", meta["url"]) + assert.Equal(t, "7", meta["github_run_id"]) +} + +func TestStatus_MalformedBuildID(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + + _, _, err := r.Status(context.Background(), entity.BuildID{ID: "not-a-number"}) + require.Error(t, err) +} + +func TestStatus_NotFound(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusNotFound) + })) + + _, _, err := r.Status(context.Background(), entity.BuildID{ID: platformgithubactions.EncodeRunID(99)}) + require.Error(t, err) +} + +// --- Cancel --- + +func TestCancel_CallsGitHub(t *testing.T) { + var cancelledRunID int64 + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + require.True(t, req.Method == http.MethodPost && strings.HasSuffix(req.URL.Path, "/cancel")) + parts := strings.Split(req.URL.Path, "/") + id, err := strconv.ParseInt(parts[len(parts)-2], 10, 64) + require.NoError(t, err) + cancelledRunID = id + w.WriteHeader(http.StatusAccepted) + })) + + require.NoError(t, r.Cancel(context.Background(), entity.BuildID{ID: platformgithubactions.EncodeRunID(5)})) + assert.Equal(t, int64(5), cancelledRunID) +} + +func TestCancel_AlreadyTerminal_Noop(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusUnprocessableEntity) + })) + + require.NoError(t, r.Cancel(context.Background(), entity.BuildID{ID: platformgithubactions.EncodeRunID(5)})) +} + +func TestCancel_MalformedBuildID(t *testing.T) { + r := newTestRunner(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + + require.Error(t, r.Cancel(context.Background(), entity.BuildID{ID: "not-a-number"})) +} From 3c10122b7f2ea497fea2ff0393fc10a017e2767a Mon Sep 17 00:00:00 2001 From: "chenghan.ying" Date: Mon, 20 Jul 2026 22:59:24 +0000 Subject: [PATCH 7/7] fix(stovepipe): match reversed BuildRunner.Trigger arg order in backends --- .../extension/buildrunner/buildkite/buildkite.go | 2 +- .../buildrunner/buildkite/buildkite_test.go | 10 +++++----- .../buildrunner/githubactions/githubactions.go | 2 +- .../buildrunner/githubactions/githubactions_test.go | 12 ++++++------ 4 files changed, 13 insertions(+), 13 deletions(-) diff --git a/stovepipe/extension/buildrunner/buildkite/buildkite.go b/stovepipe/extension/buildrunner/buildkite/buildkite.go index db890e08..c0572da1 100644 --- a/stovepipe/extension/buildrunner/buildkite/buildkite.go +++ b/stovepipe/extension/buildrunner/buildkite/buildkite.go @@ -106,7 +106,7 @@ func newRunner(cfg buildrunner.Config, c *platformbuildkite.Client, logger *zap. // Trigger calls the Buildkite API to create the build and returns the // Buildkite build number as the build ID. Errors are propagated to the // caller so the queue consumer can nack and retry. -func (r *runner) Trigger(ctx context.Context, headURI, baseURI string, metadata entity.BuildMetadata) (entity.BuildID, error) { +func (r *runner) Trigger(ctx context.Context, baseURI, headURI string, metadata entity.BuildMetadata) (entity.BuildID, error) { env := map[string]string{ EnvKeyHeadURI: headURI, EnvKeyBaseURI: baseURI, diff --git a/stovepipe/extension/buildrunner/buildkite/buildkite_test.go b/stovepipe/extension/buildrunner/buildkite/buildkite_test.go index 06093714..ff127411 100644 --- a/stovepipe/extension/buildrunner/buildkite/buildkite_test.go +++ b/stovepipe/extension/buildrunner/buildkite/buildkite_test.go @@ -79,7 +79,7 @@ func TestTrigger_SubmitsCorrectPayloadAndReturnsBuildkiteNumber(t *testing.T) { _, _ = w.Write(buildJSON(42, "scheduled", "https://buildkite.com/test-org/my-pipeline/builds/42")) })) - id, err := r.Trigger(context.Background(), "github://repo/head/aaa", "github://repo/base/bbb", nil) + id, err := r.Trigger(context.Background(), "github://repo/base/bbb", "github://repo/head/aaa", nil) require.NoError(t, err) assert.Equal(t, platformbuildkite.EncodeBuildNumber(42), id.ID) @@ -100,7 +100,7 @@ func TestTrigger_EmptyBaseURI_FullBuild(t *testing.T) { _, _ = w.Write(buildJSON(1, "scheduled", "")) })) - _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + _, err := r.Trigger(context.Background(), "", "github://repo/head/aaa", nil) require.NoError(t, err) var req platformbuildkite.CreateBuildRequest @@ -113,7 +113,7 @@ func TestTrigger_BuildkiteError_ReturnsError(t *testing.T) { w.WriteHeader(http.StatusInternalServerError) })) - _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + _, err := r.Trigger(context.Background(), "", "github://repo/head/aaa", nil) require.Error(t, err) } @@ -126,7 +126,7 @@ func TestTrigger_WithMetadata_SetsEnvVar(t *testing.T) { })) metadata := entity.BuildMetadata{"requester": "alice", "ticket": "SQ-42"} - _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", metadata) + _, err := r.Trigger(context.Background(), "", "github://repo/head/aaa", metadata) require.NoError(t, err) var req platformbuildkite.CreateBuildRequest @@ -146,7 +146,7 @@ func TestTrigger_NilMetadata_NoMetadataEnvVar(t *testing.T) { _, _ = w.Write(buildJSON(11, "scheduled", "")) })) - _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + _, err := r.Trigger(context.Background(), "", "github://repo/head/aaa", nil) require.NoError(t, err) var req platformbuildkite.CreateBuildRequest diff --git a/stovepipe/extension/buildrunner/githubactions/githubactions.go b/stovepipe/extension/buildrunner/githubactions/githubactions.go index 548cc8ea..ea4d7401 100644 --- a/stovepipe/extension/buildrunner/githubactions/githubactions.go +++ b/stovepipe/extension/buildrunner/githubactions/githubactions.go @@ -126,7 +126,7 @@ func newRunner(cfg buildrunner.Config, ref string, extraInputs map[string]string // Trigger dispatches the configured GitHub Actions workflow and returns the // GitHub workflow run ID as the build ID. Errors are propagated to the caller // so the queue consumer can nack and retry. -func (r *runner) Trigger(ctx context.Context, headURI, baseURI string, metadata entity.BuildMetadata) (entity.BuildID, error) { +func (r *runner) Trigger(ctx context.Context, baseURI, headURI string, metadata entity.BuildMetadata) (entity.BuildID, error) { inputs := make(map[string]string, len(r.extraInputs)+4) for k, v := range r.extraInputs { inputs[k] = v diff --git a/stovepipe/extension/buildrunner/githubactions/githubactions_test.go b/stovepipe/extension/buildrunner/githubactions/githubactions_test.go index 47030f08..bd20e7ac 100644 --- a/stovepipe/extension/buildrunner/githubactions/githubactions_test.go +++ b/stovepipe/extension/buildrunner/githubactions/githubactions_test.go @@ -71,7 +71,7 @@ func TestTrigger_DispatchesWorkflowAndReturnsRunID(t *testing.T) { _ = json.NewEncoder(w).Encode(platformgithubactions.DispatchWorkflowResponse{WorkflowRunID: 42}) })) - id, err := r.Trigger(context.Background(), "github://repo/head/aaa", "github://repo/base/bbb", nil) + id, err := r.Trigger(context.Background(), "github://repo/base/bbb", "github://repo/head/aaa", nil) require.NoError(t, err) assert.Equal(t, platformgithubactions.EncodeRunID(42), id.ID) @@ -95,7 +95,7 @@ func TestTrigger_EmptyBaseURI_FullBuild(t *testing.T) { _ = json.NewEncoder(w).Encode(platformgithubactions.DispatchWorkflowResponse{WorkflowRunID: 1}) })) - _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + _, err := r.Trigger(context.Background(), "", "github://repo/head/aaa", nil) require.NoError(t, err) var req platformgithubactions.DispatchWorkflowRequest @@ -108,7 +108,7 @@ func TestTrigger_DispatchError_ReturnsError(t *testing.T) { w.WriteHeader(http.StatusInternalServerError) })) - _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + _, err := r.Trigger(context.Background(), "", "github://repo/head/aaa", nil) require.Error(t, err) } @@ -117,7 +117,7 @@ func TestTrigger_ErrorsWhenDispatchResponseHasNoRunID(t *testing.T) { _ = json.NewEncoder(w).Encode(platformgithubactions.DispatchWorkflowResponse{}) })) - _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + _, err := r.Trigger(context.Background(), "", "github://repo/head/aaa", nil) require.Error(t, err) assert.Contains(t, err.Error(), "response missing workflow_run_id") } @@ -130,7 +130,7 @@ func TestTrigger_WithMetadata_SetsInput(t *testing.T) { })) metadata := entity.BuildMetadata{"requester": "alice", "ticket": "SQ-42"} - _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", metadata) + _, err := r.Trigger(context.Background(), "", "github://repo/head/aaa", metadata) require.NoError(t, err) var req platformgithubactions.DispatchWorkflowRequest @@ -149,7 +149,7 @@ func TestTrigger_NilMetadata_NoMetadataInput(t *testing.T) { _ = json.NewEncoder(w).Encode(platformgithubactions.DispatchWorkflowResponse{WorkflowRunID: 11}) })) - _, err := r.Trigger(context.Background(), "github://repo/head/aaa", "", nil) + _, err := r.Trigger(context.Background(), "", "github://repo/head/aaa", nil) require.NoError(t, err) var req platformgithubactions.DispatchWorkflowRequest