Files
1Panel/agent/app/service/container.go
ssongliu 7915230121 refactor: rebuild firewall management (#13628)
* refactor(firewall): rebuild rule management foundation

* refactor(firewall): streamline rule checks and inventory

* feat(firewall): improve native rule inventory

* refactor(firewall): refine rule management

* feat: add Docker port guard

* feat(firewall): support native nftables

* feat(firewall): add configurable firewall selection

* feat(firewall): support nftables docker port guard

* refactor(firewall): complete v2 rule management and migration

* refactor(firewall): align state and API contracts

* refactor(firewall): unify rule management operations

* feat: refine firewall v2 rules and forwarding

* fix(firewall): harden dual-stack rule management

* refactor(firewall): consolidate rule validation and persistence
2026-08-24 12:51:34 +08:00

2094 lines
59 KiB
Go

package service
import (
"archive/tar"
"bufio"
"bytes"
"context"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"net/http"
"net/url"
"os"
"os/exec"
"path"
"path/filepath"
"regexp"
"sort"
"strconv"
"strings"
"sync"
"syscall"
"time"
"github.com/1Panel-dev/1Panel/agent/app/dto"
"github.com/1Panel-dev/1Panel/agent/app/repo"
"github.com/1Panel-dev/1Panel/agent/app/task"
"github.com/1Panel-dev/1Panel/agent/buserr"
"github.com/1Panel-dev/1Panel/agent/constant"
"github.com/1Panel-dev/1Panel/agent/global"
"github.com/1Panel-dev/1Panel/agent/i18n"
"github.com/1Panel-dev/1Panel/agent/utils/cmd"
"github.com/1Panel-dev/1Panel/agent/utils/common"
"github.com/1Panel-dev/1Panel/agent/utils/docker"
"github.com/docker/docker/api/types"
"github.com/docker/docker/api/types/build"
"github.com/docker/docker/api/types/container"
"github.com/docker/docker/api/types/filters"
"github.com/docker/docker/api/types/image"
"github.com/docker/docker/api/types/mount"
"github.com/docker/docker/api/types/network"
"github.com/docker/docker/api/types/registry"
"github.com/docker/docker/api/types/volume"
"github.com/docker/docker/client"
"github.com/docker/docker/pkg/stdcopy"
"github.com/docker/go-connections/nat"
"github.com/gin-gonic/gin"
v1 "github.com/opencontainers/image-spec/specs-go/v1"
"github.com/shirou/gopsutil/v4/cpu"
"github.com/shirou/gopsutil/v4/mem"
)
type ContainerService struct{}
var containerLogAnsiRegex = regexp.MustCompile("\x1b\\[[0-9;?]*[A-Za-z]|\x1b=|\x1b>")
type IContainerService interface {
Page(req dto.PageContainer) (int64, interface{}, error)
List() []dto.ContainerOptions
ListByImage(imageName string) []dto.ContainerOptions
LoadStatus() (dto.ContainerStatus, error)
PageNetwork(req dto.SearchWithPage) (int64, interface{}, error)
ListNetwork() ([]dto.Options, error)
PageVolume(req dto.SearchWithPage) (int64, interface{}, error)
ListVolume() ([]dto.Options, error)
PageCompose(req dto.SearchWithPage) (int64, interface{}, error)
LoadComposeEnv(name string) (string, error)
CreateCompose(req dto.ComposeCreate) error
ComposeOperation(req dto.ComposeOperation) error
TestCompose(req dto.ComposeCreate) (bool, error)
ComposeUpdate(req dto.ComposeUpdate) error
ComposePin(req dto.ComposePin) error
ComposeLogClean(req dto.ComposeLogClean) error
ContainerCreate(req dto.ContainerOperate, inThread bool) error
ContainerUpdate(req dto.ContainerOperate) error
ContainerUpgrade(req dto.ContainerUpgrade) error
ContainerInfo(req dto.OperationWithName) (*dto.ContainerOperate, error)
ContainerListStats() ([]dto.ContainerListStats, error)
ContainerItemStats(req dto.OperationWithName) (dto.ContainerItemStats, error)
LoadResourceLimit() (*dto.ResourceLimit, error)
ContainerRename(req dto.ContainerRename) error
ContainerCommit(req dto.ContainerCommit) error
ContainerLogClean(req dto.OperationWithName) error
ContainerOperation(req dto.ContainerOperation) error
DownloadContainerLogs(containerType, container, since, tail string, timestamp bool, c *gin.Context) error
ContainerStats(id string) (*dto.ContainerStats, error)
Inspect(req dto.InspectReq) (string, error)
DeleteNetwork(req dto.BatchDelete) error
CreateNetwork(req dto.NetworkCreate) error
DeleteVolume(req dto.BatchDelete) error
CreateVolume(req dto.VolumeCreate) error
Prune(req dto.ContainerPrune) error
LoadUsers(req dto.OperationWithName) []string
ListContainerFiles(req dto.ContainerFileReq) ([]dto.ContainerFileInfo, error)
UploadContainerFile(req dto.ContainerFileReq, fileName string, fileSize int64, file io.Reader) error
GetContainerFileContent(req dto.ContainerFileReq) (*dto.ContainerFileContent, error)
GetContainerFileSize(req dto.ContainerFileReq) (int64, error)
DeleteContainerFile(req dto.ContainerFileBatchDeleteReq) error
DownloadContainerFile(req dto.ContainerFileReq) (io.ReadCloser, string, string, error)
StreamLogs(ctx *gin.Context, params dto.StreamLog)
}
func NewIContainerService() IContainerService {
return &ContainerService{}
}
func (u *ContainerService) Page(req dto.PageContainer) (int64, interface{}, error) {
client, err := docker.NewDockerClient()
if err != nil {
return 0, nil, err
}
defer func() { _ = client.Close() }()
options := container.ListOptions{All: true}
if len(req.Filters) != 0 {
options.Filters = filters.NewArgs()
options.Filters.Add("label", req.Filters)
}
containers, err := client.ContainerList(context.Background(), options)
if err != nil {
return 0, nil, err
}
records := searchWithFilter(req, containers)
var backData []dto.ContainerInfo
total, start, end := len(records), (req.Page-1)*req.PageSize, req.Page*req.PageSize
if start > total {
backData = make([]dto.ContainerInfo, 0)
} else {
if end >= total {
end = total
}
backData = records[start:end]
}
for i := 0; i < len(backData); i++ {
install, _ := appInstallRepo.GetFirst(appInstallRepo.WithContainerName(backData[i].Name))
if install.ID > 0 {
backData[i].AppInstallName = install.Name
backData[i].AppName = install.App.Name
websites, _ := websiteRepo.GetBy(websiteRepo.WithAppInstallId(install.ID))
for _, website := range websites {
backData[i].Websites = append(backData[i].Websites, website.PrimaryDomain)
}
}
}
return int64(total), backData, nil
}
func (u *ContainerService) List() []dto.ContainerOptions {
var options []dto.ContainerOptions
client, err := docker.NewDockerClient()
if err != nil {
global.LOG.Errorf("load docker client for contianer list failed, err: %v", err)
return nil
}
defer client.Close()
containers, err := client.ContainerList(context.Background(), container.ListOptions{All: true})
if err != nil {
global.LOG.Errorf("load container list failed, err: %v", err)
return nil
}
for _, container := range containers {
for _, name := range container.Names {
if len(name) != 0 {
options = append(options, dto.ContainerOptions{Name: strings.TrimPrefix(name, "/"), State: container.State})
}
}
}
return options
}
func (u *ContainerService) ListByImage(imageName string) []dto.ContainerOptions {
var options []dto.ContainerOptions
client, err := docker.NewDockerClient()
if err != nil {
global.LOG.Errorf("load docker client for contianer list failed, err: %v", err)
return nil
}
defer client.Close()
containers, err := client.ContainerList(context.Background(), container.ListOptions{All: true})
if err != nil {
global.LOG.Errorf("load container list failed, err: %v", err)
return nil
}
for _, container := range containers {
if container.Image != imageName {
continue
}
for _, name := range container.Names {
if len(name) != 0 {
options = append(options, dto.ContainerOptions{Name: strings.TrimPrefix(name, "/"), State: container.State})
}
}
}
return options
}
func (u *ContainerService) LoadStatus() (dto.ContainerStatus, error) {
var data dto.ContainerStatus
client, err := docker.NewDockerClient()
if err != nil {
return data, err
}
defer client.Close()
c := context.Background()
images, _ := client.ImageList(c, image.ListOptions{All: true})
data.ImageCount = len(images)
repo, _ := imageRepoRepo.List()
data.RepoCount = len(repo)
templates, _ := composeRepo.List()
data.ComposeTemplateCount = len(templates)
networks, _ := client.NetworkList(c, network.ListOptions{})
data.NetworkCount = len(networks)
volumes, _ := client.VolumeList(c, volume.ListOptions{})
data.VolumeCount = len(volumes.Volumes)
data.ComposeCount = loadComposeCount(client)
containers, _ := client.ContainerList(c, container.ListOptions{All: true})
data.ContainerCount = len(containers)
for _, item := range containers {
switch item.State {
case "created":
data.Created++
case "running":
data.Running++
case "paused":
data.Paused++
case "restarting":
data.Restarting++
case "dead":
data.Dead++
case "exited":
data.Exited++
case "removing":
data.Removing++
}
}
return data, nil
}
func (u *ContainerService) ContainerItemStats(req dto.OperationWithName) (dto.ContainerItemStats, error) {
var data dto.ContainerItemStats
client, err := docker.NewDockerClient()
if err != nil {
return data, err
}
if req.Name != "system" {
defer client.Close()
containerInfo, _, err := client.ContainerInspectWithRaw(context.Background(), req.Name, true)
if err != nil {
return data, err
}
data.SizeRw = *containerInfo.SizeRw
data.SizeRootFs = *containerInfo.SizeRootFs
return data, nil
}
usage, err := client.DiskUsage(context.Background(), types.DiskUsageOptions{})
if err != nil {
return data, err
}
for _, item := range usage.Images {
data.ImageUsage += item.Size
if item.Containers < 1 {
data.ImageReclaimable += item.Size
}
}
for _, item := range usage.Containers {
data.ContainerUsage += item.SizeRw
if item.State != "running" {
data.ContainerReclaimable += item.SizeRw
}
}
for _, item := range usage.Volumes {
data.VolumeUsage += item.UsageData.Size
if item.UsageData.RefCount == 0 {
data.VolumeReclaimable += item.UsageData.Size
}
}
for _, item := range usage.BuildCache {
data.BuildCacheUsage += item.Size
}
return data, nil
}
func (u *ContainerService) ContainerListStats() ([]dto.ContainerListStats, error) {
client, err := docker.NewDockerClient()
if err != nil {
return nil, err
}
defer client.Close()
list, err := client.ContainerList(context.Background(), container.ListOptions{All: true})
if err != nil {
return nil, err
}
datas := make([]dto.ContainerListStats, len(list))
var wg sync.WaitGroup
wg.Add(len(list))
for i := 0; i < len(list); i++ {
go func(index int, item container.Summary) {
datas[index] = loadCpuAndMem(client, item.ID)
wg.Done()
}(i, list[i])
}
wg.Wait()
return datas, nil
}
func (u *ContainerService) Inspect(req dto.InspectReq) (string, error) {
client, err := docker.NewDockerClient()
if err != nil {
return "", err
}
defer client.Close()
var inspectInfo interface{}
switch req.Type {
case "container":
inspectInfo, err = client.ContainerInspect(context.Background(), req.ID)
case "compose":
filePath := ""
cli, err := docker.NewDockerClient()
if err != nil {
return "", err
}
defer cli.Close()
options := container.ListOptions{All: true}
options.Filters = filters.NewArgs()
options.Filters.Add("label", fmt.Sprintf("%s=%s", composeProjectLabel, req.ID))
containers, err := cli.ContainerList(context.Background(), options)
if err != nil {
return "", err
}
for _, container := range containers {
config := container.Labels[composeConfigLabel]
if len(req.Detail) != 0 && strings.Contains(config, req.Detail) {
config = req.Detail
}
workdir := container.Labels[composeWorkdirLabel]
if len(config) != 0 && len(workdir) != 0 && strings.Contains(config, workdir) {
filePath = config
break
} else {
filePath = workdir
break
}
}
if len(containers) == 0 {
composeItem, _ := composeRepo.GetRecord(repo.WithByName(req.ID))
filePath = composeItem.Path
}
if _, err := os.Stat(filePath); err != nil {
return "", err
}
content, err := os.ReadFile(filePath)
if err != nil {
return "", err
}
return string(content), nil
case "image":
inspectInfo, _, err = client.ImageInspectWithRaw(context.Background(), req.ID)
case "network":
inspectInfo, err = client.NetworkInspect(context.TODO(), req.ID, network.InspectOptions{})
case "volume":
inspectInfo, err = client.VolumeInspect(context.TODO(), req.ID)
}
if err != nil {
return "", err
}
bytes, err := json.Marshal(inspectInfo)
if err != nil {
return "", err
}
return string(bytes), nil
}
func (u *ContainerService) Prune(req dto.ContainerPrune) error {
client, err := docker.NewDockerClient()
if err != nil {
return err
}
defer client.Close()
name := ""
switch req.PruneType {
case "container":
name = "Container"
case "image":
name = "Image"
case "volume":
name = "Volume"
case "buildcache":
name = "BuildCache"
case "network":
name = "Network"
}
taskItem, err := task.NewTaskWithOps(i18n.GetMsgByKey(name), task.TaskClean, task.TaskScopeContainer, req.TaskID, 1)
if err != nil {
global.LOG.Errorf("new task for create container failed, err: %v", err)
return err
}
taskItem.AddSubTask(i18n.GetMsgByKey("TaskClean"), func(t *task.Task) error {
pruneFilters := filters.NewArgs()
if req.WithTagAll {
pruneFilters.Add("dangling", "false")
if req.PruneType != "image" {
pruneFilters.Add("until", "24h")
}
}
taskItem.Log(i18n.GetMsgByKey("PruneStart"))
SpaceReclaimed := 0
switch req.PruneType {
case "container":
rep, err := client.ContainersPrune(context.Background(), pruneFilters)
if err != nil {
return err
}
SpaceReclaimed = int(rep.SpaceReclaimed)
case "image":
rep, err := client.ImagesPrune(context.Background(), pruneFilters)
if err != nil {
return err
}
SpaceReclaimed = int(rep.SpaceReclaimed)
case "network":
_, err := client.NetworksPrune(context.Background(), pruneFilters)
if err != nil {
return err
}
case "volume":
versions, err := client.ServerVersion(context.Background())
if err != nil {
return err
}
if common.ComparePanelVersion(versions.APIVersion, "1.42") {
pruneFilters.Add("all", "true")
}
rep, err := client.VolumesPrune(context.Background(), pruneFilters)
if err != nil {
return err
}
SpaceReclaimed = int(rep.SpaceReclaimed)
case "buildcache":
opts := build.CachePruneOptions{}
opts.All = true
rep, err := client.BuildCachePrune(context.Background(), opts)
if err != nil {
return err
}
SpaceReclaimed = int(rep.SpaceReclaimed)
}
taskItem.Log(i18n.GetMsgWithMap("PruneHelper", map[string]interface{}{"name": i18n.GetMsgByKey(name), "size": common.LoadSizeUnit2F(float64(SpaceReclaimed))}))
return nil
}, nil)
go func() {
_ = taskItem.Execute()
}()
return nil
}
func (u *ContainerService) LoadResourceLimit() (*dto.ResourceLimit, error) {
cpuCounts, err := cpu.Counts(true)
if err != nil {
return nil, fmt.Errorf("load cpu limit failed, err: %v", err)
}
memoryInfo, err := mem.VirtualMemory()
if err != nil {
return nil, fmt.Errorf("load memory limit failed, err: %v", err)
}
data := dto.ResourceLimit{
CPU: cpuCounts,
Memory: memoryInfo.Total,
}
return &data, nil
}
func (u *ContainerService) ContainerCreate(req dto.ContainerOperate, inThread bool) error {
client, err := docker.NewDockerClient()
if err != nil {
return err
}
unlock := containerOperationLock.lock(req.Name)
ctx := context.Background()
newContainer, _ := client.ContainerInspect(ctx, req.Name)
if newContainer.ContainerJSONBase != nil {
unlock()
_ = client.Close()
return buserr.New("ErrContainerName")
}
taskItem, err := task.NewTaskWithOps(req.Name, task.TaskCreate, task.TaskScopeContainer, req.TaskID, 1)
if err != nil {
unlock()
_ = client.Close()
global.LOG.Errorf("new task for create container failed, err: %v", err)
return err
}
taskItem.AddSubTask(i18n.GetWithName("ContainerImagePull", req.Image), func(t *task.Task) error {
if !checkImageExist(client, req.Image) || req.ForcePull {
if err := pullImages(taskItem, client, req.Image); err != nil {
if !req.ForcePull {
return err
}
}
}
return nil
}, nil)
taskItem.AddSubTask(i18n.GetMsgByKey("ContainerImageCheck"), func(t *task.Task) error {
imageInfo, _, err := client.ImageInspectWithRaw(ctx, req.Image)
if err != nil {
return err
}
if len(req.Entrypoint) == 0 {
req.Entrypoint = imageInfo.Config.Entrypoint
}
if len(req.Cmd) == 0 {
req.Cmd = imageInfo.Config.Cmd
}
return nil
}, nil)
taskItem.AddSubTask(i18n.GetWithName("ContainerCreate", req.Name), func(t *task.Task) error {
config, hostConf, networkConf, err := loadConfigInfo(true, req, nil)
taskItem.LogWithStatus(i18n.GetMsgByKey("ContainerLoadInfo"), err)
if err != nil {
return err
}
normalizeContainerEndpointSettings(ctx, client, networkConf, nil)
con, err := client.ContainerCreate(ctx, config, hostConf, networkConf, &v1.Platform{}, req.Name)
if err != nil {
taskItem.Log(i18n.GetMsgByKey("ContainerCreateFailed"))
if con.ID != "" {
_ = client.ContainerRemove(ctx, con.ID, container.RemoveOptions{RemoveVolumes: true, Force: true})
}
return err
}
err = client.ContainerStart(ctx, con.ID, container.StartOptions{})
taskItem.LogWithStatus(i18n.GetMsgByKey("ContainerStartCheck"), err)
if err != nil {
taskItem.Log(i18n.GetMsgByKey("ContainerCreateFailed"))
_ = client.ContainerRemove(ctx, con.ID, container.RemoveOptions{RemoveVolumes: true, Force: true})
return fmt.Errorf("create successful but start failed, err: %v", err)
}
return nil
}, nil)
if inThread {
go func() {
defer unlock()
defer client.Close()
if err := taskItem.Execute(); err != nil {
global.LOG.Error(err.Error())
}
}()
return nil
}
defer unlock()
defer client.Close()
return taskItem.Execute()
}
func (u *ContainerService) ContainerInfo(req dto.OperationWithName) (*dto.ContainerOperate, error) {
client, err := docker.NewDockerClient()
if err != nil {
return nil, err
}
defer client.Close()
ctx := context.Background()
oldContainer, err := client.ContainerInspect(ctx, req.Name)
if err != nil {
return nil, err
}
var data dto.ContainerOperate
data.Name = strings.ReplaceAll(oldContainer.Name, "/", "")
data.Image = oldContainer.Config.Image
if oldContainer.NetworkSettings != nil {
for net, val := range oldContainer.NetworkSettings.Networks {
data.Networks = append(data.Networks, loadContainerNetworkInfo(net, val))
}
}
exposePorts, _ := loadPortByInspect(oldContainer.ID, client)
data.ExposedPorts = loadContainerPortForInfo(exposePorts)
data.Hostname = oldContainer.Config.Hostname
data.DNS = oldContainer.HostConfig.DNS
for _, item := range oldContainer.HostConfig.ExtraHosts {
parts := strings.SplitN(item, ":", 2)
if len(parts) != 2 || len(parts[0]) == 0 || len(parts[1]) == 0 {
continue
}
data.ExtraHosts = append(data.ExtraHosts, dto.ExtraHost{
Hostname: parts[0],
IP: parts[1],
})
}
data.DomainName = oldContainer.Config.Domainname
data.Cmd = oldContainer.Config.Cmd
data.WorkingDir = oldContainer.Config.WorkingDir
data.User = oldContainer.Config.User
data.OpenStdin = oldContainer.Config.OpenStdin
data.Tty = oldContainer.Config.Tty
data.Entrypoint = oldContainer.Config.Entrypoint
data.Env = oldContainer.Config.Env
data.CPUShares = oldContainer.HostConfig.CPUShares
for key, val := range oldContainer.Config.Labels {
data.Labels = append(data.Labels, fmt.Sprintf("%s=%s", key, val))
}
data.AutoRemove = oldContainer.HostConfig.AutoRemove
data.Privileged = oldContainer.HostConfig.Privileged
data.PublishAllPorts = oldContainer.HostConfig.PublishAllPorts
data.RestartPolicy = string(oldContainer.HostConfig.RestartPolicy.Name)
if oldContainer.HostConfig.NanoCPUs != 0 {
data.NanoCPUs = float64(oldContainer.HostConfig.NanoCPUs) / 1000000000
}
if oldContainer.HostConfig.Memory != 0 {
data.Memory = float64(oldContainer.HostConfig.Memory) / 1024 / 1024
}
data.Volumes = loadVolumeBinds(oldContainer.Mounts)
return &data, nil
}
func loadContainerNetworkInfo(name string, endpoint *network.EndpointSettings) dto.ContainerNetwork {
item := dto.ContainerNetwork{Network: name}
if endpoint == nil {
return item
}
item.MacAddr = endpoint.MacAddress
item.Links = append([]string(nil), endpoint.Links...)
item.Aliases = append([]string(nil), endpoint.Aliases...)
item.DriverOpts = cloneStringMap(endpoint.DriverOpts)
item.GwPriority = endpoint.GwPriority
if endpoint.IPAMConfig != nil {
item.LinkLocalIPs = append([]string(nil), endpoint.IPAMConfig.LinkLocalIPs...)
}
if name != "bridge" {
if endpoint.IPAMConfig != nil {
item.Ipv4 = endpoint.IPAMConfig.IPv4Address
item.Ipv6 = endpoint.IPAMConfig.IPv6Address
} else {
item.Ipv4 = endpoint.IPAddress
item.Ipv6 = endpoint.GlobalIPv6Address
}
}
return item
}
func cloneStringMap(source map[string]string) map[string]string {
if len(source) == 0 {
return nil
}
result := make(map[string]string, len(source))
for key, value := range source {
result[key] = value
}
return result
}
func (u *ContainerService) ContainerRename(req dto.ContainerRename) error {
ctx := context.Background()
client, err := docker.NewDockerClient()
if err != nil {
return err
}
defer client.Close()
unlock := containerOperationLock.lock(req.Name, req.NewName)
defer unlock()
newContainer, _ := client.ContainerInspect(ctx, req.NewName)
if newContainer.ContainerJSONBase != nil {
return buserr.New("ErrContainerName")
}
return client.ContainerRename(ctx, req.Name, req.NewName)
}
func (u *ContainerService) ContainerCommit(req dto.ContainerCommit) error {
ctx := context.Background()
client, err := docker.NewDockerClient()
if err != nil {
return err
}
defer client.Close()
options := container.CommitOptions{
Reference: req.NewImageName,
Comment: req.Comment,
Author: req.Author,
Changes: nil,
Pause: req.Pause,
Config: nil,
}
taskItem, err := task.NewTaskWithOps(req.NewImageName, task.TaskCommit, task.TaskScopeContainer, req.TaskID, 1)
if err != nil {
return fmt.Errorf("new task for container commit failed, err: %v", err)
}
go func() {
taskItem.AddSubTask(i18n.GetWithName("TaskCommit", req.NewImageName), func(t *task.Task) error {
res, err := client.ContainerCommit(ctx, req.ContainerId, options)
if err != nil {
return fmt.Errorf("failed to commit container, err: %v", err)
}
taskItem.Log(res.ID)
return nil
}, nil)
_ = taskItem.Execute()
}()
return nil
}
func (u *ContainerService) ContainerOperation(req dto.ContainerOperation) error {
ctx := context.Background()
client, err := docker.NewDockerClient()
if err != nil {
return err
}
taskItem, err := task.NewTaskWithOps(strings.Join(req.Names, " "), req.Operation, task.TaskScopeContainer, req.TaskID, 1)
if err != nil {
_ = client.Close()
return fmt.Errorf("new task for container commit failed, err: %v", err)
}
for _, item := range req.Names {
item := item
taskItem.AddSubTask(item, func(t *task.Task) error {
unlock := containerOperationLock.lock(item)
defer unlock()
var operationErr error
switch req.Operation {
case constant.ContainerOpStart:
operationErr = client.ContainerStart(ctx, item, container.StartOptions{})
case constant.ContainerOpStop:
operationErr = client.ContainerStop(ctx, item, container.StopOptions{})
case constant.ContainerOpRestart:
operationErr = client.ContainerRestart(ctx, item, container.StopOptions{})
case constant.ContainerOpKill:
operationErr = client.ContainerKill(ctx, item, "SIGKILL")
case constant.ContainerOpPause:
operationErr = client.ContainerPause(ctx, item)
case constant.ContainerOpUnpause:
operationErr = client.ContainerUnpause(ctx, item)
case constant.ContainerOpRemove:
operationErr = client.ContainerRemove(ctx, item, container.RemoveOptions{RemoveVolumes: true, Force: true})
}
return operationErr
}, nil)
}
go func() {
defer client.Close()
_ = taskItem.Execute()
}()
return nil
}
func (u *ContainerService) ContainerLogClean(req dto.OperationWithName) error {
client, err := docker.NewDockerClient()
if err != nil {
return err
}
defer client.Close()
unlock := containerOperationLock.lock(req.Name)
defer unlock()
ctx := context.Background()
containerItem, err := client.ContainerInspect(ctx, req.Name)
if err != nil {
return err
}
if err := client.ContainerStop(ctx, containerItem.ID, container.StopOptions{}); err != nil {
return err
}
file, err := os.OpenFile(containerItem.LogPath, os.O_RDWR|os.O_CREATE, constant.FilePerm)
if err != nil {
return err
}
defer file.Close()
if err = file.Truncate(0); err != nil {
return err
}
_, _ = file.Seek(0, 0)
files, _ := filepath.Glob(fmt.Sprintf("%s.*", containerItem.LogPath))
for _, file := range files {
_ = os.Remove(file)
}
if err := client.ContainerStart(ctx, containerItem.ID, container.StartOptions{}); err != nil {
return err
}
return nil
}
func (u *ContainerService) StreamLogs(ctx *gin.Context, params dto.StreamLog) {
messageChan := make(chan string, 1024)
errorChan := make(chan error, 1)
doneChan := make(chan struct{})
go func() {
<-ctx.Request.Context().Done()
close(doneChan)
}()
go collectLogs(doneChan, params, messageChan, errorChan)
ctx.Stream(func(w io.Writer) bool {
select {
case msg, ok := <-messageChan:
if !ok {
return false
}
_, err := fmt.Fprintf(w, "data: %s\n\n", msg)
if err != nil {
return false
}
return true
case err := <-errorChan:
if err != nil {
_, _ = fmt.Fprintf(w, "event: error\ndata: %v\n\n", err.Error())
}
return false
}
})
}
func collectLogs(done <-chan struct{}, params dto.StreamLog, messageChan chan<- string, errorChan chan<- error) {
defer close(messageChan)
defer close(errorChan)
var cmdArgs []string
cmdArgs = append(cmdArgs, "logs")
if params.Follow {
cmdArgs = append(cmdArgs, "-f")
}
if params.Timestamp {
cmdArgs = append(cmdArgs, "-t")
}
if params.Tail != "0" {
cmdArgs = append(cmdArgs, "--tail", params.Tail)
}
if params.Since != "all" {
cmdArgs = append(cmdArgs, "--since", params.Since)
}
if params.Container != "" {
cmdArgs = append(cmdArgs, params.Container)
}
var dockerCmd *exec.Cmd
if params.Type == "compose" {
dockerComposeCmd := common.GetDockerComposeCommand()
var yamlFiles []string
for _, item := range strings.Split(params.Compose, ",") {
if len(item) != 0 {
yamlFiles = append(yamlFiles, "-f", item)
}
}
if dockerComposeCmd == "docker-compose" {
newCmdArgs := append(yamlFiles, cmdArgs...)
dockerCmd = exec.Command(dockerComposeCmd, newCmdArgs...)
} else {
newCmdArgs := append(append([]string{"compose"}, yamlFiles...), cmdArgs...)
dockerCmd = exec.Command("docker", newCmdArgs...)
}
} else {
dockerCmd = exec.Command("docker", cmdArgs...)
}
dockerCmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true}
stdout, err := dockerCmd.StdoutPipe()
if err != nil {
errorChan <- fmt.Errorf("failed to get stdout pipe: %v", err)
return
}
dockerCmd.Stderr = dockerCmd.Stdout
if err = dockerCmd.Start(); err != nil {
errorChan <- fmt.Errorf("failed to start docker logs command: %v", err)
return
}
defer func() {
if dockerCmd.Process != nil {
if pgid, err := syscall.Getpgid(dockerCmd.Process.Pid); err == nil {
_ = syscall.Kill(-pgid, syscall.SIGKILL)
}
_ = dockerCmd.Process.Kill()
_ = dockerCmd.Wait()
}
}()
reader := bufio.NewReader(stdout)
processKilled := false
go func() {
<-done
if !processKilled && dockerCmd.Process != nil {
processKilled = true
if pgid, err := syscall.Getpgid(dockerCmd.Process.Pid); err == nil {
_ = syscall.Kill(-pgid, syscall.SIGKILL)
}
_ = dockerCmd.Process.Kill()
}
}()
for {
line, err := reader.ReadString('\n')
if err != nil {
if err == io.EOF {
if len(line) > 0 {
line = strings.TrimSuffix(line, "\n")
select {
case messageChan <- line:
case <-done:
return
}
}
break
}
errorChan <- fmt.Errorf("reader error: %v", err)
return
}
line = strings.TrimSuffix(line, "\n")
select {
case messageChan <- line:
case <-done:
return
}
}
_ = dockerCmd.Wait()
}
func (u *ContainerService) DownloadContainerLogs(containerType, container, since, tail string, timestamp bool, c *gin.Context) error {
if cmd.CheckIllegal(container, since, tail) {
return buserr.New("ErrCmdIllegal")
}
ctx := c.Request.Context()
commandArg := []string{"logs", container}
dockerCommand := global.CONF.DockerConfig.Command
if containerType == "compose" {
var yamlFiles []string
for _, item := range strings.Split(container, ",") {
if len(item) != 0 {
yamlFiles = append(yamlFiles, "-f", item)
}
}
if dockerCommand == "docker-compose" {
commandArg = append(yamlFiles, "logs")
} else {
commandArg = append(append([]string{"compose"}, yamlFiles...), "logs")
}
}
if tail != "0" {
commandArg = append(commandArg, "--tail")
commandArg = append(commandArg, tail)
}
if since != "all" {
commandArg = append(commandArg, "--since")
commandArg = append(commandArg, since)
}
if timestamp {
commandArg = append(commandArg, "-t")
}
var dockerCmd *exec.Cmd
if containerType == "compose" && dockerCommand == "docker-compose" {
dockerCmd = exec.CommandContext(ctx, "docker-compose", commandArg...)
} else {
dockerCmd = exec.CommandContext(ctx, "docker", commandArg...)
}
dockerCmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true}
stdout, err := dockerCmd.StdoutPipe()
if err != nil {
return err
}
dockerCmd.Stderr = dockerCmd.Stdout
if err := dockerCmd.Start(); err != nil {
return err
}
done := make(chan struct{})
go func() {
select {
case <-ctx.Done():
killContainerLogProcess(dockerCmd)
case <-done:
}
}()
defer close(done)
tempFile, err := os.CreateTemp("", "cmd_output_*.txt")
if err != nil {
killContainerLogProcess(dockerCmd)
_ = dockerCmd.Wait()
return err
}
defer tempFile.Close()
defer func() {
if err := os.Remove(tempFile.Name()); err != nil {
global.LOG.Errorf("os.Remove() failed: %v", err)
}
}()
copyErr := copyContainerLogOutput(tempFile, stdout)
waitErr := dockerCmd.Wait()
if copyErr != nil {
return copyErr
}
if waitErr != nil {
if ctx.Err() != nil {
return ctx.Err()
}
return waitErr
}
if _, err := tempFile.Seek(0, io.SeekStart); err != nil {
return err
}
info, _ := tempFile.Stat()
c.Header("Content-Length", strconv.FormatInt(info.Size(), 10))
c.Header("Content-Disposition", "attachment; filename*=utf-8''"+url.PathEscape(info.Name()))
http.ServeContent(c.Writer, c.Request, info.Name(), info.ModTime(), tempFile)
return nil
}
func copyContainerLogOutput(dst io.Writer, src io.Reader) error {
reader := bufio.NewReader(src)
for {
line, err := reader.ReadString('\n')
if len(line) > 0 {
cleanLine := containerLogAnsiRegex.ReplaceAllString(line, "")
if _, writeErr := io.WriteString(dst, cleanLine); writeErr != nil {
return writeErr
}
}
if err != nil {
if err == io.EOF {
return nil
}
return err
}
}
}
func killContainerLogProcess(command *exec.Cmd) {
if command == nil || command.Process == nil {
return
}
if pgid, err := syscall.Getpgid(command.Process.Pid); err == nil {
_ = syscall.Kill(-pgid, syscall.SIGKILL)
return
}
_ = command.Process.Kill()
}
func (u *ContainerService) ContainerStats(id string) (*dto.ContainerStats, error) {
client, err := docker.NewDockerClient()
if err != nil {
return nil, err
}
defer client.Close()
res, err := client.ContainerStats(context.TODO(), id, false)
if err != nil {
return nil, err
}
defer res.Body.Close()
body, err := io.ReadAll(res.Body)
if err != nil {
return nil, err
}
var stats *container.StatsResponse
if err := json.Unmarshal(body, &stats); err != nil {
return nil, err
}
var data dto.ContainerStats
data.CPUPercent = calculateCPUPercentUnix(stats)
data.IORead, data.IOWrite = calculateBlockIO(stats.BlkioStats)
data.Memory = float64(stats.MemoryStats.Usage) / 1024 / 1024
if cache, ok := stats.MemoryStats.Stats["cache"]; ok {
data.Cache = float64(cache) / 1024 / 1024
}
data.NetworkRX, data.NetworkTX = calculateNetwork(stats.Networks)
data.ShotTime = stats.Read
return &data, nil
}
func (u *ContainerService) LoadUsers(req dto.OperationWithName) []string {
var users []string
std, err := cmd.RunDockerExecWithStdout(20*time.Second, req.Name, "cat", "/etc/passwd")
if err != nil {
return users
}
lines := strings.Split(string(std), "\n")
for _, line := range lines {
if strings.Contains(line, ":") {
users = append(users, strings.Split(line, ":")[0])
}
}
return users
}
func (u *ContainerService) ListContainerFiles(req dto.ContainerFileReq) ([]dto.ContainerFileInfo, error) {
if len(req.Path) == 0 {
req.Path = "/"
}
cli, err := docker.NewDockerClient()
if err != nil {
return nil, err
}
defer cli.Close()
ctx := context.Background()
stat, err := cli.ContainerStatPath(ctx, req.ContainerID, req.Path)
if err != nil {
return nil, normalizeContainerFileError(err)
}
isDir := stat.Mode.IsDir()
isLink := stat.Mode&os.ModeSymlink != 0
if isLink && !isDir {
linkDir, linkErr := isContainerDir(cli, req.ContainerID, req.Path)
if linkErr == nil {
isDir = linkDir
}
}
if !isDir {
return []dto.ContainerFileInfo{toContainerFileInfo(req.Path, stat, isDir)}, nil
}
output, err := runContainerCommand(cli, req.ContainerID, []string{"ls", "-1A", "--", req.Path})
if err != nil {
return nil, normalizeContainerFileError(err)
}
lines := strings.Split(strings.TrimSpace(output), "\n")
files := make([]dto.ContainerFileInfo, 0, len(lines))
for _, line := range lines {
name := strings.TrimSpace(line)
if len(name) == 0 || name == "." || name == ".." {
continue
}
childPath := req.Path
if childPath == "/" {
childPath = "/" + name
} else {
childPath = strings.TrimSuffix(childPath, "/") + "/" + name
}
childStat, statErr := cli.ContainerStatPath(ctx, req.ContainerID, childPath)
if statErr != nil {
continue
}
childIsDir := childStat.Mode.IsDir()
if childStat.Mode&os.ModeSymlink != 0 && !childIsDir {
linkDir, linkErr := isContainerDir(cli, req.ContainerID, childPath)
if linkErr == nil {
childIsDir = linkDir
}
}
files = append(files, toContainerFileInfo(childPath, childStat, childIsDir))
}
sort.Slice(files, func(i, j int) bool {
if files[i].IsDir != files[j].IsDir {
return files[i].IsDir
}
return strings.ToLower(files[i].Name) < strings.ToLower(files[j].Name)
})
return files, nil
}
func (u *ContainerService) DeleteContainerFile(req dto.ContainerFileBatchDeleteReq) error {
for _, item := range req.Paths {
if strings.TrimSpace(item) == "/" {
return buserr.New("ErrPathNotDelete")
}
}
cli, err := docker.NewDockerClient()
if err != nil {
return err
}
defer cli.Close()
command := []string{"rm", "-rf", "--"}
command = append(command, req.Paths...)
_, err = runContainerCommand(cli, req.ContainerID, command)
return err
}
func (u *ContainerService) UploadContainerFile(req dto.ContainerFileReq, fileName string, fileSize int64, file io.Reader) error {
if len(req.Path) == 0 {
req.Path = "/"
}
safeName := path.Base(fileName)
if safeName == "." || safeName == "/" || len(safeName) == 0 {
return buserr.New("ErrInvalidChar")
}
cli, err := docker.NewDockerClient()
if err != nil {
return err
}
defer cli.Close()
ctx := context.Background()
stat, err := cli.ContainerStatPath(ctx, req.ContainerID, req.Path)
if err != nil {
if _, mkErr := runContainerCommand(cli, req.ContainerID, []string{"mkdir", "-p", "--", req.Path}); mkErr != nil {
return mkErr
}
stat, err = cli.ContainerStatPath(ctx, req.ContainerID, req.Path)
if err != nil {
return err
}
}
if !stat.Mode.IsDir() {
return fmt.Errorf("path %s is not directory", req.Path)
}
pipeReader, pipeWriter := io.Pipe()
writeErr := make(chan error, 1)
go func() {
tw := tar.NewWriter(pipeWriter)
header := &tar.Header{
Name: safeName,
Mode: 0644,
Size: fileSize,
ModTime: time.Now(),
}
if err := tw.WriteHeader(header); err != nil {
_ = tw.Close()
_ = pipeWriter.CloseWithError(err)
writeErr <- err
return
}
if _, err := io.Copy(tw, file); err != nil {
_ = tw.Close()
_ = pipeWriter.CloseWithError(err)
writeErr <- err
return
}
if err := tw.Close(); err != nil {
_ = pipeWriter.CloseWithError(err)
writeErr <- err
return
}
_ = pipeWriter.Close()
writeErr <- nil
}()
err = cli.CopyToContainer(ctx, req.ContainerID, req.Path, pipeReader, container.CopyToContainerOptions{
CopyUIDGID: true,
})
if err != nil {
_ = pipeReader.CloseWithError(err)
_ = pipeWriter.CloseWithError(err)
<-writeErr
return err
}
if err := <-writeErr; err != nil {
return err
}
return nil
}
func (u *ContainerService) GetContainerFileContent(req dto.ContainerFileReq) (*dto.ContainerFileContent, error) {
if len(req.Path) == 0 {
return nil, buserr.New("ErrInvalidChar")
}
cli, err := docker.NewDockerClient()
if err != nil {
return nil, err
}
defer cli.Close()
stat, err := cli.ContainerStatPath(context.Background(), req.ContainerID, req.Path)
if err != nil {
return nil, normalizeContainerFileError(err)
}
if stat.Mode.IsDir() {
return nil, fmt.Errorf("path %s is directory", req.Path)
}
content := &dto.ContainerFileContent{Size: stat.Size}
headBytes, err := runContainerCommandRaw(cli, req.ContainerID, []string{"head", "-c", "4096", "--", req.Path})
if err != nil {
return nil, err
}
if bytes.IndexByte(headBytes, 0) >= 0 {
content.IsBinary = true
return content, nil
}
const inlinePreviewMax = 512 * 1024
if stat.Size <= inlinePreviewMax {
raw, err := runContainerCommandRaw(cli, req.ContainerID, []string{"cat", "--", req.Path})
if err != nil {
return nil, err
}
content.Content = string(raw)
return content, nil
}
raw, err := runContainerCommandRaw(cli, req.ContainerID, []string{"tail", "-n", "300", "--", req.Path})
if err != nil {
return nil, err
}
content.Content = string(raw)
content.Truncated = true
return content, nil
}
func (u *ContainerService) GetContainerFileSize(req dto.ContainerFileReq) (int64, error) {
if len(req.Path) == 0 {
return 0, buserr.New("ErrInvalidChar")
}
cli, err := docker.NewDockerClient()
if err != nil {
return 0, err
}
defer cli.Close()
stat, err := cli.ContainerStatPath(context.Background(), req.ContainerID, req.Path)
if err != nil {
return 0, normalizeContainerFileError(err)
}
if !stat.Mode.IsDir() {
return stat.Size, nil
}
output, err := runContainerCommand(cli, req.ContainerID, []string{"du", "-sb", "--", req.Path})
if err != nil {
return 0, err
}
parts := strings.Fields(output)
if len(parts) == 0 {
return 0, fmt.Errorf("invalid du output")
}
size, err := strconv.ParseInt(parts[0], 10, 64)
if err != nil {
return 0, err
}
return size, nil
}
func (u *ContainerService) DownloadContainerFile(req dto.ContainerFileReq) (io.ReadCloser, string, string, error) {
if len(req.Path) == 0 {
req.Path = "/"
}
cli, err := docker.NewDockerClient()
if err != nil {
return nil, "", "", err
}
ctx := context.Background()
stat, err := cli.ContainerStatPath(ctx, req.ContainerID, req.Path)
if err != nil {
_ = cli.Close()
return nil, "", "", normalizeContainerFileError(err)
}
fileName := stat.Name
if len(fileName) == 0 {
fileName = "container-file"
}
if stat.Mode.IsDir() {
if _, err := runContainerCommand(cli, req.ContainerID, []string{"tar", "--help"}); err != nil {
_ = cli.Close()
return nil, "", "", fmt.Errorf("tar command not found in container")
}
targetPath := path.Clean(req.Path)
parentPath := path.Dir(targetPath)
targetName := path.Base(targetPath)
if parentPath == "." || parentPath == "" {
parentPath = "/"
}
tarStream, err := runContainerCommandStream(cli, req.ContainerID, []string{
"tar", "-czf", "-", "-C", parentPath, "--", targetName,
})
if err != nil {
_ = cli.Close()
return nil, "", "", err
}
if !strings.HasSuffix(fileName, ".tar.gz") {
fileName += ".tar.gz"
}
return &closeHookReader{
ReadCloser: tarStream,
onClose: cli.Close,
}, fileName, "application/gzip", nil
}
fileStream, err := runContainerCommandStream(cli, req.ContainerID, []string{"cat", "--", req.Path})
if err != nil {
_ = cli.Close()
return nil, "", "", err
}
return &closeHookReader{
ReadCloser: fileStream,
onClose: cli.Close,
}, fileName, "application/octet-stream", nil
}
func normalizeContainerFileError(err error) error {
if err == nil {
return nil
}
message := strings.ToLower(err.Error())
if strings.Contains(message, "no such file or directory") || strings.Contains(message, "not found") {
return buserr.New("ErrPathNotFound")
}
return err
}
func runContainerCommand(cli *client.Client, containerID string, command []string) (string, error) {
raw, err := runContainerCommandRaw(cli, containerID, command)
if err != nil {
return "", err
}
return strings.TrimSpace(string(raw)), nil
}
type closeHookReader struct {
io.ReadCloser
onClose func() error
}
func (r *closeHookReader) Close() error {
var closeErr error
if r.ReadCloser != nil {
closeErr = r.ReadCloser.Close()
}
if r.onClose != nil {
if err := r.onClose(); err != nil && closeErr == nil {
closeErr = err
}
}
return closeErr
}
func runContainerCommandRaw(cli *client.Client, containerID string, command []string) ([]byte, error) {
ctx := context.Background()
resp, err := cli.ContainerExecCreate(ctx, containerID, container.ExecOptions{
Cmd: command,
AttachStdout: true,
AttachStderr: true,
})
if err != nil {
return nil, err
}
hijack, err := cli.ContainerExecAttach(ctx, resp.ID, container.ExecAttachOptions{})
if err != nil {
return nil, err
}
defer hijack.Close()
raw, err := io.ReadAll(hijack.Reader)
if err != nil {
return nil, err
}
var stdout bytes.Buffer
var stderr bytes.Buffer
if _, err := stdcopy.StdCopy(&stdout, &stderr, bytes.NewReader(raw)); err != nil {
return nil, err
}
output := strings.TrimSpace(stdout.String())
errorOutput := strings.TrimSpace(stderr.String())
info, err := cli.ContainerExecInspect(ctx, resp.ID)
if err != nil {
return nil, err
}
if info.ExitCode != 0 {
if len(errorOutput) != 0 {
return nil, fmt.Errorf("%s", errorOutput)
}
if len(output) == 0 {
return nil, fmt.Errorf("container command failed with exit code %d", info.ExitCode)
}
return nil, fmt.Errorf("%s", output)
}
return stdout.Bytes(), nil
}
func runContainerCommandStream(cli *client.Client, containerID string, command []string) (io.ReadCloser, error) {
ctx := context.Background()
resp, err := cli.ContainerExecCreate(ctx, containerID, container.ExecOptions{
Cmd: command,
AttachStdout: true,
AttachStderr: true,
})
if err != nil {
return nil, err
}
hijack, err := cli.ContainerExecAttach(ctx, resp.ID, container.ExecAttachOptions{})
if err != nil {
return nil, err
}
pipeReader, pipeWriter := io.Pipe()
go func() {
defer hijack.Close()
var stderr bytes.Buffer
_, copyErr := stdcopy.StdCopy(pipeWriter, &stderr, hijack.Reader)
if copyErr != nil {
_ = pipeWriter.CloseWithError(copyErr)
return
}
info, inspectErr := cli.ContainerExecInspect(ctx, resp.ID)
if inspectErr != nil {
_ = pipeWriter.CloseWithError(inspectErr)
return
}
if info.ExitCode != 0 {
msg := strings.TrimSpace(stderr.String())
if len(msg) == 0 {
msg = fmt.Sprintf("container command failed with exit code %d", info.ExitCode)
}
_ = pipeWriter.CloseWithError(fmt.Errorf("%s", msg))
return
}
_ = pipeWriter.Close()
}()
return pipeReader, nil
}
func toContainerFileInfo(filePath string, stat container.PathStat, isDir bool) dto.ContainerFileInfo {
name := stat.Name
if len(name) == 0 {
items := strings.Split(strings.TrimSuffix(filePath, "/"), "/")
name = items[len(items)-1]
}
isLink := stat.Mode&os.ModeSymlink != 0
return dto.ContainerFileInfo{
Name: name,
Path: filePath,
IsDir: isDir,
IsLink: isLink,
LinkTo: stat.LinkTarget,
Size: stat.Size,
Mode: stat.Mode.String(),
ModTime: stat.Mtime.Format(constant.DateTimeLayout),
}
}
func isContainerDir(cli *client.Client, containerID, targetPath string) (bool, error) {
checkPath := strings.TrimSuffix(targetPath, "/") + "/."
_, err := runContainerCommand(cli, containerID, []string{"ls", "-d", "--", checkPath})
if err != nil {
return false, err
}
return true, nil
}
func stringsToMap(list []string) map[string]string {
var labelMap = make(map[string]string)
for _, label := range list {
if strings.Contains(label, "=") {
sps := strings.SplitN(label, "=", 2)
labelMap[sps[0]] = sps[1]
}
}
return labelMap
}
func stringsToMap2(list []string) map[string]*string {
var labelMap = make(map[string]*string)
for _, label := range list {
if strings.Contains(label, "=") {
sps := strings.SplitN(label, "=", 2)
labelMap[sps[0]] = &sps[1]
}
}
return labelMap
}
func calculateCPUPercentUnix(stats *container.StatsResponse) float64 {
cpuPercent := 0.0
cpuDelta := float64(stats.CPUStats.CPUUsage.TotalUsage) - float64(stats.PreCPUStats.CPUUsage.TotalUsage)
systemDelta := float64(stats.CPUStats.SystemUsage) - float64(stats.PreCPUStats.SystemUsage)
if systemDelta > 0.0 && cpuDelta > 0.0 {
cpuPercent = (cpuDelta / systemDelta) * 100.0
if len(stats.CPUStats.CPUUsage.PercpuUsage) != 0 {
cpuPercent = cpuPercent * float64(len(stats.CPUStats.CPUUsage.PercpuUsage))
}
}
return cpuPercent
}
func calculateMemPercentUnix(memStats container.MemoryStats) float64 {
memPercent := 0.0
memUsage := calculateMemUsageUnixNoCache(memStats)
memLimit := float64(memStats.Limit)
if memUsage > 0.0 && memLimit > 0.0 {
memPercent = (float64(memUsage) / float64(memLimit)) * 100.0
}
return memPercent
}
func calculateMemUsageUnixNoCache(mem container.MemoryStats) uint64 {
if v, isCgroup1 := mem.Stats["total_inactive_file"]; isCgroup1 && v < mem.Usage {
return mem.Usage - v
}
if v := mem.Stats["inactive_file"]; v < mem.Usage {
return mem.Usage - v
}
return mem.Usage
}
func calculateBlockIO(blkio container.BlkioStats) (blkRead float64, blkWrite float64) {
for _, bioEntry := range blkio.IoServiceBytesRecursive {
switch strings.ToLower(bioEntry.Op) {
case "read":
blkRead = (blkRead + float64(bioEntry.Value)) / 1024 / 1024
case "write":
blkWrite = (blkWrite + float64(bioEntry.Value)) / 1024 / 1024
}
}
return
}
func calculateNetwork(network map[string]container.NetworkStats) (float64, float64) {
var rx, tx float64
for _, v := range network {
rx += float64(v.RxBytes) / 1024
tx += float64(v.TxBytes) / 1024
}
return rx, tx
}
func checkImageExist(client *client.Client, imageItem string) bool {
if client == nil {
var err error
client, err = docker.NewDockerClient()
if err != nil {
return false
}
}
images, err := client.ImageList(context.Background(), image.ListOptions{})
if err != nil {
return false
}
for _, img := range images {
for _, tag := range img.RepoTags {
if tag == imageItem || tag == imageItem+":latest" {
return true
}
}
}
return false
}
func checkImageLike(client *client.Client, imageName string) bool {
if client == nil {
var err error
client, err = docker.NewDockerClient()
if err != nil {
return false
}
}
images, err := client.ImageList(context.Background(), image.ListOptions{})
if err != nil {
return false
}
for _, img := range images {
for _, tag := range img.RepoTags {
parts := strings.Split(tag, "/")
imageNameWithTag := parts[len(parts)-1]
if imageNameWithTag == imageName {
return true
}
}
}
return false
}
func pullImages(task *task.Task, client *client.Client, imageName string) error {
dockerCli := docker.NewClientWithExist(client)
options := image.PullOptions{}
repos, _ := imageRepoRepo.List()
if len(repos) != 0 {
for _, repo := range repos {
if strings.HasPrefix(imageName, repo.DownloadUrl) && repo.Auth {
authConfig := registry.AuthConfig{
Username: repo.Username,
Password: repo.Password,
}
encodedJSON, err := json.Marshal(authConfig)
if err != nil {
return err
}
authStr := base64.URLEncoding.EncodeToString(encodedJSON)
options.RegistryAuth = authStr
}
}
} else {
hasAuth, authStr := loadAuthInfo(imageName)
if hasAuth {
options.RegistryAuth = authStr
}
}
return dockerCli.PullImageWithProcessAndOptions(task, imageName, options)
}
func loadCpuAndMem(client *client.Client, containerItem string) dto.ContainerListStats {
data := dto.ContainerListStats{
ContainerID: containerItem,
}
res, err := client.ContainerStats(context.Background(), containerItem, false)
if err != nil {
return data
}
defer res.Body.Close()
body, err := io.ReadAll(res.Body)
if err != nil {
return data
}
var stats *container.StatsResponse
if err := json.Unmarshal(body, &stats); err != nil {
return data
}
data.CPUTotalUsage = stats.CPUStats.CPUUsage.TotalUsage - stats.PreCPUStats.CPUUsage.TotalUsage
data.SystemUsage = stats.CPUStats.SystemUsage - stats.PreCPUStats.SystemUsage
data.CPUPercent = calculateCPUPercentUnix(stats)
data.PercpuUsage = len(stats.CPUStats.CPUUsage.PercpuUsage)
data.MemoryCache = stats.MemoryStats.Stats["cache"]
data.MemoryUsage = calculateMemUsageUnixNoCache(stats.MemoryStats)
data.MemoryLimit = stats.MemoryStats.Limit
data.MemoryPercent = calculateMemPercentUnix(stats.MemoryStats)
return data
}
func checkPortStats(ports []dto.PortHelper, checkInUse bool) (nat.PortMap, error) {
portMap := make(nat.PortMap)
if len(ports) == 0 {
return portMap, nil
}
for _, port := range ports {
if strings.Contains(port.ContainerPort, "-") {
if !strings.Contains(port.HostPort, "-") {
return portMap, buserr.New("ErrPortRules")
}
hostStart, _ := strconv.Atoi(strings.Split(port.HostPort, "-")[0])
hostEnd, _ := strconv.Atoi(strings.Split(port.HostPort, "-")[1])
containerStart, _ := strconv.Atoi(strings.Split(port.ContainerPort, "-")[0])
containerEnd, _ := strconv.Atoi(strings.Split(port.ContainerPort, "-")[1])
if (hostEnd-hostStart) <= 0 || (containerEnd-containerStart) <= 0 {
return portMap, buserr.New("ErrPortRules")
}
if (containerEnd - containerStart) != (hostEnd - hostStart) {
return portMap, buserr.New("ErrPortRules")
}
for i := 0; i <= hostEnd-hostStart; i++ {
bindItem := nat.PortBinding{HostPort: strconv.Itoa(hostStart + i), HostIP: port.HostIP}
portMap[nat.Port(fmt.Sprintf("%d/%s", containerStart+i, port.Protocol))] = []nat.PortBinding{bindItem}
}
for i := hostStart; i <= hostEnd; i++ {
if checkInUse && common.ScanPortWithIP(port.HostIP, i) {
return portMap, buserr.WithDetail("ErrPortInUsed", i, nil)
}
}
} else {
portItem := 0
if strings.Contains(port.HostPort, "-") {
portItem, _ = strconv.Atoi(strings.Split(port.HostPort, "-")[0])
} else {
portItem, _ = strconv.Atoi(port.HostPort)
}
if checkInUse && common.ScanPortWithIP(port.HostIP, portItem) {
return portMap, buserr.WithDetail("ErrPortInUsed", portItem, nil)
}
bindItem := nat.PortBinding{HostPort: strconv.Itoa(portItem), HostIP: port.HostIP}
portMap[nat.Port(fmt.Sprintf("%s/%s", port.ContainerPort, port.Protocol))] = []nat.PortBinding{bindItem}
}
}
return portMap, nil
}
func loadConfigInfo(isCreate bool, req dto.ContainerOperate, oldContainer *container.InspectResponse) (*container.Config, *container.HostConfig, *network.NetworkingConfig, error) {
var config container.Config
var hostConf container.HostConfig
if !isCreate {
config = *oldContainer.Config
hostConf = *oldContainer.HostConfig
}
var networkConf network.NetworkingConfig
portMap, err := checkPortStats(req.ExposedPorts, isCreate)
if err != nil {
return nil, nil, nil, err
}
exposed := make(nat.PortSet)
for port := range portMap {
exposed[port] = struct{}{}
}
config.Image = req.Image
config.Cmd = req.Cmd
config.Entrypoint = req.Entrypoint
config.Env = req.Env
config.Labels = stringsToMap(req.Labels)
config.ExposedPorts = exposed
config.OpenStdin = req.OpenStdin
config.Tty = req.Tty
config.Hostname = req.Hostname
config.Domainname = req.DomainName
config.User = req.User
config.WorkingDir = req.WorkingDir
if len(req.Networks) != 0 {
networkConf.EndpointsConfig = make(map[string]*network.EndpointSettings)
for _, item := range req.Networks {
switch item.Network {
case "host", "none", "bridge":
hostConf.NetworkMode = container.NetworkMode(item.Network)
}
endpoint := &network.EndpointSettings{
Links: append([]string(nil), item.Links...),
Aliases: append([]string(nil), item.Aliases...),
DriverOpts: cloneStringMap(item.DriverOpts),
GwPriority: item.GwPriority,
MacAddress: item.MacAddr,
}
if item.Ipv4 != "" || item.Ipv6 != "" || len(item.LinkLocalIPs) != 0 {
endpoint.IPAMConfig = &network.EndpointIPAMConfig{
IPv4Address: item.Ipv4,
IPv6Address: item.Ipv6,
LinkLocalIPs: append([]string(nil), item.LinkLocalIPs...),
}
}
networkConf.EndpointsConfig[item.Network] = endpoint
}
} else {
return nil, nil, nil, fmt.Errorf("please set up the network")
}
hostConf.Privileged = req.Privileged
hostConf.AutoRemove = req.AutoRemove
hostConf.CPUShares = req.CPUShares
hostConf.PublishAllPorts = req.PublishAllPorts
hostConf.RestartPolicy = container.RestartPolicy{Name: container.RestartPolicyMode(req.RestartPolicy)}
if req.RestartPolicy == "on-failure" {
hostConf.RestartPolicy.MaximumRetryCount = 5
}
hostConf.NanoCPUs = int64(req.NanoCPUs * 1000000000)
hostConf.Memory = int64(req.Memory * 1024 * 1024)
hostConf.MemorySwap = 0
hostConf.PortBindings = portMap
hostConf.Binds = []string{}
hostConf.Mounts = []mount.Mount{}
hostConf.DNS = req.DNS
hostConf.ExtraHosts = []string{}
for _, item := range req.ExtraHosts {
if len(item.Hostname) == 0 || len(item.IP) == 0 {
continue
}
hostConf.ExtraHosts = append(hostConf.ExtraHosts, fmt.Sprintf("%s:%s", item.Hostname, item.IP))
}
config.Volumes = make(map[string]struct{})
for _, volume := range req.Volumes {
item := mount.Mount{
Type: mount.Type(volume.Type),
Source: volume.SourceDir,
Target: volume.ContainerDir,
ReadOnly: volume.Mode == "ro",
}
if volume.Type == "bind" {
item.BindOptions = &mount.BindOptions{
Propagation: mount.Propagation(volume.Shared),
}
}
hostConf.Mounts = append(hostConf.Mounts, item)
config.Volumes[volume.ContainerDir] = struct{}{}
}
return &config, &hostConf, &networkConf, nil
}
func loadVolumeBinds(binds []container.MountPoint) []dto.VolumeHelper {
var datas []dto.VolumeHelper
for _, bind := range binds {
var volumeItem dto.VolumeHelper
volumeItem.Type = string(bind.Type)
if bind.Type == "volume" {
volumeItem.SourceDir = bind.Name
} else {
volumeItem.SourceDir = bind.Source
}
volumeItem.ContainerDir = bind.Destination
volumeItem.Mode = "ro"
if bind.RW {
volumeItem.Mode = "rw"
}
volumeItem.Shared = string(bind.Propagation)
datas = append(datas, volumeItem)
}
return datas
}
func loadPortByInspect(id string, client *client.Client) ([]container.Port, error) {
containerItem, err := client.ContainerInspect(context.Background(), id)
if err != nil {
return nil, err
}
var itemPorts []container.Port
for key, val := range containerItem.ContainerJSONBase.HostConfig.PortBindings {
if !strings.Contains(string(key), "/") {
continue
}
item := strings.Split(string(key), "/")
itemPort, _ := strconv.ParseUint(item[0], 10, 16)
for _, itemVal := range val {
publicPort, _ := strconv.ParseUint(itemVal.HostPort, 10, 16)
itemPorts = append(itemPorts, container.Port{PrivatePort: uint16(itemPort), Type: item[1], PublicPort: uint16(publicPort), IP: itemVal.HostIP})
}
}
return itemPorts, nil
}
func transPortToStr(ports []container.Port) []string {
return docker.SimplifyPorts(ports)
}
func loadComposeCount(client *client.Client) int {
options := container.ListOptions{All: true}
options.Filters = filters.NewArgs()
options.Filters.Add("label", composeProjectLabel)
list, err := client.ContainerList(context.Background(), options)
if err != nil {
return 0
}
composeCreatedByLocal, _ := composeRepo.ListRecord()
composeMap := make(map[string]struct{})
for _, container := range list {
if name, ok := container.Labels[composeProjectLabel]; ok {
if _, has := composeMap[name]; !has {
composeMap[name] = struct{}{}
}
}
}
for _, compose := range composeCreatedByLocal {
if len(compose.Path) == 0 {
continue
}
if _, has := composeMap[compose.Name]; !has {
composeMap[compose.Name] = struct{}{}
}
}
return len(composeMap)
}
func loadContainerPortForInfo(itemPorts []container.Port) []dto.PortHelper {
var exposedPorts []dto.PortHelper
samePortMap := make(map[string]dto.PortHelper)
ports := transPortToStr(itemPorts)
for _, item := range ports {
itemStr := strings.Split(item, "->")
if len(itemStr) < 2 {
continue
}
var itemPort dto.PortHelper
lastIndex := strings.LastIndex(itemStr[0], ":")
if lastIndex == -1 {
itemPort.HostPort = itemStr[0]
} else {
itemPort.HostIP = itemStr[0][0:lastIndex]
itemPort.HostPort = itemStr[0][lastIndex+1:]
}
itemContainer := strings.Split(itemStr[1], "/")
if len(itemContainer) != 2 {
continue
}
itemPort.ContainerPort = itemContainer[0]
itemPort.Protocol = itemContainer[1]
keyItem := fmt.Sprintf("%s->%s/%s", itemPort.HostPort, itemPort.ContainerPort, itemPort.Protocol)
if val, ok := samePortMap[keyItem]; ok {
val.HostIP = ""
samePortMap[keyItem] = val
} else {
samePortMap[keyItem] = itemPort
}
}
for _, val := range samePortMap {
exposedPorts = append(exposedPorts, val)
}
return exposedPorts
}
func searchWithFilter(req dto.PageContainer, containers []container.Summary) []dto.ContainerInfo {
var (
records []dto.ContainerInfo
list []container.Summary
)
if req.ExcludeAppStore {
for _, item := range containers {
if created, ok := item.Labels[composeCreatedBy]; ok && created == "Apps" {
continue
}
list = append(list, item)
}
} else {
list = containers
}
if len(req.Name) != 0 {
length, count := len(list), 0
for count < length {
if !strings.Contains(list[count].Names[0][1:], req.Name) && !strings.Contains(list[count].Image, req.Name) {
list = append(list[:count], list[(count+1):]...)
length--
} else {
count++
}
}
}
if req.State != "all" {
length, count := len(list), 0
for count < length {
if list[count].State != req.State {
list = append(list[:count], list[(count+1):]...)
length--
} else {
count++
}
}
}
switch req.OrderBy {
case "name":
sort.Slice(list, func(i, j int) bool {
if req.Order == constant.OrderAsc {
return list[i].Names[0][1:] < list[j].Names[0][1:]
}
return list[i].Names[0][1:] > list[j].Names[0][1:]
})
default:
sort.Slice(list, func(i, j int) bool {
if req.Order == constant.OrderAsc {
return list[i].Created < list[j].Created
}
return list[i].Created > list[j].Created
})
}
for _, item := range list {
IsFromCompose := false
if _, ok := item.Labels[composeProjectLabel]; ok {
IsFromCompose = true
}
IsFromApp := false
if created, ok := item.Labels[composeCreatedBy]; ok && created == "Apps" {
IsFromApp = true
}
exposePorts := transPortToStr(item.Ports)
info := dto.ContainerInfo{
ContainerID: item.ID,
CreateTime: time.Unix(item.Created, 0).Format(constant.DateTimeLayout),
Name: item.Names[0][1:],
Ports: exposePorts,
ImageId: strings.Split(item.ImageID, ":")[1],
ImageName: item.Image,
State: item.State,
RunTime: item.Status,
IsFromApp: IsFromApp,
IsFromCompose: IsFromCompose,
}
if item.NetworkSettings != nil && len(item.NetworkSettings.Networks) > 0 {
networks := make([]string, 0, len(item.NetworkSettings.Networks))
for key := range item.NetworkSettings.Networks {
networks = append(networks, item.NetworkSettings.Networks[key].IPAddress)
}
sort.Strings(networks)
info.Network = networks
}
records = append(records, info)
}
descriptions, _ := settingRepo.GetDescriptionList(repo.WithByType("container"))
for i := 0; i < len(records); i++ {
for _, desc := range descriptions {
if desc.ID == records[i].ContainerID {
records[i].Description = desc.Description
records[i].IsPinned = desc.IsPinned
}
}
}
sort.Slice(records, func(i, j int) bool {
return records[i].IsPinned && !records[j].IsPinned
})
return records
}