Files
zhuangxiu/services/api/app/integrations/storage.py
T

111 lines
4.0 KiB
Python

from dataclasses import dataclass
from collections.abc import Mapping
from io import BytesIO
from typing import Any
import boto3
from botocore.client import Config
from app.config import Settings
@dataclass(frozen=True)
class PresignedUpload:
method: str
url: str
object_key: str
expires_in: int
class S3Storage:
def __init__(self, settings: Settings | Mapping[str, Any]) -> None:
self.settings = settings
def _value(self, key: str, default: Any = "") -> Any:
if isinstance(self.settings, Mapping):
return self.settings.get(key, default)
return getattr(self.settings, key, default)
@property
def configured(self) -> bool:
values = (
self._value("s3_endpoint"),
self._value("s3_access_key_id"),
self._value("s3_secret_access_key"),
)
return all(values) and not any(value.startswith("replace_with_") for value in values)
def _client(self, public: bool = False):
endpoint = (
self._value("s3_public_endpoint")
if public and self._value("s3_public_endpoint")
else self._value("s3_endpoint")
)
return boto3.client(
"s3",
endpoint_url=endpoint,
aws_access_key_id=self._value("s3_access_key_id"),
aws_secret_access_key=self._value("s3_secret_access_key"),
region_name=self._value("s3_region", "us-east-1"),
config=Config(
signature_version="s3v4",
s3={"addressing_style": "path" if self._value("s3_force_path_style", True) else "auto"},
),
)
def presign_input_upload(self, object_key: str, content_type: str) -> PresignedUpload:
if not self.configured:
raise RuntimeError("MinIO/S3 is not configured.")
url = self._client(public=True).generate_presigned_url(
"put_object",
Params={
"Bucket": self._value("s3_bucket_inputs"),
"Key": object_key,
"ContentType": content_type,
},
ExpiresIn=int(self._value("s3_presigned_url_ttl_seconds", 900)),
)
return PresignedUpload(
method="PUT",
url=url,
object_key=object_key,
expires_in=int(self._value("s3_presigned_url_ttl_seconds", 900)),
)
def put_input(self, object_key: str, content: bytes, content_type: str) -> None:
self._put(self._value("s3_bucket_inputs"), object_key, content, content_type)
def put_output(self, object_key: str, content: bytes, content_type: str) -> None:
self._put(self._value("s3_bucket_derived"), object_key, content, content_type)
def get_output(self, object_key: str) -> tuple[bytes, str]:
if not self.configured:
raise RuntimeError("MinIO/S3 is not configured.")
response = self._client().get_object(
Bucket=self._value("s3_bucket_derived"),
Key=object_key,
)
return response["Body"].read(), response.get("ContentType", "application/octet-stream")
def put_render(self, object_key: str, content: bytes, content_type: str) -> None:
self._put(self._value("s3_bucket_renders"), object_key, content, content_type)
def get_render(self, object_key: str) -> tuple[bytes, str]:
if not self.configured:
raise RuntimeError("MinIO/S3 is not configured.")
response = self._client().get_object(
Bucket=self._value("s3_bucket_renders"),
Key=object_key,
)
return response["Body"].read(), response.get("ContentType", "application/octet-stream")
def _put(self, bucket: str, object_key: str, content: bytes, content_type: str) -> None:
if not self.configured:
raise RuntimeError("MinIO/S3 is not configured.")
self._client().upload_fileobj(
BytesIO(content),
bucket,
object_key,
ExtraArgs={"ContentType": content_type},
)