x
This commit is contained in:
@@ -52,7 +52,8 @@ class RustFSManager:
|
||||
)
|
||||
|
||||
# 测试连接
|
||||
self.client.list_buckets()
|
||||
import asyncio
|
||||
await asyncio.to_thread(self.client.list_buckets)
|
||||
|
||||
self.is_connected = True
|
||||
logger.info(f"RustFS 连接成功: {endpoint}")
|
||||
@@ -77,9 +78,17 @@ class RustFSManager:
|
||||
|
||||
async def _ensure_buckets(self):
|
||||
"""确保项目存储桶存在"""
|
||||
try:
|
||||
import asyncio
|
||||
|
||||
def _check_and_create():
|
||||
if not self.client.bucket_exists(self.bucket_name):
|
||||
self.client.make_bucket(self.bucket_name)
|
||||
return True
|
||||
return False
|
||||
|
||||
try:
|
||||
created = await asyncio.to_thread(_check_and_create)
|
||||
if created:
|
||||
logger.info(f"创建项目存储桶: {self.bucket_name}")
|
||||
else:
|
||||
logger.debug(f"项目存储桶已存在: {self.bucket_name}")
|
||||
@@ -111,6 +120,8 @@ class RustFSManager:
|
||||
original_filename: str,
|
||||
metadata: Optional[Dict] = None) -> Dict[str, Any]:
|
||||
"""上传文件到 RustFS"""
|
||||
import asyncio
|
||||
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
@@ -124,8 +135,8 @@ class RustFSManager:
|
||||
object_key = self._generate_object_key(original_filename, file_type)
|
||||
|
||||
# 上传文件
|
||||
try:
|
||||
result = self.client.fput_object(
|
||||
def _upload():
|
||||
return self.client.fput_object(
|
||||
self.bucket_name,
|
||||
object_key,
|
||||
str(file_path),
|
||||
@@ -133,6 +144,9 @@ class RustFSManager:
|
||||
metadata=metadata or {}
|
||||
)
|
||||
|
||||
try:
|
||||
result = await asyncio.to_thread(_upload)
|
||||
|
||||
logger.info(f"文件上传成功 RustFS: {self.bucket_name}/{object_key}")
|
||||
|
||||
# 获取文件大小
|
||||
@@ -154,6 +168,8 @@ class RustFSManager:
|
||||
json_data: Dict[str, Any],
|
||||
file_hash: str) -> Dict[str, Any]:
|
||||
"""上传JSON数据到 RustFS"""
|
||||
import asyncio
|
||||
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
@@ -168,8 +184,8 @@ class RustFSManager:
|
||||
import json
|
||||
json_bytes = json.dumps(json_data, ensure_ascii=False).encode('utf-8')
|
||||
|
||||
try:
|
||||
result = self.client.put_object(
|
||||
def _upload():
|
||||
return self.client.put_object(
|
||||
self.bucket_name,
|
||||
object_key,
|
||||
BytesIO(json_bytes),
|
||||
@@ -177,6 +193,9 @@ class RustFSManager:
|
||||
content_type='application/json'
|
||||
)
|
||||
|
||||
try:
|
||||
result = await asyncio.to_thread(_upload)
|
||||
|
||||
logger.info(f"JSON数据上传成功 RustFS: {self.bucket_name}/{object_key}")
|
||||
|
||||
return {
|
||||
@@ -192,18 +211,23 @@ class RustFSManager:
|
||||
|
||||
async def download_file(self, file_type: str, object_key: str) -> bytes:
|
||||
"""从 RustFS 下载文件"""
|
||||
import asyncio
|
||||
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
if file_type not in self.file_types:
|
||||
raise ValueError(f"未知的文件类型: {file_type}")
|
||||
|
||||
try:
|
||||
def _download():
|
||||
response = self.client.get_object(self.bucket_name, object_key)
|
||||
data = response.read()
|
||||
response.close()
|
||||
response.release_conn()
|
||||
return data
|
||||
|
||||
try:
|
||||
data = await asyncio.to_thread(_download)
|
||||
logger.debug(f"文件下载成功: {self.bucket_name}/{object_key}")
|
||||
return data
|
||||
|
||||
@@ -213,14 +237,19 @@ class RustFSManager:
|
||||
|
||||
async def get_file_info(self, file_type: str, object_key: str) -> Dict[str, Any]:
|
||||
"""获取文件信息"""
|
||||
import asyncio
|
||||
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
if file_type not in self.file_types:
|
||||
raise ValueError(f"未知的文件类型: {file_type}")
|
||||
|
||||
def _get_stat():
|
||||
return self.client.stat_object(self.bucket_name, object_key)
|
||||
|
||||
try:
|
||||
stat = self.client.stat_object(self.bucket_name, object_key)
|
||||
stat = await asyncio.to_thread(_get_stat)
|
||||
return {
|
||||
'size': stat.size,
|
||||
'etag': stat.etag,
|
||||
@@ -234,14 +263,19 @@ class RustFSManager:
|
||||
|
||||
async def delete_file(self, file_type: str, object_key: str):
|
||||
"""删除 RustFS 中的文件"""
|
||||
import asyncio
|
||||
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
if file_type not in self.file_types:
|
||||
raise ValueError(f"未知的文件类型: {file_type}")
|
||||
|
||||
try:
|
||||
def _delete():
|
||||
self.client.remove_object(self.bucket_name, object_key)
|
||||
|
||||
try:
|
||||
await asyncio.to_thread(_delete)
|
||||
logger.info(f"文件删除成功: {self.bucket_name}/{object_key}")
|
||||
|
||||
except S3Error as e:
|
||||
@@ -250,6 +284,8 @@ class RustFSManager:
|
||||
|
||||
async def list_files(self, file_type: str, prefix: str = '') -> list:
|
||||
"""列出存储桶中的文件"""
|
||||
import asyncio
|
||||
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
@@ -262,8 +298,11 @@ class RustFSManager:
|
||||
if prefix:
|
||||
full_prefix += prefix
|
||||
|
||||
def _list():
|
||||
return list(self.client.list_objects(self.bucket_name, prefix=full_prefix, recursive=True))
|
||||
|
||||
try:
|
||||
objects = self.client.list_objects(self.bucket_name, prefix=full_prefix, recursive=True)
|
||||
objects = await asyncio.to_thread(_list)
|
||||
return [
|
||||
{
|
||||
'object_key': obj.object_name,
|
||||
@@ -283,18 +322,23 @@ class RustFSManager:
|
||||
expires: int = 3600,
|
||||
method: str = 'GET') -> str:
|
||||
"""生成预签名URL(临时访问链接)"""
|
||||
import asyncio
|
||||
|
||||
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(
|
||||
def _generate_url():
|
||||
return self.client.presigned_get_object(
|
||||
self.bucket_name,
|
||||
object_key,
|
||||
expires=timedelta(seconds=expires)
|
||||
)
|
||||
|
||||
try:
|
||||
url = await asyncio.to_thread(_generate_url)
|
||||
return url
|
||||
|
||||
except S3Error as e:
|
||||
@@ -311,14 +355,16 @@ class RustFSManager:
|
||||
|
||||
async def get_storage_stats(self) -> Dict[str, Any]:
|
||||
"""获取存储统计信息"""
|
||||
try:
|
||||
import asyncio
|
||||
|
||||
def _get_stats():
|
||||
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)
|
||||
objects = list(self.client.list_objects(bucket.name, recursive=True))
|
||||
bucket_count = 0
|
||||
bucket_size = 0
|
||||
|
||||
@@ -340,6 +386,9 @@ class RustFSManager:
|
||||
'namespace_stats': namespace_stats
|
||||
}
|
||||
|
||||
try:
|
||||
return await asyncio.to_thread(_get_stats)
|
||||
|
||||
except S3Error as e:
|
||||
logger.error(f"RustFS 获取统计信息失败: {e}")
|
||||
raise
|
||||
|
||||
Reference in New Issue
Block a user