新增: - blob_store 二进制去重存储 + blob_dedup 测试 其它: - 后端 settings/auth/media_library/upload/operations 服务与 handler 调整 - 前端品牌(BrandLockup/BrandCropModal/BrandSeoPanel)、注册、媒体库、分析详情页、Header/SiteChrome 等更新
305 lines
9.4 KiB
Go
305 lines
9.4 KiB
Go
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
|
||
}
|