Files
jiang13-bbs/backend/service/operations_storage.go
freefire 46e0cdc0b3 feat: 媒体库 blob 去重存储与品牌/注册等前端调整
新增:
- blob_store 二进制去重存储 + blob_dedup 测试

其它:
- 后端 settings/auth/media_library/upload/operations 服务与 handler 调整
- 前端品牌(BrandLockup/BrandCropModal/BrandSeoPanel)、注册、媒体库、分析详情页、Header/SiteChrome 等更新
2026-10-03 02:38:11 +08:00

305 lines
9.4 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package service
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"github.com/freefire/jiang13-bbs/model"
"github.com/minio/minio-go/v7"
"github.com/minio/minio-go/v7/pkg/credentials"
"io"
"net/http"
"net/url"
"os"
"path/filepath"
"strings"
"time"
"gorm.io/gorm"
)
func (o *Operations) s3(c StorageConfig) (*minio.Client, error) {
u, e := url.Parse(c.Endpoint)
if e != nil {
return nil, errors.New("存储地址无效")
}
access, e := o.open(c.AccessKey, "storage:access_key")
if e != nil {
return nil, e
}
secret, e := o.open(c.SecretKey, "storage:secret_key")
if e != nil {
return nil, e
}
lookup := minio.BucketLookupDNS
if c.PathStyle {
lookup = minio.BucketLookupPath
}
tr := &http.Transport{DialContext: safeDial, TLSHandshakeTimeout: 10 * time.Second, ResponseHeaderTimeout: 15 * time.Second, DisableKeepAlives: true}
return minio.New(u.Host, &minio.Options{Creds: credentials.NewStaticV4(access, secret, ""), Secure: u.Scheme == "https", Region: c.Region, BucketLookup: lookup, Transport: boundedS3Transport{base: tr, host: func() string {
if c.PathStyle {
return u.Host
}
return c.Bucket + "." + u.Host
}(), scheme: u.Scheme}})
}
func (o *Operations) TestStorage(ctx context.Context, raw json.RawMessage, clear []string) error {
v, _, _, e := o.draft(o.db, "storage", raw, clear)
if e != nil {
return e
}
c := *v.(*StorageConfig)
if c.Backend == "local" {
return o.LocalWritable()
}
client, e := o.s3(c)
if e != nil {
return e
}
ctx, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel()
key := c.Prefix + "tests/" + randomID()
payload := []byte("jiang13 storage probe")
if _, e = client.PutObject(ctx, c.Bucket, key, bytes.NewReader(payload), int64(len(payload)), minio.PutObjectOptions{ContentType: "text/plain", DisableMultipart: true}); e != nil {
return errors.New("存储写入失败(连接、凭据或权限)")
}
cleanup := func() error {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
return client.RemoveObject(ctx, c.Bucket, key, minio.RemoveObjectOptions{})
}
object, e := client.GetObject(ctx, c.Bucket, key, minio.GetObjectOptions{})
if e != nil {
_ = cleanup()
return errors.New("存储读取失败")
}
b, e := io.ReadAll(io.LimitReader(object, 128))
_ = object.Close()
cleanErr := cleanup()
if e != nil || !bytes.Equal(b, payload) {
return errors.New("存储读取校验失败")
}
if cleanErr != nil {
return errors.New("存储清理失败;测试对象保留于应用 tests 前缀")
}
return nil
}
func (o *Operations) LocalWritable() error {
dir := filepath.Join(o.cfg.DataDir, "uploads")
f, e := os.CreateTemp(dir, ".probe-")
if e != nil {
return errors.New("本地上传目录不可写")
}
name := f.Name()
_, e = f.Write([]byte("probe"))
_ = f.Close()
removeErr := os.Remove(name)
if e != nil || removeErr != nil {
return errors.New("本地存储读写或清理失败")
}
return nil
}
// Uploads are validated locally first, then moved to the selected target. No fallback.
func (o *Operations) StoreFile(path, mimeType string, public bool) (string, error) {
v, version, e := o.read(o.db, "storage")
if e != nil {
return "", errors.New("读取存储配置失败")
}
c := *v.(*StorageConfig)
if c.Backend == "local" {
return "", nil
}
client, e := o.s3(c)
if e != nil {
return "", e
}
ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
defer cancel()
f, e := os.Open(path)
if e != nil {
return "", e
}
defer f.Close()
stat, e := f.Stat()
if e != nil {
return "", e
}
id := randomID()
key := c.Prefix + "objects/" + id
if _, e = client.PutObject(ctx, c.Bucket, key, f, stat.Size(), minio.PutObjectOptions{ContentType: mimeType, DisableMultipart: true}); e != nil {
return "", errors.New("上传到 S3 失败,当前目标未自动切换")
}
obj := model.StoredObject{ID: id, ConfigName: fmt.Sprintf("storage-%d", version), Key: key, MIME: mimeType, Public: public}
if e = o.db.Create(&obj).Error; e != nil {
_ = client.RemoveObject(ctx, c.Bucket, key, minio.RemoveObjectOptions{})
return "", errors.New("保存文件记录失败")
}
return id, nil
}
// blobObjectID 内容寻址对象的稳定 ID:公开/私有加不同前缀,同字节也不共享可见性。
func blobObjectID(hash string, public bool) string {
if public {
return "pub_" + hash
}
return "prv_" + hash
}
// StoreFileHashed 按内容 hash 幂等上传:同一 (hash, public) 只存一份,重复上传跳过传输。
// 本地后端返回空 ID(文件留在本地 blob 目录,由调用方做引用计数)。
func (o *Operations) StoreFileHashed(path, mimeType string, public bool, hash string) (string, error) {
v, version, e := o.read(o.db, "storage")
if e != nil {
return "", errors.New("读取存储配置失败")
}
c := *v.(*StorageConfig)
if c.Backend == "local" {
return "", nil
}
id := blobObjectID(hash, public)
// hash 相同字节必然相同:对象已存在时直接复用,省掉整次上传
var existing model.StoredObject
if e := o.db.First(&existing, "id = ?", id).Error; e == nil {
return id, nil
} else if !errors.Is(e, gorm.ErrRecordNotFound) {
return "", e
}
client, e := o.s3(c)
if e != nil {
return "", e
}
ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
defer cancel()
f, e := os.Open(path)
if e != nil {
return "", e
}
stat, e := f.Stat()
if e != nil {
_ = f.Close()
return "", e
}
key := c.Prefix + "objects/" + id
if _, e = client.PutObject(ctx, c.Bucket, key, f, stat.Size(), minio.PutObjectOptions{ContentType: mimeType, DisableMultipart: true}); e != nil {
_ = f.Close()
return "", errors.New("上传到 S3 失败,当前目标未自动切换")
}
_ = f.Close()
obj := model.StoredObject{ID: id, ConfigName: fmt.Sprintf("storage-%d", version), Key: key, MIME: mimeType, Public: public}
if e = o.db.Create(&obj).Error; e != nil {
if isUniqueConflict(e) { // 并发:另一请求已写入同内容对象,直接复用
return id, nil
}
_ = client.RemoveObject(ctx, c.Bucket, key, minio.RemoveObjectOptions{})
return "", errors.New("保存文件记录失败")
}
return id, nil
}
func (o *Operations) OpenObject(ctx context.Context, id string, requirePublic bool) (io.ReadCloser, string, error) {
var obj model.StoredObject
if e := o.db.First(&obj, "id = ?", id).Error; e != nil {
return nil, "", errors.New("文件不存在")
}
if requirePublic && !obj.Public {
return nil, "", errors.New("文件不存在")
}
var row model.ModuleConfig
if e := o.db.First(&row, "name = ?", obj.ConfigName).Error; e != nil {
return nil, "", errors.New("历史存储配置不可用")
}
var c StorageConfig
if e := json.Unmarshal([]byte(row.Data), &c); e != nil {
return nil, "", errors.New("历史存储配置无效")
}
client, e := o.s3(c)
if e != nil {
return nil, "", e
}
object, e := client.GetObject(ctx, c.Bucket, obj.Key, minio.GetObjectOptions{})
if e != nil {
return nil, "", errors.New("读取存储失败")
}
if _, e = object.Stat(); e != nil {
_ = object.Close()
return nil, "", errors.New("文件暂不可用")
}
return object, obj.MIME, nil
}
func (o *Operations) StorageReferences() ([]map[string]any, error) {
type row struct {
ConfigName string
Count int64
}
var rows []row
e := o.db.Model(&model.StoredObject{}).Select("config_name, count(*) as count").Group("config_name").Scan(&rows).Error
out := []map[string]any{}
for _, r := range rows {
out = append(out, map[string]any{"revision": r.ConfigName, "references": r.Count})
}
return out, e
}
func RemoteObjectID(url string) string { return strings.TrimPrefix(url, "/api/media/") }
type boundedS3Transport struct {
base http.RoundTripper
host, scheme string
}
func (t boundedS3Transport) RoundTrip(r *http.Request) (*http.Response, error) {
if !strings.EqualFold(r.URL.Host, t.host) || r.URL.Scheme != t.scheme {
return nil, errors.New("拒绝存储重定向到未配置目标")
}
return t.base.RoundTrip(r)
}
func (o *Operations) PublicObjectLocation(id string) (string, error) {
var object model.StoredObject
if e := o.db.First(&object, "id = ? AND public = true", id).Error; e != nil {
return "", errors.New("文件不存在")
}
var row model.ModuleConfig
if e := o.db.First(&row, "name = ?", object.ConfigName).Error; e != nil {
return "", e
}
var c StorageConfig
if e := json.Unmarshal([]byte(row.Data), &c); e != nil {
return "", e
}
if c.CDN == "" {
return "", nil
}
return strings.TrimRight(c.CDN, "/") + "/" + object.Key, nil
}
// RemoveObject revokes public access before best-effort remote cleanup.
// Failed cleanup retains the historic configuration reference for manual retry.
func (o *Operations) RemoveObject(id string) error {
var obj model.StoredObject
if e := o.db.First(&obj, "id = ?", id).Error; e != nil {
return e
}
if e := o.db.Model(&obj).Update("public", false).Error; e != nil {
return e
}
var row model.ModuleConfig
if e := o.db.First(&row, "name = ?", obj.ConfigName).Error; e != nil {
return e
}
var c StorageConfig
if e := json.Unmarshal([]byte(row.Data), &c); e != nil {
return e
}
client, e := o.s3(c)
if e != nil {
return e
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
if e = client.RemoveObject(ctx, c.Bucket, obj.Key, minio.RemoveObjectOptions{}); e != nil {
return errors.New("远程文件清理失败,已撤销公开访问")
}
return o.db.Delete(&obj).Error
}