ci_mock_test.go
⎇
Raw
1package e2e
2
3import (
4 "archive/tar"
5 "bytes"
6 "encoding/binary"
7 "encoding/json"
8 "io"
9 "net"
10 "net/http"
11 "os"
12 "path"
13 "path/filepath"
14 "regexp"
15 "slices"
16 "strconv"
17 "strings"
18 "sync"
19 "testing"
20 "time"
21)
22
23// execResp is one programmed answer for a container exec.
24type execResp struct {
25 output string
26 exitCode int
27 // delay holds the response open, so a step timeout can fire.
28 delay time.Duration
29}
30
31// ciUpload is one recorded PUT /containers/*/archive.
32type ciUpload struct {
33 path string
34 body []byte
35}
36
37// ciVolume is one recorded POST /volumes/create.
38type ciVolume struct {
39 name string
40 labels map[string]string
41}
42
43// ciVolumeUsage is what GET /system/df reports for a volume.
44type ciVolumeUsage struct {
45 Size int64
46 RefCount int
47}
48
49// mockDocker is a fake Docker Engine API on a unix socket. It records what
50// the CI runner did and lets a test program the exec results.
51type mockDocker struct {
52 mu sync.Mutex
53 sock string
54
55 // uploadError makes every PUT /archive answer 500 with this text.
56 uploadError string
57 uploads []ciUpload
58 pulls []string
59 volumesCreated []ciVolume
60 volumesDeleted []string
61 // volumesOnHost is what GET /volumes reports.
62 volumesOnHost []string
63 // volumeUsage is what GET /system/df reports, by volume name.
64 volumeUsage map[string]ciVolumeUsage
65 // lastCreateBody is the body of the last POST /containers/create.
66 lastCreateBody map[string]any
67 // requests logs container creates, deletes and archive downloads.
68 requests []string
69 // archiveDelay holds every GET /archive open, so a cancel can land
70 // during artifact collection.
71 archiveDelay time.Duration
72 // execMap holds the response of an exec id, assigned at creation.
73 execMap map[string]execResp
74 // execQueue is consumed in order as execs are created.
75 execQueue []execResp
76 // execCmds is every command run inside a container, in order.
77 execCmds [][]string
78 execCounter int
79
80 // archives answers GET /archive for these paths instead of the
81 // one-file default.
82 archives map[string][]byte
83 // localImages is what GET /images/*/json finds.
84 localImages []string
85 // buildCreateBody is the body of the last build VM container create.
86 buildCreateBody map[string]any
87 // buildLog and buildExit are the build VM container's output and exit
88 // code.
89 buildLog string
90 buildExit int
91 // buildDelay holds the build VM log open, so a step timeout can fire.
92 buildDelay time.Duration
93 // imageBuilds records POST /build: the tag and the context file names.
94 imageBuilds []ciImageBuild
95 // imageBuildError makes POST /build report this error in its stream.
96 imageBuildError string
97}
98
99// ciImageBuild is one recorded POST /build.
100type ciImageBuild struct {
101 tag string
102 files []string
103}
104
105var (
106 reContainerStart = regexp.MustCompile(`/containers/[^/]+/start$`)
107 reContainerExec = regexp.MustCompile(`/containers/[^/]+/exec$`)
108 reExecStart = regexp.MustCompile(`/exec/([^/]+)/start$`)
109 reExecJSON = regexp.MustCompile(`/exec/([^/]+)/json$`)
110 reContainerArchive = regexp.MustCompile(`/containers/[^/]+/archive`)
111 reContainerLogs = regexp.MustCompile(`/containers/[^/]+/logs$`)
112 reContainerWait = regexp.MustCompile(`/containers/[^/]+/wait$`)
113 reImageInspect = regexp.MustCompile(`^/v1.47/images/(.+)/json$`)
114)
115
116// newMockDocker starts the mock on a unix socket. The socket lives in a short
117// temp path: a unix socket path is limited to about 104 bytes.
118func newMockDocker(t *testing.T) *mockDocker {
119 t.Helper()
120 dir, err := os.MkdirTemp("", "hf")
121 if err != nil {
122 t.Fatal(err)
123 }
124 m := &mockDocker{sock: filepath.Join(dir, "d.sock")}
125 m.reset()
126 ln, err := net.Listen("unix", m.sock)
127 if err != nil {
128 t.Fatal(err)
129 }
130 srv := &http.Server{Handler: m, ReadHeaderTimeout: 10 * time.Second}
131 go srv.Serve(ln)
132 t.Cleanup(func() {
133 srv.Close()
134 os.RemoveAll(dir)
135 })
136 return m
137}
138
139// queueExec programs the next exec that the runner creates.
140func (m *mockDocker) queueExec(r execResp) {
141 m.mu.Lock()
142 defer m.mu.Unlock()
143 m.execQueue = append(m.execQueue, r)
144}
145
146// reset clears every recording and programmed response.
147func (m *mockDocker) reset() {
148 m.mu.Lock()
149 defer m.mu.Unlock()
150 m.execMap = map[string]execResp{}
151 m.execQueue = nil
152 m.execCmds = nil
153 m.execCounter = 0
154 m.uploads = nil
155 m.uploadError = ""
156 m.pulls = nil
157 m.volumesCreated = nil
158 m.volumesDeleted = nil
159 m.volumesOnHost = nil
160 m.volumeUsage = map[string]ciVolumeUsage{}
161 m.lastCreateBody = nil
162 m.requests = nil
163 m.archiveDelay = 0
164 m.archives = map[string][]byte{}
165 m.localImages = nil
166 m.buildCreateBody = nil
167 m.buildLog = ""
168 m.buildExit = 0
169 m.buildDelay = 0
170 m.imageBuilds = nil
171 m.imageBuildError = ""
172}
173
174func (m *mockDocker) builtImages() []ciImageBuild {
175 m.mu.Lock()
176 defer m.mu.Unlock()
177 return append([]ciImageBuild(nil), m.imageBuilds...)
178}
179
180func (m *mockDocker) failImageBuild(msg string) {
181 m.mu.Lock()
182 defer m.mu.Unlock()
183 m.imageBuildError = msg
184}
185
186// programBuild sets what the next build VM container prints, exits with,
187// and leaves at /out/image.tar. A nil image leaves nothing.
188func (m *mockDocker) programBuild(log string, exit int, image []byte) {
189 m.mu.Lock()
190 defer m.mu.Unlock()
191 m.buildLog, m.buildExit = log, exit
192 if image != nil {
193 m.archives["/out/image.tar"] = makeTar("image.tar", string(image))
194 }
195}
196
197func (m *mockDocker) delayBuild(d time.Duration) {
198 m.mu.Lock()
199 defer m.mu.Unlock()
200 m.buildDelay = d
201}
202
203func (m *mockDocker) setArchive(path string, data []byte) {
204 m.mu.Lock()
205 defer m.mu.Unlock()
206 m.archives[path] = data
207}
208
209func (m *mockDocker) setLocalImages(names ...string) {
210 m.mu.Lock()
211 defer m.mu.Unlock()
212 m.localImages = names
213}
214
215func (m *mockDocker) buildBody() map[string]any {
216 m.mu.Lock()
217 defer m.mu.Unlock()
218 return m.buildCreateBody
219}
220
221// uploadTo returns the last archive uploaded to path.
222func (m *mockDocker) uploadTo(path string) []byte {
223 m.mu.Lock()
224 defer m.mu.Unlock()
225 for i := len(m.uploads) - 1; i >= 0; i-- {
226 if m.uploads[i].path == path {
227 return m.uploads[i].body
228 }
229 }
230 return nil
231}
232
233func (m *mockDocker) containerRequests() []string {
234 m.mu.Lock()
235 defer m.mu.Unlock()
236 return append([]string(nil), m.requests...)
237}
238
239func (m *mockDocker) setVolumesOnHost(names ...string) {
240 m.mu.Lock()
241 defer m.mu.Unlock()
242 m.volumesOnHost = names
243}
244
245func (m *mockDocker) setVolumeUsage(usage map[string]ciVolumeUsage) {
246 m.mu.Lock()
247 defer m.mu.Unlock()
248 m.volumeUsage = usage
249}
250
251func (m *mockDocker) uploadedPaths() []string {
252 m.mu.Lock()
253 defer m.mu.Unlock()
254 out := make([]string, 0, len(m.uploads))
255 for _, u := range m.uploads {
256 out = append(out, u.path)
257 }
258 return out
259}
260
261// uploadsInto lists uploads whose tar holds at least one entry under dir.
262// Directories are created by the upload itself, so the PUT path is "/" and
263// the interesting part is inside the archive.
264func (m *mockDocker) uploadsInto(t *testing.T, dir string) []ciUpload {
265 t.Helper()
266 prefix := strings.TrimPrefix(dir, "/") + "/"
267 m.mu.Lock()
268 defer m.mu.Unlock()
269 var out []ciUpload
270 for _, u := range m.uploads {
271 for _, h := range tarHeaders(t, u.body) {
272 if strings.HasPrefix(h.name, prefix) && h.name != prefix {
273 out = append(out, u)
274 break
275 }
276 }
277 }
278 return out
279}
280
281func (m *mockDocker) uploadCount() int {
282 m.mu.Lock()
283 defer m.mu.Unlock()
284 return len(m.uploads)
285}
286
287func (m *mockDocker) pulledImages() []string {
288 m.mu.Lock()
289 defer m.mu.Unlock()
290 return append([]string(nil), m.pulls...)
291}
292
293func (m *mockDocker) createdVolumes() []ciVolume {
294 m.mu.Lock()
295 defer m.mu.Unlock()
296 return append([]ciVolume(nil), m.volumesCreated...)
297}
298
299func (m *mockDocker) deletedVolumes() []string {
300 m.mu.Lock()
301 defer m.mu.Unlock()
302 return append([]string(nil), m.volumesDeleted...)
303}
304
305func (m *mockDocker) createBody() map[string]any {
306 m.mu.Lock()
307 defer m.mu.Unlock()
308 return m.lastCreateBody
309}
310
311func (m *mockDocker) commands() [][]string {
312 m.mu.Lock()
313 defer m.mu.Unlock()
314 return append([][]string(nil), m.execCmds...)
315}
316
317// allCommandText joins every command the runner ran, for a substring check.
318func (m *mockDocker) allCommandText() string {
319 var parts []string
320 for _, c := range m.commands() {
321 parts = append(parts, strings.Join(c, " "))
322 }
323 return strings.Join(parts, " ")
324}
325
326func writeJSON(w http.ResponseWriter, v any) {
327 w.Header().Set("Content-Type", "application/json")
328 _ = json.NewEncoder(w).Encode(v)
329}
330
331func (m *mockDocker) ServeHTTP(w http.ResponseWriter, r *http.Request) {
332 p := r.URL.Path
333 qs := r.URL.Query()
334
335 switch {
336 // Health check.
337 case r.Method == http.MethodGet && p == "/v1.47/info":
338 writeJSON(w, map[string]string{"ServerVersion": "mock"})
339
340 // Inspect image.
341 case r.Method == http.MethodGet && reImageInspect.MatchString(p):
342 m.mu.Lock()
343 found := slices.Contains(m.localImages, reImageInspect.FindStringSubmatch(p)[1])
344 m.mu.Unlock()
345 if !found {
346 http.Error(w, "no such image", http.StatusNotFound)
347 return
348 }
349 writeJSON(w, map[string]string{"Id": "sha256:mock"})
350
351 // Build an image. A success makes it local, like the real engine.
352 case r.Method == http.MethodPost && p == "/v1.47/build":
353 body, _ := io.ReadAll(r.Body)
354 tag := qs.Get("t")
355 m.mu.Lock()
356 m.imageBuilds = append(m.imageBuilds, ciImageBuild{tag: tag, files: tarNames(body)})
357 fail := m.imageBuildError
358 if fail == "" {
359 m.localImages = append(m.localImages, tag)
360 }
361 m.mu.Unlock()
362 writeJSON(w, map[string]string{"stream": "Step 1/20 : ARG ALPINE\n"})
363 if fail != "" {
364 writeJSON(w, map[string]string{"error": fail})
365 }
366
367 // Pull image.
368 case r.Method == http.MethodPost && strings.HasPrefix(p, "/v1.47/images/create"):
369 m.mu.Lock()
370 m.pulls = append(m.pulls, qs.Get("fromImage"))
371 m.mu.Unlock()
372 _, _ = w.Write([]byte("{\"status\":\"Pull complete\"}\n"))
373
374 // Create container.
375 case r.Method == http.MethodPost && strings.Contains(p, "/containers/create"):
376 var body map[string]any
377 _ = json.NewDecoder(r.Body).Decode(&body)
378 name := qs.Get("name")
379 if name == "" {
380 name = "mock-ctr-001"
381 }
382 m.mu.Lock()
383 m.requests = append(m.requests, "POST "+path.Base(p)+"?name="+qs.Get("name"))
384 m.mu.Unlock()
385 // A copy creates its own source container, so the run's container
386 // must keep its identity.
387 switch {
388 case strings.HasSuffix(name, "-build"):
389 m.mu.Lock()
390 m.buildCreateBody = body
391 m.mu.Unlock()
392 case !strings.Contains(name, "-copy-"):
393 m.mu.Lock()
394 m.lastCreateBody = body
395 m.mu.Unlock()
396 }
397 writeJSON(w, map[string]string{"Id": name})
398
399 // Build VM container output.
400 case r.Method == http.MethodGet && reContainerLogs.MatchString(p):
401 m.mu.Lock()
402 out, delay := m.buildLog, m.buildDelay
403 m.mu.Unlock()
404 if out != "" {
405 _, _ = w.Write(muxFrame(out))
406 }
407 if delay > 0 {
408 w.(http.Flusher).Flush()
409 select {
410 case <-time.After(delay):
411 case <-r.Context().Done():
412 }
413 }
414
415 case r.Method == http.MethodPost && reContainerWait.MatchString(p):
416 m.mu.Lock()
417 code := m.buildExit
418 m.mu.Unlock()
419 writeJSON(w, map[string]int{"StatusCode": code})
420
421 // Start container.
422 case r.Method == http.MethodPost && reContainerStart.MatchString(p):
423 w.WriteHeader(http.StatusNoContent)
424
425 // Create exec: take the next queued response and bind it to this id.
426 case r.Method == http.MethodPost && reContainerExec.MatchString(p):
427 var body struct {
428 Cmd []string
429 }
430 _ = json.NewDecoder(r.Body).Decode(&body)
431 m.mu.Lock()
432 m.execCounter++
433 id := "mock-exec-" + strconv.Itoa(m.execCounter)
434 m.execCmds = append(m.execCmds, body.Cmd)
435 resp := execResp{}
436 if len(m.execQueue) > 0 {
437 resp, m.execQueue = m.execQueue[0], m.execQueue[1:]
438 }
439 m.execMap[id] = resp
440 m.mu.Unlock()
441 writeJSON(w, map[string]string{"Id": id})
442
443 // Start exec: return the programmed output as a mux stream.
444 case r.Method == http.MethodPost && reExecStart.MatchString(p):
445 resp := m.execFor(reExecStart.FindStringSubmatch(p)[1])
446 if resp.delay > 0 {
447 select {
448 case <-time.After(resp.delay):
449 case <-r.Context().Done():
450 return
451 }
452 }
453 if resp.output != "" {
454 _, _ = w.Write(muxFrame(resp.output))
455 }
456
457 // Inspect exec: return the exit code.
458 case r.Method == http.MethodGet && reExecJSON.MatchString(p):
459 resp := m.execFor(reExecJSON.FindStringSubmatch(p)[1])
460 writeJSON(w, map[string]int{"ExitCode": resp.exitCode})
461
462 // Archive upload: the checkout and the [[copy]] sources.
463 case r.Method == http.MethodPut && reContainerArchive.MatchString(p):
464 body, _ := io.ReadAll(r.Body)
465 m.mu.Lock()
466 m.uploads = append(m.uploads, ciUpload{path: qs.Get("path"), body: body})
467 fail := m.uploadError
468 m.mu.Unlock()
469 if fail != "" {
470 http.Error(w, fail, http.StatusInternalServerError)
471 return
472 }
473 w.WriteHeader(http.StatusOK)
474
475 // Archive download: artifact collection and [[copy]].
476 case r.Method == http.MethodGet && reContainerArchive.MatchString(p):
477 m.mu.Lock()
478 delay := m.archiveDelay
479 m.requests = append(m.requests, "GET archive")
480 m.mu.Unlock()
481 if delay > 0 {
482 select {
483 case <-time.After(delay):
484 case <-r.Context().Done():
485 return
486 }
487 }
488 filePath := qs.Get("path")
489 m.mu.Lock()
490 programmed, ok := m.archives[filePath]
491 m.mu.Unlock()
492 if ok {
493 w.Header().Set("Content-Type", "application/x-tar")
494 _, _ = w.Write(programmed)
495 return
496 }
497 if filePath == "" {
498 filePath = "file.txt"
499 }
500 w.Header().Set("Content-Type", "application/x-tar")
501 _, _ = w.Write(makeTar(path.Base(filePath), "artifact-content-123"))
502
503 // Delete container.
504 case r.Method == http.MethodDelete && strings.Contains(p, "/containers/"):
505 m.mu.Lock()
506 m.requests = append(m.requests, "DELETE "+path.Base(p))
507 m.mu.Unlock()
508 w.WriteHeader(http.StatusNoContent)
509
510 // Volume create, used for cache volumes.
511 case r.Method == http.MethodPost && p == "/v1.47/volumes/create":
512 var body struct {
513 Name string
514 Labels map[string]string
515 }
516 _ = json.NewDecoder(r.Body).Decode(&body)
517 if body.Labels == nil {
518 body.Labels = map[string]string{}
519 }
520 m.mu.Lock()
521 m.volumesCreated = append(m.volumesCreated, ciVolume{name: body.Name, labels: body.Labels})
522 m.mu.Unlock()
523 writeJSON(w, map[string]string{"Name": body.Name})
524
525 // Disk usage, used by the cache size caps.
526 case r.Method == http.MethodGet && p == "/v1.47/system/df":
527 type entry struct {
528 Name string
529 UsageData ciVolumeUsage
530 }
531 m.mu.Lock()
532 vols := make([]entry, 0, len(m.volumeUsage))
533 for name, usage := range m.volumeUsage {
534 vols = append(vols, entry{Name: name, UsageData: usage})
535 }
536 m.mu.Unlock()
537 writeJSON(w, map[string]any{"Volumes": vols})
538
539 // Volume list, used by prune and purge.
540 case r.Method == http.MethodGet && p == "/v1.47/volumes":
541 type entry struct{ Name string }
542 m.mu.Lock()
543 vols := make([]entry, 0, len(m.volumesOnHost))
544 for _, name := range m.volumesOnHost {
545 vols = append(vols, entry{Name: name})
546 }
547 m.mu.Unlock()
548 writeJSON(w, map[string]any{"Volumes": vols})
549
550 // Volume delete.
551 case r.Method == http.MethodDelete && strings.Contains(p, "/volumes/"):
552 m.mu.Lock()
553 m.volumesDeleted = append(m.volumesDeleted, path.Base(p))
554 m.mu.Unlock()
555 w.WriteHeader(http.StatusNoContent)
556
557 default:
558 http.Error(w, "Not found", http.StatusNotFound)
559 }
560}
561
562func (m *mockDocker) execFor(id string) execResp {
563 m.mu.Lock()
564 defer m.mu.Unlock()
565 return m.execMap[id]
566}
567
568// muxFrame builds one Docker multiplexed stream frame.
569func muxFrame(text string) []byte {
570 hdr := make([]byte, 8)
571 hdr[0] = 1
572 binary.BigEndian.PutUint32(hdr[4:], uint32(len(text)))
573 return append(hdr, text...)
574}
575
576// makeTar builds a tar archive holding one file.
577func makeTar(filename, content string) []byte {
578 var buf bytes.Buffer
579 tw := tar.NewWriter(&buf)
580 _ = tw.WriteHeader(&tar.Header{
581 Name: filename, Mode: 0o644, Size: int64(len(content)), Typeflag: tar.TypeReg,
582 })
583 _, _ = tw.Write([]byte(content))
584 _ = tw.Close()
585 return buf.Bytes()
586}
587
588// tarEntry is one header of an uploaded archive.
589type tarEntry struct {
590 name string
591 uid int
592}
593
594// tarHeaders lists the entries of an uncompressed tar.
595func tarHeaders(t *testing.T, data []byte) []tarEntry {
596 t.Helper()
597 tr := tar.NewReader(bytes.NewReader(data))
598 var out []tarEntry
599 for {
600 h, err := tr.Next()
601 if err == io.EOF {
602 break
603 }
604 if err != nil {
605 t.Fatalf("read tar: %v", err)
606 }
607 out = append(out, tarEntry{name: h.Name, uid: h.Uid})
608 }
609 return out
610}
611
612// tarNames lists the entry names of a tar, as far as it reads.
613func tarNames(data []byte) []string {
614 var out []string
615 tr := tar.NewReader(bytes.NewReader(data))
616 for {
617 h, err := tr.Next()
618 if err != nil {
619 return out
620 }
621 out = append(out, h.Name)
622 }
623}
624
625// tarEntryNames lists the file names of an uncompressed tar.
626func tarEntryNames(t *testing.T, data []byte) []string {
627 t.Helper()
628 var out []string
629 for _, h := range tarHeaders(t, data) {
630 out = append(out, h.name)
631 }
632 return out
633}
634