feat(cluster): 支持远程源服务器集中备份

为 Master 本地磁盘增加可选流式中转,远程 Agent 可直接把产物写入中央存储并通过反向通道恢复。

持久化并校验传输模式,兼容既有 Agent 本机磁盘目标,补齐鉴权、完整性、配额、访问保护、前端配置及双向链路测试。

Closes #101
This commit is contained in:
Awuqing
2026-08-07 05:54:34 +08:00
parent 05dd1baa61
commit cc50637b4b
28 changed files with 1012 additions and 180 deletions
+62 -4
View File
@@ -5,6 +5,7 @@ import (
"context"
"crypto/tls"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
@@ -125,10 +126,11 @@ type TaskSpec struct {
// StorageTargetConfig 与 service.AgentStorageTargetConfig 对齐
type StorageTargetConfig struct {
ID uint `json:"id"`
Type string `json:"type"`
Name string `json:"name"`
Config json.RawMessage `json:"config"`
ID uint `json:"id"`
Type string `json:"type"`
Name string `json:"name"`
Config json.RawMessage `json:"config"`
TransferMode string `json:"transferMode"`
}
// GetTaskSpec 拉取任务规格
@@ -149,6 +151,7 @@ type RecordUpdate struct {
Checksum string `json:"checksum,omitempty"`
StoragePath string `json:"storagePath,omitempty"`
StorageTargetID uint `json:"storageTargetId,omitempty"`
StorageTransferMode string `json:"storageTransferMode,omitempty"`
StorageUploadResults []StorageResultItem `json:"storageUploadResults,omitempty"`
ErrorMessage string `json:"errorMessage,omitempty"`
LogAppend string `json:"logAppend,omitempty"`
@@ -160,6 +163,7 @@ type StorageResultItem struct {
Status string `json:"status"`
StoragePath string `json:"storagePath,omitempty"`
FileSize int64 `json:"fileSize,omitempty"`
TransferMode string `json:"transferMode,omitempty"`
Error string `json:"error,omitempty"`
}
@@ -169,6 +173,39 @@ func (c *MasterClient) UpdateRecord(ctx context.Context, recordID uint, update R
return c.do(ctx, http.MethodPost, path, update, nil)
}
// UploadArtifact streams an artifact through the Master for storage targets
// that are not directly reachable from the Agent.
func (c *MasterClient) UploadArtifact(ctx context.Context, recordID, targetID uint, objectKey string, size int64, checksum string, reader io.Reader) error {
path := fmt.Sprintf("/api/agent/records/%d/artifacts/%d", recordID, targetID)
req, err := http.NewRequestWithContext(ctx, http.MethodPut, c.baseURL+path, reader)
if err != nil {
return err
}
// The executor owns and closes the artifact file. Prevent net/http from
// closing that underlying reader when it finishes the request body.
req.Body = io.NopCloser(reader)
req.ContentLength = size
req.Header.Set("Content-Type", "application/octet-stream")
req.Header.Set("X-Agent-Token", c.token)
req.Header.Set("X-BackupX-Object-Key", objectKey)
req.Header.Set("X-BackupX-SHA256", checksum)
client := *c.httpClient
client.Timeout = 0
resp, err := client.Do(req)
if err != nil {
return fmt.Errorf("relay artifact to Master: %w", err)
}
data, readErr := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
closeErr := resp.Body.Close()
if readErr != nil || closeErr != nil {
return errors.Join(readErr, closeErr)
}
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("relay artifact to Master: http %d: %s", resp.StatusCode, string(data))
}
return nil
}
// RestoreSpec 与 service.AgentRestoreSpec 对齐
type RestoreSpec struct {
RestoreRecordID uint `json:"restoreRecordId"`
@@ -210,6 +247,27 @@ func (c *MasterClient) GetRestoreSpec(ctx context.Context, restoreRecordID uint)
return &spec, nil
}
func (c *MasterClient) DownloadRestoreArtifact(ctx context.Context, restoreRecordID uint) (io.ReadCloser, error) {
path := fmt.Sprintf("/api/agent/restores/%d/artifact", restoreRecordID)
req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.baseURL+path, nil)
if err != nil {
return nil, err
}
req.Header.Set("X-Agent-Token", c.token)
client := *c.httpClient
client.Timeout = 0
resp, err := client.Do(req)
if err != nil {
return nil, fmt.Errorf("download relayed artifact from Master: %w", err)
}
if resp.StatusCode >= 200 && resp.StatusCode < 300 {
return resp.Body, nil
}
data, readErr := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
closeErr := resp.Body.Close()
return nil, errors.Join(fmt.Errorf("download relayed artifact from Master: http %d: %s", resp.StatusCode, string(data)), readErr, closeErr)
}
// UpdateRestore 上报恢复记录的状态/日志
func (c *MasterClient) UpdateRestore(ctx context.Context, restoreRecordID uint, update RestoreUpdate) error {
path := fmt.Sprintf("/api/agent/restores/%d", restoreRecordID)
+52 -37
View File
@@ -5,6 +5,7 @@ import (
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"os"
@@ -137,13 +138,15 @@ func (e *Executor) ExecuteRunTask(ctx context.Context, taskID, recordID uint) er
}
uploadResults := make([]StorageResultItem, 0, len(spec.StorageTargets))
selectedStorageTargetID := uint(0)
selectedStorageTransferMode := ""
var uploadErrors []string
for _, target := range spec.StorageTargets {
if err := e.uploadToTarget(ctx, recordID, target, finalPath, storagePath, fileSize, spec.TaskID); err != nil {
if err := e.uploadToTarget(ctx, recordID, target, finalPath, storagePath, fileSize, checksum, spec.TaskID); err != nil {
uploadResults = append(uploadResults, StorageResultItem{
StorageTargetID: target.ID,
StorageTargetName: target.Name,
Status: "failed",
TransferMode: target.TransferMode,
Error: err.Error(),
})
uploadErrors = append(uploadErrors, fmt.Sprintf("%s: %v", target.Name, err))
@@ -152,6 +155,7 @@ func (e *Executor) ExecuteRunTask(ctx context.Context, taskID, recordID uint) er
}
if selectedStorageTargetID == 0 {
selectedStorageTargetID = target.ID
selectedStorageTransferMode = target.TransferMode
}
uploadResults = append(uploadResults, StorageResultItem{
StorageTargetID: target.ID,
@@ -159,6 +163,7 @@ func (e *Executor) ExecuteRunTask(ctx context.Context, taskID, recordID uint) er
Status: "success",
StoragePath: storagePath,
FileSize: fileSize,
TransferMode: target.TransferMode,
})
e.appendLog(ctx, recordID, fmt.Sprintf("[agent] 已上传到存储目标 %s\n", target.Name))
}
@@ -179,34 +184,40 @@ func (e *Executor) ExecuteRunTask(ctx context.Context, taskID, recordID uint) er
Checksum: checksum,
StoragePath: storagePath,
StorageTargetID: selectedStorageTargetID,
StorageTransferMode: selectedStorageTransferMode,
StorageUploadResults: uploadResults,
LogAppend: fmt.Sprintf("[agent] 任务完成,总计 %d 字节\n", fileSize),
})
}
// uploadToTarget 上传单个目标。为保持简化不做上传级重试(rclone 本身已有 low-level 重试)。
func (e *Executor) uploadToTarget(ctx context.Context, recordID uint, target StorageTargetConfig, filePath, objectKey string, fileSize int64, taskID uint) error {
var rawConfig map[string]any
if len(target.Config) > 0 {
// DecodeRawConfig 通过 json 解析
if err := jsonUnmarshalMap(target.Config, &rawConfig); err != nil {
return fmt.Errorf("parse storage config: %w", err)
}
}
provider, err := e.storageRegistry.Create(ctx, target.Type, rawConfig)
if err != nil {
return fmt.Errorf("create provider: %w", err)
}
func (e *Executor) uploadToTarget(ctx context.Context, recordID uint, target StorageTargetConfig, filePath, objectKey string, fileSize int64, checksum string, taskID uint) error {
f, err := os.Open(filePath)
if err != nil {
return fmt.Errorf("open artifact: %w", err)
}
defer f.Close()
if target.TransferMode == storage.TransferModeMasterRelay {
uploadErr := e.client.UploadArtifact(ctx, recordID, target.ID, objectKey, fileSize, checksum, f)
return errors.Join(uploadErr, f.Close())
}
var rawConfig map[string]any
if len(target.Config) > 0 {
// DecodeRawConfig 通过 json 解析
if err := jsonUnmarshalMap(target.Config, &rawConfig); err != nil {
return errors.Join(fmt.Errorf("parse storage config: %w", err), f.Close())
}
}
provider, err := e.storageRegistry.Create(ctx, target.Type, rawConfig)
if err != nil {
closeErr := f.Close()
return errors.Join(fmt.Errorf("create provider: %w", err), closeErr)
}
meta := map[string]string{
"taskId": fmt.Sprintf("%d", taskID),
"recordId": fmt.Sprintf("%d", recordID),
}
return provider.Upload(ctx, objectKey, f, fileSize, meta)
uploadErr := provider.Upload(ctx, objectKey, f, fileSize, meta)
return errors.Join(uploadErr, f.Close())
}
// appendLog 追加日志到 Master 记录(尽力而为,失败不中断主流程)
@@ -328,7 +339,7 @@ func (e *Executor) DeleteStorageObject(ctx context.Context, targetType string, t
// ExecuteRestore 处理 restore_record 命令:拉规格 → 下载 → 解压 → 执行 runner.Restore → 上报结果。
//
// 与 ExecuteRunTask 对称,但方向相反:
// - 下载:通过 spec.Storage 创建 provider → Download(spec.StoragePath)
// - 下载:直连共享存储,或通过 Master 中转其本地磁盘对象
// - 解密:当前 Agent 不支持加密恢复(密钥未下发),spec.Encrypt=true 会直接失败
// - 执行:backup.Registry.Runner(spec.Type).Restore
// - 上报:通过 UpdateRestorestatus/logAppend
@@ -357,28 +368,31 @@ func (e *Executor) ExecuteRestore(ctx context.Context, restoreRecordID uint) err
}
defer os.RemoveAll(tmpDir)
// 1) 创建 storage provider
var rawConfig map[string]any
if len(spec.Storage.Config) > 0 {
if err := jsonUnmarshalMap(spec.Storage.Config, &rawConfig); err != nil {
e.reportRestoreFailure(ctx, restoreRecordID, fmt.Sprintf("解析存储配置失败: %v", err))
return err
}
}
provider, err := e.storageRegistry.Create(ctx, spec.Storage.Type, rawConfig)
if err != nil {
e.reportRestoreFailure(ctx, restoreRecordID, fmt.Sprintf("创建存储客户端失败: %v", err))
return err
}
// 2) 下载
// 1) 下载
fileName := spec.FileName
if strings.TrimSpace(fileName) == "" {
fileName = filepath.Base(spec.StoragePath)
}
artifactPath := filepath.Join(tmpDir, filepath.Base(fileName))
e.appendRestoreLog(ctx, restoreRecordID, fmt.Sprintf("[agent] 下载备份文件 %s\n", spec.StoragePath))
reader, err := provider.Download(ctx, spec.StoragePath)
var reader io.ReadCloser
if spec.Storage.TransferMode == storage.TransferModeMasterRelay {
reader, err = e.client.DownloadRestoreArtifact(ctx, restoreRecordID)
} else {
var rawConfig map[string]any
if len(spec.Storage.Config) > 0 {
if err := jsonUnmarshalMap(spec.Storage.Config, &rawConfig); err != nil {
e.reportRestoreFailure(ctx, restoreRecordID, fmt.Sprintf("解析存储配置失败: %v", err))
return err
}
}
provider, providerErr := e.storageRegistry.Create(ctx, spec.Storage.Type, rawConfig)
if providerErr != nil {
e.reportRestoreFailure(ctx, restoreRecordID, fmt.Sprintf("创建存储客户端失败: %v", providerErr))
return providerErr
}
reader, err = provider.Download(ctx, spec.StoragePath)
}
if err != nil {
e.reportRestoreFailure(ctx, restoreRecordID, fmt.Sprintf("下载备份失败: %v", err))
return err
@@ -489,8 +503,10 @@ func buildRestoreBackupTaskSpec(spec *RestoreSpec, startedAt time.Time, tempDir
}
// writeReaderToLocal 把 reader 写到本地文件(Agent 侧工具函数)。
func writeReaderToLocal(targetPath string, reader io.ReadCloser) error {
defer reader.Close()
func writeReaderToLocal(targetPath string, reader io.ReadCloser) (err error) {
defer func() {
err = errors.Join(err, reader.Close())
}()
if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil {
return err
}
@@ -498,9 +514,8 @@ func writeReaderToLocal(targetPath string, reader io.ReadCloser) error {
if err != nil {
return err
}
defer file.Close()
_, err = io.Copy(file, reader)
return err
_, copyErr := io.Copy(file, reader)
return errors.Join(copyErr, file.Close())
}
// 辅助函数
+146
View File
@@ -1,7 +1,10 @@
package agent
import (
"archive/tar"
"bytes"
"context"
"crypto/sha256"
"encoding/json"
"fmt"
"io"
@@ -109,6 +112,149 @@ func TestExecuteRunTaskRecordsPerTargetUploadResults(t *testing.T) {
}
}
func TestExecuteRunTaskRelaysMasterLocalDiskTarget(t *testing.T) {
sourceDir := t.TempDir()
if err := os.WriteFile(filepath.Join(sourceDir, "index.html"), []byte("centralize me"), 0o644); err != nil {
t.Fatalf("WriteFile returned error: %v", err)
}
var relayed []byte
var finalUpdate RecordUpdate
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch {
case r.Method == http.MethodGet && r.URL.Path == "/api/agent/tasks/1":
writeAgentEnvelope(t, w, TaskSpec{
TaskID: 1,
Name: "remote-source",
Type: "file",
SourcePath: sourceDir,
Compression: "gzip",
StorageTargets: []StorageTargetConfig{{
ID: 11, Name: "master-disk", Type: storage.TypeLocalDisk, TransferMode: storage.TransferModeMasterRelay,
}},
})
case r.Method == http.MethodPut && r.URL.Path == "/api/agent/records/99/artifacts/11":
body, err := io.ReadAll(r.Body)
if err != nil {
t.Fatalf("ReadAll relayed body: %v", err)
}
digest := sha256.Sum256(body)
if got := r.Header.Get("X-BackupX-SHA256"); got != fmt.Sprintf("%x", digest[:]) {
t.Fatalf("relay checksum header = %q", got)
}
if r.Header.Get("X-BackupX-Object-Key") == "" || r.ContentLength != int64(len(body)) {
t.Fatalf("invalid relay metadata: key=%q length=%d body=%d", r.Header.Get("X-BackupX-Object-Key"), r.ContentLength, len(body))
}
relayed = append([]byte(nil), body...)
writeAgentEnvelope(t, w, map[string]string{"status": "ok"})
case r.Method == http.MethodPost && r.URL.Path == "/api/agent/records/99":
var update RecordUpdate
if err := json.NewDecoder(r.Body).Decode(&update); err != nil {
t.Fatalf("Decode update returned error: %v", err)
}
if update.Status != "" {
finalUpdate = update
}
writeAgentEnvelope(t, w, map[string]string{"status": "ok"})
default:
http.NotFound(w, r)
}
}))
defer server.Close()
executor := NewExecutor(NewMasterClient(server.URL, "token", false), filepath.Join(t.TempDir(), "tmp"))
if err := executor.ExecuteRunTask(context.Background(), 1, 99); err != nil {
t.Fatalf("ExecuteRunTask returned error: %v", err)
}
if len(relayed) == 0 {
t.Fatal("expected artifact bytes to be streamed through Master")
}
if finalUpdate.Status != "success" || finalUpdate.StorageTransferMode != storage.TransferModeMasterRelay {
t.Fatalf("unexpected final relay update: %#v", finalUpdate)
}
if len(finalUpdate.StorageUploadResults) != 1 || finalUpdate.StorageUploadResults[0].TransferMode != storage.TransferModeMasterRelay {
t.Fatalf("unexpected relay target result: %#v", finalUpdate.StorageUploadResults)
}
}
func TestExecuteRestoreDownloadsMasterRelayedArtifact(t *testing.T) {
var archive bytes.Buffer
tarWriter := tar.NewWriter(&archive)
content := []byte("restored through Master")
header := &tar.Header{Name: "site/index.html", Mode: 0o644, Size: int64(len(content)), Typeflag: tar.TypeReg}
if err := tarWriter.WriteHeader(header); err != nil {
t.Fatalf("WriteHeader returned error: %v", err)
}
if _, err := tarWriter.Write(content); err != nil {
t.Fatalf("Write returned error: %v", err)
}
if err := tarWriter.Close(); err != nil {
t.Fatalf("Close returned error: %v", err)
}
artifact := append([]byte(nil), archive.Bytes()...)
digest := sha256.Sum256(artifact)
restoreRoot := t.TempDir()
restoreSource := filepath.Join(restoreRoot, "site")
artifactRequests := 0
var finalUpdate RestoreUpdate
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch {
case r.Method == http.MethodGet && r.URL.Path == "/api/agent/restores/77/spec":
writeAgentEnvelope(t, w, RestoreSpec{
RestoreRecordID: 77,
BackupRecordID: 99,
TaskID: 1,
TaskName: "remote-source",
Type: "file",
SourcePath: restoreSource,
Storage: StorageTargetConfig{
ID: 11, Name: "master-disk", Type: storage.TypeLocalDisk, TransferMode: storage.TransferModeMasterRelay,
},
StoragePath: "BackupX/file/site.tar",
FileName: "site.tar",
Checksum: fmt.Sprintf("%x", digest[:]),
})
case r.Method == http.MethodGet && r.URL.Path == "/api/agent/restores/77/artifact":
artifactRequests++
if r.Header.Get("X-Agent-Token") != "token" {
t.Fatalf("missing Agent token on relay download")
}
w.Header().Set("Content-Length", fmt.Sprintf("%d", len(artifact)))
_, _ = w.Write(artifact)
case r.Method == http.MethodPost && r.URL.Path == "/api/agent/restores/77":
var update RestoreUpdate
if err := json.NewDecoder(r.Body).Decode(&update); err != nil {
t.Fatalf("Decode update returned error: %v", err)
}
if update.Status != "" {
finalUpdate = update
}
writeAgentEnvelope(t, w, map[string]string{"status": "ok"})
default:
http.NotFound(w, r)
}
}))
defer server.Close()
executor := NewExecutor(NewMasterClient(server.URL, "token", false), filepath.Join(t.TempDir(), "tmp"))
if err := executor.ExecuteRestore(context.Background(), 77); err != nil {
t.Fatalf("ExecuteRestore returned error: %v", err)
}
if artifactRequests != 1 {
t.Fatalf("expected one relay artifact request, got %d", artifactRequests)
}
restored, err := os.ReadFile(filepath.Join(restoreRoot, "site", "index.html"))
if err != nil {
t.Fatalf("ReadFile returned error: %v", err)
}
if string(restored) != string(content) {
t.Fatalf("restored content = %q, want %q", restored, content)
}
if finalUpdate.Status != "success" {
t.Fatalf("unexpected final restore update: %#v", finalUpdate)
}
}
func TestExecuteRunTaskReportsPerTargetUploadResultsWhenAllTargetsFail(t *testing.T) {
sourceDir := t.TempDir()
if err := os.WriteFile(filepath.Join(sourceDir, "index.html"), []byte("hello"), 0o644); err != nil {
+1 -1
View File
@@ -135,7 +135,7 @@ func New(ctx context.Context, cfg config.Config, version string) (*Application,
// Agent 协议服务:命令队列 + 任务下发 + 记录上报
agentCmdRepo := repository.NewAgentCommandRepository(db)
nodeService.SetAgentCommandRepository(agentCmdRepo)
agentService := service.NewAgentService(nodeRepo, backupTaskRepo, backupRecordRepo, storageTargetRepo, agentCmdRepo, configCipher)
agentService := service.NewAgentService(nodeRepo, backupTaskRepo, backupRecordRepo, storageTargetRepo, agentCmdRepo, configCipher, storageRegistry)
agentService.SetRestoreRepository(restoreRecordRepo)
agentService.StartCommandTimeoutMonitor(ctx, 30*time.Second, 10*time.Minute)
+66
View File
@@ -156,6 +156,44 @@ func (h *AgentHandler) UpdateRecord(c *gin.Context) {
response.Success(c, gin.H{"status": "ok"})
}
// UploadArtifact streams a remote source artifact into storage mounted only on
// the Master. The request body is never buffered as a whole in memory or disk.
func (h *AgentHandler) UploadArtifact(c *gin.Context) {
node, err := h.agentService.AuthenticatedNode(c.Request.Context(), extractToken(c))
if err != nil {
response.Error(c, err)
return
}
recordID, err := strconv.ParseUint(c.Param("id"), 10, 32)
if err != nil {
response.Error(c, err)
return
}
targetID, err := strconv.ParseUint(c.Param("targetId"), 10, 32)
if err != nil {
response.Error(c, err)
return
}
if c.Request.ContentLength < 0 {
c.JSON(stdhttp.StatusLengthRequired, gin.H{"code": "CONTENT_LENGTH_REQUIRED", "message": "artifact content length is required"})
return
}
if err := h.agentService.UploadArtifact(
c.Request.Context(),
node,
uint(recordID),
uint(targetID),
c.GetHeader("X-BackupX-Object-Key"),
c.Request.ContentLength,
c.GetHeader("X-BackupX-SHA256"),
c.Request.Body,
); err != nil {
response.Error(c, err)
return
}
response.Success(c, gin.H{"status": "ok"})
}
// GetRestoreSpec Agent 拉取恢复规格。
func (h *AgentHandler) GetRestoreSpec(c *gin.Context) {
if h.restoreService == nil {
@@ -208,6 +246,34 @@ func (h *AgentHandler) UpdateRestore(c *gin.Context) {
response.Success(c, gin.H{"status": "ok"})
}
// DownloadRestoreArtifact streams a Master-local backup back to its source
// Agent for restore without exposing the local storage configuration.
func (h *AgentHandler) DownloadRestoreArtifact(c *gin.Context) {
if h.restoreService == nil {
c.JSON(stdhttp.StatusServiceUnavailable, gin.H{"code": "RESTORE_SERVICE_DISABLED", "message": "restore service is not enabled"})
return
}
node, err := h.agentService.AuthenticatedNode(c.Request.Context(), extractToken(c))
if err != nil {
response.Error(c, err)
return
}
restoreID, err := strconv.ParseUint(c.Param("id"), 10, 32)
if err != nil {
response.Error(c, err)
return
}
artifact, err := h.restoreService.DownloadAgentArtifact(c.Request.Context(), node, uint(restoreID))
if err != nil {
response.Error(c, err)
return
}
c.DataFromReader(stdhttp.StatusOK, artifact.Size, "application/octet-stream", artifact.Reader, nil)
if err := artifact.Reader.Close(); err != nil {
_ = c.Error(err)
}
}
// Self 返回当前 Agent token 所属节点的状态,供安装脚本末尾探活。
func (h *AgentHandler) Self(c *gin.Context) {
node, err := h.agentService.AuthenticatedNode(c.Request.Context(), extractToken(c))
+2
View File
@@ -322,7 +322,9 @@ func NewRouter(deps RouterDependencies) *gin.Engine {
agent.POST("/commands/:id/result", agentHandler.SubmitCommandResult)
agent.GET("/tasks/:id", agentHandler.GetTaskSpec)
agent.POST("/records/:id", agentHandler.UpdateRecord)
agent.PUT("/records/:id/artifacts/:targetId", agentHandler.UploadArtifact)
agent.GET("/restores/:id/spec", agentHandler.GetRestoreSpec)
agent.GET("/restores/:id/artifact", agentHandler.DownloadRestoreArtifact)
agent.POST("/restores/:id", agentHandler.UpdateRestore)
// Agent v1(安装脚本探活用),仅 Self 端点
+10 -8
View File
@@ -22,14 +22,16 @@ type BackupRecord struct {
Task BackupTask `json:"task,omitempty"`
StorageTargetID uint `gorm:"column:storage_target_id;index;not null" json:"storageTargetId"`
StorageTarget StorageTarget `json:"storageTarget,omitempty"`
// NodeID 执行该次备份的节点(0 = 本机 Master)。用于集群中识别 local_disk 类型
// 存储的归属节点,避免 Master 端试图跨节点访问远程 Agent 的本地存储
NodeID uint `gorm:"column:node_id;index;default:0" json:"nodeId"`
Status string `gorm:"size:20;index;not null" json:"status"`
FileName string `gorm:"column:file_name;size:255" json:"fileName"`
FileSize int64 `gorm:"column:file_size;not null;default:0" json:"fileSize"`
Checksum string `gorm:"column:checksum;size:64" json:"checksum"`
StoragePath string `gorm:"column:storage_path;size:500" json:"storagePath"`
// NodeID 执行该次备份的节点(0 = 本机 Master)。StorageTransferMode 进一步
// 区分远程 Agent 直写与 Master 中转,避免在错误节点访问 local_disk
NodeID uint `gorm:"column:node_id;index;default:0" json:"nodeId"`
Status string `gorm:"size:20;index;not null" json:"status"`
FileName string `gorm:"column:file_name;size:255" json:"fileName"`
FileSize int64 `gorm:"column:file_size;not null;default:0" json:"fileSize"`
Checksum string `gorm:"column:checksum;size:64" json:"checksum"`
StoragePath string `gorm:"column:storage_path;size:500" json:"storagePath"`
// 空值表示旧版 Agent 直写;direct / master_relay 记录新协议的实际数据路径。
StorageTransferMode string `gorm:"column:storage_transfer_mode;size:20" json:"storageTransferMode,omitempty"`
StorageUploadResults string `gorm:"column:storage_upload_results;type:text" json:"-"`
DurationSeconds int `gorm:"column:duration_seconds;not null;default:0" json:"durationSeconds"`
// Locked 保留锁定(法律保留):为 true 时该备份不参与保留期/数量自动清理,
+180 -8
View File
@@ -2,15 +2,19 @@ package service
import (
"context"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"path"
"strings"
"time"
"backupx/server/internal/apperror"
"backupx/server/internal/model"
"backupx/server/internal/repository"
"backupx/server/internal/storage"
"backupx/server/internal/storage/codec"
)
@@ -23,6 +27,7 @@ type AgentService struct {
storageRepo repository.StorageTargetRepository
cmdRepo repository.AgentCommandRepository
restoreRepo repository.RestoreRecordRepository
registry *storage.Registry
cipher *codec.ConfigCipher
}
@@ -33,6 +38,7 @@ func NewAgentService(
storageRepo repository.StorageTargetRepository,
cmdRepo repository.AgentCommandRepository,
cipher *codec.ConfigCipher,
registry *storage.Registry,
) *AgentService {
return &AgentService{
nodeRepo: nodeRepo,
@@ -40,6 +46,7 @@ func NewAgentService(
recordRepo: recordRepo,
storageRepo: storageRepo,
cmdRepo: cmdRepo,
registry: registry,
cipher: cipher,
}
}
@@ -145,10 +152,11 @@ type AgentTaskSpec struct {
// AgentStorageTargetConfig 存储目标配置(已解密)
type AgentStorageTargetConfig struct {
ID uint `json:"id"`
Type string `json:"type"`
Name string `json:"name"`
Config json.RawMessage `json:"config"`
ID uint `json:"id"`
Type string `json:"type"`
Name string `json:"name"`
Config json.RawMessage `json:"config"`
TransferMode string `json:"transferMode"`
}
// GetTaskSpec 返回 Agent 执行任务所需的完整规格。
@@ -187,11 +195,22 @@ func (s *AgentService) GetTaskSpec(ctx context.Context, node *model.Node, taskID
if err != nil {
return nil, fmt.Errorf("decrypt storage config: %w", err)
}
transferMode := storage.TransferModeDirect
if strings.EqualFold(target.Type, storage.TypeLocalDisk) {
var localConfig storage.LocalDiskConfig
if err := json.Unmarshal(configRaw, &localConfig); err != nil {
return nil, fmt.Errorf("decode local disk config: %w", err)
}
if localConfig.MasterRelay {
transferMode = storage.TransferModeMasterRelay
}
}
storageTargets = append(storageTargets, AgentStorageTargetConfig{
ID: target.ID,
Type: target.Type,
Name: target.Name,
Config: json.RawMessage(configRaw),
ID: target.ID,
Type: target.Type,
Name: target.Name,
Config: json.RawMessage(configRaw),
TransferMode: transferMode,
})
}
return &AgentTaskSpec{
@@ -214,6 +233,101 @@ func (s *AgentService) GetTaskSpec(ctx context.Context, node *model.Node, taskID
}, nil
}
// UploadArtifact receives a remote Agent artifact as a stream and writes it
// with a provider created on the Master. The first supported use is local_disk,
// whose configured path belongs to the Master rather than the source Agent.
func (s *AgentService) UploadArtifact(ctx context.Context, node *model.Node, recordID, targetID uint, objectKey string, size int64, checksum string, reader io.Reader) error {
if node == nil || reader == nil || s.registry == nil {
return apperror.BadRequest("AGENT_ARTIFACT_INVALID", "中转上传参数不完整", nil)
}
record, err := s.recordRepo.FindByID(ctx, recordID)
if err != nil {
return err
}
if record == nil {
return apperror.New(404, "BACKUP_RECORD_NOT_FOUND", "记录不存在", nil)
}
task, err := s.taskRepo.FindByID(ctx, record.TaskID)
if err != nil {
return err
}
if task == nil || !recordBelongsToNode(record, task, node.ID) {
return apperror.Unauthorized("BACKUP_RECORD_FORBIDDEN", "记录不属于当前节点", nil)
}
if isBackupRecordTerminal(record.Status) {
return apperror.BadRequest("BACKUP_RECORD_TERMINAL", "备份记录已结束,不能继续上传产物", nil)
}
allowedTarget := false
for _, configuredTargetID := range collectTargetIDs(task) {
if configuredTargetID == targetID {
allowedTarget = true
break
}
}
if !allowedTarget {
return apperror.Unauthorized("BACKUP_STORAGE_TARGET_FORBIDDEN", "存储目标不属于该任务", nil)
}
target, err := s.storageRepo.FindByID(ctx, targetID)
if err != nil {
return err
}
if target == nil || !strings.EqualFold(target.Type, storage.TypeLocalDisk) {
return apperror.BadRequest("AGENT_ARTIFACT_RELAY_UNSUPPORTED", "仅 Master 本地磁盘目标需要中转上传", nil)
}
configMap := map[string]any{}
if err := s.cipher.DecryptJSON(target.ConfigCiphertext, &configMap); err != nil {
return fmt.Errorf("decrypt storage config: %w", err)
}
masterRelay, _ := configMap["masterRelay"].(bool)
if !masterRelay {
return apperror.BadRequest("AGENT_ARTIFACT_RELAY_UNSUPPORTED", "该本地磁盘目标配置为 Agent 直接写入", nil)
}
cleanKey := path.Clean(strings.TrimSpace(objectKey))
if cleanKey == "." || path.IsAbs(cleanKey) || strings.HasPrefix(cleanKey, "../") || cleanKey != objectKey || strings.Contains(objectKey, "\\") {
return apperror.BadRequest("AGENT_ARTIFACT_INVALID_PATH", "中转上传对象路径不安全", nil)
}
checksumBytes, checksumErr := hex.DecodeString(strings.TrimSpace(checksum))
if size < 0 || checksumErr != nil || len(checksumBytes) != 32 {
return apperror.BadRequest("AGENT_ARTIFACT_INVALID", "中转上传需要有效的大小和 SHA-256", checksumErr)
}
if target.QuotaBytes > 0 {
usage, usageErr := s.recordRepo.StorageUsage(ctx)
if usageErr != nil {
return fmt.Errorf("read storage usage: %w", usageErr)
}
currentUsed := int64(0)
for _, item := range usage {
if item.StorageTargetID == targetID {
currentUsed = item.TotalSize
break
}
}
if currentUsed+size > target.QuotaBytes {
return apperror.BadRequest("BACKUP_STORAGE_QUOTA_EXCEEDED", fmt.Sprintf("超出存储目标配额(%d + %d > %d", currentUsed, size, target.QuotaBytes), nil)
}
}
provider, err := s.registry.Create(ctx, target.Type, configMap)
if err != nil {
return fmt.Errorf("create master relay provider: %w", err)
}
limited := io.LimitReader(reader, size+1)
hashed := newHashingReader(limited)
metadata := map[string]string{
"taskId": fmt.Sprintf("%d", task.ID),
"recordId": fmt.Sprintf("%d", record.ID),
"sourceNodeId": fmt.Sprintf("%d", node.ID),
"transferMode": storage.TransferModeMasterRelay,
}
if err := provider.Upload(ctx, cleanKey, hashed, size, metadata); err != nil {
return errors.Join(fmt.Errorf("relay artifact to master storage: %w", err), provider.Delete(ctx, cleanKey))
}
if hashed.n != size || !strings.EqualFold(hashed.Sum(), checksum) {
deleteErr := provider.Delete(ctx, cleanKey)
return errors.Join(fmt.Errorf("relayed artifact integrity mismatch: received %d of %d bytes", hashed.n, size), deleteErr)
}
return nil
}
func (s *AgentService) ensureTaskSpecAccess(ctx context.Context, node *model.Node, task *model.BackupTask) error {
if task.NodeID == node.ID {
return nil
@@ -236,6 +350,7 @@ type AgentRecordUpdate struct {
Checksum string `json:"checksum,omitempty"`
StoragePath string `json:"storagePath,omitempty"`
StorageTargetID uint `json:"storageTargetId,omitempty"`
StorageTransferMode string `json:"storageTransferMode,omitempty"`
StorageUploadResults []StorageUploadResultItem `json:"storageUploadResults,omitempty"`
ErrorMessage string `json:"errorMessage,omitempty"`
LogAppend string `json:"logAppend,omitempty"` // 增量日志,追加到 record.log_content
@@ -260,6 +375,60 @@ func (s *AgentService) UpdateRecord(ctx context.Context, node *model.Node, recor
if isBackupRecordTerminal(record.Status) {
return nil
}
allowedTargets := make(map[uint]struct{})
for _, targetID := range collectTargetIDs(task) {
allowedTargets[targetID] = struct{}{}
}
targetCache := make(map[uint]*model.StorageTarget)
validateTransferMode := func(targetID uint, transferMode string) error {
if _, ok := allowedTargets[targetID]; !ok {
return apperror.Unauthorized("BACKUP_STORAGE_TARGET_FORBIDDEN", "存储目标不属于该任务", nil)
}
if transferMode == "" {
return nil
}
target := targetCache[targetID]
if target == nil {
var findErr error
target, findErr = s.storageRepo.FindByID(ctx, targetID)
if findErr != nil {
return findErr
}
if target == nil {
return apperror.BadRequest("BACKUP_STORAGE_TARGET_INVALID", "存储目标不存在", nil)
}
targetCache[targetID] = target
}
expectedMode := storage.TransferModeDirect
if strings.EqualFold(target.Type, storage.TypeLocalDisk) {
var localConfig storage.LocalDiskConfig
if err := s.cipher.DecryptJSON(target.ConfigCiphertext, &localConfig); err != nil {
return fmt.Errorf("decrypt storage config: %w", err)
}
if localConfig.MasterRelay {
expectedMode = storage.TransferModeMasterRelay
}
}
if transferMode != expectedMode {
return apperror.BadRequest("AGENT_STORAGE_TRANSFER_MODE_INVALID", "Agent 上报的存储传输模式与目标配置不一致", nil)
}
return nil
}
if update.StorageTargetID > 0 {
if _, ok := allowedTargets[update.StorageTargetID]; !ok {
return apperror.Unauthorized("BACKUP_STORAGE_TARGET_FORBIDDEN", "存储目标不属于该任务", nil)
}
if err := validateTransferMode(update.StorageTargetID, update.StorageTransferMode); err != nil {
return err
}
} else if update.StorageTransferMode != "" {
return apperror.BadRequest("AGENT_STORAGE_TRANSFER_MODE_INVALID", "传输模式缺少对应的存储目标", nil)
}
for _, result := range update.StorageUploadResults {
if err := validateTransferMode(result.StorageTargetID, result.TransferMode); err != nil {
return err
}
}
if update.Status != "" {
record.Status = update.Status
}
@@ -278,6 +447,9 @@ func (s *AgentService) UpdateRecord(ctx context.Context, node *model.Node, recor
if update.StorageTargetID > 0 {
record.StorageTargetID = update.StorageTargetID
}
if update.StorageTransferMode != "" {
record.StorageTransferMode = update.StorageTransferMode
}
if len(update.StorageUploadResults) > 0 {
if resultsJSON, marshalErr := json.Marshal(update.StorageUploadResults); marshalErr == nil {
record.StorageUploadResults = string(resultsJSON)
+88 -12
View File
@@ -1,8 +1,12 @@
package service
import (
"bytes"
"context"
"crypto/sha256"
"errors"
"fmt"
"os"
"path/filepath"
"strings"
"testing"
@@ -13,7 +17,9 @@ import (
"backupx/server/internal/logger"
"backupx/server/internal/model"
"backupx/server/internal/repository"
"backupx/server/internal/storage"
"backupx/server/internal/storage/codec"
storageRclone "backupx/server/internal/storage/rclone"
"gorm.io/gorm"
)
@@ -42,7 +48,7 @@ func newAgentServicePoolTestHarness(t *testing.T) (*AgentService, *gorm.DB, repo
if err := nodeRepo.Create(context.Background(), other); err != nil {
t.Fatalf("create other node: %v", err)
}
targetConfig, err := cipher.EncryptJSON(map[string]any{"basePath": t.TempDir()})
targetConfig, err := cipher.EncryptJSON(map[string]any{"basePath": t.TempDir(), "masterRelay": true})
if err != nil {
t.Fatalf("EncryptJSON returned error: %v", err)
}
@@ -76,7 +82,8 @@ func newAgentServicePoolTestHarness(t *testing.T) (*AgentService, *gorm.DB, repo
if err := recordRepo.Create(context.Background(), record); err != nil {
t.Fatalf("create record: %v", err)
}
return NewAgentService(nodeRepo, taskRepo, recordRepo, storageRepo, cmdRepo, cipher), db, recordRepo, cmdRepo, owner, other
storageRegistry := storage.NewRegistry(storageRclone.NewLocalDiskFactory())
return NewAgentService(nodeRepo, taskRepo, recordRepo, storageRepo, cmdRepo, cipher, storageRegistry), db, recordRepo, cmdRepo, owner, other
}
func TestAgentServicePooledTaskUsesRecordNodeForSpecAndRecordUpdates(t *testing.T) {
@@ -90,19 +97,22 @@ func TestAgentServicePooledTaskUsesRecordNodeForSpecAndRecordUpdates(t *testing.
if spec.TaskID != 1 || len(spec.StorageTargets) != 1 {
t.Fatalf("unexpected spec: %#v", spec)
}
if spec.StorageTargets[0].TransferMode != storage.TransferModeMasterRelay {
t.Fatalf("expected local disk to use Master relay, got %#v", spec.StorageTargets[0])
}
if _, err := svc.GetTaskSpec(ctx, other, 1); err == nil {
t.Fatal("expected non-owner node to be forbidden from pooled task spec")
}
if err := svc.UpdateRecord(ctx, owner, 1, AgentRecordUpdate{
Status: model.BackupRecordStatusSuccess,
FileName: "backup.tar.gz",
FileSize: 123,
StoragePath: "tasks/1/backup.tar.gz",
StorageTargetID: 2,
Status: model.BackupRecordStatusSuccess,
FileName: "backup.tar.gz",
FileSize: 123,
StoragePath: "tasks/1/backup.tar.gz",
StorageTargetID: 1,
StorageTransferMode: storage.TransferModeMasterRelay,
StorageUploadResults: []StorageUploadResultItem{
{StorageTargetID: 1, StorageTargetName: "first", Status: "failed", Error: "boom"},
{StorageTargetID: 2, StorageTargetName: "second", Status: "success", StoragePath: "tasks/1/backup.tar.gz", FileSize: 123},
{StorageTargetID: 1, StorageTargetName: "local", Status: "success", StoragePath: "tasks/1/backup.tar.gz", FileSize: 123, TransferMode: storage.TransferModeMasterRelay},
},
}); err != nil {
t.Fatalf("owner UpdateRecord returned error: %v", err)
@@ -114,10 +124,13 @@ func TestAgentServicePooledTaskUsesRecordNodeForSpecAndRecordUpdates(t *testing.
if updated.Status != model.BackupRecordStatusSuccess || updated.NodeID != owner.ID {
t.Fatalf("unexpected updated record: %#v", updated)
}
if updated.StorageTargetID != 2 {
t.Fatalf("expected successful storage target id 2, got %d", updated.StorageTargetID)
if updated.StorageTargetID != 1 {
t.Fatalf("expected successful storage target id 1, got %d", updated.StorageTargetID)
}
if !strings.Contains(updated.StorageUploadResults, `"storageTargetName":"second"`) {
if updated.StorageTransferMode != storage.TransferModeMasterRelay {
t.Fatalf("expected Master relay transfer mode, got %q", updated.StorageTransferMode)
}
if !strings.Contains(updated.StorageUploadResults, `"storageTargetName":"local"`) {
t.Fatalf("expected upload results to be persisted, got %q", updated.StorageUploadResults)
}
if err := svc.UpdateRecord(ctx, other, 1, AgentRecordUpdate{LogAppend: "bad"}); err == nil {
@@ -125,6 +138,69 @@ func TestAgentServicePooledTaskUsesRecordNodeForSpecAndRecordUpdates(t *testing.
}
}
func TestAgentServiceRelaysRemoteArtifactToMasterLocalDisk(t *testing.T) {
svc, _, _, _, owner, other := newAgentServicePoolTestHarness(t)
ctx := context.Background()
payload := []byte("artifact from remote source server")
digest := sha256.Sum256(payload)
checksum := fmt.Sprintf("%x", digest[:])
objectKey := "file/2026/08/06/remote-source.tar"
if err := svc.UploadArtifact(ctx, owner, 1, 1, objectKey, int64(len(payload)), checksum, bytes.NewReader(payload)); err != nil {
t.Fatalf("UploadArtifact returned error: %v", err)
}
target, err := svc.storageRepo.FindByID(ctx, 1)
if err != nil || target == nil {
t.Fatalf("FindByID target: target=%#v err=%v", target, err)
}
config := map[string]any{}
if err := svc.cipher.DecryptJSON(target.ConfigCiphertext, &config); err != nil {
t.Fatalf("DecryptJSON target config: %v", err)
}
basePath, _ := config["basePath"].(string)
stored, err := os.ReadFile(filepath.Join(basePath, filepath.FromSlash(objectKey)))
if err != nil {
t.Fatalf("read relayed artifact: %v", err)
}
if !bytes.Equal(stored, payload) {
t.Fatalf("relayed artifact differs: got %q", stored)
}
if err := svc.UploadArtifact(ctx, other, 1, 1, "file/forbidden.tar", int64(len(payload)), checksum, bytes.NewReader(payload)); err == nil {
t.Fatal("expected a different node to be forbidden from relaying the artifact")
}
}
func TestAgentServiceKeepsExistingLocalDiskTargetsAgentLocal(t *testing.T) {
svc, _, _, _, owner, _ := newAgentServicePoolTestHarness(t)
ctx := context.Background()
target, err := svc.storageRepo.FindByID(ctx, 1)
if err != nil || target == nil {
t.Fatalf("FindByID target: target=%#v err=%v", target, err)
}
legacyConfig, err := svc.cipher.EncryptJSON(map[string]any{"basePath": t.TempDir()})
if err != nil {
t.Fatalf("EncryptJSON legacy target: %v", err)
}
target.ConfigCiphertext = legacyConfig
if err := svc.storageRepo.Update(ctx, target); err != nil {
t.Fatalf("Update legacy target: %v", err)
}
spec, err := svc.GetTaskSpec(ctx, owner, 1)
if err != nil {
t.Fatalf("GetTaskSpec returned error: %v", err)
}
if len(spec.StorageTargets) != 1 || spec.StorageTargets[0].TransferMode != storage.TransferModeDirect {
t.Fatalf("expected legacy local disk to stay Agent-local, got %#v", spec.StorageTargets)
}
payload := []byte("must not be relayed")
digest := sha256.Sum256(payload)
err = svc.UploadArtifact(ctx, owner, 1, 1, "file/legacy.tar", int64(len(payload)), fmt.Sprintf("%x", digest[:]), bytes.NewReader(payload))
if err == nil {
t.Fatal("expected relay upload to be rejected for an Agent-local target")
}
}
func TestAgentServiceUpdateRecordRefreshesTaskSummaryOnTerminalStatus(t *testing.T) {
for _, status := range []string{model.BackupRecordStatusSuccess, model.BackupRecordStatusFailed} {
t.Run(status, func(t *testing.T) {
@@ -51,6 +51,7 @@ type StorageUploadResultItem struct {
Status string `json:"status"`
StoragePath string `json:"storagePath,omitempty"`
FileSize int64 `json:"fileSize,omitempty"`
TransferMode string `json:"transferMode,omitempty"`
Error string `json:"error,omitempty"`
}
@@ -410,6 +411,9 @@ func (s *BackupExecutionService) deleteRemoteLocalDiskObject(ctx context.Context
if strings.TrimSpace(record.StoragePath) == "" || s.nodeRepo == nil {
return false, nil
}
if record.StorageTransferMode == storage.TransferModeMasterRelay {
return false, nil
}
node, err := s.nodeRepo.FindByID(ctx, record.NodeID)
if err != nil || node == nil || node.IsLocal {
return false, nil
@@ -429,6 +429,56 @@ func TestBackupExecutionServiceRestoreRecordRejectsRemoteLocalDisk(t *testing.T)
}
}
func TestBackupExecutionServiceDownloadsMasterRelayedLocalDiskRecord(t *testing.T) {
executionService, _, tasks, _, records, _, storageDir := newExecutionTestServices(t)
ctx := context.Background()
executionService.SetClusterDependencies(&nodeRepoStub{nodes: []model.Node{
{ID: 10, Name: "edge-a", Token: "edge-a-token", Status: model.NodeStatusOnline},
}}, &fakeDispatcher{})
task, err := tasks.FindByID(ctx, 1)
if err != nil {
t.Fatalf("FindByID task returned error: %v", err)
}
storagePath := "file/2026/05/09/relayed.tar"
artifactPath := filepath.Join(storageDir, filepath.FromSlash(storagePath))
if err := os.MkdirAll(filepath.Dir(artifactPath), 0o755); err != nil {
t.Fatalf("MkdirAll artifact parent returned error: %v", err)
}
content := []byte("stored on Master")
if err := os.WriteFile(artifactPath, content, 0o600); err != nil {
t.Fatalf("WriteFile artifact returned error: %v", err)
}
completedAt := time.Now().UTC()
record := &model.BackupRecord{
TaskID: task.ID,
StorageTargetID: task.StorageTargetID,
NodeID: 10,
Status: model.BackupRecordStatusSuccess,
FileName: "relayed.tar",
FileSize: int64(len(content)),
StoragePath: storagePath,
StorageTransferMode: storage.TransferModeMasterRelay,
StartedAt: completedAt.Add(-time.Second),
CompletedAt: &completedAt,
}
if err := records.Create(ctx, record); err != nil {
t.Fatalf("Create record returned error: %v", err)
}
download, err := executionService.DownloadRecord(ctx, record.ID)
if err != nil {
t.Fatalf("DownloadRecord returned error: %v", err)
}
got, readErr := io.ReadAll(download.Reader)
closeErr := download.Reader.Close()
if readErr != nil || closeErr != nil {
t.Fatalf("read relayed artifact: read=%v close=%v", readErr, closeErr)
}
if !bytes.Equal(got, content) {
t.Fatalf("downloaded content = %q, want %q", got, content)
}
}
func TestBackupExecutionServiceRecordsFirstSuccessfulStorageTarget(t *testing.T) {
executionService, _, tasks, targets, records, _, _ := newExecutionTestServices(t)
ctx := context.Background()
@@ -23,22 +23,23 @@ type BackupRecordListInput struct {
}
type BackupRecordSummary struct {
ID uint `json:"id"`
TaskID uint `json:"taskId"`
TaskName string `json:"taskName"`
StorageTargetID uint `json:"storageTargetId"`
StorageTargetName string `json:"storageTargetName"`
Status string `json:"status"`
FileName string `json:"fileName"`
FileSize int64 `json:"fileSize"`
Checksum string `json:"checksum"`
StoragePath string `json:"storagePath"`
DurationSeconds int `json:"durationSeconds"`
ErrorMessage string `json:"errorMessage"`
StartedAt time.Time `json:"startedAt"`
CompletedAt *time.Time `json:"completedAt,omitempty"`
Locked bool `json:"locked"`
BackupKind string `json:"backupKind"`
ID uint `json:"id"`
TaskID uint `json:"taskId"`
TaskName string `json:"taskName"`
StorageTargetID uint `json:"storageTargetId"`
StorageTargetName string `json:"storageTargetName"`
Status string `json:"status"`
FileName string `json:"fileName"`
FileSize int64 `json:"fileSize"`
Checksum string `json:"checksum"`
StoragePath string `json:"storagePath"`
StorageTransferMode string `json:"storageTransferMode,omitempty"`
DurationSeconds int `json:"durationSeconds"`
ErrorMessage string `json:"errorMessage"`
StartedAt time.Time `json:"startedAt"`
CompletedAt *time.Time `json:"completedAt,omitempty"`
Locked bool `json:"locked"`
BackupKind string `json:"backupKind"`
}
type BackupRecordDetail struct {
@@ -184,22 +185,23 @@ func (s *BackupRecordService) SetLock(ctx context.Context, id uint, locked bool)
func toBackupRecordSummary(item *model.BackupRecord) BackupRecordSummary {
return BackupRecordSummary{
ID: item.ID,
TaskID: item.TaskID,
TaskName: item.Task.Name,
StorageTargetID: item.StorageTargetID,
StorageTargetName: item.StorageTarget.Name,
Status: item.Status,
FileName: item.FileName,
FileSize: item.FileSize,
Checksum: item.Checksum,
StoragePath: item.StoragePath,
DurationSeconds: item.DurationSeconds,
ErrorMessage: item.ErrorMessage,
StartedAt: item.StartedAt,
CompletedAt: item.CompletedAt,
Locked: item.Locked,
BackupKind: item.BackupKind,
ID: item.ID,
TaskID: item.TaskID,
TaskName: item.Task.Name,
StorageTargetID: item.StorageTargetID,
StorageTargetName: item.StorageTarget.Name,
Status: item.Status,
FileName: item.FileName,
FileSize: item.FileSize,
Checksum: item.Checksum,
StoragePath: item.StoragePath,
StorageTransferMode: item.StorageTransferMode,
DurationSeconds: item.DurationSeconds,
ErrorMessage: item.ErrorMessage,
StartedAt: item.StartedAt,
CompletedAt: item.CompletedAt,
Locked: item.Locked,
BackupKind: item.BackupKind,
}
}
@@ -140,6 +140,11 @@ func validateCrossNodeLocalDisk(ctx context.Context, nodeRepo repository.NodeRep
if record == nil || record.NodeID == 0 || nodeRepo == nil {
return nil
}
// 中转模式的对象实际落在 Master 配置的本地磁盘,Master 可以安全访问。
// 空值和 direct 均按旧版 Agent 本地落盘处理,保持升级兼容。
if record.StorageTransferMode == storage.TransferModeMasterRelay {
return nil
}
node, err := nodeRepo.FindByID(ctx, record.NodeID)
if err != nil || node == nil || node.IsLocal {
return nil
+74 -8
View File
@@ -4,6 +4,7 @@ import (
"context"
"encoding/json"
"fmt"
"io"
"os"
"path/filepath"
"strings"
@@ -601,15 +602,22 @@ func (s *RestoreService) GetAgentRestoreSpec(ctx context.Context, node *model.No
if target == nil {
return nil, apperror.BadRequest("BACKUP_STORAGE_TARGET_INVALID", "存储目标不存在", nil)
}
configRaw, err := s.cipher.Decrypt(target.ConfigCiphertext)
if err != nil {
return nil, fmt.Errorf("decrypt storage config: %w", err)
}
// 拆开 sourcePaths
sourcePaths := []string{}
if strings.TrimSpace(task.SourcePaths) != "" {
_ = json.Unmarshal([]byte(task.SourcePaths), &sourcePaths)
}
transferMode := storage.TransferModeDirect
if backupRecord.StorageTransferMode == storage.TransferModeMasterRelay {
transferMode = storage.TransferModeMasterRelay
}
var configRaw []byte
if transferMode == storage.TransferModeDirect {
configRaw, err = s.cipher.Decrypt(target.ConfigCiphertext)
if err != nil {
return nil, fmt.Errorf("decrypt storage config: %w", err)
}
}
return &AgentRestoreSpec{
RestoreRecordID: restore.ID,
BackupRecordID: backupRecord.ID,
@@ -628,10 +636,11 @@ func (s *RestoreService) GetAgentRestoreSpec(ctx context.Context, node *model.No
Compression: task.Compression,
Encrypt: task.Encrypt,
Storage: AgentStorageTargetConfig{
ID: target.ID,
Type: target.Type,
Name: target.Name,
Config: json.RawMessage(configRaw),
ID: target.ID,
Type: target.Type,
Name: target.Name,
Config: json.RawMessage(configRaw),
TransferMode: transferMode,
},
StoragePath: backupRecord.StoragePath,
FileName: backupRecord.FileName,
@@ -639,6 +648,63 @@ func (s *RestoreService) GetAgentRestoreSpec(ctx context.Context, node *model.No
}, nil
}
type AgentArtifactDownload struct {
Reader io.ReadCloser
Size int64
}
// DownloadAgentArtifact opens a Master-local object for authenticated streaming
// back to the Agent that owns the restore record.
func (s *RestoreService) DownloadAgentArtifact(ctx context.Context, node *model.Node, restoreID uint) (*AgentArtifactDownload, error) {
if node == nil {
return nil, apperror.Unauthorized("RESTORE_RECORD_FORBIDDEN", "恢复记录不属于当前节点", nil)
}
restore, err := s.restores.FindByID(ctx, restoreID)
if err != nil {
return nil, err
}
if restore == nil {
return nil, apperror.New(404, "RESTORE_RECORD_NOT_FOUND", "恢复记录不存在", nil)
}
if restore.NodeID != node.ID {
return nil, apperror.Unauthorized("RESTORE_RECORD_FORBIDDEN", "恢复记录不属于当前节点", nil)
}
if isRestoreRecordTerminal(restore.Status) {
return nil, apperror.BadRequest("RESTORE_RECORD_TERMINAL", "恢复记录已结束,不能继续下载产物", nil)
}
record, err := s.records.FindByID(ctx, restore.BackupRecordID)
if err != nil {
return nil, err
}
if record == nil {
return nil, apperror.New(404, "BACKUP_RECORD_NOT_FOUND", "源备份记录不存在", nil)
}
target, err := s.targets.FindByID(ctx, record.StorageTargetID)
if err != nil {
return nil, err
}
if target == nil || !strings.EqualFold(target.Type, storage.TypeLocalDisk) || record.StorageTransferMode != storage.TransferModeMasterRelay {
return nil, apperror.BadRequest("AGENT_ARTIFACT_RELAY_UNSUPPORTED", "该存储目标应由 Agent 直接下载", nil)
}
configMap := map[string]any{}
if err := s.cipher.DecryptJSON(target.ConfigCiphertext, &configMap); err != nil {
return nil, fmt.Errorf("decrypt storage config: %w", err)
}
provider, err := s.storageRegistry.Create(ctx, target.Type, configMap)
if err != nil {
return nil, fmt.Errorf("create master relay provider: %w", err)
}
reader, err := provider.Download(ctx, record.StoragePath)
if err != nil {
return nil, fmt.Errorf("open master relay artifact: %w", err)
}
size := record.FileSize
if size <= 0 {
size = -1
}
return &AgentArtifactDownload{Reader: reader, Size: size}, nil
}
// UpdateAgentRestore Agent 回传状态/日志。
func (s *RestoreService) UpdateAgentRestore(ctx context.Context, node *model.Node, restoreID uint, update AgentRestoreUpdate) error {
restore, err := s.restores.FindByID(ctx, restoreID)
@@ -1,8 +1,10 @@
package service
import (
"bytes"
"context"
"encoding/json"
"io"
"os"
"path/filepath"
"strings"
@@ -427,16 +429,27 @@ func TestRestoreServiceAgentRestoreAccessUsesRestoreRecordNode(t *testing.T) {
}
startedAt := time.Now().UTC()
completedAt := startedAt.Add(time.Second)
artifact := []byte("central backup artifact")
storagePath := "file/2026/05/09/remote.tar.gz"
artifactPath := filepath.Join(h.storageDir, filepath.FromSlash(storagePath))
if err := os.MkdirAll(filepath.Dir(artifactPath), 0o755); err != nil {
t.Fatalf("MkdirAll artifact parent: %v", err)
}
if err := os.WriteFile(artifactPath, artifact, 0o600); err != nil {
t.Fatalf("WriteFile artifact: %v", err)
}
backupRecord := &model.BackupRecord{
TaskID: task.ID,
StorageTargetID: task.StorageTargetID,
NodeID: owner.ID,
Status: model.BackupRecordStatusSuccess,
FileName: "remote.tar.gz",
StoragePath: "file/2026/05/09/remote.tar.gz",
Checksum: "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef",
StartedAt: startedAt,
CompletedAt: &completedAt,
TaskID: task.ID,
StorageTargetID: task.StorageTargetID,
NodeID: owner.ID,
Status: model.BackupRecordStatusSuccess,
FileName: "remote.tar.gz",
StoragePath: storagePath,
FileSize: int64(len(artifact)),
StorageTransferMode: storage.TransferModeMasterRelay,
Checksum: "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef",
StartedAt: startedAt,
CompletedAt: &completedAt,
}
if err := h.records.Create(ctx, backupRecord); err != nil {
t.Fatalf("Create backup record: %v", err)
@@ -464,6 +477,21 @@ func TestRestoreServiceAgentRestoreAccessUsesRestoreRecordNode(t *testing.T) {
if spec.Checksum != backupRecord.Checksum {
t.Fatalf("expected spec.Checksum=%q, got %q", backupRecord.Checksum, spec.Checksum)
}
if spec.Storage.TransferMode != storage.TransferModeMasterRelay {
t.Fatalf("expected Master relay restore, got %#v", spec.Storage)
}
download, err := h.service.DownloadAgentArtifact(ctx, owner, restore.ID)
if err != nil {
t.Fatalf("DownloadAgentArtifact returned error: %v", err)
}
downloaded, readErr := io.ReadAll(download.Reader)
closeErr := download.Reader.Close()
if readErr != nil || closeErr != nil {
t.Fatalf("read relayed restore artifact: read=%v close=%v", readErr, closeErr)
}
if !bytes.Equal(downloaded, artifact) {
t.Fatalf("relayed restore artifact differs: %q", downloaded)
}
if _, err := h.service.GetAgentRestoreSpec(ctx, other, restore.ID); err == nil {
t.Fatal("expected non-owner node to be forbidden from restore spec")
}
+10 -1
View File
@@ -34,6 +34,14 @@ const (
TypeFTP = string(ProviderTypeFTP)
)
const (
// TransferModeDirect lets an Agent write to a network-accessible backend.
TransferModeDirect = "direct"
// TransferModeMasterRelay streams an artifact through the authenticated
// Agent API so a remote source can use storage mounted only on the Master.
TransferModeMasterRelay = "master_relay"
)
type ObjectInfo struct {
Key string `json:"key"`
Size int64 `json:"size"`
@@ -99,7 +107,8 @@ func ParseProviderType(value string) ProviderType {
}
type LocalDiskConfig struct {
BasePath string `json:"basePath"`
BasePath string `json:"basePath"`
MasterRelay bool `json:"masterRelay"`
}
type S3Config struct {