Parent directory

http.go

7038 bytes
  1package main
  2
  3import (
  4	"encoding/json"
  5	"errors"
  6	"io"
  7	"log/slog"
  8	"net/http"
  9	"path/filepath"
 10	"regexp"
 11	"strings"
 12	"sync"
 13	"time"
 14)
 15
 16const (
 17	maxRequestBodySize       = 1024
 18	maxStatusRequestBodySize = 4096
 19	maxStatusRepositories    = 100
 20)
 21
 22var (
 23	commitSHA = regexp.MustCompile(`^[0-9a-f]{40}([0-9a-f]{24})?$`)
 24	dockerTag = regexp.MustCompile(`^[A-Za-z0-9_][A-Za-z0-9_.-]{0,127}$`)
 25)
 26
 27type buildRequest struct {
 28	Repository string `json:"repository"`
 29	SHA        string `json:"sha"`
 30	Ref        string `json:"ref"`
 31}
 32
 33type buildStatus struct {
 34	SHA      string    `json:"sha"`
 35	State    string    `json:"state"`
 36	QueuedAt time.Time `json:"queued_at"`
 37}
 38
 39type buildStatusRequest struct {
 40	Repositories []string `json:"repositories"`
 41}
 42
 43type buildStatusResponse struct {
 44	Statuses map[string]buildStatus `json:"statuses"`
 45}
 46
 47type buildStatuses struct {
 48	mu       sync.RWMutex
 49	database *database
 50	current  map[string]buildStatus
 51	buildIDs map[string]string
 52}
 53
 54func newBuildStatuses(database *database) *buildStatuses {
 55	return &buildStatuses{
 56		database: database,
 57		current:  make(map[string]buildStatus),
 58		buildIDs: make(map[string]string),
 59	}
 60}
 61
 62func (statuses *buildStatuses) queue(job buildJob) error {
 63	queuedAt := time.Now().UTC()
 64	if err := statuses.database.queue(job, queuedAt); err != nil {
 65		return err
 66	}
 67
 68	statuses.mu.Lock()
 69	defer statuses.mu.Unlock()
 70	statuses.buildIDs[job.Request.Repository] = job.ID
 71	statuses.current[job.Request.Repository] = buildStatus{
 72		SHA:      job.Request.SHA,
 73		State:    "queued",
 74		QueuedAt: queuedAt,
 75	}
 76	return nil
 77}
 78
 79func (statuses *buildStatuses) update(job buildJob, state string) bool {
 80	if err := statuses.database.update(job, state, time.Now().UTC()); err != nil {
 81		slog.Error("save build state", "build_id", job.ID, "state", state, "error", err)
 82	}
 83
 84	statuses.mu.Lock()
 85	defer statuses.mu.Unlock()
 86	if statuses.buildIDs[job.Request.Repository] != job.ID {
 87		return false
 88	}
 89	status := statuses.current[job.Request.Repository]
 90	status.State = state
 91	statuses.current[job.Request.Repository] = status
 92	return true
 93}
 94
 95func (statuses *buildStatuses) requested(repositories []string) map[string]buildStatus {
 96	statuses.mu.RLock()
 97	defer statuses.mu.RUnlock()
 98	result := make(map[string]buildStatus, len(repositories))
 99	for _, repository := range repositories {
100		if status, ok := statuses.current[repository]; ok {
101			result[repository] = status
102		}
103	}
104	return result
105}
106
107func buildHandler(jobs chan<- buildJob, statuses *buildStatuses) http.HandlerFunc {
108	return func(writer http.ResponseWriter, request *http.Request) {
109		request.Body = http.MaxBytesReader(writer, request.Body, maxRequestBodySize)
110		defer func() {
111			if err := request.Body.Close(); err != nil {
112				slog.Warn("close build request", "error", err)
113			}
114		}()
115
116		var buildRequest buildRequest
117		decoder := json.NewDecoder(request.Body)
118		decoder.DisallowUnknownFields()
119		if err := decoder.Decode(&buildRequest); err != nil || !validBuildRequest(buildRequest) {
120			slog.Warn("build request rejected", "reason", "invalid request", "remote_address", request.RemoteAddr)
121			http.Error(writer, "invalid build request", http.StatusBadRequest)
122			return
123		}
124		if err := ensureSingleJSONValue(decoder); err != nil {
125			slog.Warn("build request rejected", "reason", "invalid request", "remote_address", request.RemoteAddr)
126			http.Error(writer, "invalid build request", http.StatusBadRequest)
127			return
128		}
129		if !isBuildRef(buildRequest.Ref) {
130			slog.Info("build ignored", "repository", buildRequest.Repository, "sha", buildRequest.SHA, "ref", buildRequest.Ref)
131			writer.WriteHeader(http.StatusNoContent)
132			return
133		}
134
135		job, err := newBuildJob(buildRequest)
136		if err != nil {
137			slog.Error("generate build ID", "error", err)
138			http.Error(writer, "unable to queue build", http.StatusInternalServerError)
139			return
140		}
141		select {
142		case jobs <- job:
143			if err := statuses.queue(job); err != nil {
144				job.ready <- false
145				slog.Error("save queued build", "repository", buildRequest.Repository, "sha", buildRequest.SHA, "error", err)
146				http.Error(writer, "unable to queue build", http.StatusInternalServerError)
147				return
148			}
149			job.ready <- true
150			slog.Info("build queued", "repository", buildRequest.Repository, "sha", buildRequest.SHA, "ref", buildRequest.Ref)
151			writer.WriteHeader(http.StatusAccepted)
152		default:
153			slog.Error("build queue full", "repository", buildRequest.Repository, "sha", buildRequest.SHA, "ref", buildRequest.Ref)
154			http.Error(writer, "build queue is full", http.StatusServiceUnavailable)
155		}
156	}
157}
158
159func buildStatusHandler(statuses *buildStatuses) http.HandlerFunc {
160	return func(writer http.ResponseWriter, request *http.Request) {
161		request.Body = http.MaxBytesReader(writer, request.Body, maxStatusRequestBodySize)
162		defer func() {
163			if err := request.Body.Close(); err != nil {
164				slog.Warn("close build status request", "error", err)
165			}
166		}()
167
168		var statusRequest buildStatusRequest
169		decoder := json.NewDecoder(request.Body)
170		decoder.DisallowUnknownFields()
171		if err := decoder.Decode(&statusRequest); err != nil || !validBuildStatusRequest(statusRequest) || ensureSingleJSONValue(decoder) != nil {
172			slog.Warn("build status request rejected", "reason", "invalid request", "remote_address", request.RemoteAddr)
173			http.Error(writer, "invalid build status request", http.StatusBadRequest)
174			return
175		}
176
177		slog.Info("build statuses requested", "repository_count", len(statusRequest.Repositories), "remote_address", request.RemoteAddr)
178		writer.Header().Set("Content-Type", "application/json")
179		if err := json.NewEncoder(writer).Encode(buildStatusResponse{Statuses: statuses.requested(statusRequest.Repositories)}); err != nil {
180			slog.Warn("write build statuses", "error", err)
181		}
182	}
183}
184
185func validBuildRequest(request buildRequest) bool {
186	return validRepository(request.Repository) && commitSHA.MatchString(request.SHA) && validBuildRef(request.Ref)
187}
188
189func validBuildRef(ref string) bool {
190	if strings.HasPrefix(ref, "refs/heads/") {
191		return strings.TrimPrefix(ref, "refs/heads/") != ""
192	}
193	return strings.HasPrefix(ref, "refs/tags/") && dockerTag.MatchString(strings.TrimPrefix(ref, "refs/tags/"))
194}
195
196func isBuildRef(ref string) bool {
197	return ref == "refs/heads/main" || strings.HasPrefix(ref, "refs/tags/")
198}
199
200func validBuildStatusRequest(request buildStatusRequest) bool {
201	if len(request.Repositories) == 0 || len(request.Repositories) > maxStatusRepositories {
202		return false
203	}
204	for _, repository := range request.Repositories {
205		if !validRepository(repository) {
206			return false
207		}
208	}
209	return true
210}
211
212func validRepository(repository string) bool {
213	cleaned := filepath.Clean(repository)
214	return repository != "." && !filepath.IsAbs(repository) && cleaned == repository && cleaned != ".." && !strings.HasPrefix(cleaned, ".."+string(filepath.Separator))
215}
216
217func ensureSingleJSONValue(decoder *json.Decoder) error {
218	var extra any
219	if err := decoder.Decode(&extra); !errors.Is(err, io.EOF) {
220		return errors.New("request has extra JSON values")
221	}
222	return nil
223}