- 单用户资料接口补上距离:此前只有推荐/附近列表会算距离,资料页和 聊天头部因此无内容可显示。沿用同一套 haversine 与隐私开关。 - 语音/图片消息可限定会员发送,两个开关在管理端「运营配置」中修改 (迁移 032)。校验放在 persistMessageContext,HTTP 与 WebSocket 两条发送路径都覆盖;文本消息永不受限。 - 短信服务关闭时注册不再要求验证码:关掉之后没人能拿到验证码,继续 要求就等于关闭注册通道。重置密码不做同样放宽,那里缺验证码等于 凭手机号夺号。app/config 增加 smsVerification 供客户端决定表单形态。 - 修复 AI 托管账号之间不回复:原规则按「发送方是否托管账号」拦截, 把真人操作测试号的正常对话也挡了。改为标记 worker 自己写入的回复, 只对 AI 生成的消息跳过入队。 - ai.default_model_id 同时接受模型 ID 与名称,填名称时不再被 MySQL 静默转成 0 而使配置失效。 - 聊天媒体留存管理与清理任务(迁移 033,两台线上均已应用)。 新增集成测试均针对真实 MySQL:会员限制、免短信注册、AI 入队规则、 默认模型解析、资料距离。 Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
294 lines
12 KiB
Go
294 lines
12 KiB
Go
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)
|
|
}
|
|
}
|