Compare commits

...

14 Commits

Author SHA1 Message Date
github-actions[bot] 22b7290ee1 chore: bump version to 0.0.46 [skip ci] 2026-08-27 07:42:16 +00:00
truewhile 14037d5dea 5 2026-08-27 15:42:00 +08:00
github-actions[bot] 152db3fb9f chore: bump version to 0.0.45 [skip ci] 2026-08-27 06:20:25 +00:00
truewhile 0413d123da 4 2026-08-27 14:20:05 +08:00
github-actions[bot] 5f6bd7b5cd chore: bump version to 0.0.44 [skip ci] 2026-08-27 05:41:04 +00:00
truewhile 3c25bb5d61 3 2026-08-27 13:40:46 +08:00
github-actions[bot] 659b91b000 chore: bump version to 0.0.43 [skip ci] 2026-08-27 03:34:05 +00:00
truewhile fe5b3bd56a 2 2026-08-27 11:33:50 +08:00
github-actions[bot] 82bbb116ae chore: bump version to 0.0.42 [skip ci] 2026-08-27 03:07:35 +00:00
truewhile 9f5ff7e6f0 1 2026-08-27 11:07:18 +08:00
github-actions[bot] a00504080a chore: bump version to 0.0.41 [skip ci] 2026-08-27 02:09:16 +00:00
truewhile fc6e2e6f10 优化 strm 同步记录:支持删除记录并展示上传数量统计 2026-08-27 10:08:54 +08:00
github-actions[bot] 65c5f3e4bf chore: bump version to 0.0.40 [skip ci] 2026-08-26 15:35:37 +00:00
truewhile 07e340251b 优化 2026-08-26 23:32:02 +08:00
22 changed files with 931 additions and 211 deletions
+117
View File
@@ -27,6 +27,9 @@ permissions:
jobs:
version-and-publish:
runs-on: ubuntu-latest
outputs:
new_version: ${{ steps.bump_version.outputs.new_version }}
tag: ${{ steps.bump_version.outputs.tag }}
steps:
- uses: actions/checkout@v4
with:
@@ -149,3 +152,117 @@ jobs:
VERSION=${{ steps.bump_version.outputs.new_version }}
cache-from: type=gha
cache-to: type=gha,mode=max
# 单文件可执行构建:把前端打包进二进制(go:embed),交叉编译 Windows /
# Linux / macOS 的 amd64 / arm64 产物,作为 GitHub Release 附件发布。
build-frontend:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with:
node-version: '20'
cache: 'npm'
cache-dependency-path: web/package-lock.json
- name: Install
working-directory: web
run: npm ci
- name: Build SPA
working-directory: web
run: npm run build
- name: Upload web/dist
uses: actions/upload-artifact@v4
with:
name: web-dist
path: web/dist
retention-days: 1
# 先创建(幂等)空的 GitHub Release,供后续 build-binaries 并行上传附件,
# 也避免矩阵各 job 并发 upload 时 release 尚不存在而互相竞争。
publish-create-release:
needs: [version-and-publish]
runs-on: ubuntu-latest
permissions:
contents: write
steps:
- uses: actions/checkout@v4
- name: Create release
env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
RELEASE_TAG: ${{ needs.version-and-publish.outputs.tag }}
run: |
set -eux
# tag 已由 version-and-publish 推送;若 release 已存在则忽略(--verify-tag 幂等)
gh release create "$RELEASE_TAG" \
--title "MMTL ${{ needs.version-and-publish.outputs.new_version }}" \
--notes "自动化发布 ${{ needs.version-and-publish.outputs.new_version }}" \
--verify-tag --latest || true
build-binaries:
needs: [version-and-publish, build-frontend, publish-create-release]
runs-on: ubuntu-latest
permissions:
contents: write
strategy:
fail-fast: false
matrix:
include:
- goos: linux
goarch: amd64
ext: ""
- goos: linux
goarch: arm64
ext: ""
- goos: windows
goarch: amd64
ext: .exe
- goos: windows
goarch: arm64
ext: .exe
- goos: darwin
goarch: amd64
ext: ""
- goos: darwin
goarch: arm64
ext: ""
steps:
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
with:
go-version: '1.25'
cache: true
- name: Download web/dist
uses: actions/download-artifact@v4
with:
name: web-dist
path: web/dist
- name: Build binary
run: |
CGO_ENABLED=0 GOOS=${{ matrix.goos }} GOARCH=${{ matrix.goarch }} \
go build -trimpath -ldflags="-s -w -X main.version=${{ needs.version-and-publish.outputs.tag }}" \
-o "dist/mmtl-${{ matrix.goos }}-${{ matrix.goarch }}${{ matrix.ext }}" ./cmd/server
- name: Package
run: |
mkdir -p package/mmtl
cp "dist/mmtl-${{ matrix.goos }}-${{ matrix.goarch }}${{ matrix.ext }}" package/mmtl/mmtl${{ matrix.ext }}
cp README.md package/mmtl/ 2>/dev/null || true
if [ "${{ matrix.goos }}" = "windows" ]; then
(cd package && zip -r "../mmtl_${{ matrix.goos }}_${{ matrix.goarch }}.zip" mmtl)
else
tar -czf "mmtl_${{ matrix.goos }}_${{ matrix.goarch }}.tar.gz" -C package mmtl
fi
- name: Upload to GitHub Release
env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
RELEASE_TAG: ${{ needs.version-and-publish.outputs.tag }}
run: |
set -eux
PKG="mmtl_${{ matrix.goos }}_${{ matrix.goarch }}.zip"
TAR="mmtl_${{ matrix.goos }}_${{ matrix.goarch }}.tar.gz"
# 并发上传到同一 release 各自文件,--clobber 幂等覆盖
if [ -f "$PKG" ]; then
for i in 1 2 3; do gh release upload "$RELEASE_TAG" "$PKG" --clobber && break || sleep 5; done
fi
if [ -f "$TAR" ]; then
for i in 1 2 3; do gh release upload "$RELEASE_TAG" "$TAR" --clobber && break || sleep 5; done
fi
+13
View File
@@ -18,6 +18,19 @@ jobs:
go-version: '1.25'
cache: true
# The binary embeds the SPA (web/dist) via go:embed, so the dist must exist
# before the Go toolchain touches the `web` package.
- uses: actions/setup-node@v4
with:
node-version: '20'
cache: 'npm'
cache-dependency-path: web/package-lock.json
- name: Build SPA
working-directory: web
run: |
npm ci
npm run build
- name: go vet
run: go vet ./...
+22 -5
View File
@@ -360,15 +360,17 @@ docker compose -f docker-compose.search.yml up -d --no-deps mmtl
本地开发需要 Go、Node.js 和 npm。
后端会将 `web/dist` 通过 `go:embed` 编进二进制,因此**在编译 / 运行后端之前要先构建前端**,否则 `web` 包会因为缺少嵌入资源而编译失败。
```bash
# 前端依赖与构建(必须先做,产物被 go:embed 打进二进制)
npm --prefix web ci
npm --prefix web run build
# 后端测试
go test ./...
# 前端依赖与构建
npm --prefix web install
npm --prefix web run build
# 本地运行后端
# 本地运行后端(二进制自带前端界面,无需额外 web 目录)
go run ./cmd/server
# 本地运行前端开发服务器
@@ -387,6 +389,21 @@ http://127.0.0.1:3000
http://127.0.0.1:8080/api/health
```
### 交叉编译单文件发布物
CI(`.github/workflows/Auto-docker-publish.yml`)每次发布会自动为 Windows / Linux(含 Debian) / macOS 交叉编译 amd64 + arm64 的单文件可执行程序,并上传到对应的 GitHub Release。你可以在 Releases 页面下载 `.zip`(Windows)或 `.tar.gz`(Linux / macOS)附件,解压后直接运行其中的 `mmtl`(Windows 为 `mmtl.exe`),无需额外携带前端目录。
本地手动交叉编译某个平台:
```bash
# 先构建前端
npm --prefix web ci && npm --prefix web run build
# 例如:构建 Linux amd64 单文件
CGO_ENABLED=0 GOOS=linux GOARCH=amd64 \
go build -trimpath -ldflags="-s -w" -o mmtl-linux-amd64 ./cmd/server
```
## 贡献与反馈
提交 Bug、功能建议或 Pull Request 前,请先阅读 [贡献规范](CONTRIBUTING.md)。
+1 -1
View File
@@ -1 +1 @@
0.0.39
0.0.46
+3 -3
View File
@@ -52,7 +52,7 @@ func TestServeSPANoCachesIndexAndServesRoutes(t *testing.T) {
}
router := gin.New()
serveSPA(router, webDir)
serveSPA(router, os.DirFS(webDir))
for _, path := range []string{"/", "/login", "/library/e1c3507e-2878-40ae-a0e1-6b6e44b7fa7a", "/media/abc"} {
req := httptest.NewRequest(http.MethodGet, path, nil)
@@ -93,7 +93,7 @@ func TestServeSPAServesAssetsImmutableAndBypassesAPIRoutes(t *testing.T) {
}
router := gin.New()
serveSPA(router, webDir)
serveSPA(router, os.DirFS(webDir))
assetReq := httptest.NewRequest(http.MethodGet, "/assets/app.js", nil)
assetResp := httptest.NewRecorder()
@@ -155,7 +155,7 @@ func TestServeSPAServesAssetsImmutableAndBypassesAPIRoutes(t *testing.T) {
func TestServeSPAMissingIndexReportsExplicit404(t *testing.T) {
gin.SetMode(gin.TestMode)
router := gin.New()
serveSPA(router, t.TempDir())
serveSPA(router, os.DirFS(t.TempDir()))
req := httptest.NewRequest(http.MethodGet, "/", nil)
w := httptest.NewRecorder()
+59 -23
View File
@@ -1,6 +1,8 @@
package main
import (
"io/fs"
"mime"
"net/http"
"os"
"path/filepath"
@@ -13,6 +15,8 @@ import (
"github.com/ShukeBta/MMTL/internal/handler"
"github.com/ShukeBta/MMTL/internal/middleware"
"github.com/ShukeBta/MMTL/internal/service"
"github.com/ShukeBta/MMTL/web"
)
func buildRouter(cfg *config.Config, logger *zap.Logger, svc *service.Container) *gin.Engine {
@@ -29,31 +33,41 @@ func buildRouter(cfg *config.Config, logger *zap.Logger, svc *service.Container)
handler.Register(r, cfg, logger, svc)
if cfg.App.WebDir != "" {
serveSPA(r, cfg.App.WebDir)
// Prefer a directory on disk when configured explicitly (e.g. the Docker image
// mounts web/dist from the build stage, or an operator overrides app.web_dir
// with a custom skin). Otherwise fall back to the SPA embedded into the binary,
// which is what makes the cross-platform single-file artifacts work.
uiFS := webui.DistFS()
if dir := cfg.App.WebDir; dir != "" {
disk := os.DirFS(dir)
if index, err := fs.Stat(disk, "index.html"); err == nil && !index.IsDir() {
uiFS = disk
}
}
serveSPA(r, uiFS)
return r
}
// serveSPA serves the React build artifacts and falls back to index.html for
// non-API, non-asset paths so client-side routing keeps working.
func serveSPA(r *gin.Engine, webDir string) {
// non-API, non-asset paths so client-side routing keeps working. The UI tree
// comes from root, which is either the compiled-in SPA or an on-disk web dir.
func serveSPA(r *gin.Engine, root fs.FS) {
assets := r.Group("/assets")
assets.Use(func(c *gin.Context) {
c.Header("Cache-Control", "public, max-age=31536000, immutable")
c.Next()
})
assets.Static("/", filepath.Join(webDir, "assets"))
assets.GET("/*filepath", serveFSDir(root, "assets"))
brand := r.Group("/brand")
brand.Use(func(c *gin.Context) {
setNoCacheHeaders(c)
c.Next()
})
brand.Static("/", filepath.Join(webDir, "brand"))
brand.GET("/*filepath", serveFSDir(root, "brand"))
for _, rootFile := range []string{"/favicon.ico", "/favicon.svg", "/artwork-cache-sw.js"} {
filePath := filepath.Join(webDir, strings.TrimPrefix(rootFile, "/"))
r.GET(rootFile, serveNoCacheFile(filePath))
r.HEAD(rootFile, serveNoCacheFile(filePath))
name := strings.TrimPrefix(rootFile, "/")
r.GET(rootFile, serveFSFile(root, name))
r.HEAD(rootFile, serveFSFile(root, name))
}
r.NoRoute(func(c *gin.Context) {
path := c.Request.URL.Path
@@ -61,28 +75,50 @@ func serveSPA(r *gin.Engine, webDir string) {
c.Status(http.StatusNotFound)
return
}
serveSPAIndex(c, filepath.Join(webDir, "index.html"))
setNoCacheHeaders(c)
data, err := fs.ReadFile(root, "index.html")
if err != nil {
c.String(http.StatusNotFound, "MMTL web UI not found")
return
}
c.Data(http.StatusOK, "text/html; charset=utf-8", data)
})
}
func serveNoCacheFile(filePath string) gin.HandlerFunc {
// serveFSDir serves a static subdirectory of root. A missing asset returns 404.
func serveFSDir(root fs.FS, dir string) gin.HandlerFunc {
sub, err := fs.Sub(root, dir)
if err != nil {
return func(c *gin.Context) { c.Status(http.StatusNotFound) }
}
handler := http.StripPrefix("/"+dir, http.FileServerFS(sub))
return func(c *gin.Context) {
setNoCacheHeaders(c)
if _, err := os.Stat(filePath); err != nil {
c.Status(http.StatusNotFound)
return
}
c.File(filePath)
handler.ServeHTTP(c.Writer, c.Request)
}
}
func serveSPAIndex(c *gin.Context, indexPath string) {
setNoCacheHeaders(c)
if _, err := os.Stat(indexPath); err != nil {
c.String(http.StatusNotFound, "MMTL web UI not found: %s", indexPath)
return
// serveFSFile serves a single root-level file (favicon / service worker) with
// no-cache headers. It reads from root, which may be the embedded SPA or disk.
func serveFSFile(root fs.FS, name string) gin.HandlerFunc {
return func(c *gin.Context) {
setNoCacheHeaders(c)
data, err := fs.ReadFile(root, name)
if err != nil {
c.Status(http.StatusNotFound)
return
}
c.Data(http.StatusOK, mimeTypeByName(name), data)
}
}
// mimeTypeByName returns an HTTP content type guessed from a file extension.
func mimeTypeByName(name string) string {
switch mime.TypeByExtension(filepath.Ext(name)) {
case "":
return "application/octet-stream"
default:
return mime.TypeByExtension(filepath.Ext(name))
}
c.File(indexPath)
}
func setNoCacheHeaders(c *gin.Context) {
+2
View File
@@ -43,6 +43,8 @@ func registerAdminStrmRoutes(admin *gin.RouterGroup, svc *service.Container) {
admin.POST("/strm/paths/:id/sync", startStrmSyncHandler(svc))
admin.POST("/strm/paths/:id/cancel", cancelStrmSyncHandler(svc))
admin.GET("/strm/records", listStrmSyncRecordsHandler(svc))
admin.DELETE("/strm/records/:id", deleteStrmSyncRecordHandler(svc))
admin.DELETE("/strm/records", clearStrmSyncRecordsHandler(svc))
admin.GET("/strm/local-dirs", listStrmLocalDirsHandler(svc))
admin.GET("/strm/downloads", downloadQueueHandler(svc))
+25
View File
@@ -300,6 +300,31 @@ func listStrmSyncRecordsHandler(svc *service.Container) gin.HandlerFunc {
}
}
func deleteStrmSyncRecordHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
if c.Param("id") == "" {
c.JSON(http.StatusBadRequest, gin.H{"error": "缺少记录 ID"})
return
}
if err := svc.Strm.DeleteSyncRecord(c.Request.Context(), c.Param("id")); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"ok": true})
}
}
func clearStrmSyncRecordsHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
deleted, err := svc.Strm.ClearSyncRecords(c.Request.Context(), c.Query("path_id"))
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"ok": true, "deleted": deleted})
}
}
// ─── 下载/上传队列 ─────────────────────────────────────────────────────────────
func downloadQueueHandler(svc *service.Container) gin.HandlerFunc {
+18
View File
@@ -170,6 +170,24 @@ func (r *StrmSyncRecordRepository) List(ctx context.Context, syncPathID string,
return rows, err
}
// Delete 删除单条同步记录(物理删除)。
func (r *StrmSyncRecordRepository) Delete(ctx context.Context, id string) error {
return withSQLiteBusyRetry(ctx, func() error {
return r.db.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.StrmSyncRecord{}).Error
})
}
// DeleteBySyncPathID 删除某同步目录下的全部同步记录(删除同步目录时级联清理)。
func (r *StrmSyncRecordRepository) DeleteBySyncPathID(ctx context.Context, syncPathID string) (int64, error) {
var count int64
err := withSQLiteBusyRetry(ctx, func() error {
res := r.db.WithContext(ctx).Unscoped().Where("sync_path_id = ?", syncPathID).Delete(&model.StrmSyncRecord{})
count = res.RowsAffected
return res.Error
})
return count, err
}
// ─── StrmDownloadTask ──────────────────────────────────────────────────────────
// StrmDownloadTaskRepository persists model.StrmDownloadTask.
+1 -1
View File
@@ -281,7 +281,7 @@ func IsThrottleCode(code int) bool {
func isTokenCode(code int) bool {
switch code {
case AccessTokenAuthFail, AccessAuthInvalid, AccessTokenExpiryCode, RefreshTokenInvalid:
case AccessTokenAuthFail, AccessAuthInvalid, AccessTokenExpiryCode, AccessTokenFormatInvalid, RefreshTokenInvalid:
return true
}
return false
+8 -7
View File
@@ -21,13 +21,14 @@ var (
const (
// 业务错误码
AccessTokenAuthFail = 40140126 // 访问过期,需刷新
AccessTokenExpiryCode = 40140125 // 访问过期,需刷新
AccessAuthInvalid = 40140124 // 访问无效,需刷新
RefreshTokenInvalid = 40140116 // 需重新授权
TokenRefreshFail = 40140121 // 刷新失败,可重试
RequestMaxLimitCode = 770004 // 访问频率过高
RequestRateLimitCode = 406 // 达到访问上限
AccessTokenAuthFail = 40140126 // 访问过期,需刷新
AccessTokenExpiryCode = 40140125 // 访问过期,需刷新
AccessAuthInvalid = 40140124 // 访问无效,需刷新
AccessTokenFormatInvalid = 40140123 // access_token 格式错误,需刷新
RefreshTokenInvalid = 40140116 // 需重新授权
TokenRefreshFail = 40140121 // 刷新失败,可重试
RequestMaxLimitCode = 770004 // 访问频率过高
RequestRateLimitCode = 406 // 达到访问上限
// 刷新 token 的错误码
RefreshTokenFormatInvalid = 40140114
+22 -2
View File
@@ -5,6 +5,8 @@ package cloud115
import (
"context"
"encoding/base64"
"encoding/json"
"errors"
"fmt"
"io"
@@ -181,6 +183,24 @@ func (u *OSSMultipartUploader) UploadFileWithResult(ctx context.Context, input O
return completeParts[i].PartNumber < completeParts[j].PartNumber
})
// 115 下发的 callback / callback_var 是 JSON 字符串,而 OSS CompleteMultipartUpload
// 要求 callback 参数为 Base64 编码后的 JSON,否则报 "The callback configuration is
// not base64 encoded"。这里把两者转为 Base64 后再提交(参考 QMediaSync 的
// BuildOSSCallbackHeaders)。
cb := input.Callback
cbVar := input.CallbackVar
if cb == "" {
return OSSMultipartUploadResult{}, errors.New("OSS callback 为空")
}
if !json.Valid([]byte(cb)) {
return OSSMultipartUploadResult{}, errors.New("解析 callback 失败:不是合法 JSON")
}
if cbVar == "" {
cbVar = "{}"
}
if !json.Valid([]byte(cbVar)) {
return OSSMultipartUploadResult{}, errors.New("解析 callback_var 失败:不是合法 JSON")
}
completeResult, err := u.client.CompleteMultipartUpload(ctx, &oss.CompleteMultipartUploadRequest{
Bucket: oss.Ptr(input.Bucket),
Key: oss.Ptr(input.Object),
@@ -188,8 +208,8 @@ func (u *OSSMultipartUploader) UploadFileWithResult(ctx context.Context, input O
CompleteMultipartUpload: &oss.CompleteMultipartUpload{
Parts: completeParts,
},
Callback: oss.Ptr(input.Callback),
CallbackVar: oss.Ptr(input.CallbackVar),
Callback: oss.Ptr(base64.StdEncoding.EncodeToString([]byte(cb))),
CallbackVar: oss.Ptr(base64.StdEncoding.EncodeToString([]byte(cbVar))),
})
if err != nil {
return OSSMultipartUploadResult{}, fmt.Errorf("完成 OSS multipart 失败:%w", err)
+3 -1
View File
@@ -40,7 +40,9 @@ func FileSHA1Partial(path string, start, end int64) (string, error) {
}
length := end - start + 1
h := sha1.New()
if _, err := io.CopyN(h, f, length); err != nil {
// io.CopyN 在文件不足 length 字节时会返回 io.EOF,导致小文件(如小于 128 KiB 的
// 元数据图片)无法上传。这里只拷贝实际读到的字节,文件尾对齐到区间终点即可。
if _, err := io.CopyN(h, f, length); err != nil && err != io.EOF {
return "", err
}
return hex.EncodeToString(h.Sum(nil)), nil
+103
View File
@@ -1,9 +1,14 @@
package cloud115
import (
"context"
"encoding/base64"
"io"
"os"
"path/filepath"
"testing"
"github.com/aliyun/alibabacloud-oss-go-sdk-v2/oss"
)
func TestFileSHA1(t *testing.T) {
@@ -39,6 +44,26 @@ func TestFileSHA1Partial(t *testing.T) {
}
}
// TestFileSHA1PartialSmallerThanWindow 回归测试:经典 bug 是 io.CopyN 在文件不足
// length 字节时返回 io.EOF。115 上传固定用 [0,128*1024-1] 窗口计算 preid,导致所有
// 小于 128 KiB 的元数据文件(如海报/缩略图)上传必然失败。
func TestFileSHA1PartialSmallerThanWindow(t *testing.T) {
dir := t.TempDir()
path := filepath.Join(dir, "small.bin")
// 6 字节小文件,不足 128 KiB 窗口
if err := os.WriteFile(path, []byte("abcdef"), 0o644); err != nil {
t.Fatal(err)
}
sum, err := FileSHA1Partial(path, 0, 128*1024-1)
if err != nil {
t.Fatalf("compute partial sha1 for small file should not fail: %v", err)
}
// 应等于整个文件(6 字节)的 sha1
if sum != "1f8ac10f23c5b5bc1167bda84b833e5c057a77d2" {
t.Errorf("unexpected partial sha1: %s", sum)
}
}
func TestParseSignCheckRange(t *testing.T) {
rng, err := parseSignCheckRange("0-131071")
if err != nil {
@@ -92,3 +117,81 @@ func TestBaseNameOf(t *testing.T) {
t.Errorf("got %s", got)
}
}
// fakeCallbackOSSClient 捕获 CompleteMultipartUpload 收到的 callback / callback_var,
// 用于断言已经 Base64 编码(116 要求 callback 必须是 Base64 后的 JSON,否则报
// "The callback configuration is not base64 encoded")。
type fakeCallbackOSSClient struct {
capturedCallback string
capturedCallbackVar string
}
func (c *fakeCallbackOSSClient) InitiateMultipartUpload(_ context.Context, _ *oss.InitiateMultipartUploadRequest, _ ...func(*oss.Options)) (*oss.InitiateMultipartUploadResult, error) {
return &oss.InitiateMultipartUploadResult{UploadId: oss.Ptr("upload-new")}, nil
}
func (c *fakeCallbackOSSClient) UploadPart(_ context.Context, r *oss.UploadPartRequest, _ ...func(*oss.Options)) (*oss.UploadPartResult, error) {
if r.Body != nil {
_, _ = io.Copy(io.Discard, r.Body)
}
return &oss.UploadPartResult{ETag: oss.Ptr("etag-1")}, nil
}
func (c *fakeCallbackOSSClient) ListParts(context.Context, *oss.ListPartsRequest, ...func(*oss.Options)) (*oss.ListPartsResult, error) {
return &oss.ListPartsResult{}, nil
}
func (c *fakeCallbackOSSClient) CompleteMultipartUpload(_ context.Context, r *oss.CompleteMultipartUploadRequest, _ ...func(*oss.Options)) (*oss.CompleteMultipartUploadResult, error) {
c.capturedCallback = *r.Callback
c.capturedCallbackVar = *r.CallbackVar
return &oss.CompleteMultipartUploadResult{
CallbackResult: map[string]any{
"state": true,
"data": map[string]any{"file_id": "file-1", "pick_code": "pick-1"},
},
}, nil
}
func (c *fakeCallbackOSSClient) AbortMultipartUpload(context.Context, *oss.AbortMultipartUploadRequest, ...func(*oss.Options)) (*oss.AbortMultipartUploadResult, error) {
return &oss.AbortMultipartUploadResult{}, nil
}
// TestCompleteMultipartUploadCallbackBase64 回归测试:OSS CompleteMultipartUpload 的
// callback 必须 Base64 编码,否则报 "The callback configuration is not base64 encoded",
// 导致大于 128 KiB 的元数据文件上传失败。
func TestCompleteMultipartUploadCallbackBase64(t *testing.T) {
dir := t.TempDir()
path := filepath.Join(dir, "big.bin")
data := make([]byte, 8) // 8 字节,PartSize=8 → 1 part
if err := os.WriteFile(path, data, 0o644); err != nil {
t.Fatal(err)
}
fake := &fakeCallbackOSSClient{}
uploader := &OSSMultipartUploader{client: fake}
callback := `{"callbackUrl":"http://uplb.115.com/3.0/completeupload.php"}`
callbackVar := `{"x:pick_code":"abc"}`
_, err := uploader.UploadFileWithResult(context.Background(), OSSMultipartUploadInput{
Bucket: "bucket-1",
Object: "object-1",
Callback: callback,
CallbackVar: callbackVar,
FilePath: path,
FileSize: int64(len(data)),
PartSize: 8,
})
if err != nil {
t.Fatalf("multipart 上传失败:%v", err)
}
// 捕获的 callback 必须是合法 Base64,且解码后与原 JSON 一致
cbBytes, err := base64.StdEncoding.DecodeString(fake.capturedCallback)
if err != nil {
t.Fatalf("callback 未 Base64 编码:%v (raw=%q)", err, fake.capturedCallback)
}
if string(cbBytes) != callback {
t.Errorf("callback 解码后 = %s,期望 %s", cbBytes, callback)
}
cbvBytes, err := base64.StdEncoding.DecodeString(fake.capturedCallbackVar)
if err != nil {
t.Fatalf("callback_var 未 Base64 编码:%v (raw=%q)", err, fake.capturedCallbackVar)
}
if string(cbvBytes) != callbackVar {
t.Errorf("callback_var 解码后 = %s,期望 %s", cbvBytes, callbackVar)
}
}
+77 -21
View File
@@ -49,7 +49,7 @@ func findMovieNFO(mediaPath, libraryRoot string) (*nfoDocument, string, error) {
continue
}
seen[key] = struct{}{}
if doc, _, err := readNFO(path); err == nil {
if doc, _, err := decodeNFOFile(path); err == nil && doc != nil {
return doc, path, nil
} else if err != nil && !errors.Is(err, os.ErrNotExist) {
return nil, "", err
@@ -62,7 +62,7 @@ func findMovieNFO(mediaPath, libraryRoot string) (*nfoDocument, string, error) {
for _, match := range matches {
baseKey := strings.ToLower(strings.ReplaceAll(strings.TrimSuffix(filepath.Base(match), filepath.Ext(match)), "-", ""))
if strings.Contains(baseKey, codeKey) || strings.Contains(codeKey, baseKey) {
if doc, _, err := readNFO(match); err == nil {
if doc, _, err := decodeNFOFile(match); err == nil && doc != nil {
return doc, match, nil
} else if err != nil && !errors.Is(err, os.ErrNotExist) {
return nil, "", err
@@ -71,7 +71,7 @@ func findMovieNFO(mediaPath, libraryRoot string) (*nfoDocument, string, error) {
}
}
if len(matches) == 1 {
if doc, _, err := readNFO(matches[0]); err == nil {
if doc, _, err := decodeNFOFile(matches[0]); err == nil && doc != nil {
return doc, matches[0], nil
} else if err != nil && !errors.Is(err, os.ErrNotExist) {
return nil, "", err
@@ -84,21 +84,33 @@ func findMovieNFO(mediaPath, libraryRoot string) (*nfoDocument, string, error) {
func readSeriesMetadata(mediaPath, libraryRoot string) (*LocalMetadata, error) {
var meta *LocalMetadata
showBaseDir := ""
if showDoc, showPath, err := findShowNFO(mediaPath, libraryRoot); err == nil && showDoc != nil {
if showDoc, showPath, partial, err := findShowNFO(mediaPath, libraryRoot); err == nil && showDoc != nil {
showBaseDir = filepath.Dir(showPath)
meta = metadataFromDoc(showDoc, showBaseDir, true)
// A truncated show NFO still yields a usable title/poster (decodePartialNFO
// only surfaces docs with recoverable fields). Keep HasNFO so the recovered
// title participates in series grouping; a fully partial episode match
// below can still demote it.
meta.HasNFO = !partial || meta.HasNFO
} else if err != nil && !errors.Is(err, os.ErrNotExist) {
return nil, err
// A damaged show NFO must not discard the whole series: readLocalScanMetadata
// treats this as a hard failure and leaves every episode without local
// metadata. When the show level fails we still merge the episode NFO and
// local artwork below instead of returning early.
meta = nil
}
if episodeDoc, episodePath, err := readNFO(nfoPath(mediaPath)); err == nil {
if episodeDoc, episodePath, episodePartial, err := readEpisodeNFO(nfoPath(mediaPath)); err == nil && episodeDoc != nil {
episodeMeta := metadataFromDoc(episodeDoc, filepath.Dir(episodePath), true)
if meta == nil {
meta = &LocalMetadata{}
}
mergeEpisodeMetadata(meta, episodeMeta, episodeDoc)
} else if err != nil && !errors.Is(err, os.ErrNotExist) {
return nil, err
// Only a fully valid episode NFO can confirm the match; a truncated one
// keeps whatever the show NFO supplied but must not force scrape_status.
if episodePartial {
meta.HasNFO = false
}
}
if meta == nil {
meta = metadataFromArtwork(mediaPath, showBaseDir)
@@ -108,19 +120,49 @@ func readSeriesMetadata(mediaPath, libraryRoot string) (*LocalMetadata, error) {
return meta, nil
}
func readNFO(path string) (*nfoDocument, string, error) {
// readEpisodeNFO reads the sidecar NFO next to an episode file and reports both
// the recovered document and whether it came from a truncated/malformed file.
func readEpisodeNFO(path string) (*nfoDocument, string, bool, error) {
body, err := os.ReadFile(path) // #nosec G304 -- path is a discovered NFO sidecar under the configured library root.
if err != nil {
return nil, "", err
return nil, "", false, err
}
var doc nfoDocument
if err := xml.Unmarshal(body, &doc); err != nil {
return nil, "", err
doc, partial, err := decodePartialNFO(body)
if err != nil {
return nil, "", false, err
}
return &doc, path, nil
return doc, path, partial, nil
}
func findShowNFO(mediaPath, libraryRoot string) (*nfoDocument, string, error) {
func decodePartialNFO(body []byte) (*nfoDocument, bool, error) {
var doc nfoDocument
err := xml.Unmarshal(body, &doc)
if err == nil {
return &doc, false, nil
}
if !isLikelyTruncatedXMLError(err) {
// A real parse error unrelated to truncation: don't trust partial fields.
return nil, false, err
}
// Unmarshal still fills the elements it closed before hitting the cut; treat
// those as partial show metadata instead of dropping everything.
if doc.Title == "" && doc.OriginalTitle == "" && len(doc.Thumbs) == 0 &&
doc.Premiered == "" && doc.Plot == "" {
return nil, false, err
}
return &doc, true, nil
}
func isLikelyTruncatedXMLError(err error) bool {
for _, msg := range []string{"unexpected EOF", "EOF"} {
if strings.Contains(strings.ToLower(err.Error()), strings.ToLower(msg)) {
return true
}
}
return false
}
func findShowNFO(mediaPath, libraryRoot string) (*nfoDocument, string, bool, error) {
dir := filepath.Dir(mediaPath)
root := filepath.Clean(libraryRoot)
for {
@@ -133,19 +175,33 @@ func findShowNFO(mediaPath, libraryRoot string) (*nfoDocument, string, error) {
names = append(names, base+".nfo")
for _, name := range names {
path := filepath.Join(dir, name)
if doc, _, err := readNFO(path); err == nil {
return doc, path, nil
} else if err != nil && !errors.Is(err, os.ErrNotExist) {
return nil, "", err
doc, partial, err := decodeNFOFile(path)
if doc != nil {
return doc, path, partial, nil
}
if err != nil && !errors.Is(err, os.ErrNotExist) {
return nil, "", false, err
}
}
if samePath(dir, root) {
return nil, "", os.ErrNotExist
return nil, "", false, os.ErrNotExist
}
parent := filepath.Dir(dir)
if parent == dir {
return nil, "", os.ErrNotExist
return nil, "", false, os.ErrNotExist
}
dir = parent
}
}
func decodeNFOFile(path string) (*nfoDocument, bool, error) {
body, err := os.ReadFile(path) // #nosec G304 -- path is a discovered NFO sidecar under the configured library root.
if err != nil {
return nil, false, err
}
doc, partial, err := decodePartialNFO(body)
if err != nil {
return nil, false, err
}
return doc, partial, nil
}
+108
View File
@@ -267,3 +267,111 @@ func TestReadLocalMetadataWithoutNFOStillFindsArtwork(t *testing.T) {
t.Fatalf("unexpected artwork metadata: %+v", got)
}
}
// TestReadLocalMetadataRecoversFromTruncatedShowNFO mirrors a real-world issue:
// some anime tvshow.nfo files are truncated mid-URL (unexpected EOF). Before the
// fix this discarded the whole series, leaving episodes pending with per-episode
// titles and no poster. The recoverable fields (title/year) and the matching
// episode NFO + local artwork must still be applied.
func TestReadLocalMetadataRecoversFromTruncatedShowNFO(t *testing.T) {
root := t.TempDir()
showDir := filepath.Join(root, "夏日重现 (2022)")
seasonDir := filepath.Join(showDir, "Season 1")
if err := os.MkdirAll(seasonDir, 0o755); err != nil {
t.Fatal(err)
}
mediaPath := filepath.Join(seasonDir, "S01E01.mkv")
if err := os.WriteFile(mediaPath, []byte("x"), 0o644); err != nil {
t.Fatal(err)
}
// tvshow.nfo cut off inside a <thumb> URL, like the broken real files.
truncated := `<?xml version="1.0" encoding="UTF-8" standalone="yes"?>
<tvshow>
<title>夏日重现</title>
<originaltitle>サマータイムレンダ</originaltitle>
<year>2022</year>
<plot>听闻自己青梅竹马死讯。</plot>
<thumb aspect="poster">https://image.tmdb.org/t/p/original/2koyWLm6iVn5OTEExTjKzVms5Iz.jpg</thumb>
<fanart>
<thumb>https://image.tmdb.org/t/p/original/p2eZlGwd8OjkWpwD2hSoBiIlHBZ.jpg</thu`
if err := os.WriteFile(filepath.Join(showDir, "tvshow.nfo"), []byte(truncated), 0o644); err != nil {
t.Fatal(err)
}
// The episode sidecar NFO is complete and should still be merged.
if err := os.WriteFile(nfoPath(mediaPath), []byte(`<episodedetails><title>再见了夏日</title><season>1</season><episode>1</episode></episodedetails>`), 0o644); err != nil {
t.Fatal(err)
}
// Local artwork next to the show folder.
poster := filepath.Join(showDir, "poster.jpg")
if err := os.WriteFile(poster, []byte("jpg"), 0o644); err != nil {
t.Fatal(err)
}
got, err := ReadLocalMetadata(mediaPath, root, true)
if err != nil {
t.Fatalf("ReadLocalMetadata returned error on truncated show NFO: %v", err)
}
if got == nil {
t.Fatal("metadata is nil; truncated show NFO discarded the series")
}
// The recovered show title takes precedence over the episode title.
if got.Title != "夏日重现" {
t.Fatalf("Title = %q, want recovered show title 夏日重现", got.Title)
}
if got.Year != 2022 {
t.Fatalf("Year = %d, want 2022", got.Year)
}
if got.EpisodeTitle != "再见了夏日" || got.SeasonNum != 1 || got.EpisodeNum != 1 {
t.Fatalf("episode metadata not preserved: %+v", got)
}
// Prior to the fix the episode metadata was dropped entirely; the poster comes
// from the local poster.jpg next to the show.
if got.PosterURL != poster {
t.Fatalf("PosterURL = %q, want local poster %q", got.PosterURL, poster)
}
// A truncated show NFO must not be treated as an authoritative match by
// itself; since the episode NFO is valid we still mark it matched so the
// recovered series title participates in grouping.
if !got.HasNFO {
t.Fatalf("HasNFO = false, want true (episode NFO is valid)")
}
}
// TestReadLocalMetadataKeepsArtworkOnlyWhenNoUsableNFO verifies that when a show
// NFO is truncated AND yields no recoverable fields, we still fall back to local
// artwork instead of returning an error.
func TestReadLocalMetadataArtworkFallbackOnGarbageShowNFO(t *testing.T) {
root := t.TempDir()
showDir := filepath.Join(root, "Some Show")
seasonDir := filepath.Join(showDir, "Season 1")
if err := os.MkdirAll(seasonDir, 0o755); err != nil {
t.Fatal(err)
}
mediaPath := filepath.Join(seasonDir, "S01E01.mkv")
if err := os.WriteFile(mediaPath, []byte("x"), 0o644); err != nil {
t.Fatal(err)
}
// Severely truncated: no recoverable title/fields at all.
if err := os.WriteFile(filepath.Join(showDir, "tvshow.nfo"), []byte(`<tvshow><title>半截`), 0o644); err != nil {
t.Fatal(err)
}
// Episode-level backdrop sits next to the episode file.
backdrop := filepath.Join(seasonDir, "S01E01-backdrop.jpg")
if err := os.WriteFile(backdrop, []byte("jpg"), 0o644); err != nil {
t.Fatal(err)
}
got, err := ReadLocalMetadata(mediaPath, root, true)
if err != nil {
t.Fatalf("ReadLocalMetadata should not error when NFO is unusable: %v", err)
}
if got == nil {
t.Fatal("metadata is nil; artwork fallback missing")
}
if got.HasNFO {
t.Fatalf("HasNFO = true, want false for unusable NFO")
}
if got.BackdropURL != backdrop {
t.Fatalf("expected episode artwork fallback, got %+v", got)
}
}
+32
View File
@@ -449,6 +449,38 @@ func (s *StrmService) ListSyncRecords(ctx context.Context, pathID string, limit
return s.repo.StrmSyncRecord.List(ctx, pathID, limit)
}
// DeleteSyncRecord 删除单条同步记录。
func (s *StrmService) DeleteSyncRecord(ctx context.Context, id string) error {
if err := s.repo.StrmSyncRecord.Delete(ctx, id); err != nil {
return err
}
return nil
}
// ClearSyncRecords 清空某同步目录(pathID 为空则全部)的同步记录,返回删除条数。
func (s *StrmService) ClearSyncRecords(ctx context.Context, pathID string) (int64, error) {
if pathID != "" {
return s.repo.StrmSyncRecord.DeleteBySyncPathID(ctx, pathID)
}
var total int64
// 全量清空:分页拉取物理删除所有记录
for {
rows, err := s.repo.StrmSyncRecord.List(ctx, "", 200)
if err != nil {
return total, err
}
if len(rows) == 0 {
return total, nil
}
for _, rec := range rows {
if err := s.repo.StrmSyncRecord.Delete(ctx, rec.ID); err != nil {
return total, err
}
}
total += int64(len(rows))
}
}
// CreateSyncPath 校验并创建同步目录。
func (s *StrmService) CreateSyncPath(ctx context.Context, p *model.StrmSyncPath) (*model.StrmSyncPath, error) {
if err := s.validateSyncPath(ctx, p); err != nil {
+28 -12
View File
@@ -235,8 +235,8 @@ func (s *StrmService) finishSync(p *model.StrmSyncPath, rec *model.StrmSyncRecor
if rec.SyncType == model.StrmSyncTypeFull {
syncTypeLabel = "全量"
}
p.LastSyncMessage = fmt.Sprintf("[%s] 完成:新增/更新 %d 个 strm,跳过 %d 个,下载 %d 个元数据,清理 %d 个文件",
syncTypeLabel, rec.NewStrm, rec.Skipped, rec.NewMeta, rec.Pruned)
p.LastSyncMessage = fmt.Sprintf("[%s] 完成:新增/更新 %d 个 strm,跳过 %d 个,下载 %d 个元数据,上传 %d 个元数据,清理 %d 个文件",
syncTypeLabel, rec.NewStrm, rec.Skipped, rec.NewMeta, rec.Uploaded, rec.Pruned)
}
if err := s.repo.StrmSyncPath.Update(context.Background(), p); err != nil {
s.log.Warn("update strm sync path failed", zap.Error(err))
@@ -244,7 +244,7 @@ func (s *StrmService) finishSync(p *model.StrmSyncPath, rec *model.StrmSyncRecor
s.log.Info("strm sync finished",
zap.String("path_id", p.ID), zap.String("sync_type", rec.SyncType), zap.String("status", status),
zap.Int64("new_strm", rec.NewStrm), zap.Int64("skipped", rec.Skipped), zap.Int64("new_meta", rec.NewMeta),
zap.Int64("pruned", rec.Pruned), zap.String("message", message))
zap.Int64("uploaded", rec.Uploaded), zap.Int64("pruned", rec.Pruned), zap.String("message", message))
}
func (st *strmSyncState) run() error {
@@ -483,7 +483,8 @@ func cleanDirRel(rel string) string {
// 极大地降低 API 请求次数并支持毫秒级/秒级增量同步。
func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error {
defer st.flushPendingDownloads()
ctx := st.ctx
ctx, cancel := context.WithCancel(st.ctx)
defer cancel()
rootCID := strings.TrimSpace(st.p.RemotePath)
if rootCID == "" {
rootCID = "0"
@@ -618,6 +619,8 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error {
dirWorkers = 8
doneDirs atomic.Int64
totalDirs = len(pidList)
errMu sync.Mutex
firstErr error
)
if len(pidList) < dirWorkers {
dirWorkers = len(pidList)
@@ -641,10 +644,19 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error {
}
detail, err := open115.GetFsDetailByCid(ctx, pid)
if err != nil {
st.s.log.Warn("115: 获取目录详情失败", zap.String("pid", pid), zap.Error(err))
// 目录详情解析失败会导致下游文件 rel 无法还原真实父路径,
// seen key 与磁盘路径对不上,增量 prune 会误删本地文件,标记本次扫描不完整。
// seen key 与磁盘路径对不上:增量 prune 会误删本地文件、上传会
// 误传本地未变文件、下载会重复下载。这里不是降级容错,而是
// 直接中止整个同步——宁可本次同步失败,也不带着损坏的相对路径
// 继续执行造成大规模误删/误传/重下(参考用户反馈"云盘没动却重下重传")。
errMu.Lock()
if firstErr == nil {
firstErr = fmt.Errorf("115: 解析目录树失败(file_id=%s):%w", pid, err)
}
errMu.Unlock()
st.scanIncomplete.Store(true)
cancel()
return
} else if detail != nil {
// 解析相对路径
relPath := cleanDirRel(detail.RelativePath(rootCID))
@@ -681,6 +693,12 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error {
}()
}
pwg.Wait()
if firstErr != nil {
// 目录树解析失败会导致 rel 塌缩,若继续处理会让大量本地文件
// 被错误判定为"云端不存在"而重复下载/上传,并可能误删本地文件。
// 中止本次同步,避免在损坏的相对路径上执行任何写操作。
return firstErr
}
}
st.updateSyncMessage(fmt.Sprintf("正在生成 STRM 与同步文件 (共 %d 个)...", len(allFiles)))
@@ -698,12 +716,10 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error {
if parentVal, ok := st.dirCache.Load(f.Pid); ok && parentVal.(string) != "" {
rel = cleanDirRel(parentVal.(string)) + "/" + cleanName
} else {
// 父目录不在目录缓存(目录详情先前解析失败),无法还原真实相对路径。
// 该文件会落到根/错误路径,seen key 与磁盘路径不符,增量 prune 会误删,标记扫描不完整。
if st.syncType == model.StrmSyncTypeIncremental {
st.scanIncomplete.Store(true)
}
rel = cleanName
// 父目录不在目录缓存,无法还原真实相对路径。若继续用塌缩后的
// 根路径处理,该文件会被错误判定,导致重复下载/上传或误删本地文件。
// 目录树不完整时宁可中止本次同步,也不带着损坏的 rel 继续执行。
return fmt.Errorf("115: 文件 %s 的父目录未解析成功,目录树不完整,中止同步以防误删/误传", cleanName)
}
}
entry := cloud.FileEntry{
+208 -133
View File
@@ -2,6 +2,8 @@ package service
import (
"context"
"net/http"
"net/http/httptest"
"net/url"
"os"
"path/filepath"
@@ -19,6 +21,7 @@ import (
"github.com/ShukeBta/MMTL/internal/model"
"github.com/ShukeBta/MMTL/internal/repository"
"github.com/ShukeBta/MMTL/internal/service/cloud"
"github.com/ShukeBta/MMTL/internal/service/cloud115"
)
// testStrmService 构建带内存库的 StrmService。
@@ -35,10 +38,10 @@ func testStrmService(t *testing.T) *StrmService {
sqlDB.SetMaxOpenConns(4)
t.Cleanup(func() { _ = sqlDB.Close() })
}
if err := db.AutoMigrate(&model.StrmAccount{}, &model.StrmSyncPath{}, &model.StrmSyncRecord{},
&model.StrmDownloadTask{}, &model.StrmUploadTask{}, &model.StrmDirCache{}, &model.Setting{}); err != nil {
t.Fatal(err)
}
if err := db.AutoMigrate(&model.StrmAccount{}, &model.StrmSyncPath{}, &model.StrmSyncRecord{},
&model.StrmDownloadTask{}, &model.StrmUploadTask{}, &model.StrmDirCache{}, &model.Setting{}); err != nil {
t.Fatal(err)
}
repos := repository.New(db)
ctx := context.Background()
if err := repos.Setting.Set(ctx, StrmSettingBaseURL, "http://test.local:8096"); err != nil {
@@ -210,7 +213,6 @@ func TestStrmFullAndIncrementalSync(t *testing.T) {
}
}
// TestStrmCronMatches cron 表达式匹配。
func TestStrmCronMatches(t *testing.T) {
cases := []struct {
@@ -494,142 +496,215 @@ func TestWalkRemoteConcurrent(t *testing.T) {
if walkErr != nil {
t.Fatal(walkErr)
}
if strmCount != 5 {
t.Errorf("生成的 .strm 数量 = %d,期望 5", strmCount)
}
if strmCount != 5 {
t.Errorf("生成的 .strm 数量 = %d,期望 5", strmCount)
}
}
// TestStrmBatchEnqueueAndConcurrentClaim 测试大规模批量入库及多协程并发认领无死锁
func TestStrmBatchEnqueueAndConcurrentClaim(t *testing.T) {
svc := testStrmService(t)
ctx := context.Background()
// 1. 批量插入 200 个下载任务
tasks := make([]*model.StrmDownloadTask, 0, 200)
for i := 0; i < 200; i++ {
tasks = append(tasks, &model.StrmDownloadTask{
SyncPathID: "test-sync-path",
AccountID: "test-acct",
Provider: model.StrmProvider115,
FileName: filepath.Base(string(rune('a'+i%26))) + ".nfo",
LocalPath: filepath.Join(t.TempDir(), string(rune('a'+i%26)), "test.nfo"),
Status: model.StrmTaskPending,
})
}
if err := svc.repo.StrmDownload.CreateInBatches(ctx, tasks, 50); err != nil {
t.Fatalf("CreateInBatches failed: %v", err)
}
// TestStrmBatchEnqueueAndConcurrentClaim 测试大规模批量入库及多协程并发认领无死锁
func TestStrmBatchEnqueueAndConcurrentClaim(t *testing.T) {
svc := testStrmService(t)
ctx := context.Background()
// 2. 验证 ActiveLocalPathMap
activeMap, err := svc.repo.StrmDownload.GetActiveLocalPathMap(ctx, "test-sync-path")
if err != nil {
t.Fatalf("GetActiveLocalPathMap failed: %v", err)
}
if len(activeMap) == 0 {
t.Fatal("expected active local path map to have entries")
}
// 1. 批量插入 200 个下载任务
tasks := make([]*model.StrmDownloadTask, 0, 200)
for i := 0; i < 200; i++ {
tasks = append(tasks, &model.StrmDownloadTask{
SyncPathID: "test-sync-path",
AccountID: "test-acct",
Provider: model.StrmProvider115,
FileName: filepath.Base(string(rune('a'+i%26))) + ".nfo",
LocalPath: filepath.Join(t.TempDir(), string(rune('a'+i%26)), "test.nfo"),
Status: model.StrmTaskPending,
})
}
if err := svc.repo.StrmDownload.CreateInBatches(ctx, tasks, 50); err != nil {
t.Fatalf("CreateInBatches failed: %v", err)
}
// 2. 验证 ActiveLocalPathMap
activeMap, err := svc.repo.StrmDownload.GetActiveLocalPathMap(ctx, "test-sync-path")
if err != nil {
t.Fatalf("GetActiveLocalPathMap failed: %v", err)
}
if len(activeMap) == 0 {
t.Fatal("expected active local path map to have entries")
}
// 3. 模拟 6 个 worker 并发 ClaimPendingDownload
claimedCount := 0
var claimMu sync.Mutex
var wg sync.WaitGroup
for w := 0; w < 6; w++ {
wg.Add(1)
go func() {
defer wg.Done()
for {
batch, err := svc.repo.StrmDownload.ClaimPendingDownload(ctx, 10)
if err != nil {
t.Errorf("concurrent ClaimPendingDownload failed: %v", err)
return
}
if len(batch) == 0 {
return
}
claimMu.Lock()
claimedCount += len(batch)
claimMu.Unlock()
// 3. 模拟 6 个 worker 并发 ClaimPendingDownload
claimedCount := 0
var claimMu sync.Mutex
var wg sync.WaitGroup
for w := 0; w < 6; w++ {
wg.Add(1)
go func() {
defer wg.Done()
for {
batch, err := svc.repo.StrmDownload.ClaimPendingDownload(ctx, 10)
if err != nil {
t.Errorf("concurrent ClaimPendingDownload failed: %v", err)
return
}
}()
}
wg.Wait()
if claimedCount != 200 {
t.Fatalf("expected all 200 tasks claimed, got %d", claimedCount)
if len(batch) == 0 {
return
}
claimMu.Lock()
claimedCount += len(batch)
claimMu.Unlock()
}
}
}()
}
wg.Wait()
// TestStrmDuplicateFileConflictResolution 测试远端存在多个同名不同大小文件时,本地确定性仲裁,避免增量死循环
func TestStrmDuplicateFileConflictResolution(t *testing.T) {
svc := testStrmService(t)
localDir := t.TempDir()
if claimedCount != 200 {
t.Fatalf("expected all 200 tasks claimed, got %d", claimedCount)
}
}
p := &model.StrmSyncPath{
Base: model.Base{ID: "dup-test-path"},
Provider: model.StrmProvider115,
RemotePath: "root",
LocalPath: localDir,
DownloadMeta: true,
}
// TestStrmDuplicateFileConflictResolution 测试远端存在多个同名不同大小文件时,本地确定性仲裁,避免增量死循环
func TestStrmDuplicateFileConflictResolution(t *testing.T) {
svc := testStrmService(t)
localDir := t.TempDir()
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)
}
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)
}
}
// TestWalk115FlatAbortsOnDirResolveFailure 回归测试:115 开放平台 token 失效/目录详情
// 解析失败时,同步必须中止而不是带着塌缩的 rel 继续处理,否则会导致本地大量元数据
// 被误判为"云端不存在"而重复下载/上传,甚至误删本地文件(用户反馈"云盘没动却重下重传")。
func TestWalk115FlatAbortsOnDirResolveFailure(t *testing.T) {
svc := testStrmService(t)
localDir := t.TempDir()
acct := &model.StrmAccount{
Name: "fake115",
Provider: "cloud115",
Config: "{}",
Enabled: true,
}
if err := svc.repo.StrmAccount.Create(context.Background(), acct); err != nil {
t.Fatal(err)
}
p := &model.StrmSyncPath{
Base: model.Base{ID: "abort-path"},
AccountID: acct.ID,
Provider: model.StrmProvider115,
RemotePath: "0",
LocalPath: localDir,
}
// 115 mock:文件列表返回一个视频(父目录 999 不在缓存,需要 get_info),
// get_info 恒返回 access_token 格式错误(40140123)→ 目录树解析失败。
api := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/open/ufile/files":
w.Write([]byte(`{"state":true,"count":1,"data":[{"fid":"100","pid":"999","fc":1,"fn":"movie.mkv","pc":"pc1","upt":1700000000,"fs":1024}]}`))
case "/open/folder/get_info":
w.Write([]byte(`{"state":false,"code":40140123,"message":"access_token 格式错误"}`))
default:
t.Errorf("unexpected path %s", r.URL.Path)
}
}))
defer api.Close()
oldPro := cloud115.ProAPIBase
cloud115.ProAPIBase = api.URL
defer func() { cloud115.ProAPIBase = oldPro }()
oc := cloud115.NewOpenClient("app", "at", "rt")
st := &strmSyncState{
s: svc,
ctx: context.Background(),
p: p,
provider: cloud.NewOpenAPI115("app", "at", "rt"),
cfg: &strmPathConfig{VideoExt: []string{"mkv"}, MetaExt: []string{"nfo"}, AddPath: 1, DownloadMeta: false},
rec: &model.StrmSyncRecord{},
syncType: model.StrmSyncTypeFull,
dirCache: sync.Map{},
seenVideo: map[string]bool{},
seenMeta: map[string]bool{},
remoteMeta: map[string]int64{},
}
err := st.walk115Flat(oc)
if err == nil {
t.Fatal("expected walk115Flat to abort on dir-resolve failure, got nil error")
}
// 中止后不允许产生任何部分写入(本地不允许生成 .strm 文件)。
var strmCount int
_ = filepath.WalkDir(localDir, func(path string, d os.DirEntry, err error) error {
if err == nil && !d.IsDir() && strings.HasSuffix(d.Name(), ".strm") {
strmCount++
}
return nil
})
if strmCount != 0 {
t.Fatalf("expected no .strm written after abort, got %d", strmCount)
}
}
+23
View File
@@ -0,0 +1,23 @@
// Package webui embeds the built React SPA so a single binary can serve the
// UI without requiring a separate web/dist directory on disk. The embed
// happens at compile time, so `web/dist` must exist before `go build` runs
// (the CI pipeline builds it via `npm run build` first).
package webui
import (
"embed"
"io/fs"
)
//go:embed all:dist
var distFS embed.FS
// DistFS returns the embedded SPA build artifacts rooted at the directory
// containing index.html (i.e. web/dist).
func DistFS() fs.FS {
sub, err := fs.Sub(distFS, "dist")
if err != nil {
panic(err)
}
return sub
}
+7
View File
@@ -117,6 +117,13 @@ export const strmAPI = {
.get<StrmSyncRecord[]>('/admin/strm/records', { params: pathId ? { path_id: pathId } : {} })
.then((r) => r.data),
deleteRecord: (id: string) => api.delete(`/admin/strm/records/${id}`).then((r) => r.data),
clearRecords: (pathId?: string) =>
api
.delete<{ deleted: number }>('/admin/strm/records', { params: pathId ? { path_id: pathId } : {} })
.then((r) => r.data),
// ── 本地目录浏览(同步目录选择器) ────────────────────────
listLocalDirs: (path?: string) =>
api
+51 -2
View File
@@ -211,7 +211,7 @@ export function StrmManagePage() {
onCancel={cancelSync}
/>
<RecordSection records={records} />
<RecordSection records={records} onDeleted={refresh} />
</>
)}
@@ -453,13 +453,48 @@ function SyncPathSection({
// ─── 同步记录 ────────────────────────────────────────────────────────────────
function RecordSection({ records }: { records: StrmSyncRecord[] }) {
function RecordSection({ records, onDeleted }: { records: StrmSyncRecord[]; onDeleted: () => void }) {
const [deletingId, setDeletingId] = useState<string | null>(null)
const deleteRecord = async (record: StrmSyncRecord) => {
const ok = await confirmAction({ message: '确定删除这条同步记录?', confirmText: '删除' })
if (!ok) return
setDeletingId(record.id)
try {
await strmAPI.deleteRecord(record.id)
toast.success('已删除同步记录')
onDeleted()
} catch (err) {
toast.error(apiErrorMessage(err))
} finally {
setDeletingId(null)
}
}
const clearRecords = async () => {
const ok = await confirmAction({ message: '确定清空全部同步记录?此操作不可恢复。', confirmText: '清空' })
if (!ok) return
try {
const res = await strmAPI.clearRecords()
toast.success(`已清空 ${res.deleted} 条同步记录`)
onDeleted()
} catch (err) {
toast.error(apiErrorMessage(err))
}
}
return (
<section className="glass-panel space-y-3 p-5">
<div className="flex items-center gap-2">
<History size={18} className="text-brand-500" />
<h2 className="font-display text-lg font-semibold text-ink-600">同步记录</h2>
<span className="rounded-full bg-gray-100 px-2 py-0.5 text-[11px] text-sand-500">{records.length}</span>
{records.length > 0 && (
<button type="button" onClick={clearRecords} className={iconButtonCls + ' ml-auto'}>
<Trash2 size={14} />
清空
</button>
)}
</div>
{records.length === 0 ? (
<p className="rounded-xl bg-gray-50 px-4 py-6 text-center text-sm text-sand-500">还没有同步记录</p>
@@ -475,8 +510,10 @@ function RecordSection({ records }: { records: StrmSyncRecord[] }) {
<th className="px-3 py-2 text-right">新增/更新</th>
<th className="px-3 py-2 text-right">跳过</th>
<th className="px-3 py-2 text-right">下载元数据</th>
<th className="px-3 py-2 text-right">上传元数据</th>
<th className="px-3 py-2 text-right">清理</th>
<th className="px-3 py-2">说明</th>
<th className="px-3 py-2"></th>
</tr>
</thead>
<tbody>
@@ -502,8 +539,20 @@ function RecordSection({ records }: { records: StrmSyncRecord[] }) {
<td className="px-3 py-2 text-right text-brand-500">{record.new_strm}</td>
<td className="px-3 py-2 text-right text-gray-500">{record.skipped}</td>
<td className="px-3 py-2 text-right">{record.new_meta}</td>
<td className="px-3 py-2 text-right">{record.uploaded ?? 0}</td>
<td className="px-3 py-2 text-right">{record.pruned}</td>
<td className="max-w-[260px] truncate px-3 py-2 text-xs text-sand-500">{record.message}</td>
<td className="px-3 py-2 text-right">
<button
type="button"
onClick={() => deleteRecord(record)}
disabled={deletingId === record.id}
title="删除记录"
className="rounded-md p-1 text-sand-400 transition hover:bg-rose-50 hover:text-rose-500 disabled:opacity-40"
>
<Trash2 size={15} />
</button>
</td>
</tr>
)
})}