Files
xiaoxia-saas/apps/api/app/core/storage.py
T
Xiaoxia AI 9656588bf2
CI/CD Pipeline / Validate Code Quality And Tests (push) Successful in 22m33s
feat(project-d): implement real video generation pipeline with FFmpeg
2026-06-19 10:54:20 +08:00

144 lines
4.7 KiB
Python

"""MinIO storage service for file uploads."""
from datetime import timedelta
from typing import BinaryIO
from urllib.parse import urlparse
import os
from minio import Minio
from minio.error import S3Error
from app.config import get_settings
class MinIOService:
"""MinIO storage service."""
def __init__(self):
settings = get_settings()
self.client = Minio(
settings.MINIO_ENDPOINT,
access_key=settings.MINIO_ACCESS_KEY,
secret_key=settings.MINIO_SECRET_KEY,
secure=settings.MINIO_SECURE,
)
self.bucket_name = settings.MINIO_BUCKET
self.public_url = settings.MINIO_PUBLIC_URL
self._ensure_bucket()
def _ensure_bucket(self):
try:
if not self.client.bucket_exists(self.bucket_name):
self.client.make_bucket(self.bucket_name)
policy = {
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Principal": {"AWS": "*"},
"Action": ["s3:GetObject"],
"Resource": [f"arn:aws:s3:::{self.bucket_name}/*"],
}
],
}
import json
self.client.set_bucket_policy(self.bucket_name, json.dumps(policy))
except S3Error as error:
print(f"Error ensuring bucket: {error}")
def upload_file(
self,
file_or_path,
storage_key: str,
content_type: str = "application/octet-stream",
) -> str:
"""
上传文件
Args:
file_or_path: 文件对象或本地文件路径
storage_key: 存储键
content_type: 内容类型
Returns:
文件 URL
"""
try:
# 如果是字符串路径,打开文件
if isinstance(file_or_path, str):
with open(file_or_path, "rb") as f:
file_size = os.path.getsize(file_or_path)
self.client.put_object(
self.bucket_name,
storage_key,
f,
file_size,
content_type=content_type,
)
else:
# 文件对象
file_or_path.seek(0, os.SEEK_END)
file_size = file_or_path.tell()
file_or_path.seek(0)
self.client.put_object(
self.bucket_name,
storage_key,
file_or_path,
file_size,
content_type=content_type,
)
return f"{self.public_url}/{self.bucket_name}/{storage_key}"
except S3Error as error:
raise Exception(f"Failed to upload file: {error}")
def get_url(self, storage_key: str) -> str:
return f"{self.public_url}/{self.bucket_name}/{storage_key}"
def get_download_url(self, storage_key_or_url: str, expires_seconds: int = 3600) -> str:
storage_key = self._normalize_storage_key(storage_key_or_url)
try:
return self.client.presigned_get_object(
self.bucket_name,
storage_key,
expires=timedelta(seconds=expires_seconds),
)
except S3Error:
return self.get_url(storage_key)
def _normalize_storage_key(self, storage_key_or_url: str) -> str:
prefix = f"/{self.bucket_name}/"
if storage_key_or_url.startswith("http://") or storage_key_or_url.startswith("https://"):
parsed = urlparse(storage_key_or_url)
if prefix in parsed.path:
return parsed.path.split(prefix, 1)[1]
return storage_key_or_url.lstrip("/")
def download_file(self, storage_key: str, local_path: str):
"""
下载文件到本地
Args:
storage_key: 存储键
local_path: 本地文件路径
"""
try:
os.makedirs(os.path.dirname(local_path), exist_ok=True)
self.client.fget_object(self.bucket_name, storage_key, local_path)
except S3Error as error:
raise Exception(f"Failed to download file: {error}")
def delete_file(self, storage_key: str):
try:
self.client.remove_object(self.bucket_name, storage_key)
except S3Error as error:
print(f"Error deleting file: {error}")
_minio_service = None
def get_minio_service() -> MinIOService:
global _minio_service
if _minio_service is None:
_minio_service = MinIOService()
return _minio_service