mirror of
https://github.com/truewhile/MeBox.git
synced 2026-09-28 11:16:37 +08:00
Compare commits
9 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 1ea4724261 | |||
| 3d372f039e | |||
| 0384017e98 | |||
| 0332579d5f | |||
| 6aefe18caa | |||
| ef72fc8d83 | |||
| 4764c09572 | |||
| 98ca766a37 | |||
| 13c9035b76 |
@@ -50,12 +50,14 @@ func registerAdminStrmRoutes(admin *gin.RouterGroup, svc *service.Container) {
|
||||
admin.POST("/strm/downloads/:id/retry", retryStrmDownloadHandler(svc))
|
||||
admin.POST("/strm/downloads/clear-done", clearDoneDownloadsHandler(svc))
|
||||
admin.POST("/strm/downloads/clear-finished", clearFinishedDownloadsHandler(svc))
|
||||
admin.POST("/strm/downloads/clear-canceled", clearCanceledDownloadsHandler(svc))
|
||||
admin.POST("/strm/downloads/retry-failed", retryAllFailedDownloadsHandler(svc))
|
||||
admin.POST("/strm/downloads/cancel-pending", cancelPendingDownloadsHandler(svc))
|
||||
admin.GET("/strm/uploads", uploadQueueHandler(svc))
|
||||
admin.POST("/strm/uploads/:id/cancel", cancelStrmUploadHandler(svc))
|
||||
admin.POST("/strm/uploads/:id/retry", retryStrmUploadHandler(svc))
|
||||
admin.POST("/strm/uploads/cancel-pending", cancelPendingUploadsHandler(svc))
|
||||
admin.POST("/strm/uploads/clear-canceled", clearCanceledUploadsHandler(svc))
|
||||
}
|
||||
|
||||
func registerAdminUserRoutes(admin *gin.RouterGroup, svc *service.Container) {
|
||||
|
||||
@@ -392,6 +392,28 @@ func clearFinishedDownloadsHandler(svc *service.Container) gin.HandlerFunc {
|
||||
}
|
||||
}
|
||||
|
||||
func clearCanceledDownloadsHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
n, err := svc.Strm.ClearCanceledDownloadTasks(c.Request.Context())
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"deleted": n})
|
||||
}
|
||||
}
|
||||
|
||||
func clearCanceledUploadsHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
n, err := svc.Strm.ClearCanceledUploadTasks(c.Request.Context())
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"deleted": n})
|
||||
}
|
||||
}
|
||||
|
||||
func retryAllFailedDownloadsHandler(svc *service.Container) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
n, err := svc.Strm.RetryAllFailedDownloadTasks(c.Request.Context())
|
||||
|
||||
@@ -50,10 +50,16 @@ func TestStrmAdminRoutesAreRegistered(t *testing.T) {
|
||||
"GET /api/admin/strm/downloads",
|
||||
"POST /api/admin/strm/downloads/:id/cancel",
|
||||
"POST /api/admin/strm/downloads/:id/retry",
|
||||
"GET /api/admin/strm/uploads",
|
||||
"POST /api/admin/strm/uploads/:id/cancel",
|
||||
"POST /api/admin/strm/uploads/:id/retry",
|
||||
"GET /api/strm/play/:provider/:file",
|
||||
"POST /api/admin/strm/downloads/clear-finished",
|
||||
"POST /api/admin/strm/downloads/clear-canceled",
|
||||
"POST /api/admin/strm/downloads/retry-failed",
|
||||
"POST /api/admin/strm/downloads/cancel-pending",
|
||||
"GET /api/admin/strm/uploads",
|
||||
"POST /api/admin/strm/uploads/:id/cancel",
|
||||
"POST /api/admin/strm/uploads/:id/retry",
|
||||
"POST /api/admin/strm/uploads/cancel-pending",
|
||||
"POST /api/admin/strm/uploads/clear-canceled",
|
||||
"GET /api/strm/play/:provider/:file",
|
||||
} {
|
||||
if !routes[want] {
|
||||
t.Fatalf("%s route is not registered", want)
|
||||
|
||||
@@ -163,7 +163,7 @@ func historyDeleteHandler(svc *service.Container) gin.HandlerFunc {
|
||||
c.JSON(http.StatusBadRequest, gin.H{"error": "status must be completed or incomplete"})
|
||||
return
|
||||
}
|
||||
res := q.Delete(&model.PlaybackHistory{})
|
||||
res := q.Unscoped().Delete(&model.PlaybackHistory{})
|
||||
if err := res.Error; err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
|
||||
@@ -62,9 +62,9 @@ func (r *ApiConfigRepository) Update(ctx context.Context, c *model.ApiConfig) er
|
||||
}).Error
|
||||
}
|
||||
|
||||
// Delete removes an API config.
|
||||
// Delete 物理删除 API 配置。
|
||||
func (r *ApiConfigRepository) Delete(ctx context.Context, provider string) error {
|
||||
return r.db.WithContext(ctx).Where("provider = ?", provider).Delete(&model.ApiConfig{}).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Where("provider = ?", provider).Delete(&model.ApiConfig{}).Error
|
||||
}
|
||||
|
||||
// UpdateTestResult 更新测试结果。
|
||||
|
||||
@@ -23,7 +23,7 @@ func (r *FavoriteRepository) Toggle(ctx context.Context, userID, mediaID string)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return false, r.db.WithContext(ctx).Delete(&f).Error
|
||||
return false, r.db.WithContext(ctx).Unscoped().Delete(&f).Error
|
||||
}
|
||||
|
||||
// ListByUser returns all favourite media IDs for a user.
|
||||
|
||||
@@ -79,10 +79,9 @@ func (r *LibraryRepository) FindByID(ctx context.Context, id string) (*model.Lib
|
||||
return &l, nil
|
||||
}
|
||||
|
||||
// Delete removes a library and (soft) cascades to its media via repository
|
||||
// callers; we do not run CASCADE here to keep this method narrow.
|
||||
// Delete 物理删除媒体库。
|
||||
func (r *LibraryRepository) Delete(ctx context.Context, id string) error {
|
||||
return r.db.WithContext(ctx).Delete(&model.Library{}, "id = ?", id).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Delete(&model.Library{}, "id = ?", id).Error
|
||||
}
|
||||
|
||||
func (r *LibraryRepository) ListRoots(ctx context.Context, libraryID string) ([]model.LibraryRoot, error) {
|
||||
@@ -149,7 +148,7 @@ func (r *LibraryRepository) DeleteRoot(ctx context.Context, libraryID, rootID st
|
||||
if !r.hasLibraryRootsTable() {
|
||||
return nil
|
||||
}
|
||||
return r.db.WithContext(ctx).Where("library_id = ?", libraryID).Delete(&model.LibraryRoot{}, "id = ?", rootID).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Where("library_id = ?", libraryID).Delete(&model.LibraryRoot{}, "id = ?", rootID).Error
|
||||
}
|
||||
|
||||
func (r *LibraryRepository) hasLibraryRootsTable() bool {
|
||||
|
||||
@@ -116,12 +116,12 @@ func (r *MediaRepository) ListByLibrariesFiltered(ctx context.Context, libraryID
|
||||
|
||||
// DeleteByLibrary purges all media tied to a library.
|
||||
func (r *MediaRepository) DeleteByLibrary(ctx context.Context, libraryID string) error {
|
||||
// FTS 行由 media 表上的触发器同步清理(软删/硬删都覆盖)。
|
||||
return r.db.WithContext(ctx).Where("library_id = ?", libraryID).Delete(&model.Media{}).Error
|
||||
// FTS 行由 media 表上的触发器同步清理(物理删除触发 FTS 清理)。
|
||||
return r.db.WithContext(ctx).Unscoped().Where("library_id = ?", libraryID).Delete(&model.Media{}).Error
|
||||
}
|
||||
|
||||
func (r *MediaRepository) DeleteByLibraryRoot(ctx context.Context, libraryID, rootID string) error {
|
||||
return r.db.WithContext(ctx).
|
||||
return r.db.WithContext(ctx).Unscoped().
|
||||
Where("library_id = ? AND library_root_id = ?", libraryID, rootID).
|
||||
Delete(&model.Media{}).Error
|
||||
}
|
||||
|
||||
@@ -51,9 +51,9 @@ func (r *PermissionRepository) Upsert(ctx context.Context, p *model.UserPermissi
|
||||
})
|
||||
}
|
||||
|
||||
// Delete removes a permission record.
|
||||
// Delete 物理删除权限记录。
|
||||
func (r *PermissionRepository) Delete(ctx context.Context, userID string) error {
|
||||
return withSQLiteBusyRetry(ctx, func() error {
|
||||
return r.db.WithContext(ctx).Where("user_id = ?", userID).Delete(&model.UserPermission{}).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Where("user_id = ?", userID).Delete(&model.UserPermission{}).Error
|
||||
})
|
||||
}
|
||||
|
||||
@@ -59,9 +59,9 @@ func (r *PlayProfileRepository) Update(ctx context.Context, id string, patch map
|
||||
Where("id = ?", id).Updates(patch).Error
|
||||
}
|
||||
|
||||
// Delete soft-deletes a profile.
|
||||
// Delete 物理删除播放档案。
|
||||
func (r *PlayProfileRepository) Delete(ctx context.Context, id string) error {
|
||||
return r.db.WithContext(ctx).Delete(&model.PlayProfile{}, "id = ?", id).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Delete(&model.PlayProfile{}, "id = ?", id).Error
|
||||
}
|
||||
|
||||
// ClearDefaultsFor resets is_default for all of a user's profiles.
|
||||
|
||||
@@ -72,10 +72,10 @@ func (r *RefreshTokenRepository) RevokeOldestActiveByUserID(ctx context.Context,
|
||||
})
|
||||
}
|
||||
|
||||
// DeleteExpired removes all expired refresh tokens.
|
||||
// DeleteExpired 物理清理所有过期的 refresh tokens。
|
||||
func (r *RefreshTokenRepository) DeleteExpired(ctx context.Context) error {
|
||||
return withSQLiteBusyRetry(ctx, func() error {
|
||||
return r.db.WithContext(ctx).Where("expires_at < ?", time.Now()).Delete(&model.RefreshToken{}).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Where("expires_at < ?", time.Now()).Delete(&model.RefreshToken{}).Error
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -29,9 +29,9 @@ func (r *SettingRepository) Set(ctx context.Context, key, value string) error {
|
||||
return r.db.WithContext(ctx).Save(&s).Error
|
||||
}
|
||||
|
||||
// Delete removes a setting key.
|
||||
// Delete 物理删除设置键。
|
||||
func (r *SettingRepository) Delete(ctx context.Context, key string) error {
|
||||
return r.db.WithContext(ctx).Where("key = ?", key).Delete(&model.Setting{}).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Where("key = ?", key).Delete(&model.Setting{}).Error
|
||||
}
|
||||
|
||||
// All returns every key/value pair (used by the admin UI).
|
||||
|
||||
@@ -66,9 +66,9 @@ func (r *StorageConfigRepository) Upsert(ctx context.Context, c *model.StorageCo
|
||||
}).Error
|
||||
}
|
||||
|
||||
// Delete removes a storage config by ID.
|
||||
// Delete 物理删除存储配置。
|
||||
func (r *StorageConfigRepository) Delete(ctx context.Context, id string) error {
|
||||
return r.db.WithContext(ctx).Where("id = ?", id).Delete(&model.StorageConfig{}).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.StorageConfig{}).Error
|
||||
}
|
||||
|
||||
// FindByID returns a storage config by ID.
|
||||
|
||||
@@ -59,7 +59,7 @@ func (r *StrmAccountRepository) Update(ctx context.Context, a *model.StrmAccount
|
||||
|
||||
func (r *StrmAccountRepository) Delete(ctx context.Context, id string) error {
|
||||
return withSQLiteBusyRetry(ctx, func() error {
|
||||
return r.db.WithContext(ctx).Where("id = ?", id).Delete(&model.StrmAccount{}).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.StrmAccount{}).Error
|
||||
})
|
||||
}
|
||||
|
||||
@@ -123,7 +123,7 @@ func (r *StrmSyncPathRepository) Update(ctx context.Context, p *model.StrmSyncPa
|
||||
|
||||
func (r *StrmSyncPathRepository) Delete(ctx context.Context, id string) error {
|
||||
return withSQLiteBusyRetry(ctx, func() error {
|
||||
return r.db.WithContext(ctx).Where("id = ?", id).Delete(&model.StrmSyncPath{}).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.StrmSyncPath{}).Error
|
||||
})
|
||||
}
|
||||
|
||||
@@ -286,7 +286,7 @@ func (r *StrmDownloadTaskRepository) Update(ctx context.Context, t *model.StrmDo
|
||||
|
||||
func (r *StrmDownloadTaskRepository) Delete(ctx context.Context, id string) error {
|
||||
return withSQLiteBusyRetry(ctx, func() error {
|
||||
return r.db.WithContext(ctx).Where("id = ?", id).Delete(&model.StrmDownloadTask{}).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.StrmDownloadTask{}).Error
|
||||
})
|
||||
}
|
||||
|
||||
@@ -294,7 +294,7 @@ func (r *StrmDownloadTaskRepository) Delete(ctx context.Context, id string) erro
|
||||
func (r *StrmDownloadTaskRepository) ClearDone(ctx context.Context) (int64, error) {
|
||||
var count int64
|
||||
err := withSQLiteBusyRetry(ctx, func() error {
|
||||
res := r.db.WithContext(ctx).Where("status = ?", model.StrmTaskDone).Delete(&model.StrmDownloadTask{})
|
||||
res := r.db.WithContext(ctx).Unscoped().Where("status = ?", model.StrmTaskDone).Delete(&model.StrmDownloadTask{})
|
||||
count = res.RowsAffected
|
||||
return res.Error
|
||||
})
|
||||
@@ -305,7 +305,7 @@ func (r *StrmDownloadTaskRepository) ClearDone(ctx context.Context) (int64, erro
|
||||
func (r *StrmDownloadTaskRepository) ClearFinished(ctx context.Context) (int64, error) {
|
||||
var count int64
|
||||
err := withSQLiteBusyRetry(ctx, func() error {
|
||||
res := r.db.WithContext(ctx).Where("status IN ?", []string{model.StrmTaskDone, model.StrmTaskFailed}).
|
||||
res := r.db.WithContext(ctx).Unscoped().Where("status IN ?", []string{model.StrmTaskDone, model.StrmTaskFailed, model.StrmTaskCanceled}).
|
||||
Delete(&model.StrmDownloadTask{})
|
||||
count = res.RowsAffected
|
||||
return res.Error
|
||||
@@ -313,6 +313,17 @@ func (r *StrmDownloadTaskRepository) ClearFinished(ctx context.Context) (int64,
|
||||
return count, err
|
||||
}
|
||||
|
||||
// ClearCanceled 清空全部已取消下载任务。
|
||||
func (r *StrmDownloadTaskRepository) ClearCanceled(ctx context.Context) (int64, error) {
|
||||
var count int64
|
||||
err := withSQLiteBusyRetry(ctx, func() error {
|
||||
res := r.db.WithContext(ctx).Unscoped().Where("status = ?", model.StrmTaskCanceled).Delete(&model.StrmDownloadTask{})
|
||||
count = res.RowsAffected
|
||||
return res.Error
|
||||
})
|
||||
return count, err
|
||||
}
|
||||
|
||||
// RetryAllFailed 把所有失败任务重置回待处理,清空错误与重试计数。
|
||||
func (r *StrmDownloadTaskRepository) RetryAllFailed(ctx context.Context) (int64, error) {
|
||||
var count int64
|
||||
@@ -381,7 +392,7 @@ func (r *StrmDownloadTaskRepository) GetActiveLocalPathMap(ctx context.Context,
|
||||
|
||||
func (r *StrmDownloadTaskRepository) DeleteFinishedOlderThan(ctx context.Context, before time.Time) error {
|
||||
return withSQLiteBusyRetry(ctx, func() error {
|
||||
return r.db.WithContext(ctx).Where("status IN ? AND finished_at < ?",
|
||||
return r.db.WithContext(ctx).Unscoped().Where("status IN ? AND finished_at < ?",
|
||||
[]string{model.StrmTaskDone, model.StrmTaskFailed, model.StrmTaskCanceled}, before).
|
||||
Delete(&model.StrmDownloadTask{}).Error
|
||||
})
|
||||
@@ -522,10 +533,21 @@ func (r *StrmUploadTaskRepository) Update(ctx context.Context, t *model.StrmUplo
|
||||
|
||||
func (r *StrmUploadTaskRepository) Delete(ctx context.Context, id string) error {
|
||||
return withSQLiteBusyRetry(ctx, func() error {
|
||||
return r.db.WithContext(ctx).Where("id = ?", id).Delete(&model.StrmUploadTask{}).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.StrmUploadTask{}).Error
|
||||
})
|
||||
}
|
||||
|
||||
// ClearCanceled 清空全部已取消上传任务。
|
||||
func (r *StrmUploadTaskRepository) ClearCanceled(ctx context.Context) (int64, error) {
|
||||
var count int64
|
||||
err := withSQLiteBusyRetry(ctx, func() error {
|
||||
res := r.db.WithContext(ctx).Unscoped().Where("status = ?", model.StrmTaskCanceled).Delete(&model.StrmUploadTask{})
|
||||
count = res.RowsAffected
|
||||
return res.Error
|
||||
})
|
||||
return count, err
|
||||
}
|
||||
|
||||
// CancelPending 批量取消所有排队中和进行中的任务。
|
||||
func (r *StrmUploadTaskRepository) CancelPending(ctx context.Context) (int64, error) {
|
||||
now := time.Now()
|
||||
@@ -573,7 +595,7 @@ func (r *StrmUploadTaskRepository) GetActiveLocalPathMap(ctx context.Context, sy
|
||||
|
||||
func (r *StrmUploadTaskRepository) DeleteFinishedOlderThan(ctx context.Context, before time.Time) error {
|
||||
return withSQLiteBusyRetry(ctx, func() error {
|
||||
return r.db.WithContext(ctx).Where("status IN ? AND finished_at < ?",
|
||||
return r.db.WithContext(ctx).Unscoped().Where("status IN ? AND finished_at < ?",
|
||||
[]string{model.StrmTaskDone, model.StrmTaskFailed, model.StrmTaskCanceled}, before).
|
||||
Delete(&model.StrmUploadTask{}).Error
|
||||
})
|
||||
@@ -614,7 +636,7 @@ func (r *StrmDirCacheRepository) Set(ctx context.Context, syncPathID, dirID, pat
|
||||
|
||||
func (r *StrmDirCacheRepository) DeleteBySyncPathID(ctx context.Context, syncPathID string) error {
|
||||
return withSQLiteBusyRetry(ctx, func() error {
|
||||
return r.db.WithContext(ctx).Where("sync_path_id = ?", syncPathID).Delete(&model.StrmDirCache{}).Error
|
||||
return r.db.WithContext(ctx).Unscoped().Where("sync_path_id = ?", syncPathID).Delete(&model.StrmDirCache{}).Error
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -133,26 +133,17 @@ func (r *UserRepository) TouchLogin(ctx context.Context, id string) error {
|
||||
})
|
||||
}
|
||||
|
||||
// Delete removes a user (soft-delete via gorm.DeletedAt), releases the unique
|
||||
// username, and drops Telegram bindings so future re-created users bind cleanly.
|
||||
// Delete 物理删除用户并级联清理其关联记录。
|
||||
func (r *UserRepository) Delete(ctx context.Context, id string) error {
|
||||
return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
var user model.User
|
||||
if err := tx.Where("id = ?", id).First(&user).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
released := user.Username + "__deleted__" + time.Now().Format("20060102150405.000000000")
|
||||
if len(released) > 64 {
|
||||
sum := sha256.Sum256([]byte(user.ID + user.Username))
|
||||
base := user.Username
|
||||
if len(base) > 43 {
|
||||
base = base[:43]
|
||||
}
|
||||
released = base + "__deleted__" + hex.EncodeToString(sum[:])[:10]
|
||||
}
|
||||
if err := tx.Model(&model.User{}).Where("id = ?", id).Update("username", released).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Delete(&model.User{}, "id = ?", id).Error
|
||||
return withSQLiteBusyRetry(ctx, func() error {
|
||||
return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
_ = tx.Unscoped().Where("user_id = ?", id).Delete(&model.RefreshToken{})
|
||||
_ = tx.Unscoped().Where("user_id = ?", id).Delete(&model.UserPermission{})
|
||||
_ = tx.Unscoped().Where("user_id = ?", id).Delete(&model.PlayProfile{})
|
||||
_ = tx.Unscoped().Where("user_id = ?", id).Delete(&model.PlaybackHistory{})
|
||||
_ = tx.Unscoped().Where("user_id = ?", id).Delete(&model.Favorite{})
|
||||
_ = tx.Unscoped().Where("user_id = ?", id).Delete(&model.UserDevice{})
|
||||
return tx.Unscoped().Delete(&model.User{}, "id = ?", id).Error
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
@@ -43,10 +43,10 @@ func (s *MediaService) DeleteLibrary(ctx context.Context, id string) error {
|
||||
if err := tx.Unscoped().Where("library_id = ?", id).Delete(&model.Media{}).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
if err := hardDeleteLibraryRoots(ctx, tx, id); err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Delete(&model.Library{}, "id = ?", id).Error
|
||||
if err := hardDeleteLibraryRoots(ctx, tx, id); err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Unscoped().Delete(&model.Library{}, "id = ?", id).Error
|
||||
})
|
||||
if err == nil {
|
||||
s.invalidateMediaCache(ctx)
|
||||
|
||||
@@ -10,25 +10,10 @@ import (
|
||||
|
||||
const maxRecycleBinRecords = 200
|
||||
|
||||
// SoftDelete moves a media row to the recycle bin (gorm soft delete).
|
||||
// The on-disk file is kept; admins can purge it later.
|
||||
// SoftDelete 物理删除媒体记录(统一硬删除以降低 SQLite 存储与索引压力)。
|
||||
func (s *MediaService) SoftDelete(ctx context.Context, id string) error {
|
||||
media, err := s.repo.Media.FindByID(ctx, id)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if media != nil && isCloudMediaPath(media.Path) {
|
||||
err := s.repo.DB.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.Media{}).Error
|
||||
if err == nil {
|
||||
s.invalidateMediaCache(ctx)
|
||||
}
|
||||
return err
|
||||
}
|
||||
err = s.repo.DB.WithContext(ctx).Where("id = ?", id).Delete(&model.Media{}).Error
|
||||
err := s.repo.DB.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.Media{}).Error
|
||||
if err == nil {
|
||||
if pruneErr := pruneRecycleBinRows(ctx, s.repo.DB, maxRecycleBinRecords); pruneErr != nil {
|
||||
return pruneErr
|
||||
}
|
||||
s.invalidateMediaCache(ctx)
|
||||
}
|
||||
return err
|
||||
|
||||
@@ -229,9 +229,9 @@ func (o *OrganizerService) replaceVersions(ctx context.Context, src string, exis
|
||||
o.log.Warn("organize replace remove existing failed",
|
||||
zap.String("path", e), zap.Error(err))
|
||||
}
|
||||
if o.repo != nil && o.repo.DB != nil {
|
||||
_ = o.repo.DB.WithContext(ctx).Where("path = ?", e).Delete(&model.Media{}).Error
|
||||
}
|
||||
if o.repo != nil && o.repo.DB != nil {
|
||||
_ = o.repo.DB.WithContext(ctx).Unscoped().Where("path = ?", e).Delete(&model.Media{}).Error
|
||||
}
|
||||
}
|
||||
// Move staged file + sidecars into the final path.
|
||||
if err := os.Rename(stage, dst); err != nil {
|
||||
|
||||
@@ -97,7 +97,7 @@ func (o *OrganizerService) deleteMediaRowForPath(ctx context.Context, path strin
|
||||
if o == nil || o.repo == nil || o.repo.DB == nil {
|
||||
return
|
||||
}
|
||||
_ = o.repo.DB.WithContext(ctx).Where("path = ?", path).Delete(&model.Media{}).Error
|
||||
_ = o.repo.DB.WithContext(ctx).Unscoped().Where("path = ?", path).Delete(&model.Media{}).Error
|
||||
}
|
||||
|
||||
func (o *OrganizerService) mediaPathExists(ctx context.Context, path string) bool {
|
||||
|
||||
@@ -196,18 +196,18 @@ func (p *PlaybackService) AddToPlaylist(ctx context.Context, playlistID, mediaID
|
||||
return p.repo.DB.Create(item).Error
|
||||
}
|
||||
|
||||
// RemoveFromPlaylist removes a media item from a playlist (idempotent).
|
||||
// RemoveFromPlaylist 物理删除播放列表项(幂等)。
|
||||
func (p *PlaybackService) RemoveFromPlaylist(ctx context.Context, playlistID, mediaID string) error {
|
||||
return p.repo.DB.
|
||||
return p.repo.DB.WithContext(ctx).Unscoped().
|
||||
Where("playlist_id = ? AND media_id = ?", playlistID, mediaID).
|
||||
Delete(&model.PlaylistItem{}).Error
|
||||
}
|
||||
|
||||
// DeletePlaylist removes a playlist and all of its items.
|
||||
// DeletePlaylist 物理删除播放列表及其全部条目。
|
||||
func (p *PlaybackService) DeletePlaylist(ctx context.Context, playlistID string) error {
|
||||
if err := p.repo.DB.Where("playlist_id = ?", playlistID).
|
||||
if err := p.repo.DB.WithContext(ctx).Unscoped().Where("playlist_id = ?", playlistID).
|
||||
Delete(&model.PlaylistItem{}).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
return p.repo.DB.Where("id = ?", playlistID).Delete(&model.Playlist{}).Error
|
||||
return p.repo.DB.WithContext(ctx).Unscoped().Where("id = ?", playlistID).Delete(&model.Playlist{}).Error
|
||||
}
|
||||
|
||||
@@ -11,13 +11,12 @@ import (
|
||||
"github.com/ShukeBta/MMTL/internal/model"
|
||||
)
|
||||
|
||||
// RemovePath deletes the media row for a path that has disappeared from disk
|
||||
// (incremental delete used by the watcher on Remove/Rename events).
|
||||
// RemovePath 物理删除磁盘上已不存在的媒体记录。
|
||||
func (s *ScannerService) RemovePath(ctx context.Context, path string) (int64, error) {
|
||||
if _, err := os.Stat(path); err == nil {
|
||||
return 0, nil // still exists; nothing to remove
|
||||
}
|
||||
res := s.repo.DB.WithContext(ctx).
|
||||
res := s.repo.DB.WithContext(ctx).Unscoped().
|
||||
Where("path = ?", path).
|
||||
Delete(&model.Media{})
|
||||
if res.Error == nil && res.RowsAffected > 0 {
|
||||
@@ -55,7 +54,7 @@ func (s *ScannerService) pruneMissingMedia(ctx context.Context, libraryID string
|
||||
}
|
||||
stale = append(stale, row.ID)
|
||||
}
|
||||
return s.deleteMediaByIDs(ctx, stale, false)
|
||||
return s.deleteMediaByIDs(ctx, stale, true)
|
||||
}
|
||||
|
||||
func (s *ScannerService) pruneMissingMediaForRoot(ctx context.Context, libraryID, rootID, rootPath string, seen map[string]struct{}) (int64, error) {
|
||||
@@ -92,7 +91,7 @@ func (s *ScannerService) pruneMissingMediaForRoot(ctx context.Context, libraryID
|
||||
}
|
||||
stale = append(stale, row.ID)
|
||||
}
|
||||
return s.deleteMediaByIDs(ctx, stale, false)
|
||||
return s.deleteMediaByIDs(ctx, stale, true)
|
||||
}
|
||||
|
||||
func pathBelongsToRoot(pathValue, rootPath string) bool {
|
||||
|
||||
@@ -501,6 +501,16 @@ func (s *StrmService) ClearFinishedDownloadTasks(ctx context.Context) (int64, er
|
||||
return s.repo.StrmDownload.ClearFinished(ctx)
|
||||
}
|
||||
|
||||
// ClearCanceledDownloadTasks 清空全部已取消的下载记录,返回删除数量。
|
||||
func (s *StrmService) ClearCanceledDownloadTasks(ctx context.Context) (int64, error) {
|
||||
return s.repo.StrmDownload.ClearCanceled(ctx)
|
||||
}
|
||||
|
||||
// ClearCanceledUploadTasks 清空全部已取消的上传记录,返回删除数量。
|
||||
func (s *StrmService) ClearCanceledUploadTasks(ctx context.Context) (int64, error) {
|
||||
return s.repo.StrmUpload.ClearCanceled(ctx)
|
||||
}
|
||||
|
||||
// RetryAllFailedDownloadTasks 批量重试所有失败下载任务,返回重新入队数量。
|
||||
func (s *StrmService) RetryAllFailedDownloadTasks(ctx context.Context) (int64, error) {
|
||||
return s.repo.StrmDownload.RetryAllFailed(ctx)
|
||||
|
||||
@@ -113,6 +113,7 @@ func NewStrmService(cfg *config.Config, log *zap.Logger, repos *repository.Conta
|
||||
// Start 启动下载/上传队列 worker、定时同步巡检、115 token 刷新与队列清理。
|
||||
func (s *StrmService) Start(ctx context.Context) {
|
||||
s.sync115RelayKey(ctx)
|
||||
s.recoverInterruptedSyncs(ctx)
|
||||
downloadThreads := s.strmIntSetting(ctx, StrmSettingDownloadThreads, 3)
|
||||
if downloadThreads < 1 {
|
||||
downloadThreads = 1
|
||||
@@ -141,6 +142,21 @@ func (s *StrmService) Start(ctx context.Context) {
|
||||
zap.Int("upload_threads", uploadThreads))
|
||||
}
|
||||
|
||||
// recoverInterruptedSyncs 在服务启动时自愈重置因服务重启遗留的 running 状态。
|
||||
func (s *StrmService) recoverInterruptedSyncs(ctx context.Context) {
|
||||
paths, err := s.repo.StrmSyncPath.List(ctx)
|
||||
if err == nil {
|
||||
for i := range paths {
|
||||
p := &paths[i]
|
||||
if p.LastSyncStatus == model.StrmSyncRecordRunning {
|
||||
p.LastSyncStatus = model.StrmSyncRecordCanceled
|
||||
p.LastSyncMessage = "服务重启,已重置同步状态"
|
||||
_ = s.repo.StrmSyncPath.Update(ctx, p)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *StrmService) Stop() {
|
||||
s.stopOnce.Do(func() { close(s.stopCh) })
|
||||
}
|
||||
|
||||
+123
-47
@@ -40,8 +40,10 @@ type strmSyncState struct {
|
||||
seenVideo map[string]bool // "v:"+去掉扩展名的相对路径 → 远端存在该视频
|
||||
seenMeta map[string]bool // "m:"+相对路径 → 远端存在该元数据
|
||||
remoteMeta map[string]int64 // 远端元数据大小(上传比对用)
|
||||
activeDownloadPaths map[string]bool // 本地已在排队/进行的下载任务路径(内存去重)
|
||||
activeUploadPaths map[string]bool // 本地已在排队/进行的上传任务路径(内存去重)
|
||||
seenMetaTarget map[string]cloud.FileEntry
|
||||
seenVideoTarget map[string]cloud.FileEntry
|
||||
activeDownloadPaths map[string]bool // 本地已在排队/进行的下载任务路径(内存去重)
|
||||
activeUploadPaths map[string]bool // 本地已在排队/进行的上传任务路径(内存去重)
|
||||
pendingDownloads []*model.StrmDownloadTask
|
||||
pendingUploads []*model.StrmUploadTask
|
||||
dirCache sync.Map // dirID (string) -> relativePath (string)
|
||||
@@ -104,15 +106,27 @@ func (s *StrmService) StartSync(ctx context.Context, pathID string, syncType ...
|
||||
return nil
|
||||
}
|
||||
|
||||
// CancelSync 取消正在进行的同步。
|
||||
// CancelSync 取消正在进行的同步(若为僵尸运行状态则直接自愈重置)。
|
||||
func (s *StrmService) CancelSync(ctx context.Context, pathID string) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
cancel, exists := s.running[pathID]
|
||||
if !exists {
|
||||
return errors.New("该目录当前没有进行中的同步")
|
||||
if exists {
|
||||
delete(s.running, pathID)
|
||||
}
|
||||
s.mu.Unlock()
|
||||
|
||||
if exists && cancel != nil {
|
||||
cancel()
|
||||
}
|
||||
|
||||
// 无论内存中是否活跃,确保同步目录状态正确重置为已取消
|
||||
if p, err := s.repo.StrmSyncPath.FindByID(ctx, pathID); err == nil && p != nil {
|
||||
if p.LastSyncStatus == model.StrmSyncRecordRunning {
|
||||
p.LastSyncStatus = model.StrmSyncRecordCanceled
|
||||
p.LastSyncMessage = "已取消"
|
||||
_ = s.repo.StrmSyncPath.Update(ctx, p)
|
||||
}
|
||||
}
|
||||
cancel()
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -160,15 +174,17 @@ func (s *StrmService) runSync(ctx context.Context, p *model.StrmSyncPath, rec *m
|
||||
return
|
||||
}
|
||||
st := &strmSyncState{
|
||||
s: s,
|
||||
ctx: ctx,
|
||||
p: p,
|
||||
cfg: cfg,
|
||||
rec: rec,
|
||||
syncType: rec.SyncType,
|
||||
seenVideo: map[string]bool{},
|
||||
seenMeta: map[string]bool{},
|
||||
remoteMeta: map[string]int64{},
|
||||
s: s,
|
||||
ctx: ctx,
|
||||
p: p,
|
||||
cfg: cfg,
|
||||
rec: rec,
|
||||
syncType: rec.SyncType,
|
||||
seenVideo: map[string]bool{},
|
||||
seenMeta: map[string]bool{},
|
||||
remoteMeta: map[string]int64{},
|
||||
seenMetaTarget: map[string]cloud.FileEntry{},
|
||||
seenVideoTarget: map[string]cloud.FileEntry{},
|
||||
}
|
||||
if p.Provider != model.StrmProviderLocal {
|
||||
acct, err := s.repo.StrmAccount.FindByID(ctx, p.AccountID)
|
||||
@@ -440,15 +456,23 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error {
|
||||
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)
|
||||
} else {
|
||||
// 增量同步:预加载历史目录缓存(过滤历史一对多塌陷冲突的脏数据以自愈刷新)
|
||||
cached, err := st.s.repo.StrmDirCache.ListBySyncPathID(ctx, st.p.ID)
|
||||
if err == nil {
|
||||
pathCounts := make(map[string]int, len(cached))
|
||||
for _, item := range cached {
|
||||
pathCounts[item.Path]++
|
||||
}
|
||||
for _, item := range cached {
|
||||
// 若同一个 path 对应了多个不同 dir_id,说明包含历史层级塌陷的脏数据,不预加载,让后续步骤重新向 115 获取精确路径
|
||||
if pathCounts[item.Path] > 1 {
|
||||
continue
|
||||
}
|
||||
st.dirCache.Store(item.DirID, item.Path)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 2. 探测文件总数
|
||||
const pageSize = 1150
|
||||
@@ -457,6 +481,8 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error {
|
||||
return fmt.Errorf("115: 获取文件列表失败:%w", err)
|
||||
}
|
||||
|
||||
st.updateSyncMessage(fmt.Sprintf("正在拉取远端文件列表 (共 %d 个文件)...", totalCount))
|
||||
|
||||
allFiles := make([]cloud115.RemoteFile, 0, totalCount)
|
||||
allFiles = append(allFiles, firstBatch...)
|
||||
|
||||
@@ -484,7 +510,7 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error {
|
||||
}
|
||||
close(taskCh)
|
||||
|
||||
workers := 4
|
||||
workers := 8
|
||||
if len(pageTasks) < workers {
|
||||
workers = len(pageTasks)
|
||||
}
|
||||
@@ -549,12 +575,16 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error {
|
||||
|
||||
var (
|
||||
pwg sync.WaitGroup
|
||||
dirWorkers = 4
|
||||
dirWorkers = 8
|
||||
doneDirs atomic.Int64
|
||||
totalDirs = len(pidList)
|
||||
)
|
||||
if len(pidList) < dirWorkers {
|
||||
dirWorkers = len(pidList)
|
||||
}
|
||||
|
||||
st.updateSyncMessage(fmt.Sprintf("正在解析目录树 (0/%d)...", totalDirs))
|
||||
|
||||
for i := 0; i < dirWorkers; i++ {
|
||||
pwg.Add(1)
|
||||
go func() {
|
||||
@@ -563,47 +593,55 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error {
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
if _, loaded := st.dirCache.Load(pid); loaded {
|
||||
if n := doneDirs.Add(1); n%20 == 0 || n == int64(totalDirs) {
|
||||
st.updateSyncMessage(fmt.Sprintf("正在解析目录树 (%d/%d)...", n, totalDirs))
|
||||
}
|
||||
continue
|
||||
}
|
||||
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)
|
||||
} else if detail != nil {
|
||||
// 解析相对路径
|
||||
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 {
|
||||
// 顺便解析并缓存 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,
|
||||
FileName: ancestor.Name,
|
||||
Paths: nil,
|
||||
}
|
||||
for _, p := range detail.Paths {
|
||||
subDetail.Paths = append(subDetail.Paths, p)
|
||||
if p.FileId == ancestor.FileId {
|
||||
break
|
||||
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)
|
||||
}
|
||||
ancestorRel := subDetail.RelativePath(rootCID)
|
||||
st.dirCache.Store(ancestor.FileId, ancestorRel)
|
||||
_ = st.s.repo.StrmDirCache.Set(ctx, st.p.ID, ancestor.FileId, ancestorRel)
|
||||
}
|
||||
}
|
||||
if n := doneDirs.Add(1); n%10 == 0 || n == int64(totalDirs) {
|
||||
st.updateSyncMessage(fmt.Sprintf("正在解析目录树 (%d/%d)...", n, totalDirs))
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
pwg.Wait()
|
||||
}
|
||||
|
||||
st.updateSyncMessage(fmt.Sprintf("正在生成 STRM 与同步文件 (共 %d 个)...", len(allFiles)))
|
||||
|
||||
// 5. 分类处理所有文件
|
||||
for _, f := range allFiles {
|
||||
if ctx.Err() != nil {
|
||||
@@ -652,6 +690,18 @@ func (st *strmSyncState) handleVideo(entry cloud.FileEntry, rel, ext string) {
|
||||
return
|
||||
}
|
||||
|
||||
st.mu.Lock()
|
||||
if st.seenVideoTarget == nil {
|
||||
st.seenVideoTarget = map[string]cloud.FileEntry{}
|
||||
}
|
||||
if _, exists := st.seenVideoTarget[target]; exists {
|
||||
st.mu.Unlock()
|
||||
st.touchProgress()
|
||||
return
|
||||
}
|
||||
st.seenVideoTarget[target] = entry
|
||||
st.mu.Unlock()
|
||||
|
||||
// 增量同步模式快速检查:本地 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 {
|
||||
@@ -801,6 +851,20 @@ func (st *strmSyncState) handleMeta(entry cloud.FileEntry, rel, ext string) {
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
st.mu.Lock()
|
||||
if st.seenMetaTarget == nil {
|
||||
st.seenMetaTarget = map[string]cloud.FileEntry{}
|
||||
}
|
||||
if _, exists := st.seenMetaTarget[target]; exists {
|
||||
// 该本地目标路径在当前批次中已被处理(存在同名/重名冲突),直接忽略重复项,避免多份不同大小的文件在本地交替覆盖导致增量死循环
|
||||
st.mu.Unlock()
|
||||
st.touchProgress()
|
||||
return
|
||||
}
|
||||
st.seenMetaTarget[target] = entry
|
||||
st.mu.Unlock()
|
||||
|
||||
if info, err := os.Stat(target); err == nil && info.Size() == entry.Size {
|
||||
st.touchProgress()
|
||||
return
|
||||
@@ -1130,6 +1194,18 @@ func (st *strmSyncState) flushProgress() {
|
||||
}
|
||||
}
|
||||
|
||||
// updateSyncMessage 实时更新同步阶段提示信息,让前端界面清晰了解当前进度。
|
||||
func (st *strmSyncState) updateSyncMessage(msg string) {
|
||||
st.mu.Lock()
|
||||
st.rec.Message = msg
|
||||
st.p.LastSyncMessage = msg
|
||||
rec := *st.rec
|
||||
p := *st.p
|
||||
st.mu.Unlock()
|
||||
_ = st.s.repo.StrmSyncRecord.Update(st.ctx, &rec)
|
||||
_ = st.s.repo.StrmSyncPath.Update(st.ctx, &p)
|
||||
}
|
||||
|
||||
// ─── 定时同步巡检 ──────────────────────────────────────────────────────────────
|
||||
|
||||
func (s *StrmService) cronLoop(ctx context.Context) {
|
||||
|
||||
@@ -554,8 +554,82 @@ func TestWalkRemoteConcurrent(t *testing.T) {
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
if claimedCount != 200 {
|
||||
t.Fatalf("expected all 200 tasks claimed, got %d", claimedCount)
|
||||
if claimedCount != 200 {
|
||||
t.Fatalf("expected all 200 tasks claimed, got %d", claimedCount)
|
||||
}
|
||||
}
|
||||
|
||||
// TestStrmDuplicateFileConflictResolution 测试远端存在多个同名不同大小文件时,本地确定性仲裁,避免增量死循环
|
||||
func TestStrmDuplicateFileConflictResolution(t *testing.T) {
|
||||
svc := testStrmService(t)
|
||||
localDir := t.TempDir()
|
||||
|
||||
p := &model.StrmSyncPath{
|
||||
Base: model.Base{ID: "dup-test-path"},
|
||||
Provider: model.StrmProvider115,
|
||||
RemotePath: "root",
|
||||
LocalPath: localDir,
|
||||
DownloadMeta: true,
|
||||
}
|
||||
|
||||
st := &strmSyncState{
|
||||
s: svc,
|
||||
ctx: context.Background(),
|
||||
p: p,
|
||||
cfg: &strmPathConfig{DownloadMeta: true, MetaExt: []string{"nfo"}},
|
||||
rec: &model.StrmSyncRecord{},
|
||||
seenMeta: map[string]bool{},
|
||||
remoteMeta: map[string]int64{},
|
||||
seenMetaTarget: map[string]cloud.FileEntry{},
|
||||
seenVideoTarget: map[string]cloud.FileEntry{},
|
||||
}
|
||||
|
||||
// 模拟远端同目录下存在两个同名不同大小的 nfo 文件 (115 历史重复上传)
|
||||
// entry1: 较早文件 (MTime: 1000, Size: 100)
|
||||
entry1 := cloud.FileEntry{ID: "f1", Name: "test.nfo", Size: 100, MTime: 1000, PickCode: "p1"}
|
||||
// entry2: 较新文件 (MTime: 2000, Size: 200)
|
||||
entry2 := cloud.FileEntry{ID: "f2", Name: "test.nfo", Size: 200, MTime: 2000, PickCode: "p2"}
|
||||
|
||||
// 第一次全量处理:两者都在列表中
|
||||
st.handleMeta(entry1, "test.nfo", ".nfo")
|
||||
st.handleMeta(entry2, "test.nfo", ".nfo")
|
||||
st.flushPendingDownloads()
|
||||
|
||||
// 验证仲裁结果:最终只产生 1 个下载任务,且使用的是首个匹配项 (Size 100/p1)
|
||||
tasks, _, err := svc.repo.StrmDownload.List(context.Background(), "", 1, 10)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(tasks) != 1 {
|
||||
t.Fatalf("expected 1 download task after conflict resolution, got %d", len(tasks))
|
||||
}
|
||||
if tasks[0].Size != 100 || tasks[0].RemoteRef != "p1" {
|
||||
t.Fatalf("expected task with size 100/p1, got size=%d ref=%s", tasks[0].Size, tasks[0].RemoteRef)
|
||||
}
|
||||
|
||||
// 模拟该任务下载落盘完成
|
||||
writeFile(t, filepath.Join(localDir, "test.nfo"), strings.Repeat("x", 100))
|
||||
|
||||
// 第二次增量同步:两者再次依次扫描
|
||||
st2 := &strmSyncState{
|
||||
s: svc,
|
||||
ctx: context.Background(),
|
||||
p: p,
|
||||
cfg: &strmPathConfig{DownloadMeta: true, MetaExt: []string{"nfo"}},
|
||||
rec: &model.StrmSyncRecord{},
|
||||
seenMeta: map[string]bool{},
|
||||
remoteMeta: map[string]int64{},
|
||||
seenMetaTarget: map[string]cloud.FileEntry{},
|
||||
seenVideoTarget: map[string]cloud.FileEntry{},
|
||||
}
|
||||
st2.handleMeta(entry1, "test.nfo", ".nfo")
|
||||
st2.handleMeta(entry2, "test.nfo", ".nfo")
|
||||
st2.flushPendingDownloads()
|
||||
|
||||
// 验证:不会新增任何下载任务,NewMeta 为 0,增量跳过
|
||||
if st2.rec.NewMeta != 0 {
|
||||
t.Fatalf("expected 0 new meta on incremental sync, got %d", st2.rec.NewMeta)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -143,6 +143,9 @@ export const strmAPI = {
|
||||
clearFinishedDownloads: () =>
|
||||
api.post<{ deleted: number }>('/admin/strm/downloads/clear-finished').then((r) => r.data),
|
||||
|
||||
clearCanceledDownloads: () =>
|
||||
api.post<{ deleted: number }>('/admin/strm/downloads/clear-canceled').then((r) => r.data),
|
||||
|
||||
retryFailedDownloads: () =>
|
||||
api.post<{ retried: number }>('/admin/strm/downloads/retry-failed').then((r) => r.data),
|
||||
|
||||
@@ -162,6 +165,9 @@ export const strmAPI = {
|
||||
cancelPendingUploads: () =>
|
||||
api.post<{ canceled: number }>('/admin/strm/uploads/cancel-pending').then((r) => r.data),
|
||||
|
||||
clearCanceledUploads: () =>
|
||||
api.post<{ deleted: number }>('/admin/strm/uploads/clear-canceled').then((r) => r.data),
|
||||
|
||||
retryUpload: (id: string) =>
|
||||
api.post(`/admin/strm/uploads/${id}/retry`).then((r) => r.data),
|
||||
}
|
||||
@@ -127,6 +127,19 @@ function StrmQueuePanel({ kind }: { kind: 'download' | 'upload' }) {
|
||||
'border-amber-200 text-amber-600 hover:bg-amber-50',
|
||||
cancelAllPendingAction,
|
||||
)
|
||||
if (filter === 'canceled')
|
||||
return batchBtn(
|
||||
'清空已取消记录',
|
||||
'trash',
|
||||
'border-gray-200 text-rose-500 hover:bg-rose-50',
|
||||
() =>
|
||||
runBatch(
|
||||
isDownload
|
||||
? () => strmAPI.clearCanceledDownloads()
|
||||
: () => strmAPI.clearCanceledUploads(),
|
||||
`确定清空所有已取消的${isDownload ? '下载' : '上传'}记录?`,
|
||||
),
|
||||
)
|
||||
if (!isDownload) return null
|
||||
if (filter === 'done')
|
||||
return batchBtn(
|
||||
|
||||
Reference in New Issue
Block a user