init
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
# storage/__init__.py
|
||||
from .rustfs_storage import RustFSManager, rustfs_manager
|
||||
|
||||
__all__ = ['RustFSManager', 'rustfs_manager']
|
||||
|
||||
@@ -0,0 +1,107 @@
|
||||
# storage/init_storage.py
|
||||
"""初始化 RustFS 对象存储"""
|
||||
import asyncio
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
# 添加项目根目录和 src 目录到 Python 路径
|
||||
project_root = Path(__file__).parent.parent.parent
|
||||
src_root = Path(__file__).parent.parent
|
||||
sys.path.insert(0, str(project_root))
|
||||
sys.path.insert(0, str(src_root))
|
||||
|
||||
from storage.rustfs_storage import rustfs_manager
|
||||
from config.settings import settings
|
||||
from utils.logger import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
async def init_rustfs_storage():
|
||||
"""初始化 RustFS 对象存储"""
|
||||
try:
|
||||
# 连接到 RustFS (S3v4 API)
|
||||
await rustfs_manager.connect(
|
||||
endpoint=settings.RUSTFS_ENDPOINT,
|
||||
access_key=settings.RUSTFS_ACCESS_KEY,
|
||||
secret_key=settings.RUSTFS_SECRET_KEY,
|
||||
timeout=settings.RUSTFS_TIMEOUT
|
||||
)
|
||||
|
||||
logger.info("RustFS 对象存储初始化完成")
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"RustFS 对象存储初始化失败: {e}")
|
||||
return False
|
||||
|
||||
|
||||
async def test_storage():
|
||||
"""测试 RustFS 对象存储功能"""
|
||||
try:
|
||||
import json
|
||||
|
||||
# 测试上传 JSON
|
||||
test_data = {"test": True, "timestamp": "2024-01-01", "storage": "rustfs"}
|
||||
result = await rustfs_manager.upload_json_data(
|
||||
file_type='stp_files',
|
||||
json_data=test_data,
|
||||
file_hash='test-hash'
|
||||
)
|
||||
|
||||
logger.info(f"RustFS 测试上传成功: {result['object_key']}")
|
||||
|
||||
# 测试下载
|
||||
downloaded_bytes = await rustfs_manager.download_file(
|
||||
file_type='stp_files',
|
||||
object_key=result['object_key']
|
||||
)
|
||||
downloaded_data = json.loads(downloaded_bytes.decode('utf-8'))
|
||||
logger.info(f"RustFS 测试下载成功: {downloaded_data}")
|
||||
|
||||
# 测试预签名 URL
|
||||
url = await rustfs_manager.generate_presigned_url(
|
||||
file_type='stp_files',
|
||||
object_key=result['object_key'],
|
||||
expires=3600
|
||||
)
|
||||
logger.info(f"RustFS 预签名URL: {url}")
|
||||
|
||||
# 清理测试文件
|
||||
await rustfs_manager.delete_file(
|
||||
file_type='stp_files',
|
||||
object_key=result['object_key']
|
||||
)
|
||||
logger.info("RustFS 测试文件已清理")
|
||||
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"RustFS 存储测试失败: {e}")
|
||||
return False
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
async def main():
|
||||
try:
|
||||
print("=== 初始化 RustFS 对象存储 ===")
|
||||
# 初始化存储
|
||||
init_result = await init_rustfs_storage()
|
||||
if init_result:
|
||||
print("[OK] RustFS 连接成功")
|
||||
|
||||
print("\n=== 测试 RustFS 功能 ===")
|
||||
# 运行测试
|
||||
test_result = await test_storage()
|
||||
if test_result:
|
||||
print("[OK] RustFS 测试全部通过")
|
||||
else:
|
||||
print("[FAIL] RustFS 测试失败")
|
||||
|
||||
finally:
|
||||
# 关闭连接
|
||||
await rustfs_manager.close()
|
||||
print("\n=== 连接已关闭 ===")
|
||||
|
||||
# 运行主函数
|
||||
asyncio.run(main())
|
||||
@@ -0,0 +1,361 @@
|
||||
# storage/object_storage.py
|
||||
"""MinIO/S3 对象存储服务"""
|
||||
from minio import Minio
|
||||
from minio.error import S3Error
|
||||
from pathlib import Path
|
||||
from typing import Optional, BinaryIO
|
||||
from io import BytesIO
|
||||
from utils.logger import get_logger
|
||||
import hashlib
|
||||
import uuid
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
class ObjectStorageManager:
|
||||
"""对象存储管理器 - MinIO/S3兼容"""
|
||||
|
||||
def __init__(self):
|
||||
self.client: Optional[Minio] = None
|
||||
self.is_connected = False
|
||||
|
||||
# 桶名称
|
||||
self.buckets = {
|
||||
'stp_files': 'moldinsight-stp-files', # STP/STEP文件
|
||||
'geometry_data': 'moldinsight-geometry', # 几何数据JSON
|
||||
'mold_cavities': 'moldinsight-mold-cavities', # 模具型腔数据
|
||||
'html_files': 'moldinsight-html', # HTML报告文件
|
||||
'user_files': 'moldinsight-user-files' # 用户上传的其他文件
|
||||
}
|
||||
|
||||
async def connect(self, endpoint: str, access_key: str, secret_key: str,
|
||||
secure: bool = False):
|
||||
"""连接到MinIO/S3服务"""
|
||||
try:
|
||||
self.client = Minio(
|
||||
endpoint,
|
||||
access_key=access_key,
|
||||
secret_key=secret_key,
|
||||
secure=secure
|
||||
)
|
||||
|
||||
# 测试连接
|
||||
self.client.list_buckets()
|
||||
|
||||
self.is_connected = True
|
||||
logger.info(f"对象存储连接成功: {endpoint}")
|
||||
|
||||
# 确保所有桶都存在
|
||||
await self._ensure_buckets()
|
||||
|
||||
except S3Error as e:
|
||||
logger.error(f"对象存储连接失败: {e}")
|
||||
self.is_connected = False
|
||||
raise
|
||||
|
||||
async def _ensure_buckets(self):
|
||||
"""确保所有必要的桶都存在"""
|
||||
for bucket_name in self.buckets.values():
|
||||
try:
|
||||
if not self.client.bucket_exists(bucket_name):
|
||||
self.client.make_bucket(bucket_name)
|
||||
logger.info(f"创建存储桶: {bucket_name}")
|
||||
else:
|
||||
logger.debug(f"存储桶已存在: {bucket_name}")
|
||||
except S3Error as e:
|
||||
logger.error(f"创建存储桶失败 {bucket_name}: {e}")
|
||||
|
||||
def _generate_object_key(self, original_filename: str, prefix: str = '') -> str:
|
||||
"""生成对象存储的唯一键名"""
|
||||
# 提取文件扩展名
|
||||
ext = Path(original_filename).suffix
|
||||
|
||||
# 生成唯一ID
|
||||
unique_id = str(uuid.uuid4())
|
||||
|
||||
# 生成键名: prefix/unique_id + original_ext
|
||||
if prefix:
|
||||
return f"{prefix}/{unique_id}{ext}"
|
||||
return f"{unique_id}{ext}"
|
||||
|
||||
async def upload_stp_file(self, file_path: Path,
|
||||
original_filename: str) -> dict:
|
||||
"""上传STP文件到对象存储"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("对象存储未连接")
|
||||
|
||||
bucket_name = self.buckets['stp_files']
|
||||
|
||||
# 计算文件哈希
|
||||
file_hash = self._calculate_file_hash(file_path)
|
||||
|
||||
# 检查是否已存在
|
||||
existing_key = await self._find_file_by_hash(bucket_name, file_hash)
|
||||
if existing_key:
|
||||
logger.info(f"文件已存在,跳过上传: {existing_key}")
|
||||
return {
|
||||
'object_key': existing_key,
|
||||
'file_hash': file_hash,
|
||||
'already_exists': True
|
||||
}
|
||||
|
||||
# 生成唯一键名
|
||||
object_key = self._generate_object_key(
|
||||
original_filename,
|
||||
prefix='stp'
|
||||
)
|
||||
|
||||
# 上传文件
|
||||
try:
|
||||
result = self.client.fput_object(
|
||||
bucket_name,
|
||||
object_key,
|
||||
str(file_path),
|
||||
content_type='application/octet-stream'
|
||||
)
|
||||
|
||||
logger.info(f"STP文件上传成功: {object_key}")
|
||||
|
||||
return {
|
||||
'object_key': object_key,
|
||||
'file_hash': file_hash,
|
||||
'file_size': result.size,
|
||||
'etag': result.etag,
|
||||
'already_exists': False
|
||||
}
|
||||
except S3Error as e:
|
||||
logger.error(f"STP文件上传失败: {e}")
|
||||
raise
|
||||
|
||||
async def upload_geometry_data(self, geometry_json: dict,
|
||||
file_hash: str) -> dict:
|
||||
"""上传几何数据JSON到对象存储"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("对象存储未连接")
|
||||
|
||||
bucket_name = self.buckets['geometry_data']
|
||||
|
||||
# 使用文件哈希作为键名的一部分
|
||||
object_key = f"geometry/{file_hash}.json"
|
||||
|
||||
# 转换为字节
|
||||
import json
|
||||
json_bytes = json.dumps(geometry_json, ensure_ascii=False).encode('utf-8')
|
||||
|
||||
# 上传
|
||||
try:
|
||||
result = self.client.put_object(
|
||||
bucket_name,
|
||||
object_key,
|
||||
BytesIO(json_bytes),
|
||||
length=len(json_bytes),
|
||||
content_type='application/json'
|
||||
)
|
||||
|
||||
logger.info(f"几何数据上传成功: {object_key}")
|
||||
|
||||
return {
|
||||
'object_key': object_key,
|
||||
'file_size': result.size,
|
||||
'etag': result.etag
|
||||
}
|
||||
except S3Error as e:
|
||||
logger.error(f"几何数据上传失败: {e}")
|
||||
raise
|
||||
|
||||
async def upload_mold_cavity_data(self, cavity_json: dict,
|
||||
file_hash: str) -> dict:
|
||||
"""上传模具型腔数据到对象存储"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("对象存储未连接")
|
||||
|
||||
bucket_name = self.buckets['mold_cavities']
|
||||
object_key = f"mold-cavity/{file_hash}.json"
|
||||
|
||||
import json
|
||||
json_bytes = json.dumps(cavity_json, ensure_ascii=False).encode('utf-8')
|
||||
|
||||
try:
|
||||
result = self.client.put_object(
|
||||
bucket_name,
|
||||
object_key,
|
||||
BytesIO(json_bytes),
|
||||
length=len(json_bytes),
|
||||
content_type='application/json'
|
||||
)
|
||||
|
||||
logger.info(f"模具型腔数据上传成功: {object_key}")
|
||||
|
||||
return {
|
||||
'object_key': object_key,
|
||||
'file_size': result.size,
|
||||
'etag': result.etag
|
||||
}
|
||||
except S3Error as e:
|
||||
logger.error(f"模具型腔数据上传失败: {e}")
|
||||
raise
|
||||
|
||||
async def upload_html_file(self, html_content: str,
|
||||
original_filename: str,
|
||||
file_hash: str) -> dict:
|
||||
"""上传HTML文件到对象存储"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("对象存储未连接")
|
||||
|
||||
bucket_name = self.buckets['html_files']
|
||||
object_key = f"html/{file_hash}.html"
|
||||
|
||||
html_bytes = html_content.encode('utf-8')
|
||||
|
||||
try:
|
||||
result = self.client.put_object(
|
||||
bucket_name,
|
||||
object_key,
|
||||
BytesIO(html_bytes),
|
||||
length=len(html_bytes),
|
||||
content_type='text/html; charset=utf-8'
|
||||
)
|
||||
|
||||
logger.info(f"HTML文件上传成功: {object_key}")
|
||||
|
||||
return {
|
||||
'object_key': object_key,
|
||||
'file_size': result.size,
|
||||
'etag': result.etag
|
||||
}
|
||||
except S3Error as e:
|
||||
logger.error(f"HTML文件上传失败: {e}")
|
||||
raise
|
||||
|
||||
async def download_file(self, bucket_type: str,
|
||||
object_key: str) -> bytes:
|
||||
"""从对象存储下载文件"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("对象存储未连接")
|
||||
|
||||
bucket_name = self.buckets.get(bucket_type)
|
||||
if not bucket_name:
|
||||
raise ValueError(f"未知的桶类型: {bucket_type}")
|
||||
|
||||
try:
|
||||
response = self.client.get_object(bucket_name, object_key)
|
||||
data = response.read()
|
||||
response.close()
|
||||
response.release_conn()
|
||||
|
||||
logger.debug(f"文件下载成功: {object_key}")
|
||||
return data
|
||||
except S3Error as e:
|
||||
logger.error(f"文件下载失败 {object_key}: {e}")
|
||||
raise
|
||||
|
||||
async def get_presigned_url(self, bucket_type: str,
|
||||
object_key: str,
|
||||
expires: int = 3600) -> str:
|
||||
"""生成预签名URL(临时访问链接)"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("对象存储未连接")
|
||||
|
||||
bucket_name = self.buckets.get(bucket_type)
|
||||
if not bucket_name:
|
||||
raise ValueError(f"未知的桶类型: {bucket_type}")
|
||||
|
||||
try:
|
||||
url = self.client.presigned_get_object(
|
||||
bucket_name,
|
||||
object_key,
|
||||
expires=expires
|
||||
)
|
||||
return url
|
||||
except S3Error as e:
|
||||
logger.error(f"生成预签名URL失败: {e}")
|
||||
raise
|
||||
|
||||
async def delete_file(self, bucket_type: str, object_key: str):
|
||||
"""删除对象存储中的文件"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("对象存储未连接")
|
||||
|
||||
bucket_name = self.buckets.get(bucket_type)
|
||||
if not bucket_name:
|
||||
raise ValueError(f"未知的桶类型: {bucket_type}")
|
||||
|
||||
try:
|
||||
self.client.remove_object(bucket_name, object_key)
|
||||
logger.info(f"文件删除成功: {object_key}")
|
||||
except S3Error as e:
|
||||
logger.error(f"文件删除失败 {object_key}: {e}")
|
||||
raise
|
||||
|
||||
def _calculate_file_hash(self, file_path: Path) -> str:
|
||||
"""计算文件的SHA256哈希"""
|
||||
sha256_hash = hashlib.sha256()
|
||||
with open(file_path, 'rb') as f:
|
||||
for byte_block in iter(lambda: f.read(4096), b""):
|
||||
sha256_hash.update(byte_block)
|
||||
return sha256_hash.hexdigest()
|
||||
|
||||
async def _find_file_by_hash(self, bucket_name: str,
|
||||
file_hash: str) -> Optional[str]:
|
||||
"""根据哈希查找已存在的文件"""
|
||||
try:
|
||||
objects = self.client.list_objects(bucket_name, recursive=True)
|
||||
for obj in objects:
|
||||
# 从对象键中提取哈希(如果有)
|
||||
if file_hash in obj.object_name:
|
||||
return obj.object_name
|
||||
return None
|
||||
except S3Error as e:
|
||||
logger.warning(f"查找文件哈希失败: {e}")
|
||||
return None
|
||||
|
||||
async def get_file_info(self, bucket_type: str,
|
||||
object_key: str) -> dict:
|
||||
"""获取文件信息"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("对象存储未连接")
|
||||
|
||||
bucket_name = self.buckets.get(bucket_type)
|
||||
if not bucket_name:
|
||||
raise ValueError(f"未知的桶类型: {bucket_type}")
|
||||
|
||||
try:
|
||||
stat = self.client.stat_object(bucket_name, object_key)
|
||||
return {
|
||||
'size': stat.size,
|
||||
'etag': stat.etag,
|
||||
'content_type': stat.content_type,
|
||||
'last_modified': stat.last_modified
|
||||
}
|
||||
except S3Error as e:
|
||||
logger.error(f"获取文件信息失败: {e}")
|
||||
raise
|
||||
|
||||
async def list_files(self, bucket_type: str,
|
||||
prefix: str = '') -> list:
|
||||
"""列出存储桶中的文件"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("对象存储未连接")
|
||||
|
||||
bucket_name = self.buckets.get(bucket_type)
|
||||
if not bucket_name:
|
||||
raise ValueError(f"未知的桶类型: {bucket_type}")
|
||||
|
||||
try:
|
||||
objects = self.client.list_objects(bucket_name, prefix=prefix)
|
||||
return [
|
||||
{
|
||||
'object_key': obj.object_name,
|
||||
'size': obj.size,
|
||||
'etag': obj.etag,
|
||||
'last_modified': obj.last_modified
|
||||
}
|
||||
for obj in objects
|
||||
]
|
||||
except S3Error as e:
|
||||
logger.error(f"列出文件失败: {e}")
|
||||
raise
|
||||
|
||||
|
||||
# 全局对象存储管理器实例
|
||||
storage_manager = ObjectStorageManager()
|
||||
@@ -0,0 +1,348 @@
|
||||
# storage/rustfs_storage.py
|
||||
"""RustFS 对象存储服务 (S3v4 API 兼容)"""
|
||||
from minio import Minio
|
||||
from minio.error import S3Error
|
||||
from pathlib import Path
|
||||
from typing import Optional, Dict, Any
|
||||
from io import BytesIO
|
||||
from utils.logger import get_logger
|
||||
from datetime import timedelta
|
||||
import hashlib
|
||||
import uuid
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
class RustFSManager:
|
||||
"""RustFS 对象存储管理器 (使用 MinIO S3 客户端)"""
|
||||
|
||||
def __init__(self, project_name: str = "moldinsight"):
|
||||
self.client: Optional[Minio] = None
|
||||
self.is_connected = False
|
||||
self.project_name = project_name
|
||||
|
||||
# 使用单个项目桶,按类型组织文件
|
||||
self.bucket_name = f"{project_name}"
|
||||
|
||||
# 文件类型前缀(子目录结构)
|
||||
self.file_types = {
|
||||
'stp_files': 'stp-files',
|
||||
'geometry_data': 'geometry',
|
||||
'mold_cavities': 'mold-cavities',
|
||||
'html_files': 'html',
|
||||
'user_files': 'user-files'
|
||||
}
|
||||
|
||||
async def connect(self, endpoint: str, access_key: str, secret_key: str, timeout: int = 30):
|
||||
"""连接到 RustFS 服务"""
|
||||
try:
|
||||
# 提取端口号和主机
|
||||
from urllib.parse import urlparse
|
||||
parsed = urlparse(endpoint)
|
||||
host = parsed.netloc or parsed.path
|
||||
|
||||
# 创建 MinIO 客户端(S3v4 兼容)
|
||||
self.client = Minio(
|
||||
host,
|
||||
access_key=access_key,
|
||||
secret_key=secret_key,
|
||||
secure=False, # HTTP 而不是 HTTPS
|
||||
region='us-east-1'
|
||||
)
|
||||
|
||||
# 测试连接
|
||||
self.client.list_buckets()
|
||||
|
||||
self.is_connected = True
|
||||
logger.info(f"RustFS 连接成功: {endpoint}")
|
||||
|
||||
# 确保所有桶都存在
|
||||
await self._ensure_buckets()
|
||||
|
||||
except S3Error as e:
|
||||
logger.error(f"RustFS 连接失败: {e}")
|
||||
self.is_connected = False
|
||||
raise
|
||||
except Exception as e:
|
||||
logger.error(f"RustFS 初始化失败: {e}")
|
||||
self.is_connected = False
|
||||
raise
|
||||
|
||||
async def close(self):
|
||||
"""关闭连接"""
|
||||
# MinIO 客户端不需要显式关闭
|
||||
self.is_connected = False
|
||||
logger.info("RustFS 连接已关闭")
|
||||
|
||||
async def _ensure_buckets(self):
|
||||
"""确保项目存储桶存在"""
|
||||
try:
|
||||
if not self.client.bucket_exists(self.bucket_name):
|
||||
self.client.make_bucket(self.bucket_name)
|
||||
logger.info(f"创建项目存储桶: {self.bucket_name}")
|
||||
else:
|
||||
logger.debug(f"项目存储桶已存在: {self.bucket_name}")
|
||||
except S3Error as e:
|
||||
logger.error(f"创建存储桶失败 {self.bucket_name}: {e}")
|
||||
|
||||
def _generate_object_key(self, original_filename: str, file_type: str = '') -> str:
|
||||
"""生成对象存储的唯一键名"""
|
||||
ext = Path(original_filename).suffix
|
||||
unique_id = str(uuid.uuid4())
|
||||
|
||||
# 格式: {文件类型}/{唯一ID}.扩展名 (去掉项目名前缀)
|
||||
if file_type and file_type in self.file_types:
|
||||
type_prefix = self.file_types[file_type]
|
||||
return f"{type_prefix}/{unique_id}{ext}"
|
||||
|
||||
# 默认格式
|
||||
return f"misc/{unique_id}{ext}"
|
||||
|
||||
def _calculate_file_hash(self, file_path: Path) -> str:
|
||||
"""计算文件的SHA256哈希"""
|
||||
sha256_hash = hashlib.sha256()
|
||||
with open(file_path, 'rb') as f:
|
||||
for byte_block in iter(lambda: f.read(4096), b""):
|
||||
sha256_hash.update(byte_block)
|
||||
return sha256_hash.hexdigest()
|
||||
|
||||
async def upload_file(self, file_type: str, file_path: Path,
|
||||
original_filename: str,
|
||||
metadata: Optional[Dict] = None) -> Dict[str, Any]:
|
||||
"""上传文件到 RustFS"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
if file_type not in self.file_types:
|
||||
raise ValueError(f"未知的文件类型: {file_type}")
|
||||
|
||||
# 计算文件哈希
|
||||
file_hash = self._calculate_file_hash(file_path)
|
||||
|
||||
# 生成唯一键名
|
||||
object_key = self._generate_object_key(original_filename, file_type)
|
||||
|
||||
# 上传文件
|
||||
try:
|
||||
result = self.client.fput_object(
|
||||
self.bucket_name,
|
||||
object_key,
|
||||
str(file_path),
|
||||
content_type='application/octet-stream',
|
||||
metadata=metadata or {}
|
||||
)
|
||||
|
||||
logger.info(f"文件上传成功 RustFS: {self.bucket_name}/{object_key}")
|
||||
|
||||
# 获取文件大小
|
||||
file_size = file_path.stat().st_size
|
||||
|
||||
return {
|
||||
'object_key': object_key,
|
||||
'bucket': self.bucket_name,
|
||||
'file_hash': file_hash,
|
||||
'file_size': file_size,
|
||||
'etag': result.etag if hasattr(result, 'etag') else None
|
||||
}
|
||||
|
||||
except S3Error as e:
|
||||
logger.error(f"RustFS 上传失败: {e}")
|
||||
raise
|
||||
|
||||
async def upload_json_data(self, file_type: str,
|
||||
json_data: Dict[str, Any],
|
||||
file_hash: str) -> Dict[str, Any]:
|
||||
"""上传JSON数据到 RustFS"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
if file_type not in self.file_types:
|
||||
raise ValueError(f"未知的文件类型: {file_type}")
|
||||
|
||||
# 格式: {文件类型}/{文件哈希}.json (去掉项目名前缀)
|
||||
type_prefix = self.file_types[file_type]
|
||||
object_key = f"{type_prefix}/{file_hash}.json"
|
||||
|
||||
# 转换为字节
|
||||
import json
|
||||
json_bytes = json.dumps(json_data, ensure_ascii=False).encode('utf-8')
|
||||
|
||||
try:
|
||||
result = self.client.put_object(
|
||||
self.bucket_name,
|
||||
object_key,
|
||||
BytesIO(json_bytes),
|
||||
length=len(json_bytes),
|
||||
content_type='application/json'
|
||||
)
|
||||
|
||||
logger.info(f"JSON数据上传成功 RustFS: {self.bucket_name}/{object_key}")
|
||||
|
||||
return {
|
||||
'object_key': object_key,
|
||||
'bucket': self.bucket_name,
|
||||
'file_size': len(json_bytes),
|
||||
'etag': result.etag if hasattr(result, 'etag') else None
|
||||
}
|
||||
|
||||
except S3Error as e:
|
||||
logger.error(f"RustFS JSON上传失败: {e}")
|
||||
raise
|
||||
|
||||
async def download_file(self, file_type: str, object_key: str) -> bytes:
|
||||
"""从 RustFS 下载文件"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
if file_type not in self.file_types:
|
||||
raise ValueError(f"未知的文件类型: {file_type}")
|
||||
|
||||
try:
|
||||
response = self.client.get_object(self.bucket_name, object_key)
|
||||
data = response.read()
|
||||
response.close()
|
||||
response.release_conn()
|
||||
|
||||
logger.debug(f"文件下载成功: {self.bucket_name}/{object_key}")
|
||||
return data
|
||||
|
||||
except S3Error as e:
|
||||
logger.error(f"RustFS 下载失败: {e}")
|
||||
raise
|
||||
|
||||
async def get_file_info(self, file_type: str, object_key: str) -> Dict[str, Any]:
|
||||
"""获取文件信息"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
if file_type not in self.file_types:
|
||||
raise ValueError(f"未知的文件类型: {file_type}")
|
||||
|
||||
try:
|
||||
stat = self.client.stat_object(self.bucket_name, object_key)
|
||||
return {
|
||||
'size': stat.size,
|
||||
'etag': stat.etag,
|
||||
'content_type': stat.content_type,
|
||||
'last_modified': stat.last_modified
|
||||
}
|
||||
|
||||
except S3Error as e:
|
||||
logger.error(f"RustFS 获取文件信息失败: {e}")
|
||||
raise
|
||||
|
||||
async def delete_file(self, file_type: str, object_key: str):
|
||||
"""删除 RustFS 中的文件"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
if file_type not in self.file_types:
|
||||
raise ValueError(f"未知的文件类型: {file_type}")
|
||||
|
||||
try:
|
||||
self.client.remove_object(self.bucket_name, object_key)
|
||||
logger.info(f"文件删除成功: {self.bucket_name}/{object_key}")
|
||||
|
||||
except S3Error as e:
|
||||
logger.error(f"RustFS 删除失败: {e}")
|
||||
raise
|
||||
|
||||
async def list_files(self, file_type: str, prefix: str = '') -> list:
|
||||
"""列出存储桶中的文件"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
if file_type not in self.file_types:
|
||||
raise ValueError(f"未知的文件类型: {file_type}")
|
||||
|
||||
# 构建完整前缀:{文件类型}/... (去掉项目名前缀)
|
||||
type_prefix = self.file_types[file_type]
|
||||
full_prefix = f"{type_prefix}/"
|
||||
if prefix:
|
||||
full_prefix += prefix
|
||||
|
||||
try:
|
||||
objects = self.client.list_objects(self.bucket_name, prefix=full_prefix, recursive=True)
|
||||
return [
|
||||
{
|
||||
'object_key': obj.object_name,
|
||||
'size': obj.size,
|
||||
'etag': obj.etag,
|
||||
'last_modified': obj.last_modified
|
||||
}
|
||||
for obj in objects
|
||||
]
|
||||
|
||||
except S3Error as e:
|
||||
logger.error(f"RustFS 列出文件失败: {e}")
|
||||
raise
|
||||
|
||||
async def generate_presigned_url(self, file_type: str,
|
||||
object_key: str,
|
||||
expires: int = 3600,
|
||||
method: str = 'GET') -> str:
|
||||
"""生成预签名URL(临时访问链接)"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
if file_type not in self.file_types:
|
||||
raise ValueError(f"未知的文件类型: {file_type}")
|
||||
|
||||
try:
|
||||
url = self.client.presigned_get_object(
|
||||
self.bucket_name,
|
||||
object_key,
|
||||
expires=timedelta(seconds=expires)
|
||||
)
|
||||
return url
|
||||
|
||||
except S3Error as e:
|
||||
logger.error(f"RustFS 生成预签名URL失败: {e}")
|
||||
raise
|
||||
|
||||
async def file_exists(self, file_type: str, object_key: str) -> bool:
|
||||
"""检查文件是否存在"""
|
||||
try:
|
||||
await self.get_file_info(file_type, object_key)
|
||||
return True
|
||||
except:
|
||||
return False
|
||||
|
||||
async def get_storage_stats(self) -> Dict[str, Any]:
|
||||
"""获取存储统计信息"""
|
||||
try:
|
||||
buckets = self.client.list_buckets()
|
||||
total_objects = 0
|
||||
total_size = 0
|
||||
namespace_stats = {}
|
||||
|
||||
for bucket in buckets:
|
||||
objects = self.client.list_objects(bucket.name, recursive=True)
|
||||
bucket_count = 0
|
||||
bucket_size = 0
|
||||
|
||||
for obj in objects:
|
||||
bucket_count += 1
|
||||
bucket_size += obj.size
|
||||
|
||||
namespace_stats[bucket.name] = {
|
||||
'object_count': bucket_count,
|
||||
'total_size': bucket_size
|
||||
}
|
||||
|
||||
total_objects += bucket_count
|
||||
total_size += bucket_size
|
||||
|
||||
return {
|
||||
'total_objects': total_objects,
|
||||
'total_size': total_size,
|
||||
'namespace_stats': namespace_stats
|
||||
}
|
||||
|
||||
except S3Error as e:
|
||||
logger.error(f"RustFS 获取统计信息失败: {e}")
|
||||
raise
|
||||
|
||||
|
||||
# 全局 RustFS 管理器实例
|
||||
rustfs_manager = RustFSManager()
|
||||
Reference in New Issue
Block a user