From e59607d46b6006ff1065a847ba283e97d87c76ce Mon Sep 17 00:00:00 2001 From: Mitchell Hashimoto Date: Fri, 1 May 2026 14:57:23 -0700 Subject: [PATCH] buildkite --- AGENTS.md | 6 +- README.md | 48 ++- http.go | 85 ++++- internal/buildkite/buildkite.go | 346 +++++++++++++++++ internal/buildkite/buildkite_test.go | 207 ++++++++++ knot.go | 2 +- main.go | 97 ++++- provider.go | 6 +- provider_buildkite.go | 552 +++++++++++++++++++++++++++ provider_buildkite_test.go | 405 ++++++++++++++++++++ provider_fake.go | 1 + store.go | 103 +++++ store_migrate.go | 26 ++ 13 files changed, 1845 insertions(+), 39 deletions(-) create mode 100644 internal/buildkite/buildkite.go create mode 100644 internal/buildkite/buildkite_test.go create mode 100644 provider_buildkite.go create mode 100644 provider_buildkite_test.go diff --git a/AGENTS.md b/AGENTS.md index e3a3650..cef8a3c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -9,7 +9,7 @@ - Comment heavily, but not redundantly. Explain the "why" behind decisions clearly, don't repeat the "what." - Try to limit line length below 100 characters, but don't be afraid to break - this rule if it improves readability. + this rule if it improves readability. - Do not make lines too short to follow this rule either. Try to use as much of the max characters as possible without sacrificing readability. @@ -18,3 +18,7 @@ - Add compile-time interface checks whenever a type implements an interface we care about. For example: `var _ io.Reader = (*MyType)(nil)`. + +## Providers + +- Buildkite API: diff --git a/README.md b/README.md index 45cfa92..63eb5d6 100644 --- a/README.md +++ b/README.md @@ -35,19 +35,35 @@ status, counts, etc. can show up inline in things like pull requests. go run . -addr :8080 ``` -## Endpoints (planned) - -- `GET /events` — WebSocket stream of pipeline status events, - consumed by the Tangled appview. -- `POST /webhooks/buildkite` — Buildkite webhook receiver. -- `POST /xrpc/sh.tangled.pipeline.cancelPipeline` — cancel a running build. - -## Configuration (planned) - -| Env var | Description | -| ---------------------- | ------------------------------------ | -| `TACK_BUILDKITE_TOKEN` | Buildkite API token | -| `TACK_BUILDKITE_ORG` | Buildkite organization slug | -| `TACK_JETSTREAM_URL` | Tangled Jetstream WebSocket URL | -| `TACK_DB_PATH` | Local SQLite path for the event log | -| `TACK_OWNER_DID` | DID of the spindle operator | +## Configuration + +### Required + +| Env var | Description | +| ---------------- | ----------------------------------------------------------- | +| `TACK_HOSTNAME` | This spindle's hostname (matches `sh.tangled.repo.spindle`) | +| `TACK_OWNER_DID` | DID of the spindle operator | + +### Optional + +| Env var | Description | +| -------------------- | -------------------------------------------------------- | +| `TACK_LISTEN_ADDR` | HTTP listen address (default `:8080`) | +| `TACK_DB_PATH` | Local SQLite path (default `tack.db`) | +| `TACK_JETSTREAM_URL` | Tangled Jetstream WebSocket URL | +| `TACK_DEV` | Use `ws://` for knot event-streams (any non-empty value) | + +### Buildkite + +Setting `TACK_BUILDKITE_TOKEN` enables Buildkite mode; when unset, tack +runs the in-process fake provider for local development. When +Buildkite mode is enabled, every other variable in this section is +required. + +| Env var | Description | +| ------------------------------- | ------------------------------------------------------------------------------ | +| `TACK_BUILDKITE_TOKEN` | Buildkite API token (enables Buildkite mode) | +| `TACK_BUILDKITE_ORG` | Buildkite organization slug | +| `TACK_BUILDKITE_PIPELINE` | Buildkite pipeline slug to fire builds on | +| `TACK_BUILDKITE_WEBHOOK_SECRET` | Shared secret for `/webhooks/buildkite` auth | +| `TACK_BUILDKITE_WEBHOOK_MODE` | `token` (default) or `signature` — must match Buildkite's notification setting | diff --git a/http.go b/http.go index 6d679ba..0341b6b 100644 --- a/http.go +++ b/http.go @@ -22,6 +22,7 @@ import ( "encoding/json" "errors" "fmt" + "io" "log/slog" "net/http" "strconv" @@ -29,6 +30,8 @@ import ( "github.com/gorilla/websocket" "tangled.org/core/api/tangled" + + "github.com/mitchellh/tack/internal/buildkite" ) // runHTTP starts the spindle's HTTP server and blocks until ctx is @@ -37,8 +40,12 @@ import ( // // The logger is read from ctx via loggerFrom. The broker is the // in-process pub/sub used by /events to fan published records out to -// connected websocket subscribers. -func runHTTP(ctx context.Context, cfg config, br *broker, provider Provider) error { +// connected websocket subscribers. bkProvider may be nil — when a +// deployment runs the fake provider, /webhooks/buildkite still +// registers but responds 503, so a misdirected Buildkite webhook +// gets a clear "this spindle isn't accepting Buildkite events" rather +// than a misleading 200. +func runHTTP(ctx context.Context, cfg config, br *broker, provider Provider, bkProvider *buildkiteProvider) error { logger := loggerFrom(ctx) mux := http.NewServeMux() @@ -46,7 +53,7 @@ func runHTTP(ctx context.Context, cfg config, br *broker, provider Provider) err mux.HandleFunc("GET /events", eventsHandler(logger, br)) mux.HandleFunc("GET /logs/{knot}/{pipelineRkey}/{workflow}", logsHandler(logger, provider)) mux.HandleFunc("GET /xrpc/"+tangled.OwnerNSID, ownerHandler(logger, cfg.OwnerDID)) - mux.HandleFunc("POST /webhooks/buildkite", buildkiteWebhookHandler()) + mux.HandleFunc("POST /webhooks/buildkite", buildkiteWebhookHandler(logger, bkProvider)) srv := &http.Server{ Addr: cfg.Addr, @@ -96,11 +103,75 @@ func ownerHandler(logger *slog.Logger, owner string) http.HandlerFunc { } } -// buildkiteWebhookHandler is a placeholder until we implement Buildkite -> -// pipeline.status translation. -func buildkiteWebhookHandler() http.HandlerFunc { +// buildkiteWebhookHandler receives Buildkite Pipelines webhook events, +// authenticates the request against whichever scheme the provider was +// configured with, and hands the decoded payload to the provider for +// translation into a sh.tangled.pipeline.status publish. +// +// Authentication is intentionally fail-closed: when bk is nil (no +// Buildkite provider configured) we 503 instead of accepting events +// silently. The body is buffered up front because signature mode +// HMACs the raw bytes — we can't rely on the JSON decoder reading +// the request body before verification. +// +// Acknowledgement contract with Buildkite: we 200 on any well-formed +// event we accepted (including events we deliberately ignore, like +// job.* or builds we don't track), and 5xx only on internal failure +// the operator should look at. A 4xx/5xx makes Buildkite retry, +// which we don't want for "this isn't an event we care about". +func buildkiteWebhookHandler(logger *slog.Logger, bk *buildkiteProvider) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { - http.Error(w, "not implemented", http.StatusNotImplemented) + if bk == nil { + http.Error(w, "buildkite provider not configured", + http.StatusServiceUnavailable) + return + } + + // Cap body size so a malicious sender can't exhaust + // memory; Buildkite payloads in practice are well under + // 64 KiB but a generous-but-bounded ceiling is the + // right shape here. + body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20)) + if err != nil { + logger.Warn("buildkite webhook: read body", "err", err) + http.Error(w, "read body", http.StatusBadRequest) + return + } + + if err := bk.VerifyWebhook(r.Header, body); err != nil { + logger.Warn("buildkite webhook: verify failed", + "err", err, + "remote", r.RemoteAddr, + ) + http.Error(w, "unauthorized", http.StatusUnauthorized) + return + } + + var payload buildkite.WebhookPayload + if err := json.Unmarshal(body, &payload); err != nil { + logger.Warn("buildkite webhook: decode body", "err", err) + http.Error(w, "bad payload", http.StatusBadRequest) + return + } + // The X-Buildkite-Event header is authoritative for the + // event name; the body field is convenience but doesn't + // always match exactly. Prefer the header. + if h := r.Header.Get("X-Buildkite-Event"); h != "" { + payload.Event = h + } + + // Translate + publish on the request context so a slow + // store/broker doesn't outlive an aborted webhook + // connection. + if err := bk.HandleWebhook(r.Context(), payload); err != nil { + logger.Error("buildkite webhook: handle", "err", err, + "event", payload.Event, + "build_uuid", payload.Build.ID, + ) + http.Error(w, "internal error", http.StatusInternalServerError) + return + } + w.WriteHeader(http.StatusOK) } } diff --git a/internal/buildkite/buildkite.go b/internal/buildkite/buildkite.go new file mode 100644 index 0000000..406bf98 --- /dev/null +++ b/internal/buildkite/buildkite.go @@ -0,0 +1,346 @@ +// Package buildkite is a small Buildkite REST + webhook client tack +// uses to drive its Buildkite-backed Provider implementation. +// +// The package deliberately covers a tiny slice of the upstream API +// (create build, get build, fetch job log, decode + authenticate +// webhook payloads). It exists as its own package so the rest of +// tack — particularly the Provider implementation that translates +// Tangled triggers into Buildkite builds — can stay focused on +// translation rather than HTTP plumbing. +// +// Naming convention: types here are *not* prefixed with "Buildkite". +// Imported as `buildkite.Client`, `buildkite.Build`, etc., the package +// path supplies the disambiguation already. +package buildkite + +import ( + "bytes" + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "strconv" + "strings" + "time" +) + +// APIBase is the public Buildkite REST API root. Exported as a var +// (not a const) so tests can swap it for an httptest server URL +// without hooking up a real Buildkite account. +var APIBase = "https://api.buildkite.com" + +// ErrNotFound is returned by Get* methods when the upstream returns +// 404. Callers translate it to whatever shape they need (the Provider +// maps it onto its own ErrLogsNotFound for the /logs handler). +var ErrNotFound = errors.New("buildkite: not found") + +// Client is a thin wrapper around net/http carrying API credentials +// + organization context so call sites don't repeat them. Safe for +// concurrent use; the embedded http.Client is goroutine-safe. +type Client struct { + http *http.Client + token string + org string +} + +// NewClient builds a Client with sensible defaults. The 30s timeout +// covers individual requests, not the whole client lifetime — +// long-poll-style endpoints aren't used here, so a generous-but-bounded +// per-request timeout is the right default. +func NewClient(token, org string) *Client { + return &Client{ + http: &http.Client{Timeout: 30 * time.Second}, + token: token, + org: org, + } +} + +// Job is the subset of a Buildkite job object we care about: the ID +// and name we need to fetch logs for it, plus its state for +// surfacing in webhook handling. +// +// Buildkite jobs come in several "type" values (script, waiter, +// manual) — only "script" jobs have logs, but we decode the slice +// as-is and let the caller decide whether to skip non-script entries. +type Job struct { + ID string `json:"id"` + Type string `json:"type"` + Name string `json:"name"` + State string `json:"state"` +} + +// Build is the subset of a Buildkite build object we care about. +// Fields not present here are dropped silently by the JSON decoder — +// keep this list lean so additions to the upstream schema don't +// force us to touch this file. +type Build struct { + ID string `json:"id"` + Number int64 `json:"number"` + State string `json:"state"` + WebURL string `json:"web_url"` + Commit string `json:"commit"` + Branch string `json:"branch"` + Message string `json:"message"` + MetaData map[string]string `json:"meta_data"` + Jobs []Job `json:"jobs"` + Pipeline map[string]interface{} `json:"pipeline"` +} + +// CreateBuildRequest is the request body for POST /builds. Only the +// fields callers actively use are exposed — the upstream API accepts +// many more, but we'd just be passing through dead options. +// +// IgnorePipelineBranchFilters defaults to false on the wire (omitempty +// elides the zero value); callers that want pipeline-level branch +// filters bypassed should set it to true. +type CreateBuildRequest struct { + Commit string `json:"commit"` + Branch string `json:"branch"` + Message string `json:"message,omitempty"` + Env map[string]string `json:"env,omitempty"` + MetaData map[string]string `json:"meta_data,omitempty"` + IgnorePipelineBranchFilters bool `json:"ignore_pipeline_branch_filters,omitempty"` +} + +// CreateBuild fires a build on the named pipeline. Returns the +// decoded response so the caller can persist build_uuid + number for +// later webhook lookup. +// +// Buildkite returns 201 on success; anything else is wrapped into an +// error that includes the response body so a misconfigured pipeline +// (e.g. wrong slug, missing branch) surfaces useful diagnostics into +// the caller's log. +func (c *Client) CreateBuild( + ctx context.Context, + pipelineSlug string, + req CreateBuildRequest, +) (*Build, error) { + body, err := json.Marshal(req) + if err != nil { + return nil, fmt.Errorf("marshal build request: %w", err) + } + + url := fmt.Sprintf("%s/v2/organizations/%s/pipelines/%s/builds", + APIBase, c.org, pipelineSlug, + ) + httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(body)) + if err != nil { + return nil, fmt.Errorf("build request: %w", err) + } + httpReq.Header.Set("Authorization", "Bearer "+c.token) + httpReq.Header.Set("Content-Type", "application/json") + + resp, err := c.http.Do(httpReq) + if err != nil { + return nil, fmt.Errorf("create build: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusCreated { + raw, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + return nil, fmt.Errorf("create build: status %d: %s", + resp.StatusCode, strings.TrimSpace(string(raw)), + ) + } + + var out Build + if err := json.NewDecoder(resp.Body).Decode(&out); err != nil { + return nil, fmt.Errorf("decode build response: %w", err) + } + return &out, nil +} + +// GetBuild fetches the full build record by number, including the +// jobs slice. Used by callers that need the current set of jobs for +// a known (pipelineSlug, buildNumber) pair. +// +// Returns ErrNotFound when Buildkite responds 404. +func (c *Client) GetBuild( + ctx context.Context, + pipelineSlug string, + buildNumber int64, +) (*Build, error) { + url := fmt.Sprintf("%s/v2/organizations/%s/pipelines/%s/builds/%d", + APIBase, c.org, pipelineSlug, buildNumber, + ) + httpReq, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) + if err != nil { + return nil, fmt.Errorf("build request: %w", err) + } + httpReq.Header.Set("Authorization", "Bearer "+c.token) + + resp, err := c.http.Do(httpReq) + if err != nil { + return nil, fmt.Errorf("get build: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode == http.StatusNotFound { + return nil, ErrNotFound + } + if resp.StatusCode != http.StatusOK { + raw, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + return nil, fmt.Errorf("get build: status %d: %s", + resp.StatusCode, strings.TrimSpace(string(raw)), + ) + } + + var out Build + if err := json.NewDecoder(resp.Body).Decode(&out); err != nil { + return nil, fmt.Errorf("decode build: %w", err) + } + return &out, nil +} + +// GetJobLog fetches the plain-text log for a single job. Buildkite +// supports several formats; we ask for text/plain explicitly so the +// response is one big string the caller can split on newlines. +// +// Returns ErrNotFound when Buildkite responds 404 (typical for a +// job that hasn't started yet — it has no log to serve). +func (c *Client) GetJobLog( + ctx context.Context, + pipelineSlug string, + buildNumber int64, + jobID string, +) (string, error) { + url := fmt.Sprintf("%s/v2/organizations/%s/pipelines/%s/builds/%d/jobs/%s/log", + APIBase, c.org, pipelineSlug, buildNumber, jobID, + ) + httpReq, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) + if err != nil { + return "", fmt.Errorf("build request: %w", err) + } + httpReq.Header.Set("Authorization", "Bearer "+c.token) + httpReq.Header.Set("Accept", "text/plain") + + resp, err := c.http.Do(httpReq) + if err != nil { + return "", fmt.Errorf("get job log: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode == http.StatusNotFound { + return "", ErrNotFound + } + if resp.StatusCode != http.StatusOK { + raw, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + return "", fmt.Errorf("get job log: status %d: %s", + resp.StatusCode, strings.TrimSpace(string(raw)), + ) + } + + body, err := io.ReadAll(resp.Body) + if err != nil { + return "", fmt.Errorf("read job log: %w", err) + } + return string(body), nil +} + +// WebhookPayload is the small slice of the webhook body callers +// actually decode. Build events all wrap a "build" object; job +// events wrap a "job" — keeping both around lets callers add +// job-level mapping later without changing the decoder. +type WebhookPayload struct { + Event string `json:"event"` + Build Build `json:"build"` + Job Job `json:"job"` +} + +// WebhookMode selects how an inbound webhook request is +// authenticated. The two values correspond directly to the two +// settings Buildkite's notification service exposes for this: +// WebhookModeToken sends the secret in the X-Buildkite-Token header +// in plain text; WebhookModeSignature sends an HMAC-SHA256 of the +// body in X-Buildkite-Signature. +type WebhookMode string + +const ( + WebhookModeToken WebhookMode = "token" + WebhookModeSignature WebhookMode = "signature" +) + +// VerifySignature validates the X-Buildkite-Signature header against +// secret using the documented "." HMAC-SHA256 +// scheme. Returns nil when the header is well-formed and the digest +// matches; any other condition returns an error. +// +// We deliberately do NOT enforce a freshness window on timestamp: +// callers in practice consume idempotent or storage-deduplicated +// events, so a replayed event is at worst a duplicate publish. +// Callers that need stricter freshness should layer it on top. +// +// The header format is "timestamp=,signature=". +func VerifySignature(header, secret string, body []byte) error { + if header == "" { + return errors.New("missing X-Buildkite-Signature header") + } + if secret == "" { + // A misconfigured server is a programmer bug, but we'd + // rather fail closed than silently accept any signature. + return errors.New("server has no webhook secret configured") + } + + var ts, sig string + for _, part := range strings.Split(header, ",") { + k, v, ok := strings.Cut(strings.TrimSpace(part), "=") + if !ok { + continue + } + switch k { + case "timestamp": + ts = v + case "signature": + sig = v + } + } + if ts == "" || sig == "" { + return errors.New("malformed signature header") + } + // Sanity-check the timestamp is a parseable int. The value + // itself isn't validated against the clock (see comment above), + // but a non-numeric timestamp is structurally invalid. + if _, err := strconv.ParseInt(ts, 10, 64); err != nil { + return fmt.Errorf("invalid timestamp: %w", err) + } + + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write([]byte(ts)) + mac.Write([]byte(".")) + mac.Write(body) + expected := hex.EncodeToString(mac.Sum(nil)) + + // Compare in constant time to keep the verifier from leaking + // the expected digest through timing. + if !hmac.Equal([]byte(expected), []byte(sig)) { + return errors.New("signature mismatch") + } + return nil +} + +// VerifyToken handles the simpler X-Buildkite-Token mode: the +// configured secret is sent verbatim in the header. Constant-time +// comparison keeps a brute-forcing attacker from learning the token +// one byte at a time. +// +// Returns nil on match. Two error cases: +// - missing header on the request (caller should 401), +// - server has no expected token configured (caller should 500; +// fail-closed, never accept-anything). +func VerifyToken(header, expected string) error { + if header == "" { + return errors.New("missing X-Buildkite-Token header") + } + if expected == "" { + return errors.New("server has no webhook token configured") + } + if !hmac.Equal([]byte(header), []byte(expected)) { + return errors.New("token mismatch") + } + return nil +} diff --git a/internal/buildkite/buildkite_test.go b/internal/buildkite/buildkite_test.go new file mode 100644 index 0000000..5322e34 --- /dev/null +++ b/internal/buildkite/buildkite_test.go @@ -0,0 +1,207 @@ +package buildkite + +// Tests for the Buildkite REST client + webhook signature/token +// verifiers. Provider-level (Tangled translation) tests live with +// the provider in the main package. + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" +) + +// TestVerifySignature covers the HMAC mode end to end. The reference +// digest is computed with the documented "." +// preimage so the test pins the wire format, not just the helper. +func TestVerifySignature(t *testing.T) { + const secret = "shhh" + body := []byte(`{"event":"build.finished"}`) + const ts = "1700000000" + + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write([]byte(ts)) + mac.Write([]byte(".")) + mac.Write(body) + good := hex.EncodeToString(mac.Sum(nil)) + + cases := []struct { + name string + header string + secret string + body []byte + wantErr bool + }{ + {"valid", "timestamp=" + ts + ",signature=" + good, secret, body, false}, + {"valid with whitespace", " timestamp=" + ts + " , signature=" + good + " ", secret, body, false}, + {"empty header", "", secret, body, true}, + {"empty server secret", "timestamp=" + ts + ",signature=" + good, "", body, true}, + {"missing timestamp", "signature=" + good, secret, body, true}, + {"missing signature", "timestamp=" + ts, secret, body, true}, + {"non-numeric timestamp", "timestamp=abc,signature=" + good, secret, body, true}, + {"wrong signature", "timestamp=" + ts + ",signature=00", secret, body, true}, + {"wrong body", "timestamp=" + ts + ",signature=" + good, secret, []byte("nope"), true}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + err := VerifySignature(c.header, c.secret, c.body) + if (err != nil) != c.wantErr { + t.Fatalf("err=%v wantErr=%v", err, c.wantErr) + } + }) + } +} + +// TestVerifyToken pins the token-mode behaviour: it must be a +// constant-time exact-match check, never a prefix or substring. +func TestVerifyToken(t *testing.T) { + cases := []struct { + name string + header string + expect string + wantErr bool + }{ + {"match", "abc123", "abc123", false}, + {"empty header", "", "abc123", true}, + {"empty expected", "abc123", "", true}, + {"mismatch", "abc124", "abc123", true}, + {"prefix is not a match", "abc", "abc123", true}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + err := VerifyToken(c.header, c.expect) + if (err != nil) != c.wantErr { + t.Fatalf("err=%v wantErr=%v", err, c.wantErr) + } + }) + } +} + +// TestClientCreateBuild covers the request shape we send (auth +// header, JSON body, URL) and the response decoding for the happy +// path. +func TestClientCreateBuild(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/v2/organizations/myorg/pipelines/mypipe/builds" { + t.Errorf("bad path %s", r.URL.Path) + } + if got := r.Header.Get("Authorization"); got != "Bearer tok" { + t.Errorf("bad auth %q", got) + } + if got := r.Header.Get("Content-Type"); got != "application/json" { + t.Errorf("bad content-type %q", got) + } + var got CreateBuildRequest + if err := json.NewDecoder(r.Body).Decode(&got); err != nil { + t.Fatalf("decode: %v", err) + } + if got.Commit != "abc" || got.Branch != "main" { + t.Errorf("bad body: %+v", got) + } + if got.MetaData["k"] != "v" { + t.Errorf("missing meta_data: %+v", got.MetaData) + } + w.WriteHeader(http.StatusCreated) + _ = json.NewEncoder(w).Encode(Build{ + ID: "uuid-1", + Number: 42, + State: "scheduled", + }) + })) + defer srv.Close() + + prev := APIBase + APIBase = srv.URL + defer func() { APIBase = prev }() + + c := NewClient("tok", "myorg") + build, err := c.CreateBuild(context.Background(), "mypipe", CreateBuildRequest{ + Commit: "abc", + Branch: "main", + MetaData: map[string]string{"k": "v"}, + }) + if err != nil { + t.Fatalf("CreateBuild: %v", err) + } + if build.ID != "uuid-1" || build.Number != 42 || build.State != "scheduled" { + t.Fatalf("unexpected build: %+v", build) + } +} + +// TestClientCreateBuildError makes sure non-2xx responses surface +// the upstream error body — that text ends up in operator logs, so +// silently dropping it would make misconfigurations very painful to +// diagnose. +func TestClientCreateBuildError(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusUnprocessableEntity) + fmt.Fprint(w, `{"message":"branch is required"}`) + })) + defer srv.Close() + + prev := APIBase + APIBase = srv.URL + defer func() { APIBase = prev }() + + c := NewClient("tok", "myorg") + _, err := c.CreateBuild(context.Background(), "mypipe", CreateBuildRequest{}) + if err == nil { + t.Fatal("expected error") + } + if !strings.Contains(err.Error(), "branch is required") { + t.Fatalf("error missing upstream body: %v", err) + } +} + +// TestClientGetJobLog confirms we send the right Accept header (so +// Buildkite returns plain text, not JSON) and surface 404 as +// ErrNotFound for callers to translate. +func TestClientGetJobLog(t *testing.T) { + t.Run("ok", func(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if got := r.Header.Get("Accept"); got != "text/plain" { + t.Errorf("bad accept %q", got) + } + if !strings.HasSuffix(r.URL.Path, "/builds/7/jobs/job-1/log") { + t.Errorf("bad path %s", r.URL.Path) + } + io.WriteString(w, "line1\nline2\n") + })) + defer srv.Close() + + prev := APIBase + APIBase = srv.URL + defer func() { APIBase = prev }() + + body, err := NewClient("tok", "myorg"). + GetJobLog(context.Background(), "mypipe", 7, "job-1") + if err != nil { + t.Fatalf("GetJobLog: %v", err) + } + if body != "line1\nline2\n" { + t.Fatalf("body = %q", body) + } + }) + t.Run("404", func(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "no", http.StatusNotFound) + })) + defer srv.Close() + prev := APIBase + APIBase = srv.URL + defer func() { APIBase = prev }() + + _, err := NewClient("tok", "myorg"). + GetJobLog(context.Background(), "mypipe", 7, "job-1") + if err != ErrNotFound { + t.Fatalf("err = %v; want ErrNotFound", err) + } + }) +} diff --git a/knot.go b/knot.go index 77ece3c..708ff04 100644 --- a/knot.go +++ b/knot.go @@ -205,7 +205,7 @@ func (k *knotConsumer) process(ctx context.Context, src eventconsumer.Source, ms // Spawn is non-blocking — it fans out into provider-owned // goroutines so this worker can move on to the next event. // The provider keeps ctx around for shutdown coordination. - k.provider.Spawn(ctx, src.Key(), msg.Rkey, p.Workflows) + k.provider.Spawn(ctx, src.Key(), msg.Rkey, p.TriggerMetadata, p.Workflows) default: // Knots may publish other record types over the same stream; we diff --git a/main.go b/main.go index c8447f8..c9e88e3 100644 --- a/main.go +++ b/main.go @@ -8,12 +8,15 @@ import ( "context" "errors" "flag" + "fmt" "log/slog" "os" "os/signal" "syscall" charmlog "github.com/charmbracelet/log" + + "github.com/mitchellh/tack/internal/buildkite" ) // config is the runtime configuration, sourced from environment variables and @@ -28,16 +31,34 @@ type config struct { // Dev flips the knot event-stream scheme from wss:// to ws://. // Useful when running against a local knot during development. Dev bool + + // Buildkite-mode configuration. BuildkiteToken is the switch: + // when empty we fall back to the in-process fake provider + // (useful for local development against a real Tangled + // jetstream); when set, the other Buildkite fields are + // required and tack will refuse to start without them. + BuildkiteToken string + BuildkiteOrg string + BuildkitePipeline string + BuildkiteWebhookSecret string + BuildkiteWebhookMode buildkite.WebhookMode } func loadConfig() (config, error) { cfg := config{ - Addr: envOr("TACK_LISTEN_ADDR", ":8080"), - Hostname: os.Getenv("TACK_HOSTNAME"), - OwnerDID: os.Getenv("TACK_OWNER_DID"), - JetstreamURL: envOr("TACK_JETSTREAM_URL", "wss://jetstream1.us-west.bsky.network/subscribe"), - DBPath: envOr("TACK_DB_PATH", "tack.db"), - Dev: os.Getenv("TACK_DEV") != "", + Addr: envOr("TACK_LISTEN_ADDR", ":8080"), + Hostname: os.Getenv("TACK_HOSTNAME"), + OwnerDID: os.Getenv("TACK_OWNER_DID"), + JetstreamURL: envOr("TACK_JETSTREAM_URL", "wss://jetstream1.us-west.bsky.network/subscribe"), + DBPath: envOr("TACK_DB_PATH", "tack.db"), + Dev: os.Getenv("TACK_DEV") != "", + BuildkiteToken: os.Getenv("TACK_BUILDKITE_TOKEN"), + BuildkiteOrg: os.Getenv("TACK_BUILDKITE_ORG"), + BuildkitePipeline: os.Getenv("TACK_BUILDKITE_PIPELINE"), + BuildkiteWebhookSecret: os.Getenv("TACK_BUILDKITE_WEBHOOK_SECRET"), + BuildkiteWebhookMode: buildkite.WebhookMode( + envOr("TACK_BUILDKITE_WEBHOOK_MODE", string(buildkite.WebhookModeToken)), + ), } addrFlag := flag.String("addr", cfg.Addr, "HTTP listen address (overrides TACK_LISTEN_ADDR)") flag.Parse() @@ -56,6 +77,30 @@ func loadConfig() (config, error) { return cfg, errors.New("TACK_HOSTNAME is required") } + // If the operator opted into Buildkite mode (by supplying a + // token), every other Buildkite knob has to be present. Half- + // configured Buildkite leads to confusing failures deep in the + // provider; catch it at startup. + if cfg.BuildkiteToken != "" { + if cfg.BuildkiteOrg == "" { + return cfg, errors.New("TACK_BUILDKITE_ORG is required when TACK_BUILDKITE_TOKEN is set") + } + if cfg.BuildkitePipeline == "" { + return cfg, errors.New("TACK_BUILDKITE_PIPELINE is required when TACK_BUILDKITE_TOKEN is set") + } + if cfg.BuildkiteWebhookSecret == "" { + return cfg, errors.New("TACK_BUILDKITE_WEBHOOK_SECRET is required when TACK_BUILDKITE_TOKEN is set") + } + switch cfg.BuildkiteWebhookMode { + case buildkite.WebhookModeToken, buildkite.WebhookModeSignature: + default: + return cfg, fmt.Errorf("TACK_BUILDKITE_WEBHOOK_MODE must be %q or %q; got %q", + buildkite.WebhookModeToken, buildkite.WebhookModeSignature, + cfg.BuildkiteWebhookMode, + ) + } + } + return cfg, nil } @@ -115,12 +160,38 @@ func main() { br := newBroker(st) // Provider that turns Tangled pipeline triggers into - // pipeline.status events. The fake provider stands in for a real - // CI integration: it emits synthetic running/success heartbeats - // over the broker so the entire jetstream → knot → /events flow - // is exercisable end-to-end. Swap this for a Buildkite-backed - // implementation once that lands. - provider := newFakeProvider(br, logger) + // pipeline.status events. The Buildkite provider is the real + // integration; the fake one stands in when no Buildkite token is + // configured so the full jetstream → knot → /events flow is + // still exercisable locally without a Buildkite account. + // + // bkProvider is kept as a typed pointer separately because the + // /webhooks/buildkite handler needs the concrete *buildkiteProvider + // (for HandleWebhook + signature verification), not the abstract + // Provider surface. + var ( + provider Provider + bkProvider *buildkiteProvider + ) + if cfg.BuildkiteToken != "" { + bkProvider = newBuildkiteProvider( + br, st, + buildkite.NewClient(cfg.BuildkiteToken, cfg.BuildkiteOrg), + cfg.BuildkitePipeline, + cfg.BuildkiteWebhookSecret, + cfg.BuildkiteWebhookMode, + logger, + ) + provider = bkProvider + logger.Info("buildkite provider enabled", + "org", cfg.BuildkiteOrg, + "pipeline", cfg.BuildkitePipeline, + "webhook_mode", cfg.BuildkiteWebhookMode, + ) + } else { + provider = newFakeProvider(br, logger) + logger.Info("fake provider enabled (set TACK_BUILDKITE_TOKEN to use buildkite)") + } // Start the knot event-stream consumer first so the jetstream // loop has somewhere to register newly-observed knots into. It @@ -143,7 +214,7 @@ func main() { // Run the HTTP server. This blocks until ctx is cancelled or the // listener errors. - if err := runHTTP(ctx, cfg, br, provider); err != nil { + if err := runHTTP(ctx, cfg, br, provider, bkProvider); err != nil { logger.Error("http server error", "err", err) os.Exit(1) } diff --git a/provider.go b/provider.go index 0503e5d..cbe05b1 100644 --- a/provider.go +++ b/provider.go @@ -73,7 +73,10 @@ type Provider interface { // knot is the knot hostname the trigger arrived on; it's the // authority half of the pipeline ATURI that pipeline.status // records reference. pipelineRkey is the trigger record's rkey - // on that knot. workflows is the unmodified slice from the + // on that knot. trigger is the decoded record's TriggerMetadata + // (may be nil — the lexicon doesn't enforce its presence) and + // carries the commit/branch/PR data a real CI provider needs to + // kick off a build. workflows is the unmodified slice from the // decoded sh.tangled.pipeline record; implementations should // tolerate nil entries and zero-length names defensively, since // the lexicon doesn't enforce either. @@ -81,6 +84,7 @@ type Provider interface { ctx context.Context, knot string, pipelineRkey string, + trigger *tangled.Pipeline_TriggerMetadata, workflows []*tangled.Pipeline_Workflow, ) diff --git a/provider_buildkite.go b/provider_buildkite.go new file mode 100644 index 0000000..3602608 --- /dev/null +++ b/provider_buildkite.go @@ -0,0 +1,552 @@ +package main + +// buildkiteProvider implements Provider against a real Buildkite +// account. Spawn translates a Tangled pipeline trigger into one +// Buildkite build per workflow; status updates flow back asynchronously +// through the /webhooks/buildkite handler (see http.go), which looks +// the build UUID up in the buildkite_builds table to recover the +// (knot, pipelineRkey, workflow) tuple this provider persisted at +// Spawn time and publishes a sh.tangled.pipeline.status record on +// the in-process broker. +// +// Only one Buildkite pipeline is used per spindle (TACK_BUILDKITE_PIPELINE). +// Every Tangled workflow runs as a build on that single pipeline, with +// the workflow identity plumbed through env + meta_data. The operator +// configures their Buildkite pipeline to read those env vars and +// dispatch accordingly (e.g. via `pipeline upload`). Mapping every +// Tangled workflow to its own Buildkite pipeline would force operators +// to provision Buildkite resources for each workflow file in every +// repo that points at the spindle — friction we don't want to impose. + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "log/slog" + "net/http" + "strings" + "time" + + "tangled.org/core/api/tangled" + + "github.com/mitchellh/tack/internal/buildkite" +) + +// Buildkite-side meta_data keys carrying the Tangled identity of a +// build. Mirrored into env vars (see envFromTuple) so an operator's +// Buildkite pipeline script can also reach them via $TACK_*. They +// stay tightly namespaced so a coexisting Buildkite job that uses +// meta_data for its own purposes won't collide. +const ( + bkMetaKnot = "tack:knot" + bkMetaPipelineRkey = "tack:pipeline_rkey" + bkMetaWorkflow = "tack:workflow" +) + +// buildkiteProvider implements Provider. +// +// webhookSecret + webhookMode live on the provider rather than on +// the HTTP server because the provider is the single owner of +// "everything Buildkite-y": colocating the auth knob with the API +// client and the state translator keeps configuration drift to one +// place and makes the http.go side pure transport. +type buildkiteProvider struct { + br *broker + st *store + log *slog.Logger + client *buildkite.Client + pipelineSlug string + webhookSecret string + webhookMode buildkite.WebhookMode +} + +// Compile-time interface conformance check. +var _ Provider = (*buildkiteProvider)(nil) + +// newBuildkiteProvider wires a provider to its Buildkite client and +// to the broker it publishes pipeline.status records on. pipelineSlug +// is the Buildkite pipeline that all builds get fired on (see file +// header for why there's only one). webhookSecret/webhookMode govern +// inbound /webhooks/buildkite request authentication. +func newBuildkiteProvider( + br *broker, + st *store, + client *buildkite.Client, + pipelineSlug string, + webhookSecret string, + webhookMode buildkite.WebhookMode, + log *slog.Logger, +) *buildkiteProvider { + return &buildkiteProvider{ + br: br, + st: st, + log: log.With("component", "provider", "kind", "buildkite"), + client: client, + pipelineSlug: pipelineSlug, + webhookSecret: webhookSecret, + webhookMode: webhookMode, + } +} + +// VerifyWebhook authenticates an inbound webhook request using +// whichever mode the provider was configured with. Returns nil on +// success; the HTTP handler maps any returned error to 401. +func (p *buildkiteProvider) VerifyWebhook(headers http.Header, body []byte) error { + switch p.webhookMode { + case buildkite.WebhookModeSignature: + return buildkite.VerifySignature( + headers.Get("X-Buildkite-Signature"), + p.webhookSecret, body, + ) + default: + // Token mode is the Buildkite default and our default, so + // any unrecognised value falls through to it rather than + // fail-closed at startup. + return buildkite.VerifyToken( + headers.Get("X-Buildkite-Token"), + p.webhookSecret, + ) + } +} + +// Spawn satisfies Provider. For each workflow it fires a separate +// Buildkite build off the configured pipeline so each workflow gets +// its own status timeline. The actual API call runs on a goroutine — +// CreateBuild is one HTTP round-trip, but we still want Spawn to be +// non-blocking per the interface contract. +// +// On a successful create we persist the build UUID → (knot, rkey, +// workflow) mapping and publish a "pending" pipeline.status so the +// appview sees activity immediately, instead of waiting for the +// first webhook to land. +func (p *buildkiteProvider) Spawn( + ctx context.Context, + knot string, + pipelineRkey string, + trigger *tangled.Pipeline_TriggerMetadata, + workflows []*tangled.Pipeline_Workflow, +) { + if len(workflows) == 0 { + p.log.Warn("pipeline has no workflows; nothing to spawn", + "knot", knot, "rkey", pipelineRkey, + ) + return + } + + // Derive build inputs once. Every workflow on this trigger + // targets the same commit/branch — only the workflow name + // varies between the per-workflow goroutines below. + commit, branch := triggerCommitAndBranch(trigger) + if commit == "" { + // Buildkite's create-build API requires a commit; we'd + // rather log loudly and skip than fire builds on "HEAD" + // and silently get whatever main happens to look like. + p.log.Error("trigger has no commit; refusing to spawn", + "knot", knot, "rkey", pipelineRkey, + ) + return + } + + for _, wf := range workflows { + if wf == nil || wf.Name == "" { + continue + } + wf := wf + go p.spawnWorkflow(ctx, knot, pipelineRkey, commit, branch, wf) + } +} + +// spawnWorkflow does the per-workflow API + persistence work for +// Spawn. Errors are logged with full context but not returned — +// nothing in tack consumes the result, and a failed Spawn just +// surfaces as the absence of any status update for the affected +// workflow. +func (p *buildkiteProvider) spawnWorkflow( + ctx context.Context, + knot string, + pipelineRkey string, + commit string, + branch string, + wf *tangled.Pipeline_Workflow, +) { + logger := p.log.With( + "knot", knot, + "pipeline_rkey", pipelineRkey, + "workflow", wf.Name, + ) + + pipelineURI := pipelineATURI(knot, pipelineRkey) + meta := map[string]string{ + bkMetaKnot: knot, + bkMetaPipelineRkey: pipelineRkey, + bkMetaWorkflow: wf.Name, + } + env := envFromTuple(knot, pipelineRkey, wf) + + req := buildkite.CreateBuildRequest{ + Commit: commit, + Branch: branch, + Message: fmt.Sprintf("tangled: %s", wf.Name), + Env: env, + MetaData: meta, + IgnorePipelineBranchFilters: true, + } + + build, err := p.client.CreateBuild(ctx, p.pipelineSlug, req) + if err != nil { + logger.Error("create buildkite build", "err", err) + return + } + logger.Info("buildkite build created", + "build_uuid", build.ID, + "build_number", build.Number, + "web_url", build.WebURL, + ) + + if err := p.st.InsertBuildkiteBuild(ctx, BuildkiteBuildRef{ + BuildUUID: build.ID, + BuildNumber: build.Number, + PipelineSlug: p.pipelineSlug, + Knot: knot, + PipelineRkey: pipelineRkey, + Workflow: wf.Name, + PipelineURI: pipelineURI, + }); err != nil { + // Webhook handlers will fail to translate this build's + // events because they can't recover the tuple. Surface + // loudly and bail; we don't want a half-tracked build + // silently leaking status into the broker. + logger.Error("persist buildkite build mapping", "err", err, + "build_uuid", build.ID, + ) + return + } + + // Initial status publish so the appview shows the build as + // queued without waiting for the first webhook. This mirrors + // the upstream spindle's "schedule then run" cadence. + if err := p.publishStatus( + ctx, pipelineURI, wf.Name, "pending", build.ID, + nil, nil, + ); err != nil { + logger.Error("publish initial pending status", "err", err) + } +} + +// Logs satisfies Provider. We resolve the (knot, rkey, workflow) +// tuple to a Buildkite build via the store, fetch the current jobs +// list, then drain each job's plain-text log into the channel as one +// LogLine per output line. +// +// Per-job control frames bracket each job so the appview's renderer +// has start/end markers to lay out timing — same shape as the fake +// provider and the upstream spindle. +// +// This is a snapshot read, not a tail — finished or in-progress, we +// fetch what's there and close. Live tailing would require Buildkite +// agent log streaming, which the public REST API doesn't expose; the +// appview's repeated /logs calls during a running build give us +// "good enough" liveness without that complexity. +func (p *buildkiteProvider) Logs( + ctx context.Context, + knot string, + pipelineRkey string, + workflow string, +) (<-chan LogLine, error) { + ref, err := p.st.LookupBuildkiteBuildByTuple(ctx, knot, pipelineRkey, workflow) + if err != nil { + return nil, fmt.Errorf("lookup build for logs: %w", err) + } + if ref == nil { + return nil, ErrLogsNotFound + } + + // Fresh fetch so we get the current job set, not whatever was + // returned at create time (when most jobs are still nil). The + // upstream's not-found is mapped to the Provider-shaped one + // here because the /logs handler only knows about ErrLogsNotFound. + build, err := p.client.GetBuild(ctx, ref.PipelineSlug, ref.BuildNumber) + if err != nil { + if errors.Is(err, buildkite.ErrNotFound) { + return nil, ErrLogsNotFound + } + return nil, fmt.Errorf("get build for logs: %w", err) + } + + out := make(chan LogLine, 64) + go func() { + defer close(out) + stepID := 0 + for _, job := range build.Jobs { + // Only "script" jobs have agent-produced logs. + // Waiter / manual / trigger jobs have no body to + // fetch; skip them so we don't hit Buildkite with + // 404-bound requests. + if job.Type != "" && job.Type != "script" { + continue + } + + name := job.Name + if name == "" { + name = job.ID + } + + // Job-level start frame so the appview can bound + // timing per job. + if !sendLine(ctx, out, LogLine{ + Kind: LogKindControl, + Time: time.Now(), + Content: name, + StepId: stepID, + StepStatus: StepStatusStart, + }) { + return + } + + body, err := p.client.GetJobLog(ctx, ref.PipelineSlug, ref.BuildNumber, job.ID) + if err != nil { + p.log.Debug("fetch job log", + "err", err, + "build_uuid", ref.BuildUUID, + "job_id", job.ID, + ) + // Don't fail the whole stream on one job; + // emit the end frame and move on so the + // appview at least sees what other jobs + // produced. + body = "" + } + + for _, line := range strings.Split(strings.TrimRight(body, "\n"), "\n") { + if line == "" { + // Skip the leading empty entry that + // Split produces for empty bodies. + continue + } + if !sendLine(ctx, out, LogLine{ + Kind: LogKindData, + Time: time.Now(), + Content: line + "\n", + StepId: stepID, + Stream: "stdout", + }) { + return + } + } + + if !sendLine(ctx, out, LogLine{ + Kind: LogKindControl, + Time: time.Now(), + Content: name, + StepId: stepID, + StepStatus: StepStatusEnd, + }) { + return + } + stepID++ + } + }() + return out, nil +} + +// publishStatus assembles a tangled.PipelineStatus record and pushes +// it through the broker. buildUUID is mixed into the rkey so multiple +// status events for the same workflow don't collide on the events +// table's (rkey) uniqueness — and so an operator grepping the log +// can find every record that pertains to a given Buildkite build. +// +// errMsg/exitCode are optional; pass nil for non-failure transitions. +func (p *buildkiteProvider) publishStatus( + ctx context.Context, + pipelineURI, workflow, status, buildUUID string, + errMsg *string, + exitCode *int64, +) error { + rec := tangled.PipelineStatus{ + LexiconTypeID: tangled.PipelineStatusNSID, + Pipeline: pipelineURI, + Workflow: workflow, + Status: status, + CreatedAt: time.Now().UTC().Format(time.RFC3339), + Error: errMsg, + ExitCode: exitCode, + } + body, err := json.Marshal(rec) + if err != nil { + return fmt.Errorf("marshal pipeline.status: %w", err) + } + rkey := fmt.Sprintf("bk-%s-%s-%d", buildUUID, status, time.Now().UnixNano()) + if _, err := p.br.Publish(ctx, rkey, tangled.PipelineStatusNSID, body); err != nil { + return fmt.Errorf("publish pipeline.status: %w", err) + } + return nil +} + +// HandleWebhook applies a decoded Buildkite webhook payload: looks +// the build up in the store, translates the Buildkite state into a +// Tangled StatusKind, and publishes a pipeline.status record. Used +// by the HTTP webhook handler so both the ingress logic and the +// translation logic live next to each other. +// +// Returns nil for events we intentionally ignore (job.* events, +// build.scheduled which we already publish locally on Spawn, builds +// we don't have a mapping for) so the handler can 200 them — webhook +// retries from Buildkite on a 4xx/5xx are noisy and not what we want +// for "we just don't care about this event". +func (p *buildkiteProvider) HandleWebhook( + ctx context.Context, + payload buildkite.WebhookPayload, +) error { + // Only build.* events drive pipeline.status today. Everything + // else (job.*, agent.*, ping) is acknowledged silently. + if !strings.HasPrefix(payload.Event, "build.") { + return nil + } + + ref, err := p.st.LookupBuildkiteBuildByUUID(ctx, payload.Build.ID) + if err != nil { + return fmt.Errorf("lookup build by uuid: %w", err) + } + if ref == nil { + // Most likely: this build was triggered outside tack and + // just happens to share our webhook URL. Nothing to do. + p.log.Debug("webhook for unknown build; ignoring", + "event", payload.Event, + "build_uuid", payload.Build.ID, + ) + return nil + } + + status, ok := mapBuildkiteState(payload.Build.State) + if !ok { + // Unknown / transient state ("blocked", "skipped", + // "not_run", "waiting"…) — log so we can extend the map + // later, but don't error out the webhook. + p.log.Debug("unmapped buildkite state; ignoring", + "event", payload.Event, + "state", payload.Build.State, + "build_uuid", payload.Build.ID, + ) + return nil + } + + if err := p.publishStatus(ctx, ref.PipelineURI, ref.Workflow, + status, ref.BuildUUID, nil, nil); err != nil { + return fmt.Errorf("publish webhook status: %w", err) + } + p.log.Info("buildkite webhook → pipeline.status", + "event", payload.Event, + "state", payload.Build.State, + "status", status, + "build_uuid", payload.Build.ID, + "workflow", ref.Workflow, + ) + return nil +} + +// mapBuildkiteState translates Buildkite's build state strings into +// the Tangled spindle StatusKind enum. The mapping aligns with the +// upstream constants (StatusKindRunning/Failed/Cancelled/Success); +// states that don't have a direct analogue (blocked, skipped, +// not_run) are reported as not-mapped so the caller can decide +// whether to ignore them. +func mapBuildkiteState(state string) (string, bool) { + switch state { + case "scheduled": + return "pending", true + case "running", "failing": + return "running", true + case "passed": + return "success", true + case "failed": + return "failed", true + case "canceled", "canceling": + return "cancelled", true + default: + return "", false + } +} + +// envFromTuple builds the env block forwarded into the Buildkite +// build. These are the only handle a user's Buildkite pipeline has +// on the originating Tangled trigger: their pipeline.yml typically +// reads $TACK_WORKFLOW and dispatches based on it (e.g. running a +// `pipeline upload` against a workflow-specific YAML file). +// +// TACK_WORKFLOW_RAW carries the entire YAML body of the workflow as +// captured in the Tangled record. It can be empty if the workflow +// definition omitted it; consumers should defend. +func envFromTuple(knot, pipelineRkey string, wf *tangled.Pipeline_Workflow) map[string]string { + return map[string]string{ + "TACK_KNOT": knot, + "TACK_PIPELINE_RKEY": pipelineRkey, + "TACK_WORKFLOW": wf.Name, + "TACK_WORKFLOW_RAW": wf.Raw, + } +} + +// pipelineATURI returns the at-uri the appview joins pipeline.status +// records back to their originating pipeline on. Format mirrors the +// upstream spindle; the appview strips the `did:web:` prefix and +// treats the remainder as the knot identifier. +func pipelineATURI(knot, pipelineRkey string) string { + return fmt.Sprintf("at://did:web:%s/%s/%s", + knot, tangled.PipelineNSID, pipelineRkey, + ) +} + +// triggerCommitAndBranch extracts (commit, branch) from a Tangled +// pipeline trigger, regardless of whether it was a push, a pull +// request, or a manual run. Returns empty strings on a fully-empty +// trigger so the caller can decide whether that's fatal. +func triggerCommitAndBranch(trigger *tangled.Pipeline_TriggerMetadata) (string, string) { + if trigger == nil { + return "", "" + } + switch { + case trigger.Push != nil: + // For push events, NewSha is the commit being built and + // Ref is the full ref (e.g. "refs/heads/main") — strip + // the prefix so Buildkite's branch-aware features work. + return trigger.Push.NewSha, refToBranch(trigger.Push.Ref) + case trigger.PullRequest != nil: + // PRs build the source commit on the source branch. + // Buildkite's pipeline can opt into PR-aware behaviour + // via pull_request_id (not currently plumbed through). + return trigger.PullRequest.SourceSha, trigger.PullRequest.SourceBranch + default: + // Manual triggers and any future kinds: fall back to the + // repo default branch with no commit, which the caller + // will treat as fatal — manual triggers will need + // additional plumbing to pick a commit. + if trigger.Repo != nil { + return "", trigger.Repo.DefaultBranch + } + return "", "" + } +} + +// refToBranch strips the conventional refs/heads/ prefix from a git +// ref. Refs that don't match the prefix (tags, refs/pull/N/head) are +// returned as-is so downstream tooling can decide what to do with +// them — Buildkite happily accepts either form in `branch`. +func refToBranch(ref string) string { + const prefix = "refs/heads/" + if strings.HasPrefix(ref, prefix) { + return strings.TrimPrefix(ref, prefix) + } + return ref +} + +// sendLine pushes one LogLine into out, returning false if ctx +// fired first. Centralised so the per-job loop in Logs stays +// focused on the wire-shape decisions. +func sendLine(ctx context.Context, out chan<- LogLine, line LogLine) bool { + select { + case <-ctx.Done(): + return false + case out <- line: + return true + } +} diff --git a/provider_buildkite_test.go b/provider_buildkite_test.go new file mode 100644 index 0000000..3acdaa5 --- /dev/null +++ b/provider_buildkite_test.go @@ -0,0 +1,405 @@ +package main + +// Provider-level integration tests for the Buildkite implementation: +// Spawn → CreateBuild + persist + initial pending publish, and +// HandleWebhook → translate state + publish status. Buildkite itself +// is stubbed with httptest so the tests don't need network access. + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "io" + "log/slog" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "tangled.org/core/api/tangled" + + "github.com/mitchellh/tack/internal/buildkite" +) + +// newBuildkiteTestProvider wires a buildkiteProvider against an +// httptest server impersonating api.buildkite.com. Returns the +// store/broker so tests can inspect publishes + persistence. +func newBuildkiteTestProvider( + t *testing.T, + mode buildkite.WebhookMode, + secret string, + bkHandler http.HandlerFunc, +) (*buildkiteProvider, *store, *broker, *httptest.Server) { + t.Helper() + srv := httptest.NewServer(bkHandler) + t.Cleanup(srv.Close) + + prev := buildkite.APIBase + buildkite.APIBase = srv.URL + t.Cleanup(func() { buildkite.APIBase = prev }) + + st := newTestStore(t) + br := newBroker(st) + logger := slog.Default() + p := newBuildkiteProvider( + br, st, + buildkite.NewClient("tok", "myorg"), + "mypipe", + secret, mode, + logger, + ) + return p, st, br, srv +} + +// TestBuildkiteSpawn covers the full create-build path: trigger → +// API call → DB row → "pending" status on the broker. +func TestBuildkiteSpawn(t *testing.T) { + bk := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusCreated) + _ = json.NewEncoder(w).Encode(buildkite.Build{ + ID: "uuid-1", + Number: 7, + }) + }) + p, st, _, _ := newBuildkiteTestProvider(t, buildkite.WebhookModeToken, "secret", bk) + + trigger := &tangled.Pipeline_TriggerMetadata{ + Push: &tangled.Pipeline_PushTriggerData{ + NewSha: "abcdef0123", + Ref: "refs/heads/main", + }, + } + workflows := []*tangled.Pipeline_Workflow{ + {Name: "test.yml", Raw: "steps:\n - run: true\n"}, + } + + p.Spawn(context.Background(), "knot.example.com", "rkey-1", trigger, workflows) + + // Spawn fans out into goroutines; wait briefly for the side + // effects to land. The store row is the load-bearing artifact + // — once it's present, the publish has already happened too. + deadline := time.Now().Add(2 * time.Second) + var ref *BuildkiteBuildRef + for time.Now().Before(deadline) { + var err error + ref, err = st.LookupBuildkiteBuildByUUID(context.Background(), "uuid-1") + if err != nil { + t.Fatalf("lookup: %v", err) + } + if ref != nil { + break + } + time.Sleep(20 * time.Millisecond) + } + if ref == nil { + t.Fatal("buildkite build row not persisted within deadline") + } + if ref.Workflow != "test.yml" || ref.Knot != "knot.example.com" || ref.PipelineRkey != "rkey-1" { + t.Fatalf("ref mismatch: %+v", ref) + } + if ref.PipelineSlug != "mypipe" || ref.BuildNumber != 7 { + t.Fatalf("buildkite ref mismatch: %+v", ref) + } + + // One pending status should be on the events log. + rows, err := st.EventsAfter(context.Background(), 0) + if err != nil { + t.Fatalf("EventsAfter: %v", err) + } + if len(rows) != 1 { + t.Fatalf("got %d events, want 1", len(rows)) + } + var rec tangled.PipelineStatus + if err := json.Unmarshal(rows[0].EventJSON, &rec); err != nil { + t.Fatalf("decode status: %v", err) + } + if rec.Status != "pending" || rec.Workflow != "test.yml" { + t.Fatalf("unexpected status: %+v", rec) + } + if !strings.Contains(rec.Pipeline, "knot.example.com") || + !strings.Contains(rec.Pipeline, "rkey-1") { + t.Fatalf("pipeline ATURI wrong: %s", rec.Pipeline) + } +} + +// TestBuildkiteSpawnNoCommit confirms we don't fire a build when the +// trigger has no commit to build — kicking one off would resolve to +// whatever main looks like at agent-fetch time, which is dangerously +// surprising. +func TestBuildkiteSpawnNoCommit(t *testing.T) { + called := false + bk := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + called = true + }) + p, st, _, _ := newBuildkiteTestProvider(t, buildkite.WebhookModeToken, "secret", bk) + + p.Spawn(context.Background(), "knot.example.com", "rkey-1", + &tangled.Pipeline_TriggerMetadata{Manual: &tangled.Pipeline_ManualTriggerData{}}, + []*tangled.Pipeline_Workflow{{Name: "test.yml"}}, + ) + + // Give any rogue goroutine a moment. + time.Sleep(50 * time.Millisecond) + if called { + t.Fatal("CreateBuild called despite missing commit") + } + rows, _ := st.EventsAfter(context.Background(), 0) + if len(rows) != 0 { + t.Fatalf("got %d events, want 0", len(rows)) + } +} + +// TestBuildkiteHandleWebhook checks the translation pipeline: +// recorded build + matching webhook → success status published. +func TestBuildkiteHandleWebhook(t *testing.T) { + p, st, _, _ := newBuildkiteTestProvider(t, buildkite.WebhookModeToken, "secret", + func(w http.ResponseWriter, r *http.Request) { t.Fatal("buildkite shouldn't be called") }) + + // Pre-seed a known build mapping. + if err := st.InsertBuildkiteBuild(context.Background(), BuildkiteBuildRef{ + BuildUUID: "uuid-1", + BuildNumber: 7, + PipelineSlug: "mypipe", + Knot: "knot.example.com", + PipelineRkey: "rkey-1", + Workflow: "test.yml", + PipelineURI: "at://did:web:knot.example.com/sh.tangled.pipeline/rkey-1", + }); err != nil { + t.Fatalf("InsertBuildkiteBuild: %v", err) + } + + err := p.HandleWebhook(context.Background(), buildkite.WebhookPayload{ + Event: "build.finished", + Build: buildkite.Build{ID: "uuid-1", State: "passed"}, + }) + if err != nil { + t.Fatalf("HandleWebhook: %v", err) + } + + rows, err := st.EventsAfter(context.Background(), 0) + if err != nil { + t.Fatalf("EventsAfter: %v", err) + } + if len(rows) != 1 { + t.Fatalf("got %d events, want 1", len(rows)) + } + var rec tangled.PipelineStatus + if err := json.Unmarshal(rows[0].EventJSON, &rec); err != nil { + t.Fatalf("decode: %v", err) + } + if rec.Status != "success" || rec.Workflow != "test.yml" { + t.Fatalf("bad status: %+v", rec) + } +} + +// TestBuildkiteHandleWebhookIgnored covers the "we don't care" paths: +// non-build events and unknown builds must be no-op (no publish, no +// error) so Buildkite doesn't retry them. +func TestBuildkiteHandleWebhookIgnored(t *testing.T) { + p, st, _, _ := newBuildkiteTestProvider(t, buildkite.WebhookModeToken, "secret", + func(w http.ResponseWriter, r *http.Request) {}) + + // Non-build event: no lookup, no publish. + if err := p.HandleWebhook(context.Background(), buildkite.WebhookPayload{ + Event: "job.started", + Build: buildkite.Build{ID: "uuid-x"}, + }); err != nil { + t.Fatalf("HandleWebhook (job.started): %v", err) + } + + // Build event for unknown UUID: no publish. + if err := p.HandleWebhook(context.Background(), buildkite.WebhookPayload{ + Event: "build.finished", + Build: buildkite.Build{ID: "unknown-uuid", State: "passed"}, + }); err != nil { + t.Fatalf("HandleWebhook (unknown): %v", err) + } + + // Known build but unmapped state: no publish. + if err := st.InsertBuildkiteBuild(context.Background(), BuildkiteBuildRef{ + BuildUUID: "uuid-blocked", PipelineSlug: "mypipe", + Knot: "k", PipelineRkey: "r", Workflow: "w", + PipelineURI: "at://x", + }); err != nil { + t.Fatalf("seed: %v", err) + } + if err := p.HandleWebhook(context.Background(), buildkite.WebhookPayload{ + Event: "build.finished", + Build: buildkite.Build{ID: "uuid-blocked", State: "blocked"}, + }); err != nil { + t.Fatalf("HandleWebhook (blocked): %v", err) + } + + rows, _ := st.EventsAfter(context.Background(), 0) + if len(rows) != 0 { + t.Fatalf("got %d events, want 0", len(rows)) + } +} + +// TestBuildkiteWebhookHandlerHTTP exercises the full HTTP path +// including auth: a request signed with the wrong secret must be +// rejected, and a correctly-signed one must reach the provider. +func TestBuildkiteWebhookHandlerHTTP(t *testing.T) { + // Signature mode is the more interesting code path; we cover + // token mode in the verifier-level tests above. + const secret = "swordfish" + p, st, _, _ := newBuildkiteTestProvider(t, buildkite.WebhookModeSignature, secret, + func(w http.ResponseWriter, r *http.Request) { /* unused */ }) + + // Pre-seed so the provider's HandleWebhook can resolve the build. + if err := st.InsertBuildkiteBuild(context.Background(), BuildkiteBuildRef{ + BuildUUID: "uuid-2", + BuildNumber: 9, + PipelineSlug: "mypipe", + Knot: "knot.example.com", + PipelineRkey: "rkey-2", + Workflow: "test.yml", + PipelineURI: "at://did:web:knot.example.com/sh.tangled.pipeline/rkey-2", + }); err != nil { + t.Fatalf("seed: %v", err) + } + + body, _ := json.Marshal(map[string]any{ + "event": "build.finished", + "build": map[string]any{ + "id": "uuid-2", + "state": "failed", + }, + }) + + logger := slog.Default() + handler := buildkiteWebhookHandler(logger, p) + + // Unsigned request → 401. + t.Run("rejects unsigned", func(t *testing.T) { + req := httptest.NewRequest(http.MethodPost, "/webhooks/buildkite", + strings.NewReader(string(body))) + req.Header.Set("X-Buildkite-Event", "build.finished") + w := httptest.NewRecorder() + handler(w, req) + if w.Code != http.StatusUnauthorized { + t.Fatalf("status = %d; want 401", w.Code) + } + }) + + // Wrong-secret request → 401. + t.Run("rejects bad signature", func(t *testing.T) { + ts := fmt.Sprintf("%d", time.Now().Unix()) + mac := hmac.New(sha256.New, []byte("wrong")) + mac.Write([]byte(ts)) + mac.Write([]byte(".")) + mac.Write(body) + sig := hex.EncodeToString(mac.Sum(nil)) + + req := httptest.NewRequest(http.MethodPost, "/webhooks/buildkite", + strings.NewReader(string(body))) + req.Header.Set("X-Buildkite-Event", "build.finished") + req.Header.Set("X-Buildkite-Signature", "timestamp="+ts+",signature="+sig) + w := httptest.NewRecorder() + handler(w, req) + if w.Code != http.StatusUnauthorized { + t.Fatalf("status = %d; want 401", w.Code) + } + }) + + // Valid request → 200, status published. + t.Run("accepts valid", func(t *testing.T) { + ts := fmt.Sprintf("%d", time.Now().Unix()) + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write([]byte(ts)) + mac.Write([]byte(".")) + mac.Write(body) + sig := hex.EncodeToString(mac.Sum(nil)) + + req := httptest.NewRequest(http.MethodPost, "/webhooks/buildkite", + strings.NewReader(string(body))) + req.Header.Set("X-Buildkite-Event", "build.finished") + req.Header.Set("X-Buildkite-Signature", "timestamp="+ts+",signature="+sig) + w := httptest.NewRecorder() + handler(w, req) + if w.Code != http.StatusOK { + b, _ := io.ReadAll(w.Body) + t.Fatalf("status = %d body=%s; want 200", w.Code, string(b)) + } + rows, _ := st.EventsAfter(context.Background(), 0) + if len(rows) != 1 { + t.Fatalf("got %d events, want 1", len(rows)) + } + var rec tangled.PipelineStatus + if err := json.Unmarshal(rows[0].EventJSON, &rec); err != nil { + t.Fatalf("decode: %v", err) + } + if rec.Status != "failed" { + t.Fatalf("status = %q; want failed", rec.Status) + } + }) +} + +// TestBuildkiteWebhookHandlerNoProvider confirms the 503 branch when +// tack is running with the fake provider — a misdirected webhook +// must get a clear "not configured here" instead of a misleading +// 200 OK that silently throws the event away. +func TestBuildkiteWebhookHandlerNoProvider(t *testing.T) { + handler := buildkiteWebhookHandler(slog.Default(), nil) + req := httptest.NewRequest(http.MethodPost, "/webhooks/buildkite", + strings.NewReader("{}")) + w := httptest.NewRecorder() + handler(w, req) + if w.Code != http.StatusServiceUnavailable { + t.Fatalf("status = %d; want 503", w.Code) + } +} + +// TestTriggerCommitAndBranch pins the trigger-shape mapping. Each +// case pairs an input trigger with the (commit, branch) tuple a +// real CI provider would feed into its build-creation API. +func TestTriggerCommitAndBranch(t *testing.T) { + cases := []struct { + name string + in *tangled.Pipeline_TriggerMetadata + wantCommit string + wantBranch string + }{ + {"nil", nil, "", ""}, + {"push refs/heads", + &tangled.Pipeline_TriggerMetadata{ + Push: &tangled.Pipeline_PushTriggerData{NewSha: "abc", Ref: "refs/heads/main"}, + }, + "abc", "main", + }, + {"push tag ref preserved", + &tangled.Pipeline_TriggerMetadata{ + Push: &tangled.Pipeline_PushTriggerData{NewSha: "abc", Ref: "refs/tags/v1"}, + }, + "abc", "refs/tags/v1", + }, + {"pull request", + &tangled.Pipeline_TriggerMetadata{ + PullRequest: &tangled.Pipeline_PullRequestTriggerData{ + SourceSha: "def", SourceBranch: "feature", + }, + }, + "def", "feature", + }, + {"manual with default branch", + &tangled.Pipeline_TriggerMetadata{ + Manual: &tangled.Pipeline_ManualTriggerData{}, + Repo: &tangled.Pipeline_TriggerRepo{DefaultBranch: "main"}, + }, + "", "main", + }, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + gotC, gotB := triggerCommitAndBranch(c.in) + if gotC != c.wantCommit || gotB != c.wantBranch { + t.Fatalf("got (%q,%q); want (%q,%q)", + gotC, gotB, c.wantCommit, c.wantBranch) + } + }) + } +} diff --git a/provider_fake.go b/provider_fake.go index 99cc764..4cd57b8 100644 --- a/provider_fake.go +++ b/provider_fake.go @@ -81,6 +81,7 @@ func (p *fakeProvider) Spawn( ctx context.Context, knot string, pipelineRkey string, + _ *tangled.Pipeline_TriggerMetadata, workflows []*tangled.Pipeline_Workflow, ) { if len(workflows) == 0 { diff --git a/store.go b/store.go index 341d5dd..a94518c 100644 --- a/store.go +++ b/store.go @@ -322,6 +322,109 @@ func (s *store) InsertEvent(ctx context.Context, rkey, nsid string, eventJSON [] return id, nil } +// BuildkiteBuildRef is the persisted mapping from one Buildkite build +// to the Tangled pipeline tuple that spawned it. It's the row written +// by the Buildkite provider at Spawn time and read back from two +// places: the webhook handler (by build UUID) when an event arrives, +// and the /logs handler (by knot+rkey+workflow) when an appview +// client asks for output. +type BuildkiteBuildRef struct { + BuildUUID string + BuildNumber int64 + PipelineSlug string + Knot string + PipelineRkey string + Workflow string + PipelineURI string +} + +// InsertBuildkiteBuild records that a Buildkite build was created on +// behalf of the given (knot, pipelineRkey, workflow) tuple. Uses +// INSERT OR REPLACE so that an unlikely build-uuid collision (or a +// Buildkite-side rebuild that re-fires us) just refreshes the row +// instead of failing. +func (s *store) InsertBuildkiteBuild(ctx context.Context, ref BuildkiteBuildRef) error { + _, err := s.db.ExecContext(ctx, + `INSERT INTO buildkite_builds ( + build_uuid, build_number, pipeline_slug, + knot, pipeline_rkey, workflow, + pipeline_uri, created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(build_uuid) DO UPDATE SET + build_number = excluded.build_number, + pipeline_slug = excluded.pipeline_slug, + knot = excluded.knot, + pipeline_rkey = excluded.pipeline_rkey, + workflow = excluded.workflow, + pipeline_uri = excluded.pipeline_uri, + created_at = excluded.created_at`, + ref.BuildUUID, ref.BuildNumber, ref.PipelineSlug, + ref.Knot, ref.PipelineRkey, ref.Workflow, + ref.PipelineURI, time.Now().UTC().Format(time.RFC3339Nano), + ) + if err != nil { + return fmt.Errorf("insert buildkite_build: %w", err) + } + return nil +} + +// LookupBuildkiteBuildByUUID returns the saved mapping for the given +// Buildkite build UUID, or nil when no such build is recorded. +// Returning a nil pointer rather than a sentinel error keeps the +// webhook handler's "we don't know about this build" branch a simple +// nil check. +func (s *store) LookupBuildkiteBuildByUUID(ctx context.Context, buildUUID string) (*BuildkiteBuildRef, error) { + var ref BuildkiteBuildRef + err := s.db.QueryRowContext(ctx, + `SELECT build_uuid, build_number, pipeline_slug, + knot, pipeline_rkey, workflow, pipeline_uri + FROM buildkite_builds WHERE build_uuid = ?`, + buildUUID, + ).Scan( + &ref.BuildUUID, &ref.BuildNumber, &ref.PipelineSlug, + &ref.Knot, &ref.PipelineRkey, &ref.Workflow, &ref.PipelineURI, + ) + if errors.Is(err, sql.ErrNoRows) { + return nil, nil + } + if err != nil { + return nil, fmt.Errorf("lookup buildkite_build by uuid: %w", err) + } + return &ref, nil +} + +// LookupBuildkiteBuildByTuple finds the most recently created build +// for (knot, pipelineRkey, workflow). Returns nil when no build has +// been recorded for that tuple — used by /logs to translate the +// appview's path-based identity back into something Buildkite knows. +// +// "Most recent" matters because a workflow may have multiple builds +// over time (rebuilds, re-triggers). We always serve logs for the +// latest run; older runs are still queryable by build UUID directly +// if anyone ever wants that. +func (s *store) LookupBuildkiteBuildByTuple(ctx context.Context, knot, pipelineRkey, workflow string) (*BuildkiteBuildRef, error) { + var ref BuildkiteBuildRef + err := s.db.QueryRowContext(ctx, + `SELECT build_uuid, build_number, pipeline_slug, + knot, pipeline_rkey, workflow, pipeline_uri + FROM buildkite_builds + WHERE knot = ? AND pipeline_rkey = ? AND workflow = ? + ORDER BY created_at DESC + LIMIT 1`, + knot, pipelineRkey, workflow, + ).Scan( + &ref.BuildUUID, &ref.BuildNumber, &ref.PipelineSlug, + &ref.Knot, &ref.PipelineRkey, &ref.Workflow, &ref.PipelineURI, + ) + if errors.Is(err, sql.ErrNoRows) { + return nil, nil + } + if err != nil { + return nil, fmt.Errorf("lookup buildkite_build by tuple: %w", err) + } + return &ref, nil +} + // EventsAfter returns every event row with `created` strictly greater // than cursor, in cursor order. Used by /events to backfill a // reconnecting subscriber and to drain newly-published rows on each diff --git a/store_migrate.go b/store_migrate.go index b1f589c..517bc8d 100644 --- a/store_migrate.go +++ b/store_migrate.go @@ -78,6 +78,32 @@ CREATE TABLE IF NOT EXISTS events ( event_json TEXT NOT NULL, inserted_at TEXT NOT NULL ); + +-- Mapping from a Buildkite build back to the Tangled pipeline that +-- spawned it. The Buildkite webhook receiver only knows the build +-- UUID; everything we need to publish a pipeline.status record +-- (knot, pipeline rkey, workflow name, full pipeline ATURI) lives +-- on this row. +-- +-- pipeline_uri is denormalized off (knot, pipeline_rkey) so the +-- webhook handler doesn't have to recompute the at:// string on +-- every event — it's a constant for the lifetime of the build and +-- the webhook is the hot path for status fan-out. +-- +-- The (knot, pipeline_rkey, workflow) index supports the /logs +-- handler, which only knows that tuple at request time. +CREATE TABLE IF NOT EXISTS buildkite_builds ( + build_uuid TEXT PRIMARY KEY, + build_number INTEGER NOT NULL, + pipeline_slug TEXT NOT NULL, + knot TEXT NOT NULL, + pipeline_rkey TEXT NOT NULL, + workflow TEXT NOT NULL, + pipeline_uri TEXT NOT NULL, + created_at TEXT NOT NULL +); +CREATE INDEX IF NOT EXISTS buildkite_builds_lookup + ON buildkite_builds (knot, pipeline_rkey, workflow); ` // migrate applies the schema. Safe to call repeatedly. -- 2.51.2