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 }