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