优化
This commit is contained in:
truewhile
2026-08-24 15:00:58 +08:00
parent a59cb5b858
commit afae22ffe2
6 changed files with 221 additions and 1 deletions
+6 -1
View File
@@ -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
}
+74
View File
@@ -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)
}
@@ -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)
}
}
+60
View File
@@ -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")
}
+46
View File
@@ -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:<!doctypehtml>...访问被阻断"), 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)
}
}
}
+4
View File
@@ -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{