mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-08 06:16:37 +08:00
优化strm同步
优化strm同步
This commit is contained in:
@@ -178,6 +178,7 @@ type strmSyncPathReq struct {
|
|||||||
DeleteDir *bool `json:"delete_dir"`
|
DeleteDir *bool `json:"delete_dir"`
|
||||||
Cron string `json:"cron"`
|
Cron string `json:"cron"`
|
||||||
EnableCron *bool `json:"enable_cron"`
|
EnableCron *bool `json:"enable_cron"`
|
||||||
|
SyncMode string `json:"sync_mode"`
|
||||||
Enabled *bool `json:"enabled"`
|
Enabled *bool `json:"enabled"`
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -261,7 +262,16 @@ func deleteStrmSyncPathHandler(svc *service.Container) gin.HandlerFunc {
|
|||||||
|
|
||||||
func startStrmSyncHandler(svc *service.Container) gin.HandlerFunc {
|
func startStrmSyncHandler(svc *service.Container) gin.HandlerFunc {
|
||||||
return func(c *gin.Context) {
|
return func(c *gin.Context) {
|
||||||
if err := svc.Strm.StartSync(c.Request.Context(), c.Param("id")); err != nil {
|
mode := c.Query("mode")
|
||||||
|
if mode == "" {
|
||||||
|
var body struct {
|
||||||
|
Mode string `json:"mode"`
|
||||||
|
}
|
||||||
|
if err := c.ShouldBindJSON(&body); err == nil && body.Mode != "" {
|
||||||
|
mode = body.Mode
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if err := svc.Strm.StartSync(c.Request.Context(), c.Param("id"), mode); err != nil {
|
||||||
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
|
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -458,6 +468,7 @@ func strmSyncPathFromReq(req strmSyncPathReq) *model.StrmSyncPath {
|
|||||||
DeleteDir: boolValue(req.DeleteDir, false),
|
DeleteDir: boolValue(req.DeleteDir, false),
|
||||||
Cron: strings.TrimSpace(req.Cron),
|
Cron: strings.TrimSpace(req.Cron),
|
||||||
EnableCron: boolValue(req.EnableCron, false),
|
EnableCron: boolValue(req.EnableCron, false),
|
||||||
|
SyncMode: strings.TrimSpace(req.SyncMode),
|
||||||
Enabled: boolValue(req.Enabled, true),
|
Enabled: boolValue(req.Enabled, true),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -54,7 +54,8 @@ func AllModels() []interface{} {
|
|||||||
&StrmAccount{},
|
&StrmAccount{},
|
||||||
&StrmSyncPath{},
|
&StrmSyncPath{},
|
||||||
&StrmSyncRecord{},
|
&StrmSyncRecord{},
|
||||||
&StrmDownloadTask{},
|
&StrmDownloadTask{},
|
||||||
&StrmUploadTask{},
|
&StrmUploadTask{},
|
||||||
}
|
&StrmDirCache{},
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -48,12 +48,19 @@ type StrmSyncPath struct {
|
|||||||
DeleteDir bool `json:"delete_dir"` // 清理多余文件时删除空目录
|
DeleteDir bool `json:"delete_dir"` // 清理多余文件时删除空目录
|
||||||
Cron string `gorm:"size:128" json:"cron"` // 5 段 cron 表达式(可选)
|
Cron string `gorm:"size:128" json:"cron"` // 5 段 cron 表达式(可选)
|
||||||
EnableCron bool `json:"enable_cron"` // 是否按 Cron 定时同步
|
EnableCron bool `json:"enable_cron"` // 是否按 Cron 定时同步
|
||||||
|
SyncMode string `gorm:"size:32;default:'incremental'" json:"sync_mode"` // 默认同步模式:incremental / full
|
||||||
Enabled bool `gorm:"default:true" json:"enabled"`
|
Enabled bool `gorm:"default:true" json:"enabled"`
|
||||||
LastSyncAt *time.Time `json:"last_sync_at"`
|
LastSyncAt *time.Time `json:"last_sync_at"`
|
||||||
LastSyncStatus string `gorm:"size:16" json:"last_sync_status"` // idle/running/ok/error/canceled
|
LastSyncStatus string `gorm:"size:16" json:"last_sync_status"` // idle/running/ok/error/canceled
|
||||||
LastSyncMessage string `gorm:"size:1024" json:"last_sync_message"`
|
LastSyncMessage string `gorm:"size:1024" json:"last_sync_message"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// STRM 同步类型。
|
||||||
|
const (
|
||||||
|
StrmSyncTypeIncremental = "incremental"
|
||||||
|
StrmSyncTypeFull = "full"
|
||||||
|
)
|
||||||
|
|
||||||
// StrmSyncRecord 是一次同步执行的记录。
|
// StrmSyncRecord 是一次同步执行的记录。
|
||||||
const (
|
const (
|
||||||
StrmSyncRecordPending = "pending"
|
StrmSyncRecordPending = "pending"
|
||||||
@@ -66,6 +73,7 @@ const (
|
|||||||
type StrmSyncRecord struct {
|
type StrmSyncRecord struct {
|
||||||
Base
|
Base
|
||||||
SyncPathID string `gorm:"size:36;index" json:"sync_path_id"`
|
SyncPathID string `gorm:"size:36;index" json:"sync_path_id"`
|
||||||
|
SyncType string `gorm:"size:32;default:'incremental'" json:"sync_type"` // incremental / full
|
||||||
Status string `gorm:"size:16;index" json:"status"`
|
Status string `gorm:"size:16;index" json:"status"`
|
||||||
Total int64 `json:"total"` // 远端发现的文件总数
|
Total int64 `json:"total"` // 远端发现的文件总数
|
||||||
NewStrm int64 `json:"new_strm"` // 本次新建/更新的 strm 数
|
NewStrm int64 `json:"new_strm"` // 本次新建/更新的 strm 数
|
||||||
@@ -123,3 +131,12 @@ type StrmUploadTask struct {
|
|||||||
StartedAt *time.Time `json:"started_at"`
|
StartedAt *time.Time `json:"started_at"`
|
||||||
FinishedAt *time.Time `json:"finished_at"`
|
FinishedAt *time.Time `json:"finished_at"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// StrmDirCache 缓存远端网盘目录 ID 与相对路径映射(支持 115 增量同步秒级寻址)。
|
||||||
|
type StrmDirCache struct {
|
||||||
|
Base
|
||||||
|
SyncPathID string `gorm:"size:36;index:idx_strm_dir_cache,priority:1" json:"sync_path_id"`
|
||||||
|
DirID string `gorm:"size:128;index:idx_strm_dir_cache,priority:2" json:"dir_id"`
|
||||||
|
Path string `gorm:"size:1024" json:"path"` // 相对根目录的路径
|
||||||
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -31,6 +31,7 @@ type Container struct {
|
|||||||
StrmSyncRecord *StrmSyncRecordRepository
|
StrmSyncRecord *StrmSyncRecordRepository
|
||||||
StrmDownload *StrmDownloadTaskRepository
|
StrmDownload *StrmDownloadTaskRepository
|
||||||
StrmUpload *StrmUploadTaskRepository
|
StrmUpload *StrmUploadTaskRepository
|
||||||
|
StrmDirCache *StrmDirCacheRepository
|
||||||
}
|
}
|
||||||
|
|
||||||
// New 将每个 repository 连接到单个 *gorm.DB。
|
// New 将每个 repository 连接到单个 *gorm.DB。
|
||||||
@@ -58,5 +59,6 @@ func New(db *gorm.DB) *Container {
|
|||||||
StrmSyncRecord: &StrmSyncRecordRepository{db: db},
|
StrmSyncRecord: &StrmSyncRecordRepository{db: db},
|
||||||
StrmDownload: &StrmDownloadTaskRepository{db: db},
|
StrmDownload: &StrmDownloadTaskRepository{db: db},
|
||||||
StrmUpload: &StrmUploadTaskRepository{db: db},
|
StrmUpload: &StrmUploadTaskRepository{db: db},
|
||||||
|
StrmDirCache: &StrmDirCacheRepository{db: db},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -96,11 +96,12 @@ func (r *StrmSyncPathRepository) Update(ctx context.Context, p *model.StrmSyncPa
|
|||||||
"add_path": p.AddPath,
|
"add_path": p.AddPath,
|
||||||
"download_meta": p.DownloadMeta,
|
"download_meta": p.DownloadMeta,
|
||||||
"upload_meta": p.UploadMeta,
|
"upload_meta": p.UploadMeta,
|
||||||
"delete_dir": p.DeleteDir,
|
"delete_dir": p.DeleteDir,
|
||||||
"cron": p.Cron,
|
"cron": p.Cron,
|
||||||
"enable_cron": p.EnableCron,
|
"enable_cron": p.EnableCron,
|
||||||
"enabled": p.Enabled,
|
"sync_mode": p.SyncMode,
|
||||||
"last_sync_at": p.LastSyncAt,
|
"enabled": p.Enabled,
|
||||||
|
"last_sync_at": p.LastSyncAt,
|
||||||
"last_sync_status": p.LastSyncStatus,
|
"last_sync_status": p.LastSyncStatus,
|
||||||
"last_sync_message": p.LastSyncMessage,
|
"last_sync_message": p.LastSyncMessage,
|
||||||
"updated_at": time.Now(),
|
"updated_at": time.Now(),
|
||||||
@@ -122,6 +123,7 @@ func (r *StrmSyncRecordRepository) Create(ctx context.Context, rec *model.StrmSy
|
|||||||
|
|
||||||
func (r *StrmSyncRecordRepository) Update(ctx context.Context, rec *model.StrmSyncRecord) error {
|
func (r *StrmSyncRecordRepository) Update(ctx context.Context, rec *model.StrmSyncRecord) error {
|
||||||
return r.db.WithContext(ctx).Model(&model.StrmSyncRecord{}).Where("id = ?", rec.ID).Updates(map[string]any{
|
return r.db.WithContext(ctx).Model(&model.StrmSyncRecord{}).Where("id = ?", rec.ID).Updates(map[string]any{
|
||||||
|
"sync_type": rec.SyncType,
|
||||||
"status": rec.Status,
|
"status": rec.Status,
|
||||||
"total": rec.Total,
|
"total": rec.Total,
|
||||||
"new_strm": rec.NewStrm,
|
"new_strm": rec.NewStrm,
|
||||||
@@ -436,3 +438,39 @@ func (r *StrmUploadTaskRepository) DeleteFinishedOlderThan(ctx context.Context,
|
|||||||
[]string{model.StrmTaskDone, model.StrmTaskFailed, model.StrmTaskCanceled}, before).
|
[]string{model.StrmTaskDone, model.StrmTaskFailed, model.StrmTaskCanceled}, before).
|
||||||
Delete(&model.StrmUploadTask{}).Error
|
Delete(&model.StrmUploadTask{}).Error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ─── StrmDirCache ─────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
// StrmDirCacheRepository persists model.StrmDirCache.
|
||||||
|
type StrmDirCacheRepository struct{ db *gorm.DB }
|
||||||
|
|
||||||
|
func (r *StrmDirCacheRepository) ListBySyncPathID(ctx context.Context, syncPathID string) ([]model.StrmDirCache, error) {
|
||||||
|
var rows []model.StrmDirCache
|
||||||
|
err := r.db.WithContext(ctx).Where("sync_path_id = ?", syncPathID).Find(&rows).Error
|
||||||
|
return rows, err
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *StrmDirCacheRepository) Set(ctx context.Context, syncPathID, dirID, path string) error {
|
||||||
|
var row model.StrmDirCache
|
||||||
|
err := r.db.WithContext(ctx).Where("sync_path_id = ? AND dir_id = ?", syncPathID, dirID).First(&row).Error
|
||||||
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||||
|
row = model.StrmDirCache{
|
||||||
|
SyncPathID: syncPathID,
|
||||||
|
DirID: dirID,
|
||||||
|
Path: path,
|
||||||
|
}
|
||||||
|
return r.db.WithContext(ctx).Create(&row).Error
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
return r.db.WithContext(ctx).Model(&model.StrmDirCache{}).Where("id = ?", row.ID).Updates(map[string]any{
|
||||||
|
"path": path,
|
||||||
|
"updated_at": time.Now(),
|
||||||
|
}).Error
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *StrmDirCacheRepository) DeleteBySyncPathID(ctx context.Context, syncPathID string) error {
|
||||||
|
return r.db.WithContext(ctx).Where("sync_path_id = ?", syncPathID).Delete(&model.StrmDirCache{}).Error
|
||||||
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -37,10 +37,11 @@ var ErrUnsupported = errors.New("unsupported cloud provider")
|
|||||||
|
|
||||||
// FileEntry is one item in a cloud directory listing.
|
// FileEntry is one item in a cloud directory listing.
|
||||||
type FileEntry struct {
|
type FileEntry struct {
|
||||||
ID string `json:"id"` // provider-native file id
|
ID string `json:"id"` // provider-native file id
|
||||||
Name string `json:"name"`
|
Name string `json:"name"`
|
||||||
IsDir bool `json:"is_dir"`
|
IsDir bool `json:"is_dir"`
|
||||||
Size int64 `json:"size"`
|
Size int64 `json:"size"`
|
||||||
|
MTime int64 `json:"mtime,omitempty"`
|
||||||
// PickCode is 115-specific; other providers use ID directly.
|
// PickCode is 115-specific; other providers use ID directly.
|
||||||
PickCode string `json:"pick_code,omitempty"`
|
PickCode string `json:"pick_code,omitempty"`
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -22,6 +22,12 @@ import (
|
|||||||
"github.com/ShukeBta/MMTL/internal/service/cloud115"
|
"github.com/ShukeBta/MMTL/internal/service/cloud115"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// OpenAPI115Provider 暴露 115 开放平台驱动接口。
|
||||||
|
type OpenAPI115Provider interface {
|
||||||
|
Provider
|
||||||
|
OpenClient() *cloud115.OpenClient
|
||||||
|
}
|
||||||
|
|
||||||
// openAPI115Provider 实现 Provider 接口:List 列目录、Resolve 用 pickcode
|
// openAPI115Provider 实现 Provider 接口:List 列目录、Resolve 用 pickcode
|
||||||
// 换下载直链(302 offload,无需代理)、Ping 探测根目录。
|
// 换下载直链(302 offload,无需代理)、Ping 探测根目录。
|
||||||
type openAPI115Provider struct {
|
type openAPI115Provider struct {
|
||||||
@@ -55,15 +61,16 @@ func (p *openAPI115Provider) List(ctx context.Context, dirID string) ([]FileEntr
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
for _, f := range files {
|
for _, f := range files {
|
||||||
out = append(out, FileEntry{
|
out = append(out, FileEntry{
|
||||||
ID: f.FileId,
|
ID: f.FileId,
|
||||||
Name: f.FileName,
|
Name: f.FileName,
|
||||||
IsDir: f.Category == cloud115.TypeDir,
|
IsDir: f.Category == cloud115.TypeDir,
|
||||||
Size: f.FileSize,
|
Size: f.FileSize,
|
||||||
PickCode: f.PickCode,
|
MTime: f.Utime,
|
||||||
})
|
PickCode: f.PickCode,
|
||||||
}
|
})
|
||||||
|
}
|
||||||
if len(files) < pageSize {
|
if len(files) < pageSize {
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -90,6 +90,7 @@ type RespBase struct {
|
|||||||
Errno int `json:"errno"`
|
Errno int `json:"errno"`
|
||||||
Message string `json:"message"`
|
Message string `json:"message"`
|
||||||
Error string `json:"error"`
|
Error string `json:"error"`
|
||||||
|
Count int64 `json:"count"`
|
||||||
Data json.RawMessage `json:"data"`
|
Data json.RawMessage `json:"data"`
|
||||||
Raw json.RawMessage `json:"-"` // 原始响应体(外层附加字段用)
|
Raw json.RawMessage `json:"-"` // 原始响应体(外层附加字段用)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -89,6 +89,36 @@ func (c *OpenClient) GetFsList(ctx context.Context, cid string, offset, limit in
|
|||||||
return files, strings.Join(pathStr, "/"), nil
|
return files, strings.Join(pathStr, "/"), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// GetFsListFlat 递归扁平化列出 cid 下的所有文件(跨越所有子目录,不包含文件夹节点),并返回文件列表与该树下的总文件数。
|
||||||
|
// 类似于 QMediaSync 的 115 扁平化批量拉取机制,极大地降低多层级子目录下的 API 请求次数。
|
||||||
|
func (c *OpenClient) GetFsListFlat(ctx context.Context, cid string, offset, limit int) ([]RemoteFile, int64, error) {
|
||||||
|
if cid == "" {
|
||||||
|
cid = "0"
|
||||||
|
}
|
||||||
|
if limit <= 0 {
|
||||||
|
limit = 1150
|
||||||
|
}
|
||||||
|
params := map[string]string{
|
||||||
|
"cid": cid,
|
||||||
|
"limit": fmt.Sprint(limit),
|
||||||
|
"offset": fmt.Sprint(offset),
|
||||||
|
"cur": "0",
|
||||||
|
"show_dir": "0",
|
||||||
|
}
|
||||||
|
resp, err := c.doAuthJSON(ctx, "GET", ProAPIBase+"/open/ufile/files", params, 2)
|
||||||
|
if err != nil {
|
||||||
|
return nil, 0, err
|
||||||
|
}
|
||||||
|
if !resp.State {
|
||||||
|
return nil, 0, NewOpenAPIResponseError(resp.Code, resp.Errno, resp.Message, resp.Error, "115 接口调用失败")
|
||||||
|
}
|
||||||
|
files, err := openList[RemoteFile](resp.Data)
|
||||||
|
if err != nil {
|
||||||
|
return nil, 0, fmt.Errorf("115: 解析文件列表失败:%w", err)
|
||||||
|
}
|
||||||
|
return files, resp.Count, nil
|
||||||
|
}
|
||||||
|
|
||||||
// GetFsDetailByCid 查询文件(夹)详情。
|
// GetFsDetailByCid 查询文件(夹)详情。
|
||||||
func (c *OpenClient) GetFsDetailByCid(ctx context.Context, fileId string) (*RemoteFileDetail, error) {
|
func (c *OpenClient) GetFsDetailByCid(ctx context.Context, fileId string) (*RemoteFileDetail, error) {
|
||||||
params := map[string]string{"file_id": fileId}
|
params := map[string]string{"file_id": fileId}
|
||||||
@@ -113,6 +143,38 @@ type RemoteFileDetail struct {
|
|||||||
} `json:"paths"`
|
} `json:"paths"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// RelativePath 计算该目录相对于根同步目录(rootCID)的相对路径。
|
||||||
|
func (d *RemoteFileDetail) RelativePath(rootCID string) string {
|
||||||
|
if d == nil || len(d.Paths) == 0 {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
if rootCID == "" {
|
||||||
|
rootCID = "0"
|
||||||
|
}
|
||||||
|
rootIdx := -1
|
||||||
|
for i, p := range d.Paths {
|
||||||
|
if p.FileId == rootCID {
|
||||||
|
rootIdx = i
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
var segments []string
|
||||||
|
start := 0
|
||||||
|
if rootIdx >= 0 {
|
||||||
|
start = rootIdx + 1
|
||||||
|
} else if len(d.Paths) > 0 && (d.Paths[0].FileId == "0" || d.Paths[0].FileId == "") {
|
||||||
|
start = 1
|
||||||
|
}
|
||||||
|
for i := start; i < len(d.Paths); i++ {
|
||||||
|
name := strings.TrimSpace(d.Paths[i].Name)
|
||||||
|
if name != "" {
|
||||||
|
segments = append(segments, name)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return strings.Join(segments, "/")
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
// ─── 下载直链 ──────────────────────────────────────────────────────────────────
|
// ─── 下载直链 ──────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
type downloadURLData struct {
|
type downloadURLData struct {
|
||||||
|
|||||||
@@ -404,6 +404,9 @@ func (s *StrmService) CreateSyncPath(ctx context.Context, p *model.StrmSyncPath)
|
|||||||
if strings.TrimSpace(p.Name) == "" {
|
if strings.TrimSpace(p.Name) == "" {
|
||||||
p.Name = "同步目录 " + time.Now().Format("01-02 15:04")
|
p.Name = "同步目录 " + time.Now().Format("01-02 15:04")
|
||||||
}
|
}
|
||||||
|
if p.SyncMode == "" {
|
||||||
|
p.SyncMode = model.StrmSyncTypeIncremental
|
||||||
|
}
|
||||||
if p.EnableCron && strings.TrimSpace(p.Cron) == "" {
|
if p.EnableCron && strings.TrimSpace(p.Cron) == "" {
|
||||||
return nil, errors.New("启用定时同步需要填写 cron 表达式")
|
return nil, errors.New("启用定时同步需要填写 cron 表达式")
|
||||||
}
|
}
|
||||||
@@ -430,6 +433,12 @@ func (s *StrmService) UpdateSyncPath(ctx context.Context, id string, p *model.St
|
|||||||
p.LastSyncAt = existing.LastSyncAt
|
p.LastSyncAt = existing.LastSyncAt
|
||||||
p.LastSyncStatus = existing.LastSyncStatus
|
p.LastSyncStatus = existing.LastSyncStatus
|
||||||
p.LastSyncMessage = existing.LastSyncMessage
|
p.LastSyncMessage = existing.LastSyncMessage
|
||||||
|
if p.SyncMode == "" {
|
||||||
|
p.SyncMode = existing.SyncMode
|
||||||
|
if p.SyncMode == "" {
|
||||||
|
p.SyncMode = model.StrmSyncTypeIncremental
|
||||||
|
}
|
||||||
|
}
|
||||||
if p.EnableCron && strings.TrimSpace(p.Cron) == "" {
|
if p.EnableCron && strings.TrimSpace(p.Cron) == "" {
|
||||||
return nil, errors.New("启用定时同步需要填写 cron 表达式")
|
return nil, errors.New("启用定时同步需要填写 cron 表达式")
|
||||||
}
|
}
|
||||||
|
|||||||
+297
-27
@@ -20,6 +20,7 @@ import (
|
|||||||
|
|
||||||
"github.com/ShukeBta/MMTL/internal/model"
|
"github.com/ShukeBta/MMTL/internal/model"
|
||||||
"github.com/ShukeBta/MMTL/internal/service/cloud"
|
"github.com/ShukeBta/MMTL/internal/service/cloud"
|
||||||
|
"github.com/ShukeBta/MMTL/internal/service/cloud115"
|
||||||
)
|
)
|
||||||
|
|
||||||
// strmSyncState 是一次同步执行的上下文。
|
// strmSyncState 是一次同步执行的上下文。
|
||||||
@@ -31,16 +32,19 @@ type strmSyncState struct {
|
|||||||
provider cloud.Provider // local 提供方为 nil
|
provider cloud.Provider // local 提供方为 nil
|
||||||
cfg *strmPathConfig
|
cfg *strmPathConfig
|
||||||
rec *model.StrmSyncRecord
|
rec *model.StrmSyncRecord
|
||||||
|
syncType string
|
||||||
|
|
||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
processed int // 已处理文件计数(用于定期落库进度)
|
processed int // 已处理文件计数(用于定期落库进度)
|
||||||
seenVideo map[string]bool // "v:"+去掉扩展名的相对路径 → 远端存在该视频
|
seenVideo map[string]bool // "v:"+去掉扩展名的相对路径 → 远端存在该视频
|
||||||
seenMeta map[string]bool // "m:"+相对路径 → 远端存在该元数据
|
seenMeta map[string]bool // "m:"+相对路径 → 远端存在该元数据
|
||||||
remoteMeta map[string]int64 // 远端元数据大小(上传比对用)
|
remoteMeta map[string]int64 // 远端元数据大小(上传比对用)
|
||||||
|
dirCache sync.Map // dirID (string) -> relativePath (string)
|
||||||
}
|
}
|
||||||
|
|
||||||
// StartSync 启动一次同步(异步执行,同一目录同时只允许一个任务)。
|
// StartSync 启动一次同步(异步执行,同一目录同时只允许一个任务)。
|
||||||
func (s *StrmService) StartSync(ctx context.Context, pathID string) error {
|
// syncType 支持 "incremental"(默认增量)和 "full"(全量同步)。
|
||||||
|
func (s *StrmService) StartSync(ctx context.Context, pathID string, syncType ...string) error {
|
||||||
p, err := s.repo.StrmSyncPath.FindByID(ctx, pathID)
|
p, err := s.repo.StrmSyncPath.FindByID(ctx, pathID)
|
||||||
if err != nil || p == nil {
|
if err != nil || p == nil {
|
||||||
return errNotFoundOr(err, "同步目录不存在")
|
return errNotFoundOr(err, "同步目录不存在")
|
||||||
@@ -64,9 +68,20 @@ func (s *StrmService) StartSync(ctx context.Context, pathID string) error {
|
|||||||
s.running[pathID] = cancel
|
s.running[pathID] = cancel
|
||||||
s.mu.Unlock()
|
s.mu.Unlock()
|
||||||
|
|
||||||
|
mode := model.StrmSyncTypeIncremental
|
||||||
|
if len(syncType) > 0 && syncType[0] != "" {
|
||||||
|
mode = syncType[0]
|
||||||
|
} else if p.SyncMode != "" {
|
||||||
|
mode = p.SyncMode
|
||||||
|
}
|
||||||
|
if mode != model.StrmSyncTypeFull {
|
||||||
|
mode = model.StrmSyncTypeIncremental
|
||||||
|
}
|
||||||
|
|
||||||
now := time.Now()
|
now := time.Now()
|
||||||
rec := &model.StrmSyncRecord{
|
rec := &model.StrmSyncRecord{
|
||||||
SyncPathID: pathID,
|
SyncPathID: pathID,
|
||||||
|
SyncType: mode,
|
||||||
Status: model.StrmSyncRecordRunning,
|
Status: model.StrmSyncRecordRunning,
|
||||||
StartedAt: &now,
|
StartedAt: &now,
|
||||||
}
|
}
|
||||||
@@ -145,6 +160,7 @@ func (s *StrmService) runSync(ctx context.Context, p *model.StrmSyncPath, rec *m
|
|||||||
p: p,
|
p: p,
|
||||||
cfg: cfg,
|
cfg: cfg,
|
||||||
rec: rec,
|
rec: rec,
|
||||||
|
syncType: rec.SyncType,
|
||||||
seenVideo: map[string]bool{},
|
seenVideo: map[string]bool{},
|
||||||
seenMeta: map[string]bool{},
|
seenMeta: map[string]bool{},
|
||||||
remoteMeta: map[string]int64{},
|
remoteMeta: map[string]int64{},
|
||||||
@@ -191,15 +207,19 @@ func (s *StrmService) finishSync(p *model.StrmSyncPath, rec *model.StrmSyncRecor
|
|||||||
p.LastSyncStatus = status
|
p.LastSyncStatus = status
|
||||||
p.LastSyncMessage = message
|
p.LastSyncMessage = message
|
||||||
if status != model.StrmSyncRecordFailed && message == "" {
|
if status != model.StrmSyncRecordFailed && message == "" {
|
||||||
p.LastSyncMessage = fmt.Sprintf("完成:新增/更新 %d 个 strm,下载 %d 个元数据,清理 %d 个文件",
|
syncTypeLabel := "增量"
|
||||||
rec.NewStrm, rec.NewMeta, rec.Pruned)
|
if rec.SyncType == model.StrmSyncTypeFull {
|
||||||
|
syncTypeLabel = "全量"
|
||||||
|
}
|
||||||
|
p.LastSyncMessage = fmt.Sprintf("[%s] 完成:新增/更新 %d 个 strm,跳过 %d 个,下载 %d 个元数据,清理 %d 个文件",
|
||||||
|
syncTypeLabel, rec.NewStrm, rec.Skipped, rec.NewMeta, rec.Pruned)
|
||||||
}
|
}
|
||||||
if err := s.repo.StrmSyncPath.Update(context.Background(), p); err != nil {
|
if err := s.repo.StrmSyncPath.Update(context.Background(), p); err != nil {
|
||||||
s.log.Warn("update strm sync path failed", zap.Error(err))
|
s.log.Warn("update strm sync path failed", zap.Error(err))
|
||||||
}
|
}
|
||||||
s.log.Info("strm sync finished",
|
s.log.Info("strm sync finished",
|
||||||
zap.String("path_id", p.ID), zap.String("status", status),
|
zap.String("path_id", p.ID), zap.String("sync_type", rec.SyncType), zap.String("status", status),
|
||||||
zap.Int64("new_strm", rec.NewStrm), zap.Int64("new_meta", rec.NewMeta),
|
zap.Int64("new_strm", rec.NewStrm), zap.Int64("skipped", rec.Skipped), zap.Int64("new_meta", rec.NewMeta),
|
||||||
zap.Int64("pruned", rec.Pruned), zap.String("message", message))
|
zap.Int64("pruned", rec.Pruned), zap.String("message", message))
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -208,8 +228,14 @@ func (st *strmSyncState) run() error {
|
|||||||
return fmt.Errorf("创建输出目录失败:%w", err)
|
return fmt.Errorf("创建输出目录失败:%w", err)
|
||||||
}
|
}
|
||||||
if st.provider != nil {
|
if st.provider != nil {
|
||||||
if err := st.walkRemote(); err != nil {
|
if open115, ok := st.provider.(cloud.OpenAPI115Provider); ok && st.p.Provider == model.StrmProvider115 {
|
||||||
return err
|
if err := st.walk115Flat(open115.OpenClient()); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
if err := st.walkRemote(); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
if err := st.walkLocalSource(); err != nil {
|
if err := st.walkLocalSource(); err != nil {
|
||||||
@@ -374,6 +400,215 @@ func (st *strmSyncState) isMetaExt(ext string) bool {
|
|||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// walk115Flat 使用 115 开放平台扁平化分页批量拉取机制与目录拓扑缓存(参考 QMediaSync)。
|
||||||
|
// 极大地降低 API 请求次数并支持毫秒级/秒级增量同步。
|
||||||
|
func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error {
|
||||||
|
ctx := st.ctx
|
||||||
|
rootCID := strings.TrimSpace(st.p.RemotePath)
|
||||||
|
if rootCID == "" {
|
||||||
|
rootCID = "0"
|
||||||
|
}
|
||||||
|
|
||||||
|
// 1. 目录拓扑缓存处理
|
||||||
|
st.dirCache.Store(rootCID, "")
|
||||||
|
if st.syncType == model.StrmSyncTypeFull {
|
||||||
|
// 全量同步:清空本路径的历史目录缓存
|
||||||
|
if err := st.s.repo.StrmDirCache.DeleteBySyncPathID(ctx, st.p.ID); err != nil {
|
||||||
|
st.s.log.Warn("delete strm dir cache failed", zap.Error(err))
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
// 增量同步:预加载历史目录缓存
|
||||||
|
cached, err := st.s.repo.StrmDirCache.ListBySyncPathID(ctx, st.p.ID)
|
||||||
|
if err == nil {
|
||||||
|
for _, item := range cached {
|
||||||
|
st.dirCache.Store(item.DirID, item.Path)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 2. 探测文件总数
|
||||||
|
const pageSize = 1150
|
||||||
|
firstBatch, totalCount, err := open115.GetFsListFlat(ctx, rootCID, 0, pageSize)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("115: 获取文件列表失败:%w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
allFiles := make([]cloud115.RemoteFile, 0, totalCount)
|
||||||
|
allFiles = append(allFiles, firstBatch...)
|
||||||
|
|
||||||
|
// 3. 并发分页拉取剩余文件
|
||||||
|
if totalCount > int64(len(firstBatch)) {
|
||||||
|
totalPages := int((totalCount + pageSize - 1) / pageSize)
|
||||||
|
type pageTask struct {
|
||||||
|
offset int
|
||||||
|
}
|
||||||
|
pageTasks := make([]pageTask, 0, totalPages-1)
|
||||||
|
for page := 1; page < totalPages; page++ {
|
||||||
|
pageTasks = append(pageTasks, pageTask{offset: page * pageSize})
|
||||||
|
}
|
||||||
|
|
||||||
|
var (
|
||||||
|
filesMu sync.Mutex
|
||||||
|
wg sync.WaitGroup
|
||||||
|
taskCh = make(chan pageTask, len(pageTasks))
|
||||||
|
errMu sync.Mutex
|
||||||
|
fetchErr error
|
||||||
|
)
|
||||||
|
|
||||||
|
for _, t := range pageTasks {
|
||||||
|
taskCh <- t
|
||||||
|
}
|
||||||
|
close(taskCh)
|
||||||
|
|
||||||
|
workers := 4
|
||||||
|
if len(pageTasks) < workers {
|
||||||
|
workers = len(pageTasks)
|
||||||
|
}
|
||||||
|
|
||||||
|
for i := 0; i < workers; i++ {
|
||||||
|
wg.Add(1)
|
||||||
|
go func() {
|
||||||
|
defer wg.Done()
|
||||||
|
for t := range taskCh {
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
files, _, err := open115.GetFsListFlat(ctx, rootCID, t.offset, pageSize)
|
||||||
|
if err != nil {
|
||||||
|
errMu.Lock()
|
||||||
|
if fetchErr == nil {
|
||||||
|
fetchErr = err
|
||||||
|
}
|
||||||
|
errMu.Unlock()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
filesMu.Lock()
|
||||||
|
allFiles = append(allFiles, files...)
|
||||||
|
filesMu.Unlock()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
wg.Wait()
|
||||||
|
if fetchErr != nil {
|
||||||
|
return fmt.Errorf("115: 分页拉取失败:%w", fetchErr)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
return ctx.Err()
|
||||||
|
}
|
||||||
|
|
||||||
|
// 4. 收集所有未在缓存中的父目录 ID (file.Pid)
|
||||||
|
missingPids := make(map[string]struct{})
|
||||||
|
for _, f := range allFiles {
|
||||||
|
pid := f.Pid
|
||||||
|
if pid == "" || pid == rootCID {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if _, ok := st.dirCache.Load(pid); !ok {
|
||||||
|
missingPids[pid] = struct{}{}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 并发补全未知目录详情与祖先链
|
||||||
|
if len(missingPids) > 0 {
|
||||||
|
pidList := make([]string, 0, len(missingPids))
|
||||||
|
for pid := range missingPids {
|
||||||
|
pidList = append(pidList, pid)
|
||||||
|
}
|
||||||
|
|
||||||
|
pidCh := make(chan string, len(pidList))
|
||||||
|
for _, pid := range pidList {
|
||||||
|
pidCh <- pid
|
||||||
|
}
|
||||||
|
close(pidCh)
|
||||||
|
|
||||||
|
var (
|
||||||
|
pwg sync.WaitGroup
|
||||||
|
dirWorkers = 4
|
||||||
|
)
|
||||||
|
if len(pidList) < dirWorkers {
|
||||||
|
dirWorkers = len(pidList)
|
||||||
|
}
|
||||||
|
|
||||||
|
for i := 0; i < dirWorkers; i++ {
|
||||||
|
pwg.Add(1)
|
||||||
|
go func() {
|
||||||
|
defer pwg.Done()
|
||||||
|
for pid := range pidCh {
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
detail, err := open115.GetFsDetailByCid(ctx, pid)
|
||||||
|
if err != nil {
|
||||||
|
st.s.log.Warn("115: 获取目录详情失败", zap.String("pid", pid), zap.Error(err))
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if detail == nil {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
// 解析相对路径
|
||||||
|
relPath := detail.RelativePath(rootCID)
|
||||||
|
st.dirCache.Store(pid, relPath)
|
||||||
|
_ = st.s.repo.StrmDirCache.Set(ctx, st.p.ID, pid, relPath)
|
||||||
|
|
||||||
|
// 顺便解析并缓存 detail.Paths 中包含的中间各层级目录
|
||||||
|
for _, ancestor := range detail.Paths {
|
||||||
|
if ancestor.FileId == "0" || ancestor.FileId == rootCID {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if _, loaded := st.dirCache.Load(ancestor.FileId); !loaded {
|
||||||
|
subDetail := &cloud115.RemoteFileDetail{
|
||||||
|
FileId: ancestor.FileId,
|
||||||
|
Paths: nil,
|
||||||
|
}
|
||||||
|
for _, p := range detail.Paths {
|
||||||
|
subDetail.Paths = append(subDetail.Paths, p)
|
||||||
|
if p.FileId == ancestor.FileId {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
ancestorRel := subDetail.RelativePath(rootCID)
|
||||||
|
st.dirCache.Store(ancestor.FileId, ancestorRel)
|
||||||
|
_ = st.s.repo.StrmDirCache.Set(ctx, st.p.ID, ancestor.FileId, ancestorRel)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
pwg.Wait()
|
||||||
|
}
|
||||||
|
|
||||||
|
// 5. 分类处理所有文件
|
||||||
|
for _, f := range allFiles {
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
return ctx.Err()
|
||||||
|
}
|
||||||
|
cleanName := cleanEntryName(f.FileName, false)
|
||||||
|
var rel string
|
||||||
|
if f.Pid == "" || f.Pid == rootCID {
|
||||||
|
rel = cleanName
|
||||||
|
} else {
|
||||||
|
if parentVal, ok := st.dirCache.Load(f.Pid); ok && parentVal.(string) != "" {
|
||||||
|
rel = parentVal.(string) + "/" + cleanName
|
||||||
|
} else {
|
||||||
|
rel = cleanName
|
||||||
|
}
|
||||||
|
}
|
||||||
|
entry := cloud.FileEntry{
|
||||||
|
ID: f.FileId,
|
||||||
|
Name: f.FileName,
|
||||||
|
IsDir: false,
|
||||||
|
Size: f.FileSize,
|
||||||
|
MTime: f.Utime,
|
||||||
|
PickCode: f.PickCode,
|
||||||
|
}
|
||||||
|
st.processRemoteFile(entry, rel)
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
// handleVideo 生成/更新 .strm 文件。
|
// handleVideo 生成/更新 .strm 文件。
|
||||||
func (st *strmSyncState) handleVideo(entry cloud.FileEntry, rel, ext string) {
|
func (st *strmSyncState) handleVideo(entry cloud.FileEntry, rel, ext string) {
|
||||||
relSansExt := rel[:len(rel)-len(ext)]
|
relSansExt := rel[:len(rel)-len(ext)]
|
||||||
@@ -391,6 +626,18 @@ func (st *strmSyncState) handleVideo(entry cloud.FileEntry, rel, ext string) {
|
|||||||
st.s.log.Warn("strm target path out of root", zap.String("rel", targetRel), zap.Error(err))
|
st.s.log.Warn("strm target path out of root", zap.String("rel", targetRel), zap.Error(err))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 增量同步模式快速检查:本地 strm 文件存在、非空且修改时间与远端 mtime 一致,直接跳过无需读磁盘
|
||||||
|
if st.syncType == model.StrmSyncTypeIncremental && entry.MTime > 0 {
|
||||||
|
if info, err := os.Stat(target); err == nil && info.Size() > 0 && info.ModTime().Unix() == entry.MTime {
|
||||||
|
st.mu.Lock()
|
||||||
|
st.rec.Skipped++
|
||||||
|
st.mu.Unlock()
|
||||||
|
st.touchProgress()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
content, err := st.strmContent(entry, rel, ext)
|
content, err := st.strmContent(entry, rel, ext)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
// 并发 worker 下 rec.Message 无锁写会有数据竞争,这里仅记录日志;
|
// 并发 worker 下 rec.Message 无锁写会有数据竞争,这里仅记录日志;
|
||||||
@@ -403,6 +650,11 @@ func (st *strmSyncState) handleVideo(entry cloud.FileEntry, rel, ext string) {
|
|||||||
existing = string(data)
|
existing = string(data)
|
||||||
}
|
}
|
||||||
if existing == content {
|
if existing == content {
|
||||||
|
// 对齐本地 strm 修改时间为远端 mtime,便于后续秒级比对
|
||||||
|
if entry.MTime > 0 {
|
||||||
|
mTime := time.Unix(entry.MTime, 0)
|
||||||
|
_ = os.Chtimes(target, mTime, mTime)
|
||||||
|
}
|
||||||
st.mu.Lock()
|
st.mu.Lock()
|
||||||
st.rec.Skipped++
|
st.rec.Skipped++
|
||||||
st.mu.Unlock()
|
st.mu.Unlock()
|
||||||
@@ -423,6 +675,10 @@ func (st *strmSyncState) handleVideo(entry cloud.FileEntry, rel, ext string) {
|
|||||||
st.s.log.Warn("rename strm failed", zap.String("file", target), zap.Error(err))
|
st.s.log.Warn("rename strm failed", zap.String("file", target), zap.Error(err))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
if entry.MTime > 0 {
|
||||||
|
mTime := time.Unix(entry.MTime, 0)
|
||||||
|
_ = os.Chtimes(target, mTime, mTime)
|
||||||
|
}
|
||||||
st.mu.Lock()
|
st.mu.Lock()
|
||||||
st.rec.NewStrm++
|
st.rec.NewStrm++
|
||||||
st.mu.Unlock()
|
st.mu.Unlock()
|
||||||
@@ -576,29 +832,43 @@ func (st *strmSyncState) walkLocalSource() error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
target, err := joinLocalRel(st.p.LocalPath, relSansExt+".strm")
|
target, err := joinLocalRel(st.p.LocalPath, relSansExt+".strm")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
if data, err := os.ReadFile(target); err == nil && string(data) == content {
|
mTime := info.ModTime()
|
||||||
|
if st.syncType == model.StrmSyncTypeIncremental {
|
||||||
|
if tInfo, err := os.Stat(target); err == nil && tInfo.Size() > 0 && tInfo.ModTime().Unix() == mTime.Unix() {
|
||||||
|
st.mu.Lock()
|
||||||
|
st.rec.Skipped++
|
||||||
|
st.mu.Unlock()
|
||||||
|
st.touchProgress()
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if data, err := os.ReadFile(target); err == nil && string(data) == content {
|
||||||
|
_ = os.Chtimes(target, mTime, mTime)
|
||||||
|
st.mu.Lock()
|
||||||
|
st.rec.Skipped++
|
||||||
|
st.mu.Unlock()
|
||||||
|
st.touchProgress()
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if err := os.MkdirAll(filepath.Dir(target), 0o755); err != nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
tmp := target + ".tmp"
|
||||||
|
if err := os.WriteFile(tmp, []byte(content), 0o644); err == nil {
|
||||||
|
_ = os.Rename(tmp, target)
|
||||||
|
_ = os.Chtimes(target, mTime, mTime)
|
||||||
|
} else {
|
||||||
|
_ = os.Remove(tmp)
|
||||||
|
}
|
||||||
st.mu.Lock()
|
st.mu.Lock()
|
||||||
st.rec.Skipped++
|
st.rec.NewStrm++
|
||||||
st.mu.Unlock()
|
st.mu.Unlock()
|
||||||
|
st.touchProgress()
|
||||||
return nil
|
return nil
|
||||||
}
|
|
||||||
if err := os.MkdirAll(filepath.Dir(target), 0o755); err != nil {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
tmp := target + ".tmp"
|
|
||||||
if err := os.WriteFile(tmp, []byte(content), 0o644); err == nil {
|
|
||||||
_ = os.Rename(tmp, target)
|
|
||||||
} else {
|
|
||||||
_ = os.Remove(tmp)
|
|
||||||
}
|
|
||||||
st.mu.Lock()
|
|
||||||
st.rec.NewStrm++
|
|
||||||
st.mu.Unlock()
|
|
||||||
return nil
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -34,10 +34,10 @@ func testStrmService(t *testing.T) *StrmService {
|
|||||||
sqlDB.SetMaxOpenConns(4)
|
sqlDB.SetMaxOpenConns(4)
|
||||||
t.Cleanup(func() { _ = sqlDB.Close() })
|
t.Cleanup(func() { _ = sqlDB.Close() })
|
||||||
}
|
}
|
||||||
if err := db.AutoMigrate(&model.StrmAccount{}, &model.StrmSyncPath{}, &model.StrmSyncRecord{},
|
if err := db.AutoMigrate(&model.StrmAccount{}, &model.StrmSyncPath{}, &model.StrmSyncRecord{},
|
||||||
&model.StrmDownloadTask{}, &model.StrmUploadTask{}, &model.Setting{}); err != nil {
|
&model.StrmDownloadTask{}, &model.StrmUploadTask{}, &model.StrmDirCache{}, &model.Setting{}); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
repos := repository.New(db)
|
repos := repository.New(db)
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
if err := repos.Setting.Set(ctx, StrmSettingBaseURL, "http://test.local:8096"); err != nil {
|
if err := repos.Setting.Set(ctx, StrmSettingBaseURL, "http://test.local:8096"); err != nil {
|
||||||
@@ -159,6 +159,57 @@ func TestLocalStrmSync(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestStrmFullAndIncrementalSync 测试增量同步与全量同步模式切换及记录
|
||||||
|
func TestStrmFullAndIncrementalSync(t *testing.T) {
|
||||||
|
svc := testStrmService(t)
|
||||||
|
src := t.TempDir()
|
||||||
|
out := t.TempDir()
|
||||||
|
|
||||||
|
writeFile(t, filepath.Join(src, "电影", "星际穿越.mkv"), "fake-video-data")
|
||||||
|
|
||||||
|
p := syncPathRecord(t, svc, model.StrmProviderLocal, src, out, true)
|
||||||
|
|
||||||
|
// 1. 默认触发增量同步
|
||||||
|
if err := svc.StartSync(context.Background(), p.ID, model.StrmSyncTypeIncremental); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
record := waitSyncDone(t, svc, p.ID, 10*time.Second)
|
||||||
|
if record.Status != model.StrmSyncRecordDone {
|
||||||
|
t.Fatalf("sync status = %s, message = %s", record.Status, record.Message)
|
||||||
|
}
|
||||||
|
if record.SyncType != model.StrmSyncTypeIncremental {
|
||||||
|
t.Fatalf("expected sync_type = incremental, got %s", record.SyncType)
|
||||||
|
}
|
||||||
|
if record.NewStrm != 1 {
|
||||||
|
t.Fatalf("expected 1 new strm, got %d", record.NewStrm)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 2. 再次执行增量同步,应当跳过
|
||||||
|
if err := svc.StartSync(context.Background(), p.ID, model.StrmSyncTypeIncremental); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
record = waitSyncDone(t, svc, p.ID, 10*time.Second)
|
||||||
|
if record.SyncType != model.StrmSyncTypeIncremental {
|
||||||
|
t.Fatalf("expected sync_type = incremental, got %s", record.SyncType)
|
||||||
|
}
|
||||||
|
if record.Skipped != 1 {
|
||||||
|
t.Fatalf("expected 1 skipped, got %d", record.Skipped)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 3. 执行全量同步
|
||||||
|
if err := svc.StartSync(context.Background(), p.ID, model.StrmSyncTypeFull); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
record = waitSyncDone(t, svc, p.ID, 10*time.Second)
|
||||||
|
if record.SyncType != model.StrmSyncTypeFull {
|
||||||
|
t.Fatalf("expected sync_type = full, got %s", record.SyncType)
|
||||||
|
}
|
||||||
|
if record.Status != model.StrmSyncRecordDone {
|
||||||
|
t.Fatalf("full sync failed: status = %s, message = %s", record.Status, record.Message)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
// TestStrmCronMatches cron 表达式匹配。
|
// TestStrmCronMatches cron 表达式匹配。
|
||||||
func TestStrmCronMatches(t *testing.T) {
|
func TestStrmCronMatches(t *testing.T) {
|
||||||
cases := []struct {
|
cases := []struct {
|
||||||
|
|||||||
+2
-1
@@ -107,7 +107,8 @@ export const strmAPI = {
|
|||||||
|
|
||||||
deletePath: (id: string) => api.delete(`/admin/strm/paths/${id}`).then((r) => r.data),
|
deletePath: (id: string) => api.delete(`/admin/strm/paths/${id}`).then((r) => r.data),
|
||||||
|
|
||||||
startSync: (id: string) => api.post(`/admin/strm/paths/${id}/sync`).then((r) => r.data),
|
startSync: (id: string, mode: 'incremental' | 'full' = 'incremental') =>
|
||||||
|
api.post(`/admin/strm/paths/${id}/sync`, null, { params: { mode } }).then((r) => r.data),
|
||||||
|
|
||||||
cancelSync: (id: string) => api.post(`/admin/strm/paths/${id}/cancel`).then((r) => r.data),
|
cancelSync: (id: string) => api.post(`/admin/strm/paths/${id}/cancel`).then((r) => r.data),
|
||||||
|
|
||||||
|
|||||||
@@ -674,6 +674,7 @@ export function StrmSyncPathDialog({
|
|||||||
delete_dir: existing?.delete_dir ?? false,
|
delete_dir: existing?.delete_dir ?? false,
|
||||||
cron: existing?.cron ?? '',
|
cron: existing?.cron ?? '',
|
||||||
enable_cron: existing?.enable_cron ?? false,
|
enable_cron: existing?.enable_cron ?? false,
|
||||||
|
sync_mode: existing?.sync_mode ?? 'incremental',
|
||||||
enabled: existing?.enabled ?? true,
|
enabled: existing?.enabled ?? true,
|
||||||
}))
|
}))
|
||||||
const [saving, setSaving] = useState(false)
|
const [saving, setSaving] = useState(false)
|
||||||
@@ -831,7 +832,7 @@ export function StrmSyncPathDialog({
|
|||||||
<input className={inputCls} value={form.exclude_name ?? ''} placeholder="sample,trailer" onChange={(e) => set('exclude_name', e.target.value)} />
|
<input className={inputCls} value={form.exclude_name ?? ''} placeholder="sample,trailer" onChange={(e) => set('exclude_name', e.target.value)} />
|
||||||
</Field>
|
</Field>
|
||||||
</div>
|
</div>
|
||||||
<div className="grid gap-3 md:grid-cols-2">
|
<div className="grid gap-3 md:grid-cols-3">
|
||||||
<Field label="STRM 链接 path 参数">
|
<Field label="STRM 链接 path 参数">
|
||||||
<select className={inputCls} value={form.add_path ?? 1} onChange={(e) => set('add_path', Number(e.target.value))}>
|
<select className={inputCls} value={form.add_path ?? 1} onChange={(e) => set('add_path', Number(e.target.value))}>
|
||||||
<option value={1}>完整远端路径</option>
|
<option value={1}>完整远端路径</option>
|
||||||
@@ -839,7 +840,13 @@ export function StrmSyncPathDialog({
|
|||||||
<option value={3}>不带 path</option>
|
<option value={3}>不带 path</option>
|
||||||
</select>
|
</select>
|
||||||
</Field>
|
</Field>
|
||||||
<Field label="定时同步 Cron" hint="5 段表达式,如 0 */6 * * *(每 6 小时)">
|
<Field label="默认同步模式" hint="定时触发或快速同步时的策略">
|
||||||
|
<select className={inputCls} value={form.sync_mode ?? 'incremental'} onChange={(e) => set('sync_mode', e.target.value as 'incremental' | 'full')}>
|
||||||
|
<option value="incremental">增量同步(快速)</option>
|
||||||
|
<option value="full">全量同步(全量校验)</option>
|
||||||
|
</select>
|
||||||
|
</Field>
|
||||||
|
<Field label="定时同步 Cron" hint="5 段表达式,如 0 */6 * * *">
|
||||||
<input className={inputCls} value={form.cron ?? ''} placeholder="0 */6 * * *" onChange={(e) => set('cron', e.target.value)} />
|
<input className={inputCls} value={form.cron ?? ''} placeholder="0 */6 * * *" onChange={(e) => set('cron', e.target.value)} />
|
||||||
</Field>
|
</Field>
|
||||||
</div>
|
</div>
|
||||||
|
|||||||
@@ -96,11 +96,11 @@ export function StrmManagePage() {
|
|||||||
return () => clearInterval(timer)
|
return () => clearInterval(timer)
|
||||||
}, [paths, refresh])
|
}, [paths, refresh])
|
||||||
|
|
||||||
const startSync = async (path: StrmSyncPath) => {
|
const startSync = async (path: StrmSyncPath, mode: 'incremental' | 'full' = 'incremental') => {
|
||||||
setActingPath(path.id)
|
setActingPath(path.id)
|
||||||
try {
|
try {
|
||||||
await strmAPI.startSync(path.id)
|
await strmAPI.startSync(path.id, mode)
|
||||||
toast.success(`已开始同步「${path.name}」`)
|
toast.success(`已开始${mode === 'full' ? '全量' : '增量'}同步「${path.name}」`)
|
||||||
await refresh()
|
await refresh()
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
toast.error(apiErrorMessage(err))
|
toast.error(apiErrorMessage(err))
|
||||||
@@ -344,7 +344,7 @@ function SyncPathSection({
|
|||||||
onAdd: () => void
|
onAdd: () => void
|
||||||
onEdit: (path: StrmSyncPath) => void
|
onEdit: (path: StrmSyncPath) => void
|
||||||
onDelete: (path: StrmSyncPath) => void
|
onDelete: (path: StrmSyncPath) => void
|
||||||
onStart: (path: StrmSyncPath) => void
|
onStart: (path: StrmSyncPath, mode?: 'incremental' | 'full') => void
|
||||||
onCancel: (path: StrmSyncPath) => void
|
onCancel: (path: StrmSyncPath) => void
|
||||||
}) {
|
}) {
|
||||||
return (
|
return (
|
||||||
@@ -401,15 +401,28 @@ function SyncPathSection({
|
|||||||
取消
|
取消
|
||||||
</button>
|
</button>
|
||||||
) : (
|
) : (
|
||||||
<button
|
<>
|
||||||
type="button"
|
<button
|
||||||
disabled={actingPath === path.id || !path.enabled}
|
type="button"
|
||||||
onClick={() => onStart(path)}
|
disabled={actingPath === path.id || !path.enabled}
|
||||||
className={`${iconButtonCls} disabled:opacity-40`}
|
onClick={() => onStart(path, 'incremental')}
|
||||||
>
|
className={`${iconButtonCls} text-brand-600 font-medium disabled:opacity-40`}
|
||||||
{actingPath === path.id ? <Loader2 size={14} className="animate-spin" /> : <Play size={14} />}
|
title="增量同步:基于目录缓存快速同步新增与更新文件"
|
||||||
立即同步
|
>
|
||||||
</button>
|
{actingPath === path.id ? <Loader2 size={14} className="animate-spin" /> : <Play size={14} />}
|
||||||
|
增量同步
|
||||||
|
</button>
|
||||||
|
<button
|
||||||
|
type="button"
|
||||||
|
disabled={actingPath === path.id || !path.enabled}
|
||||||
|
onClick={() => onStart(path, 'full')}
|
||||||
|
className={`${iconButtonCls} text-sand-600 disabled:opacity-40`}
|
||||||
|
title="全量同步:重置目录缓存并全量比对所有文件"
|
||||||
|
>
|
||||||
|
<RefreshCw size={14} />
|
||||||
|
全量同步
|
||||||
|
</button>
|
||||||
|
</>
|
||||||
)}
|
)}
|
||||||
<button type="button" onClick={() => onEdit(path)} className={`${iconButtonCls}`}>
|
<button type="button" onClick={() => onEdit(path)} className={`${iconButtonCls}`}>
|
||||||
<Pencil size={14} />
|
<Pencil size={14} />
|
||||||
@@ -456,9 +469,11 @@ function RecordSection({ records }: { records: StrmSyncRecord[] }) {
|
|||||||
<thead className="border-b border-gray-200 text-xs uppercase tracking-wider text-sand-500">
|
<thead className="border-b border-gray-200 text-xs uppercase tracking-wider text-sand-500">
|
||||||
<tr>
|
<tr>
|
||||||
<th className="px-3 py-2">时间</th>
|
<th className="px-3 py-2">时间</th>
|
||||||
|
<th className="px-3 py-2">类型</th>
|
||||||
<th className="px-3 py-2">状态</th>
|
<th className="px-3 py-2">状态</th>
|
||||||
<th className="px-3 py-2 text-right">扫描文件</th>
|
<th className="px-3 py-2 text-right">扫描文件</th>
|
||||||
<th className="px-3 py-2 text-right">新增 strm</th>
|
<th className="px-3 py-2 text-right">新增/更新</th>
|
||||||
|
<th className="px-3 py-2 text-right">跳过</th>
|
||||||
<th className="px-3 py-2 text-right">下载元数据</th>
|
<th className="px-3 py-2 text-right">下载元数据</th>
|
||||||
<th className="px-3 py-2 text-right">清理</th>
|
<th className="px-3 py-2 text-right">清理</th>
|
||||||
<th className="px-3 py-2">说明</th>
|
<th className="px-3 py-2">说明</th>
|
||||||
@@ -467,11 +482,17 @@ function RecordSection({ records }: { records: StrmSyncRecord[] }) {
|
|||||||
<tbody>
|
<tbody>
|
||||||
{records.map((record) => {
|
{records.map((record) => {
|
||||||
const meta = RECORD_STATUS_META[record.status] ?? RECORD_STATUS_META.pending
|
const meta = RECORD_STATUS_META[record.status] ?? RECORD_STATUS_META.pending
|
||||||
|
const isFull = record.sync_type === 'full'
|
||||||
return (
|
return (
|
||||||
<tr key={record.id} className="border-t border-gray-100">
|
<tr key={record.id} className="border-t border-gray-100">
|
||||||
<td className="whitespace-nowrap px-3 py-2 text-xs text-ink-50">
|
<td className="whitespace-nowrap px-3 py-2 text-xs text-ink-50">
|
||||||
{formatTime(record.started_at ?? record.created_at)}
|
{formatTime(record.started_at ?? record.created_at)}
|
||||||
</td>
|
</td>
|
||||||
|
<td className="px-3 py-2">
|
||||||
|
<span className={`rounded-full px-2 py-0.5 text-[11px] font-medium ${isFull ? 'bg-amber-50 text-amber-600 border border-amber-200' : 'bg-brand-50 text-brand-600 border border-brand-200'}`}>
|
||||||
|
{isFull ? '全量' : '增量'}
|
||||||
|
</span>
|
||||||
|
</td>
|
||||||
<td className="px-3 py-2">
|
<td className="px-3 py-2">
|
||||||
<span className={'rounded-full px-2 py-0.5 text-[11px] font-semibold ' + meta.cls}>
|
<span className={'rounded-full px-2 py-0.5 text-[11px] font-semibold ' + meta.cls}>
|
||||||
{meta.label}
|
{meta.label}
|
||||||
@@ -479,6 +500,7 @@ function RecordSection({ records }: { records: StrmSyncRecord[] }) {
|
|||||||
</td>
|
</td>
|
||||||
<td className="px-3 py-2 text-right">{record.total}</td>
|
<td className="px-3 py-2 text-right">{record.total}</td>
|
||||||
<td className="px-3 py-2 text-right text-brand-500">{record.new_strm}</td>
|
<td className="px-3 py-2 text-right text-brand-500">{record.new_strm}</td>
|
||||||
|
<td className="px-3 py-2 text-right text-gray-500">{record.skipped}</td>
|
||||||
<td className="px-3 py-2 text-right">{record.new_meta}</td>
|
<td className="px-3 py-2 text-right">{record.new_meta}</td>
|
||||||
<td className="px-3 py-2 text-right">{record.pruned}</td>
|
<td className="px-3 py-2 text-right">{record.pruned}</td>
|
||||||
<td className="max-w-[260px] truncate px-3 py-2 text-xs text-sand-500">{record.message}</td>
|
<td className="max-w-[260px] truncate px-3 py-2 text-xs text-sand-500">{record.message}</td>
|
||||||
|
|||||||
@@ -49,6 +49,7 @@ export interface StrmSyncPath {
|
|||||||
delete_dir: boolean
|
delete_dir: boolean
|
||||||
cron: string
|
cron: string
|
||||||
enable_cron: boolean
|
enable_cron: boolean
|
||||||
|
sync_mode?: 'incremental' | 'full'
|
||||||
enabled: boolean
|
enabled: boolean
|
||||||
created_at: string
|
created_at: string
|
||||||
last_sync_at?: string | null
|
last_sync_at?: string | null
|
||||||
@@ -75,12 +76,14 @@ export interface StrmSyncPathInput {
|
|||||||
delete_dir?: boolean
|
delete_dir?: boolean
|
||||||
cron?: string
|
cron?: string
|
||||||
enable_cron?: boolean
|
enable_cron?: boolean
|
||||||
|
sync_mode?: 'incremental' | 'full'
|
||||||
enabled?: boolean
|
enabled?: boolean
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface StrmSyncRecord {
|
export interface StrmSyncRecord {
|
||||||
id: string
|
id: string
|
||||||
sync_path_id: string
|
sync_path_id: string
|
||||||
|
sync_type?: 'incremental' | 'full'
|
||||||
status: 'pending' | 'running' | 'done' | 'failed' | 'canceled'
|
status: 'pending' | 'running' | 'done' | 'failed' | 'canceled'
|
||||||
total: number
|
total: number
|
||||||
new_strm: number
|
new_strm: number
|
||||||
|
|||||||
Reference in New Issue
Block a user