实现 MinIO 对象存储图片流水线(可迁移 OSS)

- 新增 S3 兼容 Storage 抽象与 MinIO 实现,compose 启动 MinIO
- 上传接口:校验 → 缩放 → WebP → 主图/缩略图入库
- 消息图片 content 改为对象 URL,拒绝 base64
- 工作台/访客端改为先上传再发消息
This commit is contained in:
yml2213
2026-07-15 12:21:31 +08:00
parent 85b407fddf
commit ec2f12c6ed
20 changed files with 907 additions and 108 deletions
+22 -12
View File
@@ -1,17 +1,24 @@
package handler
import (
"encoding/base64"
"errors"
"strings"
"unicode/utf8"
"kefu-sys/server/internal/storage"
)
const (
maxTextMessageLength = 2000
maxImageMessageSize = 5 * 1024 * 1024
)
// imagePublicBase 由 SetupRoutes 注入,用于校验图片消息 URL 归属
var imagePublicBase string
func SetImagePublicBase(base string) {
imagePublicBase = storage.NormalizePublicBase(base)
}
func validateMessageContent(messageType, content string) (string, error) {
switch messageType {
case "text":
@@ -24,19 +31,22 @@ func validateMessageContent(messageType, content string) (string, error) {
}
return content, nil
case "image":
parts := strings.SplitN(content, ",", 2)
if len(parts) != 2 {
return "", errors.New("图片格式无效")
content = strings.TrimSpace(content)
if content == "" {
return "", errors.New("图片地址不能为空")
}
if parts[0] != "data:image/jpeg;base64" && parts[0] != "data:image/png;base64" && parts[0] != "data:image/gif;base64" {
return "", errors.New("仅支持 jpg、png、gif 图片")
// 开发阶段直接切换为对象存储 URL,不再接受 base64
if strings.HasPrefix(content, "data:") {
return "", errors.New("请先上传图片,勿直接发送 base64")
}
bytes, err := base64.StdEncoding.DecodeString(parts[1])
if err != nil || len(bytes) == 0 {
return "", errors.New("图片内容无效")
if !strings.HasPrefix(content, "http://") && !strings.HasPrefix(content, "https://") {
return "", errors.New("图片地址无效")
}
if len(bytes) > maxImageMessageSize {
return "", errors.New("图片不能超过 5 MB")
if imagePublicBase != "" && !storage.IsAllowedObjectURL(imagePublicBase, content) {
return "", errors.New("图片地址不在允许的存储域名下")
}
if len(content) > 2000 {
return "", errors.New("图片地址过长")
}
return content, nil
default:
@@ -1,22 +1,25 @@
package handler
import (
"strings"
"testing"
)
import "testing"
func TestValidateMessageContent(t *testing.T) {
validImage := "data:image/png;base64,aGVsbG8="
if content, err := validateMessageContent("image", validImage); err != nil || content != validImage {
t.Fatalf("合法图片校验失败: content=%q err=%v", content, err)
SetImagePublicBase("http://localhost:9000/kefu")
if _, err := validateMessageContent("text", "hello"); err != nil {
t.Fatalf("文本消息应通过: %v", err)
}
if _, err := validateMessageContent("image", "data:image/webp;base64,aGVsbG8="); err == nil {
t.Fatal("不支持的图片格式未被拦截")
if _, err := validateMessageContent("text", ""); err == nil {
t.Fatal("空文本应失败")
}
if _, err := validateMessageContent("text", " "); err == nil {
t.Fatal("空白文本未被拦截")
validURL := "http://localhost:9000/kefu/tenants/1/chat/a.webp"
if content, err := validateMessageContent("image", validURL); err != nil || content != validURL {
t.Fatalf("合法图片 URL 校验失败: content=%q err=%v", content, err)
}
if _, err := validateMessageContent("text", strings.Repeat("字", maxTextMessageLength+1)); err == nil {
t.Fatal("超长文本未被拦截")
if _, err := validateMessageContent("image", "data:image/png;base64,aGVsbG8="); err == nil {
t.Fatal("base64 图片应被拒绝")
}
if _, err := validateMessageContent("image", "https://evil.example.com/a.webp"); err == nil {
t.Fatal("外域图片 URL 应被拒绝")
}
}
+10 -1
View File
@@ -2,10 +2,12 @@ package handler
import (
"github.com/gin-gonic/gin"
"kefu-sys/server/internal/config"
"kefu-sys/server/internal/middleware"
"kefu-sys/server/internal/storage"
)
func SetupRoutes(r *gin.Engine) {
func SetupRoutes(r *gin.Engine, store storage.ObjectStorage, storageCfg config.StorageConfig) {
auth := NewAuthHandler()
session := NewSessionHandler()
customer := NewCustomerHandler()
@@ -16,6 +18,9 @@ func SetupRoutes(r *gin.Engine) {
settings := NewSettingsHandler()
ws := NewWsHandler()
widget := NewWidgetHandler()
upload := NewUploadHandler(store, storageCfg)
SetImagePublicBase(storageCfg.PublicBase)
api := r.Group("/api")
@@ -31,6 +36,7 @@ func SetupRoutes(r *gin.Engine) {
widgetApi.GET("/messages", widget.GetMessages)
widgetApi.GET("/ws", widget.Connect)
widgetApi.POST("/rating", widget.SubmitRating)
widgetApi.POST("/upload", upload.WidgetUploadImage)
// 需要认证的接口
authRequired := api.Group("")
@@ -39,6 +45,9 @@ func SetupRoutes(r *gin.Engine) {
// WebSocket
authRequired.GET("/ws", ws.Connect)
// 上传
authRequired.POST("/uploads", upload.UploadImage)
// 会话管理
authRequired.GET("/agents/available", session.ListAvailableAgents)
sessions := authRequired.Group("/sessions")
@@ -2,8 +2,10 @@ package handler_test
import (
"bytes"
"encoding/base64"
"encoding/json"
"fmt"
"mime/multipart"
"net/http"
"net/http/httptest"
"testing"
@@ -13,9 +15,11 @@ import (
"golang.org/x/crypto/bcrypt"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"kefu-sys/server/internal/config"
"kefu-sys/server/internal/handler"
"kefu-sys/server/internal/middleware"
"kefu-sys/server/internal/model"
"kefu-sys/server/internal/storage"
)
func setupRouter(t *testing.T) *gin.Engine {
@@ -32,10 +36,45 @@ func setupRouter(t *testing.T) *gin.Engine {
model.DB = db
middleware.InitJWT("test-secret")
router := gin.New()
handler.SetupRoutes(router)
mem := storage.NewMemory("http://localhost:9000/kefu")
handler.SetupRoutes(router, mem, config.StorageConfig{
PublicBase: "http://localhost:9000/kefu",
MaxUploadMB: 10,
MaxImageEdge: 1920,
WebPQuality: 80,
})
return router
}
// tinyPNG 1x1 像素 PNG
func tinyPNG() []byte {
b, _ := base64.StdEncoding.DecodeString("iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR42mP8z8BQDwAEhQGAhKmMIQAAAABJRU5ErkJggg==")
return b
}
func multipartImageRequest(t *testing.T, method, target string, fileField string, filename string, data []byte, user model.User) *http.Request {
t.Helper()
var body bytes.Buffer
w := multipart.NewWriter(&body)
part, err := w.CreateFormFile(fileField, filename)
if err != nil {
t.Fatalf("创建 form file 失败: %v", err)
}
if _, err := part.Write(data); err != nil {
t.Fatalf("写入文件失败: %v", err)
}
_ = w.Close()
req := httptest.NewRequest(method, target, &body)
req.Header.Set("Content-Type", w.FormDataContentType())
token, err := middleware.GenerateToken(user.ID, user.TenantID, user.Role)
if err != nil {
t.Fatalf("生成令牌失败: %v", err)
}
req.Header.Set("Authorization", "Bearer "+token)
return req
}
func createTenant(t *testing.T, name, status string) model.Tenant {
t.Helper()
tenant := model.Tenant{Name: name, Status: status, ExpireAt: time.Now().AddDate(1, 0, 0)}
@@ -417,8 +456,22 @@ func TestWorkbenchSessionLifecycleUnreadNotesTransferAndImage(t *testing.T) {
t.Fatalf("转接后原客服仍可回复: %d %s", oldAgentMessageRecorder.Code, oldAgentMessageRecorder.Body.String())
}
uploadRecorder := httptest.NewRecorder()
router.ServeHTTP(uploadRecorder, multipartImageRequest(t, http.MethodPost, "/api/uploads", "file", "dot.png", tinyPNG(), agentTwo))
if uploadRecorder.Code != http.StatusOK {
t.Fatalf("上传图片失败: %d %s", uploadRecorder.Code, uploadRecorder.Body.String())
}
var uploadResp struct {
Data struct {
URL string `json:"url"`
} `json:"data"`
}
if err := json.Unmarshal(uploadRecorder.Body.Bytes(), &uploadResp); err != nil || uploadResp.Data.URL == "" {
t.Fatalf("解析上传响应失败: %v body=%s", err, uploadRecorder.Body.String())
}
imageRecorder := httptest.NewRecorder()
imageBody := []byte(`{"content":"data:image/png;base64,aGVsbG8=","type":"image"}`)
imageBody := []byte(fmt.Sprintf(`{"content":%q,"type":"image"}`, uploadResp.Data.URL))
router.ServeHTTP(imageRecorder, bearerRequest(t, http.MethodPost, fmt.Sprintf("/api/sessions/%d/messages", session.ID), imageBody, agentTwo))
if imageRecorder.Code != http.StatusOK {
t.Fatalf("发送图片消息失败: %d %s", imageRecorder.Code, imageRecorder.Body.String())
+156
View File
@@ -0,0 +1,156 @@
package handler
import (
"bytes"
"fmt"
"io"
"net/http"
"path"
"strings"
"github.com/gin-gonic/gin"
"kefu-sys/server/internal/config"
"kefu-sys/server/internal/media"
"kefu-sys/server/internal/middleware"
"kefu-sys/server/internal/storage"
)
type UploadHandler struct {
store storage.ObjectStorage
cfg config.StorageConfig
}
func NewUploadHandler(store storage.ObjectStorage, cfg config.StorageConfig) *UploadHandler {
return &UploadHandler{store: store, cfg: cfg}
}
// UploadImage 客服端上传:multipart file → 处理 → 存对象存储
func (h *UploadHandler) UploadImage(c *gin.Context) {
h.handleUpload(c, fmt.Sprintf("tenants/%d/chat", middleware.GetTenantID(c)))
}
// WidgetUploadImage 访客端上传:需 visitor token + session_id
func (h *UploadHandler) WidgetUploadImage(c *gin.Context) {
sessionID := c.PostForm("session_id")
if sessionID == "" {
sessionID = c.Query("session_id")
}
if sessionID == "" {
c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "缺少 session_id"})
return
}
token := visitorTokenFromRequest(c, c.PostForm("visitor_token"))
var sid uint
if _, err := fmt.Sscanf(sessionID, "%d", &sid); err != nil || sid == 0 {
c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "会话参数错误"})
return
}
session, ok := loadVisitorSession(c, sid, token)
if !ok {
return
}
if session.Status == "ended" || session.Status == "archived" {
c.JSON(http.StatusConflict, gin.H{"code": 409, "message": "会话已结束"})
return
}
h.handleUpload(c, fmt.Sprintf("tenants/%d/widget/%d", session.TenantID, session.ID))
}
func (h *UploadHandler) handleUpload(c *gin.Context, keyPrefix string) {
if h.store == nil {
c.JSON(http.StatusServiceUnavailable, gin.H{"code": 503, "message": "对象存储未配置"})
return
}
maxBytes := int64(h.cfg.MaxUploadMB) * 1024 * 1024
if maxBytes <= 0 {
maxBytes = 10 * 1024 * 1024
}
c.Request.Body = http.MaxBytesReader(c.Writer, c.Request.Body, maxBytes+512)
file, header, err := c.Request.FormFile("file")
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "请上传文件字段 file"})
return
}
defer file.Close()
if header.Size > maxBytes {
c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": fmt.Sprintf("文件不能超过 %d MB", h.cfg.MaxUploadMB)})
return
}
// 嗅探 MIME
head := make([]byte, 512)
n, _ := io.ReadFull(file, head)
contentType := http.DetectContentType(head[:n])
// 复位读取:拼回 head + 剩余
reader := io.MultiReader(bytes.NewReader(head[:n]), file)
// 部分浏览器 Content-Type 不准,也允许 header 中的 type
if !media.AllowedUploadMIME(contentType) {
ct := header.Header.Get("Content-Type")
if media.AllowedUploadMIME(ct) {
contentType = ct
} else {
// 按扩展名兜底
ext := strings.ToLower(path.Ext(header.Filename))
switch ext {
case ".jpg", ".jpeg":
contentType = "image/jpeg"
case ".png":
contentType = "image/png"
case ".gif":
contentType = "image/gif"
case ".webp":
contentType = "image/webp"
default:
c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "仅支持 jpg/png/gif/webp 图片"})
return
}
}
}
_ = contentType
raw, err := io.ReadAll(reader)
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "读取文件失败"})
return
}
if int64(len(raw)) > maxBytes {
c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": "文件过大"})
return
}
processed, err := media.ProcessImage(bytes.NewReader(raw), media.ProcessOptions{
MaxEdge: h.cfg.MaxImageEdge,
ThumbEdge: 400,
WebPQuality: h.cfg.WebPQuality,
})
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"code": 400, "message": err.Error()})
return
}
mainKey := storage.NewObjectKey(keyPrefix, processed.MainExt)
thumbKey := storage.NewObjectKey(keyPrefix+"/thumbs", processed.ThumbExt)
mainURL, err := h.store.Put(c.Request.Context(), mainKey, bytes.NewReader(processed.Main), int64(len(processed.Main)), processed.MainType)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "上传失败: " + err.Error()})
return
}
thumbURL, err := h.store.Put(c.Request.Context(), thumbKey, bytes.NewReader(processed.Thumb), int64(len(processed.Thumb)), processed.ThumbType)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "message": "缩略图上传失败: " + err.Error()})
return
}
middleware.JSON(c, gin.H{
"url": mainURL,
"thumb_url": thumbURL,
"content_type": processed.MainType,
"width": processed.Width,
"height": processed.Height,
"size": len(processed.Main),
})
}