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}