package app import ( "context" "fmt" "io" "net/http" "net/url" "os" "path/filepath" "regexp" "strings" "time" aliyunoss "github.com/aliyun/alibabacloud-oss-go-sdk-v2/oss" aliyuncredentials "github.com/aliyun/alibabacloud-oss-go-sdk-v2/oss/credentials" huaweiobs "github.com/huaweicloud/huaweicloud-sdk-go-obs/obs" qiniuauth "github.com/qiniu/go-sdk/v7/auth/qbox" qiniustorage "github.com/qiniu/go-sdk/v7/storage" qiniucredentials "github.com/qiniu/go-sdk/v7/storagev2/credentials" qiniuhttpclient "github.com/qiniu/go-sdk/v7/storagev2/http_client" qiniuuploader "github.com/qiniu/go-sdk/v7/storagev2/uploader" tencentcos "github.com/tencentyun/cos-go-sdk-v5" ) var ( storageBucketPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]{1,126}[A-Za-z0-9]$`) storagePrefixPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9/_-]{0,119}$`) ) type mediaObjectStorage interface { Put(context.Context, string, string, int64, io.Reader) error Delete(context.Context, string) error Bucket() string } func storageProviderName(provider string) string { switch provider { case "local": return "本地存储" case "aliyun_oss": return "阿里云 OSS" case "tencent_cos": return "腾讯云 COS" case "qiniu": return "七牛云存储" case "huawei_obs": return "华为云 OBS" case "huawei_flexus": return "华为云 Flexus 对象存储" default: return provider } } func parseStorageHTTPSURL(raw string, allowPath bool) (*url.URL, error) { parsed, err := url.Parse(strings.TrimSpace(raw)) if err != nil || parsed.Scheme != "https" || parsed.Host == "" || parsed.User != nil || parsed.RawQuery != "" || parsed.Fragment != "" { return nil, fmt.Errorf("必须是无账号、查询参数和片段的有效 HTTPS 地址") } if !allowPath && strings.Trim(parsed.EscapedPath(), "/") != "" { return nil, fmt.Errorf("地址不能包含路径") } return parsed, nil } func validateLocalStorageDirectory(raw string) error { directory := strings.TrimSpace(raw) if directory == "" { return fmt.Errorf("请填写本地存储目录") } abs, err := filepath.Abs(filepath.Clean(directory)) if err != nil { return fmt.Errorf("本地存储目录无效") } root := filepath.VolumeName(abs) + string(os.PathSeparator) if strings.EqualFold(filepath.Clean(abs), filepath.Clean(root)) { return fmt.Errorf("本地存储目录不能是磁盘根目录") } if info, statErr := os.Stat(abs); statErr == nil && !info.IsDir() { return fmt.Errorf("本地存储目录指向了文件") } else if statErr != nil && !os.IsNotExist(statErr) { return fmt.Errorf("无法访问本地存储目录") } return nil } func (a *App) validateStorageProviderConfig(ctx context.Context) error { provider := a.configPlain(ctx, "storage.provider", "local") providers := []string{"local", "aliyun_oss", "tencent_cos", "qiniu", "huawei_obs", "huawei_flexus"} if !containsString(providers, provider) { return fmt.Errorf("不支持的文件存储厂商 %q", provider) } for _, spec := range integrationSpecs["storage"] { if !spec.Required || (len(spec.Providers) > 0 && !containsString(spec.Providers, provider)) { continue } if strings.TrimSpace(a.configPlain(ctx, spec.Key, "")) == "" { return fmt.Errorf("请填写%s", spec.Label) } } prefix := strings.Trim(strings.TrimSpace(a.configPlain(ctx, "storage.object_prefix", "media")), "/") if !storagePrefixPattern.MatchString(prefix) || strings.Contains(prefix, "//") || strings.Contains(prefix, "..") { return fmt.Errorf("云端对象前缀格式无效") } if provider == "local" { if err := validateLocalStorageDirectory(a.configPlain(ctx, "storage.local.directory", a.config.MediaDir)); err != nil { return err } baseURL := strings.TrimSpace(a.configPlain(ctx, "storage.local.public_base_url", "")) if baseURL != "" { parsed, err := url.Parse(baseURL) if err != nil || parsed.Host == "" || parsed.User != nil || parsed.RawQuery != "" || parsed.Fragment != "" || (parsed.Scheme != "http" && parsed.Scheme != "https") { return fmt.Errorf("本地公开访问地址无效") } if a.config.Environment == "production" && parsed.Scheme != "https" { return fmt.Errorf("生产环境本地公开访问地址必须使用 HTTPS") } } return nil } keyPrefix := "storage." + provider + "." if provider == "qiniu" { if !storageBucketPattern.MatchString(a.configPlain(ctx, keyPrefix+"bucket", "")) { return fmt.Errorf("七牛云空间名称格式无效") } } else { if _, err := parseStorageHTTPSURL(a.configPlain(ctx, keyPrefix+"endpoint", ""), false); err != nil { return fmt.Errorf("%s Endpoint %v", storageProviderName(provider), err) } if !storageBucketPattern.MatchString(a.configPlain(ctx, keyPrefix+"bucket", "")) { return fmt.Errorf("%s Bucket 名称格式无效", storageProviderName(provider)) } } if _, err := parseStorageHTTPSURL(a.configPlain(ctx, keyPrefix+"public_base_url", ""), true); err != nil { return fmt.Errorf("%s文件访问域名%v", storageProviderName(provider), err) } return nil } func storagePublicURL(baseURL, objectKey string) string { return strings.TrimRight(strings.TrimSpace(baseURL), "/") + "/" + strings.TrimLeft(objectKey, "/") } func storageHTTPClient() *http.Client { return &http.Client{ Timeout: 30 * time.Second, CheckRedirect: func(_ *http.Request, _ []*http.Request) error { return http.ErrUseLastResponse }, } } type aliyunOSSStorage struct { bucket string client *aliyunoss.Client } func (s *aliyunOSSStorage) Bucket() string { return s.bucket } func (s *aliyunOSSStorage) Put(ctx context.Context, key, contentType string, size int64, body io.Reader) error { _, err := s.client.PutObject(ctx, &aliyunoss.PutObjectRequest{ Bucket: aliyunoss.Ptr(s.bucket), Key: aliyunoss.Ptr(key), Body: body, ContentType: aliyunoss.Ptr(contentType), ContentLength: aliyunoss.Ptr(size), }) return err } func (s *aliyunOSSStorage) Delete(ctx context.Context, key string) error { _, err := s.client.DeleteObject(ctx, &aliyunoss.DeleteObjectRequest{Bucket: aliyunoss.Ptr(s.bucket), Key: aliyunoss.Ptr(key)}) return err } type tencentCOSStorage struct { bucket string client *tencentcos.Client } func (s *tencentCOSStorage) Bucket() string { return s.bucket } func (s *tencentCOSStorage) Put(ctx context.Context, key, contentType string, size int64, body io.Reader) error { _, err := s.client.Object.Put(ctx, key, body, &tencentcos.ObjectPutOptions{ObjectPutHeaderOptions: &tencentcos.ObjectPutHeaderOptions{ContentType: contentType, ContentLength: size}}) return err } func (s *tencentCOSStorage) Delete(ctx context.Context, key string) error { _, err := s.client.Object.Delete(ctx, key) return err } type qiniuStorage struct { bucket string uploader *qiniuuploader.UploadManager deleteMac *qiniuauth.Mac } func (s *qiniuStorage) Bucket() string { return s.bucket } func (s *qiniuStorage) Put(ctx context.Context, key, contentType string, _ int64, body io.Reader) error { return s.uploader.UploadReader(ctx, body, &qiniuuploader.ObjectOptions{BucketName: s.bucket, ObjectName: &key, FileName: filepath.Base(key), ContentType: contentType}, nil) } func (s *qiniuStorage) Delete(_ context.Context, key string) error { manager := qiniustorage.NewBucketManager(s.deleteMac, &qiniustorage.Config{UseHTTPS: true}) return manager.Delete(s.bucket, key) } type huaweiOBSStorage struct { bucket string client *huaweiobs.ObsClient } func (s *huaweiOBSStorage) Bucket() string { return s.bucket } func (s *huaweiOBSStorage) Put(_ context.Context, key, contentType string, size int64, body io.Reader) error { _, err := s.client.PutObject(&huaweiobs.PutObjectInput{PutObjectBasicInput: huaweiobs.PutObjectBasicInput{ ObjectOperationInput: huaweiobs.ObjectOperationInput{Bucket: s.bucket, Key: key}, HttpHeader: huaweiobs.HttpHeader{ContentType: contentType}, ContentLength: size, }, Body: body}) return err } func (s *huaweiOBSStorage) Delete(_ context.Context, key string) error { _, err := s.client.DeleteObject(&huaweiobs.DeleteObjectInput{Bucket: s.bucket, Key: key}) return err } func (a *App) newMediaObjectStorage(ctx context.Context) (mediaObjectStorage, func(), error) { if err := a.validateStorageProviderConfig(ctx); err != nil { return nil, func() {}, err } provider := a.configPlain(ctx, "storage.provider", "local") return a.newMediaObjectStorageForProvider(ctx, provider, "") } // Historical media must be deleted through the provider that originally // stored it, even when uploads have since moved to another provider. The // credentials remain in system_configs and the recorded bucket protects us // from deleting an object with the same key from a newly configured bucket. func (a *App) newMediaObjectStorageForProvider(ctx context.Context, provider, recordedBucket string) (mediaObjectStorage, func(), error) { provider = strings.TrimSpace(provider) if provider == "local" { directory := a.configPlain(ctx, "storage.local.directory", a.config.MediaDir) if err := validateLocalStorageDirectory(directory); err != nil { return nil, func() {}, err } return localMediaStorage{directory: directory}, func() {}, nil } keyPrefix := "storage." + provider + "." configuredBucket := strings.TrimSpace(a.configPlain(ctx, keyPrefix+"bucket", "")) if configuredBucket == "" { return nil, func() {}, fmt.Errorf("%s Bucket 未配置", storageProviderName(provider)) } if recordedBucket != "" && configuredBucket != recordedBucket { return nil, func() {}, fmt.Errorf("记录 Bucket %q 与当前配置 %q 不一致", recordedBucket, configuredBucket) } switch provider { case "aliyun_oss": if a.configPlain(ctx, keyPrefix+"access_key_id", "") == "" || a.configPlain(ctx, keyPrefix+"access_key_secret", "") == "" { return nil, func() {}, fmt.Errorf("阿里云 OSS 凭证未配置") } config := aliyunoss.LoadDefaultConfig(). WithCredentialsProvider(aliyuncredentials.NewStaticCredentialsProvider(a.configPlain(ctx, keyPrefix+"access_key_id", ""), a.configPlain(ctx, keyPrefix+"access_key_secret", ""))). WithRegion(a.configPlain(ctx, keyPrefix+"region", "")). WithEndpoint(a.configPlain(ctx, keyPrefix+"endpoint", "")). WithConnectTimeout(5 * time.Second). WithReadWriteTimeout(30 * time.Second). WithRetryMaxAttempts(3) return &aliyunOSSStorage{bucket: configuredBucket, client: aliyunoss.NewClient(config)}, func() {}, nil case "tencent_cos": endpoint, err := url.Parse(a.configPlain(ctx, keyPrefix+"endpoint", "")) if err != nil || endpoint.Host == "" || a.configPlain(ctx, keyPrefix+"secret_id", "") == "" || a.configPlain(ctx, keyPrefix+"secret_key", "") == "" { return nil, func() {}, fmt.Errorf("腾讯云 COS 地址或凭证未配置") } client := tencentcos.NewClient(&tencentcos.BaseURL{BucketURL: endpoint}, &http.Client{ Timeout: 30 * time.Second, Transport: &tencentcos.AuthorizationTransport{SecretID: a.configPlain(ctx, keyPrefix+"secret_id", ""), SecretKey: a.configPlain(ctx, keyPrefix+"secret_key", "")}, }) return &tencentCOSStorage{bucket: configuredBucket, client: client}, func() {}, nil case "qiniu": accessKey, secretKey := a.configPlain(ctx, keyPrefix+"access_key", ""), a.configPlain(ctx, keyPrefix+"secret_key", "") if accessKey == "" || secretKey == "" { return nil, func() {}, fmt.Errorf("七牛云凭证未配置") } manager := qiniuuploader.NewUploadManager(&qiniuuploader.UploadManagerOptions{Options: qiniuhttpclient.Options{Credentials: qiniucredentials.NewCredentials(accessKey, secretKey), BasicHTTPClient: storageHTTPClient()}, MultiPartsThreshold: 8 << 20, PartSize: 4 << 20, Concurrency: 2}) return &qiniuStorage{bucket: configuredBucket, uploader: manager, deleteMac: qiniuauth.NewMac(accessKey, secretKey)}, func() {}, nil case "huawei_obs", "huawei_flexus": if a.configPlain(ctx, keyPrefix+"access_key", "") == "" || a.configPlain(ctx, keyPrefix+"secret_key", "") == "" { return nil, func() {}, fmt.Errorf("%s凭证未配置", storageProviderName(provider)) } client, err := huaweiobs.New(a.configPlain(ctx, keyPrefix+"access_key", ""), a.configPlain(ctx, keyPrefix+"secret_key", ""), a.configPlain(ctx, keyPrefix+"endpoint", ""), huaweiobs.WithConnectTimeout(5), huaweiobs.WithSocketTimeout(30), huaweiobs.WithMaxRetryCount(2)) if err != nil { return nil, func() {}, err } return &huaweiOBSStorage{bucket: configuredBucket, client: client}, client.Close, nil default: return nil, func() {}, fmt.Errorf("不支持的文件存储厂商 %q", provider) } }