PCNUSMSE/transcript_service
0
1"""OSS云存储服务模块2 3提供阿里云OSS文件上传、下载和管理功能。4"""5 6import os7import uuid8from datetime import datetime, timedelta9from pathlib import Path10from typing import List, Optional, Tuple11import asyncio12import aiohttp13import oss214from oss2.exceptions import OssError15 16from ..core.config import get_config17from ..utils.logger import get_task_logger18 19 20class OSSService:21 """OSS云存储服务"""22 23 def __init__(self):24 """初始化OSS服务"""25 self.config = get_config()26 self.oss_config = self.config.oss27 28 # 初始化OSS客户端29 auth = oss2.Auth(30 self.oss_config.access_key_id,31 self.oss_config.access_key_secret32 )33 self.bucket = oss2.Bucket(34 auth,35 self.oss_config.endpoint,36 self.oss_config.bucket_name37 )38 39 self.logger = get_task_logger(logger_name="transcript_service.oss")40 41 def _generate_object_key(self, filename: str, task_id: str) -> str:42 """生成OSS对象键名43 44 Args:45 filename: 原始文件名46 task_id: 任务ID47 48 Returns:49 OSS对象键名50 """51 now = datetime.now()52 date_path = now.strftime("%Y/%m/%d")53 timestamp = now.strftime("%Y%m%d_%H%M%S")54 55 # 获取文件扩展名56 file_ext = Path(filename).suffix57 safe_filename = f"{timestamp}_{task_id}_{uuid.uuid4().hex[:8]}{file_ext}"58 59 return f"{self.oss_config.temp_prefix}/{date_path}/{safe_filename}"60 61 async def upload_file(self, file_path: Path, task_id: str) -> Tuple[bool, str, Optional[str]]:62 """上传文件到OSS63 64 Args:65 file_path: 本地文件路径66 task_id: 任务ID67 68 Returns:69 (是否成功, 公网URL或错误信息, 对象键名)70 """71 try:72 self.logger.info(f"开始上传文件到OSS: {file_path.name}")73 74 # 生成对象键名75 object_key = self._generate_object_key(file_path.name, task_id)76 77 # 上传文件并设置公共读取权限78 try:79 # 首先上传文件80 self.bucket.put_object_from_file(object_key, str(file_path))81 82 # 设置对象ACL为公共读取83 self.bucket.put_object_acl(object_key, oss2.OBJECT_ACL_PUBLIC_READ)84 85 # 生成公网访问URL86 url = self._generate_public_url(object_key)87 self.logger.info(f"文件上传成功: {object_key}, URL: {url}")88 return True, url, object_key89 90 except oss2.exceptions.OssError as oss_err:91 # 如果设置ACL失败,尝试使用签名URL92 if 'public-read' in str(oss_err).lower():93 self.logger.warning(f"ACL设置失败,使用签名URL: {oss_err}")94 url = self._generate_signed_url(object_key)95 self.logger.info(f"文件上传成功: {object_key}, URL: {url}")96 return True, url, object_key97 else:98 raise99 100 except OssError as e:101 error_msg = f"OSS错误: {str(e)}"102 self.logger.error(error_msg)103 return False, error_msg, None104 except Exception as e:105 error_msg = f"上传文件时发生未知错误: {str(e)}"106 self.logger.exception(error_msg)107 return False, error_msg, None108 109 async def upload_multiple_files(self, file_paths: List[Path], task_id: str) -> List[Tuple[str, bool, str, Optional[str]]]:110 """批量上传文件到OSS111 112 Args:113 file_paths: 本地文件路径列表114 task_id: 任务ID115 116 Returns:117 [(文件名, 是否成功, URL或错误信息, 对象键名), ...]118 """119 results = []120 121 # 创建异步任务122 tasks = []123 for file_path in file_paths:124 task = self._upload_single_file_async(file_path, task_id)125 tasks.append((file_path.name, task))126 127 # 等待所有上传完成128 for filename, task in tasks:129 success, url_or_error, object_key = await task130 results.append((filename, success, url_or_error, object_key))131 132 return results133 134 async def _upload_single_file_async(self, file_path: Path, task_id: str) -> Tuple[bool, str, Optional[str]]:135 """异步上传单个文件"""136 return await asyncio.get_event_loop().run_in_executor(137 None, 138 lambda: asyncio.run(self.upload_file(file_path, task_id))139 )140 141 def _generate_public_url(self, object_key: str) -> str:142 """生成公网访问URL143 144 Args:145 object_key: OSS对象键名146 147 Returns:148 公网访问URL149 """150 # 生成简单的公网访问URL(不带签名)151 # 正确的格式: https://bucket-name.endpoint/object-key152 # 注意: endpoint不能包含协议前缀153 endpoint = self.oss_config.endpoint154 if endpoint.startswith('http://'):155 endpoint = endpoint[7:]156 elif endpoint.startswith('https://'):157 endpoint = endpoint[8:]158 159 # 构造公网URL - 注意这里的格式必须正确160 url = f"https://{self.oss_config.bucket_name}.{endpoint}/{object_key}"161 162 # 记录生成的URL以便调试163 self.logger.debug(f"生成公网URL: {url}")164 165 return url166 167 def _generate_signed_url(self, object_key: str) -> str:168 """生成签名URL(备用方案)169 170 Args:171 object_key: OSS对象键名172 173 Returns:174 签名URL175 """176 # 生成有时效性的签名URL177 expire_time = int((datetime.now() + timedelta(hours=self.oss_config.url_expire_hours)).timestamp())178 url = self.bucket.sign_url('GET', object_key, expire_time)179 return url180 181 def delete_file(self, object_key: str) -> bool:182 """删除OSS文件183 184 Args:185 object_key: OSS对象键名186 187 Returns:188 是否删除成功189 """190 try:191 self.bucket.delete_object(object_key)192 self.logger.info(f"文件删除成功: {object_key}")193 return True194 except OssError as e:195 self.logger.error(f"删除文件失败: {object_key}, 错误: {str(e)}")196 return False197 except Exception as e:198 self.logger.exception(f"删除文件时发生未知错误: {object_key}, 错误: {str(e)}")199 return False200 201 def cleanup_old_files(self, days: Optional[int] = None) -> int:202 """清理过期的临时文件203 204 Args:205 days: 保留天数,默认使用配置中的值206 207 Returns:208 删除的文件数量209 """210 cleanup_days = days or self.oss_config.auto_cleanup_days211 cutoff_date = datetime.now() - timedelta(days=cleanup_days)212 213 deleted_count = 0214 prefix = self.oss_config.temp_prefix215 216 try:217 # 列出所有临时文件218 for obj in oss2.ObjectIterator(self.bucket, prefix=prefix):219 # 检查文件最后修改时间220 if obj.last_modified.replace(tzinfo=None) < cutoff_date:221 if self.delete_file(obj.key):222 deleted_count += 1223 224 self.logger.info(f"清理完成,删除了 {deleted_count} 个过期文件")225 return deleted_count226 227 except Exception as e:228 self.logger.exception(f"清理过期文件时发生错误: {str(e)}")229 return deleted_count230 231 def get_file_info(self, object_key: str) -> Optional[dict]:232 """获取文件信息233 234 Args:235 object_key: OSS对象键名236 237 Returns:238 文件信息字典239 """240 try:241 info = self.bucket.head_object(object_key)242 return {243 'size': info.content_length,244 'last_modified': info.last_modified,245 'etag': info.etag,246 'content_type': info.content_type247 }248 except OssError as e:249 self.logger.error(f"获取文件信息失败: {object_key}, 错误: {str(e)}")250 return None251 252 def check_bucket_exists(self) -> bool:253 """检查存储桶是否存在254 255 Returns:256 存储桶是否存在257 """258 try:259 return self.bucket.bucket_exists()260 except Exception as e:261 self.logger.error(f"检查存储桶失败: {str(e)}")262 return False263 264 def get_bucket_info(self) -> Optional[dict]:265 """获取存储桶信息266 267 Returns:268 存储桶信息269 """270 try:271 info = self.bucket.get_bucket_info()272 return {273 'name': info.name,274 'location': info.location,275 'creation_date': info.creation_date,276 'storage_class': info.storage_class277 }278 except Exception as e:279 self.logger.error(f"获取存储桶信息失败: {str(e)}")280 return None281 282 283# 全局OSS服务实例284oss_service = OSSService()285 286 287def get_oss_service() -> OSSService:288 """获取OSS服务实例289 290 Returns:291 OSS服务实例292 """293 return oss_service