From 7d7f3cc75882665155f3490f5f202bfc3dfd9f66 Mon Sep 17 00:00:00 2001 From: ShukeBta Date: Tue, 9 Jun 2026 19:03:28 +0800 Subject: [PATCH] feat: add cloud transfer and cache optimizations --- internal/handler/cloud.go | 91 +++- internal/handler/routes_admin.go | 2 + internal/handler/storage_config.go | 21 + internal/handler/system_extra.go | 14 + internal/service/cloud/cloud_test.go | 78 ++++ internal/service/cloud/pan115.go | 97 ++-- internal/service/cloud/quark.go | 81 ++-- internal/service/emby_compat.go | 37 +- internal/service/emby_compat_test.go | 33 ++ internal/service/image_proxy.go | 46 +- internal/service/image_proxy_test.go | 43 ++ internal/service/media.go | 64 +++ internal/service/organizer_directory.go | 11 +- internal/service/organizer_directory_test.go | 38 ++ internal/service/playback.go | 21 +- internal/service/scanner.go | 188 ++++++++ internal/service/scanner_cloud_test.go | 117 +++++ internal/service/scheduler.go | 144 ++++++ internal/service/scheduler_test.go | 62 ++- internal/service/service.go | 3 +- internal/service/storage_config.go | 2 +- internal/service/storage_upload.go | 457 +++++++++++++++++++ internal/service/storage_upload_test.go | 161 +++++++ internal/service/watcher.go | 3 + web/src/api/storage_config.ts | 35 ++ web/src/pages/SettingsPage.tsx | 64 +++ web/src/pages/StorageConfigPage.tsx | 176 ++++++- 27 files changed, 1972 insertions(+), 117 deletions(-) create mode 100644 internal/service/scanner_cloud_test.go create mode 100644 internal/service/storage_upload.go create mode 100644 internal/service/storage_upload_test.go diff --git a/internal/handler/cloud.go b/internal/handler/cloud.go index 6f73d72..1cad468 100644 --- a/internal/handler/cloud.go +++ b/internal/handler/cloud.go @@ -5,9 +5,12 @@ package handler import ( "io" "net/http" + "net/url" + "strings" "github.com/gin-gonic/gin" + "github.com/ShukeBta/MediaStationGo/internal/model" "github.com/ShukeBta/MediaStationGo/internal/service" "github.com/ShukeBta/MediaStationGo/internal/service/cloud" ) @@ -48,6 +51,84 @@ func cloudImportHandler(svc *service.Container) gin.HandlerFunc { } } +// cloudMountHandler creates or reuses a cloud:// media library for a cloud +// directory, then scans it recursively so cloud files become playable STRM/302 +// media rows without copying bytes to local disk. +func cloudMountHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + typ := c.Param("type") + var in struct { + Dir string `json:"dir"` + Name string `json:"name"` + MediaType string `json:"media_type"` + } + _ = c.ShouldBindJSON(&in) + if !cloud.IsCloudType(typ) { + c.JSON(http.StatusBadRequest, gin.H{"error": "unsupported cloud provider"}) + return + } + if _, err := svc.StorageCfg.CloudProvider(c.Request.Context(), typ); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + path := "cloud://" + typ + if dir := strings.TrimSpace(in.Dir); dir != "" { + path += "/" + url.PathEscape(dir) + } + name := strings.TrimSpace(in.Name) + if name == "" { + name = cloudMountLibraryName(typ, strings.TrimSpace(in.Dir)) + } + mediaType := strings.TrimSpace(in.MediaType) + if mediaType == "" { + mediaType = "movie" + } + libs, err := svc.Repo.Library.List(c.Request.Context()) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + var lib *model.Library + for i := range libs { + if libs[i].Path == path { + lib = &libs[i] + break + } + } + if lib == nil { + lib = &model.Library{Name: name, Path: path, Type: mediaType, Enabled: true} + if err := svc.Repo.Library.Create(c.Request.Context(), lib); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + } + var scan any + if svc.Scan != nil { + res, err := svc.Scan.ScanLibrary(c.Request.Context(), lib.ID) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error(), "library": lib}) + return + } + scan = res + } + c.JSON(http.StatusOK, gin.H{"library": lib, "scan": scan}) + } +} + +func cloudMountLibraryName(typ, dir string) string { + base := typ + switch typ { + case cloud.TypeQuark: + base = "夸克网盘" + case cloud.Type115: + base = "115 网盘" + } + if dir == "" || dir == "0" { + return base + } + return base + " · " + dir +} + // cloud115QRStartHandler begins a 115 QR-code login and returns the session + // QR image URL for the frontend to render. func cloud115QRStartHandler(svc *service.Container) gin.HandlerFunc { @@ -102,7 +183,11 @@ func cloudPlayHandler(svc *service.Container) gin.HandlerFunc { } // Proxy mode: the direct link needs auth headers the browser cannot // carry. Stream through with Range forwarding. - req, err := http.NewRequestWithContext(c.Request.Context(), http.MethodGet, link.URL, nil) + method := c.Request.Method + if method == "" { + method = http.MethodGet + } + req, err := http.NewRequestWithContext(c.Request.Context(), method, link.URL, nil) if err != nil { c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()}) return @@ -125,6 +210,8 @@ func cloudPlayHandler(svc *service.Container) gin.HandlerFunc { } } c.Status(resp.StatusCode) - _, _ = io.Copy(c.Writer, resp.Body) + if c.Request.Method != http.MethodHead { + _, _ = io.Copy(c.Writer, resp.Body) + } } } diff --git a/internal/handler/routes_admin.go b/internal/handler/routes_admin.go index 18c72aa..12156d8 100644 --- a/internal/handler/routes_admin.go +++ b/internal/handler/routes_admin.go @@ -35,10 +35,12 @@ func registerAdminRoutes(api *gin.RouterGroup, cfg *config.Config, svc *service. admin.GET("/storage/:type", getStorageConfigHandler(svc)) admin.PUT("/storage/:type", saveStorageConfigHandler(svc)) admin.POST("/storage/:type/test", testStorageConfigHandler(svc)) + admin.POST("/storage/:type/upload-local", storageUploadLocalHandler(svc)) // Cloud disk (115 / 夸克) browsing, QR login and 302 import. admin.GET("/cloud/:type/list", cloudListHandler(svc)) admin.POST("/cloud/:type/import", cloudImportHandler(svc)) + admin.POST("/cloud/:type/mount", cloudMountHandler(svc)) admin.POST("/cloud/:type/qr/start", cloud115QRStartHandler(svc)) admin.POST("/cloud/:type/qr/poll", cloud115QRPollHandler(svc)) diff --git a/internal/handler/storage_config.go b/internal/handler/storage_config.go index 9b937cb..2e2cd73 100644 --- a/internal/handler/storage_config.go +++ b/internal/handler/storage_config.go @@ -73,3 +73,24 @@ func testStorageConfigHandler(svc *service.Container) gin.HandlerFunc { c.JSON(http.StatusOK, gin.H{"ok": true}) } } + +func storageUploadLocalHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + var req service.CloudUploadInput + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + req.Type = c.Param("type") + res, err := svc.StorageCfg.UploadLocal(c.Request.Context(), req) + if err != nil && res == nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + if err != nil { + c.JSON(http.StatusOK, gin.H{"result": res, "error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"result": res}) + } +} diff --git a/internal/handler/system_extra.go b/internal/handler/system_extra.go index 3a3b03f..17bf2e7 100644 --- a/internal/handler/system_extra.go +++ b/internal/handler/system_extra.go @@ -98,6 +98,20 @@ func schemaHandler(_ *service.Container) gin.HandlerFunc { {"key": "scrape.delay_max_ms", "type": "number", "label": "刮削最大间隔毫秒"}, }, }, + { + "key": "cloud-upload", + "label": "网盘转存", + "items": []gin.H{ + {"key": "cloud.upload_auto_enabled", "type": "toggle", "label": "启用自动转存"}, + {"key": "cloud.upload_provider", "type": "select", "label": "转存目标"}, + {"key": "cloud.upload_source_dir", "type": "text", "label": "本地源目录"}, + {"key": "cloud.upload_dest_path", "type": "text", "label": "网盘目标目录"}, + {"key": "cloud.upload_recursive", "type": "toggle", "label": "递归扫描源目录"}, + {"key": "cloud.upload_sidecars", "type": "toggle", "label": "同步 NFO / 海报 / 字幕"}, + {"key": "cloud.upload_overwrite", "type": "toggle", "label": "覆盖远端同名文件"}, + {"key": "cloud.upload_interval_seconds", "type": "number", "label": "自动转存间隔秒数"}, + }, + }, { "key": "adult", "label": "Adult / NSFW", diff --git a/internal/service/cloud/cloud_test.go b/internal/service/cloud/cloud_test.go index dc634b1..572cf80 100644 --- a/internal/service/cloud/cloud_test.go +++ b/internal/service/cloud/cloud_test.go @@ -3,8 +3,10 @@ package cloud import ( "context" "encoding/base64" + "fmt" "net/http" "net/http/httptest" + "strconv" "strings" "testing" "time" @@ -62,6 +64,47 @@ func TestQuarkListAndResolve(t *testing.T) { } } +func TestQuarkListPaginates(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/file/sort" { + t.Fatalf("unexpected path %s", r.URL.Path) + } + page, _ := strconv.Atoi(r.URL.Query().Get("_page")) + w.Write([]byte(`{"status":200,"code":0,"data":{"list":[` + quarkPagePayload(page) + `]}}`)) + })) + defer srv.Close() + + p, err := New(TypeQuark, map[string]any{"cookie": "kps=abc", "base": srv.URL}, srv.Client()) + if err != nil { + t.Fatal(err) + } + entries, err := p.List(context.Background(), "0") + if err != nil { + t.Fatalf("list: %v", err) + } + if len(entries) != 101 { + t.Fatalf("entries = %d, want 101", len(entries)) + } + if entries[100].ID != "f100" || entries[100].Name != "Movie.100.mkv" { + t.Fatalf("last entry wrong: %#v", entries[100]) + } +} + +func quarkPagePayload(page int) string { + count := 100 + offset := 0 + if page > 1 { + count = 1 + offset = 100 + } + items := make([]string, 0, count) + for i := 0; i < count; i++ { + n := offset + i + items = append(items, fmt.Sprintf(`{"fid":"f%d","file_name":"Movie.%03d.mkv","dir":false,"size":%d}`, n, n, n)) + } + return strings.Join(items, ",") +} + func TestQuarkForce302(t *testing.T) { p := newQuark(map[string]any{"cookie": "c", "force_302": "true"}, http.DefaultClient) if p.proxy { @@ -128,6 +171,41 @@ func Test115ListAndResolve(t *testing.T) { } } +func Test115ListPaginates(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/files" { + t.Fatalf("unexpected path %s", r.URL.Path) + } + offset, _ := strconv.Atoi(r.URL.Query().Get("offset")) + count := 100 + if offset > 0 { + count = 1 + } + items := make([]string, 0, count) + for i := 0; i < count; i++ { + n := offset + i + items = append(items, fmt.Sprintf(`{"fid":"%d","n":"Movie.%03d.mkv","s":%d,"pc":"pick%d"}`, n, n, n, n)) + } + w.Write([]byte(`{"state":true,"data":[` + strings.Join(items, ",") + `]}`)) + })) + defer srv.Close() + + p, err := New(Type115, map[string]any{"cookie": "UID=1; CID=2", "base": srv.URL}, srv.Client()) + if err != nil { + t.Fatal(err) + } + entries, err := p.List(context.Background(), "0") + if err != nil { + t.Fatalf("list: %v", err) + } + if len(entries) != 101 { + t.Fatalf("entries = %d, want 101", len(entries)) + } + if entries[100].ID != "100" || entries[100].PickCode != "pick100" { + t.Fatalf("last entry wrong: %#v", entries[100]) + } +} + // Test115DownURLEndpointAndError exercises the live fetchDownURLPayload path: // it must POST an m115-encrypted `data` body to /app/chrome/downurl?t=... and // surface 115's error when state=false (no decryption needed for that branch). diff --git a/internal/service/cloud/pan115.go b/internal/service/cloud/pan115.go index 75e2669..117a45d 100644 --- a/internal/service/cloud/pan115.go +++ b/internal/service/cloud/pan115.go @@ -91,52 +91,59 @@ func (p *pan115Provider) List(ctx context.Context, dirID string) ([]FileEntry, e if dirID == "" { dirID = "0" } - q := url.Values{} - q.Set("aid", "1") - q.Set("cid", dirID) - q.Set("o", "user_ptime") - q.Set("asc", "0") - q.Set("offset", "0") - q.Set("show_dir", "1") - q.Set("limit", "100") - q.Set("format", "json") - resp, err := p.get(ctx, p.webBase+"/files?"+q.Encode()) - if err != nil { - return nil, err - } - defer resp.Body.Close() - var r struct { - State bool `json:"state"` - Error string `json:"error"` - Data []struct { - Fid string `json:"fid"` // file id (files only) - Cid string `json:"cid"` // category id (dirs use this) - N string `json:"n"` // name - S json.Number `json:"s"` // size - Pc string `json:"pc"` // pickcode - } `json:"data"` - } - if err := json.NewDecoder(resp.Body).Decode(&r); err != nil { - return nil, fmt.Errorf("115: decode list: %w", err) - } - if !r.State { - return nil, fmt.Errorf("115: list failed: %s", r.Error) - } - out := make([]FileEntry, 0, len(r.Data)) - for _, it := range r.Data { - isDir := it.Fid == "" - id := it.Fid - if isDir { - id = it.Cid + const pageSize = 100 + out := make([]FileEntry, 0, pageSize) + for offset := 0; ; offset += pageSize { + q := url.Values{} + q.Set("aid", "1") + q.Set("cid", dirID) + q.Set("o", "user_ptime") + q.Set("asc", "0") + q.Set("offset", strconv.Itoa(offset)) + q.Set("show_dir", "1") + q.Set("limit", strconv.Itoa(pageSize)) + q.Set("format", "json") + resp, err := p.get(ctx, p.webBase+"/files?"+q.Encode()) + if err != nil { + return nil, err + } + var r struct { + State bool `json:"state"` + Error string `json:"error"` + Data []struct { + Fid string `json:"fid"` // file id (files only) + Cid string `json:"cid"` // category id (dirs use this) + N string `json:"n"` // name + S json.Number `json:"s"` // size + Pc string `json:"pc"` // pickcode + } `json:"data"` + } + err = json.NewDecoder(resp.Body).Decode(&r) + _ = resp.Body.Close() + if err != nil { + return nil, fmt.Errorf("115: decode list: %w", err) + } + if !r.State { + return nil, fmt.Errorf("115: list failed: %s", r.Error) + } + for _, it := range r.Data { + isDir := it.Fid == "" + id := it.Fid + if isDir { + id = it.Cid + } + size, _ := it.S.Int64() + out = append(out, FileEntry{ + ID: id, + Name: it.N, + IsDir: isDir, + Size: size, + PickCode: it.Pc, + }) + } + if len(r.Data) < pageSize { + break } - size, _ := it.S.Int64() - out = append(out, FileEntry{ - ID: id, - Name: it.N, - IsDir: isDir, - Size: size, - PickCode: it.Pc, - }) } return out, nil } diff --git a/internal/service/cloud/quark.go b/internal/service/cloud/quark.go index b498998..aa37392 100644 --- a/internal/service/cloud/quark.go +++ b/internal/service/cloud/quark.go @@ -7,6 +7,7 @@ import ( "fmt" "io" "net/http" + "net/url" "strings" ) @@ -85,38 +86,54 @@ func (q *quarkProvider) List(ctx context.Context, dirID string) ([]FileEntry, er if dirID == "" { dirID = "0" } - path := fmt.Sprintf("/file/sort?pr=ucpro&fr=pc&uc_param_str=&pdir_fid=%s&_page=1&_size=100&_fetch_total=1&_sort=file_type:asc,updated_at:desc", dirID) - resp, err := q.do(ctx, http.MethodGet, path, nil) - if err != nil { - return nil, err - } - defer resp.Body.Close() - var r quarkResp - if err := json.NewDecoder(resp.Body).Decode(&r); err != nil { - return nil, fmt.Errorf("quark: decode list: %w", err) - } - if r.Code != 0 && r.Status != 200 { - return nil, fmt.Errorf("quark: list failed: %s", r.Message) - } - var data struct { - List []struct { - Fid string `json:"fid"` - FileName string `json:"file_name"` - Dir bool `json:"dir"` - Size int64 `json:"size"` - } `json:"list"` - } - if err := json.Unmarshal(r.Data, &data); err != nil { - return nil, fmt.Errorf("quark: decode list data: %w", err) - } - out := make([]FileEntry, 0, len(data.List)) - for _, it := range data.List { - out = append(out, FileEntry{ - ID: it.Fid, - Name: it.FileName, - IsDir: it.Dir, - Size: it.Size, - }) + const pageSize = 100 + out := make([]FileEntry, 0, pageSize) + for page := 1; ; page++ { + query := url.Values{} + query.Set("pr", "ucpro") + query.Set("fr", "pc") + query.Set("uc_param_str", "") + query.Set("pdir_fid", dirID) + query.Set("_page", fmt.Sprint(page)) + query.Set("_size", fmt.Sprint(pageSize)) + query.Set("_fetch_total", "1") + query.Set("_sort", "file_type:asc,updated_at:desc") + path := "/file/sort?" + query.Encode() + resp, err := q.do(ctx, http.MethodGet, path, nil) + if err != nil { + return nil, err + } + var r quarkResp + err = json.NewDecoder(resp.Body).Decode(&r) + _ = resp.Body.Close() + if err != nil { + return nil, fmt.Errorf("quark: decode list: %w", err) + } + if r.Code != 0 && r.Status != 200 { + return nil, fmt.Errorf("quark: list failed: %s", r.Message) + } + var data struct { + List []struct { + Fid string `json:"fid"` + FileName string `json:"file_name"` + Dir bool `json:"dir"` + Size int64 `json:"size"` + } `json:"list"` + } + if err := json.Unmarshal(r.Data, &data); err != nil { + return nil, fmt.Errorf("quark: decode list data: %w", err) + } + for _, it := range data.List { + out = append(out, FileEntry{ + ID: it.Fid, + Name: it.FileName, + IsDir: it.Dir, + Size: it.Size, + }) + } + if len(data.List) < pageSize { + break + } } return out, nil } diff --git a/internal/service/emby_compat.go b/internal/service/emby_compat.go index b4d5761..7ac9d81 100644 --- a/internal/service/emby_compat.go +++ b/internal/service/emby_compat.go @@ -423,14 +423,25 @@ func (e *EmbyService) episodeItems(ctx context.Context, rows []model.Media, p It func (e *EmbyService) payloadsForMedia(ctx context.Context, rows []model.Media, userID string) ([]map[string]any, error) { userFavs := map[string]bool{} userPos := map[string]int64{} - if userID != "" { + if userID != "" && len(rows) > 0 { + mediaIDs := make([]string, 0, len(rows)) + for _, row := range rows { + if strings.TrimSpace(row.ID) != "" { + mediaIDs = append(mediaIDs, row.ID) + } + } + if len(mediaIDs) == 0 { + mediaIDs = []string{"__none__"} + } var favs []model.Favorite - _ = e.repo.DB.WithContext(ctx).Where("user_id = ?", userID).Find(&favs).Error + favQuery := e.repo.DB.WithContext(ctx).Where("user_id = ?", userID).Where("media_id IN ?", mediaIDs) + _ = favQuery.Find(&favs).Error for _, f := range favs { userFavs[f.MediaID] = true } var hist []model.PlaybackHistory - _ = e.repo.DB.WithContext(ctx).Where("user_id = ?", userID).Find(&hist).Error + histQuery := e.repo.DB.WithContext(ctx).Where("user_id = ?", userID).Where("media_id IN ?", mediaIDs) + _ = histQuery.Find(&hist).Error for _, h := range hist { userPos[h.MediaID] = h.PositionMs } @@ -522,9 +533,18 @@ func (e *EmbyService) LatestItems(ctx context.Context, userID, parentID string, return nil, err } favs := map[string]bool{} - if userID != "" { + if userID != "" && len(rows) > 0 { + mediaIDs := make([]string, 0, len(rows)) + for _, row := range rows { + if strings.TrimSpace(row.ID) != "" { + mediaIDs = append(mediaIDs, row.ID) + } + } + if len(mediaIDs) == 0 { + mediaIDs = []string{"__none__"} + } var fr []model.Favorite - _ = e.repo.DB.WithContext(ctx).Where("user_id = ?", userID).Find(&fr).Error + _ = e.repo.DB.WithContext(ctx).Where("user_id = ? AND media_id IN ?", userID, mediaIDs).Find(&fr).Error for _, f := range fr { favs[f.MediaID] = true } @@ -1231,9 +1251,12 @@ func (e *EmbyService) mediaSource(m *model.Media, asEmbedded, directOnly bool) m } } if strings.TrimSpace(m.STRMURL) != "" { - // STRM 重定向:客户端直接拉远端,跳过我们这一层。 + // STRM / cloud:// media still plays through /Videos/{id}/stream. + // That route delegates to StreamService, which appends the caller's + // token to internal /api/cloud/play redirects and only then 302s to the + // provider/CDN. Returning m.STRMURL directly here would make Emby/Yamby + // clients hit /api/cloud/play without an auth token and fail with 401. src["IsRemote"] = true - src["DirectStreamUrl"] = m.STRMURL src["Path"] = m.STRMURL } return src diff --git a/internal/service/emby_compat_test.go b/internal/service/emby_compat_test.go index 3085f22..c1c7c41 100644 --- a/internal/service/emby_compat_test.go +++ b/internal/service/emby_compat_test.go @@ -249,6 +249,39 @@ func TestEmbyPlaybackInfoRespectsDirectPlayOnly(t *testing.T) { } } +func TestEmbyPlaybackInfoKeepsSTRMBehindStreamEndpoint(t *testing.T) { + svc := newTestEmbyService(t) + lib := model.Library{Name: "夸克网盘", Path: `cloud://quark/0`, Type: "movie", Enabled: true} + if err := svc.repo.Library.Create(t.Context(), &lib); err != nil { + t.Fatalf("create library: %v", err) + } + media := model.Media{ + Base: model.Base{ID: "cloud-1"}, + LibraryID: lib.ID, + Title: "Cloud Movie", + Path: `cloud://quark/f1`, + STRMURL: `/api/cloud/play/quark?ref=f1`, + } + if err := svc.repo.DB.Create(&media).Error; err != nil { + t.Fatalf("create media: %v", err) + } + + pb, err := svc.PlaybackInfo(t.Context(), "cloud-1", "user-1") + if err != nil { + t.Fatalf("playback info: %v", err) + } + src := pb["MediaSources"].([]map[string]any)[0] + if src["IsRemote"] != true { + t.Fatalf("strm media should be marked remote: %#v", src) + } + if src["DirectStreamUrl"] != "/Videos/cloud-1/stream" { + t.Fatalf("strm playback must stay behind token-aware stream endpoint: %#v", src) + } + if src["Path"] != "/api/cloud/play/quark?ref=f1" { + t.Fatalf("path should expose the strm target for diagnostics: %#v", src) + } +} + func newTestEmbyService(t *testing.T) *EmbyService { t.Helper() db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) diff --git a/internal/service/image_proxy.go b/internal/service/image_proxy.go index 7377ddc..4f68115 100644 --- a/internal/service/image_proxy.go +++ b/internal/service/image_proxy.go @@ -84,6 +84,12 @@ type ImageProxy struct { libRootsAt time.Time } +const ( + imageBrowserCacheControl = "public, max-age=2592000, immutable" + imagePlaceholderCacheControl = "public, max-age=3600" + imageNegativeCacheTTL = 6 * time.Hour +) + // NewImageProxy is the constructor. func NewImageProxy(cfg *config.Config, log *zap.Logger) *ImageProxy { // Honor HTTP(S)_PROXY env vars so deployments behind GFW can pull @@ -222,6 +228,13 @@ func servePlaceholder(w http.ResponseWriter) { _, _ = w.Write(transparent1x1PNG) } +func serveCachedPlaceholder(w http.ResponseWriter) { + w.Header().Set("Content-Type", "image/png") + w.Header().Set("Cache-Control", imagePlaceholderCacheControl) + w.WriteHeader(http.StatusOK) + _, _ = w.Write(transparent1x1PNG) +} + // Serve writes the requested image to w. Caller is expected to validate // the JWT before invoking it. func (p *ImageProxy) Serve(ctx context.Context, w http.ResponseWriter, r *http.Request, raw string) error { @@ -244,7 +257,7 @@ func (p *ImageProxy) Serve(ctx context.Context, w http.ResponseWriter, r *http.R modTime = stat.ModTime() } w.Header().Set("Content-Type", detectContentType(data)) - w.Header().Set("Cache-Control", "public, max-age=604800") + w.Header().Set("Cache-Control", imageBrowserCacheControl) http.ServeContent(w, r, filepath.Base(path), modTime, bytes.NewReader(data)) return nil } @@ -261,11 +274,12 @@ func (p *ImageProxy) Serve(ctx context.Context, w http.ResponseWriter, r *http.R sum := sha1.Sum([]byte(raw)) key := hex.EncodeToString(sum[:]) cachePath := filepath.Join(p.cacheDir, key) + failPath := cachePath + ".fail" // Cache hit. if data, err := os.ReadFile(cachePath); err == nil && len(data) > 0 { w.Header().Set("Content-Type", detectContentType(data)) - w.Header().Set("Cache-Control", "public, max-age=604800") + w.Header().Set("Cache-Control", imageBrowserCacheControl) stat, _ := os.Stat(cachePath) modTime := time.Now() if stat != nil { @@ -274,6 +288,12 @@ func (p *ImageProxy) Serve(ctx context.Context, w http.ResponseWriter, r *http.R http.ServeContent(w, r, key, modTime, bytes.NewReader(data)) return nil } + if stat, err := os.Stat(failPath); err == nil && time.Since(stat.ModTime()) < imageNegativeCacheTTL { + serveCachedPlaceholder(w) + return nil + } else if err == nil { + _ = os.Remove(failPath) + } // Cache miss → fetch upstream. if err := os.MkdirAll(p.cacheDir, 0o755); err != nil { @@ -298,14 +318,16 @@ func (p *ImageProxy) Serve(ctx context.Context, w http.ResponseWriter, r *http.R if err != nil { p.log.Warn("imageproxy: upstream fetch failed", zap.String("host", host), zap.Error(err)) - servePlaceholder(w) + p.markImageFetchFailed(failPath) + serveCachedPlaceholder(w) return nil } defer resp.Body.Close() if resp.StatusCode >= 400 { p.log.Warn("imageproxy: upstream returned non-OK", zap.String("host", host), zap.String("status", resp.Status)) - servePlaceholder(w) + p.markImageFetchFailed(failPath) + serveCachedPlaceholder(w) return nil } @@ -313,7 +335,8 @@ func (p *ImageProxy) Serve(ctx context.Context, w http.ResponseWriter, r *http.R if err != nil || len(data) == 0 { p.log.Warn("imageproxy: read upstream body failed", zap.String("host", host), zap.Error(err)) - servePlaceholder(w) + p.markImageFetchFailed(failPath) + serveCachedPlaceholder(w) return nil } @@ -325,6 +348,8 @@ func (p *ImageProxy) Serve(ctx context.Context, w http.ResponseWriter, r *http.R tmp.Close() if rerr := os.Rename(tmp.Name(), cachePath); rerr != nil { _ = os.Remove(tmp.Name()) + } else { + _ = os.Remove(failPath) } } else { tmp.Close() @@ -347,11 +372,20 @@ func (p *ImageProxy) Serve(ctx context.Context, w http.ResponseWriter, r *http.R if v := resp.Header.Get("Last-Modified"); v != "" { w.Header().Set("Last-Modified", v) } - w.Header().Set("Cache-Control", "public, max-age=604800") + w.Header().Set("Cache-Control", imageBrowserCacheControl) http.ServeContent(w, r, key, time.Now(), bytes.NewReader(data)) return nil } +func (p *ImageProxy) markImageFetchFailed(failPath string) { + if err := os.MkdirAll(filepath.Dir(failPath), 0o755); err != nil { + return + } + p.mu.Lock() + defer p.mu.Unlock() + _ = os.WriteFile(failPath, []byte(time.Now().Format(time.RFC3339Nano)), 0o644) +} + // Fetch 拉取远程图片并返回字节和 Content-Type(带缓存)。 func (p *ImageProxy) Fetch(ctx context.Context, raw string) ([]byte, string, error) { u, err := p.validateURL(raw) diff --git a/internal/service/image_proxy_test.go b/internal/service/image_proxy_test.go index 930232e..91f8e51 100644 --- a/internal/service/image_proxy_test.go +++ b/internal/service/image_proxy_test.go @@ -1,10 +1,13 @@ package service import ( + "io" "net/http" "net/http/httptest" "os" "path/filepath" + "strings" + "sync/atomic" "testing" "go.uber.org/zap" @@ -78,6 +81,46 @@ func TestImageProxyServesPosterUnderLibraryRoot(t *testing.T) { } } +func TestImageProxyCachesFailedRemoteImageFetch(t *testing.T) { + var calls int32 + proxy := NewImageProxy(&config.Config{Cache: config.CacheConfig{CacheDir: filepath.Join(t.TempDir(), "cache")}}, zap.NewNop()) + proxy.client = &http.Client{Transport: imageRoundTripFunc(func(req *http.Request) (*http.Response, error) { + atomic.AddInt32(&calls, 1) + return &http.Response{ + StatusCode: http.StatusBadGateway, + Status: "502 Bad Gateway", + Header: make(http.Header), + Body: io.NopCloser(strings.NewReader("upstream unavailable")), + Request: req, + }, nil + })} + raw := "https://image.tmdb.org/t/p/w500/poster.jpg" + for i := 0; i < 2; i++ { + rec := httptest.NewRecorder() + if err := proxy.Serve(t.Context(), rec, httptest.NewRequest(http.MethodGet, "/api/img", nil), raw); err != nil { + t.Fatal(err) + } + if rec.Code != http.StatusOK { + t.Fatalf("status = %d, want 200", rec.Code) + } + if rec.Body.Len() != len(transparent1x1PNG) { + t.Fatalf("body length = %d, want placeholder %d", rec.Body.Len(), len(transparent1x1PNG)) + } + if got := rec.Header().Get("Cache-Control"); got != imagePlaceholderCacheControl { + t.Fatalf("Cache-Control = %q, want %q", got, imagePlaceholderCacheControl) + } + } + if got := atomic.LoadInt32(&calls); got != 1 { + t.Fatalf("upstream calls = %d, want 1 due to negative cache", got) + } +} + +type imageRoundTripFunc func(*http.Request) (*http.Response, error) + +func (f imageRoundTripFunc) RoundTrip(req *http.Request) (*http.Response, error) { + return f(req) +} + func TestIsPrivateHost(t *testing.T) { blocked := []string{"127.0.0.1", "10.0.0.5", "192.168.1.10", "169.254.169.254", "0.0.0.0", "::1", ""} for _, h := range blocked { diff --git a/internal/service/media.go b/internal/service/media.go index 33eb291..1bd5b5b 100644 --- a/internal/service/media.go +++ b/internal/service/media.go @@ -92,6 +92,70 @@ func resolveAccessibleLibraryPath(path string) (string, error) { return "", fmt.Errorf("path is not an accessible directory: %s", abs) } +func resolveAccessibleMappedPath(path string) (string, os.FileInfo, error) { + input := strings.TrimSpace(path) + if input == "" { + return "", nil, errors.New("path required") + } + candidates := mappedPathCandidates(input) + for _, candidate := range candidates { + if info, err := os.Stat(candidate); err == nil { + return filepath.Clean(candidate), info, nil + } + } + abs, err := filepath.Abs(input) + if err != nil { + return "", nil, fmt.Errorf("invalid path: %w", err) + } + return "", nil, fmt.Errorf("path is not accessible: %s", abs) +} + +func resolveMappedDestinationPath(path string) string { + path = strings.TrimSpace(path) + if path == "" { + return "" + } + clean := filepath.Clean(path) + if _, err := os.Stat(clean); err == nil { + return clean + } + for _, candidate := range mappedPathCandidates(clean) { + if candidate == clean { + continue + } + return filepath.Clean(candidate) + } + return clean +} + +func mappedPathCandidates(input string) []string { + var candidates []string + add := func(candidate string) { + candidate = filepath.Clean(filepath.FromSlash(strings.TrimSpace(candidate))) + if candidate == "" || candidate == "." { + return + } + for _, existing := range candidates { + if sameLibraryPath(existing, candidate) { + return + } + } + candidates = append(candidates, candidate) + } + clean := filepath.Clean(input) + add(clean) + for _, candidate := range dockerVolumePathCandidates(clean) { + add(candidate) + } + if abs, err := filepath.Abs(input); err == nil { + add(abs) + for _, candidate := range dockerVolumePathCandidates(abs) { + add(candidate) + } + } + return candidates +} + func isAccessibleDir(path string) bool { info, err := os.Stat(path) return err == nil && info.IsDir() diff --git a/internal/service/organizer_directory.go b/internal/service/organizer_directory.go index 1c70e22..f6e0add 100644 --- a/internal/service/organizer_directory.go +++ b/internal/service/organizer_directory.go @@ -108,16 +108,15 @@ func (o *OrganizerService) defaultDestRoot(ctx context.Context, override string) // OrganizeDirectory organizes every video file found under opts.SourcePath into // the destination root, applying dedup + 洗版 (resolution replacement). func (o *OrganizerService) OrganizeDirectory(ctx context.Context, opts OrganizeOptions) (*OrganizeResult, error) { - source := strings.TrimSpace(o.defaultSourceRoot(ctx, opts.SourcePath)) - if source == "" { + requestedSource := strings.TrimSpace(o.defaultSourceRoot(ctx, opts.SourcePath)) + if requestedSource == "" { return nil, errors.New("source path required") } - source = filepath.Clean(source) - info, statErr := os.Stat(source) + source, info, statErr := resolveAccessibleMappedPath(requestedSource) if statErr != nil { - return nil, fmt.Errorf("source directory not accessible: %s", source) + return nil, fmt.Errorf("source directory not accessible: %s", filepath.Clean(requestedSource)) } - dest := filepath.Clean(o.defaultDestRoot(ctx, opts.DestPath)) + dest := resolveMappedDestinationPath(o.defaultDestRoot(ctx, opts.DestPath)) if dest == "" || dest == "." { return nil, errors.New("destination path required") } diff --git a/internal/service/organizer_directory_test.go b/internal/service/organizer_directory_test.go index 376e5cf..48ced49 100644 --- a/internal/service/organizer_directory_test.go +++ b/internal/service/organizer_directory_test.go @@ -100,6 +100,44 @@ func TestOrganizeDirectoryUsesConfiguredSourceWhenRequestSourceEmpty(t *testing. } } +func TestOrganizeDirectoryMapsConfiguredHostPathsToContainerPaths(t *testing.T) { + root := t.TempDir() + hostDownloads := filepath.Join(root, "nas-host", "downloads") + hostMedia := filepath.Join(root, "nas-host", "media") + containerDownloads := filepath.Join(root, "container", "downloads") + containerMedia := filepath.Join(root, "container", "media") + containerSource := filepath.Join(containerDownloads, "国产剧") + writeOrgFile(t, filepath.Join(containerSource, "Some Show S01E01 2024 1080p.mkv"), "show-e01") + + t.Setenv("MEDIASTATION_DOWNLOAD_DIR", hostDownloads) + t.Setenv("MEDIASTATION_DOWNLOAD_CONTAINER_DIR", containerDownloads) + t.Setenv("MEDIASTATION_MEDIA_DIR", hostMedia) + t.Setenv("MEDIASTATION_MEDIA_CONTAINER_DIR", containerMedia) + + repos := newOrganizerTestRepo(t) + for key, value := range map[string]string{ + "organize.source_dir": filepath.Join(hostDownloads, "国产剧"), + "organize.target_dir": hostMedia, + "organize.transfer_mode": "copy", + } { + if err := repos.Setting.Set(t.Context(), key, value); err != nil { + t.Fatal(err) + } + } + org := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) + res, err := org.OrganizeDirectory(t.Context(), OrganizeOptions{}) + if err != nil { + t.Fatalf("organize mapped host paths: %v", err) + } + if res.SourcePath != filepath.Clean(containerSource) || res.DestPath != filepath.Clean(containerMedia) { + t.Fatalf("result paths = source %q dest %q, want %q -> %q", res.SourcePath, res.DestPath, containerSource, containerMedia) + } + want := filepath.Join(containerMedia, "电视剧", "国产剧", "Some Show", "Season 01", "Some Show - S01E01.mkv") + if _, err := os.Stat(want); err != nil { + t.Fatalf("expected organized file at %q: %v", want, err) + } +} + func TestOrganizeDirectoryAcceptsSingleVideoFileSource(t *testing.T) { root := t.TempDir() src := filepath.Join(root, "downloads", "Dune 2021 2160p WEB-DL.mkv") diff --git a/internal/service/playback.go b/internal/service/playback.go index 289ddc0..3031720 100644 --- a/internal/service/playback.go +++ b/internal/service/playback.go @@ -64,11 +64,26 @@ func (p *PlaybackService) RecentHistory(ctx context.Context, userID string, limi if err != nil { return nil, err } + mediaIDs := make([]string, 0, len(rows)) + for i := range rows { + if rows[i].MediaID != "" { + mediaIDs = append(mediaIDs, rows[i].MediaID) + } + } + mediaByID := map[string]model.Media{} + if len(mediaIDs) > 0 { + var mediaRows []model.Media + if err := p.repo.DB.WithContext(ctx).Where("id IN ?", mediaIDs).Find(&mediaRows).Error; err == nil { + for _, media := range mediaRows { + mediaByID[media.ID] = media + } + } + } items := make([]HistoryItem, 0, len(rows)) for i := range rows { - var m model.Media - if err := p.repo.DB.Where("id = ?", rows[i].MediaID).First(&m).Error; err == nil { - items = append(items, HistoryItem{PlaybackHistory: rows[i], Media: &m}) + if m, ok := mediaByID[rows[i].MediaID]; ok { + media := m + items = append(items, HistoryItem{PlaybackHistory: rows[i], Media: &media}) } else { items = append(items, HistoryItem{PlaybackHistory: rows[i]}) } diff --git a/internal/service/scanner.go b/internal/service/scanner.go index f52bdf4..3ea430c 100644 --- a/internal/service/scanner.go +++ b/internal/service/scanner.go @@ -12,6 +12,8 @@ package service import ( "context" + "fmt" + "net/url" "os" "path/filepath" "strings" @@ -49,6 +51,7 @@ type ScannerService struct { hub *Hub probe *FFprobeService scraper *ScraperService + storage *StorageConfigService } // NewScannerService is the constructor. @@ -66,6 +69,13 @@ func NewScannerService( } } +// SetStorageConfig wires cloud-disk storage access into the scanner. It is set +// after service construction because StorageConfigService depends on Crypto, +// while the scanner is needed earlier by watcher/download services. +func (s *ScannerService) SetStorageConfig(storage *StorageConfigService) { + s.storage = storage +} + // ScanResult summarises a scan run. type ScanResult struct { LibraryID string `json:"library_id"` @@ -88,6 +98,9 @@ func (s *ScannerService) scanLibrary(ctx context.Context, libraryID string, auto if err != nil || lib == nil { return nil, err } + if typ, dirID, ok := parseCloudLibraryPath(lib.Path); ok { + return s.scanCloudLibrary(ctx, lib, typ, dirID, autoScrape) + } res := &ScanResult{LibraryID: lib.ID} seen := make(map[string]struct{}) seenInodes := make(map[string]string) @@ -161,6 +174,124 @@ func (s *ScannerService) IngestPath(ctx context.Context, libraryID, path string) return res.Added+res.Updated > 0, nil } +func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Library, typ, rootDir string, autoScrape bool) (*ScanResult, error) { + res := &ScanResult{LibraryID: lib.ID} + if s.storage == nil { + return res, fmt.Errorf("cloud storage service unavailable") + } + seen := make(map[string]struct{}) + visitedDirs := map[string]struct{}{} + var walkCloud func(string) error + walkCloud = func(dirID string) error { + if _, ok := visitedDirs[dirID]; ok { + return nil + } + visitedDirs[dirID] = struct{}{} + entries, err := s.storage.CloudList(ctx, typ, dirID) + if err != nil { + return err + } + for _, entry := range entries { + select { + case <-ctx.Done(): + return ctx.Err() + default: + } + if entry.IsDir { + if strings.TrimSpace(entry.ID) != "" { + if err := walkCloud(entry.ID); err != nil { + return err + } + } + continue + } + ext := strings.ToLower(filepath.Ext(entry.Name)) + if _, ok := videoExtensions[ext]; !ok { + continue + } + ref := cloudEntryRef(typ, entry.ID, entry.PickCode) + if ref == "" { + res.Skipped++ + continue + } + path := cloudMediaPath(typ, ref) + seen[path] = struct{}{} + s.ingestCloudFile(ctx, lib, typ, ref, entry.Name, entry.Size, res) + } + return nil + } + if err := walkCloud(rootDir); err != nil { + return res, err + } + removed, err := s.pruneMissingCloudMedia(ctx, lib.ID, seen) + if err != nil { + s.log.Warn("prune missing cloud media failed", zap.String("library_id", lib.ID), zap.Error(err)) + } else { + res.Removed = removed + } + s.hub.Publish("scan", map[string]any{ + "library_id": lib.ID, + "finished": true, + "visited": res.Visited, + "added": res.Added, + "updated": res.Updated, + "removed": res.Removed, + "cloud": true, + }) + if autoScrape && s.scraper != nil && s.scraper.AnyEnabled() && s.autoScrapeEnabled(ctx) { + go func(libID string) { + if _, err := s.scraper.EnrichLibrary(context.Background(), libID); err != nil { + s.log.Warn("scraper enrich failed", zap.Error(err)) + } + }(lib.ID) + } + return res, nil +} + +func (s *ScannerService) ingestCloudFile(ctx context.Context, lib *model.Library, typ, ref, name string, size int64, res *ScanResult) { + res.Visited++ + ext := strings.ToLower(filepath.Ext(name)) + title, year := CleanQuery(name) + if title == "" { + title = strings.TrimSuffix(filepath.Base(name), ext) + } + if title == "" { + title = ref + } + path := cloudMediaPath(typ, ref) + isNewMedia := !s.mediaPathExists(ctx, path) + m := &model.Media{ + LibraryID: lib.ID, + Title: title, + Year: year, + Path: path, + SizeBytes: size, + Container: strings.TrimPrefix(ext, "."), + STRMURL: "/api/cloud/play/" + typ + "?ref=" + url.QueryEscape(ref), + ScrapeStatus: "pending", + } + parsedSeason, parsedEpisode := ParseEpisode(name) + m.SeasonNum = parsedSeason + m.EpisodeNum = parsedEpisode + if err := s.repo.Media.Upsert(ctx, m); err != nil { + s.log.Warn("upsert cloud media failed", zap.String("path", path), zap.Error(err)) + return + } + if isNewMedia { + res.Added++ + } else { + res.Updated++ + } + s.hub.Publish("scan", map[string]any{ + "library_id": lib.ID, + "path": path, + "visited": res.Visited, + "added": res.Added, + "updated": res.Updated, + "cloud": true, + }) +} + // RemovePath deletes the media row for a path that has disappeared from disk // (incremental delete used by the watcher on Remove/Rename events). func (s *ScannerService) RemovePath(ctx context.Context, path string) (int64, error) { @@ -326,6 +457,63 @@ func (s *ScannerService) pruneMissingMedia(ctx context.Context, libraryID string return removed, nil } +func (s *ScannerService) pruneMissingCloudMedia(ctx context.Context, libraryID string, seen map[string]struct{}) (int64, error) { + var rows []model.Media + if err := s.repo.DB.WithContext(ctx). + Where("library_id = ? AND path LIKE ?", libraryID, "cloud://%"). + Find(&rows).Error; err != nil { + return 0, err + } + var removed int64 + for _, row := range rows { + if _, ok := seen[row.Path]; ok { + continue + } + res := s.repo.DB.WithContext(ctx). + Where("id = ?", row.ID). + Delete(&model.Media{}) + if res.Error != nil { + return removed, res.Error + } + removed += res.RowsAffected + } + return removed, nil +} + +func parseCloudLibraryPath(raw string) (typ, dirID string, ok bool) { + raw = strings.TrimSpace(raw) + if !strings.HasPrefix(strings.ToLower(raw), "cloud://") { + return "", "", false + } + u, err := url.Parse(raw) + if err != nil || strings.ToLower(u.Scheme) != "cloud" { + return "", "", false + } + typ = strings.TrimSpace(u.Host) + if typ == "" { + return "", "", false + } + dirID = strings.Trim(strings.TrimSpace(u.Path), "/") + if qDir := strings.TrimSpace(u.Query().Get("dir")); qDir != "" { + dirID = qDir + } + if decoded, err := url.PathUnescape(dirID); err == nil { + dirID = decoded + } + return typ, dirID, true +} + +func cloudEntryRef(typ, id, pickCode string) string { + if typ == "cloud115" && strings.TrimSpace(pickCode) != "" { + return strings.TrimSpace(pickCode) + } + return strings.TrimSpace(id) +} + +func cloudMediaPath(typ, ref string) string { + return "cloud://" + strings.TrimSpace(typ) + "/" + strings.TrimSpace(ref) +} + func applyLocalMetadata(m *model.Media, local *LocalMetadata) { if local.Title != "" { m.Title = local.Title diff --git a/internal/service/scanner_cloud_test.go b/internal/service/scanner_cloud_test.go new file mode 100644 index 0000000..704ebd2 --- /dev/null +++ b/internal/service/scanner_cloud_test.go @@ -0,0 +1,117 @@ +package service + +import ( + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/glebarez/sqlite" + "go.uber.org/zap" + "gorm.io/gorm" + + "github.com/ShukeBta/MediaStationGo/internal/config" + "github.com/ShukeBta/MediaStationGo/internal/model" + "github.com/ShukeBta/MediaStationGo/internal/repository" +) + +func TestScanCloudLibraryImportsRecursivePlayableMedia(t *testing.T) { + empty := false + upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/file/sort" { + t.Fatalf("unexpected path %s", r.URL.Path) + } + w.Header().Set("Content-Type", "application/json") + if empty { + _, _ = w.Write([]byte(`{"status":200,"code":0,"data":{"list":[]}}`)) + return + } + switch r.URL.Query().Get("pdir_fid") { + case "0": + _, _ = w.Write([]byte(`{"status":200,"code":0,"data":{"list":[ + {"fid":"d1","file_name":"Movies","dir":true,"size":0}, + {"fid":"f1","file_name":"Root.Movie.2024.mkv","dir":false,"size":123} + ]}}`)) + case "d1": + _, _ = w.Write([]byte(`{"status":200,"code":0,"data":{"list":[ + {"fid":"f2","file_name":"Nested.Show.S01E02.mp4","dir":false,"size":456} + ]}}`)) + default: + t.Fatalf("unexpected pdir_fid %q", r.URL.Query().Get("pdir_fid")) + } + })) + defer upstream.Close() + + db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&model.Library{}, &model.Media{}, &model.Setting{}, &model.StorageConfig{}); err != nil { + t.Fatal(err) + } + repos := repository.New(db) + log := zap.NewNop() + storage := NewStorageConfigService(log, repos, NewCryptoService("", log)) + if _, err := storage.Save(t.Context(), StorageInput{ + Type: "quark", + Config: map[string]any{ + "cookie": "kps=test", + "base": upstream.URL, + }, + }); err != nil { + t.Fatal(err) + } + lib := model.Library{Name: "夸克网盘", Path: "cloud://quark/0", Type: "tv", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + scanner := NewScannerService(&config.Config{}, log, repos, NewHub(log), nil, nil) + scanner.SetStorageConfig(storage) + + res, err := scanner.ScanLibrary(t.Context(), lib.ID) + if err != nil { + t.Fatalf("scan cloud: %v", err) + } + if res.Visited != 2 || res.Added != 2 { + t.Fatalf("scan result = %#v, want visited=2 added=2", res) + } + var rows []model.Media + if err := repos.DB.Order("path").Find(&rows).Error; err != nil { + t.Fatal(err) + } + if len(rows) != 2 { + t.Fatalf("media rows = %d, want 2: %#v", len(rows), rows) + } + if rows[0].Path != "cloud://quark/f1" || rows[0].STRMURL != "/api/cloud/play/quark?ref=f1" { + t.Fatalf("root media path/strm wrong: path=%q strm=%q", rows[0].Path, rows[0].STRMURL) + } + if rows[1].SeasonNum != 1 || rows[1].EpisodeNum != 2 || !strings.Contains(rows[1].STRMURL, "ref=f2") { + t.Fatalf("nested episode metadata wrong: %#v", rows[1]) + } + + empty = true + res, err = scanner.ScanLibrary(t.Context(), lib.ID) + if err != nil { + t.Fatalf("rescan cloud: %v", err) + } + if res.Removed != 2 { + t.Fatalf("removed = %d, want 2", res.Removed) + } + if got := countMedia(t, repos); got != 0 { + t.Fatalf("media count after prune = %d, want 0", got) + } +} + +func TestCloudLibraryPathParsing(t *testing.T) { + typ, dir, ok := parseCloudLibraryPath("cloud://cloud115/abc%20123?ignored=1") + if !ok || typ != "cloud115" || dir != "abc 123" { + t.Fatalf("parse path got typ=%q dir=%q ok=%v", typ, dir, ok) + } + typ, dir, ok = parseCloudLibraryPath("cloud://quark?dir=0") + if !ok || typ != "quark" || dir != "0" { + t.Fatalf("parse query got typ=%q dir=%q ok=%v", typ, dir, ok) + } + if ref := cloudEntryRef("cloud115", "fid", "pick"); ref != "pick" { + t.Fatalf("115 ref = %q, want pick", ref) + } +} diff --git a/internal/service/scheduler.go b/internal/service/scheduler.go index 7890a49..3f76655 100644 --- a/internal/service/scheduler.go +++ b/internal/service/scheduler.go @@ -44,6 +44,7 @@ type SchedulerService struct { scanner *ScannerService transcoder *TranscoderService organizer *OrganizerService + storageCfg *StorageConfigService hub *Hub cacheDir string @@ -70,6 +71,7 @@ func NewSchedulerService( scanner *ScannerService, transcoder *TranscoderService, organizer *OrganizerService, + storageCfg *StorageConfigService, hub *Hub, cacheDir string, ) *SchedulerService { @@ -79,6 +81,7 @@ func NewSchedulerService( scanner: scanner, transcoder: transcoder, organizer: organizer, + storageCfg: storageCfg, hub: hub, cacheDir: cacheDir, stopCh: make(chan struct{}), @@ -93,6 +96,16 @@ func (s *SchedulerService) Start(ctx context.Context) { interval: 60 * time.Minute, run: s.jobScanLibraries, }, + { + name: "cloud_sync", + interval: s.cloudSyncInterval(ctx), + run: s.jobSyncCloudLibraries, + }, + { + name: "cloud_upload", + interval: s.cloudUploadInterval(ctx), + run: s.jobUploadLocalToCloud, + }, { name: "organize_source", interval: s.organizeSourceInterval(ctx), @@ -227,6 +240,137 @@ func (s *SchedulerService) jobScanLibraries(ctx context.Context) error { return nil } +// jobUploadLocalToCloud copies local media files into the configured external +// storage backend. It is opt-in and never deletes the local source files. +func (s *SchedulerService) jobUploadLocalToCloud(ctx context.Context) error { + manual, _ := ctx.Value(schedulerManualRunKey{}).(bool) + if s.storageCfg == nil || (!manual && !s.autoCloudUploadEnabled(ctx)) { + return nil + } + input := s.cloudUploadInput(ctx) + if strings.TrimSpace(input.Type) == "" || strings.TrimSpace(input.SourcePath) == "" { + return nil + } + res, err := s.storageCfg.UploadLocal(ctx, input) + if s.log != nil && res != nil { + s.log.Info("cloud upload finished", + zap.String("type", input.Type), + zap.String("source", res.SourcePath), + zap.String("dest", res.DestPath), + zap.Int("uploaded", res.Uploaded), + zap.Int("skipped", res.Skipped), + zap.Int64("bytes", res.Bytes), + zap.Int("errors", len(res.Errors)), + ) + } + return err +} + +func (s *SchedulerService) cloudUploadInput(ctx context.Context) CloudUploadInput { + get := func(key string) string { + if s.repo == nil || s.repo.Setting == nil { + return "" + } + v, _ := s.repo.Setting.Get(ctx, key) + return strings.TrimSpace(v) + } + return CloudUploadInput{ + Type: get(CloudUploadProviderKey), + SourcePath: get(CloudUploadSourceDirKey), + DestPath: get(CloudUploadDestPathKey), + Recursive: parseBoolSetting(get(CloudUploadRecursiveKey), true), + IncludeSidecars: parseBoolSetting(get(CloudUploadSidecarsKey), true), + Overwrite: parseBoolSetting(get(CloudUploadOverwriteKey), false), + } +} + +func (s *SchedulerService) autoCloudUploadEnabled(ctx context.Context) bool { + if s.repo == nil || s.repo.Setting == nil { + return false + } + v, err := s.repo.Setting.Get(ctx, CloudUploadAutoEnabledKey) + if err != nil { + return false + } + return parseBoolSetting(v, false) +} + +func (s *SchedulerService) cloudUploadInterval(ctx context.Context) time.Duration { + const fallback = time.Hour + if s.repo == nil || s.repo.Setting == nil { + return fallback + } + v, err := s.repo.Setting.Get(ctx, CloudUploadIntervalSecondsKey) + if err != nil { + return fallback + } + seconds, err := strconv.Atoi(strings.TrimSpace(v)) + if err != nil || seconds <= 0 { + return fallback + } + if seconds < 300 { + seconds = 300 + } + return time.Duration(seconds) * time.Second +} + +// jobSyncCloudLibraries keeps mounted cloud:// libraries refreshed without +// enabling full disk scans. It imports remote cloud files as STRM-backed media +// rows; the actual bytes stay on the provider and playback continues through +// /api/cloud/play 302/proxy. +func (s *SchedulerService) jobSyncCloudLibraries(ctx context.Context) error { + manual, _ := ctx.Value(schedulerManualRunKey{}).(bool) + if s.scanner == nil || (!manual && !s.autoCloudSyncEnabled(ctx)) { + return nil + } + libs, err := s.repo.Library.List(ctx) + if err != nil { + return err + } + for _, l := range libs { + if !l.Enabled { + continue + } + if _, _, ok := parseCloudLibraryPath(l.Path); !ok { + continue + } + if _, err := s.scanner.ScanLibrary(ctx, l.ID); err != nil { + s.log.Warn("cloud sync failed", zap.String("library", l.ID), zap.Error(err)) + } + } + return nil +} + +func (s *SchedulerService) autoCloudSyncEnabled(ctx context.Context) bool { + if s.repo == nil || s.repo.Setting == nil { + return true + } + v, err := s.repo.Setting.Get(ctx, "cloud.auto_sync_enabled") + if err != nil { + return true + } + return parseBoolSetting(v, true) +} + +func (s *SchedulerService) cloudSyncInterval(ctx context.Context) time.Duration { + const fallback = 30 * time.Minute + if s.repo == nil || s.repo.Setting == nil { + return fallback + } + v, err := s.repo.Setting.Get(ctx, "cloud.sync_interval_seconds") + if err != nil { + return fallback + } + seconds, err := strconv.Atoi(strings.TrimSpace(v)) + if err != nil || seconds <= 0 { + return fallback + } + if seconds < 300 { + seconds = 300 + } + return time.Duration(seconds) * time.Second +} + // periodicScanEnabled reports whether the operator opted into periodic full // library re-scans. Defaults to false so the incremental watcher is the only // thing touching the disk under normal operation. diff --git a/internal/service/scheduler_test.go b/internal/service/scheduler_test.go index d3bc96c..d4e8c77 100644 --- a/internal/service/scheduler_test.go +++ b/internal/service/scheduler_test.go @@ -1,14 +1,20 @@ package service import ( + "net/http" + "net/http/httptest" "os" "path/filepath" "testing" "time" + "github.com/glebarez/sqlite" "go.uber.org/zap" + "gorm.io/gorm" "github.com/ShukeBta/MediaStationGo/internal/config" + "github.com/ShukeBta/MediaStationGo/internal/model" + "github.com/ShukeBta/MediaStationGo/internal/repository" ) func TestSchedulerOrganizeSourceDisabledByDefault(t *testing.T) { @@ -30,7 +36,7 @@ func TestSchedulerOrganizeSourceDisabledByDefault(t *testing.T) { } organizer := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) - scheduler := NewSchedulerService(zap.NewNop(), repos, nil, nil, organizer, NewHub(zap.NewNop()), "") + scheduler := NewSchedulerService(zap.NewNop(), repos, nil, nil, organizer, nil, NewHub(zap.NewNop()), "") if err := scheduler.jobOrganizeSource(t.Context()); err != nil { t.Fatalf("disabled organize source job should be a no-op: %v", err) } @@ -60,7 +66,7 @@ func TestSchedulerOrganizeSourceUsesConfiguredSourceAndDestination(t *testing.T) } organizer := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) - scheduler := NewSchedulerService(zap.NewNop(), repos, nil, nil, organizer, NewHub(zap.NewNop()), "") + scheduler := NewSchedulerService(zap.NewNop(), repos, nil, nil, organizer, nil, NewHub(zap.NewNop()), "") if err := scheduler.jobOrganizeSource(t.Context()); err != nil { t.Fatalf("organize source job: %v", err) } @@ -89,7 +95,7 @@ func TestSchedulerRunNowOrganizeSourceBypassesDisabledSwitch(t *testing.T) { } organizer := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) - scheduler := NewSchedulerService(zap.NewNop(), repos, nil, nil, organizer, NewHub(zap.NewNop()), "") + scheduler := NewSchedulerService(zap.NewNop(), repos, nil, nil, organizer, nil, NewHub(zap.NewNop()), "") scheduler.jobs = []*scheduledJob{{ name: "organize_source", interval: time.Minute, @@ -105,3 +111,53 @@ func TestSchedulerRunNowOrganizeSourceBypassesDisabledSwitch(t *testing.T) { t.Fatalf("expected run-now organize output at %q: %v", want, err) } } + +func TestSchedulerCloudSyncImportsMountedCloudLibrary(t *testing.T) { + upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/file/sort" || r.URL.Query().Get("pdir_fid") != "0" { + t.Fatalf("unexpected cloud list request %s?%s", r.URL.Path, r.URL.RawQuery) + } + _, _ = w.Write([]byte(`{"status":200,"code":0,"data":{"list":[ + {"fid":"f1","file_name":"Cloud.Movie.2026.mkv","dir":false,"size":1024} + ]}}`)) + })) + defer upstream.Close() + + db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&model.Library{}, &model.Media{}, &model.Setting{}, &model.StorageConfig{}); err != nil { + t.Fatal(err) + } + repos := repository.New(db) + log := zap.NewNop() + storage := NewStorageConfigService(log, repos, NewCryptoService("", log)) + if _, err := storage.Save(t.Context(), StorageInput{ + Type: "quark", + Config: map[string]any{ + "cookie": "kps=test", + "base": upstream.URL, + }, + }); err != nil { + t.Fatal(err) + } + lib := model.Library{Name: "夸克网盘", Path: "cloud://quark/0", Type: "movie", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + scanner := NewScannerService(&config.Config{}, log, repos, NewHub(log), nil, nil) + scanner.SetStorageConfig(storage) + scheduler := NewSchedulerService(log, repos, scanner, nil, nil, storage, NewHub(log), "") + + if err := scheduler.jobSyncCloudLibraries(t.Context()); err != nil { + t.Fatalf("cloud sync: %v", err) + } + var media model.Media + if err := repos.DB.First(&media, "path = ?", "cloud://quark/f1").Error; err != nil { + t.Fatalf("cloud media not imported: %v", err) + } + if media.STRMURL != "/api/cloud/play/quark?ref=f1" { + t.Fatalf("strm url = %q", media.STRMURL) + } +} diff --git a/internal/service/service.go b/internal/service/service.go index a58eab1..a27c80e 100644 --- a/internal/service/service.go +++ b/internal/service/service.go @@ -113,10 +113,11 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont playProfiles := NewPlayProfileService(log, repos) permissions := NewPermissionService(log, repos) storageCfg := NewStorageConfigService(log, repos, crypto) + scanner.SetStorageConfig(storageCfg) downloadClients := NewDownloadClientService(log, repos) assistant := NewAssistantService(log, repos, ai) douban := NewDoubanProvider(cfg, log) - scheduler := NewSchedulerService(log, repos, scanner, transcoder, organizer, hub, cfg.Cache.CacheDir) + scheduler := NewSchedulerService(log, repos, scanner, transcoder, organizer, storageCfg, hub, cfg.Cache.CacheDir) // 初始化认证相关服务 tokenSvc := NewTokenService(cfg, log, repos) diff --git a/internal/service/storage_config.go b/internal/service/storage_config.go index dfe712b..51bd00d 100644 --- a/internal/service/storage_config.go +++ b/internal/service/storage_config.go @@ -319,7 +319,7 @@ func (s *StorageConfigService) CloudImport(ctx context.Context, typ, fileRef, na m := &model.Media{ LibraryID: lib.ID, Title: title, - Path: "cloud://" + typ + "/" + fileRef, + Path: cloudMediaPath(typ, fileRef), SizeBytes: size, Container: container, STRMURL: "/api/cloud/play/" + typ + "?ref=" + url.QueryEscape(fileRef), diff --git a/internal/service/storage_upload.go b/internal/service/storage_upload.go new file mode 100644 index 0000000..ebff62f --- /dev/null +++ b/internal/service/storage_upload.go @@ -0,0 +1,457 @@ +package service + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "net/url" + "os" + "path" + "path/filepath" + "strings" +) + +const ( + CloudUploadAutoEnabledKey = "cloud.upload_auto_enabled" + CloudUploadProviderKey = "cloud.upload_provider" + CloudUploadSourceDirKey = "cloud.upload_source_dir" + CloudUploadDestPathKey = "cloud.upload_dest_path" + CloudUploadRecursiveKey = "cloud.upload_recursive" + CloudUploadSidecarsKey = "cloud.upload_sidecars" + CloudUploadOverwriteKey = "cloud.upload_overwrite" + CloudUploadIntervalSecondsKey = "cloud.upload_interval_seconds" + CloudUploadUnsupportedProvider = "本地文件直传目前支持 Alist / WebDAV;115/夸克原生上传需要各自的分片上传私有接口,建议先用 Alist 挂载 115/夸克后选择 Alist 转存。" +) + +type CloudUploadInput struct { + Type string `json:"type"` + SourcePath string `json:"source_path"` + DestPath string `json:"dest_path"` + Recursive bool `json:"recursive"` + IncludeSidecars bool `json:"include_sidecars"` + Overwrite bool `json:"overwrite"` +} + +type CloudUploadResult struct { + SourcePath string `json:"source_path"` + DestPath string `json:"dest_path"` + Uploaded int `json:"uploaded"` + Skipped int `json:"skipped"` + Bytes int64 `json:"bytes"` + Errors []string `json:"errors,omitempty"` + Items []CloudUploadResultItem `json:"items,omitempty"` +} + +type CloudUploadResultItem struct { + Source string `json:"source"` + Target string `json:"target"` + Action string `json:"action"` // upload / skip / error + Size int64 `json:"size,omitempty"` + Reason string `json:"reason,omitempty"` +} + +type storageUploader interface { + ensureDir(ctx context.Context, remoteDir string) error + exists(ctx context.Context, remotePath string) (bool, error) + upload(ctx context.Context, localPath, remotePath string, size int64) error +} + +var cloudUploadSidecarExtensions = map[string]struct{}{ + ".nfo": {}, ".jpg": {}, ".jpeg": {}, ".png": {}, ".webp": {}, + ".srt": {}, ".ass": {}, ".ssa": {}, ".vtt": {}, ".sub": {}, ".idx": {}, +} + +// UploadLocal copies local media files into an external storage backend. It is +// intentionally conservative: it never deletes local files, skips existing +// remote targets unless Overwrite is set, and only uploads video + common +// sidecar metadata files. +func (s *StorageConfigService) UploadLocal(ctx context.Context, in CloudUploadInput) (*CloudUploadResult, error) { + in.Type = strings.TrimSpace(in.Type) + in.SourcePath = strings.TrimSpace(in.SourcePath) + in.DestPath = normalizeRemotePath(in.DestPath) + if in.SourcePath == "" { + return nil, errors.New("source_path required") + } + uploader, err := s.uploader(ctx, in.Type) + if err != nil { + return nil, err + } + info, err := os.Stat(in.SourcePath) + if err != nil { + return nil, fmt.Errorf("source path not accessible: %w", err) + } + result := &CloudUploadResult{SourcePath: in.SourcePath, DestPath: in.DestPath} + if !info.IsDir() { + s.uploadOne(ctx, uploader, in, in.SourcePath, filepath.Base(in.SourcePath), info.Size(), result) + return result, firstUploadError(result) + } + root := filepath.Clean(in.SourcePath) + walkFn := func(localPath string, entryInfo os.FileInfo, walkErr error) error { + if walkErr != nil { + addUploadError(result, localPath, "", walkErr) + return nil + } + if entryInfo == nil || entryInfo.IsDir() { + if !in.Recursive && filepath.Clean(localPath) != root { + return filepath.SkipDir + } + return nil + } + if !eligibleCloudUploadFile(localPath, in.IncludeSidecars) { + return nil + } + rel, err := filepath.Rel(root, localPath) + if err != nil { + addUploadError(result, localPath, "", err) + return nil + } + s.uploadOne(ctx, uploader, in, localPath, filepath.ToSlash(rel), entryInfo.Size(), result) + return nil + } + if err := filepath.Walk(in.SourcePath, walkFn); err != nil { + return result, err + } + return result, firstUploadError(result) +} + +func (s *StorageConfigService) uploader(ctx context.Context, typ string) (storageUploader, error) { + view, err := s.Get(ctx, typ) + if err != nil { + return nil, err + } + if view == nil || !view.Enabled { + return nil, fmt.Errorf("%s storage not configured", typ) + } + switch typ { + case "alist": + return newAlistUploader(view.Config), nil + case "webdav": + return newWebDAVUploader(view.Config), nil + case "s3": + return nil, errors.New("s3 local upload is not implemented yet") + case "cloud115", "quark": + return nil, errors.New(CloudUploadUnsupportedProvider) + default: + return nil, fmt.Errorf("unsupported storage type %q", typ) + } +} + +func (s *StorageConfigService) uploadOne(ctx context.Context, uploader storageUploader, in CloudUploadInput, localPath, rel string, size int64, result *CloudUploadResult) { + remotePath := joinRemotePath(in.DestPath, rel) + if err := uploader.ensureDir(ctx, path.Dir(remotePath)); err != nil { + addUploadError(result, localPath, remotePath, err) + return + } + if !in.Overwrite { + exists, err := uploader.exists(ctx, remotePath) + if err != nil { + addUploadError(result, localPath, remotePath, err) + return + } + if exists { + result.Skipped++ + addUploadItem(result, CloudUploadResultItem{Source: localPath, Target: remotePath, Action: "skip", Size: size, Reason: "remote exists"}) + return + } + } + if err := uploader.upload(ctx, localPath, remotePath, size); err != nil { + addUploadError(result, localPath, remotePath, err) + return + } + result.Uploaded++ + result.Bytes += size + addUploadItem(result, CloudUploadResultItem{Source: localPath, Target: remotePath, Action: "upload", Size: size}) +} + +func eligibleCloudUploadFile(localPath string, includeSidecars bool) bool { + ext := strings.ToLower(filepath.Ext(localPath)) + if _, ok := videoExtensions[ext]; ok { + return true + } + if includeSidecars { + _, ok := cloudUploadSidecarExtensions[ext] + return ok + } + return false +} + +func addUploadError(result *CloudUploadResult, source, target string, err error) { + result.Errors = append(result.Errors, fmt.Sprintf("%s: %v", source, err)) + addUploadItem(result, CloudUploadResultItem{Source: source, Target: target, Action: "error", Reason: err.Error()}) +} + +func addUploadItem(result *CloudUploadResult, item CloudUploadResultItem) { + if len(result.Items) < 200 { + result.Items = append(result.Items, item) + } +} + +func firstUploadError(result *CloudUploadResult) error { + if result.Uploaded > 0 || len(result.Errors) == 0 { + return nil + } + return errors.New(result.Errors[0]) +} + +func normalizeRemotePath(p string) string { + p = strings.ReplaceAll(strings.TrimSpace(p), "\\", "/") + if p == "" || p == "." { + return "/" + } + if !strings.HasPrefix(p, "/") { + p = "/" + p + } + return path.Clean(p) +} + +func joinRemotePath(base, rel string) string { + parts := []string{normalizeRemotePath(base)} + for _, part := range strings.Split(strings.ReplaceAll(rel, "\\", "/"), "/") { + part = strings.TrimSpace(part) + if part != "" && part != "." { + parts = append(parts, part) + } + } + return path.Clean(path.Join(parts...)) +} + +type alistUploader struct { + server string + token string + client *http.Client +} + +func newAlistUploader(cfg map[string]any) *alistUploader { + return &alistUploader{ + server: strings.TrimRight(strr(cfg["server"]), "/"), + token: strr(cfg["token"]), + client: &http.Client{}, + } +} + +func (a *alistUploader) ensureDir(ctx context.Context, remoteDir string) error { + if a.server == "" { + return errors.New("alist missing server") + } + remoteDir = normalizeRemotePath(remoteDir) + if remoteDir == "/" { + return nil + } + current := "" + for _, part := range strings.Split(strings.Trim(remoteDir, "/"), "/") { + current = normalizeRemotePath(path.Join(current, part)) + payload, _ := json.Marshal(map[string]string{"path": current}) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, a.server+"/api/fs/mkdir", bytes.NewReader(payload)) + if err != nil { + return err + } + a.auth(req) + req.Header.Set("Content-Type", "application/json") + resp, err := a.client.Do(req) + if err != nil { + return err + } + err = a.checkJSON(resp, "alist mkdir") + if err != nil && !isAlreadyExistsMessage(err.Error()) { + return err + } + } + return nil +} + +func (a *alistUploader) exists(ctx context.Context, remotePath string) (bool, error) { + payload, _ := json.Marshal(map[string]string{"path": normalizeRemotePath(remotePath)}) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, a.server+"/api/fs/get", bytes.NewReader(payload)) + if err != nil { + return false, err + } + a.auth(req) + req.Header.Set("Content-Type", "application/json") + resp, err := a.client.Do(req) + if err != nil { + return false, err + } + defer resp.Body.Close() + if resp.StatusCode == http.StatusNotFound { + return false, nil + } + var out struct { + Code int `json:"code"` + Message string `json:"message"` + } + _ = json.NewDecoder(resp.Body).Decode(&out) + return resp.StatusCode >= 200 && resp.StatusCode < 300 && out.Code == 200, nil +} + +func (a *alistUploader) upload(ctx context.Context, localPath, remotePath string, size int64) error { + f, err := os.Open(localPath) + if err != nil { + return err + } + defer f.Close() + req, err := http.NewRequestWithContext(ctx, http.MethodPut, a.server+"/api/fs/put", f) + if err != nil { + return err + } + a.auth(req) + req.ContentLength = size + req.Header.Set("Content-Type", "application/octet-stream") + req.Header.Set("File-Path", url.PathEscape(normalizeRemotePath(remotePath))) + resp, err := a.client.Do(req) + if err != nil { + return err + } + return a.checkJSON(resp, "alist upload") +} + +func (a *alistUploader) auth(req *http.Request) { + if a.token != "" { + req.Header.Set("Authorization", a.token) + } +} + +func (a *alistUploader) checkJSON(resp *http.Response, op string) error { + defer resp.Body.Close() + body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return fmt.Errorf("%s: http %d: %s", op, resp.StatusCode, strings.TrimSpace(string(body))) + } + var out struct { + Code int `json:"code"` + Message string `json:"message"` + } + if len(bytes.TrimSpace(body)) == 0 { + return nil + } + if err := json.Unmarshal(body, &out); err != nil { + return nil + } + if out.Code != 0 && out.Code != 200 { + return fmt.Errorf("%s: %s", op, out.Message) + } + return nil +} + +type webDAVUploader struct { + base *url.URL + username string + password string + client *http.Client +} + +func newWebDAVUploader(cfg map[string]any) *webDAVUploader { + u, _ := url.Parse(strings.TrimRight(strr(cfg["url"]), "/")) + return &webDAVUploader{ + base: u, + username: strr(cfg["username"]), + password: strr(cfg["password"]), + client: &http.Client{}, + } +} + +func (w *webDAVUploader) ensureDir(ctx context.Context, remoteDir string) error { + if w.base == nil || w.base.Scheme == "" || w.base.Host == "" { + return errors.New("webdav missing url") + } + remoteDir = normalizeRemotePath(remoteDir) + if remoteDir == "/" { + return nil + } + current := "" + for _, part := range strings.Split(strings.Trim(remoteDir, "/"), "/") { + current = normalizeRemotePath(path.Join(current, part)) + req, err := http.NewRequestWithContext(ctx, "MKCOL", w.urlFor(current), nil) + if err != nil { + return err + } + w.auth(req) + resp, err := w.client.Do(req) + if err != nil { + return err + } + _, _ = io.Copy(io.Discard, resp.Body) + _ = resp.Body.Close() + if resp.StatusCode >= 200 && resp.StatusCode < 300 { + continue + } + if resp.StatusCode == http.StatusMethodNotAllowed || resp.StatusCode == http.StatusConflict { + continue + } + return fmt.Errorf("webdav mkdir %s: http %d", current, resp.StatusCode) + } + return nil +} + +func (w *webDAVUploader) exists(ctx context.Context, remotePath string) (bool, error) { + req, err := http.NewRequestWithContext(ctx, http.MethodHead, w.urlFor(remotePath), nil) + if err != nil { + return false, err + } + w.auth(req) + resp, err := w.client.Do(req) + if err != nil { + return false, err + } + _, _ = io.Copy(io.Discard, resp.Body) + _ = resp.Body.Close() + if resp.StatusCode == http.StatusNotFound { + return false, nil + } + return resp.StatusCode >= 200 && resp.StatusCode < 300, nil +} + +func (w *webDAVUploader) upload(ctx context.Context, localPath, remotePath string, size int64) error { + f, err := os.Open(localPath) + if err != nil { + return err + } + defer f.Close() + req, err := http.NewRequestWithContext(ctx, http.MethodPut, w.urlFor(remotePath), f) + if err != nil { + return err + } + w.auth(req) + req.ContentLength = size + resp, err := w.client.Do(req) + if err != nil { + return err + } + _, _ = io.Copy(io.Discard, resp.Body) + _ = resp.Body.Close() + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return fmt.Errorf("webdav upload %s: http %d", remotePath, resp.StatusCode) + } + return nil +} + +func (w *webDAVUploader) auth(req *http.Request) { + if w.username != "" { + req.SetBasicAuth(w.username, w.password) + } +} + +func (w *webDAVUploader) urlFor(remotePath string) string { + u := *w.base + basePath := strings.TrimRight(u.EscapedPath(), "/") + segments := make([]string, 0) + if basePath != "" && basePath != "/" { + segments = append(segments, strings.Trim(basePath, "/")) + } + for _, part := range strings.Split(strings.Trim(normalizeRemotePath(remotePath), "/"), "/") { + if part != "" { + segments = append(segments, url.PathEscape(part)) + } + } + u.RawPath = "" + u.Path = "/" + strings.Join(segments, "/") + return u.String() +} + +func isAlreadyExistsMessage(message string) bool { + message = strings.ToLower(message) + return strings.Contains(message, "exist") || strings.Contains(message, "已存在") +} diff --git a/internal/service/storage_upload_test.go b/internal/service/storage_upload_test.go new file mode 100644 index 0000000..01bc432 --- /dev/null +++ b/internal/service/storage_upload_test.go @@ -0,0 +1,161 @@ +package service + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "net/url" + "os" + "path/filepath" + "sort" + "strings" + "testing" + + "github.com/glebarez/sqlite" + "go.uber.org/zap" + "gorm.io/gorm" + + "github.com/ShukeBta/MediaStationGo/internal/model" + "github.com/ShukeBta/MediaStationGo/internal/repository" +) + +func TestStorageConfigUploadLocalToAlist(t *testing.T) { + var uploaded []string + var authHeaders []string + alist := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/api/fs/mkdir": + _ = json.NewDecoder(r.Body).Decode(&map[string]string{}) + _, _ = w.Write([]byte(`{"code":200,"message":"success"}`)) + case "/api/fs/get": + w.WriteHeader(http.StatusNotFound) + _, _ = w.Write([]byte(`{"code":404,"message":"not found"}`)) + case "/api/fs/put": + authHeaders = append(authHeaders, r.Header.Get("Authorization")) + decoded, err := url.PathUnescape(r.Header.Get("File-Path")) + if err != nil { + t.Fatalf("decode file path: %v", err) + } + uploaded = append(uploaded, decoded) + _, _ = w.Write([]byte(`{"code":200,"message":"success"}`)) + default: + t.Fatalf("unexpected alist path %s", r.URL.Path) + } + })) + defer alist.Close() + + repos, storage := newStorageUploadTestService(t) + if _, err := storage.Save(t.Context(), StorageInput{ + Type: "alist", + Config: map[string]any{ + "server": alist.URL, + "token": "alist-token", + }, + }); err != nil { + t.Fatal(err) + } + source := t.TempDir() + if err := os.WriteFile(filepath.Join(source, "Movie.2026.mkv"), []byte("movie"), 0o644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(source, "Movie.2026.nfo"), []byte("nfo"), 0o644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(source, "ignore.txt"), []byte("txt"), 0o644); err != nil { + t.Fatal(err) + } + + res, err := storage.UploadLocal(t.Context(), CloudUploadInput{ + Type: "alist", + SourcePath: source, + DestPath: "/backup", + Recursive: true, + IncludeSidecars: true, + }) + if err != nil { + t.Fatalf("upload local: %v", err) + } + if res.Uploaded != 2 || res.Skipped != 0 || len(res.Errors) != 0 { + t.Fatalf("result = %+v", res) + } + sort.Strings(uploaded) + want := []string{"/backup/Movie.2026.mkv", "/backup/Movie.2026.nfo"} + if strings.Join(uploaded, "\n") != strings.Join(want, "\n") { + t.Fatalf("uploaded = %#v, want %#v", uploaded, want) + } + for _, header := range authHeaders { + if header != "alist-token" { + t.Fatalf("authorization header = %q", header) + } + } + if got, _ := repos.StorageConfig.Get(t.Context(), "alist"); got == nil { + t.Fatalf("storage config should remain saved") + } +} + +func TestSchedulerCloudUploadUsesConfiguredLocalSource(t *testing.T) { + var uploaded []string + alist := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/api/fs/mkdir": + _, _ = w.Write([]byte(`{"code":200}`)) + case "/api/fs/get": + w.WriteHeader(http.StatusNotFound) + case "/api/fs/put": + decoded, _ := url.PathUnescape(r.Header.Get("File-Path")) + uploaded = append(uploaded, decoded) + _, _ = w.Write([]byte(`{"code":200}`)) + default: + t.Fatalf("unexpected alist path %s", r.URL.Path) + } + })) + defer alist.Close() + + repos, storage := newStorageUploadTestService(t) + if _, err := storage.Save(t.Context(), StorageInput{ + Type: "alist", + Config: map[string]any{ + "server": alist.URL, + "token": "token", + }, + }); err != nil { + t.Fatal(err) + } + source := t.TempDir() + if err := os.WriteFile(filepath.Join(source, "Show.S01E01.mkv"), []byte("episode"), 0o644); err != nil { + t.Fatal(err) + } + for key, value := range map[string]string{ + CloudUploadAutoEnabledKey: "true", + CloudUploadProviderKey: "alist", + CloudUploadSourceDirKey: source, + CloudUploadDestPathKey: "/cloud-media", + CloudUploadRecursiveKey: "true", + CloudUploadSidecarsKey: "false", + } { + if err := repos.Setting.Set(t.Context(), key, value); err != nil { + t.Fatal(err) + } + } + scheduler := NewSchedulerService(zap.NewNop(), repos, nil, nil, nil, storage, NewHub(zap.NewNop()), "") + if err := scheduler.jobUploadLocalToCloud(t.Context()); err != nil { + t.Fatalf("cloud upload job: %v", err) + } + if len(uploaded) != 1 || uploaded[0] != "/cloud-media/Show.S01E01.mkv" { + t.Fatalf("uploaded = %#v", uploaded) + } +} + +func newStorageUploadTestService(t *testing.T) (*repository.Container, *StorageConfigService) { + t.Helper() + db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&model.StorageConfig{}, &model.Setting{}); err != nil { + t.Fatal(err) + } + repos := repository.New(db) + log := zap.NewNop() + return repos, NewStorageConfigService(log, repos, NewCryptoService("", log)) +} diff --git a/internal/service/watcher.go b/internal/service/watcher.go index d7f5f1d..1b1e22a 100644 --- a/internal/service/watcher.go +++ b/internal/service/watcher.go @@ -99,6 +99,9 @@ func (w *WatcherService) Refresh(ctx context.Context) error { if !l.Enabled { continue } + if _, _, ok := parseCloudLibraryPath(l.Path); ok { + continue + } for _, dir := range listDirsForWatch(l.Path) { current[dir] = l.ID } diff --git a/web/src/api/storage_config.ts b/web/src/api/storage_config.ts index 0008641..3669258 100644 --- a/web/src/api/storage_config.ts +++ b/web/src/api/storage_config.ts @@ -32,6 +32,22 @@ export interface StorageConfig { updated_at: string } +export interface CloudUploadResult { + source_path: string + dest_path: string + uploaded: number + skipped: number + bytes: number + errors?: string[] + items?: Array<{ + source: string + target: string + action: 'upload' | 'skip' | 'error' + size?: number + reason?: string + }> +} + export const storageAPI = { status: () => api @@ -53,6 +69,20 @@ export const storageAPI = { config, }) .then((r) => r.data), + + uploadLocal: ( + type: StorageType, + input: { + source_path: string + dest_path: string + recursive: boolean + include_sidecars: boolean + overwrite: boolean + }, + ) => + api + .post<{ result: CloudUploadResult; error?: string }>(`/admin/storage/${type}/upload-local`, input) + .then((r) => r.data), } // cloudAPI drives 网盘 browsing, QR login and 302 import. @@ -69,6 +99,11 @@ export const cloudAPI = { .post(`/admin/cloud/${type}/import`, { ref, name, size }) .then((r) => r.data), + mount: (type: StorageType, dir = '', name = '', media_type = 'movie') => + api + .post(`/admin/cloud/${type}/mount`, { dir, name, media_type }) + .then((r) => r.data), + qrStart: (type: StorageType) => api.post(`/admin/cloud/${type}/qr/start`).then((r) => r.data), diff --git a/web/src/pages/SettingsPage.tsx b/web/src/pages/SettingsPage.tsx index bd99d5a..5ca9d4c 100644 --- a/web/src/pages/SettingsPage.tsx +++ b/web/src/pages/SettingsPage.tsx @@ -281,6 +281,70 @@ const GROUPS: SettingGroup[] = [ }, ], }, + { + key: 'cloud-upload', + label: '网盘转存', + description: '把本地媒体复制上传到外部存储。推荐:将 115/夸克挂载到 Alist 后,使用 Alist 作为转存目标。', + items: [ + { + key: 'cloud.upload_auto_enabled', + label: '启用自动转存', + type: 'toggle', + hint: '开启后后台会按间隔扫描本地源目录,把视频、NFO、海报、字幕复制到目标存储;不会删除本地源文件。', + defaultValue: 'false', + }, + { + key: 'cloud.upload_provider', + label: '转存目标', + type: 'select', + defaultValue: 'alist', + options: [ + { value: 'alist', label: 'Alist(推荐,可桥接 115/夸克)' }, + { value: 'webdav', label: 'WebDAV' }, + { value: 'cloud115', label: '115 原生(待接分片上传)' }, + { value: 'quark', label: '夸克原生(待接分片上传)' }, + ], + }, + { + key: 'cloud.upload_source_dir', + label: '本地源目录', + type: 'text', + placeholder: '/media/电影 或 F:\\media\\Movies', + }, + { + key: 'cloud.upload_dest_path', + label: '网盘目标目录', + type: 'text', + defaultValue: '/MediaStationGo', + placeholder: '/MediaStationGo', + }, + { + key: 'cloud.upload_recursive', + label: '递归扫描源目录', + type: 'toggle', + defaultValue: 'true', + }, + { + key: 'cloud.upload_sidecars', + label: '同步 NFO / 海报 / 字幕', + type: 'toggle', + defaultValue: 'true', + }, + { + key: 'cloud.upload_overwrite', + label: '覆盖远端同名文件', + type: 'toggle', + defaultValue: 'false', + }, + { + key: 'cloud.upload_interval_seconds', + label: '自动转存间隔秒数', + type: 'number', + hint: '最小 300 秒,建议 3600 秒或更高,避免频繁读盘和触发网盘风控。', + defaultValue: '3600', + }, + ], + }, { key: 'adult', label: 'Adult / NSFW', diff --git a/web/src/pages/StorageConfigPage.tsx b/web/src/pages/StorageConfigPage.tsx index 47f3db6..8726e38 100644 --- a/web/src/pages/StorageConfigPage.tsx +++ b/web/src/pages/StorageConfigPage.tsx @@ -1,5 +1,5 @@ import { FormEvent, useEffect, useMemo, useState } from 'react' -import { Cloud, FileVideo, Folder, Loader2, QrCode, Save, Send } from 'lucide-react' +import { Cloud, FileVideo, Folder, Loader2, QrCode, Save, Send, Upload } from 'lucide-react' import toast from 'react-hot-toast' import { @@ -212,11 +212,125 @@ function StorageForm({ type }: { type: StorageType }) { 保存 + {isCloud(type) && } ) } +function StorageUploadPanel({ type }: { type: StorageType }) { + const [sourcePath, setSourcePath] = useState('') + const [destPath, setDestPath] = useState('/MediaStationGo') + const [recursive, setRecursive] = useState(true) + const [includeSidecars, setIncludeSidecars] = useState(true) + const [overwrite, setOverwrite] = useState(false) + const [busy, setBusy] = useState(false) + + const supported = type === 'alist' || type === 'webdav' + + const submit = async () => { + if (!supported) { + toast.error('本地直传目前支持 Alist / WebDAV。115/夸克建议先挂载到 Alist,再通过 Alist 转存。') + return + } + if (!sourcePath.trim()) { + toast.error('请填写本地源目录或文件路径') + return + } + setBusy(true) + try { + const { result, error } = await storageAPI.uploadLocal(type, { + source_path: sourcePath.trim(), + dest_path: destPath.trim() || '/', + recursive, + include_sidecars: includeSidecars, + overwrite, + }) + const errText = error || (result.errors && result.errors.length > 0 ? ` · 错误 ${result.errors.length}` : '') + toast.success(`转存完成:上传 ${result.uploaded} · 跳过 ${result.skipped} · ${fmtBytes(result.bytes)}${errText}`) + } catch (err: unknown) { + toast.error((err as { response?: { data?: { error?: string } } })?.response?.data?.error ?? '转存失败') + } finally { + setBusy(false) + } + } + + return ( +
+
+
+

