Diff
1diff --git a/Dockerfile b/Dockerfile
2index b1e6993094749a726d354e5a21b06f917cf21dc0..d7b9aee033bfe600b67d7f4ab71fcda2def3b7cf 100644
3--- a/Dockerfile
4+++ b/Dockerfile
5@@ -1,7 +1,9 @@
6 FROM golang:1.27-alpine AS worker-build
7
8 WORKDIR /src
9-COPY go.mod main.go ./
10+COPY go.mod go.sum ./
11+RUN go mod download
12+COPY *.go ./
13 RUN CGO_ENABLED=0 go build -trimpath -ldflags='-s -w' -o /worker .
14
15 FROM moby/buildkit:rootless AS buildkit
16diff --git a/config.go b/config.go
17index 5ff3ea8fb2348962f0398d21094424d0daaf4047..9e2a9ba3488b815777c8e34f12f9411258b3243d 100644
18--- a/config.go
19+++ b/config.go
20@@ -8,6 +8,7 @@ import (
21 type config struct {
22 registry string
23 buildkitAddress string
24+ databasePath string
25 }
26
27 func loadConfig() (config, error) {
28@@ -24,6 +25,7 @@ func loadConfig() (config, error) {
29 return config{
30 registry: registry,
31 buildkitAddress: buildkitAddress,
32+ databasePath: environmentOrDefault("DATABASE_PATH", "/work/builds.db"),
33 }, nil
34 }
35
36@@ -34,3 +36,10 @@ func requiredEnvironment(name string) (string, error) {
37 }
38 return value, nil
39 }
40+
41+func environmentOrDefault(name, defaultValue string) string {
42+ if value := os.Getenv(name); value != "" {
43+ return value
44+ }
45+ return defaultValue
46+}
47diff --git a/db.go b/db.go
48new file mode 100644
49index 0000000000000000000000000000000000000000..1e19b8fc0761a4d1c2176e02adfb0f40f0902aa7
50--- /dev/null
51+++ b/db.go
52@@ -0,0 +1,90 @@
53+package main
54+
55+import (
56+ "context"
57+ "database/sql"
58+ "fmt"
59+ "os"
60+ "path/filepath"
61+ "time"
62+
63+ _ "modernc.org/sqlite"
64+)
65+
66+type database struct {
67+ connection *sql.DB
68+}
69+
70+func openDatabase(path string) (*database, error) {
71+ if err := os.MkdirAll(filepath.Dir(path), 0o750); err != nil {
72+ return nil, fmt.Errorf("create database directory: %w", err)
73+ }
74+
75+ connection, err := sql.Open("sqlite", path)
76+ if err != nil {
77+ return nil, fmt.Errorf("open database: %w", err)
78+ }
79+ connection.SetMaxOpenConns(1)
80+
81+ database := &database{connection: connection}
82+ if err := database.bootstrap(context.Background()); err != nil {
83+ _ = connection.Close()
84+ return nil, err
85+ }
86+ return database, nil
87+}
88+
89+func (database *database) bootstrap(ctx context.Context) error {
90+ _, err := database.connection.ExecContext(ctx, `
91+ CREATE TABLE IF NOT EXISTS builds (
92+ build_id TEXT PRIMARY KEY,
93+ repository TEXT NOT NULL,
94+ sha TEXT NOT NULL,
95+ state TEXT NOT NULL,
96+ queued_at INTEGER NOT NULL,
97+ started_at INTEGER,
98+ finished_at INTEGER
99+ );
100+ CREATE INDEX IF NOT EXISTS builds_repository_queued_at ON builds (repository, queued_at DESC);
101+ `)
102+ if err != nil {
103+ return fmt.Errorf("bootstrap database: %w", err)
104+ }
105+ return nil
106+}
107+
108+func (database *database) close() error {
109+ return database.connection.Close()
110+}
111+
112+func (database *database) queue(job buildJob, queuedAt time.Time) error {
113+ _, err := database.connection.Exec(
114+ `INSERT INTO builds (build_id, repository, sha, state, queued_at) VALUES (?, ?, ?, ?, ?)`,
115+ job.ID,
116+ job.Request.Repository,
117+ job.Request.SHA,
118+ "queued",
119+ queuedAt.Unix(),
120+ )
121+ if err != nil {
122+ return fmt.Errorf("save queued build: %w", err)
123+ }
124+ return nil
125+}
126+
127+func (database *database) update(job buildJob, state string, updatedAt time.Time) error {
128+ var query string
129+ switch state {
130+ case "running":
131+ query = `UPDATE builds SET state = ?, started_at = ? WHERE build_id = ?`
132+ case "passed", "failed":
133+ query = `UPDATE builds SET state = ?, finished_at = ? WHERE build_id = ?`
134+ default:
135+ return fmt.Errorf("unsupported build state %q", state)
136+ }
137+
138+ if _, err := database.connection.Exec(query, state, updatedAt.Unix(), job.ID); err != nil {
139+ return fmt.Errorf("save build state: %w", err)
140+ }
141+ return nil
142+}
143diff --git a/go.mod b/go.mod
144index ef929c55e85bd722009cac9ae288acc79fe5a8c0..b8ee982797380d724a456900e172bd7ca504a639 100644
145--- a/go.mod
146+++ b/go.mod
147@@ -1,3 +1,17 @@
148 module soft-worker
149
150 go 1.27
151+
152+require modernc.org/sqlite v1.57.0
153+
154+require (
155+ github.com/dustin/go-humanize v1.0.1 // indirect
156+ github.com/google/uuid v1.6.0 // indirect
157+ github.com/mattn/go-isatty v0.0.24 // indirect
158+ github.com/ncruces/go-strftime v1.0.0 // indirect
159+ github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect
160+ golang.org/x/sys v0.47.0 // indirect
161+ modernc.org/libc v1.74.4 // indirect
162+ modernc.org/mathutil v1.7.1 // indirect
163+ modernc.org/memory v1.11.0 // indirect
164+)
165diff --git a/go.sum b/go.sum
166new file mode 100644
167index 0000000000000000000000000000000000000000..3efffe8fe14ac5af05bc7c0f7fc5824af42b7a3a
168--- /dev/null
169+++ b/go.sum
170@@ -0,0 +1,50 @@
171+github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY=
172+github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto=
173+github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3 h1:LMLX+LgTNWpfvCBdFebv6EsYotImrt/Ppc5cXIriCSo=
174+github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3/go.mod h1:jl5iWTm0/hd5PjEYEOuwAJ57L/CibdZfrqZ5XA5GrCk=
175+github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
176+github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
177+github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k=
178+github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM=
179+github.com/mattn/go-isatty v0.0.24 h1:tGZZoVgT/KiqK1c8ocVLeDS8BSWMRd47J3Lbz7vsReI=
180+github.com/mattn/go-isatty v0.0.24/go.mod h1:nMCL3Zebbrt45jsMDgnfIwz6ydEQApk5oEI3HqDio6A=
181+github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w=
182+github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls=
183+github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE=
184+github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo=
185+golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ=
186+golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0=
187+golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM=
188+golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
189+golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
190+golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
191+golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q=
192+golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA=
193+modernc.org/cc/v4 v4.29.1 h1:MKgdCV3WykTSPqpVrnxdEDS0HEd2FHpKZDzxzU5LyeI=
194+modernc.org/cc/v4 v4.29.1/go.mod h1:OnovgIhbbMXMu1aISnJ0wvVD1KnW+cAUJkIrAWh+kVI=
195+modernc.org/ccgo/v4 v4.34.6 h1:sBgfIwyN0TQ9C5hwIeuqyeAKyMWnbvj2fvpF4L11uzU=
196+modernc.org/ccgo/v4 v4.34.6/go.mod h1:SZ8YcN9NG7XVsQYdm6jYBvi8PQP1qi+kqB6OhjqI3Fk=
197+modernc.org/fileutil v1.4.0 h1:j6ZzNTftVS054gi281TyLjHPp6CPHr2KCxEXjEbD6SM=
198+modernc.org/fileutil v1.4.0/go.mod h1:EqdKFDxiByqxLk8ozOxObDSfcVOv/54xDs/DUHdvCUU=
199+modernc.org/gc/v2 v2.6.5 h1:nyqdV8q46KvTpZlsw66kWqwXRHdjIlJOhG6kxiV/9xI=
200+modernc.org/gc/v2 v2.6.5/go.mod h1:YgIahr1ypgfe7chRuJi2gD7DBQiKSLMPgBQe9oIiito=
201+modernc.org/gc/v3 v3.1.4 h1:2g65LGVSmFQrXeITAw97x7hCRvZFcyE1uDP+7Vng7JI=
202+modernc.org/gc/v3 v3.1.4/go.mod h1:HFK/6AGESC7Ex+EZJhJ2Gni6cTaYpSMmU/cT9RmlfYY=
203+modernc.org/goabi0 v0.2.0 h1:HvEowk7LxcPd0eq6mVOAEMai46V+i7Jrj13t4AzuNks=
204+modernc.org/goabi0 v0.2.0/go.mod h1:CEFRnnJhKvWT1c1JTI3Avm+tgOWbkOu5oPA8eH8LnMI=
205+modernc.org/libc v1.74.4 h1:fX1Omw4o2/1C2iRkkIsrQTasJQldLhRmuPreXLoWs9k=
206+modernc.org/libc v1.74.4/go.mod h1:eeQAS9W3sZeKYMFubydxJpII9ybHWshk+7or7bLG9co=
207+modernc.org/mathutil v1.7.1 h1:GCZVGXdaN8gTqB1Mf/usp1Y/hSqgI2vAGGP4jZMCxOU=
208+modernc.org/mathutil v1.7.1/go.mod h1:4p5IwJITfppl0G4sUEDtCr4DthTaT47/N3aT6MhfgJg=
209+modernc.org/memory v1.11.0 h1:o4QC8aMQzmcwCK3t3Ux/ZHmwFPzE6hf2Y5LbkRs+hbI=
210+modernc.org/memory v1.11.0/go.mod h1:/JP4VbVC+K5sU2wZi9bHoq2MAkCnrt2r98UGeSK7Mjw=
211+modernc.org/opt v0.2.0 h1:tGyef5ApycA7FSEOMraay9SaTk5zmbx7Tu+cJs4QKZg=
212+modernc.org/opt v0.2.0/go.mod h1:03fq9lsNfvkYSfxrfUhZCWPk1lm4cq4N+Bh//bEtgns=
213+modernc.org/sortutil v1.2.1 h1:+xyoGf15mM3NMlPDnFqrteY07klSFxLElE2PVuWIJ7w=
214+modernc.org/sortutil v1.2.1/go.mod h1:7ZI3a3REbai7gzCLcotuw9AC4VZVpYMjDzETGsSMqJE=
215+modernc.org/sqlite v1.57.0 h1:qNQP6xnx5M0ISNtlnxoOX0+cD5bJ0/gr9aMmndFczzg=
216+modernc.org/sqlite v1.57.0/go.mod h1:yCJ2cmAaIkHQ25oXWrF8H4O1lIfPYPR26yCEDj2P3pQ=
217+modernc.org/strutil v1.2.1 h1:UneZBkQA+DX2Rp35KcM69cSsNES9ly8mQWD71HKlOA0=
218+modernc.org/strutil v1.2.1/go.mod h1:EHkiggD70koQxjVdSBM3JKM7k6L0FbGE5eymy9i3B9A=
219+modernc.org/token v1.1.0 h1:Xl7Ap9dKaEs5kLoOQeQmPWevfnk/DM5qcLcYlA8ys6Y=
220+modernc.org/token v1.1.0/go.mod h1:UGzOrNV1mAFSEB63lOFHIpNRUVMvYTc6yu1SMY/XTDM=
221diff --git a/http.go b/http.go
222index fdce67231949212b38ca46dc7ad91b0aa62a53e4..8f71ca0fbc38ace2dca02c41dc64bb47f072b4b3 100644
223--- a/http.go
224+++ b/http.go
225@@ -42,29 +42,41 @@ type buildStatusResponse struct {
226
227 type buildStatuses struct {
228 mu sync.RWMutex
229+ database *database
230 current map[string]buildStatus
231 buildIDs map[string]string
232 }
233
234-func newBuildStatuses() *buildStatuses {
235+func newBuildStatuses(database *database) *buildStatuses {
236 return &buildStatuses{
237+ database: database,
238 current: make(map[string]buildStatus),
239 buildIDs: make(map[string]string),
240 }
241 }
242
243-func (statuses *buildStatuses) queue(job buildJob) {
244+func (statuses *buildStatuses) queue(job buildJob) error {
245+ queuedAt := time.Now().UTC()
246+ if err := statuses.database.queue(job, queuedAt); err != nil {
247+ return err
248+ }
249+
250 statuses.mu.Lock()
251 defer statuses.mu.Unlock()
252 statuses.buildIDs[job.Request.Repository] = job.ID
253 statuses.current[job.Request.Repository] = buildStatus{
254 SHA: job.Request.SHA,
255 State: "queued",
256- QueuedAt: time.Now().UTC(),
257+ QueuedAt: queuedAt,
258 }
259+ return nil
260 }
261
262 func (statuses *buildStatuses) update(job buildJob, state string) bool {
263+ if err := statuses.database.update(job, state, time.Now().UTC()); err != nil {
264+ slog.Error("save build state", "build_id", job.ID, "state", state, "error", err)
265+ }
266+
267 statuses.mu.Lock()
268 defer statuses.mu.Unlock()
269 if statuses.buildIDs[job.Request.Repository] != job.ID {
270@@ -119,8 +131,13 @@ func buildHandler(jobs chan<- buildJob, statuses *buildStatuses) http.HandlerFun
271 }
272 select {
273 case jobs <- job:
274- statuses.queue(job)
275- close(job.ready)
276+ if err := statuses.queue(job); err != nil {
277+ job.ready <- false
278+ slog.Error("save queued build", "repository", buildRequest.Repository, "sha", buildRequest.SHA, "error", err)
279+ http.Error(writer, "unable to queue build", http.StatusInternalServerError)
280+ return
281+ }
282+ job.ready <- true
283 slog.Info("build queued", "repository", buildRequest.Repository, "sha", buildRequest.SHA)
284 writer.WriteHeader(http.StatusAccepted)
285 default:
286diff --git a/main.go b/main.go
287index e05bbd20ad46fe320151811252c4ee277edb5b9b..7516fc20c28fe33094d2ca9c6dbb88a62177c366 100644
288--- a/main.go
289+++ b/main.go
290@@ -15,14 +15,25 @@ func main() {
291 os.Exit(1)
292 }
293
294- statuses := newBuildStatuses()
295+ database, err := openDatabase(config.databasePath)
296+ if err != nil {
297+ slog.Error("open database", "error", err)
298+ os.Exit(1)
299+ }
300+ defer func() {
301+ if err := database.close(); err != nil {
302+ slog.Error("close database", "error", err)
303+ }
304+ }()
305+
306+ statuses := newBuildStatuses(database)
307 jobs := make(chan buildJob, 100)
308 go runWorker(config, jobs, statuses)
309
310 server := http.NewServeMux()
311 server.HandleFunc("POST /build", buildHandler(jobs, statuses))
312 server.HandleFunc("POST /build-statuses", buildStatusHandler(statuses))
313- slog.Info("worker started", "address", ":8080", "registry", config.registry, "buildkit_address", config.buildkitAddress)
314+ slog.Info("worker started", "address", ":8080", "registry", config.registry, "buildkit_address", config.buildkitAddress, "database_path", config.databasePath)
315 if err := http.ListenAndServe(":8080", server); err != nil {
316 slog.Error("serve HTTP", "error", err)
317 os.Exit(1)
318diff --git a/worker.go b/worker.go
319index 2de0d96b143920ea4a9f030ba50b0c261941e919..776bc588c04eec87bc5c0a7521a1dc2ebfb4ccab 100644
320--- a/worker.go
321+++ b/worker.go
322@@ -15,7 +15,7 @@ import (
323 type buildJob struct {
324 ID string
325 Request buildRequest
326- ready chan struct{}
327+ ready chan bool
328 }
329
330 func newBuildJob(request buildRequest) (buildJob, error) {
331@@ -23,7 +23,7 @@ func newBuildJob(request buildRequest) (buildJob, error) {
332 if err != nil {
333 return buildJob{}, err
334 }
335- return buildJob{ID: id, Request: request, ready: make(chan struct{})}, nil
336+ return buildJob{ID: id, Request: request, ready: make(chan bool, 1)}, nil
337 }
338
339 func newBuildID() (string, error) {
340@@ -36,7 +36,9 @@ func newBuildID() (string, error) {
341
342 func runWorker(config config, jobs <-chan buildJob, statuses *buildStatuses) {
343 for job := range jobs {
344- <-job.ready
345+ if !<-job.ready {
346+ continue
347+ }
348 request := job.Request
349 image := imageName(config.registry, request.Repository)
350 statuses.update(job, "running")