package replays import ( "archive/tar" "bytes" "context" "encoding/binary" "encoding/json" "fmt" "io" "net" "net/http" "net/url" "os" "strings" "time" ) // Docker is the container runtime's own API, over its socket: what a replay needs to raise a resolver and // ask it from programs under two C libraries — pull, network, run, wait, logs, remove — and nothing else. // Spoken directly rather than through the command line, because the Go toolchain a check runs in carries // no docker client, and a replay that needed one would be a replay that runs only on a workstation. type Docker struct { http *http.Client // Label marks everything it makes, so a replay that dies is cleaned up by it. Label string } // DockerFromEnv is the runtime at DOCKER_HOST's socket, or the usual one; ok false when there is none. func DockerFromEnv() (*Docker, bool) { sock := "/var/run/docker.sock" if h := os.Getenv("DOCKER_HOST"); strings.HasPrefix(h, "unix://") { sock = strings.TrimPrefix(h, "unix://") } if _, err := os.Stat(sock); err != nil { return nil, false } return &Docker{Label: "mesh.replay", http: &http.Client{Transport: &http.Transport{ DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) { return (&net.Dialer{}).DialContext(ctx, "unix", sock) }}}}, true } func (d *Docker) call(ctx context.Context, method, path string, body any, out any) error { var r io.Reader if body != nil { raw, err := json.Marshal(body) if err != nil { return err } r = bytes.NewReader(raw) } req, err := http.NewRequestWithContext(ctx, method, "http://docker"+path, r) if err != nil { return err } req.Header.Set("Content-Type", "application/json") res, err := d.http.Do(req) if err != nil { return err } defer res.Body.Close() raw, _ := io.ReadAll(res.Body) if res.StatusCode >= 300 { return fmt.Errorf("docker %s %s: %s %s", method, path, res.Status, strings.TrimSpace(string(raw))) } if out != nil && len(raw) > 0 { return json.Unmarshal(raw, out) } return nil } // Pull fetches an image unless the runtime has it. func (d *Docker) Pull(ctx context.Context, image string) error { if err := d.call(ctx, http.MethodGet, "/images/"+url.PathEscape(image)+"/json", nil, nil); err == nil { return nil } // A digest is the tag the runtime is asked for, as its own client asks: `name@sha256:…` cut at the // `@`; otherwise the tag after the last `:` that is not part of a registry's port. name, tag, digested := strings.Cut(image, "@") if !digested { name, tag = image, "latest" if at := strings.LastIndex(image, ":"); at > strings.LastIndex(image, "/") { name, tag = image[:at], image[at+1:] } } err := d.call(ctx, http.MethodPost, "/images/create?fromImage="+url.QueryEscape(name)+"&tag="+url.QueryEscape(tag), nil, nil) if err != nil { return err } // The runtime answers a pull it could not make with a stream that ends in an error, not a status. return d.call(ctx, http.MethodGet, "/images/"+url.PathEscape(image)+"/json", nil, nil) } // Network makes a network of its own and answers its id. func (d *Docker) Network(ctx context.Context, name string) (string, error) { var made struct{ ID string } err := d.call(ctx, http.MethodPost, "/networks/create", map[string]any{"Name": name, "Labels": map[string]string{d.Label: name}}, &made) return made.ID, err } // RemoveNetwork removes a network. func (d *Docker) RemoveNetwork(ctx context.Context, id string) { _ = d.call(ctx, http.MethodDelete, "/networks/"+id, nil, nil) } // Run is one container: an image, a command, the files it is given, its network and resolver. type Run struct { Image string Cmd []string Binds []string Network string DNS []string // Name, Labels and Restart are the container's name, its labels beside the replay's own, and the // runtime's restart policy — what the node-engine creates a module's container with. Name string Labels map[string]string Restart string } // Start creates and starts a container and answers its id and its address on its network. func (d *Docker) Start(ctx context.Context, r Run) (string, string, error) { var made struct{ ID string } host := map[string]any{"Binds": r.Binds, "NetworkMode": r.Network} if len(r.DNS) > 0 { host["Dns"] = r.DNS } if r.Restart != "" { host["RestartPolicy"] = map[string]any{"Name": r.Restart} } labels := map[string]string{d.Label: "1"} for k, v := range r.Labels { labels[k] = v } path := "/containers/create" if r.Name != "" { path += "?name=" + url.QueryEscape(r.Name) } if err := d.call(ctx, http.MethodPost, path, map[string]any{"Image": r.Image, "Cmd": r.Cmd, "Labels": labels, "HostConfig": host}, &made); err != nil { return "", "", err } if err := d.call(ctx, http.MethodPost, "/containers/"+made.ID+"/start", nil, nil); err != nil { d.Remove(ctx, made.ID) return "", "", err } var seen struct { NetworkSettings struct { Networks map[string]struct{ IPAddress string } } } if err := d.call(ctx, http.MethodGet, "/containers/"+made.ID+"/json", nil, &seen); err != nil { return made.ID, "", err } for _, n := range seen.NetworkSettings.Networks { return made.ID, n.IPAddress, nil } return made.ID, "", nil } // Wait waits for a container to end, and answers its exit code and what it said. func (d *Docker) Wait(ctx context.Context, id string) (int, string, error) { var ended struct{ StatusCode int } if err := d.call(ctx, http.MethodPost, "/containers/"+id+"/wait", nil, &ended); err != nil { return -1, "", err } return ended.StatusCode, d.Logs(ctx, id), nil } // Logs is what a container said, both streams. func (d *Docker) Logs(ctx context.Context, id string) string { req, err := http.NewRequestWithContext(ctx, http.MethodGet, "http://docker/containers/"+id+"/logs?stdout=1&stderr=1", nil) if err != nil { return "" } res, err := d.http.Do(req) if err != nil { return "" } defer res.Body.Close() var b strings.Builder head := make([]byte, 8) for { if _, err := io.ReadFull(res.Body, head); err != nil { return b.String() } n := binary.BigEndian.Uint32(head[4:]) if _, err := io.CopyN(&b, res.Body, int64(n)); err != nil { return b.String() } } } // Remove removes a container, running or not. func (d *Docker) Remove(ctx context.Context, id string) { _ = d.call(ctx, http.MethodDelete, "/containers/"+id+"?force=1", nil, nil) } // Exec runs a command in a running container and answers its exit code. func (d *Docker) Exec(ctx context.Context, id string, cmd ...string) (int, error) { var made struct{ ID string } if err := d.call(ctx, http.MethodPost, "/containers/"+id+"/exec", map[string]any{"Cmd": cmd}, &made); err != nil { return -1, err } if err := d.call(ctx, http.MethodPost, "/exec/"+made.ID+"/start", map[string]any{"Detach": true}, nil); err != nil { return -1, err } for i := 0; i < 600; i++ { var seen struct { Running bool ExitCode int } if err := d.call(ctx, http.MethodGet, "/exec/"+made.ID+"/json", nil, &seen); err != nil { return -1, err } if !seen.Running { return seen.ExitCode, nil } time.Sleep(100 * time.Millisecond) } return -1, fmt.Errorf("%v did not end", cmd) } // Build builds an image from a Dockerfile alone, tagged as given, and labelled so a replay that dies is // cleaned up by it. func (d *Docker) Build(ctx context.Context, tag, dockerfile string) error { var archive bytes.Buffer tw := tar.NewWriter(&archive) if err := tw.WriteHeader(&tar.Header{Name: "Dockerfile", Mode: 0o644, Size: int64(len(dockerfile))}); err != nil { return err } if _, err := tw.Write([]byte(dockerfile)); err != nil { return err } if err := tw.Close(); err != nil { return err } labels, _ := json.Marshal(map[string]string{d.Label: "1"}) req, err := http.NewRequestWithContext(ctx, http.MethodPost, "http://docker/build?t="+url.QueryEscape(tag)+ "&labels="+url.QueryEscape(string(labels))+"&rm=1", &archive) if err != nil { return err } req.Header.Set("Content-Type", "application/x-tar") res, err := d.http.Do(req) if err != nil { return err } defer res.Body.Close() said, _ := io.ReadAll(res.Body) if res.StatusCode >= 300 || bytes.Contains(said, []byte(`"error"`)) { return fmt.Errorf("building %s: %s %s", tag, res.Status, strings.TrimSpace(string(said))) } return nil } // RemoveImage removes an image. func (d *Docker) RemoveImage(ctx context.Context, ref string) { _ = d.call(ctx, http.MethodDelete, "/images/"+url.PathEscape(ref)+"?force=1", nil, nil) }