package containerengine import ( "context" "encoding/json" "errors" "fmt" "io" cerrdefs "github.com/containerd/errdefs" "github.com/moby/moby/api/pkg/stdcopy" "github.com/moby/moby/api/types/container" "github.com/moby/moby/api/types/jsonstream" "github.com/moby/moby/api/types/mount" "github.com/moby/moby/client" ocispec "github.com/opencontainers/image-spec/specs-go/v1" ) // MobyEngine 适配官方 Docker Engine Go 客户端,实现 Engine 接口。 type MobyEngine struct { // client 底层 Docker Engine 客户端。 client *client.Client } // NewMobyEngine 依据 Docker 官方文档定义的环境变量创建客户端。 // 即使 DOCKER_HOST 指向非默认 socket,API 版本协商仍然保持启用。 // 创建失败时返回包装后的错误。 func NewMobyEngine() (*MobyEngine, error) { apiClient, err := client.New(client.FromEnv) if err != nil { return nil, fmt.Errorf("create Docker Engine client: %w", err) } return &MobyEngine{client: apiClient}, nil } // Ping 探测 Docker Engine 是否可达,不可达时返回包装后的错误。 func (e *MobyEngine) Ping(ctx context.Context) error { if _, err := e.client.Ping(ctx, client.PingOptions{NegotiateAPIVersion: true}); err != nil { return fmt.Errorf("ping Docker Engine: %w", err) } return nil } // PullImage 从远端仓库拉取指定镜像引用。 // 拉取完成后逐条读取并校验引擎返回的 JSON 流,确保过程中无错误。 func (e *MobyEngine) PullImage(ctx context.Context, imageReference string) error { response, err := e.client.ImagePull(ctx, imageReference, client.ImagePullOptions{}) if err != nil { return fmt.Errorf("pull image %s: %w", imageReference, err) } defer response.Close() if err := decodeImageLoadResponse(response); err != nil { return fmt.Errorf("pull image %s response: %w", imageReference, err) } return nil } // LoadImage 从归档读取流加载镜像。 // input 为空时返回错误;加载完成后同样校验引擎返回的 JSON 流。 func (e *MobyEngine) LoadImage(ctx context.Context, input io.Reader) error { if input == nil { return errors.New("image archive reader is required") } response, err := e.client.ImageLoad(ctx, input) if err != nil { return fmt.Errorf("load image archive: %w", err) } defer response.Close() if err := decodeImageLoadResponse(response); err != nil { return fmt.Errorf("load image archive response: %w", err) } return nil } // decodeImageLoadResponse 逐条解码 Docker 引擎的 JSON 消息流。 // 读到 EOF 表示成功;遇到消息中的错误字段时立即返回该错误。 func decodeImageLoadResponse(input io.Reader) error { decoder := json.NewDecoder(input) for { var message jsonstream.Message if err := decoder.Decode(&message); err != nil { if errors.Is(err, io.EOF) { return nil } return fmt.Errorf("decode Docker JSON stream: %w", err) } if message.Error != nil { return message.Error } } } // InspectImage 检查指定镜像引用并返回不可变镜像信息。 // 找不到镜像时返回的错误会包装 ErrNotFound。 func (e *MobyEngine) InspectImage(ctx context.Context, reference string) (Image, error) { response, err := e.client.ImageInspect(ctx, reference) if err != nil { return Image{}, engineError("inspect image", err) } descriptorDigest := "" if response.Descriptor != nil { descriptorDigest = response.Descriptor.Digest.String() } return Image{ ID: response.ID, RepoDigests: append([]string(nil), response.RepoDigests...), DescriptorDigest: descriptorDigest, Platform: Platform{ OS: response.Os, Architecture: response.Architecture, Variant: response.Variant, }, }, nil } // CreateContainer 依据给定规格创建容器。 // 成功创建后立即通过 InspectContainer 返回容器的完整运行时状态。 func (e *MobyEngine) CreateContainer(ctx context.Context, spec ContainerSpec) (Container, error) { apiMounts := make([]mount.Mount, 0, len(spec.Mounts)) for _, item := range spec.Mounts { apiMounts = append(apiMounts, mount.Mount{ Type: mount.Type(item.Type), Source: item.Source, Target: item.Target, ReadOnly: item.ReadOnly, }) } stopTimeout := spec.StopTimeoutSeconds result, err := e.client.ContainerCreate(ctx, client.ContainerCreateOptions{ Config: &container.Config{ Env: append([]string(nil), spec.Environment...), Labels: cloneMap(spec.Labels), User: spec.User, StopTimeout: &stopTimeout, }, HostConfig: &container.HostConfig{ NetworkMode: container.NetworkMode(spec.NetworkMode), RestartPolicy: container.RestartPolicy{ Name: container.RestartPolicyMode(spec.RestartPolicy.Name), MaximumRetryCount: spec.RestartPolicy.MaximumRetryCount, }, Mounts: apiMounts, }, Platform: &ocispec.Platform{ OS: spec.Platform.OS, Architecture: spec.Platform.Architecture, Variant: spec.Platform.Variant, }, Name: spec.Name, Image: spec.ImageReference, }) if err != nil { return Container{}, engineError("create container", err) } return e.InspectContainer(ctx, result.ID) } // StartContainer 启动指定容器,参数可为容器 ID 或名称。 func (e *MobyEngine) StartContainer(ctx context.Context, idOrName string) error { if _, err := e.client.ContainerStart(ctx, idOrName, client.ContainerStartOptions{}); err != nil { return engineError("start container", err) } return nil } // ContainerLogs 跟随读取指定容器的 stdout 与 stderr 日志,将引擎的多路复用流 // 解码后合并为一个持续流返回,直到 ctx 取消或容器停止。 func (e *MobyEngine) ContainerLogs(ctx context.Context, idOrName string) (io.ReadCloser, error) { stream, err := e.client.ContainerLogs(ctx, idOrName, client.ContainerLogsOptions{ShowStdout: true, ShowStderr: true, Follow: true, Tail: "all"}) if err != nil { return nil, engineError("read container logs", err) } reader, writer := io.Pipe() go func() { _, copyErr := stdcopy.StdCopy(writer, writer, stream) _ = stream.Close() _ = writer.CloseWithError(copyErr) }() return reader, nil } // StopContainer 停止指定容器,参数可为容器 ID 或名称。 // timeoutSeconds 是发送停止信号后到强制终止前的等待秒数,传给引擎以在超时后强制终止。 func (e *MobyEngine) StopContainer(ctx context.Context, idOrName string, timeoutSeconds int) error { if _, err := e.client.ContainerStop(ctx, idOrName, client.ContainerStopOptions{Timeout: &timeoutSeconds}); err != nil { return engineError("stop container", err) } return nil } // InspectContainer 检查指定容器并返回其运行时状态。 // 各字段按引擎返回的嵌套结构逐层提取,缺失的可选部分保持零值。 func (e *MobyEngine) InspectContainer(ctx context.Context, idOrName string) (Container, error) { result, err := e.client.ContainerInspect(ctx, idOrName, client.ContainerInspectOptions{}) if err != nil { return Container{}, engineError("inspect container", err) } response := result.Container record := Container{ ID: response.ID, Name: response.Name, ImageID: response.Image, Platform: response.Platform, } if response.State != nil { record.Running = response.State.Running record.Dead = response.State.Dead record.Status = string(response.State.Status) } if response.Config != nil { record.ImageReference = response.Config.Image record.Environment = append([]string(nil), response.Config.Env...) record.Labels = cloneMap(response.Config.Labels) record.User = response.Config.User if response.Config.StopTimeout != nil { record.StopTimeoutSeconds = *response.Config.StopTimeout } } if response.HostConfig != nil { record.NetworkMode = string(response.HostConfig.NetworkMode) record.RestartPolicy = RestartPolicy{ Name: string(response.HostConfig.RestartPolicy.Name), MaximumRetryCount: response.HostConfig.RestartPolicy.MaximumRetryCount, } } for _, item := range response.Mounts { record.Mounts = append(record.Mounts, Mount{ Type: string(item.Type), Source: item.Source, Target: item.Destination, ReadOnly: !item.RW, }) } return record, nil } // RemoveContainer 移除指定容器,force 决定是否强制移除正在运行的容器。 func (e *MobyEngine) RemoveContainer(ctx context.Context, idOrName string, force bool) error { _, err := e.client.ContainerRemove(ctx, idOrName, client.ContainerRemoveOptions{Force: force}) if err != nil { return engineError("remove container", err) } return nil } // Close 关闭底层 Docker Engine 客户端连接。 func (e *MobyEngine) Close() error { return e.client.Close() } // engineError 包装引擎操作错误。 // 当底层错误属于“未找到”类别时,额外包装 ErrNotFound,以便调用方通过 errors.Is 识别。 func engineError(action string, err error) error { if cerrdefs.IsNotFound(err) { return fmt.Errorf("%s: %w: %v", action, ErrNotFound, err) } return fmt.Errorf("%s: %w", action, err) } // cloneMap 浅拷贝一个字符串映射,避免调用方后续修改影响内部数据。 // 若 source 为 nil 则返回 nil。 func cloneMap(source map[string]string) map[string]string { if source == nil { return nil } result := make(map[string]string, len(source)) for key, value := range source { result[key] = value } return result } var _ Engine = (*MobyEngine)(nil)