Spaces:
Running
Running
| """Supabase Storage integration for the Media-to-Media conversion APIs. | |
| Provides a lazy, connection-light async client that: | |
| 1. Ensures the destination bucket exists (creating it if needed). | |
| 2. Uploads converted media files. | |
| 3. Returns short-lived signed URLs (default 24 hours) for download. | |
| All storage3 calls are already async/HTTP-based, so no thread-pool offload is | |
| required here. The client is created on first use and closed during shutdown. | |
| """ | |
| from __future__ import annotations | |
| import asyncio | |
| import time | |
| from datetime import datetime, timezone | |
| from typing import Any, Optional | |
| from app.config import get_settings | |
| from app.core.logger import get_logger | |
| _logger = get_logger(__name__) | |
| _settings = get_settings() | |
| class MediaStorageError(Exception): | |
| """Raised when Supabase Storage operations fail.""" | |
| def __init__(self, message: str, status_code: int = 502) -> None: | |
| super().__init__(message) | |
| self.message = message | |
| self.status_code = status_code | |
| class MediaStorageService: | |
| """Async client for uploading converted media to Supabase Storage.""" | |
| def __init__(self) -> None: | |
| self._client: Optional[Any] = None | |
| self._client_lock = asyncio.Lock() | |
| self._known_buckets: set[str] = set() | |
| # ------------------------------------------------------------------ | |
| # Client management | |
| # ------------------------------------------------------------------ | |
| async def _get_client(self): | |
| """Lazily build the Supabase Storage async client.""" | |
| if self._client is None: | |
| async with self._client_lock: | |
| if self._client is None: | |
| url = _settings.supabase_url | |
| key = _settings.supabase_service_role_key | |
| if not url or not key: | |
| raise MediaStorageError( | |
| "Supabase is not configured. Set SUPABASE_URL and " | |
| "SUPABASE_SERVICE_ROLE_KEY to enable storage uploads.", | |
| status_code=503, | |
| ) | |
| from storage3 import AsyncStorageClient | |
| self._client = AsyncStorageClient( | |
| url, | |
| headers={ | |
| "apikey": key, | |
| "Authorization": f"Bearer {key}", | |
| }, | |
| timeout=60, | |
| ) | |
| _logger.info("Supabase Storage client initialized") | |
| return self._client | |
| async def close(self) -> None: | |
| async with self._client_lock: | |
| if self._client is not None: | |
| try: | |
| await self._client.aclose() | |
| except Exception as exc: # pragma: no cover - defensive | |
| _logger.debug("Error closing Supabase Storage client: %s", exc) | |
| self._client = None | |
| # ------------------------------------------------------------------ | |
| # Bucket management | |
| # ------------------------------------------------------------------ | |
| async def ensure_bucket(self, bucket_id: str) -> str: | |
| """Verify the bucket exists, creating it (private) if it does not.""" | |
| bucket_id = (bucket_id or _settings.supabase_storage_bucket or "media-convert").strip() | |
| if not bucket_id: | |
| raise MediaStorageError("A non-empty storage bucket name is required.", status_code=400) | |
| if bucket_id in self._known_buckets: | |
| return bucket_id | |
| client = await self._get_client() | |
| exists = False | |
| try: | |
| buckets = await client.list_buckets() | |
| exists = any(getattr(b, "id", None) == bucket_id for b in buckets) | |
| except Exception as exc: | |
| _logger.warning("Could not list Supabase buckets: %s", exc) | |
| if not exists: | |
| try: | |
| await client.create_bucket(bucket_id, name=bucket_id, options={"public": False}) | |
| _logger.info("Created Supabase Storage bucket: %s", bucket_id) | |
| except Exception as exc: | |
| raise MediaStorageError( | |
| f"Failed to create storage bucket '{bucket_id}': {exc}", status_code=502 | |
| ) from exc | |
| self._known_buckets.add(bucket_id) | |
| return bucket_id | |
| # ------------------------------------------------------------------ | |
| # Uploads & signed URLs | |
| # ------------------------------------------------------------------ | |
| async def upload_file(self, bucket_id: str, path: str, data: bytes, content_type: str) -> None: | |
| """Upload raw bytes to ``bucket_id`` under ``path``.""" | |
| client = await self._get_client() | |
| try: | |
| response = await client.from_(bucket_id).upload( | |
| path, | |
| data, | |
| file_options={"content-type": content_type, "cache-control": "3600"}, | |
| ) | |
| if response is None or getattr(response, "error", None): | |
| _logger.error("Supabase upload returned an error: %s", response) | |
| raise MediaStorageError( | |
| "Supabase Storage upload failed (unknown response).", status_code=502 | |
| ) | |
| except MediaStorageError: | |
| raise | |
| except Exception as exc: | |
| raise MediaStorageError( | |
| f"Supabase Storage upload failed for '{path}': {exc}", status_code=502 | |
| ) from exc | |
| async def create_signed_url(self, bucket_id: str, path: str, ttl_seconds: int) -> str: | |
| """Create a signed download URL valid for ``ttl_seconds``.""" | |
| client = await self._get_client() | |
| try: | |
| result = await client.from_(bucket_id).create_signed_url(path, int(ttl_seconds)) | |
| signed = getattr(result, "signed_url", None) or (result.get("signedURL") if isinstance(result, dict) else None) | |
| if not signed: | |
| raise MediaStorageError( | |
| "Supabase returned an empty signed URL.", status_code=502 | |
| ) | |
| return signed | |
| except MediaStorageError: | |
| raise | |
| except Exception as exc: | |
| raise MediaStorageError( | |
| f"Failed to create signed URL for '{path}': {exc}", status_code=502 | |
| ) from exc | |
| _storage_service: Optional[MediaStorageService] = None | |
| _storage_service_lock = asyncio.Lock() | |
| async def get_storage_service() -> MediaStorageService: | |
| """Return the shared MediaStorageService singleton.""" | |
| global _storage_service | |
| if _storage_service is None: | |
| async with _storage_service_lock: | |
| if _storage_service is None: | |
| _storage_service = MediaStorageService() | |
| return _storage_service | |
| async def close_storage_service() -> None: | |
| global _storage_service | |
| if _storage_service is not None: | |
| await _storage_service.close() | |
| _storage_service = None | |
| def iso_expiry(ttl_seconds: int) -> str: | |
| """ISO-8601 UTC timestamp ``ttl_seconds`` from now.""" | |
| return datetime.fromtimestamp(time.time() + ttl_seconds, tz=timezone.utc).isoformat() | |