diff --git a/agent/app/api/v2/terminal.go b/agent/app/api/v2/terminal.go index 193c7353d..864b21224 100644 --- a/agent/app/api/v2/terminal.go +++ b/agent/app/api/v2/terminal.go @@ -165,7 +165,7 @@ func loadMapFromDockerTop(containerID string) map[string]string { pidMap := make(map[string]string) sudo := cmd.SudoHandleCmd() - stdout, err := cmd.Execf("%s docker top %s -eo pid,command ", sudo, containerID) + stdout, err := cmd.RunDefaultWithStdoutBashCf("%s docker top %s -eo pid,command ", sudo, containerID) if err != nil { return pidMap } @@ -192,7 +192,7 @@ func killBash(containerID, comm string, pidMap map[string]string) { } } if !isOld && command == comm { - _, _ = cmd.Execf("%s kill -9 %s", sudo, pid) + _, _ = cmd.RunDefaultWithStdoutBashCf("%s kill -9 %s", sudo, pid) } } } diff --git a/agent/app/dto/cronjob.go b/agent/app/dto/cronjob.go index 6f0118b7c..a5e1ed05d 100644 --- a/agent/app/dto/cronjob.go +++ b/agent/app/dto/cronjob.go @@ -40,6 +40,8 @@ type CronjobCreate struct { SourceAccountIDs string `json:"sourceAccountIDs"` DownloadAccountID uint `json:"downloadAccountID"` RetainCopies int `json:"retainCopies" validate:"number,min=1"` + RetryTimes int `json:"retryTimes" validate:"number,min=0"` + Timeout uint `json:"timeout" validate:"number,min=1"` Secret string `json:"secret"` AlertCount uint `json:"alertCount"` @@ -72,6 +74,8 @@ type CronjobUpdate struct { SourceAccountIDs string `json:"sourceAccountIDs"` DownloadAccountID uint `json:"downloadAccountID"` RetainCopies int `json:"retainCopies" validate:"number,min=1"` + RetryTimes int `json:"retryTimes" validate:"number,min=0"` + Timeout uint `json:"timeout" validate:"number,min=1"` Secret string `json:"secret"` AlertCount uint `json:"alertCount"` @@ -122,6 +126,8 @@ type CronjobInfo struct { IsDir bool `json:"isDir"` SourceDir string `json:"sourceDir"` RetainCopies int `json:"retainCopies"` + RetryTimes int `json:"retryTimes"` + Timeout uint `json:"timeout"` SourceAccounts []string `json:"sourceAccounts"` DownloadAccount string `json:"downloadAccount"` diff --git a/agent/app/model/cronjob.go b/agent/app/model/cronjob.go index 48d04dce4..ad79b7b33 100644 --- a/agent/app/model/cronjob.go +++ b/agent/app/model/cronjob.go @@ -32,6 +32,8 @@ type Cronjob struct { SourceAccountIDs string `json:"sourceAccountIDs"` DownloadAccountID uint `json:"downloadAccountID"` + RetryTimes uint `json:"retryTimes"` + Timeout uint `json:"timeout"` RetainCopies uint64 `json:"retainCopies"` Status string `json:"status"` diff --git a/agent/app/service/ai.go b/agent/app/service/ai.go index f2ece5d5c..c1fed9732 100644 --- a/agent/app/service/ai.go +++ b/agent/app/service/ai.go @@ -73,7 +73,7 @@ func (u *AIToolService) LoadDetail(name string) (string, error) { if err != nil { return "", err } - stdout, err := cmd.Execf("docker exec %s ollama show %s", containerName, name) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker exec %s ollama show %s", containerName, name) if err != nil { return "", err } @@ -107,7 +107,8 @@ func (u *AIToolService) Create(req dto.OllamaModelName) error { } go func() { taskItem.AddSubTask(i18n.GetWithName("OllamaModelPull", req.Name), func(t *task.Task) error { - return cmd.ExecShellWithTask(taskItem, time.Hour, "docker", "exec", containerName, "ollama", "pull", info.Name) + cmdMgr := cmd.NewCommandMgr(cmd.WithTask(*taskItem), cmd.WithTimeout(time.Hour)) + return cmdMgr.Run("docker", "exec", containerName, "ollama", "pull", info.Name) }, nil) taskItem.AddSubTask(i18n.GetWithName("OllamaModelSize", req.Name), func(t *task.Task) error { itemSize, err := loadModelSize(info.Name, containerName) @@ -133,7 +134,7 @@ func (u *AIToolService) Close(name string) error { if err != nil { return err } - stdout, err := cmd.Execf("docker exec %s ollama stop %s", containerName, name) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker exec %s ollama stop %s", containerName, name) if err != nil { return fmt.Errorf("handle ollama stop %s failed, stdout: %s, err: %v", name, stdout, err) } @@ -162,7 +163,8 @@ func (u *AIToolService) Recreate(req dto.OllamaModelName) error { } go func() { taskItem.AddSubTask(i18n.GetWithName("OllamaModelPull", req.Name), func(t *task.Task) error { - return cmd.ExecShellWithTask(taskItem, time.Hour, "docker", "exec", containerName, "ollama", "pull", req.Name) + cmdMgr := cmd.NewCommandMgr(cmd.WithTask(*taskItem), cmd.WithTimeout(time.Hour)) + return cmdMgr.Run("docker", "exec", containerName, "ollama", "pull", req.Name) }, nil) taskItem.AddSubTask(i18n.GetWithName("OllamaModelSize", req.Name), func(t *task.Task) error { itemSize, err := loadModelSize(modelInfo.Name, containerName) @@ -191,7 +193,7 @@ func (u *AIToolService) Delete(req dto.ForceDelete) error { } for _, item := range ollamaList { if item.Status != constant.StatusDeleted { - stdout, err := cmd.Execf("docker exec %s ollama rm %s", containerName, item.Name) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker exec %s ollama rm %s", containerName, item.Name) if err != nil && !req.ForceDelete { return fmt.Errorf("handle ollama rm %s failed, stdout: %s, err: %v", item.Name, stdout, err) } @@ -208,7 +210,7 @@ func (u *AIToolService) Sync() ([]dto.OllamaModelDropList, error) { if err != nil { return nil, err } - stdout, err := cmd.Execf("docker exec %s ollama list", containerName) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker exec %s ollama list", containerName) if err != nil { return nil, err } @@ -380,7 +382,7 @@ func LoadContainerName() (string, error) { } func loadModelSize(name string, containerName string) (string, error) { - stdout, err := cmd.Execf("docker exec %s ollama list | grep %s", containerName, name) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker exec %s ollama list | grep %s", containerName, name) if err != nil { return "", err } diff --git a/agent/app/service/app_utils.go b/agent/app/service/app_utils.go index 6f0265047..aa9bf54b3 100644 --- a/agent/app/service/app_utils.go +++ b/agent/app/service/app_utils.go @@ -995,7 +995,9 @@ func runScript(task *task.Task, appInstall *model.AppInstall, operate string) er } logStr := i18n.GetWithName("ExecShell", operate) task.LogStart(logStr) - out, err := cmd.ExecScript(scriptPath, workDir) + + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(10*time.Minute), cmd.WithScriptPath(scriptPath)) + out, err := cmdMgr.RunWithStdout("bash") if err != nil { if out != "" { err = errors.New(out) diff --git a/agent/app/service/backup_redis.go b/agent/app/service/backup_redis.go index e05575c84..cc6736290 100644 --- a/agent/app/service/backup_redis.go +++ b/agent/app/service/backup_redis.go @@ -96,7 +96,7 @@ func handleRedisBackup(redisInfo *repo.RootInfo, parentTask *task.Task, backupDi } } - stdout, err := cmd.Execf("docker exec %s redis-cli -a %s --no-auth-warning save", redisInfo.ContainerName, redisInfo.Password) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker exec %s redis-cli -a %s --no-auth-warning save", redisInfo.ContainerName, redisInfo.Password) if err != nil { return errors.New(string(stdout)) } @@ -109,14 +109,14 @@ func handleRedisBackup(redisInfo *repo.RootInfo, parentTask *task.Task, backupDi return nil } if strings.HasSuffix(fileName, ".aof") { - stdout1, err := cmd.Execf("docker cp %s:/data/appendonly.aof %s/%s", redisInfo.ContainerName, backupDir, fileName) + stdout1, err := cmd.RunDefaultWithStdoutBashCf("docker cp %s:/data/appendonly.aof %s/%s", redisInfo.ContainerName, backupDir, fileName) if err != nil { return errors.New(string(stdout1)) } return nil } - stdout1, err1 := cmd.Execf("docker cp %s:/data/dump.rdb %s/%s", redisInfo.ContainerName, backupDir, fileName) + stdout1, err1 := cmd.RunDefaultWithStdoutBashCf("docker cp %s:/data/dump.rdb %s/%s", redisInfo.ContainerName, backupDir, fileName) if err1 != nil { return errors.New(string(stdout1)) } diff --git a/agent/app/service/backup_website.go b/agent/app/service/backup_website.go index 30ccfd7a9..fc1945dee 100644 --- a/agent/app/service/backup_website.go +++ b/agent/app/service/backup_website.go @@ -189,7 +189,7 @@ func handleWebsiteRecover(website *model.Website, recoverFile string, isRollback t.LogFailedWithErr(taskName, err) return err } - stdout, err := cmd.Execf("docker exec -i %s nginx -s reload", nginxInfo.ContainerName) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker exec -i %s nginx -s reload", nginxInfo.ContainerName) if err != nil { return errors.New(stdout) } diff --git a/agent/app/service/clam.go b/agent/app/service/clam.go index e1f0eda31..426e92bc2 100644 --- a/agent/app/service/clam.go +++ b/agent/app/service/clam.go @@ -84,7 +84,7 @@ func (c *ClamService) LoadBaseInfo() (dto.ClamBaseInfo, error) { } if baseInfo.IsActive { - version, err := cmd.Exec("clamdscan --version") + version, err := cmd.RunDefaultWithStdoutBashC("clamdscan --version") if err == nil { if strings.Contains(version, "/") { baseInfo.Version = strings.TrimPrefix(strings.Split(version, "/")[0], "ClamAV ") @@ -96,7 +96,7 @@ func (c *ClamService) LoadBaseInfo() (dto.ClamBaseInfo, error) { _ = StopAllCronJob(false) } if baseInfo.FreshIsActive { - version, err := cmd.Exec("freshclam --version") + version, err := cmd.RunDefaultWithStdoutBashC("freshclam --version") if err == nil { if strings.Contains(version, "/") { baseInfo.FreshVersion = strings.TrimPrefix(strings.Split(version, "/")[0], "ClamAV ") @@ -111,13 +111,13 @@ func (c *ClamService) LoadBaseInfo() (dto.ClamBaseInfo, error) { func (c *ClamService) Operate(operate string) error { switch operate { case "start", "restart", "stop": - stdout, err := cmd.Execf("systemctl %s %s", operate, c.serviceName) + stdout, err := cmd.RunDefaultWithStdoutBashCf("systemctl %s %s", operate, c.serviceName) if err != nil { return fmt.Errorf("%s the %s failed, err: %s", operate, c.serviceName, stdout) } return nil case "fresh-start", "fresh-restart", "fresh-stop": - stdout, err := cmd.Execf("systemctl %s %s", strings.TrimPrefix(operate, "fresh-"), freshClamService) + stdout, err := cmd.RunDefaultWithStdoutBashCf("systemctl %s %s", strings.TrimPrefix(operate, "fresh-"), freshClamService) if err != nil { return fmt.Errorf("%s the %s failed, err: %s", operate, c.serviceName, stdout) } @@ -344,7 +344,7 @@ func (c *ClamService) HandleOnce(req dto.OperateByID) error { } } global.LOG.Debugf("clamdscan --fdpass %s %s -l %s", strategy, clam.Path, logFile) - stdout, err := cmd.Execf("clamdscan --fdpass %s %s -l %s", strategy, clam.Path, logFile) + stdout, err := cmd.RunDefaultWithStdoutBashCf("clamdscan --fdpass %s %s -l %s", strategy, clam.Path, logFile) if err != nil { global.LOG.Errorf("clamdscan failed, stdout: %v, err: %v", stdout, err) } diff --git a/agent/app/service/container.go b/agent/app/service/container.go index 3f1f69f3f..94cb804c2 100644 --- a/agent/app/service/container.go +++ b/agent/app/service/container.go @@ -340,7 +340,8 @@ func (u *ContainerService) ContainerCreateByCommand(req dto.ContainerCreateByCom } go func() { taskItem.AddSubTask(i18n.GetWithName("ContainerCreate", containerName), func(t *task.Task) error { - return cmd.ExecShellWithTask(taskItem, 5*time.Minute, "bash", "-c", req.Command) + cmdMgr := cmd.NewCommandMgr(cmd.WithTask(*taskItem), cmd.WithTimeout(5*time.Minute)) + return cmdMgr.RunBashC(req.Command) }, nil) _ = taskItem.Execute() }() diff --git a/agent/app/service/cronjob_helper.go b/agent/app/service/cronjob_helper.go index 1616aeedf..558cf03b3 100644 --- a/agent/app/service/cronjob_helper.go +++ b/agent/app/service/cronjob_helper.go @@ -23,49 +23,33 @@ import ( ) func (u *CronjobService) HandleJob(cronjob *model.Cronjob) { - var ( - message []byte - err error - ) record := cronjobRepo.StartRecords(cronjob.ID, "", cronjob.Type) go func() { - switch cronjob.Type { - case "shell": - if len(cronjob.Script) == 0 { - return + var ( + message []byte + err error + ) + cronjob.RetryTimes = cronjob.RetryTimes + 1 + for i := uint(0); i < cronjob.RetryTimes; i++ { + ctx, cancel := context.WithTimeout(context.Background(), time.Duration(cronjob.Timeout)*time.Second) + done := make(chan error) + go func() { + message, err = u.handleJob(cronjob, record) + if err != nil { + global.LOG.Debugf("try handle cron job [%s] %s failed %d/%d, err: %v", cronjob.Type, cronjob.Name, i+1, cronjob.RetryTimes, err) + } + close(done) + }() + select { + case <-done: + cancel() + case <-ctx.Done(): + global.LOG.Debugf("try handle cron job [%s] %s failed %d/%d, err: timeout", cronjob.Type, cronjob.Name, i+1, cronjob.RetryTimes) + err = fmt.Errorf("handle timeout") + cancel() + continue } - err = u.handleShell(*cronjob, record.TaskID) - case "curl": - if len(cronjob.URL) == 0 { - return - } - err = u.handleCurl(*cronjob, record.TaskID) - case "ntp": - err = u.handleNtpSync(*cronjob, record.TaskID) - case "cutWebsiteLog": - var messageItem []string - messageItem, record.File, err = u.handleCutWebsiteLog(cronjob, record.StartTime) - message = []byte(strings.Join(messageItem, "\n")) - case "clean": - err = u.handleSystemClean(*cronjob, record.TaskID) - case "website": - err = u.handleWebsite(*cronjob, record.StartTime, record.TaskID) - case "app": - err = u.handleApp(*cronjob, record.StartTime, record.TaskID) - case "database": - err = u.handleDatabase(*cronjob, record.StartTime, record.TaskID) - case "directory": - if len(cronjob.SourceDir) == 0 { - return - } - err = u.handleDirectory(*cronjob, record.StartTime) - case "log": - err = u.handleSystemLog(*cronjob, record.StartTime) - case "snapshot": - _ = cronjobRepo.UpdateRecords(record.ID, map[string]interface{}{"records": record.Records}) - err = u.handleSnapshot(*cronjob, record) } - if err != nil { if len(message) != 0 { record.Records, _ = mkdirAndWriteFile(cronjob, record.StartTime, message) @@ -84,6 +68,50 @@ func (u *CronjobService) HandleJob(cronjob *model.Cronjob) { }() } +func (u *CronjobService) handleJob(cronjob *model.Cronjob, record model.JobRecords) ([]byte, error) { + var ( + message []byte + err error + ) + switch cronjob.Type { + case "shell": + if len(cronjob.Script) == 0 { + return nil, fmt.Errorf("the script content is empty and is skipped") + } + err = u.handleShell(*cronjob, record.TaskID) + case "curl": + if len(cronjob.URL) == 0 { + return nil, fmt.Errorf("the url is empty and is skipped") + } + err = u.handleCurl(*cronjob, record.TaskID) + case "ntp": + err = u.handleNtpSync(*cronjob, record.TaskID) + case "cutWebsiteLog": + var messageItem []string + messageItem, record.File, err = u.handleCutWebsiteLog(cronjob, record.StartTime) + message = []byte(strings.Join(messageItem, "\n")) + case "clean": + err = u.handleSystemClean(*cronjob, record.TaskID) + case "website": + err = u.handleWebsite(*cronjob, record.StartTime, record.TaskID) + case "app": + err = u.handleApp(*cronjob, record.StartTime, record.TaskID) + case "database": + err = u.handleDatabase(*cronjob, record.StartTime, record.TaskID) + case "directory": + if len(cronjob.SourceDir) == 0 { + return nil, fmt.Errorf("the source dir is empty and is skipped") + } + err = u.handleDirectory(*cronjob, record.StartTime) + case "log": + err = u.handleSystemLog(*cronjob, record.StartTime) + case "snapshot": + _ = cronjobRepo.UpdateRecords(record.ID, map[string]interface{}{"records": record.Records}) + err = u.handleSnapshot(*cronjob, record) + } + return message, err +} + func (u *CronjobService) handleShell(cronjob model.Cronjob, taskID string) error { taskItem, err := task.NewTaskWithOps(fmt.Sprintf("cronjob-%s", cronjob.Name), task.TaskHandle, task.TaskScopeCronjob, taskID, cronjob.ID) if err != nil { @@ -91,13 +119,14 @@ func (u *CronjobService) handleShell(cronjob model.Cronjob, taskID string) error return err } + cmdMgr := cmd.NewCommandMgr(cmd.WithTask(*taskItem), cmd.WithTimeout(24*time.Hour)) taskItem.AddSubTask(i18n.GetWithName("HandleShell", cronjob.Name), func(t *task.Task) error { if len(cronjob.ContainerName) != 0 { command := "sh" if len(cronjob.Command) != 0 { command = cronjob.Command } - return cmd.ExecShellWithTask(taskItem, 24*time.Hour, "docker", "exec", cronjob.ContainerName, command, "-c", strings.ReplaceAll(cronjob.Script, "\"", "\\\"")) + return cmdMgr.Run("docker", "exec", cronjob.ContainerName, command, "-c", strings.ReplaceAll(cronjob.Script, "\"", "\\\"")) } if len(cronjob.Executor) == 0 { cronjob.Executor = "bash" @@ -114,14 +143,14 @@ func (u *CronjobService) handleShell(cronjob model.Cronjob, taskID string) error return err } if len(cronjob.User) == 0 { - return cmd.ExecShellWithTask(taskItem, 24*time.Hour, cronjob.Executor, fileItem) + return cmdMgr.Run(cronjob.Executor, fileItem) } - return cmd.ExecShellWithTask(taskItem, 24*time.Hour, "sudo", "-u", cronjob.User, cronjob.Executor, fileItem) + return cmdMgr.Run("sudo", "-u", cronjob.User, cronjob.Executor, fileItem) } if len(cronjob.User) == 0 { - return cmd.ExecShellWithTask(taskItem, 24*time.Hour, cronjob.Executor, cronjob.Script) + return cmdMgr.Run(cronjob.Executor, cronjob.Script) } - if err := cmd.ExecShellWithTask(taskItem, 24*time.Hour, "sudo", "-u", cronjob.User, cronjob.Executor, cronjob.Script); err != nil { + if err := cmdMgr.Run("sudo", "-u", cronjob.User, cronjob.Executor, cronjob.Script); err != nil { return err } return nil @@ -139,10 +168,8 @@ func (u *CronjobService) handleCurl(cronjob model.Cronjob, taskID string) error } taskItem.AddSubTask(i18n.GetWithName("HandleShell", cronjob.Name), func(t *task.Task) error { - if err := cmd.ExecShellWithTask(taskItem, 24*time.Hour, "curl", cronjob.URL); err != nil { - return err - } - return nil + cmdMgr := cmd.NewCommandMgr(cmd.WithTask(*taskItem), cmd.WithTimeout(time.Hour)) + return cmdMgr.Run("curl", cronjob.URL) }, nil, ) @@ -214,7 +241,8 @@ func (u *CronjobService) handleCutWebsiteLog(cronjob *model.Cronjob, startTime t } func backupLogFile(dstFilePath, websiteLogDir string, fileOp files.FileOp) error { - if err := cmd.ExecCmd(fmt.Sprintf("tar -czf %s -C %s %s", dstFilePath, websiteLogDir, strings.Join([]string{"access.log", "error.log"}, " "))); err != nil { + cmdMgr := cmd.NewCommandMgr() + if err := cmdMgr.RunBashCf("tar -czf %s -C %s %s", dstFilePath, websiteLogDir, strings.Join([]string{"access.log", "error.log"}, " ")); err != nil { dstDir := pathUtils.Dir(dstFilePath) if err = fileOp.Copy(pathUtils.Join(websiteLogDir, "access.log"), dstDir); err != nil { return err @@ -222,7 +250,7 @@ func backupLogFile(dstFilePath, websiteLogDir string, fileOp files.FileOp) error if err = fileOp.Copy(pathUtils.Join(websiteLogDir, "error.log"), dstDir); err != nil { return err } - if err = cmd.ExecCmd(fmt.Sprintf("tar -czf %s -C %s %s", dstFilePath, dstDir, strings.Join([]string{"access.log", "error.log"}, " "))); err != nil { + if err = cmdMgr.RunBashCf("tar -czf %s -C %s %s", dstFilePath, dstDir, strings.Join([]string{"access.log", "error.log"}, " ")); err != nil { return err } _ = fileOp.DeleteFile(pathUtils.Join(dstDir, "access.log")) diff --git a/agent/app/service/dashboard.go b/agent/app/service/dashboard.go index f09b71c2b..6ceae4a55 100644 --- a/agent/app/service/dashboard.go +++ b/agent/app/service/dashboard.go @@ -59,7 +59,7 @@ func (u *DashboardService) Restart(operation string) error { itemCmd = fmt.Sprintf("%s systemctl restart 1panel-agent.service", cmd.SudoHandleCmd()) } go func() { - stdout, err := cmd.Exec(itemCmd) + stdout, err := cmd.RunDefaultWithStdoutBashC(itemCmd) if err != nil { global.LOG.Errorf("handle %s failed, err: %v", itemCmd, stdout) } @@ -389,9 +389,11 @@ type diskInfo struct { func loadDiskInfo() []dto.DiskInfo { var datas []dto.DiskInfo - stdout, err := cmd.ExecWithTimeOut("df -hT -P|grep '/'|grep -v tmpfs|grep -v 'snap/core'|grep -v udev", 2*time.Second) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(2 * time.Second)) + stdout, err := cmdMgr.RunWithStdoutBashC("df -hT -P|grep '/'|grep -v tmpfs|grep -v 'snap/core'|grep -v udev") if err != nil { - stdout, err = cmd.ExecWithTimeOut("df -lhT -P|grep '/'|grep -v tmpfs|grep -v 'snap/core'|grep -v udev", 1*time.Second) + cmdMgr2 := cmd.NewCommandMgr(cmd.WithTimeout(1 * time.Second)) + stdout, err = cmdMgr2.RunWithStdoutBashC("df -lhT -P|grep '/'|grep -v tmpfs|grep -v 'snap/core'|grep -v udev") if err != nil { return datas } diff --git a/agent/app/service/device.go b/agent/app/service/device.go index 53c849f2c..5d4e894ec 100644 --- a/agent/app/service/device.go +++ b/agent/app/service/device.go @@ -79,7 +79,7 @@ func (u *DeviceService) LoadBaseInfo() (dto.DeviceBaseInfo, error) { } func (u *DeviceService) LoadTimeZone() ([]string, error) { - std, err := cmd.Exec("timedatectl list-timezones") + std, err := cmd.NewCommandMgr(cmd.WithTimeout(10 * time.Minute)).RunWithStdoutBashC("timedatectl list-timezones") if err != nil { return []string{}, err } @@ -129,7 +129,7 @@ func (u *DeviceService) Update(key, value string) error { if cmd.CheckIllegal(value) { return buserr.New("ErrCmdIllegal") } - std, err := cmd.Execf("%s hostnamectl set-hostname %s", cmd.SudoHandleCmd(), value) + std, err := cmd.RunDefaultWithStdoutBashCf("%s hostnamectl set-hostname %s", cmd.SudoHandleCmd(), value) if err != nil { return errors.New(std) } @@ -220,7 +220,7 @@ func (u *DeviceService) UpdatePasswd(req dto.ChangePasswd) error { if cmd.CheckIllegal(req.User, req.Passwd) { return buserr.New("ErrCmdIllegal") } - std, err := cmd.Execf("%s echo '%s:%s' | %s chpasswd", cmd.SudoHandleCmd(), req.User, req.Passwd, cmd.SudoHandleCmd()) + std, err := cmd.RunDefaultWithStdoutBashCf("%s echo '%s:%s' | %s chpasswd", cmd.SudoHandleCmd(), req.User, req.Passwd, cmd.SudoHandleCmd()) if err != nil { if strings.Contains(err.Error(), "does not exist") { return buserr.New("ErrNotExistUser") @@ -235,7 +235,7 @@ func (u *DeviceService) UpdateSwap(req dto.SwapHelper) error { return buserr.New("ErrCmdIllegal") } if !req.IsNew { - std, err := cmd.Execf("%s swapoff %s", cmd.SudoHandleCmd(), req.Path) + std, err := cmd.RunDefaultWithStdoutBashCf("%s swapoff %s", cmd.SudoHandleCmd(), req.Path) if err != nil { return fmt.Errorf("handle swapoff %s failed, err: %s", req.Path, std) } @@ -246,17 +246,17 @@ func (u *DeviceService) UpdateSwap(req dto.SwapHelper) error { } return operateSwapWithFile(true, req) } - std1, err := cmd.Execf("%s dd if=/dev/zero of=%s bs=1024 count=%d", cmd.SudoHandleCmd(), req.Path, req.Size) + std1, err := cmd.RunDefaultWithStdoutBashCf("%s dd if=/dev/zero of=%s bs=1024 count=%d", cmd.SudoHandleCmd(), req.Path, req.Size) if err != nil { return fmt.Errorf("handle dd path %s failed, err: %s", req.Path, std1) } - std2, err := cmd.Execf("%s mkswap -f %s", cmd.SudoHandleCmd(), req.Path) + std2, err := cmd.RunDefaultWithStdoutBashCf("%s mkswap -f %s", cmd.SudoHandleCmd(), req.Path) if err != nil { return fmt.Errorf("handle dd path %s failed, err: %s", req.Path, std2) } - std3, err := cmd.Execf("%s swapon %s", cmd.SudoHandleCmd(), req.Path) + std3, err := cmd.RunDefaultWithStdoutBashCf("%s swapon %s", cmd.SudoHandleCmd(), req.Path) if err != nil { - _, _ = cmd.Execf("%s swapoff %s", cmd.SudoHandleCmd(), req.Path) + _, _ = cmd.RunDefaultWithStdoutBashCf("%s swapoff %s", cmd.SudoHandleCmd(), req.Path) return fmt.Errorf("handle dd path %s failed, err: %s", req.Path, std3) } return operateSwapWithFile(false, req) @@ -373,7 +373,7 @@ func loadHosts() []dto.HostHelper { } func loadHostname() string { - std, err := cmd.Exec("hostname") + std, err := cmd.RunDefaultWithStdoutBashC("hostname") if err != nil { return "" } @@ -381,7 +381,7 @@ func loadHostname() string { } func loadUser() string { - std, err := cmd.Exec("whoami") + std, err := cmd.RunDefaultWithStdoutBashC("whoami") if err != nil { return "" } @@ -390,7 +390,7 @@ func loadUser() string { func loadSwap() []dto.SwapHelper { var data []dto.SwapHelper - std, err := cmd.Execf("%s swapon --summary", cmd.SudoHandleCmd()) + std, err := cmd.RunDefaultWithStdoutBashCf("%s swapon --summary", cmd.SudoHandleCmd()) if err != nil { return data } diff --git a/agent/app/service/docker.go b/agent/app/service/docker.go index 6e9d7b425..10323601f 100644 --- a/agent/app/service/docker.go +++ b/agent/app/service/docker.go @@ -90,7 +90,7 @@ func (u *DockerService) LoadDockerConf() (*dto.DaemonJsonConf, error) { data.Version = itemVersion.Version } data.IsSwarm = false - stdout2, _ := cmd.Exec("docker info | grep Swarm") + stdout2, _ := cmd.RunDefaultWithStdoutBashC("docker info | grep Swarm") if string(stdout2) == " Swarm: active\n" { data.IsSwarm = true } @@ -380,7 +380,7 @@ func (u *DockerService) OperateDocker(req dto.DockerOperation) error { if req.Operation == "stop" { isSocketActive, _ := systemctl.IsActive("docker.socket") if isSocketActive { - std, err := cmd.Execf("%s systemctl stop docker.socket", sudo) + std, err := cmd.RunDefaultWithStdoutBashCf("%s systemctl stop docker.socket", sudo) if err != nil { global.LOG.Errorf("handle systemctl stop docker.socket failed, err: %v", std) } @@ -391,7 +391,7 @@ func (u *DockerService) OperateDocker(req dto.DockerOperation) error { return err } } - stdout, err := cmd.Execf("%s %s %s ", dockerCmd, req.Operation, service) + stdout, err := cmd.RunDefaultWithStdoutBashCf("%s %s %s ", dockerCmd, req.Operation, service) if err != nil { return errors.New(string(stdout)) } @@ -451,7 +451,7 @@ func validateDockerConfig() error { if !cmd.Which("dockerd") { return nil } - stdout, err := cmd.Exec("dockerd --validate") + stdout, err := cmd.RunDefaultWithStdoutBashC("dockerd --validate") if strings.Contains(stdout, "unknown flag: --validate") { return nil } @@ -462,7 +462,7 @@ func validateDockerConfig() error { } func getDockerRestartCommand() (string, error) { - stdout, err := cmd.Exec("which docker") + stdout, err := cmd.RunDefaultWithStdoutBashC("which docker") if err != nil { return "", fmt.Errorf("failed to find docker: %v", err) } @@ -478,7 +478,7 @@ func restartDocker() error { if err != nil { return err } - stdout, err := cmd.Execf("%s restart docker", restartCmd) + stdout, err := cmd.RunDefaultWithStdoutBashCf("%s restart docker", restartCmd) if err != nil { return fmt.Errorf("failed to restart Docker: %s", stdout) } diff --git a/agent/app/service/firewall.go b/agent/app/service/firewall.go index 3f612f7cf..3be87a426 100644 --- a/agent/app/service/firewall.go +++ b/agent/app/service/firewall.go @@ -191,19 +191,19 @@ func (u *FirewallService) OperateFirewall(operation string) error { _ = client.Stop() return err } - _, _ = cmd.Exec("systemctl restart docker") + _, _ = cmd.RunDefaultWithStdoutBashC("systemctl restart docker") return nil case "stop": if err := client.Stop(); err != nil { return err } - _, _ = cmd.Exec("systemctl restart docker") + _, _ = cmd.RunDefaultWithStdoutBashC("systemctl restart docker") return nil case "restart": if err := client.Restart(); err != nil { return err } - _, _ = cmd.Exec("systemctl restart docker") + _, _ = cmd.RunDefaultWithStdoutBashC("systemctl restart docker") return nil case "disablePing": return u.updatePingStatus("0") @@ -567,7 +567,7 @@ func (u *FirewallService) pingStatus() string { if _, err := os.Stat("/etc/sysctl.conf"); err != nil { return constant.StatusNone } - stdout, _ := cmd.Execf("%s sysctl -a 2>/dev/null | grep 'net.ipv4.icmp.echo_ignore_all'", cmd.SudoHandleCmd()) + stdout, _ := cmd.RunDefaultWithStdoutBashCf("%s sysctl -a 2>/dev/null | grep 'net.ipv4.icmp.echo_ignore_all'", cmd.SudoHandleCmd()) if stdout == "net.ipv4.icmp_echo_ignore_all = 1\n" { return constant.StatusEnable } @@ -582,7 +582,7 @@ func (u *FirewallService) updatePingStatus(enable string) error { files := strings.Split(string(lineBytes), "\n") var newFiles []string var hasIpv6 bool - ipv6Status, _ := cmd.Exec("sysctl -a 2>/dev/null | grep 'net.ipv6.icmp.echo_ignore_all'") + ipv6Status, _ := cmd.RunDefaultWithStdoutBashC("sysctl -a 2>/dev/null | grep 'net.ipv6.icmp.echo_ignore_all'") if len(strings.ReplaceAll(ipv6Status, "\n", "")) != 0 { hasIpv6 = true } @@ -616,9 +616,7 @@ func (u *FirewallService) updatePingStatus(enable string) error { return err } - sudo := cmd.SudoHandleCmd() - command := fmt.Sprintf("%s sysctl -p", sudo) - stdout, err := cmd.Exec(command) + stdout, err := cmd.RunDefaultWithStdoutBashCf("%s sysctl -p", cmd.SudoHandleCmd()) if err != nil { return fmt.Errorf("update ping status failed, err: %v", stdout) } diff --git a/agent/app/service/host_tool.go b/agent/app/service/host_tool.go index 04cf4c911..bf29a2099 100644 --- a/agent/app/service/host_tool.go +++ b/agent/app/service/host_tool.go @@ -72,7 +72,7 @@ func (h *HostToolService) GetToolStatus(req request.HostToolReq) (*response.Host supervisorConfig.ServiceName = serviceNameSet.Value } - versionRes, _ := cmd.Exec("supervisord -v") + versionRes, _ := cmd.RunDefaultWithStdoutBashC("supervisord -v") supervisorConfig.Version = strings.TrimSuffix(versionRes, "\n") _, ctlRrr := exec.LookPath("supervisorctl") supervisorConfig.CtlExist = ctlRrr == nil @@ -547,7 +547,8 @@ func operateSupervisorCtl(operate, name, group, includeDir, containerName string err error ) if containerName != "" { - output, err = cmd.ExecWithTimeOut(fmt.Sprintf("docker exec %s supervisorctl %s", containerName, strings.Join(processNames, " ")), 2*time.Second) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(2 * time.Second)) + output, err = cmdMgr.RunWithStdoutBashCf("docker exec %s supervisorctl %s", containerName, strings.Join(processNames, " ")) } else { var out []byte out, err = exec.Command("supervisorctl", processNames...).Output() @@ -589,8 +590,8 @@ func getProcessStatus(config *response.SupervisorProcessConfig, containerName st ) processNames = append(processNames, getProcessName(config.Name, config.Numprocs)...) if containerName != "" { - execStr := fmt.Sprintf("docker exec %s supervisorctl %s", containerName, strings.Join(processNames, " ")) - output, err = cmd.ExecWithTimeOut(execStr, 3*time.Second) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(3 * time.Second)) + output, err = cmdMgr.RunWithStdoutBashCf("docker exec %s supervisorctl %s", containerName, strings.Join(processNames, " ")) } else { var out []byte out, err = exec.Command("supervisorctl", processNames...).Output() diff --git a/agent/app/service/image_repo.go b/agent/app/service/image_repo.go index 22d7e1e30..cfe08d82c 100644 --- a/agent/app/service/image_repo.go +++ b/agent/app/service/image_repo.go @@ -91,7 +91,7 @@ func (u *ImageRepoService) Create(req dto.ImageRepoCreate) error { if err := u.handleRegistries(req.DownloadUrl, "", "create"); err != nil { return fmt.Errorf("create registry %s failed, err: %v", req.DownloadUrl, err) } - stdout, err := cmd.Exec("systemctl restart docker") + stdout, err := cmd.RunDefaultWithStdoutBashC("systemctl restart docker") if err != nil { return errors.New(string(stdout)) } @@ -105,7 +105,7 @@ func (u *ImageRepoService) Create(req dto.ImageRepoCreate) error { cancel() return errors.New("the docker service cannot be restarted") default: - stdout, err := cmd.Exec("systemctl is-active docker") + stdout, err := cmd.RunDefaultWithStdoutBashC("systemctl is-active docker") if string(stdout) == "active\n" && err == nil { global.LOG.Info("docker restart with new conf successful!") return nil @@ -172,7 +172,8 @@ func (u *ImageRepoService) Update(req dto.ImageRepoUpdate) error { } if repo.Auth != req.Auth || repo.DownloadUrl != req.DownloadUrl { if repo.Auth { - _, _ = cmd.ExecWithCheck("docker", "logout", repo.DownloadUrl) + cmdMgr := cmd.NewCommandMgr() + _, _ = cmdMgr.RunWithStdout("docker", "logout", "-i", repo.DownloadUrl) } if req.Auth { if err := u.CheckConn(req.DownloadUrl, req.Username, req.Password); err != nil { @@ -200,7 +201,8 @@ func (u *ImageRepoService) Update(req dto.ImageRepoUpdate) error { } func (u *ImageRepoService) CheckConn(host, user, password string) error { - stdout, err := cmd.ExecWithCheck("docker", "login", "-u", user, "-p", password, host) + cmdMgr := cmd.NewCommandMgr() + stdout, err := cmdMgr.RunWithStdout("docker", "login", "-u", user, "-p", password, host) if err != nil { return fmt.Errorf("stdout: %s, stderr: %v", stdout, err) } diff --git a/agent/app/service/nginx.go b/agent/app/service/nginx.go index a62735dcf..00d9e0f90 100644 --- a/agent/app/service/nginx.go +++ b/agent/app/service/nginx.go @@ -15,6 +15,7 @@ import ( "github.com/1Panel-dev/1Panel/agent/app/task" "github.com/1Panel-dev/1Panel/agent/buserr" "github.com/1Panel-dev/1Panel/agent/global" + "github.com/1Panel-dev/1Panel/agent/utils/cmd" cmd2 "github.com/1Panel-dev/1Panel/agent/utils/cmd" "github.com/subosito/gotenv" @@ -235,7 +236,8 @@ func (n NginxService) Build(req request.NginxBuildReq) error { return err } buildTask.AddSubTask("", func(t *task.Task) error { - if err = cmd2.ExecWithLogger(fmt.Sprintf("docker compose -f %s build", nginxInstall.GetComposePath()), t.Logger, 15*time.Minute); err != nil { + cmdMgr := cmd2.NewCommandMgr(cmd.WithTask(*buildTask), cmd.WithTimeout(15*time.Minute)) + if err = cmdMgr.RunBashCf("docker compose -f %s build", nginxInstall.GetComposePath()); err != nil { return err } _, err = compose.DownAndUp(nginxInstall.GetComposePath()) diff --git a/agent/app/service/nginx_utils.go b/agent/app/service/nginx_utils.go index 8b9e1345f..30579428b 100644 --- a/agent/app/service/nginx_utils.go +++ b/agent/app/service/nginx_utils.go @@ -2,7 +2,6 @@ package service import ( "errors" - "fmt" "os" "path" "strings" @@ -230,11 +229,17 @@ func getNginxParamsFromStaticFile(scope dto.NginxKey, newParams []dto.NginxParam } func opNginx(containerName, operate string) error { - nginxCmd := fmt.Sprintf("docker exec -i %s %s", containerName, "nginx -s reload") + var ( + out string + err error + ) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(20 * time.Second)) if operate == constant.NginxCheck { - nginxCmd = fmt.Sprintf("docker exec -i %s %s", containerName, "nginx -t") + out, err = cmdMgr.RunWithStdout("docker", "exec", "-i", containerName, "nginx -t") + } else { + out, err = cmdMgr.RunWithStdout("docker", "exec", "-i", containerName, "nginx -s reload") } - if out, err := cmd.ExecWithTimeOut(nginxCmd, 20*time.Second); err != nil { + if err != nil { if out != "" { return errors.New(out) } diff --git a/agent/app/service/runtime.go b/agent/app/service/runtime.go index 4da272f07..0ebfe5450 100644 --- a/agent/app/service/runtime.go +++ b/agent/app/service/runtime.go @@ -26,6 +26,7 @@ import ( "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/utils/cmd" cmd2 "github.com/1Panel-dev/1Panel/agent/utils/cmd" "github.com/1Panel-dev/1Panel/agent/utils/compose" "github.com/1Panel-dev/1Panel/agent/utils/docker" @@ -666,7 +667,9 @@ func (r *RuntimeService) OperateNodeModules(req request.NodeModuleOperateReq) er } } cmd += " " + req.Module - return cmd2.ExecContainerScript(containerName, cmd, 5*time.Minute) + + cmdMgr := cmd2.NewCommandMgr(cmd2.WithTimeout(5 * time.Minute)) + return cmdMgr.Run("docker", "exec", "-i", containerName, "bash", "-c", fmt.Sprintf("'%s'", cmd)) } func (r *RuntimeService) SyncForRestart() error { @@ -703,8 +706,8 @@ func (r *RuntimeService) GetPHPExtensions(runtimeID uint) (response.PHPExtension if err != nil { return res, err } - phpCmd := fmt.Sprintf("docker exec -i %s %s", runtime.ContainerName, "php -m") - out, err := cmd2.ExecWithTimeOut(phpCmd, 20*time.Second) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(20 * time.Second)) + out, err := cmdMgr.RunWithStdoutBashCf("docker exec -i %s php -m", runtime.ContainerName) if err != nil { if out != "" { return res, errors.New(out) @@ -750,9 +753,9 @@ func (r *RuntimeService) InstallPHPExtension(req request.PHPExtensionInstallReq) if err != nil { return err } + cmdMgr := cmd2.NewCommandMgr(cmd.WithTask(*installTask), cmd.WithTimeout(15*time.Minute)) installTask.AddSubTask("", func(t *task.Task) error { - installCmd := fmt.Sprintf("docker exec -i %s %s %s", runtime.ContainerName, "install-ext", req.Name) - err = cmd2.ExecWithLogger(installCmd, t.Logger, 15*time.Minute) + err = cmdMgr.RunBashCf("docker exec -i %s install-ext %s", runtime.ContainerName, req.Name) if err != nil { return err } @@ -763,8 +766,7 @@ func (r *RuntimeService) InstallPHPExtension(req request.PHPExtensionInstallReq) if err != nil { return err } - commitCmd := fmt.Sprintf("docker commit %s %s", runtime.ContainerName, runtime.Image) - err = cmd2.ExecWithLogger(commitCmd, t.Logger, 15*time.Minute) + err = cmdMgr.RunBashCf("docker commit %s %s", runtime.ContainerName, runtime.Image) if err != nil { return err } diff --git a/agent/app/service/runtime_utils.go b/agent/app/service/runtime_utils.go index b06727922..6dbf8b333 100644 --- a/agent/app/service/runtime_utils.go +++ b/agent/app/service/runtime_utils.go @@ -6,8 +6,6 @@ import ( "context" "encoding/json" "fmt" - "github.com/1Panel-dev/1Panel/agent/i18n" - "github.com/1Panel-dev/1Panel/agent/utils/common" "io" "os" "os/exec" @@ -16,6 +14,9 @@ import ( "strings" "time" + "github.com/1Panel-dev/1Panel/agent/i18n" + "github.com/1Panel-dev/1Panel/agent/utils/common" + "github.com/1Panel-dev/1Panel/agent/app/dto" "github.com/1Panel-dev/1Panel/agent/app/dto/request" "github.com/1Panel-dev/1Panel/agent/app/dto/response" @@ -335,9 +336,8 @@ func buildRuntime(runtime *model.Runtime, oldImageID string, oldEnv string, rebu } extensions := getRuntimeEnv(runtime.Env, "PHP_EXTENSIONS") if extensions != "" { - installCmd := fmt.Sprintf("docker exec -i %s %s %s", runtime.ContainerName, "install-ext", extensions) - err = cmd2.ExecWithLogFile(installCmd, 60*time.Minute, logPath) - if err != nil { + cmdMgr := cmd2.NewCommandMgr(cmd2.WithTimeout(60*time.Minute), cmd2.WithOutputFile(logPath)) + if err = cmdMgr.Run("docker", "exec", "-i", runtime.ContainerName, "install-ext", extensions); err != nil { runtime.Status = constant.StatusError runtime.Message = buserr.New("ErrImageBuildErr").Error() + ":" + err.Error() _ = runtimeRepo.Save(runtime) diff --git a/agent/app/service/snapshot.go b/agent/app/service/snapshot.go index e7f400d6a..af5bc5fb0 100644 --- a/agent/app/service/snapshot.go +++ b/agent/app/service/snapshot.go @@ -290,7 +290,7 @@ func loadAppImage(list []dto.DataTree) []dto.DataTree { for i := 0; i < len(list); i++ { itemAppImage := dto.DataTree{ID: uuid.NewString(), Label: "appImage"} - stdout, err := cmd.Execf("cat %s | grep image: ", path.Join(global.Dir.AppDir, list[i].Key, list[i].Name, "docker-compose.yml")) + stdout, err := cmd.RunDefaultWithStdoutBashCf("cat %s | grep image: ", path.Join(global.Dir.AppDir, list[i].Key, list[i].Name, "docker-compose.yml")) if err != nil { list[i].Children = append(list[i].Children, itemAppImage) continue diff --git a/agent/app/service/snapshot_create.go b/agent/app/service/snapshot_create.go index 0c938cfac..fbecb8f65 100644 --- a/agent/app/service/snapshot_create.go +++ b/agent/app/service/snapshot_create.go @@ -344,7 +344,7 @@ func snapAppImage(snap snapHelper, req dto.SnapshotCreate, targetDir string) err } var imageList []string - existStr, _ := cmd.Exec("docker images | awk '{print $1\":\"$2}' | grep -v REPOSITORY:TAG") + existStr, _ := cmd.RunDefaultWithStdoutBashC("docker images | awk '{print $1\":\"$2}' | grep -v REPOSITORY:TAG") existImages := strings.Split(existStr, "\n") for _, app := range req.AppData { for _, item := range app.Children { @@ -364,7 +364,7 @@ func snapAppImage(snap snapHelper, req dto.SnapshotCreate, targetDir string) err snap.Task.Log(strings.Join(imageList, " ")) if len(imageList) != 0 { snap.Task.Logf("docker save %s | gzip -c > %s", strings.Join(imageList, " "), path.Join(targetDir, "images.tar.gz")) - std, err := cmd.Execf("docker save %s | gzip -c > %s", strings.Join(imageList, " "), path.Join(targetDir, "images.tar.gz")) + std, err := cmd.RunDefaultWithStdoutBashCf("docker save %s | gzip -c > %s", strings.Join(imageList, " "), path.Join(targetDir, "images.tar.gz")) if err != nil { snap.Task.LogFailedWithErr(i18n.GetMsgByKey("SnapDockerSave"), errors.New(std)) return errors.New(std) diff --git a/agent/app/service/snapshot_recover.go b/agent/app/service/snapshot_recover.go index b211e66a2..4868afe58 100644 --- a/agent/app/service/snapshot_recover.go +++ b/agent/app/service/snapshot_recover.go @@ -314,7 +314,7 @@ func recoverAppData(src string, itemHelper *snapRecoverHelper) error { itemHelper.Task.Log(i18n.GetMsgByKey("RecoverAppEmpty")) return nil } - std, err := cmd.Execf("docker load < %s", path.Join(src, "images.tar.gz")) + std, err := cmd.RunDefaultWithStdoutBashCf("docker load < %s", path.Join(src, "images.tar.gz")) if err != nil { itemHelper.Task.LogFailedWithErr(i18n.GetMsgByKey("RecoverAppImage"), errors.New(std)) return fmt.Errorf("docker load images failed, err: %v", err) @@ -371,7 +371,7 @@ func recoverBaseData(src string, itemHelper *snapRecoverHelper) error { } } - _, _ = cmd.Exec("systemctl restart docker") + _, _ = cmd.RunDefaultWithStdoutBashC("systemctl restart docker") return nil } @@ -400,7 +400,7 @@ func restartCompose(composePath string, itemHelper *snapRecoverHelper) error { continue } upCmd := fmt.Sprintf("docker compose -f %s up -d", pathItem) - stdout, err := cmd.Exec(upCmd) + stdout, err := cmd.RunDefaultWithStdoutBashC(upCmd) if err != nil { itemHelper.Task.LogFailedWithErr(i18n.GetMsgByKey("RecoverCompose"), errors.New(stdout)) continue diff --git a/agent/app/service/ssh.go b/agent/app/service/ssh.go index 8e25e742d..5480ca8bb 100644 --- a/agent/app/service/ssh.go +++ b/agent/app/service/ssh.go @@ -125,17 +125,17 @@ func (u *SSHService) OperateSSH(operation string) error { if operation == "stop" { isSocketActive, _ := systemctl.IsActive(serviceName + ".socket") if isSocketActive { - std, err := cmd.Execf("%s systemctl stop %s", sudo, serviceName+".socket") + std, err := cmd.RunDefaultWithStdoutBashCf("%s systemctl stop %s", sudo, serviceName+".socket") if err != nil { global.LOG.Errorf("handle systemctl stop %s.socket failed, err: %v", serviceName, std) } } } - stdout, err := cmd.Execf("%s systemctl %s %s", sudo, operation, serviceName) + stdout, err := cmd.RunDefaultWithStdoutBashCf("%s systemctl %s %s", sudo, operation, serviceName) if err != nil { if strings.Contains(stdout, "alias name or linked unit file") { - stdout, err := cmd.Execf("%s systemctl %s ssh", sudo, operation) + stdout, err := cmd.RunDefaultWithStdoutBashCf("%s systemctl %s ssh", sudo, operation) if err != nil { return fmt.Errorf("%s ssh(alias name or linked unit file) failed, stdout: %s, err: %v", operation, stdout, err) } @@ -167,9 +167,9 @@ func (u *SSHService) Update(req dto.SSHUpdate) error { } sudo := cmd.SudoHandleCmd() if req.Key == "Port" { - stdout, _ := cmd.Execf("%s getenforce", sudo) + stdout, _ := cmd.RunDefaultWithStdoutBashCf("%s getenforce", sudo) if stdout == "Enforcing\n" { - _, _ = cmd.Execf("%s semanage port -a -t ssh_port_t -p tcp %s", sudo, req.NewValue) + _, _ = cmd.RunDefaultWithStdoutBashCf("%s semanage port -a -t ssh_port_t -p tcp %s", sudo, req.NewValue) } ruleItem := dto.PortRuleUpdate{ @@ -191,7 +191,7 @@ func (u *SSHService) Update(req dto.SSHUpdate) error { } } - _, _ = cmd.Execf("%s systemctl restart %s", sudo, serviceName) + _, _ = cmd.RunDefaultWithStdoutBashCf("%s systemctl restart %s", sudo, serviceName) return nil } @@ -210,7 +210,7 @@ func (u *SSHService) UpdateByFile(value string) error { return err } sudo := cmd.SudoHandleCmd() - _, _ = cmd.Execf("%s systemctl restart %s", sudo, serviceName) + _, _ = cmd.RunDefaultWithStdoutBashCf("%s systemctl restart %s", sudo, serviceName) return nil } @@ -230,7 +230,7 @@ func (u *SSHService) GenerateSSH(req dto.GenerateSSH) error { if len(req.Password) != 0 { command = fmt.Sprintf("ssh-keygen -t %s -P %s -f %s/.ssh/id_item_%s | echo y", req.EncryptionMode, req.Password, currentUser.HomeDir, req.EncryptionMode) } - stdout, err := cmd.Exec(command) + stdout, err := cmd.RunDefaultWithStdoutBashC(command) if err != nil { return fmt.Errorf("generate failed, err: %v, message: %s", err, stdout) } @@ -248,7 +248,7 @@ func (u *SSHService) GenerateSSH(req dto.GenerateSSH) error { } defer authFile.Close() } - stdout1, err := cmd.Execf("cat %s >> %s/.ssh/authorized_keys", secretPubFile, currentUser.HomeDir) + stdout1, err := cmd.RunDefaultWithStdoutBashCf("cat %s >> %s/.ssh/authorized_keys", secretPubFile, currentUser.HomeDir) if err != nil { return fmt.Errorf("generate failed, err: %v, message: %s", err, stdout1) } @@ -423,7 +423,7 @@ func loadSSHData(ctx *gin.Context, command string, showCountFrom, showCountTo, c if err != nil { return datas, 0, 0 } - stdout, err := cmd.Exec(command) + stdout, err := cmd.RunDefaultWithStdoutBashC(command) if err != nil { return datas, 0, 0 } @@ -532,7 +532,7 @@ func loadFailedSecureDatas(line string) dto.SSHHistory { } func handleGunzip(path string) error { - if _, err := cmd.Execf("gunzip %s", path); err != nil { + if _, err := cmd.RunDefaultWithStdoutBashCf("gunzip %s", path); err != nil { return err } return nil diff --git a/agent/app/service/website.go b/agent/app/service/website.go index a65857d31..02a1a384f 100644 --- a/agent/app/service/website.go +++ b/agent/app/service/website.go @@ -1539,11 +1539,8 @@ func (w WebsiteService) UpdateSitePermission(req request.WebsiteUpdateDirPermiss return err } absoluteIndexPath := GetSitePath(website, SiteIndexDir) - chownCmd := fmt.Sprintf("chown -R %s:%s %s", req.User, req.Group, absoluteIndexPath) - if cmd.HasNoPasswordSudo() { - chownCmd = fmt.Sprintf("sudo %s", chownCmd) - } - if out, err := cmd.ExecWithTimeOut(chownCmd, 10*time.Second); err != nil { + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(10 * time.Second)) + if out, err := cmdMgr.RunWithStdoutBashCf("%s chown -R %s:%s %s", cmd.SudoHandleCmd(), req.User, req.Group, absoluteIndexPath); err != nil { if out != "" { return errors.New(out) } diff --git a/agent/app/service/website_ca.go b/agent/app/service/website_ca.go index a0315d49b..087fd95f1 100644 --- a/agent/app/service/website_ca.go +++ b/agent/app/service/website_ca.go @@ -378,7 +378,8 @@ func (w WebsiteCAService) ObtainSSL(req request.WebsiteCAObtain) (*model.Website workDir = websiteSSL.Dir } logger.Println(i18n.GetMsgByKey("ExecShellStart")) - if err = cmd.ExecShellWithTimeOut(websiteSSL.Shell, workDir, logger, 30*time.Minute); err != nil { + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(30*time.Minute), cmd.WithLogger(logger), cmd.WithWorkDir(workDir)) + if err = cmdMgr.RunBashC(websiteSSL.Shell); err != nil { logger.Println(i18n.GetMsgWithMap("ErrExecShell", map[string]interface{}{"err": err.Error()})) } else { logger.Println(i18n.GetMsgByKey("ExecShellSuccess")) diff --git a/agent/app/service/website_ssl.go b/agent/app/service/website_ssl.go index 5b2a91a5e..05726f60c 100644 --- a/agent/app/service/website_ssl.go +++ b/agent/app/service/website_ssl.go @@ -354,7 +354,8 @@ func (w WebsiteSSLService) ObtainSSL(apply request.WebsiteSSLApply) error { workDir = websiteSSL.Dir } printSSLLog(logger, "ExecShellStart", nil, apply.DisableLog) - if err = cmd.ExecShellWithTimeOut(websiteSSL.Shell, workDir, logger, 30*time.Minute); err != nil { + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(30*time.Minute), cmd.WithLogger(logger), cmd.WithWorkDir(workDir)) + if err = cmdMgr.RunBashC(websiteSSL.Shell); err != nil { printSSLLog(logger, "ErrExecShell", map[string]interface{}{"err": err.Error()}, apply.DisableLog) } else { printSSLLog(logger, "ExecShellSuccess", nil, apply.DisableLog) diff --git a/agent/app/service/website_utils.go b/agent/app/service/website_utils.go index 80add476f..d78fd6295 100644 --- a/agent/app/service/website_utils.go +++ b/agent/app/service/website_utils.go @@ -981,7 +981,8 @@ func checkIsLinkApp(website model.Website) bool { } func chownRootDir(path string) error { - _, err := cmd.ExecWithTimeOut(fmt.Sprintf(`chown -R 1000:1000 "%s"`, path), 1*time.Second) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(1 * time.Second)) + _, err := cmdMgr.RunWithStdoutBashCf(`chown -R 1000:1000 "%s"`, path) if err != nil { return err } diff --git a/agent/init/migration/migrations/init.go b/agent/init/migration/migrations/init.go index 9ace5d0b9..d5ee2bfde 100644 --- a/agent/init/migration/migrations/init.go +++ b/agent/init/migration/migrations/init.go @@ -19,7 +19,7 @@ import ( ) var AddTable = &gormigrate.Migration{ - ID: "20250408-add-table", + ID: "20250413-add-table", Migrate: func(tx *gorm.DB) error { return tx.AutoMigrate( &model.AppDetail{}, diff --git a/agent/utils/ai_tools/gpu/gpu.go b/agent/utils/ai_tools/gpu/gpu.go index 1c557d03e..47d94ee4c 100644 --- a/agent/utils/ai_tools/gpu/gpu.go +++ b/agent/utils/ai_tools/gpu/gpu.go @@ -23,7 +23,8 @@ func New() (bool, NvidiaSMI) { } func (n NvidiaSMI) LoadGpuInfo() (*common.GpuInfo, error) { - itemData, err := cmd.ExecWithTimeOut("nvidia-smi -q -x", 5*time.Second) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(5 * time.Second)) + itemData, err := cmdMgr.RunWithStdoutBashC("nvidia-smi -q -x") if err != nil { return nil, fmt.Errorf("calling nvidia-smi failed, err: %w", err) } diff --git a/agent/utils/ai_tools/xpu/xpu.go b/agent/utils/ai_tools/xpu/xpu.go index 1c9b5a11a..5ae08d8dd 100644 --- a/agent/utils/ai_tools/xpu/xpu.go +++ b/agent/utils/ai_tools/xpu/xpu.go @@ -31,14 +31,15 @@ func (x XpuSMI) loadDeviceData(device Device, wg *sync.WaitGroup, res *[]XPUSimp var wgCmd sync.WaitGroup wgCmd.Add(2) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(5 * time.Second)) go func() { defer wgCmd.Done() - xpuData, xpuErr = cmd.ExecWithTimeOut(fmt.Sprintf("xpu-smi discovery -d %d -j", device.DeviceID), 5*time.Second) + xpuData, xpuErr = cmdMgr.RunWithStdoutBashCf("xpu-smi discovery -d %d -j", device.DeviceID) }() go func() { defer wgCmd.Done() - statsData, statsErr = cmd.ExecWithTimeOut(fmt.Sprintf("xpu-smi stats -d %d -j", device.DeviceID), 5*time.Second) + statsData, statsErr = cmdMgr.RunWithStdoutBashCf("xpu-smi stats -d %d -j", device.DeviceID) }() wgCmd.Wait() @@ -91,7 +92,8 @@ func (x XpuSMI) loadDeviceData(device Device, wg *sync.WaitGroup, res *[]XPUSimp } func (x XpuSMI) LoadDashData() ([]XPUSimpleInfo, error) { - data, err := cmd.ExecWithTimeOut("xpu-smi discovery -j", 5*time.Second) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(5 * time.Second)) + data, err := cmdMgr.RunWithStdoutBashC("xpu-smi discovery -j") if err != nil { return nil, fmt.Errorf("calling xpu-smi failed, err: %w", err) } @@ -119,7 +121,8 @@ func (x XpuSMI) LoadDashData() ([]XPUSimpleInfo, error) { } func (x XpuSMI) LoadGpuInfo() (*XpuInfo, error) { - data, err := cmd.ExecWithTimeOut("xpu-smi discovery -j", 5*time.Second) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(5 * time.Second)) + data, err := cmdMgr.RunWithStdoutBashC("xpu-smi discovery -j") if err != nil { return nil, fmt.Errorf("calling xpu-smi failed, err: %w", err) } @@ -141,7 +144,7 @@ func (x XpuSMI) LoadGpuInfo() (*XpuInfo, error) { wg.Wait() - processData, err := cmd.ExecWithTimeOut(fmt.Sprintf("xpu-smi ps -j"), 5*time.Second) + processData, err := cmdMgr.RunWithStdoutBashC("xpu-smi ps -j") if err != nil { return nil, fmt.Errorf("calling xpu-smi ps failed, err: %w", err) } @@ -188,14 +191,15 @@ func (x XpuSMI) loadDeviceInfo(device Device, wg *sync.WaitGroup, res *XpuInfo, var wgCmd sync.WaitGroup wgCmd.Add(2) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(5 * time.Second)) go func() { defer wgCmd.Done() - xpuData, xpuErr = cmd.ExecWithTimeOut(fmt.Sprintf("xpu-smi discovery -d %d -j", device.DeviceID), 5*time.Second) + xpuData, xpuErr = cmdMgr.RunWithStdoutBashCf("xpu-smi discovery -d %d -j", device.DeviceID) }() go func() { defer wgCmd.Done() - statsData, statsErr = cmd.ExecWithTimeOut(fmt.Sprintf("xpu-smi stats -d %d -j", device.DeviceID), 5*time.Second) + statsData, statsErr = cmdMgr.RunWithStdoutBashCf("xpu-smi stats -d %d -j", device.DeviceID) }() wgCmd.Wait() diff --git a/agent/utils/cmd/cmd.go b/agent/utils/cmd/cmd.go index d6dc9c1cd..fd1d7f46a 100644 --- a/agent/utils/cmd/cmd.go +++ b/agent/utils/cmd/cmd.go @@ -1,321 +1,10 @@ package cmd import ( - "bufio" - "bytes" - "context" - "errors" - "fmt" - "log" - "os" "os/exec" "strings" - "time" - - "github.com/1Panel-dev/1Panel/agent/app/task" - "github.com/1Panel-dev/1Panel/agent/buserr" - "github.com/1Panel-dev/1Panel/agent/constant" ) -func Exec(cmdStr string) (string, error) { - return ExecWithTimeOut(cmdStr, 20*time.Second) -} - -func handleErr(stdout, stderr bytes.Buffer, err error) (string, error) { - errMsg := "" - if len(stderr.String()) != 0 { - errMsg = fmt.Sprintf("stderr: %s", stderr.String()) - } - if len(stdout.String()) != 0 { - if len(errMsg) != 0 { - errMsg = fmt.Sprintf("%s; stdout: %s", errMsg, stdout.String()) - } else { - errMsg = fmt.Sprintf("stdout: %s", stdout.String()) - } - } - return errMsg, err -} - -func ExecWithTimeOut(cmdStr string, timeout time.Duration) (string, error) { - env := os.Environ() - cmd := exec.Command("bash", "-c", cmdStr) - var stdout, stderr bytes.Buffer - cmd.Stdout = &stdout - cmd.Stderr = &stderr - cmd.Env = env - if err := cmd.Start(); err != nil { - return "", err - } - done := make(chan error, 1) - go func() { - done <- cmd.Wait() - }() - after := time.After(timeout) - select { - case <-after: - _ = cmd.Process.Kill() - return "", buserr.New("ErrCmdTimeout") - case err := <-done: - if err != nil { - return handleErr(stdout, stderr, err) - } - } - - return stdout.String(), nil -} - -func ExecWithLogFile(cmdStr string, timeout time.Duration, outputFile string) error { - env := os.Environ() - cmd := exec.Command("bash", "-c", cmdStr) - - outFile, err := os.OpenFile(outputFile, os.O_CREATE|os.O_WRONLY|os.O_APPEND, constant.DirPerm) - if err != nil { - return err - } - defer outFile.Close() - - cmd.Stdout = outFile - cmd.Stderr = outFile - cmd.Env = env - - if err := cmd.Start(); err != nil { - return err - } - - done := make(chan error, 1) - go func() { - done <- cmd.Wait() - }() - - after := time.After(timeout) - select { - case <-after: - _ = cmd.Process.Kill() - return buserr.New("ErrCmdTimeout") - case err := <-done: - if err != nil { - return err - } - } - - return nil -} - -func ExecWithLogger(cmdStr string, logger *log.Logger, timeout time.Duration) error { - cmd := exec.Command("bash", "-c", cmdStr) - - stdoutPipe, err := cmd.StdoutPipe() - if err != nil { - return err - } - stderrPipe, err := cmd.StderrPipe() - if err != nil { - return err - } - - if err := cmd.Start(); err != nil { - return err - } - - go func() { - scanner := bufio.NewScanner(stdoutPipe) - for scanner.Scan() { - logger.Print(scanner.Text()) - } - }() - go func() { - scanner := bufio.NewScanner(stderrPipe) - for scanner.Scan() { - logger.Print(scanner.Text()) - } - }() - - done := make(chan error, 1) - go func() { - done <- cmd.Wait() - }() - - after := time.After(timeout) - select { - case <-after: - _ = cmd.Process.Kill() - return buserr.New("ErrCmdTimeout") - case err := <-done: - return err - } -} - -func ExecContainerScript(containerName, cmdStr string, timeout time.Duration) error { - cmdStr = fmt.Sprintf("docker exec -i %s bash -c '%s'", containerName, cmdStr) - out, err := ExecWithTimeOut(cmdStr, timeout) - if err != nil { - if out != "" { - return fmt.Errorf("%s; err: %v", out, err) - } - return err - } - return nil -} - -func ExecShell(outPath string, timeout time.Duration, name string, arg ...string) error { - env := os.Environ() - file, err := os.OpenFile(outPath, os.O_WRONLY|os.O_CREATE, constant.FilePerm) - if err != nil { - return err - } - defer file.Close() - - cmd := exec.Command(name, arg...) - cmd.Stdout = file - cmd.Stderr = file - cmd.Env = env - if err := cmd.Start(); err != nil { - return err - } - done := make(chan error, 1) - go func() { - done <- cmd.Wait() - }() - after := time.After(timeout) - select { - case <-after: - _ = cmd.Process.Kill() - return buserr.New("ErrCmdTimeout") - case err := <-done: - if err != nil { - return err - } - } - return nil -} - -type CustomWriter struct { - taskItem *task.Task - buffer bytes.Buffer -} - -func (cw *CustomWriter) Write(p []byte) (n int, err error) { - cw.buffer.Write(p) - lines := strings.Split(cw.buffer.String(), "\n") - - for i := 0; i < len(lines)-1; i++ { - cw.taskItem.Log(lines[i]) - } - cw.buffer.Reset() - cw.buffer.WriteString(lines[len(lines)-1]) - - return len(p), nil -} -func (cw *CustomWriter) Flush() { - if cw.buffer.Len() > 0 { - cw.taskItem.Log(cw.buffer.String()) - cw.buffer.Reset() - } -} -func ExecShellWithTask(taskItem *task.Task, timeout time.Duration, name string, arg ...string) error { - env := os.Environ() - customWriter := &CustomWriter{taskItem: taskItem} - cmd := exec.Command(name, arg...) - cmd.Stdout = customWriter - cmd.Stderr = customWriter - cmd.Env = env - if err := cmd.Start(); err != nil { - return err - } - done := make(chan error, 1) - go func() { - done <- cmd.Wait() - customWriter.Flush() - }() - after := time.After(timeout) - select { - case <-after: - _ = cmd.Process.Kill() - return buserr.New("ErrCmdTimeout") - case err := <-done: - if err != nil { - return err - } - } - return nil -} - -func Execf(cmdStr string, a ...interface{}) (string, error) { - env := os.Environ() - cmd := exec.Command("bash", "-c", fmt.Sprintf(cmdStr, a...)) - var stdout, stderr bytes.Buffer - cmd.Stdout = &stdout - cmd.Stderr = &stderr - cmd.Env = env - err := cmd.Run() - if err != nil { - return handleErr(stdout, stderr, err) - } - return stdout.String(), nil -} - -func ExecWithCheck(name string, a ...string) (string, error) { - env := os.Environ() - cmd := exec.Command(name, a...) - var stdout, stderr bytes.Buffer - cmd.Stdout = &stdout - cmd.Stderr = &stderr - cmd.Env = env - err := cmd.Run() - if err != nil { - return handleErr(stdout, stderr, err) - } - return stdout.String(), nil -} - -func ExecScript(scriptPath, workDir string) (string, error) { - env := os.Environ() - cmd := exec.Command("bash", scriptPath) - var stdout, stderr bytes.Buffer - cmd.Dir = workDir - cmd.Stdout = &stdout - cmd.Stderr = &stderr - cmd.Env = env - if err := cmd.Start(); err != nil { - return "", err - } - done := make(chan error, 1) - go func() { - done <- cmd.Wait() - }() - after := time.After(10 * time.Minute) - select { - case <-after: - _ = cmd.Process.Kill() - return "", buserr.New("ErrCmdTimeout") - case err := <-done: - if err != nil { - return handleErr(stdout, stderr, err) - } - } - - return stdout.String(), nil -} - -func ExecCmd(cmdStr string) error { - cmd := exec.Command("bash", "-c", cmdStr) - output, err := cmd.CombinedOutput() - if err != nil { - return fmt.Errorf("error : %v, output: %s", err, output) - } - return nil -} - -func ExecCmdWithDir(cmdStr, workDir string) error { - cmd := exec.Command("bash", "-c", cmdStr) - cmd.Dir = workDir - output, err := cmd.CombinedOutput() - if err != nil { - return fmt.Errorf("error : %v, output: %s", err, output) - } - return nil -} - func CheckIllegal(args ...string) bool { if args == nil { return false @@ -331,12 +20,6 @@ func CheckIllegal(args ...string) bool { return false } -func HasNoPasswordSudo() bool { - cmd2 := exec.Command("sudo", "-n", "ls") - err2 := cmd2.Run() - return err2 == nil -} - func SudoHandleCmd() string { cmd := exec.Command("sudo", "-n", "ls") if err := cmd.Run(); err == nil { @@ -346,29 +29,9 @@ func SudoHandleCmd() string { } func Which(name string) bool { - stdout, err := Execf("which %s", name) + stdout, err := RunDefaultWithStdoutBashCf("which %s", name) if err != nil || (len(strings.ReplaceAll(stdout, "\n", "")) == 0) { return false } return true } - -func ExecShellWithTimeOut(cmdStr, workdir string, logger *log.Logger, timeout time.Duration) error { - env := os.Environ() - ctx, cancel := context.WithTimeout(context.Background(), timeout) - defer cancel() - - cmd := exec.CommandContext(ctx, "bash", "-c", cmdStr) - cmd.Dir = workdir - cmd.Stdout = logger.Writer() - cmd.Stderr = logger.Writer() - cmd.Env = env - if err := cmd.Start(); err != nil { - return err - } - err := cmd.Wait() - if errors.Is(ctx.Err(), context.DeadlineExceeded) { - return buserr.New("ErrCmdTimeout") - } - return err -} diff --git a/agent/utils/cmd/cmdx.go b/agent/utils/cmd/cmdx.go new file mode 100644 index 000000000..dfb303eff --- /dev/null +++ b/agent/utils/cmd/cmdx.go @@ -0,0 +1,213 @@ +package cmd + +import ( + "bytes" + "fmt" + "log" + "os" + "os/exec" + "strings" + "time" + + "github.com/1Panel-dev/1Panel/agent/app/task" + "github.com/1Panel-dev/1Panel/agent/buserr" + "github.com/1Panel-dev/1Panel/agent/constant" +) + +type CommandHelper struct { + workDir string + outputFile string + scriptPath string + timeout time.Duration + taskItem *task.Task + logger *log.Logger +} + +type Option func(*CommandHelper) + +func NewCommandMgr(opts ...Option) *CommandHelper { + s := &CommandHelper{} + for _, opt := range opts { + opt(s) + } + return s +} + +func RunDefaultBashC(command string) error { + mgr := NewCommandMgr() + return mgr.RunBashC(command) +} +func RunDefaultBashCf(command string, arg ...interface{}) error { + mgr := NewCommandMgr() + return mgr.RunBashCf(command, arg...) +} +func RunDefaultWithStdoutBashC(command string) (string, error) { + mgr := NewCommandMgr(WithTimeout(20 * time.Second)) + return mgr.RunWithStdoutBashC(command) +} +func RunDefaultWithStdoutBashCf(command string, arg ...interface{}) (string, error) { + mgr := NewCommandMgr(WithTimeout(20 * time.Second)) + return mgr.RunWithStdoutBashCf(command, arg...) +} + +func (c *CommandHelper) Run(name string, arg ...string) error { + _, err := c.run(name, arg...) + return err +} +func (c *CommandHelper) RunBashCWithArgs(arg ...string) error { + arg = append([]string{"-c"}, arg...) + _, err := c.run("bash", arg...) + return err +} +func (c *CommandHelper) RunBashC(command string) error { + _, err := c.run("bash", "-c", command) + return err +} +func (c *CommandHelper) RunBashCf(command string, arg ...interface{}) error { + _, err := c.run("bash", "-c", fmt.Sprintf(command, arg...)) + return err +} + +func (c *CommandHelper) RunWithStdout(name string, arg ...string) (string, error) { + return c.run(name, arg...) +} +func (c *CommandHelper) RunWithStdoutBashC(command string) (string, error) { + return c.run("bash", "-c", command) +} +func (c *CommandHelper) RunWithStdoutBashCf(command string, arg ...interface{}) (string, error) { + return c.run("bash", "-c", fmt.Sprintf(command, arg...)) +} + +func (c *CommandHelper) run(name string, arg ...string) (string, error) { + cmd := exec.Command(name, arg...) + + customWriter := &CustomWriter{taskItem: c.taskItem} + var stdout, stderr bytes.Buffer + if c.taskItem != nil { + cmd.Stdout = customWriter + cmd.Stderr = customWriter + } else if c.logger != nil { + cmd.Stdout = c.logger.Writer() + cmd.Stderr = c.logger.Writer() + } else if len(c.outputFile) != 0 { + file, err := os.OpenFile(c.outputFile, os.O_WRONLY|os.O_CREATE, constant.FilePerm) + if err != nil { + return "", err + } + defer file.Close() + cmd.Stdout = file + cmd.Stderr = file + } else if len(c.scriptPath) != 0 { + cmd.Stdout = &stdout + cmd.Stderr = &stderr + cmd = exec.Command("bash", c.scriptPath) + } else { + cmd.Stdout = &stdout + cmd.Stderr = &stderr + } + env := os.Environ() + cmd.Env = env + if len(c.workDir) != 0 { + cmd.Dir = c.workDir + } + + if err := cmd.Start(); err != nil { + return "", err + } + if c.timeout != 0 { + done := make(chan error, 1) + go func() { + done <- cmd.Wait() + if c.taskItem != nil { + customWriter.Flush() + } + }() + after := time.After(c.timeout) + select { + case <-after: + _ = cmd.Process.Kill() + return "", buserr.New("ErrCmdTimeout") + case err := <-done: + if err != nil { + return handleErr(stdout, stderr, err) + } + } + return stdout.String(), nil + } + + err := cmd.Run() + if err != nil { + return handleErr(stdout, stderr, err) + } + return stdout.String(), nil +} + +func WithOutputFile(outputFile string) Option { + return func(s *CommandHelper) { + s.outputFile = outputFile + } +} +func WithTimeout(timeout time.Duration) Option { + return func(s *CommandHelper) { + s.timeout = timeout + } +} +func WithLogger(logger *log.Logger) Option { + return func(s *CommandHelper) { + s.logger = logger + } +} +func WithTask(taskItem task.Task) Option { + return func(s *CommandHelper) { + s.taskItem = &taskItem + } +} +func WithWorkDir(workDir string) Option { + return func(s *CommandHelper) { + s.workDir = workDir + } +} +func WithScriptPath(scriptPath string) Option { + return func(s *CommandHelper) { + s.scriptPath = scriptPath + } +} + +type CustomWriter struct { + taskItem *task.Task + buffer bytes.Buffer +} + +func (cw *CustomWriter) Write(p []byte) (n int, err error) { + cw.buffer.Write(p) + lines := strings.Split(cw.buffer.String(), "\n") + + for i := 0; i < len(lines)-1; i++ { + cw.taskItem.Log(lines[i]) + } + cw.buffer.Reset() + cw.buffer.WriteString(lines[len(lines)-1]) + + return len(p), nil +} +func (cw *CustomWriter) Flush() { + if cw.buffer.Len() > 0 { + cw.taskItem.Log(cw.buffer.String()) + cw.buffer.Reset() + } +} + +func handleErr(stdout, stderr bytes.Buffer, err error) (string, error) { + errMsg := "" + if len(stderr.String()) != 0 { + errMsg = fmt.Sprintf("stderr: %s", stderr.String()) + } + if len(stdout.String()) != 0 { + if len(errMsg) != 0 { + errMsg = fmt.Sprintf("%s; stdout: %s", errMsg, stdout.String()) + } else { + errMsg = fmt.Sprintf("stdout: %s", stdout.String()) + } + } + return errMsg, err +} diff --git a/agent/utils/common/common.go b/agent/utils/common/common.go index 11de0834d..af21b82f5 100644 --- a/agent/utils/common/common.go +++ b/agent/utils/common/common.go @@ -270,7 +270,7 @@ func LoadTimeZoneByCmd() string { if _, err := time.LoadLocation(loc); err != nil { loc = "Asia/Shanghai" } - std, err := cmd.Exec("timedatectl | grep 'Time zone'") + std, err := cmd.RunDefaultWithStdoutBashC("timedatectl | grep 'Time zone'") if err != nil { return loc } @@ -394,7 +394,8 @@ func RestartService(core, agent, reload bool) { default: return } - std, err := cmd.ExecWithTimeOut(command, 1*time.Minute) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(1 * time.Second)) + std, err := cmdMgr.RunWithStdoutBashC(command) if err != nil { global.LOG.Errorf("restart 1panel service failed, err: %v, std: %s", err, std) } diff --git a/agent/utils/compose/compose.go b/agent/utils/compose/compose.go index b2e5faa32..c2f73c268 100644 --- a/agent/utils/compose/compose.go +++ b/agent/utils/compose/compose.go @@ -5,35 +5,35 @@ import ( ) func Up(filePath string) (string, error) { - stdout, err := cmd.Execf("docker compose -f %s up -d", filePath) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker compose -f %s up -d", filePath) return stdout, err } func Down(filePath string) (string, error) { - stdout, err := cmd.Execf("docker compose -f %s down --remove-orphans", filePath) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker compose -f %s down --remove-orphans", filePath) return stdout, err } func Stop(filePath string) (string, error) { - stdout, err := cmd.Execf("docker compose -f %s stop", filePath) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker compose -f %s stop", filePath) return stdout, err } func Restart(filePath string) (string, error) { - stdout, err := cmd.Execf("docker compose -f %s restart", filePath) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker compose -f %s restart", filePath) return stdout, err } func Operate(filePath, operation string) (string, error) { - stdout, err := cmd.Execf("docker compose -f %s %s", filePath, operation) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker compose -f %s %s", filePath, operation) return stdout, err } func DownAndUp(filePath string) (string, error) { - stdout, err := cmd.Execf("docker compose -f %s down", filePath) + stdout, err := cmd.RunDefaultWithStdoutBashCf("docker compose -f %s down", filePath) if err != nil { return stdout, err } - stdout, err = cmd.Execf("docker compose -f %s up -d", filePath) + stdout, err = cmd.RunDefaultWithStdoutBashCf("docker compose -f %s up -d", filePath) return stdout, err } diff --git a/agent/utils/files/file_op.go b/agent/utils/files/file_op.go index 5697bc313..f017f60f2 100644 --- a/agent/utils/files/file_op.go +++ b/agent/utils/files/file_op.go @@ -117,11 +117,11 @@ func (f FileOp) DeleteFile(dst string) error { } func (f FileOp) CleanDir(dst string) error { - return cmd.ExecCmd(fmt.Sprintf("rm -rf %s/*", dst)) + return cmd.RunDefaultBashCf("rm -rf %s/*", dst) } func (f FileOp) RmRf(dst string) error { - return cmd.ExecCmd(fmt.Sprintf("rm -rf %s", dst)) + return cmd.RunDefaultBashCf("rm -rf %s", dst) } func (f FileOp) WriteFile(dst string, in io.Reader, mode fs.FileMode) error { @@ -172,14 +172,12 @@ func (f FileOp) SaveFileWithByte(dst string, content []byte, mode fs.FileMode) e } func (f FileOp) ChownR(dst string, uid string, gid string, sub bool) error { - cmdStr := fmt.Sprintf(`chown %s:%s "%s"`, uid, gid, dst) + cmdStr := fmt.Sprintf(`%s chown %s:%s "%s"`, cmd.SudoHandleCmd(), uid, gid, dst) if sub { cmdStr = fmt.Sprintf(`chown -R %s:%s "%s"`, uid, gid, dst) } - if cmd.HasNoPasswordSudo() { - cmdStr = fmt.Sprintf("sudo %s", cmdStr) - } - if msg, err := cmd.ExecWithTimeOut(cmdStr, 10*time.Second); err != nil { + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(10 * time.Second)) + if msg, err := cmdMgr.RunWithStdoutBashC(cmdStr); err != nil { if msg != "" { return errors.New(msg) } @@ -189,14 +187,12 @@ func (f FileOp) ChownR(dst string, uid string, gid string, sub bool) error { } func (f FileOp) ChmodR(dst string, mode int64, sub bool) error { - cmdStr := fmt.Sprintf(`chmod %v "%s"`, fmt.Sprintf("%04o", mode), dst) + cmdStr := fmt.Sprintf(`%s chmod %v "%s"`, cmd.SudoHandleCmd(), fmt.Sprintf("%04o", mode), dst) if sub { - cmdStr = fmt.Sprintf(`chmod -R %v "%s"`, fmt.Sprintf("%04o", mode), dst) + cmdStr = fmt.Sprintf(`%s chmod -R %v "%s"`, cmd.SudoHandleCmd(), fmt.Sprintf("%04o", mode), dst) } - if cmd.HasNoPasswordSudo() { - cmdStr = fmt.Sprintf("sudo %s", cmdStr) - } - if msg, err := cmd.ExecWithTimeOut(cmdStr, 10*time.Second); err != nil { + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(10 * time.Second)) + if msg, err := cmdMgr.RunWithStdoutBashC(cmdStr); err != nil { if msg != "" { return errors.New(msg) } @@ -206,14 +202,12 @@ func (f FileOp) ChmodR(dst string, mode int64, sub bool) error { } func (f FileOp) ChmodRWithMode(dst string, mode fs.FileMode, sub bool) error { - cmdStr := fmt.Sprintf(`chmod %v "%s"`, fmt.Sprintf("%o", mode.Perm()), dst) + cmdStr := fmt.Sprintf(`%s chmod %v "%s"`, cmd.SudoHandleCmd(), fmt.Sprintf("%o", mode.Perm()), dst) if sub { - cmdStr = fmt.Sprintf(`chmod -R %v "%s"`, fmt.Sprintf("%o", mode.Perm()), dst) + cmdStr = fmt.Sprintf(`%s chmod -R %v "%s"`, cmd.SudoHandleCmd(), fmt.Sprintf("%o", mode.Perm()), dst) } - if cmd.HasNoPasswordSudo() { - cmdStr = fmt.Sprintf("sudo %s", cmdStr) - } - if msg, err := cmd.ExecWithTimeOut(cmdStr, 10*time.Second); err != nil { + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(10 * time.Second)) + if msg, err := cmdMgr.RunWithStdoutBashC(cmdStr); err != nil { if msg != "" { return errors.New(msg) } @@ -351,8 +345,7 @@ func (f FileOp) Cut(oldPaths []string, dst, name string, cover bool) error { coverFlag = "-f" } - cmdStr := fmt.Sprintf(`mv %s '%s' '%s'`, coverFlag, p, dstPath) - if err := cmd.ExecCmd(cmdStr); err != nil { + if err := cmd.RunDefaultBashCf(`mv %s '%s' '%s'`, coverFlag, p, dstPath); err != nil { return err } } @@ -360,8 +353,7 @@ func (f FileOp) Cut(oldPaths []string, dst, name string, cover bool) error { } func (f FileOp) Mv(oldPath, dstPath string) error { - cmdStr := fmt.Sprintf(`mv '%s' '%s'`, oldPath, dstPath) - if err := cmd.ExecCmd(cmdStr); err != nil { + if err := cmd.RunDefaultBashCf(`mv '%s' '%s'`, oldPath, dstPath); err != nil { return err } return nil @@ -417,20 +409,19 @@ func (f FileOp) CopyAndReName(src, dst, name string, cover bool) error { if name != "" && !cover { dstPath = filepath.Join(dst, name) } - return cmd.ExecCmd(fmt.Sprintf(`cp -rf '%s' '%s'`, src, dstPath)) + return cmd.RunDefaultBashCf(`cp -rf '%s' '%s'`, src, dstPath) } else { dstPath := filepath.Join(dst, name) if cover { dstPath = dst } - return cmd.ExecCmd(fmt.Sprintf(`cp -f '%s' '%s'`, src, dstPath)) + return cmd.RunDefaultBashCf(`cp -f '%s' '%s'`, src, dstPath) } } func (f FileOp) CopyDirWithNewName(src, dst, newName string) error { dstDir := filepath.Join(dst, newName) - str := fmt.Sprintf(`cp -rf '%s' '%s'`, src, dstDir) - return cmd.ExecCmd(str) + return cmd.RunDefaultBashCf(`cp -rf '%s' '%s'`, src, dstDir) } func (f FileOp) CopyDir(src, dst string) error { @@ -442,7 +433,7 @@ func (f FileOp) CopyDir(src, dst string) error { if err = f.Fs.MkdirAll(dstDir, srcInfo.Mode()); err != nil { return err } - return cmd.ExecCmd(fmt.Sprintf(`cp -rf '%s' '%s'`, src, dst+"/")) + return cmd.RunDefaultBashCf(`cp -rf '%s' '%s'`, src, dst+"/") } func (f FileOp) CopyDirWithExclude(src, dst string, excludeNames []string) error { @@ -455,7 +446,7 @@ func (f FileOp) CopyDirWithExclude(src, dst string, excludeNames []string) error return err } if len(excludeNames) == 0 { - return cmd.ExecCmd(fmt.Sprintf(`cp -rf '%s' '%s'`, src, dst+"/")) + return cmd.RunDefaultBashCf(`cp -rf '%s' '%s'`, src, dst+"/") } tmpFiles, err := os.ReadDir(src) if err != nil { @@ -488,7 +479,7 @@ func (f FileOp) CopyDirWithExclude(src, dst string, excludeNames []string) error func (f FileOp) CopyFile(src, dst string) error { dst = filepath.Clean(dst) + string(filepath.Separator) - return cmd.ExecCmd(fmt.Sprintf(`cp -f '%s' '%s'`, src, dst+"/")) + return cmd.RunDefaultBashCf(`cp -f '%s' '%s'`, src, dst+"/") } func (f FileOp) GetDirSize(path string) (int64, error) { @@ -757,7 +748,9 @@ func (f FileOp) TarGzCompressPro(withDir bool, src, dst, secret, exclusionRules commands = fmt.Sprintf("tar --warning=no-file-changed --ignore-failed-read --exclude-from=<(find %s -type s -printf '%s' | sed 's|^|%s/|') -zcf %s %s %s", src, "%P\n", itemPrefix, dst, exStr, srcItem) global.LOG.Debug(commands) } - return cmd.ExecCmdWithDir(commands, workdir) + + cmdMgr := cmd.NewCommandMgr(cmd.WithWorkDir(workdir)) + return cmdMgr.RunBashC(commands) } func (f FileOp) TarGzFilesWithCompressPro(list []string, dst, secret string) error { @@ -779,7 +772,7 @@ func (f FileOp) TarGzFilesWithCompressPro(list []string, dst, secret string) err commands = fmt.Sprintf("tar --warning=no-file-changed --ignore-failed-read -zcf %s %s", dst, strings.Join(filelist, " ")) global.LOG.Debug(commands) } - return cmd.ExecCmd(commands) + return cmd.RunDefaultBashC(commands) } func (f FileOp) TarGzExtractPro(src, dst string, secret string) error { @@ -797,7 +790,8 @@ func (f FileOp) TarGzExtractPro(src, dst string, secret string) error { commands = fmt.Sprintf("tar zxvf %s", src) global.LOG.Debug(commands) } - return cmd.ExecCmdWithDir(commands, dst) + cmdMgr := cmd.NewCommandMgr(cmd.WithWorkDir(dst)) + return cmdMgr.RunBashC(commands) } func CopyCustomAppFile(srcPath, dstPath string) error { if _, err := os.Stat(srcPath); os.IsNotExist(err) { diff --git a/agent/utils/files/tar.go b/agent/utils/files/tar.go index a4b013c55..0d67e6840 100644 --- a/agent/utils/files/tar.go +++ b/agent/utils/files/tar.go @@ -1,8 +1,6 @@ package files import ( - "fmt" - "github.com/1Panel-dev/1Panel/agent/utils/cmd" ) @@ -19,7 +17,7 @@ func NewTarArchiver(compressType CompressType) ShellArchiver { } func (t TarArchiver) Extract(FilePath string, dstDir string, secret string) error { - return cmd.ExecCmd(fmt.Sprintf("%s %s \"%s\" -C \"%s\"", t.Cmd, t.getOptionStr("extract"), FilePath, dstDir)) + return cmd.RunDefaultBashCf("%s %s \"%s\" -C \"%s\"", t.Cmd, t.getOptionStr("extract"), FilePath, dstDir) } func (t TarArchiver) Compress(sourcePaths []string, dstFile string, secret string) error { diff --git a/agent/utils/files/tar_gz.go b/agent/utils/files/tar_gz.go index e8e95e41a..d9dfadd48 100644 --- a/agent/utils/files/tar_gz.go +++ b/agent/utils/files/tar_gz.go @@ -27,7 +27,7 @@ func (t TarGzArchiver) Extract(filePath, dstDir string, secret string) error { commands = fmt.Sprintf("tar -zxvf '%s' -C '%s' > /dev/null 2>&1", filePath, dstDir) global.LOG.Debug(commands) } - if err = cmd.ExecCmd(commands); err != nil { + if err = cmd.RunDefaultBashC(commands); err != nil { return err } return nil @@ -52,7 +52,7 @@ func (t TarGzArchiver) Compress(sourcePaths []string, dstFile string, secret str commands = fmt.Sprintf("tar -zcf \"%s\" -C \"%s\" %s", dstFile, aheadDir, itemDir) global.LOG.Debug(commands) } - if err := cmd.ExecCmd(commands); err != nil { + if err := cmd.RunDefaultBashC(commands); err != nil { return err } return nil diff --git a/agent/utils/files/zip.go b/agent/utils/files/zip.go index 390605af1..80704a4c6 100644 --- a/agent/utils/files/zip.go +++ b/agent/utils/files/zip.go @@ -23,7 +23,7 @@ func (z ZipArchiver) Extract(filePath, dstDir string, secret string) error { if err := checkCmdAvailability("unzip"); err != nil { return err } - return cmd.ExecCmd(fmt.Sprintf("unzip -qo %s -d %s", filePath, dstDir)) + return cmd.RunDefaultBashCf("unzip -qo %s -d %s", filePath, dstDir) } func (z ZipArchiver) Compress(sourcePaths []string, dstFile string, _ string) error { @@ -41,8 +41,8 @@ func (z ZipArchiver) Compress(sourcePaths []string, dstFile string, _ string) er for i, sp := range sourcePaths { relativePaths[i] = path.Base(sp) } - cmdStr := fmt.Sprintf("zip -qr %s %s", tmpFile, strings.Join(relativePaths, " ")) - if err = cmd.ExecCmdWithDir(cmdStr, baseDir); err != nil { + cmdMgr := cmd.NewCommandMgr(cmd.WithWorkDir(baseDir)) + if err = cmdMgr.Run("zip", "-qr", tmpFile, strings.Join(relativePaths, " ")); err != nil { return err } if err = op.Mv(tmpFile, dstFile); err != nil { diff --git a/agent/utils/firewall/client/firewalld.go b/agent/utils/firewall/client/firewalld.go index 92684ddd9..8fe816f5d 100644 --- a/agent/utils/firewall/client/firewalld.go +++ b/agent/utils/firewall/client/firewalld.go @@ -24,12 +24,12 @@ func (f *Firewall) Name() string { } func (f *Firewall) Status() (bool, error) { - stdout, _ := cmd.Exec("LANGUAGE=en_US:en firewall-cmd --state") + stdout, _ := cmd.RunDefaultWithStdoutBashC("LANGUAGE=en_US:en firewall-cmd --state") return stdout == "running\n", nil } func (f *Firewall) Version() (string, error) { - stdout, err := cmd.Exec("LANGUAGE=en_US:en firewall-cmd --version") + stdout, err := cmd.RunDefaultWithStdoutBashC("LANGUAGE=en_US:en firewall-cmd --version") if err != nil { return "", fmt.Errorf("load the firewall version failed, err: %s", stdout) } @@ -37,7 +37,7 @@ func (f *Firewall) Version() (string, error) { } func (f *Firewall) Start() error { - stdout, err := cmd.Exec("systemctl start firewalld") + stdout, err := cmd.RunDefaultWithStdoutBashC("systemctl start firewalld") if err != nil { return fmt.Errorf("enable the firewall failed, err: %s", stdout) } @@ -45,7 +45,7 @@ func (f *Firewall) Start() error { } func (f *Firewall) Stop() error { - stdout, err := cmd.Exec("systemctl stop firewalld") + stdout, err := cmd.RunDefaultWithStdoutBashC("systemctl stop firewalld") if err != nil { return fmt.Errorf("stop the firewall failed, err: %s", stdout) } @@ -53,7 +53,7 @@ func (f *Firewall) Stop() error { } func (f *Firewall) Restart() error { - stdout, err := cmd.Exec("systemctl restart firewalld") + stdout, err := cmd.RunDefaultWithStdoutBashC("systemctl restart firewalld") if err != nil { return fmt.Errorf("restart the firewall failed, err: %s", stdout) } @@ -61,7 +61,7 @@ func (f *Firewall) Restart() error { } func (f *Firewall) Reload() error { - stdout, err := cmd.Exec("firewall-cmd --reload") + stdout, err := cmd.RunDefaultWithStdoutBashC("firewall-cmd --reload") if err != nil { return fmt.Errorf("reload firewall failed, err: %s", stdout) } @@ -74,7 +74,7 @@ func (f *Firewall) ListPort() ([]FireInfo, error) { wg.Add(2) go func() { defer wg.Done() - stdout, err := cmd.Exec("firewall-cmd --zone=public --list-ports") + stdout, err := cmd.RunDefaultWithStdoutBashC("firewall-cmd --zone=public --list-ports") if err != nil { return } @@ -95,7 +95,7 @@ func (f *Firewall) ListPort() ([]FireInfo, error) { go func() { defer wg.Done() - stdout1, err := cmd.Exec("firewall-cmd --zone=public --list-rich-rules") + stdout1, err := cmd.RunDefaultWithStdoutBashC("firewall-cmd --zone=public --list-rich-rules") if err != nil { return } @@ -118,7 +118,7 @@ func (f *Firewall) ListForward() ([]FireInfo, error) { if err := f.EnableForward(); err != nil { global.LOG.Errorf("init port forward failed, err: %v", err) } - stdout, err := cmd.Exec("firewall-cmd --zone=public --list-forward-ports") + stdout, err := cmd.RunDefaultWithStdoutBashC("firewall-cmd --zone=public --list-forward-ports") if err != nil { return nil, err } @@ -147,7 +147,7 @@ func (f *Firewall) ListForward() ([]FireInfo, error) { } func (f *Firewall) ListAddress() ([]FireInfo, error) { - stdout, err := cmd.Exec("firewall-cmd --zone=public --list-rich-rules") + stdout, err := cmd.RunDefaultWithStdoutBashC("firewall-cmd --zone=public --list-rich-rules") if err != nil { return nil, err } @@ -170,7 +170,7 @@ func (f *Firewall) Port(port FireInfo, operation string) error { return buserr.New("ErrCmdIllegal") } - stdout, err := cmd.Execf("firewall-cmd --zone=public --%s-port=%s/%s --permanent", operation, port.Port, port.Protocol) + stdout, err := cmd.RunDefaultWithStdoutBashCf("firewall-cmd --zone=public --%s-port=%s/%s --permanent", operation, port.Port, port.Protocol) if err != nil { return fmt.Errorf("%s (port: %s/%s strategy: %s) failed, err: %s", operation, port.Port, port.Protocol, port.Strategy, stdout) } @@ -195,12 +195,12 @@ func (f *Firewall) RichRules(rule FireInfo, operation string) error { ruleStr += fmt.Sprintf("protocol=%s ", rule.Protocol) } ruleStr += rule.Strategy - stdout, err := cmd.Execf("firewall-cmd --zone=public --%s-rich-rule '%s' --permanent", operation, ruleStr) + stdout, err := cmd.RunDefaultWithStdoutBashCf("firewall-cmd --zone=public --%s-rich-rule '%s' --permanent", operation, ruleStr) if err != nil { return fmt.Errorf("%s rich rules (%s) failed, err: %s", operation, ruleStr, stdout) } if len(rule.Address) == 0 { - stdout1, err := cmd.Execf("firewall-cmd --zone=public --%s-rich-rule '%s' --permanent", operation, strings.ReplaceAll(ruleStr, "family=ipv4 ", "family=ipv6 ")) + stdout1, err := cmd.RunDefaultWithStdoutBashCf("firewall-cmd --zone=public --%s-rich-rule '%s' --permanent", operation, strings.ReplaceAll(ruleStr, "family=ipv4 ", "family=ipv6 ")) if err != nil { return fmt.Errorf("%s rich rules (%s) failed, err: %s", operation, strings.ReplaceAll(ruleStr, "family=ipv4 ", "family=ipv6 "), stdout1) } @@ -214,7 +214,7 @@ func (f *Firewall) PortForward(info Forward, operation string) error { ruleStr = fmt.Sprintf("firewall-cmd --zone=public --%s-forward-port=port=%s:proto=%s:toaddr=%s:toport=%s --permanent", operation, info.Port, info.Protocol, info.TargetIP, info.TargetPort) } - stdout, err := cmd.Exec(ruleStr) + stdout, err := cmd.RunDefaultWithStdoutBashC(ruleStr) if err != nil { return fmt.Errorf("%s port forward failed, err: %s", operation, stdout) } @@ -247,10 +247,10 @@ func (f *Firewall) loadInfo(line string) FireInfo { } func (f *Firewall) EnableForward() error { - stdout, err := cmd.Exec("firewall-cmd --zone=public --query-masquerade") + stdout, err := cmd.RunDefaultWithStdoutBashC("firewall-cmd --zone=public --query-masquerade") if err != nil { if strings.HasSuffix(strings.TrimSpace(stdout), "no") { - stdout, err = cmd.Exec("firewall-cmd --zone=public --add-masquerade --permanent") + stdout, err = cmd.RunDefaultWithStdoutBashC("firewall-cmd --zone=public --add-masquerade --permanent") if err != nil { return fmt.Errorf("%s: %s", err, stdout) } diff --git a/agent/utils/firewall/client/iptables.go b/agent/utils/firewall/client/iptables.go index 4aa1e1e31..16b3b4e1f 100644 --- a/agent/utils/firewall/client/iptables.go +++ b/agent/utils/firewall/client/iptables.go @@ -20,15 +20,13 @@ type Iptables struct { func NewIptables() (*Iptables, error) { iptables := new(Iptables) - if cmd.HasNoPasswordSudo() { - iptables.CmdStr = "sudo" - } + iptables.CmdStr = cmd.SudoHandleCmd() return iptables, nil } func (iptables *Iptables) run(rule string) error { - stdout, err := cmd.Execf("%s iptables -t nat %s", iptables.CmdStr, rule) + stdout, err := cmd.RunDefaultWithStdoutBashCf("%s iptables -t nat %s", iptables.CmdStr, rule) if err != nil { return fmt.Errorf("%s, %s", err, stdout) } @@ -44,7 +42,7 @@ func (iptables *Iptables) runf(rule string, a ...any) error { } func (iptables *Iptables) Check() error { - stdout, err := cmd.Exec("cat /proc/sys/net/ipv4/ip_forward") + stdout, err := cmd.RunDefaultWithStdoutBashC("cat /proc/sys/net/ipv4/ip_forward") if err != nil { return fmt.Errorf("%s, %s", err, stdout) } @@ -68,7 +66,7 @@ func (iptables *Iptables) NatList(chain ...string) ([]IptablesNatInfo, error) { if len(chain) == 1 { rule = fmt.Sprintf("%s iptables -t nat -nL %s --line", iptables.CmdStr, chain[0]) } - stdout, err := cmd.Exec(rule) + stdout, err := cmd.RunDefaultWithStdoutBashC(rule) if err != nil { return nil, err } diff --git a/agent/utils/firewall/client/ufw.go b/agent/utils/firewall/client/ufw.go index 904fa9bff..f394ee40a 100644 --- a/agent/utils/firewall/client/ufw.go +++ b/agent/utils/firewall/client/ufw.go @@ -15,11 +15,7 @@ type Ufw struct { func NewUfw() (*Ufw, error) { var ufw Ufw - if cmd.HasNoPasswordSudo() { - ufw.CmdStr = "LANGUAGE=en_US:en sudo ufw" - } else { - ufw.CmdStr = "LANGUAGE=en_US:en ufw" - } + ufw.CmdStr = fmt.Sprintf("LANGUAGE=en_US:en %s ufw", cmd.SudoHandleCmd()) return &ufw, nil } @@ -28,11 +24,11 @@ func (f *Ufw) Name() string { } func (f *Ufw) Status() (bool, error) { - stdout, _ := cmd.Execf("%s status | grep Status", f.CmdStr) + stdout, _ := cmd.RunDefaultWithStdoutBashCf("%s status | grep Status", f.CmdStr) if stdout == "Status: active\n" { return true, nil } - stdout1, _ := cmd.Execf("%s status | grep 状态", f.CmdStr) + stdout1, _ := cmd.RunDefaultWithStdoutBashCf("%s status | grep 状态", f.CmdStr) if stdout1 == "状态: 激活\n" { return true, nil } @@ -40,7 +36,7 @@ func (f *Ufw) Status() (bool, error) { } func (f *Ufw) Version() (string, error) { - stdout, err := cmd.Execf("%s version | grep ufw", f.CmdStr) + stdout, err := cmd.RunDefaultWithStdoutBashCf("%s version | grep ufw", f.CmdStr) if err != nil { return "", fmt.Errorf("load the firewall status failed, err: %s", stdout) } @@ -49,7 +45,7 @@ func (f *Ufw) Version() (string, error) { } func (f *Ufw) Start() error { - stdout, err := cmd.Execf("echo y | %s enable", f.CmdStr) + stdout, err := cmd.RunDefaultWithStdoutBashCf("echo y | %s enable", f.CmdStr) if err != nil { return fmt.Errorf("enable the firewall failed, err: %s", stdout) } @@ -57,7 +53,7 @@ func (f *Ufw) Start() error { } func (f *Ufw) Stop() error { - stdout, err := cmd.Execf("%s disable", f.CmdStr) + stdout, err := cmd.RunDefaultWithStdoutBashCf("%s disable", f.CmdStr) if err != nil { return fmt.Errorf("stop the firewall failed, err: %s", stdout) } @@ -79,7 +75,7 @@ func (f *Ufw) Reload() error { } func (f *Ufw) ListPort() ([]FireInfo, error) { - stdout, err := cmd.Execf("%s status verbose", f.CmdStr) + stdout, err := cmd.RunDefaultWithStdoutBashCf("%s status verbose", f.CmdStr) if err != nil { return nil, err } @@ -108,7 +104,7 @@ func (f *Ufw) ListForward() ([]FireInfo, error) { if err != nil { return nil, err } - panelChian, _ := cmd.Execf("%s iptables -t nat -L -n | grep 'Chain 1PANEL'", iptables.CmdStr) + panelChian, _ := cmd.RunDefaultWithStdoutBashCf("%s iptables -t nat -L -n | grep 'Chain 1PANEL'", iptables.CmdStr) if len(strings.ReplaceAll(panelChian, "\n", "")) == 0 { if err := f.EnableForward(); err != nil { global.LOG.Errorf("init port forward failed, err: %v", err) @@ -140,7 +136,7 @@ func (f *Ufw) ListForward() ([]FireInfo, error) { } func (f *Ufw) ListAddress() ([]FireInfo, error) { - stdout, err := cmd.Execf("%s status verbose", f.CmdStr) + stdout, err := cmd.RunDefaultWithStdoutBashCf("%s status verbose", f.CmdStr) if err != nil { return nil, err } @@ -190,7 +186,7 @@ func (f *Ufw) Port(port FireInfo, operation string) error { if len(port.Protocol) != 0 { command += fmt.Sprintf("/%s", port.Protocol) } - stdout, err := cmd.Exec(command) + stdout, err := cmd.RunDefaultWithStdoutBashC(command) if err != nil { return fmt.Errorf("%s (%s) failed, err: %s", operation, command, stdout) } @@ -227,10 +223,10 @@ func (f *Ufw) RichRules(rule FireInfo, operation string) error { ruleStr += fmt.Sprintf("to any port %s ", rule.Port) } - stdout, err := cmd.Exec(ruleStr) + stdout, err := cmd.RunDefaultWithStdoutBashC(ruleStr) if err != nil { if strings.Contains(stdout, "ERROR: Invalid position") || strings.Contains(stdout, "ERROR: 无效位置") { - stdout, err := cmd.Exec(strings.ReplaceAll(ruleStr, "insert 1 ", "")) + stdout, err := cmd.RunDefaultWithStdoutBashC(strings.ReplaceAll(ruleStr, "insert 1 ", "")) if err != nil { return fmt.Errorf("%s rich rules (%s), failed, err: %s", operation, ruleStr, stdout) } diff --git a/agent/utils/ntp/ntp.go b/agent/utils/ntp/ntp.go index 914b22962..c7979bf18 100644 --- a/agent/utils/ntp/ntp.go +++ b/agent/utils/ntp/ntp.go @@ -62,7 +62,7 @@ func GetRemoteTime(site string) (time.Time, error) { func UpdateSystemTime(dateTime string) error { system := runtime.GOOS if system == "linux" { - stdout2, err := cmd.Execf(`%s date -s "%s"`, cmd.SudoHandleCmd(), dateTime) + stdout2, err := cmd.RunDefaultWithStdoutBashCf(`%s date -s "%s"`, cmd.SudoHandleCmd(), dateTime) if err != nil { return fmt.Errorf("update system time failed,stdout: %s, err: %v", stdout2, err) } @@ -74,7 +74,7 @@ func UpdateSystemTime(dateTime string) error { func UpdateSystemTimeZone(timezone string) error { system := runtime.GOOS if system == "linux" { - stdout, err := cmd.Execf(`%s timedatectl set-timezone "%s"`, cmd.SudoHandleCmd(), timezone) + stdout, err := cmd.RunDefaultWithStdoutBashCf(`%s timedatectl set-timezone "%s"`, cmd.SudoHandleCmd(), timezone) if err != nil { return fmt.Errorf("update system time zone failed, stdout: %s, err: %v", stdout, err) } diff --git a/agent/utils/toolbox/fail2ban.go b/agent/utils/toolbox/fail2ban.go index 0026cd5be..881f9cc45 100644 --- a/agent/utils/toolbox/fail2ban.go +++ b/agent/utils/toolbox/fail2ban.go @@ -28,7 +28,7 @@ func NewFail2Ban() (*Fail2ban, error) { if err := initLocalFile(); err != nil { return nil, err } - stdout, err := cmd.Exec("systemctl restart fail2ban.service") + stdout, err := cmd.RunDefaultWithStdoutBashC("systemctl restart fail2ban.service") if err != nil { global.LOG.Errorf("restart fail2ban failed, err: %s", stdout) return nil, err @@ -47,7 +47,7 @@ func (f *Fail2ban) Status() (bool, bool, bool) { } func (f *Fail2ban) Version() string { - stdout, err := cmd.Exec("fail2ban-client version") + stdout, err := cmd.RunDefaultWithStdoutBashC("fail2ban-client version") if err != nil { global.LOG.Errorf("load the fail2ban version failed, err: %s", stdout) return "-" @@ -58,13 +58,13 @@ func (f *Fail2ban) Version() string { func (f *Fail2ban) Operate(operate string) error { switch operate { case "start", "restart", "stop", "enable", "disable": - stdout, err := cmd.Execf("systemctl %s fail2ban.service", operate) + stdout, err := cmd.RunDefaultWithStdoutBashCf("systemctl %s fail2ban.service", operate) if err != nil { return fmt.Errorf("%s the fail2ban.service failed, err: %s", operate, stdout) } return nil case "reload": - stdout, err := cmd.Exec("fail2ban-client reload") + stdout, err := cmd.RunDefaultWithStdoutBashC("fail2ban-client reload") if err != nil { return fmt.Errorf("fail2ban-client reload, err: %s", stdout) } @@ -76,15 +76,15 @@ func (f *Fail2ban) Operate(operate string) error { func (f *Fail2ban) ReBanIPs(ips []string) error { ipItems, _ := f.ListBanned() - stdout, err := cmd.Execf("fail2ban-client unban --all") + stdout, err := cmd.RunDefaultWithStdoutBashCf("fail2ban-client unban --all") if err != nil { - stdout1, err := cmd.Execf("fail2ban-client set sshd banip %s", strings.Join(ipItems, " ")) + stdout1, err := cmd.RunDefaultWithStdoutBashCf("fail2ban-client set sshd banip %s", strings.Join(ipItems, " ")) if err != nil { global.LOG.Errorf("rebanip after fail2ban-client unban --all failed, err: %s", stdout1) } return fmt.Errorf("fail2ban-client unban --all failed, err: %s", stdout) } - stdout1, err := cmd.Execf("fail2ban-client set sshd banip %s", strings.Join(ips, " ")) + stdout1, err := cmd.RunDefaultWithStdoutBashCf("fail2ban-client set sshd banip %s", strings.Join(ips, " ")) if err != nil { return fmt.Errorf("handle `fail2ban-client set sshd banip %s` failed, err: %s", strings.Join(ips, " "), stdout1) } @@ -93,7 +93,7 @@ func (f *Fail2ban) ReBanIPs(ips []string) error { func (f *Fail2ban) ListBanned() ([]string, error) { var lists []string - stdout, err := cmd.Exec("fail2ban-client status sshd | grep 'Banned IP list:'") + stdout, err := cmd.RunDefaultWithStdoutBashC("fail2ban-client status sshd | grep 'Banned IP list:'") if err != nil { return lists, err } @@ -113,7 +113,7 @@ func (f *Fail2ban) ListBanned() ([]string, error) { func (f *Fail2ban) ListIgnore() ([]string, error) { var lists []string - stdout, err := cmd.Exec("fail2ban-client get sshd ignoreip") + stdout, err := cmd.RunDefaultWithStdoutBashC("fail2ban-client get sshd ignoreip") if err != nil { return lists, err } diff --git a/agent/utils/toolbox/pure-ftpd.go b/agent/utils/toolbox/pure-ftpd.go index 870458139..ee9618a3e 100644 --- a/agent/utils/toolbox/pure-ftpd.go +++ b/agent/utils/toolbox/pure-ftpd.go @@ -63,7 +63,7 @@ func NewFtpClient() (*Ftp, error) { groupItem, err := user.LookupGroupId("1000") if err == nil { - stdout2, err := cmd.Execf("useradd -u 1000 -g %s %s", groupItem.Name, "1panel") + stdout2, err := cmd.RunDefaultWithStdoutBashCf("useradd -u 1000 -g %s %s", groupItem.Name, "1panel") if err != nil { return nil, errors.New(stdout2) } @@ -72,11 +72,11 @@ func NewFtpClient() (*Ftp, error) { if err.Error() != user.UnknownGroupIdError("1000").Error() { return nil, err } - stdout, err := cmd.Exec("groupadd -g 1000 1panel") + stdout, err := cmd.RunDefaultWithStdoutBashC("groupadd -g 1000 1panel") if err != nil { return nil, errors.New(string(stdout)) } - stdout2, err := cmd.Exec("useradd -u 1000 -g 1panel 1panel") + stdout2, err := cmd.RunDefaultWithStdoutBashC("useradd -u 1000 -g 1panel 1panel") if err != nil { return nil, errors.New(stdout2) } @@ -93,7 +93,7 @@ func (f *Ftp) Status() (bool, bool) { func (f *Ftp) Operate(operate string) error { switch operate { case "start", "restart", "stop": - stdout, err := cmd.Execf("systemctl %s pure-ftpd.service", operate) + stdout, err := cmd.RunDefaultWithStdoutBashCf("systemctl %s pure-ftpd.service", operate) if err != nil { return fmt.Errorf("%s the pure-ftpd.service failed, err: %s", operate, stdout) } @@ -119,7 +119,7 @@ func (f *Ftp) UserAdd(username, passwd, path string) error { return err } _ = f.Reload() - std2, err := cmd.Execf("chown -R %s:%s %s", f.DefaultUser, f.DefaultGroup, path) + std2, err := cmd.RunDefaultWithStdoutBashCf("chown -R %s:%s %s", f.DefaultUser, f.DefaultGroup, path) if err != nil { return errors.New(std2) } @@ -127,7 +127,7 @@ func (f *Ftp) UserAdd(username, passwd, path string) error { } func (f *Ftp) UserDel(username string) error { - std, err := cmd.Execf("pure-pw userdel %s", username) + std, err := cmd.RunDefaultWithStdoutBashCf("pure-pw userdel %s", username) if err != nil { return errors.New(std) } @@ -186,11 +186,11 @@ func (f *Ftp) SetPasswd(username, passwd string) error { } func (f *Ftp) SetPath(username, path string) error { - std, err := cmd.Execf("pure-pw usermod %s -d %s", username, path) + std, err := cmd.RunDefaultWithStdoutBashCf("pure-pw usermod %s -d %s", username, path) if err != nil { return errors.New(std) } - std2, err := cmd.Execf("chown -R %s:%s %s", f.DefaultUser, f.DefaultGroup, path) + std2, err := cmd.RunDefaultWithStdoutBashCf("chown -R %s:%s %s", f.DefaultUser, f.DefaultGroup, path) if err != nil { return errors.New(std2) } @@ -202,7 +202,7 @@ func (f *Ftp) SetStatus(username, status string) error { if status == constant.StatusDisable { statusItem = "1" } - std, err := cmd.Execf("pure-pw usermod %s -r %s", username, statusItem) + std, err := cmd.RunDefaultWithStdoutBashCf("pure-pw usermod %s -r %s", username, statusItem) if err != nil { return errors.New(std) } @@ -210,7 +210,7 @@ func (f *Ftp) SetStatus(username, status string) error { } func (f *Ftp) LoadList() ([]FtpList, error) { - std, err := cmd.Exec("pure-pw list") + std, err := cmd.RunDefaultWithStdoutBashC("pure-pw list") if err != nil { return nil, errors.New(std) } @@ -221,7 +221,7 @@ func (f *Ftp) LoadList() ([]FtpList, error) { if len(parts) < 2 { continue } - std2, err := cmd.Execf("pure-pw show %s | grep 'Allowed client IPs :'", parts[0]) + std2, err := cmd.RunDefaultWithStdoutBashCf("pure-pw show %s | grep 'Allowed client IPs :'", parts[0]) if err != nil { global.LOG.Errorf("handle pure-pw show %s failed, err: %v", parts[0], std2) continue @@ -237,7 +237,7 @@ func (f *Ftp) LoadList() ([]FtpList, error) { } func (f *Ftp) Reload() error { - std, err := cmd.Exec("pure-pw mkdb") + std, err := cmd.RunDefaultWithStdoutBashC("pure-pw mkdb") if err != nil { return errors.New(std) } @@ -248,7 +248,7 @@ func (f *Ftp) LoadLogs(user, operation string) ([]FtpLog, error) { var logs []FtpLog logItem := "" if _, err := os.Stat("/etc/pure-ftpd/conf"); err != nil && os.IsNotExist(err) { - std, err := cmd.Exec("cat /etc/pure-ftpd/pure-ftpd.conf | grep AltLog | grep clf:") + std, err := cmd.RunDefaultWithStdoutBashC("cat /etc/pure-ftpd/pure-ftpd.conf | grep AltLog | grep clf:") logItem = "/var/log/pureftpd.log" if err == nil && !strings.HasPrefix(std, "#") { logItem = std @@ -257,7 +257,7 @@ func (f *Ftp) LoadLogs(user, operation string) ([]FtpLog, error) { if err != nil { return logs, err } - std, err := cmd.Exec("cat /etc/pure-ftpd/conf/AltLog") + std, err := cmd.RunDefaultWithStdoutBashC("cat /etc/pure-ftpd/conf/AltLog") logItem = "/var/log/pure-ftpd/transfer.log" if err != nil && !strings.HasPrefix(std, "#") { logItem = std @@ -296,7 +296,7 @@ func (f *Ftp) LoadLogs(user, operation string) ([]FtpLog, error) { } func handleGunzip(path string) error { - if _, err := cmd.Execf("gunzip %s", path); err != nil { + if _, err := cmd.RunDefaultWithStdoutBashCf("gunzip %s", path); err != nil { return err } return nil diff --git a/agent/utils/xpack/xpack.go b/agent/utils/xpack/xpack.go index a327b347e..f4c5fd2b9 100644 --- a/agent/utils/xpack/xpack.go +++ b/agent/utils/xpack/xpack.go @@ -29,7 +29,7 @@ func LoadNodeInfo(isBase bool) (model.NodeInfo, error) { } func loadParams(param string) string { - stdout, err := cmd.Execf("grep '^%s=' /usr/local/bin/1pctl | cut -d'=' -f2", param) + stdout, err := cmd.RunDefaultWithStdoutBashCf("grep '^%s=' /usr/local/bin/1pctl | cut -d'=' -f2", param) if err != nil { panic(err) } diff --git a/core/app/service/logs.go b/core/app/service/logs.go index 418862749..c390bca89 100644 --- a/core/app/service/logs.go +++ b/core/app/service/logs.go @@ -109,5 +109,5 @@ func (u *LogService) CleanLogs(logtype string) error { } func writeLogs(version string) { - _, _ = cmd.Execf("curl -sfL %s | sh -s 1p upgrade %s", logs, version) + _, _ = cmd.RunDefaultWithStdoutBashCf("curl -sfL %s | sh -s 1p upgrade %s", logs, version) } diff --git a/core/app/service/setting.go b/core/app/service/setting.go index c60e09851..24f7b9510 100644 --- a/core/app/service/setting.go +++ b/core/app/service/setting.go @@ -171,7 +171,7 @@ func (u *SettingService) UpdateBindInfo(req dto.BindInfo) error { } go func() { time.Sleep(1 * time.Second) - _, err := cmd.Exec("systemctl restart 1panel-core.service") + _, err := cmd.RunDefaultWithStdoutBashC("systemctl restart 1panel-core.service") if err != nil { global.LOG.Errorf("restart system with new bind info failed, err: %v", err) } @@ -222,7 +222,7 @@ func (u *SettingService) UpdatePort(port uint) error { } go func() { time.Sleep(1 * time.Second) - if _, err := cmd.Exec("systemctl restart 1panel-core.service"); err != nil { + if _, err := cmd.RunDefaultWithStdoutBashC("systemctl restart 1panel-core.service"); err != nil { global.LOG.Errorf("restart system port failed, err: %v", err) } }() @@ -243,7 +243,7 @@ func (u *SettingService) UpdateSSL(c *gin.Context, req dto.SSLUpdate) error { _ = os.Remove(path.Join(secretDir, "server.key")) go func() { time.Sleep(1 * time.Second) - _, err := cmd.Exec("systemctl restart 1panel-core.service") + _, err := cmd.RunDefaultWithStdoutBashC("systemctl restart 1panel-core.service") if err != nil { global.LOG.Errorf("restart system failed, err: %v", err) } @@ -341,7 +341,7 @@ func (u *SettingService) UpdateSSL(c *gin.Context, req dto.SSLUpdate) error { if req.SSL != status { go func() { time.Sleep(1 * time.Second) - _, err := cmd.Exec("systemctl restart 1panel-core.service") + _, err := cmd.RunDefaultWithStdoutBashC("systemctl restart 1panel-core.service") if err != nil { global.LOG.Errorf("restart system failed, err: %v", err) } diff --git a/core/app/service/upgrade.go b/core/app/service/upgrade.go index 59a5664d9..c87623abd 100644 --- a/core/app/service/upgrade.go +++ b/core/app/service/upgrade.go @@ -154,7 +154,7 @@ func (u *UpgradeService) Upgrade(req dto.Upgrade) error { u.handleRollback(originalDir, 2) return } - if _, err := cmd.Execf("sed -i -e 's#BASE_DIR=.*#BASE_DIR=%s#g' /usr/local/bin/1pctl", global.CONF.Base.InstallDir); err != nil { + if _, err := cmd.RunDefaultWithStdoutBashCf("sed -i -e 's#BASE_DIR=.*#BASE_DIR=%s#g' /usr/local/bin/1pctl", global.CONF.Base.InstallDir); err != nil { global.LOG.Errorf("upgrade basedir in 1pctl failed, err: %v", err) u.handleRollback(originalDir, 2) return @@ -173,9 +173,10 @@ func (u *UpgradeService) Upgrade(req dto.Upgrade) error { global.CONF.Base.Version = req.Version _ = settingRepo.Update("SystemStatus", "Free") - _, _ = cmd.ExecWithTimeOut("systemctl daemon-reload", 30*time.Second) - _, _ = cmd.ExecWithTimeOut("systemctl restart 1panel-agent.service", 1*time.Second) - _, _ = cmd.ExecWithTimeOut("systemctl restart 1panel-core.service", 1*time.Second) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(10 * time.Second)) + _, _ = cmdMgr.RunWithStdoutBashC("systemctl daemon-reload") + _, _ = cmdMgr.RunWithStdoutBashC("systemctl restart 1panel-agent.service") + _, _ = cmdMgr.RunWithStdoutBashC("systemctl restart 1panel-core.service") }() return nil } @@ -352,7 +353,7 @@ func (u *UpgradeService) loadReleaseNotes(path string) (string, error) { } func loadArch() (string, error) { - std, err := cmd.Exec("uname -a") + std, err := cmd.RunDefaultWithStdoutBashC("uname -a") if err != nil { return "", fmt.Errorf("std: %s, err: %s", std, err.Error()) } diff --git a/core/cmd/server/cmd/restore.go b/core/cmd/server/cmd/restore.go index db94bf60b..0c3be8a34 100644 --- a/core/cmd/server/cmd/restore.go +++ b/core/cmd/server/cmd/restore.go @@ -28,7 +28,7 @@ var restoreCmd = &cobra.Command{ fmt.Println(i18n.GetMsgWithMapForCmd("SudoHelper", map[string]interface{}{"cmd": "sudo 1pctl restore"})) return nil } - stdout, err := cmdUtils.Exec("grep '^BASE_DIR=' /usr/local/bin/1pctl | cut -d'=' -f2") + stdout, err := cmdUtils.RunDefaultWithStdoutBashC("grep '^BASE_DIR=' /usr/local/bin/1pctl | cut -d'=' -f2") if err != nil { return fmt.Errorf("handle load `BASE_DIR` failed, err: %v", err) } @@ -53,16 +53,16 @@ var restoreCmd = &cobra.Command{ return err } sudo := cmdUtils.SudoHandleCmd() - _, _ = cmdUtils.Execf("%s chmod 755 /usr/local/bin/1panel-agent /usr/local/bin/1panel-core", sudo) + _, _ = cmdUtils.RunDefaultWithStdoutBashCf("%s chmod 755 /usr/local/bin/1panel-agent /usr/local/bin/1panel-core", sudo) fmt.Println(i18n.GetMsgByKeyForCmd("RestoreStep2")) if err := files.CopyFile(path.Join(tmpPath, "1pctl"), "/usr/local/bin/1pctl", true); err != nil { return err } - _, _ = cmdUtils.Execf("%s chmod 755 /usr/local/bin/1pctl", sudo) - _, _ = cmdUtils.Execf("cp -r %s /usr/local/bin", path.Join(tmpPath, "lang")) + _, _ = cmdUtils.RunDefaultWithStdoutBashCf("%s chmod 755 /usr/local/bin/1pctl", sudo) + _, _ = cmdUtils.RunDefaultWithStdoutBashCf("cp -r %s /usr/local/bin", path.Join(tmpPath, "lang")) geoPath := path.Join(global.CONF.Base.InstallDir, "1panel/geo") - _, _ = cmdUtils.Execf("mkdir %s && cp %s %s/", geoPath, path.Join(tmpPath, "GeoIP.mmdb"), geoPath) + _, _ = cmdUtils.RunDefaultWithStdoutBashCf("mkdir %s && cp %s %s/", geoPath, path.Join(tmpPath, "GeoIP.mmdb"), geoPath) fmt.Println(i18n.GetMsgByKeyForCmd("RestoreStep3")) if err := files.CopyFile(path.Join(tmpPath, "1panel-core.service"), "/etc/systemd/system/1panel-core.service", true); err != nil { diff --git a/core/cmd/server/cmd/root.go b/core/cmd/server/cmd/root.go index 847a4a5de..280eb208f 100644 --- a/core/cmd/server/cmd/root.go +++ b/core/cmd/server/cmd/root.go @@ -37,7 +37,7 @@ type setting struct { } func loadDBConn() (*gorm.DB, error) { - stdout, err := cmdUtils.Exec("grep '^BASE_DIR=' /usr/local/bin/1pctl | cut -d'=' -f2") + stdout, err := cmdUtils.RunDefaultWithStdoutBashC("grep '^BASE_DIR=' /usr/local/bin/1pctl | cut -d'=' -f2") if err != nil { return nil, fmt.Errorf("handle load `BASE_DIR` failed, err: %v", err) } diff --git a/core/cmd/server/cmd/update.go b/core/cmd/server/cmd/update.go index 14521c005..8f9329192 100644 --- a/core/cmd/server/cmd/update.go +++ b/core/cmd/server/cmd/update.go @@ -209,7 +209,7 @@ func port() { fmt.Println("\n" + i18n.GetMsgByKeyForCmd("UpdateSuccessful")) fmt.Println(i18n.GetMsgWithMapForCmd("UpdatePortResult", map[string]interface{}{"name": newPortStr})) - std, err := cmd.Exec("1pctl restart") + std, err := cmd.RunDefaultWithStdoutBashC("1pctl restart") if err != nil { fmt.Println(std) } diff --git a/core/init/hook/hook.go b/core/init/hook/hook.go index 6b78193d9..be8ca9210 100644 --- a/core/init/hook/hook.go +++ b/core/init/hook/hook.go @@ -68,8 +68,7 @@ func handleUserInfo(tags string, settingRepo repo.ISettingRepo) { } } - sudo := cmd.SudoHandleCmd() - _, _ = cmd.Execf("%s sed -i '/CHANGE_USER_INFO=%v/d' /usr/local/bin/1pctl", sudo, global.CONF.Base.ChangeUserInfo) + _, _ = cmd.RunDefaultWithStdoutBashCf("%s sed -i '/CHANGE_USER_INFO=%v/d' /usr/local/bin/1pctl", cmd.SudoHandleCmd(), global.CONF.Base.ChangeUserInfo) } func generateKey() { diff --git a/core/init/lang/lang.go b/core/init/lang/lang.go index 976e60964..d3341ab22 100644 --- a/core/init/lang/lang.go +++ b/core/init/lang/lang.go @@ -61,7 +61,7 @@ func initLang() { downloadLangFromRemote() return } - std, err := cmd.Execf("cp -r %s %s", path.Join(tmpPath, "lang"), "/usr/local/bin/") + std, err := cmd.RunDefaultWithStdoutBashCf("cp -r %s %s", path.Join(tmpPath, "lang"), "/usr/local/bin/") if err != nil { global.LOG.Errorf("load lang from package failed, std: %s, err: %v", std, err) return @@ -73,7 +73,7 @@ func initLang() { downloadGeoFromRemote(geoPath) return } - std, err := cmd.Execf("mkdir %s && cp %s %s/", path.Dir(geoPath), path.Join(tmpPath, "GeoIP.mmdb"), path.Dir(geoPath)) + std, err := cmd.RunDefaultWithStdoutBashCf("mkdir %s && cp %s %s/", path.Dir(geoPath), path.Join(tmpPath, "GeoIP.mmdb"), path.Dir(geoPath)) if err != nil { global.LOG.Errorf("load geo ip from package failed, std: %s, err: %v", std, err) return @@ -115,7 +115,7 @@ func downloadLangFromRemote() { global.LOG.Error("download lang.tar.gz failed, no such file") return } - std, err := cmd.Execf("tar zxvfC %s %s", "/usr/local/bin/lang.tar.gz", "/usr/local/bin/") + std, err := cmd.RunDefaultWithStdoutBashCf("tar zxvfC %s %s", "/usr/local/bin/lang.tar.gz", "/usr/local/bin/") if err != nil { fmt.Printf("decompress lang.tar.gz failed, std: %s, err: %v", std, err) return diff --git a/core/init/viper/viper.go b/core/init/viper/viper.go index 7653fee05..615636666 100644 --- a/core/init/viper/viper.go +++ b/core/init/viper/viper.go @@ -101,7 +101,7 @@ func Init() { } func loadParams(param string) string { - stdout, err := cmd.Execf("grep '^%s=' /usr/local/bin/1pctl | cut -d'=' -f2", param) + stdout, err := cmd.RunDefaultWithStdoutBashCf("grep '^%s=' /usr/local/bin/1pctl | cut -d'=' -f2", param) if err != nil { panic(err) } @@ -113,7 +113,7 @@ func loadParams(param string) string { } func loadChangeInfo() string { - stdout, err := cmd.Exec("grep '^CHANGE_USER_INFO=' /usr/local/bin/1pctl | cut -d'=' -f2") + stdout, err := cmd.RunDefaultWithStdoutBashC("grep '^CHANGE_USER_INFO=' /usr/local/bin/1pctl | cut -d'=' -f2") if err != nil { return "" } diff --git a/core/utils/cmd/cmd.go b/core/utils/cmd/cmd.go index 45e255366..50eaf5ee3 100644 --- a/core/utils/cmd/cmd.go +++ b/core/utils/cmd/cmd.go @@ -1,19 +1,10 @@ package cmd import ( - "bytes" - "fmt" "os/exec" "strings" - "time" - - "github.com/1Panel-dev/1Panel/core/buserr" ) -func Exec(cmdStr string) (string, error) { - return ExecWithTimeOut(cmdStr, 20*time.Second) -} - func SudoHandleCmd() string { cmd := exec.Command("sudo", "-n", "ls") if err := cmd.Run(); err == nil { @@ -23,62 +14,9 @@ func SudoHandleCmd() string { } func Which(name string) bool { - stdout, err := Execf("which %s", name) + stdout, err := RunDefaultWithStdoutBashCf("which %s", name) if err != nil || (len(strings.ReplaceAll(stdout, "\n", "")) == 0) { return false } return true } - -func Execf(cmdStr string, a ...interface{}) (string, error) { - cmd := exec.Command("bash", "-c", fmt.Sprintf(cmdStr, a...)) - var stdout, stderr bytes.Buffer - cmd.Stdout = &stdout - cmd.Stderr = &stderr - err := cmd.Run() - if err != nil { - return handleErr(stdout, stderr, err) - } - return stdout.String(), nil -} - -func ExecWithTimeOut(cmdStr string, timeout time.Duration) (string, error) { - cmd := exec.Command("bash", "-c", cmdStr) - var stdout, stderr bytes.Buffer - cmd.Stdout = &stdout - cmd.Stderr = &stderr - if err := cmd.Start(); err != nil { - return "", err - } - done := make(chan error, 1) - go func() { - done <- cmd.Wait() - }() - after := time.After(timeout) - select { - case <-after: - _ = cmd.Process.Kill() - return "", buserr.New("ErrCmdTimeout") - case err := <-done: - if err != nil { - return handleErr(stdout, stderr, err) - } - } - - return stdout.String(), nil -} - -func handleErr(stdout, stderr bytes.Buffer, err error) (string, error) { - errMsg := "" - if len(stderr.String()) != 0 { - errMsg = fmt.Sprintf("stderr: %s", stderr.String()) - } - if len(stdout.String()) != 0 { - if len(errMsg) != 0 { - errMsg = fmt.Sprintf("%s; stdout: %s", errMsg, stdout.String()) - } else { - errMsg = fmt.Sprintf("stdout: %s", stdout.String()) - } - } - return errMsg, err -} diff --git a/core/utils/cmd/cmdx.go b/core/utils/cmd/cmdx.go new file mode 100644 index 000000000..e541b1604 --- /dev/null +++ b/core/utils/cmd/cmdx.go @@ -0,0 +1,213 @@ +package cmd + +import ( + "bytes" + "fmt" + "log" + "os" + "os/exec" + "strings" + "time" + + "github.com/1Panel-dev/1Panel/core/app/task" + "github.com/1Panel-dev/1Panel/core/buserr" + "github.com/1Panel-dev/1Panel/core/constant" +) + +type CommandHelper struct { + workDir string + outputFile string + scriptPath string + timeout time.Duration + taskItem *task.Task + logger *log.Logger +} + +type Option func(*CommandHelper) + +func NewCommandMgr(opts ...Option) *CommandHelper { + s := &CommandHelper{} + for _, opt := range opts { + opt(s) + } + return s +} + +func RunDefaultBashC(command string) error { + mgr := NewCommandMgr() + return mgr.RunBashC(command) +} +func RunDefaultBashCf(command string, arg ...interface{}) error { + mgr := NewCommandMgr() + return mgr.RunBashCf(command, arg...) +} +func RunDefaultWithStdoutBashC(command string) (string, error) { + mgr := NewCommandMgr(WithTimeout(20 * time.Second)) + return mgr.RunWithStdoutBashC(command) +} +func RunDefaultWithStdoutBashCf(command string, arg ...interface{}) (string, error) { + mgr := NewCommandMgr(WithTimeout(20 * time.Second)) + return mgr.RunWithStdoutBashCf(command, arg...) +} + +func (c *CommandHelper) Run(name string, arg ...string) error { + _, err := c.run(name, arg...) + return err +} +func (c *CommandHelper) RunBashCWithArgs(arg ...string) error { + arg = append([]string{"-c"}, arg...) + _, err := c.run("bash", arg...) + return err +} +func (c *CommandHelper) RunBashC(command string) error { + _, err := c.run("bash", "-c", command) + return err +} +func (c *CommandHelper) RunBashCf(command string, arg ...interface{}) error { + _, err := c.run("bash", "-c", fmt.Sprintf(command, arg...)) + return err +} + +func (c *CommandHelper) RunWithStdout(name string, arg ...string) (string, error) { + return c.run(name, arg...) +} +func (c *CommandHelper) RunWithStdoutBashC(command string) (string, error) { + return c.run("bash", "-c", command) +} +func (c *CommandHelper) RunWithStdoutBashCf(command string, arg ...interface{}) (string, error) { + return c.run("bash", "-c", fmt.Sprintf(command, arg...)) +} + +func (c *CommandHelper) run(name string, arg ...string) (string, error) { + cmd := exec.Command(name, arg...) + + customWriter := &CustomWriter{taskItem: c.taskItem} + var stdout, stderr bytes.Buffer + if c.taskItem != nil { + cmd.Stdout = customWriter + cmd.Stderr = customWriter + } else if c.logger != nil { + cmd.Stdout = c.logger.Writer() + cmd.Stderr = c.logger.Writer() + } else if len(c.outputFile) != 0 { + file, err := os.OpenFile(c.outputFile, os.O_WRONLY|os.O_CREATE, constant.FilePerm) + if err != nil { + return "", err + } + defer file.Close() + cmd.Stdout = file + cmd.Stderr = file + } else if len(c.scriptPath) != 0 { + cmd.Stdout = &stdout + cmd.Stderr = &stderr + cmd = exec.Command("bash", c.scriptPath) + } else { + cmd.Stdout = &stdout + cmd.Stderr = &stderr + } + env := os.Environ() + cmd.Env = env + if len(c.workDir) != 0 { + cmd.Dir = c.workDir + } + + if err := cmd.Start(); err != nil { + return "", err + } + if c.timeout != 0 { + done := make(chan error, 1) + go func() { + done <- cmd.Wait() + if c.taskItem != nil { + customWriter.Flush() + } + }() + after := time.After(c.timeout) + select { + case <-after: + _ = cmd.Process.Kill() + return "", buserr.New("ErrCmdTimeout") + case err := <-done: + if err != nil { + return handleErr(stdout, stderr, err) + } + } + return stdout.String(), nil + } + + err := cmd.Run() + if err != nil { + return handleErr(stdout, stderr, err) + } + return stdout.String(), nil +} + +func WithOutputFile(outputFile string) Option { + return func(s *CommandHelper) { + s.outputFile = outputFile + } +} +func WithTimeout(timeout time.Duration) Option { + return func(s *CommandHelper) { + s.timeout = timeout + } +} +func WithLogger(logger *log.Logger) Option { + return func(s *CommandHelper) { + s.logger = logger + } +} +func WithTask(taskItem task.Task) Option { + return func(s *CommandHelper) { + s.taskItem = &taskItem + } +} +func WithWorkDir(workDir string) Option { + return func(s *CommandHelper) { + s.workDir = workDir + } +} +func WithScriptPath(scriptPath string) Option { + return func(s *CommandHelper) { + s.scriptPath = scriptPath + } +} + +type CustomWriter struct { + taskItem *task.Task + buffer bytes.Buffer +} + +func (cw *CustomWriter) Write(p []byte) (n int, err error) { + cw.buffer.Write(p) + lines := strings.Split(cw.buffer.String(), "\n") + + for i := 0; i < len(lines)-1; i++ { + cw.taskItem.Log(lines[i]) + } + cw.buffer.Reset() + cw.buffer.WriteString(lines[len(lines)-1]) + + return len(p), nil +} +func (cw *CustomWriter) Flush() { + if cw.buffer.Len() > 0 { + cw.taskItem.Log(cw.buffer.String()) + cw.buffer.Reset() + } +} + +func handleErr(stdout, stderr bytes.Buffer, err error) (string, error) { + errMsg := "" + if len(stderr.String()) != 0 { + errMsg = fmt.Sprintf("stderr: %s", stderr.String()) + } + if len(stdout.String()) != 0 { + if len(errMsg) != 0 { + errMsg = fmt.Sprintf("%s; stdout: %s", errMsg, stdout.String()) + } else { + errMsg = fmt.Sprintf("stdout: %s", stdout.String()) + } + } + return errMsg, err +} diff --git a/core/utils/common/common.go b/core/utils/common/common.go index 8b8e46031..54eb1cdd8 100644 --- a/core/utils/common/common.go +++ b/core/utils/common/common.go @@ -2,8 +2,6 @@ package common import ( "fmt" - "github.com/1Panel-dev/1Panel/core/constant" - "github.com/gin-gonic/gin" mathRand "math/rand" "net" "net/http" @@ -13,6 +11,9 @@ import ( "strings" "time" + "github.com/1Panel-dev/1Panel/core/constant" + "github.com/gin-gonic/gin" + "github.com/1Panel-dev/1Panel/core/global" "github.com/1Panel-dev/1Panel/core/utils/cmd" ) @@ -42,7 +43,7 @@ func LoadTimeZoneByCmd() string { if _, err := time.LoadLocation(loc); err != nil { loc = "Asia/Shanghai" } - std, err := cmd.Exec("timedatectl | grep 'Time zone'") + std, err := cmd.RunDefaultWithStdoutBashC("timedatectl | grep 'Time zone'") if err != nil { return loc } @@ -116,7 +117,7 @@ func SplitStr(str string, spi ...string) []string { } func LoadArch() (string, error) { - std, err := cmd.Exec("uname -a") + std, err := cmd.RunDefaultWithStdoutBashC("uname -a") if err != nil { return "", fmt.Errorf("std: %s, err: %s", std, err.Error()) } diff --git a/core/utils/files/files.go b/core/utils/files/files.go index c0f236d75..3718e2aeb 100644 --- a/core/utils/files/files.go +++ b/core/utils/files/files.go @@ -66,7 +66,7 @@ func CopyItem(isDir, withName bool, src, dst string) error { if !isDir { cmdStr = fmt.Sprintf(`cp -f %s %s`, src, dst+"/") } - stdout, err := cmd.Exec(cmdStr) + stdout, err := cmd.RunDefaultWithStdoutBashC(cmdStr) if err != nil { return fmt.Errorf("handle %s failed, stdout: %s, err: %v", cmdStr, stdout, err) } @@ -111,7 +111,8 @@ func HandleTar(sourceDir, targetDir, name, exclusionRules string, secret string) commands = fmt.Sprintf("tar --warning=no-file-changed --ignore-failed-read -zcf %s %s %s", targetDir+"/"+name, excludeRules, path) global.LOG.Debug(commands) } - stdout, err := cmd.ExecWithTimeOut(commands, 24*time.Hour) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(24 * time.Hour)) + stdout, err := cmdMgr.RunWithStdoutBashC(commands) if err != nil { if len(stdout) != 0 { global.LOG.Errorf("do handle tar failed, stdout: %s, err: %v", stdout, err) @@ -137,7 +138,8 @@ func HandleUnTar(sourceFile, targetDir string, secret string) error { global.LOG.Debug(commands) } - stdout, err := cmd.ExecWithTimeOut(commands, 24*time.Hour) + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(24 * time.Hour)) + stdout, err := cmdMgr.RunWithStdoutBashC(commands) if err != nil { global.LOG.Errorf("do handle untar failed, stdout: %s, err: %v", stdout, err) return errors.New(stdout) diff --git a/core/utils/firewall/firewall.go b/core/utils/firewall/firewall.go index b9a175404..d87cb3810 100644 --- a/core/utils/firewall/firewall.go +++ b/core/utils/firewall/firewall.go @@ -9,7 +9,7 @@ import ( func UpdatePort(oldPort, newPort string) error { firewalld := cmd.Which("firewalld") if firewalld { - status, _ := cmd.Exec("LANGUAGE=en_US:en firewall-cmd --state") + status, _ := cmd.RunDefaultWithStdoutBashC("LANGUAGE=en_US:en firewall-cmd --state") isRunning := status == "running\n" if isRunning { return firewallUpdatePort(oldPort, newPort) @@ -20,7 +20,7 @@ func UpdatePort(oldPort, newPort string) error { if !ufw { return nil } - status, _ := cmd.Exec("LANGUAGE=en_US:en ufw status | grep Status") + status, _ := cmd.RunDefaultWithStdoutBashC("LANGUAGE=en_US:en ufw status | grep Status") isRuning := status == "Status: active\n" if isRuning { return ufwUpdatePort(oldPort, newPort) @@ -29,22 +29,22 @@ func UpdatePort(oldPort, newPort string) error { } func firewallUpdatePort(oldPort, newPort string) error { - stdout, err := cmd.Execf("firewall-cmd --zone=public --add-port=%s/tcp --permanent", newPort) + stdout, err := cmd.RunDefaultWithStdoutBashCf("firewall-cmd --zone=public --add-port=%s/tcp --permanent", newPort) if err != nil { return fmt.Errorf("add (port: %s/tcp) failed, err: %s", newPort, stdout) } - _, _ = cmd.Execf("firewall-cmd --zone=public --remove-port=%s/tcp --permanent", oldPort) - _, _ = cmd.Exec("firewall-cmd --reload") + _, _ = cmd.RunDefaultWithStdoutBashCf("firewall-cmd --zone=public --remove-port=%s/tcp --permanent", oldPort) + _, _ = cmd.RunDefaultWithStdoutBashC("firewall-cmd --reload") return nil } func ufwUpdatePort(oldPort, newPort string) error { - stdout, err := cmd.Execf("ufw allow %s", newPort) + stdout, err := cmd.RunDefaultWithStdoutBashCf("ufw allow %s", newPort) if err != nil { return fmt.Errorf("add (port: %s/tcp) failed, err: %s", newPort, stdout) } - _, _ = cmd.Execf("ufw delete allow %s", oldPort) + _, _ = cmd.RunDefaultWithStdoutBashCf("ufw delete allow %s", oldPort) return nil }