+ 本地媒体转存到此存储 +

+

+ 复制本地媒体文件到外部存储,保留本地源文件;自动跳过远端已存在文件。 +

+
+ {!supported && ( + + 直传待接 + + )} +
+ {!supported && ( +

+ 115 / 夸克原生上传需要私有分片上传协议。本版本先支持 Alist / WebDAV:把 115 或夸克挂到 Alist 后,选择 Alist 即可把本地文件转存到对应网盘。 +

+ )} +
+ + +
+
+ + + + +
+
+ ) +} + +function fmtBytes(value: number): string { + if (!value) return '0 B' + const units = ['B', 'KB', 'MB', 'GB', 'TB'] + let size = value + let idx = 0 + while (size >= 1024 && idx < units.length - 1) { + size /= 1024 + idx++ + } + return `${size.toFixed(size >= 10 || idx === 0 ? 0 : 1)} ${units[idx]}` +} + // QRLoginPanel drives the 115 QR-code login: start → render image → poll → // fill the cookie field on confirmation. function QRLoginPanel({ type, onCookie }: { type: StorageType; onCookie: (c: string) => void }) { @@ -293,6 +407,8 @@ function CloudBrowser({ type }: { type: StorageType }) { const [stack, setStack] = useState<{ id: string; name: string }[]>([{ id: '', name: '根目录' }]) const [items, setItems] = useState([]) const [loading, setLoading] = useState(false) + const [mounting, setMounting] = useState(false) + const [mountMediaType, setMountMediaType] = useState('movie') const [error, setError] = useState('') const cur = stack[stack.length - 1] @@ -329,18 +445,56 @@ function CloudBrowser({ type }: { type: StorageType }) { } } + const mountCurrent = async () => { + setMounting(true) + try { + const label = TYPE_LABEL[type] ?? type + const name = cur.id ? `${label} · ${cur.name}` : label + const res = await cloudAPI.mount(type, cur.id, name, mountMediaType) + const scan = (res as { scan?: { added?: number; updated?: number; removed?: number; visited?: number } }).scan + toast.success(`已挂载为媒体库并扫描:新增 ${scan?.added ?? 0} · 更新 ${scan?.updated ?? 0} · 访问 ${scan?.visited ?? 0}`) + } catch (err: unknown) { + toast.error((err as { response?: { data?: { error?: string } } })?.response?.data?.error ?? '挂载失败') + } finally { + setMounting(false) + } + } + return (
e.preventDefault()}> -
- 网盘资源: - {stack.map((s, i) => ( - - - {i < stack.length - 1 && /} - - ))} +
+
+ 网盘资源: + {stack.map((s, i) => ( + + + {i < stack.length - 1 && /} + + ))} +
+
+ + +
{loading ? (