package e2e import ( "archive/tar" "bytes" "encoding/binary" "encoding/json" "io" "net" "net/http" "os" "path" "path/filepath" "regexp" "strconv" "strings" "sync" "testing" "time" ) // execResp is one programmed answer for a container exec. type execResp struct { output string exitCode int // delay holds the response open, so a step timeout can fire. delay time.Duration } // ciUpload is one recorded PUT /containers/*/archive. type ciUpload struct { path string body []byte } // ciVolume is one recorded POST /volumes/create. type ciVolume struct { name string labels map[string]string } // ciVolumeUsage is what GET /system/df reports for a volume. type ciVolumeUsage struct { Size int64 RefCount int } // mockDocker is a fake Docker Engine API on a unix socket. It records what // the CI runner did and lets a test program the exec results. type mockDocker struct { mu sync.Mutex sock string // uploadError makes every PUT /archive answer 500 with this text. uploadError string uploads []ciUpload pulls []string volumesCreated []ciVolume volumesDeleted []string // volumesOnHost is what GET /volumes reports. volumesOnHost []string // volumeUsage is what GET /system/df reports, by volume name. volumeUsage map[string]ciVolumeUsage // lastCreateBody is the body of the last POST /containers/create. lastCreateBody map[string]any // execMap holds the response of an exec id, assigned at creation. execMap map[string]execResp // execQueue is consumed in order as execs are created. execQueue []execResp // execCmds is every command run inside a container, in order. execCmds [][]string // execEnvs is the environment of every exec, aligned with execCmds. execEnvs [][]string execCounter int } var ( reContainerStart = regexp.MustCompile(`/containers/[^/]+/start$`) reContainerExec = regexp.MustCompile(`/containers/[^/]+/exec$`) reExecStart = regexp.MustCompile(`/exec/([^/]+)/start$`) reExecJSON = regexp.MustCompile(`/exec/([^/]+)/json$`) reContainerArchive = regexp.MustCompile(`/containers/[^/]+/archive`) ) // newMockDocker starts the mock on a unix socket. The socket lives in a short // temp path: a unix socket path is limited to about 104 bytes. func newMockDocker(t *testing.T) *mockDocker { t.Helper() dir, err := os.MkdirTemp("", "hf") if err != nil { t.Fatal(err) } m := &mockDocker{sock: filepath.Join(dir, "d.sock")} m.reset() ln, err := net.Listen("unix", m.sock) if err != nil { t.Fatal(err) } srv := &http.Server{Handler: m, ReadHeaderTimeout: 10 * time.Second} go srv.Serve(ln) t.Cleanup(func() { srv.Close() os.RemoveAll(dir) }) return m } // queueExec programs the next exec that the runner creates. func (m *mockDocker) queueExec(r execResp) { m.mu.Lock() defer m.mu.Unlock() m.execQueue = append(m.execQueue, r) } // reset clears every recording and programmed response. func (m *mockDocker) reset() { m.mu.Lock() defer m.mu.Unlock() m.execMap = map[string]execResp{} m.execQueue = nil m.execCmds = nil m.execEnvs = nil m.execCounter = 0 m.uploads = nil m.uploadError = "" m.pulls = nil m.volumesCreated = nil m.volumesDeleted = nil m.volumesOnHost = nil m.volumeUsage = map[string]ciVolumeUsage{} m.lastCreateBody = nil } func (m *mockDocker) setVolumesOnHost(names ...string) { m.mu.Lock() defer m.mu.Unlock() m.volumesOnHost = names } func (m *mockDocker) setVolumeUsage(usage map[string]ciVolumeUsage) { m.mu.Lock() defer m.mu.Unlock() m.volumeUsage = usage } func (m *mockDocker) uploadedPaths() []string { m.mu.Lock() defer m.mu.Unlock() out := make([]string, 0, len(m.uploads)) for _, u := range m.uploads { out = append(out, u.path) } return out } // uploadsInto lists uploads whose tar holds at least one entry under dir. // Directories are created by the upload itself, so the PUT path is "/" and // the interesting part is inside the archive. func (m *mockDocker) uploadsInto(t *testing.T, dir string) []ciUpload { t.Helper() prefix := strings.TrimPrefix(dir, "/") + "/" m.mu.Lock() defer m.mu.Unlock() var out []ciUpload for _, u := range m.uploads { for _, h := range tarHeaders(t, u.body) { if strings.HasPrefix(h.name, prefix) && h.name != prefix { out = append(out, u) break } } } return out } func (m *mockDocker) uploadCount() int { m.mu.Lock() defer m.mu.Unlock() return len(m.uploads) } func (m *mockDocker) pulledImages() []string { m.mu.Lock() defer m.mu.Unlock() return append([]string(nil), m.pulls...) } func (m *mockDocker) createdVolumes() []ciVolume { m.mu.Lock() defer m.mu.Unlock() return append([]ciVolume(nil), m.volumesCreated...) } func (m *mockDocker) deletedVolumes() []string { m.mu.Lock() defer m.mu.Unlock() return append([]string(nil), m.volumesDeleted...) } func (m *mockDocker) createBody() map[string]any { m.mu.Lock() defer m.mu.Unlock() return m.lastCreateBody } // envForCommand returns the environment of the first exec whose command // contains sub, or nil when no command matches. func (m *mockDocker) envForCommand(sub string) []string { m.mu.Lock() defer m.mu.Unlock() for i, c := range m.execCmds { if strings.Contains(strings.Join(c, " "), sub) { return m.execEnvs[i] } } return nil } func (m *mockDocker) commands() [][]string { m.mu.Lock() defer m.mu.Unlock() return append([][]string(nil), m.execCmds...) } // allCommandText joins every command the runner ran, for a substring check. func (m *mockDocker) allCommandText() string { var parts []string for _, c := range m.commands() { parts = append(parts, strings.Join(c, " ")) } return strings.Join(parts, " ") } func writeJSON(w http.ResponseWriter, v any) { w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(v) } func (m *mockDocker) ServeHTTP(w http.ResponseWriter, r *http.Request) { p := r.URL.Path qs := r.URL.Query() switch { // Health check. case r.Method == http.MethodGet && p == "/v1.47/info": writeJSON(w, map[string]string{"ServerVersion": "mock"}) // Pull image. case r.Method == http.MethodPost && strings.HasPrefix(p, "/v1.47/images/create"): m.mu.Lock() m.pulls = append(m.pulls, qs.Get("fromImage")) m.mu.Unlock() _, _ = w.Write([]byte("{\"status\":\"Pull complete\"}\n")) // Create container. case r.Method == http.MethodPost && strings.Contains(p, "/containers/create"): var body map[string]any _ = json.NewDecoder(r.Body).Decode(&body) name := qs.Get("name") if name == "" { name = "mock-ctr-001" } // A copy creates its own source container, so the run's container // must keep its identity. if !strings.Contains(name, "-copy-") { m.mu.Lock() m.lastCreateBody = body m.mu.Unlock() } writeJSON(w, map[string]string{"Id": name}) // Start container. case r.Method == http.MethodPost && reContainerStart.MatchString(p): w.WriteHeader(http.StatusNoContent) // Create exec: take the next queued response and bind it to this id. case r.Method == http.MethodPost && reContainerExec.MatchString(p): var body struct { Cmd []string Env []string } _ = json.NewDecoder(r.Body).Decode(&body) m.mu.Lock() m.execCounter++ id := "mock-exec-" + strconv.Itoa(m.execCounter) m.execCmds = append(m.execCmds, body.Cmd) m.execEnvs = append(m.execEnvs, body.Env) resp := execResp{} if len(m.execQueue) > 0 { resp, m.execQueue = m.execQueue[0], m.execQueue[1:] } m.execMap[id] = resp m.mu.Unlock() writeJSON(w, map[string]string{"Id": id}) // Start exec: return the programmed output as a mux stream. case r.Method == http.MethodPost && reExecStart.MatchString(p): resp := m.execFor(reExecStart.FindStringSubmatch(p)[1]) if resp.delay > 0 { select { case <-time.After(resp.delay): case <-r.Context().Done(): return } } if resp.output != "" { _, _ = w.Write(muxFrame(resp.output)) } // Inspect exec: return the exit code. case r.Method == http.MethodGet && reExecJSON.MatchString(p): resp := m.execFor(reExecJSON.FindStringSubmatch(p)[1]) writeJSON(w, map[string]int{"ExitCode": resp.exitCode}) // Archive upload: the checkout and the [[copy]] sources. case r.Method == http.MethodPut && reContainerArchive.MatchString(p): body, _ := io.ReadAll(r.Body) m.mu.Lock() m.uploads = append(m.uploads, ciUpload{path: qs.Get("path"), body: body}) fail := m.uploadError m.mu.Unlock() if fail != "" { http.Error(w, fail, http.StatusInternalServerError) return } w.WriteHeader(http.StatusOK) // Archive download: artifact collection and [[copy]]. case r.Method == http.MethodGet && reContainerArchive.MatchString(p): filePath := qs.Get("path") if filePath == "" { filePath = "file.txt" } w.Header().Set("Content-Type", "application/x-tar") _, _ = w.Write(makeTar(path.Base(filePath), "artifact-content-123")) // Delete container. case r.Method == http.MethodDelete && strings.Contains(p, "/containers/"): w.WriteHeader(http.StatusNoContent) // Volume create, used for cache volumes. case r.Method == http.MethodPost && p == "/v1.47/volumes/create": var body struct { Name string Labels map[string]string } _ = json.NewDecoder(r.Body).Decode(&body) if body.Labels == nil { body.Labels = map[string]string{} } m.mu.Lock() m.volumesCreated = append(m.volumesCreated, ciVolume{name: body.Name, labels: body.Labels}) m.mu.Unlock() writeJSON(w, map[string]string{"Name": body.Name}) // Disk usage, used by the cache size caps. case r.Method == http.MethodGet && p == "/v1.47/system/df": type entry struct { Name string UsageData ciVolumeUsage } m.mu.Lock() vols := make([]entry, 0, len(m.volumeUsage)) for name, usage := range m.volumeUsage { vols = append(vols, entry{Name: name, UsageData: usage}) } m.mu.Unlock() writeJSON(w, map[string]any{"Volumes": vols}) // Volume list, used by prune and purge. case r.Method == http.MethodGet && p == "/v1.47/volumes": type entry struct{ Name string } m.mu.Lock() vols := make([]entry, 0, len(m.volumesOnHost)) for _, name := range m.volumesOnHost { vols = append(vols, entry{Name: name}) } m.mu.Unlock() writeJSON(w, map[string]any{"Volumes": vols}) // Volume delete. case r.Method == http.MethodDelete && strings.Contains(p, "/volumes/"): m.mu.Lock() m.volumesDeleted = append(m.volumesDeleted, path.Base(p)) m.mu.Unlock() w.WriteHeader(http.StatusNoContent) default: http.Error(w, "Not found", http.StatusNotFound) } } func (m *mockDocker) execFor(id string) execResp { m.mu.Lock() defer m.mu.Unlock() return m.execMap[id] } // muxFrame builds one Docker multiplexed stream frame. func muxFrame(text string) []byte { hdr := make([]byte, 8) hdr[0] = 1 binary.BigEndian.PutUint32(hdr[4:], uint32(len(text))) return append(hdr, text...) } // makeTar builds a tar archive holding one file. func makeTar(filename, content string) []byte { var buf bytes.Buffer tw := tar.NewWriter(&buf) _ = tw.WriteHeader(&tar.Header{ Name: filename, Mode: 0o644, Size: int64(len(content)), Typeflag: tar.TypeReg, }) _, _ = tw.Write([]byte(content)) _ = tw.Close() return buf.Bytes() } // tarEntry is one header of an uploaded archive. type tarEntry struct { name string uid int } // tarHeaders lists the entries of an uncompressed tar. func tarHeaders(t *testing.T, data []byte) []tarEntry { t.Helper() tr := tar.NewReader(bytes.NewReader(data)) var out []tarEntry for { h, err := tr.Next() if err == io.EOF { break } if err != nil { t.Fatalf("read tar: %v", err) } out = append(out, tarEntry{name: h.Name, uid: h.Uid}) } return out } // tarEntryNames lists the file names of an uncompressed tar. func tarEntryNames(t *testing.T, data []byte) []string { t.Helper() var out []string for _, h := range tarHeaders(t, data) { out = append(out, h.name) } return out }