Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 24 additions & 3 deletions agent/app/api/v2/database_redis.go
Original file line number Diff line number Diff line change
Expand Up @@ -86,17 +86,38 @@ func (b *BaseApi) CheckHasCli(c *gin.Context) {

// @Tags Database Redis
// @Summary Install redis-cli
// @Success 200
// @Accept json
// @Param request body dto.RedisCliInstall true "request"
// @Success 200 {object} dto.RedisCliStatus
// @Security ApiKeyAuth
// @Security Timestamp
// @Router /databases/redis/install/cli [post]
func (b *BaseApi) InstallCli(c *gin.Context) {
if err := redisService.InstallCli(); err != nil {
var req dto.RedisCliInstall
if err := helper.CheckBindAndValidate(&req, c); err != nil {
return
}
data, err := redisService.InstallCli(req)
if err != nil {
helper.InternalServer(c, err)
return
}
helper.SuccessWithData(c, data)
}

helper.Success(c)
// @Tags Database Redis
// @Summary Load redis-cli installation status
// @Success 200 {object} dto.RedisCliStatus
// @Security ApiKeyAuth
// @Security Timestamp
// @Router /databases/redis/cli/status [get]
func (b *BaseApi) LoadRedisCliStatus(c *gin.Context) {
data, err := redisService.LoadCliStatus()
if err != nil {
helper.InternalServer(c, err)
return
}
helper.SuccessWithData(c, data)
}

// @Tags Database Redis
Expand Down
11 changes: 11 additions & 0 deletions agent/app/dto/database.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,17 @@ type DBBaseInfo struct {
Port int64 `json:"port"`
}

type RedisCliInstall struct {
TaskID string `json:"taskID" validate:"omitempty,uuid"`
}

type RedisCliStatus struct {
Installed bool `json:"installed"`
TaskID string `json:"taskID"`
Status string `json:"status"`
ErrorMsg string `json:"errorMsg"`
}

// mysql
type MysqlDBSearch struct {
PageInfo
Expand Down
14 changes: 13 additions & 1 deletion agent/app/service/container.go
Original file line number Diff line number Diff line change
Expand Up @@ -484,6 +484,10 @@ func (u *ContainerService) LoadResourceLimit() (*dto.ResourceLimit, error) {
}

func (u *ContainerService) ContainerCreate(req dto.ContainerOperate, inThread bool) error {
return u.containerCreate(req, inThread, "")
}

func (u *ContainerService) containerCreate(req dto.ContainerOperate, inThread bool, taskName string) error {
client, err := docker.NewDockerClient()
if err != nil {
return err
Expand All @@ -497,7 +501,10 @@ func (u *ContainerService) ContainerCreate(req dto.ContainerOperate, inThread bo
return buserr.New("ErrContainerName")
}

taskItem, err := task.NewTaskWithOps(req.Name, task.TaskCreate, task.TaskScopeContainer, req.TaskID, 1)
if taskName == "" {
taskName = task.GetTaskName(req.Name, task.TaskCreate, task.TaskScopeContainer)
}
taskItem, err := task.NewTask(taskName, task.TaskCreate, task.TaskScopeContainer, req.TaskID, 1)
if err != nil {
unlock()
_ = client.Close()
Expand Down Expand Up @@ -558,6 +565,11 @@ func (u *ContainerService) ContainerCreate(req dto.ContainerOperate, inThread bo
}, nil)

if inThread {
if err := taskItem.Prepare(); err != nil {
unlock()
_ = client.Close()
return err
}
go func() {
defer unlock()
defer client.Close()
Expand Down
62 changes: 58 additions & 4 deletions agent/app/service/database_redis.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,12 @@ import (
"os"
"os/exec"
"strings"
"sync"

"github.com/1Panel-dev/1Panel/agent/app/task"
"github.com/1Panel-dev/1Panel/agent/buserr"
"github.com/google/uuid"
"gorm.io/gorm"

"github.com/1Panel-dev/1Panel/agent/app/repo"
"github.com/1Panel-dev/1Panel/agent/global"
Expand All @@ -23,6 +29,11 @@ import (

type RedisService struct{}

const redisCliTaskName = "RedisCliEnable"

// The CLI container is shared by all remote Redis databases on this node.
var redisCliInstallMutex sync.Mutex

type IRedisService interface {
UpdateConf(req dto.RedisConfUpdate) error
UpdatePersistenceConf(req dto.RedisConfPersistenceUpdate) error
Expand All @@ -33,7 +44,8 @@ type IRedisService interface {
LoadPersistenceConf(req dto.LoadRedisStatus) (*dto.RedisPersistence, error)

CheckHasCli() bool
InstallCli() error
InstallCli(req dto.RedisCliInstall) (*dto.RedisCliStatus, error)
LoadCliStatus() (*dto.RedisCliStatus, error)
}

func NewIRedisService() IRedisService {
Expand Down Expand Up @@ -71,20 +83,62 @@ func (u *RedisService) CheckHasCli() bool {
return false
}
for _, item := range containerLists {
if strings.ReplaceAll(item.Names[0], "/", "") == "1Panel-redis-cli-tools" {
if len(item.Names) > 0 && strings.TrimPrefix(item.Names[0], "/") == "1Panel-redis-cli-tools" {
return true
}
}
return false
}

func (u *RedisService) InstallCli() error {
func (u *RedisService) LoadCliStatus() (*dto.RedisCliStatus, error) {
result := &dto.RedisCliStatus{}
latest, err := taskRepo.GetFirst(repo.WithByName(redisCliTaskName), repo.WithByType(task.TaskScopeContainer), repo.WithOrderDesc("created_at"))
if errors.Is(err, gorm.ErrRecordNotFound) {
result.Installed = u.CheckHasCli()
return result, nil
}
if err != nil {
return nil, err
}
result.TaskID = latest.ID
result.Status = latest.Status
result.ErrorMsg = latest.ErrorMsg
result.Installed = u.CheckHasCli()
return result, nil
}

func (u *RedisService) InstallCli(req dto.RedisCliInstall) (*dto.RedisCliStatus, error) {
if !redisCliInstallMutex.TryLock() {
return nil, buserr.New("TaskIsExecuting")
}
defer redisCliInstallMutex.Unlock()
status, err := u.LoadCliStatus()
if err != nil {
return nil, err
}
if status.Status == constant.StatusExecuting || status.Installed {
return status, nil
}
if req.TaskID == "" {
req.TaskID = uuid.NewString()
}
// Never reuse an existing task ID: doing so would truncate its log.
if _, err := taskRepo.GetFirst(taskRepo.WithByID(req.TaskID)); !errors.Is(err, gorm.ErrRecordNotFound) {
if err != nil {
return nil, err
}
return nil, buserr.New("TaskIsExecuting")
}
item := dto.ContainerOperate{
TaskID: req.TaskID,
Name: "1Panel-redis-cli-tools",
Image: "redis:7.4.4",
Networks: []dto.ContainerNetwork{{Network: "1panel-network"}},
}
return NewIContainerService().ContainerCreate(item, false)
if err := (&ContainerService{}).containerCreate(item, true, redisCliTaskName); err != nil {
return nil, err
}
return &dto.RedisCliStatus{TaskID: req.TaskID, Status: constant.StatusExecuting}, nil
}

func (u *RedisService) ChangePassword(req dto.ChangeRedisPass) error {
Expand Down
24 changes: 8 additions & 16 deletions agent/app/service/image.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import (
"context"
"encoding/base64"
"encoding/json"
"errors"
"fmt"
"io"
"os"
Expand Down Expand Up @@ -325,18 +324,16 @@ func (u *ImageService) ImageLoad(req dto.ImageLoad) error {
}

go func() {
client, err := docker.NewDockerClient()
if err != nil {
taskItem.Log("Failed to create Docker client: " + err.Error())
return
}
defer client.Close()

for _, itemPath := range req.Paths {
currentPath := itemPath
itemName := path.Base(currentPath)
taskItem.AddSubTask(i18n.GetWithName("TaskImport", itemName), func(t *task.Task) error {
taskItem.Logf("----------------- %s -----------------", itemName)
client, err := docker.NewDockerClient()
if err != nil {
return err
}
defer client.Close()
file, err := os.Open(currentPath)
if err != nil {
return err
Expand All @@ -347,14 +344,9 @@ func (u *ImageService) ImageLoad(req dto.ImageLoad) error {
return err
}
defer res.Body.Close()
content, err := io.ReadAll(res.Body)
if err != nil {
return err
}
if strings.Contains(string(content), "Error") {
return errors.New(string(content))
}
return nil
return consumeImageLoadResponse(res.Body, func(message string) {
taskItem.Log(message)
})
}, nil)
}
_ = taskItem.Execute()
Expand Down
37 changes: 37 additions & 0 deletions agent/app/service/image_load_response.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
package service

import (
"encoding/json"
"errors"
"io"
"strings"
)

// Docker may report load failures in a successful HTTP response's JSON stream.
func consumeImageLoadResponse(reader io.Reader, log func(string)) error {
decoder := json.NewDecoder(reader)
for {
var message struct {
Stream string `json:"stream"`
Error string `json:"error"`
ErrorDetail struct {
Message string `json:"message"`
} `json:"errorDetail"`
}
if err := decoder.Decode(&message); err != nil {
if errors.Is(err, io.EOF) {
return nil
}
return err
}
if message.Error != "" {
return errors.New(message.Error)
}
if message.ErrorDetail.Message != "" {
return errors.New(message.ErrorDetail.Message)
}
if text := strings.TrimSpace(message.Stream); text != "" {
log(text)
}
}
}
12 changes: 11 additions & 1 deletion agent/app/task/task.go
Original file line number Diff line number Diff line change
Expand Up @@ -283,8 +283,18 @@ func (t *Task) updateTask(task *model.Task) {
_ = t.taskRepo.Update(context.Background(), task)
}

func (t *Task) Execute() error {
// Prepare makes a task visible before dispatching it to a background worker.
func (t *Task) Prepare() error {
if err := t.taskRepo.Save(context.Background(), t.Task); err != nil {
_ = t.logFile.Close()
global.RemoveTaskCancel(t.TaskID)
return err
}
return nil
}

func (t *Task) Execute() error {
if err := t.Prepare(); err != nil {
return err
}
var err error
Expand Down
1 change: 1 addition & 0 deletions agent/router/ro_database.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ func (s *DatabaseRouter) InitRouter(Router *gin.RouterGroup) {
cmdRouter.POST("/redis/status", baseApi.LoadRedisStatus)
cmdRouter.POST("/redis/conf", baseApi.LoadRedisConf)
cmdRouter.GET("/redis/check", baseApi.CheckHasCli)
cmdRouter.GET("/redis/cli/status", baseApi.LoadRedisCliStatus)
cmdRouter.POST("/redis/install/cli", baseApi.InstallCli)
cmdRouter.POST("/redis/password", baseApi.ChangeRedisPassword)
cmdRouter.POST("/redis/conf/update", baseApi.UpdateRedisConf)
Expand Down
Loading
Loading