169 lines
4.4 KiB
Go
169 lines
4.4 KiB
Go
|
|
package service
|
||
|
|
|
||
|
|
import (
|
||
|
|
"fmt"
|
||
|
|
"io"
|
||
|
|
"os"
|
||
|
|
"time"
|
||
|
|
|
||
|
|
"seeyon-filesystem/config"
|
||
|
|
"seeyon-filesystem/model"
|
||
|
|
"seeyon-filesystem/storage"
|
||
|
|
)
|
||
|
|
|
||
|
|
// StorageTestService 存储连通性测试服务
|
||
|
|
type StorageTestService struct{}
|
||
|
|
|
||
|
|
func NewStorageTestService() *StorageTestService {
|
||
|
|
return &StorageTestService{}
|
||
|
|
}
|
||
|
|
|
||
|
|
// TestResult 测试结果
|
||
|
|
type TestResult struct {
|
||
|
|
Success bool `json:"success"`
|
||
|
|
Message string `json:"message"`
|
||
|
|
Latency int64 `json:"latency_ms"` // 耗时毫秒
|
||
|
|
Detail string `json:"detail,omitempty"`
|
||
|
|
}
|
||
|
|
|
||
|
|
// TestPolicy 测试存储策略连通性
|
||
|
|
func (s *StorageTestService) TestPolicy(policyID uint64) (*TestResult, error) {
|
||
|
|
var policy model.StoragePolicy
|
||
|
|
if err := config.DB.First(&policy, policyID).Error; err != nil {
|
||
|
|
return nil, fmt.Errorf("存储策略不存在")
|
||
|
|
}
|
||
|
|
|
||
|
|
switch policy.Type {
|
||
|
|
case "local":
|
||
|
|
return s.testLocal(policy), nil
|
||
|
|
case "minio":
|
||
|
|
return s.testMinIO(policy), nil
|
||
|
|
default:
|
||
|
|
return &TestResult{Success: false, Message: "不支持的存储类型"}, nil
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// TestNewPolicy 测试新配置(未保存)
|
||
|
|
func (s *StorageTestService) TestNewPolicy(policyType string, cfg model.PolicyConfig) *TestResult {
|
||
|
|
switch policyType {
|
||
|
|
case "local":
|
||
|
|
return s.testLocalConfig(cfg)
|
||
|
|
case "minio":
|
||
|
|
return s.testMinIOConfig(cfg)
|
||
|
|
default:
|
||
|
|
return &TestResult{Success: false, Message: "不支持的存储类型"}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// testLocal 测试本地存储
|
||
|
|
func (s *StorageTestService) testLocal(policy model.StoragePolicy) *TestResult {
|
||
|
|
return s.testLocalConfig(policy.Config)
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *StorageTestService) testLocalConfig(cfg model.PolicyConfig) *TestResult {
|
||
|
|
start := time.Now()
|
||
|
|
|
||
|
|
basePath := cfg.BasePath
|
||
|
|
if basePath == "" {
|
||
|
|
return &TestResult{Success: false, Message: "未配置存储路径"}
|
||
|
|
}
|
||
|
|
|
||
|
|
// 检查目录是否存在
|
||
|
|
info, err := os.Stat(basePath)
|
||
|
|
if err != nil {
|
||
|
|
if os.IsNotExist(err) {
|
||
|
|
// 尝试创建
|
||
|
|
if err := os.MkdirAll(basePath, 0755); err != nil {
|
||
|
|
return &TestResult{Success: false, Message: "目录不存在且无法创建: " + err.Error()}
|
||
|
|
}
|
||
|
|
return &TestResult{Success: true, Message: "目录已自动创建", Latency: time.Since(start).Milliseconds(), Detail: basePath}
|
||
|
|
}
|
||
|
|
return &TestResult{Success: false, Message: "访问目录失败: " + err.Error()}
|
||
|
|
}
|
||
|
|
|
||
|
|
if !info.IsDir() {
|
||
|
|
return &TestResult{Success: false, Message: "路径不是目录"}
|
||
|
|
}
|
||
|
|
|
||
|
|
// 测试写入权限
|
||
|
|
testFile := basePath + "/.connectivity_test"
|
||
|
|
if err := os.WriteFile(testFile, []byte("test"), 0644); err != nil {
|
||
|
|
return &TestResult{Success: false, Message: "目录无写入权限: " + err.Error()}
|
||
|
|
}
|
||
|
|
os.Remove(testFile)
|
||
|
|
|
||
|
|
return &TestResult{
|
||
|
|
Success: true,
|
||
|
|
Message: "连接成功, 目录可读写",
|
||
|
|
Latency: time.Since(start).Milliseconds(),
|
||
|
|
Detail: basePath,
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// testMinIO 测试MinIO连接
|
||
|
|
func (s *StorageTestService) testMinIO(policy model.StoragePolicy) *TestResult {
|
||
|
|
return s.testMinIOConfig(policy.Config)
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *StorageTestService) testMinIOConfig(cfg model.PolicyConfig) *TestResult {
|
||
|
|
start := time.Now()
|
||
|
|
|
||
|
|
if cfg.Endpoint == "" {
|
||
|
|
return &TestResult{Success: false, Message: "未配置Endpoint"}
|
||
|
|
}
|
||
|
|
if cfg.Bucket == "" {
|
||
|
|
return &TestResult{Success: false, Message: "未配置Bucket"}
|
||
|
|
}
|
||
|
|
|
||
|
|
// AccessKey/SecretKey 可选(公开桶不需要)
|
||
|
|
engine, err := storage.NewMinIOEngine(cfg.Endpoint, cfg.AccessKey, cfg.SecretKey, cfg.Bucket, cfg.UseSSL)
|
||
|
|
if err != nil {
|
||
|
|
return &TestResult{
|
||
|
|
Success: false,
|
||
|
|
Message: "连接失败: " + err.Error(),
|
||
|
|
Latency: time.Since(start).Milliseconds(),
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// 测试上传和下载
|
||
|
|
testKey := ".connectivity_test"
|
||
|
|
testContent := []byte("minio connectivity test")
|
||
|
|
if err := engine.Upload(testKey, bytesReader(testContent), int64(len(testContent))); err != nil {
|
||
|
|
return &TestResult{
|
||
|
|
Success: false,
|
||
|
|
Message: "上传测试失败: " + err.Error(),
|
||
|
|
Latency: time.Since(start).Milliseconds(),
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// 清理测试文件
|
||
|
|
engine.Delete(testKey)
|
||
|
|
|
||
|
|
latency := time.Since(start).Milliseconds()
|
||
|
|
return &TestResult{
|
||
|
|
Success: true,
|
||
|
|
Message: fmt.Sprintf("连接成功, Bucket=%s, 延迟=%dms", cfg.Bucket, latency),
|
||
|
|
Latency: latency,
|
||
|
|
Detail: fmt.Sprintf("%s/%s", cfg.Endpoint, cfg.Bucket),
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// bytesReader 创建bytes.Reader
|
||
|
|
func bytesReader(b []byte) *bytesReaderWrapper {
|
||
|
|
return &bytesReaderWrapper{data: b, offset: 0}
|
||
|
|
}
|
||
|
|
|
||
|
|
type bytesReaderWrapper struct {
|
||
|
|
data []byte
|
||
|
|
offset int
|
||
|
|
}
|
||
|
|
|
||
|
|
func (r *bytesReaderWrapper) Read(p []byte) (int, error) {
|
||
|
|
if r.offset >= len(r.data) {
|
||
|
|
return 0, io.EOF
|
||
|
|
}
|
||
|
|
n := copy(p, r.data[r.offset:])
|
||
|
|
r.offset += n
|
||
|
|
return n, nil
|
||
|
|
}
|