From afae22ffe277cab073de9c31ef01af093a17e1a8 Mon Sep 17 00:00:00 2001 From: truewhile <62226914+truewhile@users.noreply.github.com> Date: Mon, 24 Aug 2026 15:00:58 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 优化 --- internal/service/cloud115/open.go | 7 +- internal/service/cloud115/url_cache.go | 74 +++++++++++++++++++++ internal/service/cloud115/url_cache_test.go | 31 +++++++++ internal/service/strm_queue.go | 60 +++++++++++++++++ internal/service/strm_queue_test.go | 46 +++++++++++++ internal/service/strm_service.go | 4 ++ 6 files changed, 221 insertions(+), 1 deletion(-) create mode 100644 internal/service/cloud115/url_cache.go create mode 100644 internal/service/cloud115/url_cache_test.go create mode 100644 internal/service/strm_queue_test.go diff --git a/internal/service/cloud115/open.go b/internal/service/cloud115/open.go index 31b044a..5192add 100644 --- a/internal/service/cloud115/open.go +++ b/internal/service/cloud115/open.go @@ -124,8 +124,12 @@ type downloadURLData struct { } `json:"url"` } -// GetDownloadURL 获取下载直链(pickcode)。 +// GetDownloadURL 获取下载直链(pickcode)。命中缓存直接返回, +// 避免对同一文件反复换取直链触发 115 风控。 func (c *OpenClient) GetDownloadURL(ctx context.Context, pickCode string) (string, error) { + if cached := GetDownloadURLCache(pickCode); cached != "" { + return cached, nil + } params := map[string]string{"pick_code": pickCode} resp, err := c.doAuthJSON(ctx, "POST", ProAPIBase+"/open/ufile/downurl", params, 1) if err != nil { @@ -139,6 +143,7 @@ func (c *OpenClient) GetDownloadURL(ctx context.Context, pickCode string) (strin if first.URL.URL == "" { return "", fmt.Errorf("115: 下载地址为空(文件可能未上传完成或已被删除)") } + SetDownloadURLCache(pickCode, first.URL.URL) return first.URL.URL, nil } diff --git a/internal/service/cloud115/url_cache.go b/internal/service/cloud115/url_cache.go new file mode 100644 index 0000000..e0c0841 --- /dev/null +++ b/internal/service/cloud115/url_cache.go @@ -0,0 +1,74 @@ +// 下载直链进程级缓存。 +// +// 115 的 /open/ufile/downurl 换取结果在一段时间内有效,且换取请求是 +// WAF 风控重点关注的密集接口。缓存按 pickcode 存直链、避免对同一文件 +// 反复换取;下载失败(http 403/404 等)时由调用方主动清除对应条目, +// 下一轮重试会重新换取链接。 +package cloud115 + +import ( + "sync" + "time" +) + +// downloadURLCacheTTL 直链缓存时长。115 直链一般有效 1 小时左右, +// 取 45 分钟留出余量。QMediaSync 同类缓存使用 50 分钟。 +const downloadURLCacheTTL = 45 * time.Minute + +// maxCachedURLs 超过该数量时顺带清理过期条目,防止 map 无限膨胀。 +const maxCachedURLs = 10000 + +type urlCacheEntry struct { + url string + expiresAt time.Time +} + +var ( + urlCacheMu sync.Mutex + urlCache = map[string]urlCacheEntry{} +) + +// GetDownloadURLCache 返回未过期的缓存直链;不存在或已过期返回空串。 +func GetDownloadURLCache(pickCode string) string { + if pickCode == "" { + return "" + } + urlCacheMu.Lock() + defer urlCacheMu.Unlock() + entry, ok := urlCache[pickCode] + if !ok || time.Now().After(entry.expiresAt) { + if ok { + delete(urlCache, pickCode) + } + return "" + } + return entry.url +} + +// SetDownloadURLCache 写入直链缓存。 +func SetDownloadURLCache(pickCode, url string) { + if pickCode == "" || url == "" { + return + } + urlCacheMu.Lock() + defer urlCacheMu.Unlock() + if len(urlCache) >= maxCachedURLs { + now := time.Now() + for k, e := range urlCache { + if now.After(e.expiresAt) { + delete(urlCache, k) + } + } + } + urlCache[pickCode] = urlCacheEntry{url: url, expiresAt: time.Now().Add(downloadURLCacheTTL)} +} + +// ClearDownloadURLCache 删除指定 pickcode 的缓存(下载得到非 2xx 时调用)。 +func ClearDownloadURLCache(pickCode string) { + if pickCode == "" { + return + } + urlCacheMu.Lock() + defer urlCacheMu.Unlock() + delete(urlCache, pickCode) +} \ No newline at end of file diff --git a/internal/service/cloud115/url_cache_test.go b/internal/service/cloud115/url_cache_test.go new file mode 100644 index 0000000..a331d52 --- /dev/null +++ b/internal/service/cloud115/url_cache_test.go @@ -0,0 +1,31 @@ +package cloud115 + +import ( + "testing" + "time" +) + +func TestDownloadURLCache(t *testing.T) { + urlCacheMu.Lock() + urlCache = map[string]urlCacheEntry{} + urlCacheMu.Unlock() + + if got := GetDownloadURLCache("pick1"); got != "" { + t.Fatalf("empty cache should return empty, got %q", got) + } + SetDownloadURLCache("pick1", "http://cdn/1") + if got := GetDownloadURLCache("pick1"); got != "http://cdn/1" { + t.Fatalf("cache miss, got %q", got) + } + // 过期条目按未命中处理 + urlCacheMu.Lock() + urlCache["pick1"] = urlCacheEntry{url: "http://cdn/old", expiresAt: time.Now().Add(-time.Second)} + urlCacheMu.Unlock() + if got := GetDownloadURLCache("pick1"); got != "" { + t.Fatalf("expired entry should be a miss, got %q", got) + } + ClearDownloadURLCache("pick1") + if got := GetDownloadURLCache("pick1"); got != "" { + t.Fatalf("cleared entry should be a miss, got %q", got) + } +} \ No newline at end of file diff --git a/internal/service/strm_queue.go b/internal/service/strm_queue.go index 45f707d..a779182 100644 --- a/internal/service/strm_queue.go +++ b/internal/service/strm_queue.go @@ -12,6 +12,7 @@ import ( "net/http" "os" "path/filepath" + "strings" "time" "go.uber.org/zap" @@ -35,6 +36,12 @@ func (s *StrmService) downloadWorker(ctx context.Context) { return default: } + // 115 风控/限流熔断:冷却期间整体暂停,不给 WAF 续封机会 + if left := s.wafCooldownLeft(); left > 0 { + s.log.Debug("下载队列冷却中", zap.Duration("remaining", left)) + sleepContext(ctx, left) + continue + } tasks, err := s.repo.StrmDownload.ClaimPendingDownload(ctx, 1) if err != nil { s.log.Warn("claim strm download task failed", zap.Error(err)) @@ -73,10 +80,17 @@ func (s *StrmService) processDownloadTask(ctx context.Context, task *model.StrmD } link, err := provider.Resolve(ctx, task.RemoteRef) if err != nil { + if is115Blocked(err) { + s.triggerWAFCooldown() + } s.downloadTaskFailWithRetry(task, "解析下载地址失败:"+err.Error()) return } if err := downloadToFile(ctx, link, task.LocalPath, s.http); err != nil { + // 直链失效(403/404/410 等):清掉缓存让下一轮重新换取 + if isHTTPDownloadFailure(err) && task.Provider == model.StrmProvider115 { + cloud115.ClearDownloadURLCache(task.RemoteRef) + } s.downloadTaskFailWithRetry(task, "下载失败:"+err.Error()) return } @@ -497,3 +511,49 @@ func sleepContext(ctx context.Context, d time.Duration) { case <-time.After(d): } } + +// ─── 115 风控/限流熔断 ──────────────────────────────────────────────────────── + +// triggerWAFCooldown 检测到 115 风控/限流后触发全局冷却,冷却期间下载 worker 暂停。 +// 冷却时间取最大值,避免连续触发时缩短等待。 +func (s *StrmService) triggerWAFCooldown() { + s.mu.Lock() + defer s.mu.Unlock() + until := time.Now().Add(strmWAFCooldown) + if until.After(s.wafUntil) { + s.wafUntil = until + s.log.Warn("115 风控/限流,下载队列进入冷却", zap.Duration("cooldown", strmWAFCooldown)) + } +} + +// wafCooldownLeft 返回剩余冷却时间(0 表示无需冷却)。 +func (s *StrmService) wafCooldownLeft() time.Duration { + s.mu.Lock() + defer s.mu.Unlock() + if s.wafUntil.After(time.Now()) { + return time.Until(s.wafUntil) + } + return 0 +} + +// is115Blocked 判断错误是否来自 115 的风控/限流(WAF 405 拦截页或限流错误码)。 +func is115Blocked(err error) bool { + if err == nil { + return false + } + msg := strings.ToLower(err.Error()) + return strings.Contains(msg, "115 接口返回 http 405") || + strings.Contains(msg, "访问被阻断") || + strings.Contains(msg, "request has been blocked") || + strings.Contains(msg, "115 接口错误(770004") || + strings.Contains(msg, "115 接口错误(406") +} + +// isHTTPDownloadFailure 判断下载是否因 HTTP 状态码失败(直链失效需清缓存重取)。 +func isHTTPDownloadFailure(err error) bool { + if err == nil { + return false + } + msg := strings.ToLower(err.Error()) + return strings.Contains(msg, "http 4") || strings.Contains(msg, "http 5") +} diff --git a/internal/service/strm_queue_test.go b/internal/service/strm_queue_test.go new file mode 100644 index 0000000..03dc514 --- /dev/null +++ b/internal/service/strm_queue_test.go @@ -0,0 +1,46 @@ +package service + +import ( + "errors" + "testing" +) + +func TestIs115Blocked(t *testing.T) { + cases := []struct { + err error + want bool + }{ + {errors.New("115 接口返回 HTTP 405:...访问被阻断"), true}, + {errors.New("115 接口返回 HTTP 405"), true}, + {errors.New("115 接口错误(770004):访问频率过高"), true}, + {errors.New("115 接口错误(406):达到访问上限"), true}, + {errors.New("下载失败:http 403"), false}, + {errors.New("解析下载地址失败:115 接口调用失败"), false}, + {errors.New("database is locked (517)"), false}, + {nil, false}, + } + for _, c := range cases { + if got := is115Blocked(c.err); got != c.want { + t.Errorf("is115Blocked(%v) = %v, want %v", c.err, got, c.want) + } + } +} + +func TestIsHTTPDownloadFailure(t *testing.T) { + cases := []struct { + err error + want bool + }{ + {errors.New("http 403"), true}, + {errors.New("http 404"), true}, + {errors.New("http 410"), true}, + {errors.New("http 500"), true}, + {errors.New("Get \"https://x\": connection refused"), false}, + {nil, false}, + } + for _, c := range cases { + if got := isHTTPDownloadFailure(c.err); got != c.want { + t.Errorf("isHTTPDownloadFailure(%v) = %v, want %v", c.err, got, c.want) + } + } +} \ No newline at end of file diff --git a/internal/service/strm_service.go b/internal/service/strm_service.go index d02cab9..aae0c90 100644 --- a/internal/service/strm_service.go +++ b/internal/service/strm_service.go @@ -89,8 +89,12 @@ type StrmService struct { mu sync.Mutex running map[string]context.CancelFunc // sync path id -> cancel oauthSessions map[string]*strm115AuthSession + wafUntil time.Time // 115 风控/限流熔断截止时间(由 mu 保护) } +// strmWAFCooldown 检测到 115 风控/限流后下载队列的全局冷却时长。 +const strmWAFCooldown = 3 * time.Minute + // NewStrmService constructs the STRM service. func NewStrmService(cfg *config.Config, log *zap.Logger, repos *repository.Container, crypto *CryptoService) *StrmService { return &StrmService{