refactor: update agent and core utilities (#12621)

This commit is contained in:
ssongliu
2026-04-28 11:17:30 +08:00
committed by zhengkunwang223
parent 84a27de792
commit d55982bde0
86 changed files with 1854 additions and 924 deletions

View File

@@ -1,7 +1,6 @@
package client
import (
"bytes"
"compress/gzip"
"context"
"errors"
@@ -14,8 +13,8 @@ import (
"time"
"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"
"github.com/1Panel-dev/1Panel/agent/utils/files"
)
@@ -135,29 +134,15 @@ func (r *Local) Backup(info BackupInfo) error {
return fmt.Errorf("mkdir %s failed, err: %v", info.TargetDir, err)
}
}
outfile, err := os.OpenFile(path.Join(info.TargetDir, info.FileName), os.O_RDWR|os.O_CREATE, constant.DirPerm)
if err != nil {
return fmt.Errorf("open file %s failed, err: %v", path.Join(info.TargetDir, info.FileName), err)
}
defer outfile.Close()
global.LOG.Infof("start to pg_dump | gzip > %s.gzip", info.TargetDir+"/"+info.FileName)
cmd := exec.Command("docker", "exec", "-i", r.ContainerName,
"sh", "-c",
fmt.Sprintf("PGPASSWORD=%s pg_dump -F c -U %s -d %s", r.Password, r.Username, info.Name),
)
var stderr bytes.Buffer
cmd.Stderr = &stderr
gzipCmd := exec.Command("gzip", "-cf")
gzipCmd.Stdin, _ = cmd.StdoutPipe()
gzipCmd.Stdout = outfile
_ = gzipCmd.Start()
if err := cmd.Run(); err != nil {
return fmt.Errorf("handle backup database failed, err: %v", stderr.String())
cmdMgr := cmd.NewCommandMgr(cmd.WithOutputFile(path.Join(info.TargetDir, info.FileName)))
if _, err := cmdMgr.RunPipe(
cmd.PipeCommand{Name: "docker", Args: []string{"exec", "-i", "-e", "PGPASSWORD=" + r.Password, r.ContainerName, "pg_dump", "-F", "c", "-U", r.Username, "-d", info.Name}},
cmd.PipeCommand{Name: "gzip", Args: []string{"-cf"}},
); err != nil {
return fmt.Errorf("handle backup database failed, err: %v", err)
}
_ = gzipCmd.Wait()
return nil
}
@@ -165,8 +150,8 @@ func (r *Local) Recover(info RecoverInfo) error {
fi, _ := os.Open(info.SourceFile)
defer fi.Close()
cmd := exec.Command("docker", "exec", "-i", r.ContainerName, "sh", "-c",
fmt.Sprintf("PGPASSWORD=%s pg_restore -F c -c --if-exists --no-owner -U %s -d %s", r.Password, r.Username, info.Name),
cmd := exec.Command("docker", "exec", "-i", "-e", "PGPASSWORD="+r.Password, r.ContainerName,
"pg_restore", "-F", "c", "-c", "--if-exists", "--no-owner", "-U", r.Username, "-d", info.Name,
)
if strings.HasSuffix(info.SourceFile, ".gz") {
gzipFile, err := os.Open(info.SourceFile)

View File

@@ -2,6 +2,7 @@ package client
import (
"bufio"
"bytes"
"context"
"database/sql"
"fmt"
@@ -24,6 +25,37 @@ import (
_ "github.com/jackc/pgx/v5/stdlib"
)
const maxPgDumpStderrCapture = 64 * 1024
var pgDumpMagic = []byte("PGDMP")
type limitedBuffer struct {
buf bytes.Buffer
limit int
truncated int
}
func (b *limitedBuffer) Write(p []byte) (int, error) {
if b.limit > 0 && b.buf.Len() >= b.limit {
b.truncated += len(p)
return len(p), nil
}
if b.limit > 0 && b.buf.Len()+len(p) > b.limit {
keep := b.limit - b.buf.Len()
_, _ = b.buf.Write(p[:keep])
b.truncated += len(p) - keep
return len(p), nil
}
return b.buf.Write(p)
}
func (b *limitedBuffer) String() string {
if b.truncated == 0 {
return b.buf.String()
}
return fmt.Sprintf("%s\n... truncated %d bytes ...", b.buf.String(), b.truncated)
}
type Remote struct {
Client *sql.DB
From string
@@ -160,21 +192,51 @@ func (r *Remote) Backup(info BackupInfo) error {
}
}
fileNameItem := info.TargetDir + "/" + strings.TrimSuffix(info.FileName, ".gz")
backupCommand := exec.Command("bash", "-c",
fmt.Sprintf("docker run --rm --net=host -i %s /bin/bash -c 'PGPASSWORD='\\''%s'\\'' pg_dump -h %s -p %d --no-owner -Fc -U %s %s' > %s",
imageTag, r.Password, r.Address, r.Port, r.User, info.Name, fileNameItem))
_ = backupCommand.Run()
b := make([]byte, 5)
n := []byte{80, 71, 68, 77, 80}
backupFile, err := os.OpenFile(fileNameItem, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, os.ModePerm)
if err != nil {
return err
}
backupFileClosed := false
defer func() {
if !backupFileClosed {
_ = backupFile.Close()
}
}()
backupCommand := exec.Command(
"docker",
"run", "--rm", "--net=host", "-i",
"-e", "PGPASSWORD="+r.Password,
imageTag,
"pg_dump",
"-h", r.Address,
"-p", fmt.Sprintf("%d", r.Port),
"--no-owner",
"-Fc",
"-U", r.User,
info.Name,
)
backupCommand.Stdout = backupFile
stderr := &limitedBuffer{limit: maxPgDumpStderrCapture}
backupCommand.Stderr = stderr
if err := backupCommand.Run(); err != nil {
return fmt.Errorf("backup failed, stderr: %s, err: %v", strings.TrimSpace(stderr.String()), err)
}
if err := backupFile.Close(); err != nil {
return fmt.Errorf("close backup file failed, err: %v", err)
}
backupFileClosed = true
b := make([]byte, len(pgDumpMagic))
handle, err := os.OpenFile(fileNameItem, os.O_RDONLY, os.ModePerm)
if err != nil {
return fmt.Errorf("backup file not found,err:%v", err)
}
defer handle.Close()
_, _ = handle.Read(b)
if string(b) != string(n) {
errBytes, _ := os.ReadFile(fileNameItem)
return fmt.Errorf("backup failed, err: %s", string(errBytes))
if _, err := io.ReadFull(handle, b); err != nil {
return fmt.Errorf("read backup header failed, stderr: %s, err: %v", strings.TrimSpace(stderr.String()), err)
}
if !bytes.Equal(b, pgDumpMagic) {
return fmt.Errorf("backup failed, invalid pg dump header: %q, stderr: %s", string(b), strings.TrimSpace(stderr.String()))
}
gzipCmd := exec.Command("gzip", fileNameItem)
@@ -207,9 +269,32 @@ func (r *Remote) Recover(info RecoverInfo) error {
_, _ = gzipCmd.CombinedOutput()
}()
}
recoverCommand := exec.Command("bash", "-c",
fmt.Sprintf("docker run --rm --net=host -i %s /bin/bash -c 'PGPASSWORD='\\''%s'\\'' pg_restore -h %s -p %d --verbose --clean --no-privileges --no-owner -Fc -c --if-exists --no-owner -U %s -d %s --role=%s' < %s",
imageTag, r.Password, r.Address, r.Port, r.User, info.Name, info.Username, fileName))
restoreFile, err := os.Open(fileName)
if err != nil {
return err
}
defer restoreFile.Close()
recoverCommand := exec.Command(
"docker",
"run", "--rm", "--net=host", "-i",
"-e", "PGPASSWORD="+r.Password,
imageTag,
"pg_restore",
"-h", r.Address,
"-p", fmt.Sprintf("%d", r.Port),
"--verbose",
"--clean",
"--no-privileges",
"--no-owner",
"-Fc",
"-c",
"--if-exists",
"--no-owner",
"-U", r.User,
"-d", info.Name,
"--role="+info.Username,
)
recoverCommand.Stdin = restoreFile
pipe, _ := recoverCommand.StdoutPipe()
stderrPipe, _ := recoverCommand.StderrPipe()
defer pipe.Close()
@@ -230,6 +315,11 @@ func (r *Remote) Recover(info RecoverInfo) error {
}
global.LOG.Infof("[PostgreSQL] DB:[%s] Restoring: %s", info.Name, readString)
}
if err := recoverCommand.Wait(); err != nil {
all, _ := io.ReadAll(stderrPipe)
global.LOG.Errorf("[PostgreSQL] DB:[%s] Recover Error: %s", info.Name, string(all))
return err
}
return nil
}