mirror of
https://github.com/truewhile/MeBox.git
synced 2026-09-28 11:16:37 +08:00
Compare commits
12 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 60c815a8b3 | |||
| 41b155ea31 | |||
| 2888ae8bf7 | |||
| 7363064d89 | |||
| 1d53bf2ae1 | |||
| 618165ec31 | |||
| 87c66a9b8c | |||
| 1ea4724261 | |||
| 3d372f039e | |||
| 0384017e98 | |||
| 0332579d5f | |||
| 6aefe18caa |
@@ -16,12 +16,13 @@ import (
|
||||
)
|
||||
|
||||
type createLibraryReq struct {
|
||||
Name string `json:"name" binding:"required"`
|
||||
Path string `json:"path"`
|
||||
Paths []string `json:"paths"`
|
||||
Roots []service.LibraryRootInput `json:"roots"`
|
||||
Type string `json:"type"`
|
||||
CoverURL string `json:"cover_url"`
|
||||
Name string `json:"name"`
|
||||
Path string `json:"path"`
|
||||
Paths []string `json:"paths"`
|
||||
Roots []service.LibraryRootInput `json:"roots"`
|
||||
Type string `json:"type"`
|
||||
CoverURL string `json:"cover_url"`
|
||||
CreatePerSubfolder bool `json:"create_per_subfolder"`
|
||||
}
|
||||
|
||||
func listLibrariesHandler(svc *service.Container) gin.HandlerFunc {
|
||||
@@ -88,9 +89,38 @@ func createLibraryHandler(svc *service.Container) gin.HandlerFunc {
|
||||
}
|
||||
}
|
||||
if len(roots) == 0 && strings.TrimSpace(req.Path) != "" {
|
||||
roots = append(roots, service.LibraryRootInput{Path: req.Path})
|
||||
roots = append(roots, service.LibraryRootInput{Path: req.Path})
|
||||
}
|
||||
var l *model.Library
|
||||
if req.CreatePerSubfolder {
|
||||
parent := ""
|
||||
if len(roots) > 0 {
|
||||
parent = roots[0].Path
|
||||
} else if strings.TrimSpace(req.Path) != "" {
|
||||
parent = req.Path
|
||||
}
|
||||
l, err := svc.Media.CreateLibraryWithRootsAndCover(c.Request.Context(), req.Name, req.Type, req.CoverURL, roots)
|
||||
created, err := svc.Media.CreateLibrariesPerSubfolder(c.Request.Context(), parent, req.Type, req.CoverURL)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
uid, _ := c.Get("ctx_user_id")
|
||||
for i := range created {
|
||||
lib := &created[i]
|
||||
svc.Audit.Record(c.Request.Context(), toString(uid), "library.create", lib.ID, c.ClientIP(), lib.Path)
|
||||
if svc.Watcher != nil {
|
||||
go func() { _ = svc.Watcher.Refresh(context.Background()) }()
|
||||
}
|
||||
for _, root := range lib.Roots {
|
||||
if root.Enabled {
|
||||
queueLibraryRootScan(svc, lib.ID, root.ID)
|
||||
}
|
||||
}
|
||||
}
|
||||
c.JSON(http.StatusCreated, gin.H{"libraries": created})
|
||||
return
|
||||
}
|
||||
l, err := svc.Media.CreateLibraryWithRootsAndCover(c.Request.Context(), req.Name, req.Type, req.CoverURL, roots)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
|
||||
@@ -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
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
@@ -26,10 +26,15 @@ var (
|
||||
executorOnce sync.Once
|
||||
)
|
||||
|
||||
// GetGlobalExecutor 获取全局队列执行器单例(默认 QPS=2, QPM=120, QPH=6000,保障 115 API 调用安全不超频)。
|
||||
// GetGlobalExecutor 获取全局队列执行器单例(默认 QPS=3, QPM=200, QPH=12000,保障 115 API 调用安全不超频)。
|
||||
//
|
||||
// 历史教训:QPS 提到 8 后,下载换直链接口(/open/ufile/downurl,WAF 重点盯防对象)
|
||||
// 瞬时突发撞上 115 风控,返回阿里云 405 阻断页(HTTP 405),导致全量同步失败。
|
||||
// 因此回调到 3——这是经过实测的安全上限:宁慢勿触发风控,一旦 405 冷却 180 秒,
|
||||
// 整体吞吐反而更低。下载实际走 CDN 不受此限速影响,瓶颈仅在换链环节。
|
||||
func GetGlobalExecutor() *QueueExecutor {
|
||||
executorOnce.Do(func() {
|
||||
globalExecutor = NewQueueExecutor(2, 120, 6000)
|
||||
globalExecutor = NewQueueExecutor(3, 200, 12000)
|
||||
})
|
||||
return globalExecutor
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
|
||||
@@ -59,6 +60,46 @@ func (s *MediaService) CreateLibraryWithRootsAndCover(ctx context.Context, name,
|
||||
return lib, nil
|
||||
}
|
||||
|
||||
// CreateLibrariesPerSubfolder 为 parent 目录下的每个直接子目录各建一个媒体库,
|
||||
// 媒体库名取子目录名,路径指向该子目录。kind 为空时按子目录名推断类型。
|
||||
func (s *MediaService) CreateLibrariesPerSubfolder(ctx context.Context, parent, kind, coverURL string) ([]model.Library, error) {
|
||||
parent = strings.TrimSpace(parent)
|
||||
if parent == "" {
|
||||
return nil, errors.New("parent path required")
|
||||
}
|
||||
dir, err := resolveAccessibleLibraryPath(parent)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
entries, err := os.ReadDir(dir)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read directory failed: %w", err)
|
||||
}
|
||||
subdirs := make([]string, 0, len(entries))
|
||||
for _, entry := range entries {
|
||||
if !entry.IsDir() {
|
||||
continue
|
||||
}
|
||||
if strings.HasPrefix(entry.Name(), ".") {
|
||||
continue
|
||||
}
|
||||
subdirs = append(subdirs, filepath.Join(dir, entry.Name()))
|
||||
}
|
||||
if len(subdirs) == 0 {
|
||||
return nil, errors.New("no subfolders found")
|
||||
}
|
||||
created := make([]model.Library, 0, len(subdirs))
|
||||
for _, subdir := range subdirs {
|
||||
name := filepath.Base(subdir)
|
||||
lib, err := s.CreateLibraryWithRootsAndCover(ctx, name, kind, coverURL, []LibraryRootInput{{Path: subdir}})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create library for %s: %w", subdir, err)
|
||||
}
|
||||
created = append(created, *lib)
|
||||
}
|
||||
return created, nil
|
||||
}
|
||||
|
||||
func (s *MediaService) UpdateLibraryCover(ctx context.Context, libraryID, coverURL string) error {
|
||||
return s.repo.DB.WithContext(ctx).Model(&model.Library{}).Where("id = ?", libraryID).
|
||||
Update("cover_url", strings.TrimSpace(coverURL)).Error
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
@@ -27,7 +28,12 @@ const (
|
||||
)
|
||||
|
||||
// downloadWorker 下载队列 worker:认领 → 解析直链 → 下载 → 落盘。
|
||||
//
|
||||
// 采用「批量认领 + 全局并发限流」:一次认领数个任务,用 StrmService 上的全局信号量
|
||||
// 限制整个进程「同时换直链+下载」的并发数(与 115 换链风控匹配,见 strmDownloadSemCap),
|
||||
// 同时让下载充分并行。换链走全局令牌桶(QPS=3)兜底,下载走 CDN 不限速。
|
||||
func (s *StrmService) downloadWorker(ctx context.Context) {
|
||||
const claimBatch = 12 // 每次批量认领的任务数
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
@@ -42,7 +48,7 @@ func (s *StrmService) downloadWorker(ctx context.Context) {
|
||||
sleepContext(ctx, left)
|
||||
continue
|
||||
}
|
||||
tasks, err := s.repo.StrmDownload.ClaimPendingDownload(ctx, 1)
|
||||
tasks, err := s.repo.StrmDownload.ClaimPendingDownload(ctx, claimBatch)
|
||||
if err != nil {
|
||||
s.log.Warn("claim strm download task failed", zap.Error(err))
|
||||
sleepContext(ctx, 3*time.Second)
|
||||
@@ -52,9 +58,21 @@ func (s *StrmService) downloadWorker(ctx context.Context) {
|
||||
sleepContext(ctx, 2*time.Second)
|
||||
continue
|
||||
}
|
||||
// 并发处理本批任务:每个任务先获取全局下载槽位,槽位内部执行换链+下载。
|
||||
// 信号量与令牌桶双重限速,确保任意时刻并发换链请求不超过安全阈值。
|
||||
var wg sync.WaitGroup
|
||||
for i := range tasks {
|
||||
s.processDownloadTask(ctx, &tasks[i])
|
||||
wg.Add(1)
|
||||
go func(i int) {
|
||||
defer wg.Done()
|
||||
if !s.acquireDownloadSlot(ctx) {
|
||||
return
|
||||
}
|
||||
defer s.releaseDownloadSlot()
|
||||
s.processDownloadTask(ctx, &tasks[i])
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -501,6 +519,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)
|
||||
|
||||
@@ -90,11 +90,48 @@ type StrmService struct {
|
||||
running map[string]context.CancelFunc // sync path id -> cancel
|
||||
oauthSessions map[string]*strm115AuthSession
|
||||
wafUntil time.Time // 115 风控/限流熔断截止时间(由 mu 保护)
|
||||
|
||||
downloadSem chan struct{} // 全局下载并发信号量:限制整个进程同时进行「换直链+下载」的并发数
|
||||
downloadSemOnce sync.Once
|
||||
}
|
||||
|
||||
// strmWAFCooldown 检测到 115 风控/限流后下载队列的全局冷却时长。
|
||||
const strmWAFCooldown = 3 * time.Minute
|
||||
|
||||
// strmDownloadSemCap 全局同时进行「换直链+下载」的并发上限。
|
||||
//
|
||||
// 115 对换直链接口(/open/ufile/downurl)风控极严:过去把全局 QPS 提到 8 或让多
|
||||
// worker 高并发换链,会瞬时撞上 WAF 返回 405 阻断页并触发 180 秒冷却,反而更慢。
|
||||
// 因此用信号量把整个进程同时换直链的并发数压到 3,与令牌桶限速共同兜底:
|
||||
// 宁可下载稍慢,也绝不触发风控。下载本身走 CDN 不限速。
|
||||
const strmDownloadSemCap = 3
|
||||
|
||||
// ensureDownloadSem 惰性初始化全局共享的下载并发信号量。
|
||||
func (s *StrmService) ensureDownloadSem() {
|
||||
s.downloadSemOnce.Do(func() {
|
||||
s.downloadSem = make(chan struct{}, strmDownloadSemCap)
|
||||
})
|
||||
}
|
||||
|
||||
// acquireDownloadSlot 获取一个下载并发槽位(等待/取消安全)。
|
||||
func (s *StrmService) acquireDownloadSlot(ctx context.Context) bool {
|
||||
s.ensureDownloadSem()
|
||||
select {
|
||||
case s.downloadSem <- struct{}{}:
|
||||
return true
|
||||
case <-ctx.Done():
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// releaseDownloadSlot 释放一个下载并发槽位。
|
||||
func (s *StrmService) releaseDownloadSlot() {
|
||||
if s.downloadSem == nil {
|
||||
return
|
||||
}
|
||||
<-s.downloadSem
|
||||
}
|
||||
|
||||
// NewStrmService constructs the STRM service.
|
||||
func NewStrmService(cfg *config.Config, log *zap.Logger, repos *repository.Container, crypto *CryptoService) *StrmService {
|
||||
return &StrmService{
|
||||
|
||||
@@ -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)
|
||||
@@ -172,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)
|
||||
@@ -686,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 {
|
||||
@@ -835,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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
-25
@@ -1,25 +0,0 @@
|
||||
package main
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"github.com/ShukeBta/MMTL/internal/service"
|
||||
"github.com/ShukeBta/MMTL/internal/service/cloud115"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
func main() {
|
||||
crypto := service.NewCryptoService("test-secret", zap.NewNop())
|
||||
// Let's test with a mock RemoteFileDetail
|
||||
d := &cloud115.RemoteFileDetail{
|
||||
FileId: "3251154147730910635",
|
||||
FileName: "出包王女",
|
||||
Paths: []struct {
|
||||
FileId string
|
||||
Name string
|
||||
}{
|
||||
{FileId: "0", Name: "根目录"},
|
||||
{FileId: "3238787832374488117", Name: "影视库"},
|
||||
{FileId: "3238787913223892116", Name: "动漫"},
|
||||
},
|
||||
}
|
||||
fmt.Println("RelativePath when rootCID is 3238787832374488117:", d.RelativePath("3238787832374488117"))
|
||||
}
|
||||
@@ -102,6 +102,9 @@ export const libraryAPI = {
|
||||
createWithRoots: (name: string, type: string, roots: LibraryRootInput[], coverURL = '') =>
|
||||
api.post<Library>('/libraries', { name, type, roots, cover_url: coverURL }).then((r) => r.data),
|
||||
|
||||
createPerSubfolder: (parentPath: string, type: string, coverURL = '') =>
|
||||
api.post<{ libraries: Library[] }>('/libraries', { path: parentPath, type, cover_url: coverURL, create_per_subfolder: true }).then((r) => r.data),
|
||||
|
||||
update: (id: string, payload: { cover_url: string }) =>
|
||||
api.patch<Library>(`/libraries/${id}`, payload).then((r) => r.data),
|
||||
|
||||
|
||||
@@ -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),
|
||||
}
|
||||
@@ -13,9 +13,11 @@ export function AdminLibraryPanel() {
|
||||
type={createForm.type}
|
||||
coverURL={createForm.coverURL}
|
||||
roots={createForm.roots}
|
||||
createPerSubfolder={createForm.createPerSubfolder}
|
||||
onNameChange={createForm.setName}
|
||||
onTypeChange={createForm.setType}
|
||||
onCoverURLChange={createForm.setCoverURL}
|
||||
onCreatePerSubfolderChange={createForm.setCreatePerSubfolder}
|
||||
onRootChange={createForm.updateRoot}
|
||||
onAddRoot={createForm.addRoot}
|
||||
onRemoveRoot={createForm.removeRoot}
|
||||
|
||||
@@ -9,9 +9,11 @@ type CreateFormProps = {
|
||||
type: string
|
||||
coverURL: string
|
||||
roots: RootDraft[]
|
||||
createPerSubfolder: boolean
|
||||
onNameChange: (value: string) => void
|
||||
onTypeChange: (value: string) => void
|
||||
onCoverURLChange: (value: string) => void
|
||||
onCreatePerSubfolderChange: (value: boolean) => void
|
||||
onRootChange: (index: number, patch: Partial<RootDraft>) => void
|
||||
onAddRoot: () => void
|
||||
onRemoveRoot: (index: number) => void
|
||||
@@ -23,9 +25,11 @@ export function AdminLibraryCreateForm({
|
||||
type,
|
||||
coverURL,
|
||||
roots,
|
||||
createPerSubfolder,
|
||||
onNameChange,
|
||||
onTypeChange,
|
||||
onCoverURLChange,
|
||||
onCreatePerSubfolderChange,
|
||||
onRootChange,
|
||||
onAddRoot,
|
||||
onRemoveRoot,
|
||||
@@ -50,9 +54,9 @@ export function AdminLibraryCreateForm({
|
||||
<>
|
||||
<form onSubmit={onSubmit} className="glass-panel grid gap-3 md:grid-cols-4">
|
||||
<input
|
||||
required
|
||||
required={!createPerSubfolder}
|
||||
className="input-base"
|
||||
placeholder="名称"
|
||||
placeholder={createPerSubfolder ? '父级媒体库名(批量模式忽略)' : '名称'}
|
||||
value={name}
|
||||
onChange={(e) => onNameChange(e.target.value)}
|
||||
/>
|
||||
@@ -81,15 +85,31 @@ export function AdminLibraryCreateForm({
|
||||
onRemove={onRemoveRoot}
|
||||
/>
|
||||
))}
|
||||
<button type="button" className="inline-flex items-center gap-2 rounded-lg border px-3 py-2 text-sm" onClick={onAddRoot}>
|
||||
<Plus size={16} /> 添加路径
|
||||
</button>
|
||||
{!createPerSubfolder && (
|
||||
<button type="button" className="inline-flex items-center gap-2 rounded-lg border px-3 py-2 text-sm" onClick={onAddRoot}>
|
||||
<Plus size={16} /> 添加路径
|
||||
</button>
|
||||
)}
|
||||
</div>
|
||||
<p className="md:col-span-4 -mt-2 text-xs text-sand-500">
|
||||
支持直接点选或手动输入;名称和类型与现有媒体库一致时,会自动把这里填写的路径追加到该媒体库。
|
||||
</p>
|
||||
<label className="md:col-span-4 flex items-center gap-2 text-sm text-ink-100">
|
||||
<input
|
||||
type="checkbox"
|
||||
className="h-4 w-4 accent-brand-400"
|
||||
checked={createPerSubfolder}
|
||||
onChange={(e) => onCreatePerSubfolderChange(e.target.checked)}
|
||||
/>
|
||||
<span>按目录下每个子文件夹各建一个媒体库(媒体库名取子文件夹名)</span>
|
||||
</label>
|
||||
{createPerSubfolder && (
|
||||
<p className="md:col-span-4 -mt-2 text-xs text-sand-500">
|
||||
批处理模式:仅取上方第一个路径作为父级目录,会为其中每个子文件夹分别创建媒体库,可自选类型用于整体推断。
|
||||
</p>
|
||||
)}
|
||||
<button type="submit" className="neon-button md:col-span-4">
|
||||
新建 / 追加路径
|
||||
{createPerSubfolder ? '按目录批量创建' : '新建 / 追加路径'}
|
||||
</button>
|
||||
</form>
|
||||
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -32,20 +32,32 @@ function useCreateLibraryForm(refresh: () => Promise<void>) {
|
||||
const [roots, setRoots] = useState<RootDraft[]>([emptyRootDraft()])
|
||||
const [type, setType] = useState('movie')
|
||||
const [coverURL, setCoverURL] = useState('')
|
||||
const [createPerSubfolder, setCreatePerSubfolder] = useState(false)
|
||||
|
||||
const handleCreate = async (e: FormEvent) => {
|
||||
e.preventDefault()
|
||||
try {
|
||||
const payload = createRootPayload(roots)
|
||||
if (payload.length === 0) {
|
||||
toast.error('请至少填写一个路径')
|
||||
return
|
||||
if (createPerSubfolder) {
|
||||
const parentPath = roots[0]?.path?.trim()
|
||||
if (!parentPath) {
|
||||
toast.error('请先选择或填写父级目录')
|
||||
return
|
||||
}
|
||||
const { libraries } = await libraryAPI.createPerSubfolder(parentPath, type, coverURL.trim())
|
||||
toast.success(`已按目录创建 ${libraries.length} 个媒体库`)
|
||||
} else {
|
||||
const payload = createRootPayload(roots)
|
||||
if (payload.length === 0) {
|
||||
toast.error('请至少填写一个路径')
|
||||
return
|
||||
}
|
||||
await libraryAPI.createWithRoots(name, type, payload, coverURL.trim())
|
||||
toast.success('媒体库已保存')
|
||||
}
|
||||
await libraryAPI.createWithRoots(name, type, payload, coverURL.trim())
|
||||
toast.success('媒体库已保存')
|
||||
setName('')
|
||||
setRoots([emptyRootDraft()])
|
||||
setCoverURL('')
|
||||
setCreatePerSubfolder(false)
|
||||
await refresh()
|
||||
} catch (err: unknown) {
|
||||
toast.error(apiErrorMessage(err, '创建失败'))
|
||||
@@ -61,9 +73,11 @@ function useCreateLibraryForm(refresh: () => Promise<void>) {
|
||||
type,
|
||||
coverURL,
|
||||
roots,
|
||||
createPerSubfolder,
|
||||
setName,
|
||||
setType,
|
||||
setCoverURL,
|
||||
setCreatePerSubfolder,
|
||||
updateRoot,
|
||||
addRoot: () => setRoots((prev) => [...prev, emptyRootDraft()]),
|
||||
removeRoot: (index: number) => setRoots((prev) => (prev.length <= 1 ? prev : prev.filter((_, i) => i !== index))),
|
||||
|
||||
Reference in New Issue
Block a user