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 "strconv"
16 "strings"
17 "sync"
18 "testing"
19 "time"
20)
21
22// execResp is one programmed answer for a container exec.
23type execResp struct {
24 output string
25 exitCode int
26 // delay holds the response open, so a step timeout can fire.
27 delay time.Duration
28}
29
30// ciUpload is one recorded PUT /containers/*/archive.
31type ciUpload struct {
32 path string
33 body []byte
34}
35
36// ciVolume is one recorded POST /volumes/create.
37type ciVolume struct {
38 name string
39 labels map[string]string
40}
41
42// ciVolumeUsage is what GET /system/df reports for a volume.
43type ciVolumeUsage struct {
44 Size int64
45 RefCount int
46}
47
48// mockDocker is a fake Docker Engine API on a unix socket. It records what
49// the CI runner did and lets a test program the exec results.
50type mockDocker struct {
51 mu sync.Mutex
52 sock string
53
54 // uploadError makes every PUT /archive answer 500 with this text.
55 uploadError string
56 uploads []ciUpload
57 pulls []string
58 volumesCreated []ciVolume
59 volumesDeleted []string
60 // volumesOnHost is what GET /volumes reports.
61 volumesOnHost []string
62 // volumeUsage is what GET /system/df reports, by volume name.
63 volumeUsage map[string]ciVolumeUsage
64 // lastCreateBody is the body of the last POST /containers/create.
65 lastCreateBody map[string]any
66 // requests logs container creates, deletes and archive downloads.
67 requests []string
68 // archiveDelay holds every GET /archive open, so a cancel can land
69 // during artifact collection.
70 archiveDelay time.Duration
71 // execMap holds the response of an exec id, assigned at creation.
72 execMap map[string]execResp
73 // execQueue is consumed in order as execs are created.
74 execQueue []execResp
75 // execCmds is every command run inside a container, in order.
76 execCmds [][]string
77 // execEnvs is the environment of every exec, aligned with execCmds.
78 execEnvs [][]string
79 execCounter int
80}
81
82var (
83 reContainerStart = regexp.MustCompile(`/containers/[^/]+/start$`)
84 reContainerExec = regexp.MustCompile(`/containers/[^/]+/exec$`)
85 reExecStart = regexp.MustCompile(`/exec/([^/]+)/start$`)
86 reExecJSON = regexp.MustCompile(`/exec/([^/]+)/json$`)
87 reContainerArchive = regexp.MustCompile(`/containers/[^/]+/archive`)
88)
89
90// newMockDocker starts the mock on a unix socket. The socket lives in a short
91// temp path: a unix socket path is limited to about 104 bytes.
92func newMockDocker(t *testing.T) *mockDocker {
93 t.Helper()
94 dir, err := os.MkdirTemp("", "hf")
95 if err != nil {
96 t.Fatal(err)
97 }
98 m := &mockDocker{sock: filepath.Join(dir, "d.sock")}
99 m.reset()
100 ln, err := net.Listen("unix", m.sock)
101 if err != nil {
102 t.Fatal(err)
103 }
104 srv := &http.Server{Handler: m, ReadHeaderTimeout: 10 * time.Second}
105 go srv.Serve(ln)
106 t.Cleanup(func() {
107 srv.Close()
108 os.RemoveAll(dir)
109 })
110 return m
111}
112
113// queueExec programs the next exec that the runner creates.
114func (m *mockDocker) queueExec(r execResp) {
115 m.mu.Lock()
116 defer m.mu.Unlock()
117 m.execQueue = append(m.execQueue, r)
118}
119
120// reset clears every recording and programmed response.
121func (m *mockDocker) reset() {
122 m.mu.Lock()
123 defer m.mu.Unlock()
124 m.execMap = map[string]execResp{}
125 m.execQueue = nil
126 m.execCmds = nil
127 m.execEnvs = nil
128 m.execCounter = 0
129 m.uploads = nil
130 m.uploadError = ""
131 m.pulls = nil
132 m.volumesCreated = nil
133 m.volumesDeleted = nil
134 m.volumesOnHost = nil
135 m.volumeUsage = map[string]ciVolumeUsage{}
136 m.lastCreateBody = nil
137 m.requests = nil
138 m.archiveDelay = 0
139}
140
141func (m *mockDocker) containerRequests() []string {
142 m.mu.Lock()
143 defer m.mu.Unlock()
144 return append([]string(nil), m.requests...)
145}
146
147func (m *mockDocker) setVolumesOnHost(names ...string) {
148 m.mu.Lock()
149 defer m.mu.Unlock()
150 m.volumesOnHost = names
151}
152
153func (m *mockDocker) setVolumeUsage(usage map[string]ciVolumeUsage) {
154 m.mu.Lock()
155 defer m.mu.Unlock()
156 m.volumeUsage = usage
157}
158
159func (m *mockDocker) uploadedPaths() []string {
160 m.mu.Lock()
161 defer m.mu.Unlock()
162 out := make([]string, 0, len(m.uploads))
163 for _, u := range m.uploads {
164 out = append(out, u.path)
165 }
166 return out
167}
168
169// uploadsInto lists uploads whose tar holds at least one entry under dir.
170// Directories are created by the upload itself, so the PUT path is "/" and
171// the interesting part is inside the archive.
172func (m *mockDocker) uploadsInto(t *testing.T, dir string) []ciUpload {
173 t.Helper()
174 prefix := strings.TrimPrefix(dir, "/") + "/"
175 m.mu.Lock()
176 defer m.mu.Unlock()
177 var out []ciUpload
178 for _, u := range m.uploads {
179 for _, h := range tarHeaders(t, u.body) {
180 if strings.HasPrefix(h.name, prefix) && h.name != prefix {
181 out = append(out, u)
182 break
183 }
184 }
185 }
186 return out
187}
188
189func (m *mockDocker) uploadCount() int {
190 m.mu.Lock()
191 defer m.mu.Unlock()
192 return len(m.uploads)
193}
194
195func (m *mockDocker) pulledImages() []string {
196 m.mu.Lock()
197 defer m.mu.Unlock()
198 return append([]string(nil), m.pulls...)
199}
200
201func (m *mockDocker) createdVolumes() []ciVolume {
202 m.mu.Lock()
203 defer m.mu.Unlock()
204 return append([]ciVolume(nil), m.volumesCreated...)
205}
206
207func (m *mockDocker) deletedVolumes() []string {
208 m.mu.Lock()
209 defer m.mu.Unlock()
210 return append([]string(nil), m.volumesDeleted...)
211}
212
213func (m *mockDocker) createBody() map[string]any {
214 m.mu.Lock()
215 defer m.mu.Unlock()
216 return m.lastCreateBody
217}
218
219// envForCommand returns the environment of the first exec whose command
220// contains sub, or nil when no command matches.
221func (m *mockDocker) envForCommand(sub string) []string {
222 m.mu.Lock()
223 defer m.mu.Unlock()
224 for i, c := range m.execCmds {
225 if strings.Contains(strings.Join(c, " "), sub) {
226 return m.execEnvs[i]
227 }
228 }
229 return nil
230}
231
232func (m *mockDocker) commands() [][]string {
233 m.mu.Lock()
234 defer m.mu.Unlock()
235 return append([][]string(nil), m.execCmds...)
236}
237
238// allCommandText joins every command the runner ran, for a substring check.
239func (m *mockDocker) allCommandText() string {
240 var parts []string
241 for _, c := range m.commands() {
242 parts = append(parts, strings.Join(c, " "))
243 }
244 return strings.Join(parts, " ")
245}
246
247func writeJSON(w http.ResponseWriter, v any) {
248 w.Header().Set("Content-Type", "application/json")
249 _ = json.NewEncoder(w).Encode(v)
250}
251
252func (m *mockDocker) ServeHTTP(w http.ResponseWriter, r *http.Request) {
253 p := r.URL.Path
254 qs := r.URL.Query()
255
256 switch {
257 // Health check.
258 case r.Method == http.MethodGet && p == "/v1.47/info":
259 writeJSON(w, map[string]string{"ServerVersion": "mock"})
260
261 // Pull image.
262 case r.Method == http.MethodPost && strings.HasPrefix(p, "/v1.47/images/create"):
263 m.mu.Lock()
264 m.pulls = append(m.pulls, qs.Get("fromImage"))
265 m.mu.Unlock()
266 _, _ = w.Write([]byte("{\"status\":\"Pull complete\"}\n"))
267
268 // Create container.
269 case r.Method == http.MethodPost && strings.Contains(p, "/containers/create"):
270 var body map[string]any
271 _ = json.NewDecoder(r.Body).Decode(&body)
272 name := qs.Get("name")
273 if name == "" {
274 name = "mock-ctr-001"
275 }
276 m.mu.Lock()
277 m.requests = append(m.requests, "POST "+path.Base(p)+"?name="+qs.Get("name"))
278 m.mu.Unlock()
279 // A copy creates its own source container, so the run's container
280 // must keep its identity.
281 if !strings.Contains(name, "-copy-") {
282 m.mu.Lock()
283 m.lastCreateBody = body
284 m.mu.Unlock()
285 }
286 writeJSON(w, map[string]string{"Id": name})
287
288 // Start container.
289 case r.Method == http.MethodPost && reContainerStart.MatchString(p):
290 w.WriteHeader(http.StatusNoContent)
291
292 // Create exec: take the next queued response and bind it to this id.
293 case r.Method == http.MethodPost && reContainerExec.MatchString(p):
294 var body struct {
295 Cmd []string
296 Env []string
297 }
298 _ = json.NewDecoder(r.Body).Decode(&body)
299 m.mu.Lock()
300 m.execCounter++
301 id := "mock-exec-" + strconv.Itoa(m.execCounter)
302 m.execCmds = append(m.execCmds, body.Cmd)
303 m.execEnvs = append(m.execEnvs, body.Env)
304 resp := execResp{}
305 if len(m.execQueue) > 0 {
306 resp, m.execQueue = m.execQueue[0], m.execQueue[1:]
307 }
308 m.execMap[id] = resp
309 m.mu.Unlock()
310 writeJSON(w, map[string]string{"Id": id})
311
312 // Start exec: return the programmed output as a mux stream.
313 case r.Method == http.MethodPost && reExecStart.MatchString(p):
314 resp := m.execFor(reExecStart.FindStringSubmatch(p)[1])
315 if resp.delay > 0 {
316 select {
317 case <-time.After(resp.delay):
318 case <-r.Context().Done():
319 return
320 }
321 }
322 if resp.output != "" {
323 _, _ = w.Write(muxFrame(resp.output))
324 }
325
326 // Inspect exec: return the exit code.
327 case r.Method == http.MethodGet && reExecJSON.MatchString(p):
328 resp := m.execFor(reExecJSON.FindStringSubmatch(p)[1])
329 writeJSON(w, map[string]int{"ExitCode": resp.exitCode})
330
331 // Archive upload: the checkout and the [[copy]] sources.
332 case r.Method == http.MethodPut && reContainerArchive.MatchString(p):
333 body, _ := io.ReadAll(r.Body)
334 m.mu.Lock()
335 m.uploads = append(m.uploads, ciUpload{path: qs.Get("path"), body: body})
336 fail := m.uploadError
337 m.mu.Unlock()
338 if fail != "" {
339 http.Error(w, fail, http.StatusInternalServerError)
340 return
341 }
342 w.WriteHeader(http.StatusOK)
343
344 // Archive download: artifact collection and [[copy]].
345 case r.Method == http.MethodGet && reContainerArchive.MatchString(p):
346 m.mu.Lock()
347 delay := m.archiveDelay
348 m.requests = append(m.requests, "GET archive")
349 m.mu.Unlock()
350 if delay > 0 {
351 select {
352 case <-time.After(delay):
353 case <-r.Context().Done():
354 return
355 }
356 }
357 filePath := qs.Get("path")
358 if filePath == "" {
359 filePath = "file.txt"
360 }
361 w.Header().Set("Content-Type", "application/x-tar")
362 _, _ = w.Write(makeTar(path.Base(filePath), "artifact-content-123"))
363
364 // Delete container.
365 case r.Method == http.MethodDelete && strings.Contains(p, "/containers/"):
366 m.mu.Lock()
367 m.requests = append(m.requests, "DELETE "+path.Base(p))
368 m.mu.Unlock()
369 w.WriteHeader(http.StatusNoContent)
370
371 // Volume create, used for cache volumes.
372 case r.Method == http.MethodPost && p == "/v1.47/volumes/create":
373 var body struct {
374 Name string
375 Labels map[string]string
376 }
377 _ = json.NewDecoder(r.Body).Decode(&body)
378 if body.Labels == nil {
379 body.Labels = map[string]string{}
380 }
381 m.mu.Lock()
382 m.volumesCreated = append(m.volumesCreated, ciVolume{name: body.Name, labels: body.Labels})
383 m.mu.Unlock()
384 writeJSON(w, map[string]string{"Name": body.Name})
385
386 // Disk usage, used by the cache size caps.
387 case r.Method == http.MethodGet && p == "/v1.47/system/df":
388 type entry struct {
389 Name string
390 UsageData ciVolumeUsage
391 }
392 m.mu.Lock()
393 vols := make([]entry, 0, len(m.volumeUsage))
394 for name, usage := range m.volumeUsage {
395 vols = append(vols, entry{Name: name, UsageData: usage})
396 }
397 m.mu.Unlock()
398 writeJSON(w, map[string]any{"Volumes": vols})
399
400 // Volume list, used by prune and purge.
401 case r.Method == http.MethodGet && p == "/v1.47/volumes":
402 type entry struct{ Name string }
403 m.mu.Lock()
404 vols := make([]entry, 0, len(m.volumesOnHost))
405 for _, name := range m.volumesOnHost {
406 vols = append(vols, entry{Name: name})
407 }
408 m.mu.Unlock()
409 writeJSON(w, map[string]any{"Volumes": vols})
410
411 // Volume delete.
412 case r.Method == http.MethodDelete && strings.Contains(p, "/volumes/"):
413 m.mu.Lock()
414 m.volumesDeleted = append(m.volumesDeleted, path.Base(p))
415 m.mu.Unlock()
416 w.WriteHeader(http.StatusNoContent)
417
418 default:
419 http.Error(w, "Not found", http.StatusNotFound)
420 }
421}
422
423func (m *mockDocker) execFor(id string) execResp {
424 m.mu.Lock()
425 defer m.mu.Unlock()
426 return m.execMap[id]
427}
428
429// muxFrame builds one Docker multiplexed stream frame.
430func muxFrame(text string) []byte {
431 hdr := make([]byte, 8)
432 hdr[0] = 1
433 binary.BigEndian.PutUint32(hdr[4:], uint32(len(text)))
434 return append(hdr, text...)
435}
436
437// makeTar builds a tar archive holding one file.
438func makeTar(filename, content string) []byte {
439 var buf bytes.Buffer
440 tw := tar.NewWriter(&buf)
441 _ = tw.WriteHeader(&tar.Header{
442 Name: filename, Mode: 0o644, Size: int64(len(content)), Typeflag: tar.TypeReg,
443 })
444 _, _ = tw.Write([]byte(content))
445 _ = tw.Close()
446 return buf.Bytes()
447}
448
449// tarEntry is one header of an uploaded archive.
450type tarEntry struct {
451 name string
452 uid int
453}
454
455// tarHeaders lists the entries of an uncompressed tar.
456func tarHeaders(t *testing.T, data []byte) []tarEntry {
457 t.Helper()
458 tr := tar.NewReader(bytes.NewReader(data))
459 var out []tarEntry
460 for {
461 h, err := tr.Next()
462 if err == io.EOF {
463 break
464 }
465 if err != nil {
466 t.Fatalf("read tar: %v", err)
467 }
468 out = append(out, tarEntry{name: h.Name, uid: h.Uid})
469 }
470 return out
471}
472
473// tarEntryNames lists the file names of an uncompressed tar.
474func tarEntryNames(t *testing.T, data []byte) []string {
475 t.Helper()
476 var out []string
477 for _, h := range tarHeaders(t, data) {
478 out = append(out, h.name)
479 }
480 return out
481}
482