- 后端新增 storage 模块(minio_store),支持 MinIO 预签名 URL 上传与对象管理 - config 新增 MinIO 及存储等待相关配置项 - 算力节点新增 /compute/cache/prepare 缓存预下载接口(带校验和原子落盘) - 算力节点健康接口增加存储可用性探针 - SQL 迁移补充资源存储相关表结构 - Docker 新增 minio 服务与后端 minio 配置 - 补充 minio-compute-cache-plan 设计文档 Co-Authored-By: Claude <noreply@anthropic.com>
69 lines
2.7 KiB
Python
69 lines
2.7 KiB
Python
from __future__ import annotations
|
|
|
|
from datetime import timedelta
|
|
from functools import lru_cache
|
|
from io import BytesIO
|
|
from typing import Any
|
|
|
|
from minio import Minio
|
|
from minio.error import S3Error
|
|
|
|
from app.core.config import get_settings
|
|
|
|
|
|
class ObjectStorageError(RuntimeError):
|
|
pass
|
|
|
|
|
|
class MinioObjectStorage:
|
|
def __init__(self) -> None:
|
|
settings = get_settings()
|
|
endpoint = settings.minio_endpoint.replace("http://", "").replace("https://", "").rstrip("/")
|
|
self.client = Minio(endpoint, access_key=settings.minio_access_key, secret_key=settings.minio_secret_key, secure=settings.minio_secure)
|
|
self.bucket = settings.minio_bucket
|
|
|
|
def _ensure_enabled(self) -> None:
|
|
if not get_settings().minio_enabled:
|
|
raise ObjectStorageError("MinIO object storage is disabled")
|
|
|
|
def ensure_bucket(self) -> None:
|
|
self._ensure_enabled()
|
|
try:
|
|
if not self.client.bucket_exists(self.bucket):
|
|
self.client.make_bucket(self.bucket)
|
|
except S3Error as exc:
|
|
raise ObjectStorageError(str(exc)) from exc
|
|
|
|
def presigned_put(self, object_key: str, expires_seconds: int = 3600) -> str:
|
|
self._ensure_enabled()
|
|
self.ensure_bucket()
|
|
return self.client.presigned_put_object(self.bucket, object_key, expires=timedelta(seconds=expires_seconds))
|
|
|
|
def presigned_get(self, object_key: str, expires_seconds: int = 3600) -> str:
|
|
self._ensure_enabled()
|
|
self.ensure_bucket()
|
|
return self.client.presigned_get_object(self.bucket, object_key, expires=timedelta(seconds=expires_seconds))
|
|
|
|
def stat(self, object_key: str) -> dict[str, Any]:
|
|
self._ensure_enabled()
|
|
self.ensure_bucket()
|
|
try:
|
|
result = self.client.stat_object(self.bucket, object_key)
|
|
return {"object_key": object_key, "byte_size": result.size, "etag": result.etag, "last_modified": result.last_modified.isoformat() if result.last_modified else None}
|
|
except S3Error as exc:
|
|
raise ObjectStorageError(str(exc)) from exc
|
|
|
|
def put_bytes(self, object_key: str, content: bytes, content_type: str = "application/octet-stream") -> dict[str, Any]:
|
|
self._ensure_enabled()
|
|
self.ensure_bucket()
|
|
try:
|
|
result = self.client.put_object(self.bucket, object_key, BytesIO(content), len(content), content_type=content_type)
|
|
return {"bucket": self.bucket, "object_key": object_key, "etag": result.etag, "byte_size": len(content)}
|
|
except S3Error as exc:
|
|
raise ObjectStorageError(str(exc)) from exc
|
|
|
|
|
|
@lru_cache
|
|
def get_object_storage() -> MinioObjectStorage:
|
|
return MinioObjectStorage()
|