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