244 lines
7.4 KiB
Go
244 lines
7.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"
|
|
)
|
|
|
|
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
|
|
}
|
|
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
|
|
}
|