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