401 lines
13 KiB
Go
401 lines
13 KiB
Go
package customer
|
|
|
|
import (
|
|
"archive/tar"
|
|
"bytes"
|
|
"context"
|
|
"encoding/base64"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"net"
|
|
"net/http"
|
|
"net/url"
|
|
"path/filepath"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
type DockerClient struct {
|
|
hc *http.Client
|
|
base string
|
|
}
|
|
|
|
func NewDockerClient(raw string) (*DockerClient, error) {
|
|
raw = strings.TrimSpace(raw)
|
|
if raw == "" {
|
|
raw = "unix:///var/run/docker.sock"
|
|
}
|
|
if strings.HasPrefix(raw, "unix://") {
|
|
sock := strings.TrimPrefix(raw, "unix://")
|
|
tr := &http.Transport{DialContext: func(ctx context.Context, network, addr string) (net.Conn, error) {
|
|
var d net.Dialer
|
|
return d.DialContext(ctx, "unix", sock)
|
|
}}
|
|
return &DockerClient{hc: &http.Client{Transport: tr, Timeout: 30 * time.Second}, base: "http://docker"}, nil
|
|
}
|
|
u, err := url.Parse(raw)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if u.Scheme != "http" && u.Scheme != "https" {
|
|
return nil, errors.New("DOCKER_HOST must be unix://, http:// or https://")
|
|
}
|
|
return &DockerClient{hc: &http.Client{Timeout: 30 * time.Second}, base: strings.TrimRight(raw, "/")}, nil
|
|
}
|
|
|
|
func (d *DockerClient) req(ctx context.Context, method, path string, in, out any) error {
|
|
var body io.Reader
|
|
if in != nil {
|
|
b, err := json.Marshal(in)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
body = bytes.NewReader(b)
|
|
}
|
|
req, err := http.NewRequestWithContext(ctx, method, d.base+path, body)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if in != nil {
|
|
req.Header.Set("Content-Type", "application/json")
|
|
}
|
|
resp, err := d.hc.Do(req)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer resp.Body.Close()
|
|
b, readErr := io.ReadAll(io.LimitReader(resp.Body, 4<<20))
|
|
if readErr != nil {
|
|
return readErr
|
|
}
|
|
if resp.StatusCode/100 != 2 {
|
|
return fmt.Errorf("docker API %s %s HTTP %d: %s", method, path, resp.StatusCode, strings.TrimSpace(string(b)))
|
|
}
|
|
if out != nil && len(bytes.TrimSpace(b)) > 0 {
|
|
return json.Unmarshal(b, out)
|
|
}
|
|
return nil
|
|
}
|
|
func (d *DockerClient) Ping(ctx context.Context) error {
|
|
return d.req(ctx, http.MethodGet, "/_ping", nil, nil)
|
|
}
|
|
|
|
// ImageExists checks the local Docker image cache without pulling anything.
|
|
func (d *DockerClient) ImageExists(ctx context.Context, image string) (bool, error) {
|
|
image = strings.TrimSpace(image)
|
|
if image == "" {
|
|
return false, errors.New("worker image is empty")
|
|
}
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, d.base+"/images/"+url.PathEscape(image)+"/json", nil)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
resp, err := d.hc.Do(req)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
defer resp.Body.Close()
|
|
_, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, 1<<20))
|
|
if resp.StatusCode == http.StatusNotFound {
|
|
return false, nil
|
|
}
|
|
if resp.StatusCode/100 != 2 {
|
|
return false, fmt.Errorf("docker image inspect HTTP %d", resp.StatusCode)
|
|
}
|
|
return true, nil
|
|
}
|
|
|
|
// PullImage asks Docker Engine to pull a public/configured registry image.
|
|
// Private-registry credentials are passed as Docker's X-Registry-Auth header.
|
|
// The stream is inspected for daemon-side pull errors.
|
|
func (d *DockerClient) PullImage(ctx context.Context, image, registryAuth string) error {
|
|
image = strings.TrimSpace(image)
|
|
if image == "" {
|
|
return errors.New("worker image is empty")
|
|
}
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, d.base+"/images/create?fromImage="+url.QueryEscape(image), nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if strings.TrimSpace(registryAuth) != "" {
|
|
req.Header.Set("X-Registry-Auth", strings.TrimSpace(registryAuth))
|
|
}
|
|
pullClient := *d.hc
|
|
pullClient.Timeout = 10 * time.Minute
|
|
resp, err := pullClient.Do(req)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer resp.Body.Close()
|
|
if resp.StatusCode/100 != 2 {
|
|
b, _ := io.ReadAll(io.LimitReader(resp.Body, 64<<10))
|
|
return fmt.Errorf("docker image pull HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(b)))
|
|
}
|
|
dec := json.NewDecoder(io.LimitReader(resp.Body, 32<<20))
|
|
for {
|
|
var msg struct {
|
|
Error string `json:"error"`
|
|
}
|
|
if err := dec.Decode(&msg); err != nil {
|
|
if errors.Is(err, io.EOF) {
|
|
break
|
|
}
|
|
return fmt.Errorf("docker image pull stream: %w", err)
|
|
}
|
|
if strings.TrimSpace(msg.Error) != "" {
|
|
return fmt.Errorf("docker image pull: %s", strings.TrimSpace(msg.Error))
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// RegistryAuthHeader builds Docker Engine's X-Registry-Auth value. Use a
|
|
// registry-scoped read-only deploy token instead of a personal password.
|
|
func RegistryAuthHeader(username, password, serverAddress string) (string, error) {
|
|
username = strings.TrimSpace(username)
|
|
password = strings.TrimSpace(password)
|
|
serverAddress = strings.TrimSpace(serverAddress)
|
|
if username == "" && password == "" && serverAddress == "" {
|
|
return "", nil
|
|
}
|
|
if username == "" || password == "" {
|
|
return "", errors.New("both worker registry username and password/token are required")
|
|
}
|
|
payload := map[string]string{"username": username, "password": password}
|
|
if serverAddress != "" {
|
|
payload["serveraddress"] = serverAddress
|
|
}
|
|
b, err := json.Marshal(payload)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
return base64.URLEncoding.EncodeToString(b), nil
|
|
}
|
|
|
|
func (d *DockerClient) EnsureImage(ctx context.Context, image string, autoPull bool, registryAuth string) error {
|
|
ok, err := d.ImageExists(ctx, image)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if ok {
|
|
return nil
|
|
}
|
|
if !autoPull {
|
|
return fmt.Errorf("worker image %q is not present on the Docker host and CS_WORKER_AUTO_PULL is disabled", image)
|
|
}
|
|
if err := d.PullImage(ctx, image, registryAuth); err != nil {
|
|
return fmt.Errorf("pull worker image %q: %w", image, err)
|
|
}
|
|
ok, err = d.ImageExists(ctx, image)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !ok {
|
|
return fmt.Errorf("worker image %q is still unavailable after pull", image)
|
|
}
|
|
return nil
|
|
}
|
|
func (d *DockerClient) CreateVolume(ctx context.Context, name string) error {
|
|
var out map[string]any
|
|
return d.req(ctx, http.MethodPost, "/volumes/create", map[string]any{"Name": name, "Labels": map[string]string{"neuralhunt.managed": "true"}}, &out)
|
|
}
|
|
|
|
type WorkerContainerConfig struct {
|
|
Image, Entrypoint, Network, GameURL, RegisterURL, WorkerID, RegisterToken, TaskID, BeaconPath, Volume, Name string
|
|
}
|
|
|
|
func (d *DockerClient) CreateWorker(ctx context.Context, c WorkerContainerConfig) (string, error) {
|
|
name := url.QueryEscape(c.Name)
|
|
body := map[string]any{
|
|
"Image": c.Image,
|
|
// Named Docker volumes are root-owned when first mounted, and identities
|
|
// created by older releases may be owned by a different UID. The dedicated
|
|
// worker image starts as uid 0 only for its tiny ownership-normalization
|
|
// entrypoint and immediately drops to the unprivileged app user.
|
|
"User": "0:0",
|
|
"Cmd": []string{"-url", c.GameURL, "-identity", "/identity/identity.json", "-non-interactive", "-quiet", "-task", c.TaskID, "-beacon-path", c.BeaconPath},
|
|
"Env": []string{
|
|
"NEURALHUNT_WORKER_REGISTER_URL=" + c.RegisterURL,
|
|
"NEURALHUNT_WORKER_LEASE_URL=" + strings.TrimSuffix(c.RegisterURL, "/register") + "/lease",
|
|
"NEURALHUNT_WORKER_REGISTER_TOKEN=" + c.RegisterToken,
|
|
"NEURALHUNT_WORKER_ID=" + c.WorkerID,
|
|
},
|
|
"Labels": map[string]string{"neuralhunt.managed": "true", "neuralhunt.worker_id": c.WorkerID},
|
|
"HostConfig": map[string]any{
|
|
"Mounts": []map[string]any{{"Type": "volume", "Source": c.Volume, "Target": "/identity"}},
|
|
"NetworkMode": c.Network,
|
|
"ReadonlyRootfs": true,
|
|
// A managed unattended worker must survive process crashes and temporary
|
|
// game/network outages. Docker will restart it automatically, while an
|
|
// explicit Docker Stop from the Customer Service keeps it stopped.
|
|
"RestartPolicy": map[string]any{"Name": "unless-stopped", "MaximumRetryCount": 0},
|
|
"CapDrop": []string{"ALL"},
|
|
// Bootstrap-only capabilities: the image entrypoint fixes ownership of
|
|
// /identity and then su-exec permanently switches to uid/gid app.
|
|
"CapAdd": []string{"CHOWN", "DAC_OVERRIDE", "SETUID", "SETGID"},
|
|
"SecurityOpt": []string{"no-new-privileges"},
|
|
"PidsLimit": 128,
|
|
"Memory": 256 * 1024 * 1024,
|
|
"NanoCpus": int64(1_000_000_000),
|
|
},
|
|
}
|
|
// A dedicated worker image already declares /app/neuralhunt-client as its
|
|
// ENTRYPOINT. Leaving Entrypoint unset makes CS_WORKER_IMAGE genuinely
|
|
// pluggable. CS_WORKER_ENTRYPOINT exists only as a compatibility override
|
|
// for older monolithic images.
|
|
if strings.TrimSpace(c.Entrypoint) != "" {
|
|
body["Entrypoint"] = []string{strings.TrimSpace(c.Entrypoint)}
|
|
}
|
|
var out struct {
|
|
ID string `json:"Id"`
|
|
}
|
|
if err := d.req(ctx, http.MethodPost, "/containers/create?name="+name, body, &out); err != nil {
|
|
return "", err
|
|
}
|
|
if out.ID == "" {
|
|
return "", errors.New("docker returned empty container id")
|
|
}
|
|
return out.ID, nil
|
|
}
|
|
func (d *DockerClient) Start(ctx context.Context, id string) error {
|
|
return d.req(ctx, http.MethodPost, "/containers/"+url.PathEscape(id)+"/start", nil, nil)
|
|
}
|
|
func (d *DockerClient) Stop(ctx context.Context, id string, seconds int) error {
|
|
if id == "" {
|
|
return nil
|
|
}
|
|
if seconds < 1 {
|
|
seconds = 10
|
|
}
|
|
return d.req(ctx, http.MethodPost, "/containers/"+url.PathEscape(id)+"/stop?t="+fmt.Sprint(seconds), nil, nil)
|
|
}
|
|
func (d *DockerClient) Remove(ctx context.Context, id string) error {
|
|
if id == "" {
|
|
return nil
|
|
}
|
|
err := d.req(ctx, http.MethodDelete, "/containers/"+url.PathEscape(id)+"?force=true&v=false", nil, nil)
|
|
if err != nil && strings.Contains(err.Error(), "404") {
|
|
return nil
|
|
}
|
|
return err
|
|
}
|
|
func (d *DockerClient) Running(ctx context.Context, id string) (bool, error) {
|
|
var out struct {
|
|
State struct {
|
|
Running bool `json:"Running"`
|
|
} `json:"State"`
|
|
}
|
|
if err := d.req(ctx, http.MethodGet, "/containers/"+url.PathEscape(id)+"/json", nil, &out); err != nil {
|
|
return false, err
|
|
}
|
|
return out.State.Running, nil
|
|
}
|
|
|
|
func (d *DockerClient) GetFile(ctx context.Context, containerID, path string) ([]byte, error) {
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, d.base+"/containers/"+url.PathEscape(containerID)+"/archive?path="+url.QueryEscape(path), nil)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
resp, err := d.hc.Do(req)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer resp.Body.Close()
|
|
if resp.StatusCode/100 != 2 {
|
|
b, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
|
|
return nil, fmt.Errorf("docker archive HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(b)))
|
|
}
|
|
tr := tar.NewReader(io.LimitReader(resp.Body, 8<<20))
|
|
for {
|
|
h, err := tr.Next()
|
|
if errors.Is(err, io.EOF) {
|
|
break
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if filepath.Base(h.Name) == filepath.Base(path) && h.Typeflag == tar.TypeReg {
|
|
return io.ReadAll(io.LimitReader(tr, 2<<20))
|
|
}
|
|
}
|
|
return nil, errors.New("identity file not found in container volume")
|
|
}
|
|
func (d *DockerClient) PutFile(ctx context.Context, containerID, dir, name string, data []byte) error {
|
|
var buf bytes.Buffer
|
|
tw := tar.NewWriter(&buf)
|
|
if err := tw.WriteHeader(&tar.Header{Name: name, Mode: 0600, Size: int64(len(data)), ModTime: time.Now()}); err != nil {
|
|
return err
|
|
}
|
|
if _, err := tw.Write(data); err != nil {
|
|
return err
|
|
}
|
|
if err := tw.Close(); err != nil {
|
|
return err
|
|
}
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPut, d.base+"/containers/"+url.PathEscape(containerID)+"/archive?path="+url.QueryEscape(dir), bytes.NewReader(buf.Bytes()))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
req.Header.Set("Content-Type", "application/x-tar")
|
|
resp, err := d.hc.Do(req)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer resp.Body.Close()
|
|
if resp.StatusCode/100 != 2 {
|
|
b, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
|
|
return fmt.Errorf("docker put archive HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(b)))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (d *DockerClient) RemoveVolume(ctx context.Context, name string) error {
|
|
if name == "" {
|
|
return nil
|
|
}
|
|
err := d.req(ctx, http.MethodDelete, "/volumes/"+url.PathEscape(name)+"?force=true", nil, nil)
|
|
if err != nil && strings.Contains(err.Error(), "404") {
|
|
return nil
|
|
}
|
|
return err
|
|
}
|
|
|
|
type ManagedContainer struct {
|
|
ID string `json:"id"`
|
|
Name string `json:"name"`
|
|
Image string `json:"image"`
|
|
State string `json:"state"`
|
|
Status string `json:"status"`
|
|
WorkerID string `json:"worker_id"`
|
|
Running bool `json:"running"`
|
|
}
|
|
|
|
// ManagedContainers lists only containers created by Neural Hunt. It is used by
|
|
// the Service Controller emergency UI and never exposes unrelated Docker
|
|
// workloads on the host.
|
|
func (d *DockerClient) ManagedContainers(ctx context.Context) ([]ManagedContainer, error) {
|
|
q := url.QueryEscape(`{"label":["neuralhunt.managed=true"]}`)
|
|
var raw []struct {
|
|
ID string `json:"Id"`
|
|
Names []string `json:"Names"`
|
|
Image string `json:"Image"`
|
|
State string `json:"State"`
|
|
Status string `json:"Status"`
|
|
Labels map[string]string `json:"Labels"`
|
|
}
|
|
if err := d.req(ctx, http.MethodGet, "/containers/json?all=true&filters="+q, nil, &raw); err != nil {
|
|
return nil, err
|
|
}
|
|
out := make([]ManagedContainer, 0, len(raw))
|
|
for _, x := range raw {
|
|
name := ""
|
|
if len(x.Names) > 0 {
|
|
name = strings.TrimPrefix(x.Names[0], "/")
|
|
}
|
|
out = append(out, ManagedContainer{ID: x.ID, Name: name, Image: x.Image, State: x.State, Status: x.Status, WorkerID: x.Labels["neuralhunt.worker_id"], Running: x.State == "running"})
|
|
}
|
|
return out, nil
|
|
}
|