update: single write cache
This commit is contained in:
+46
-7
@@ -2,6 +2,7 @@ package main
|
||||
|
||||
import (
|
||||
"flag"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
@@ -12,6 +13,7 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"golang.org/x/sync/singleflight"
|
||||
)
|
||||
|
||||
// -----------------------------------------------------------------------------
|
||||
@@ -37,6 +39,11 @@ var (
|
||||
// 30 min 足够 330 KB 的 iife 与 vendor/ 下最大的 docx/pptx wasm (~50 MB) 在慢网下完成
|
||||
var fileViewerHTTPClient = &http.Client{Timeout: 30 * time.Minute}
|
||||
|
||||
// fileViewerFetchGroup 按 cachePath 去重并发 fetch:
|
||||
// 同 key 同一时刻只有一个 goroutine 真正去 CDN 拉, 其他并发请求阻塞等待,
|
||||
// fetch 完成后从磁盘读 (文件已经落盘). 避免多个请求对同一文件重复上游拉取 + 写盘竞态.
|
||||
var fileViewerFetchGroup singleflight.Group
|
||||
|
||||
// fileViewerProxy 处理 GET /file-viewer/*filepath.
|
||||
// 1. 校验路径防 ../ 越界
|
||||
// 2. 缓存命中 -> 直接读盘
|
||||
@@ -112,21 +119,49 @@ func proxyFileViewer(c *gin.Context, name string) {
|
||||
}
|
||||
}
|
||||
|
||||
// fetchAndCache: 从 CDN 拉取资源, 同时写盘 + 流回客户端.
|
||||
// fetchAndCache: 通过 singleflight 按 cachePath 去重并发. 同一个 key 只有第一个
|
||||
// 请求调 doFetchAndCache (拉 CDN + 边写盘边流回 writer 的 c.Writer); 其他并发
|
||||
// 请求阻塞等待, 拿到结果后从磁盘读 (c.File).
|
||||
//
|
||||
// 注意: singleflight 的 shared=true 表示"结果被分享给多个 caller", 并不区分
|
||||
// 首 caller 还是后续 caller. 这里用 closure 变量 isWriter 精确标识谁是写者:
|
||||
// - closure 被执行 = 写者 (响应已流回本 c.Writer, 直接返回)
|
||||
// - closure 未执行 = 等待者 (自己组装响应: 错误 -> 502, 成功 -> c.File)
|
||||
func fetchAndCache(c *gin.Context, cachePath, name string) {
|
||||
isWriter := false
|
||||
_, err, _ := fileViewerFetchGroup.Do(cachePath, func() (interface{}, error) {
|
||||
isWriter = true
|
||||
return nil, doFetchAndCache(c, cachePath, name)
|
||||
})
|
||||
|
||||
if isWriter {
|
||||
return
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.Header("Cache-Control", "public, max-age=31536000, immutable")
|
||||
c.File(cachePath)
|
||||
}
|
||||
|
||||
// doFetchAndCache: writer 路径专属 — 拉上游, 边写盘边流回传入的 c.Writer.
|
||||
// 任何错误都会写入 c.Writer (502), 并返回 error 让 singleflight 把结果传给 waiter.
|
||||
func doFetchAndCache(c *gin.Context, cachePath, name string) error {
|
||||
upstream := upstreamURL(name)
|
||||
resp, err := fileViewerHTTPClient.Get(upstream)
|
||||
if err != nil {
|
||||
log.Printf("[file-viewer] upstream fetch %s: %v", name, err)
|
||||
c.JSON(http.StatusBadGateway, gin.H{"error": "upstream fetch failed"})
|
||||
return
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
log.Printf("[file-viewer] upstream %s returned %d", name, resp.StatusCode)
|
||||
c.JSON(http.StatusBadGateway, gin.H{"error": "upstream returned " + resp.Status})
|
||||
return
|
||||
return fmt.Errorf("upstream returned %s", resp.Status)
|
||||
}
|
||||
|
||||
if ct := resp.Header.Get("Content-Type"); ct != "" {
|
||||
@@ -135,17 +170,17 @@ func fetchAndCache(c *gin.Context, cachePath, name string) {
|
||||
c.Header("Cache-Control", "public, max-age=31536000, immutable")
|
||||
c.Status(http.StatusOK)
|
||||
|
||||
// 目录创建失败 -> 降级透传, 不缓存
|
||||
// 目录/文件创建失败 -> 降级透传, 不缓存
|
||||
if err := os.MkdirAll(filepath.Dir(cachePath), 0o755); err != nil {
|
||||
log.Printf("[file-viewer] mkdir cache dir: %v", err)
|
||||
_, _ = io.Copy(c.Writer, resp.Body)
|
||||
return
|
||||
return err
|
||||
}
|
||||
f, err := os.Create(cachePath)
|
||||
if err != nil {
|
||||
log.Printf("[file-viewer] create cache file: %v", err)
|
||||
_, _ = io.Copy(c.Writer, resp.Body)
|
||||
return
|
||||
return err
|
||||
}
|
||||
|
||||
// 关键: MultiWriter 让客户端拿到的同时落盘, 用户感知零等待
|
||||
@@ -161,9 +196,13 @@ func fetchAndCache(c *gin.Context, cachePath, name string) {
|
||||
if cErr != nil {
|
||||
log.Printf("[file-viewer] close %s: %v", name, cErr)
|
||||
}
|
||||
return
|
||||
if copyErr != nil {
|
||||
return copyErr
|
||||
}
|
||||
return cErr
|
||||
}
|
||||
log.Printf("[file-viewer] cached %s (%d bytes)", name, written)
|
||||
return nil
|
||||
}
|
||||
|
||||
func upstreamURL(name string) string {
|
||||
|
||||
@@ -1,10 +1,11 @@
|
||||
module tmp-upload
|
||||
|
||||
go 1.24
|
||||
go 1.25.0
|
||||
|
||||
require (
|
||||
github.com/gin-gonic/gin v1.10.0
|
||||
github.com/robfig/cron/v3 v3.0.1
|
||||
golang.org/x/sync v0.22.0
|
||||
)
|
||||
|
||||
require (
|
||||
|
||||
@@ -72,6 +72,8 @@ golang.org/x/crypto v0.23.0 h1:dIJU/v2J8Mdglj/8rJ6UUOM3Zc9zLZxVZwwxMooUSAI=
|
||||
golang.org/x/crypto v0.23.0/go.mod h1:CKFgDieR+mRhux2Lsu27y0fO304Db0wZe70UKqHu0v8=
|
||||
golang.org/x/net v0.25.0 h1:d/OCCoBEUq33pjydKrGQhw7IlUPI2Oylr+8qLx49kac=
|
||||
golang.org/x/net v0.25.0/go.mod h1:JkAGAh7GEvH74S6FOH42FLoXpXbE/aqXSrIQjXgsiwM=
|
||||
golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek=
|
||||
golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
|
||||
golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.20.0 h1:Od9JTbYCk261bKm4M/mw7AklTlFYIa0bIp9BgSm1S8Y=
|
||||
|
||||
Reference in New Issue
Block a user