Spaces:
Running
Running
| from __future__ import annotations | |
| import base64 | |
| import json | |
| import mimetypes | |
| import re | |
| import time | |
| from pathlib import Path | |
| from typing import Any, Dict, Optional | |
| from urllib.parse import urlparse | |
| from fastapi import APIRouter, Depends, File, Form, Header, HTTPException, Query, Response, UploadFile | |
| from app.config import get_settings | |
| from app.core.logger import get_logger | |
| from app.models.schemas import ( | |
| GCSACLRequest, | |
| GCSBaseCredentialsRequest, | |
| GCSBucketIAMRequest, | |
| GCSBucketRequest, | |
| GCSBucketRestoreRequest, | |
| GCSBucketUpdateRequest, | |
| GCSComposeRequest, | |
| GCSDownloadResponse, | |
| GCSGenericResponse, | |
| GCSListResponse, | |
| GCSObjectTransferRequest, | |
| GCSObjectUpdateRequest, | |
| GCSResumableUploadRequest, | |
| GCSURLRequest, | |
| GCSURLResponse, | |
| GCSWatchRequest, | |
| ) | |
| from app.services.gcs_service import GCSCredentials, GCSCredentialsError, GCSError, GCSService | |
| router = APIRouter(prefix="/google/gcs", tags=["Google Cloud Storage"]) | |
| _logger = get_logger(__name__) | |
| _settings = get_settings() | |
| _MAX_UPLOAD_BYTES = _settings.max_upload_bytes | |
| _gcs_service = GCSService() | |
| def get_gcs_service() -> GCSService: | |
| return _gcs_service | |
| async def close_gcs_service() -> None: | |
| await _gcs_service.close() | |
| # --------------------------------------------------------------------------- | |
| # Shared helpers | |
| # --------------------------------------------------------------------------- | |
| def _elapsed_ms(start: float) -> float: | |
| return round((time.perf_counter() - start) * 1000, 2) | |
| async def _resolve( | |
| service: GCSService, | |
| *, | |
| payload: Optional[Any] = None, | |
| url: Optional[str] = None, | |
| file_bytes: Optional[bytes] = None, | |
| file_name: Optional[str] = None, | |
| ) -> GCSCredentials: | |
| try: | |
| return await service.resolve_credentials( | |
| payload=payload, url=url, file_bytes=file_bytes, file_name=file_name | |
| ) | |
| except GCSCredentialsError as exc: | |
| raise HTTPException(status_code=400, detail=exc.message) | |
| def _creds_fields(body: Any) -> Dict[str, Any]: | |
| """Extract credential kwargs from a GCS request model.""" | |
| payload = body.credentials if isinstance(body.credentials, dict) else body.credentials_json | |
| return {"payload": payload, "url": body.credentials_url} | |
| def _error_response(cls, start: float, exc: GCSError): | |
| return cls(success=False, time_ms=_elapsed_ms(start), error=exc.message) | |
| async def _read_file_limited(f: UploadFile, *, what: str) -> bytes: | |
| data = await f.read() | |
| if len(data) > _MAX_UPLOAD_BYTES: | |
| raise HTTPException( | |
| status_code=413, | |
| detail=f"{what} exceeds the {_MAX_UPLOAD_BYTES} byte upload limit", | |
| ) | |
| return data | |
| # --------------------------------------------------------------------------- | |
| # Buckets — CRUD | |
| # --------------------------------------------------------------------------- | |
| async def create_bucket( | |
| body: GCSBucketRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| bucket_body: Dict[str, Any] = dict(body.bucket or {}) | |
| if body.location: | |
| bucket_body["location"] = body.location | |
| if body.storage_class: | |
| bucket_body["storageClass"] = body.storage_class | |
| try: | |
| result = await service.create_bucket(creds, body.name, project=body.project_id, bucket_body=bucket_body) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def get_bucket( | |
| bucket: str, | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials", description="Service account JSON (or JSON string)"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url", description="URL to a service account JSON file"), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| try: | |
| result = await service.get_bucket(creds, bucket) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def list_buckets( | |
| project_id: Optional[str] = Query(None, description="GCP project id (defaults to the service account project)"), | |
| prefix: Optional[str] = Query(None, description="Filter buckets whose names start with this prefix"), | |
| max_results: Optional[int] = Query(None, ge=1, le=1000, description="Maximum number of results to return"), | |
| page_token: Optional[str] = Query(None, description="Pagination token from a previous response"), | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials", description="Service account JSON (or JSON string)"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url", description="URL to a service account JSON file"), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| try: | |
| result = await service.list_buckets(creds, project=project_id, prefix=prefix, | |
| max_results=max_results, page_token=page_token) | |
| data = result.get("data") or {} | |
| items = data.get("items", []) | |
| return GCSListResponse( | |
| success=True, time_ms=_elapsed_ms(start), count=len(items), | |
| next_page_token=data.get("nextPageToken"), | |
| items=items, prefixes=[], | |
| ) | |
| except GCSError as exc: | |
| return _error_response(GCSListResponse, start, exc) | |
| async def patch_bucket( | |
| bucket: str, | |
| body: GCSBucketUpdateRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.patch_bucket(creds, bucket, body.bucket) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def update_bucket( | |
| bucket: str, | |
| body: GCSBucketUpdateRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.update_bucket(creds, bucket, body.bucket) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def delete_bucket( | |
| bucket: str, | |
| body: Optional[Dict[str, Any]] = None, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve( | |
| service, | |
| payload=(body or {}).get("credentials_json") or (body or {}).get("credentials"), | |
| url=(body or {}).get("credentials_url"), | |
| ) | |
| try: | |
| await service.delete_bucket(creds, bucket) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data={"deleted": bucket}) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| # --------------------------------------------------------------------------- | |
| # Buckets — IAM & permissions | |
| # --------------------------------------------------------------------------- | |
| async def restore_bucket( | |
| bucket: str, | |
| body: GCSBucketRestoreRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.restore_bucket(creds, bucket, if_metageneration_match=body.if_metageneration_match) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def get_bucket_iam( | |
| bucket: str, | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url"), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| try: | |
| result = await service.get_bucket_iam(creds, bucket) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def set_bucket_iam( | |
| bucket: str, | |
| body: GCSBucketIAMRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.set_bucket_iam(creds, bucket, body.policy) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def test_bucket_permissions( | |
| bucket: str, | |
| permissions: str = Query(..., description="Comma-separated permissions to test, e.g. 'storage.buckets.get,storage.objects.list'"), | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url"), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| permission_list = [p.strip() for p in permissions.split(",") if p.strip()] | |
| if not permission_list: | |
| raise HTTPException(status_code=422, detail="At least one permission is required.") | |
| try: | |
| result = await service.test_bucket_permissions(creds, bucket, permission_list) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| # --------------------------------------------------------------------------- | |
| # Buckets — default object ACLs | |
| # --------------------------------------------------------------------------- | |
| async def list_default_object_acl( | |
| bucket: str, | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url"), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| try: | |
| result = await service.list_default_object_acl(creds, bucket) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def insert_default_object_acl( | |
| bucket: str, | |
| body: GCSACLRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| acl = {"entity": body.entity, "role": body.role} | |
| try: | |
| result = await service.insert_default_object_acl(creds, bucket, acl) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def get_default_object_acl( | |
| bucket: str, | |
| entity: str, | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url"), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| try: | |
| result = await service.get_default_object_acl(creds, bucket, entity) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def patch_default_object_acl( | |
| bucket: str, | |
| entity: str, | |
| body: GCSACLRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| acl = {"entity": body.entity, "role": body.role} | |
| try: | |
| result = await service.patch_default_object_acl(creds, bucket, entity, acl) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def update_default_object_acl( | |
| bucket: str, | |
| entity: str, | |
| body: GCSACLRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| acl = {"entity": body.entity, "role": body.role} | |
| try: | |
| result = await service.update_default_object_acl(creds, bucket, entity, acl) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def delete_default_object_acl( | |
| bucket: str, | |
| entity: str, | |
| body: Optional[Dict[str, Any]] = None, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve( | |
| service, | |
| payload=(body or {}).get("credentials_json") or (body or {}).get("credentials"), | |
| url=(body or {}).get("credentials_url"), | |
| ) | |
| try: | |
| await service.delete_default_object_acl(creds, bucket, entity) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data={"deleted": entity}) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| # --------------------------------------------------------------------------- | |
| # Objects — list / upload / download | |
| # --------------------------------------------------------------------------- | |
| async def list_objects( | |
| bucket: str, | |
| prefix: Optional[str] = Query(None, description="Filter objects whose names start with this prefix"), | |
| delimiter: Optional[str] = Query(None, description="Fold object names at this delimiter (e.g. '/')"), | |
| max_results: Optional[int] = Query(None, ge=1, le=1000, description="Maximum number of results"), | |
| page_token: Optional[str] = Query(None, description="Pagination token from a previous response"), | |
| versions: Optional[bool] = Query(None, description="Include object versions"), | |
| match_glob: Optional[str] = Query(None, description="Glob pattern to filter object names"), | |
| start_offset: Optional[str] = Query(None, description="Lexicographically first object to return"), | |
| end_offset: Optional[str] = Query(None, description="Lexicographically last object to return"), | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url"), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| try: | |
| result = await service.list_objects( | |
| creds, bucket, prefix=prefix, delimiter=delimiter, max_results=max_results, | |
| page_token=page_token, versions=versions, match_glob=match_glob, | |
| start_offset=start_offset, end_offset=end_offset, | |
| ) | |
| data = result.get("data") or {} | |
| items = data.get("items", []) | |
| return GCSListResponse( | |
| success=True, time_ms=_elapsed_ms(start), count=len(items), | |
| next_page_token=data.get("nextPageToken"), | |
| items=items, prefixes=data.get("prefixes", []), | |
| ) | |
| except GCSError as exc: | |
| return _error_response(GCSListResponse, start, exc) | |
| async def upload_object( | |
| bucket: str, | |
| name: Optional[str] = Form(None, min_length=1, max_length=1024, description="Destination object name (defaults to the uploaded file name)"), | |
| file: Optional[UploadFile] = File(None, description="Object content to upload"), | |
| file_url: Optional[str] = Form(None, description="Publicly reachable URL whose content becomes the object"), | |
| content_type: Optional[str] = Form(None, description="Content-Type of the uploaded object (auto-detected when omitted)"), | |
| metadata_json: str = Form(None, description="Optional JSON object of object metadata to set on upload"), | |
| if_generation_match: Optional[int] = Form(None, description="Generation match condition"), | |
| credentials_file: Optional[UploadFile] = File(None, description="Service account JSON file (alternative to inline credentials)"), | |
| credentials_json: Optional[str] = Form(None, description="Service account JSON as a string"), | |
| credentials_url: Optional[str] = Form(None, description="URL to a service account JSON file"), | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| if file is None and not file_url: | |
| raise HTTPException(status_code=422, detail="Either 'file' or 'file_url' must be provided to upload an object.") | |
| if file is not None and file_url: | |
| raise HTTPException(status_code=422, detail="Provide either 'file' or 'file_url', not both.") | |
| if name is None: | |
| if file is not None and file.filename: | |
| name = Path(file.filename).name | |
| elif file_url: | |
| parsed_url = urlparse(file_url) | |
| name = Path(parsed_url.path).name or "download" | |
| if not name or not name.strip(): | |
| raise HTTPException(status_code=422, detail="Object name must be non-empty and must not start with '/'.") | |
| cleaned_name = name.strip() | |
| if cleaned_name.startswith("/"): | |
| raise HTTPException(status_code=422, detail="Object name must be non-empty and must not start with '/'.") | |
| if content_type is None: | |
| if file is not None and file.content_type: | |
| content_type = file.content_type | |
| else: | |
| content_type = mimetypes.guess_type(cleaned_name)[0] or "application/octet-stream" | |
| if not re.fullmatch(r"[\w./+-]{1,255}", content_type): | |
| raise HTTPException(status_code=422, detail="content_type must be a valid media type.") | |
| if credentials_file is not None: | |
| cred_bytes = await _read_file_limited(credentials_file, what="credentials file") | |
| creds = await _resolve(service, file_bytes=cred_bytes, file_name=credentials_file.filename) | |
| else: | |
| creds = await _resolve(service, payload=credentials_json, url=credentials_url) | |
| if file_url: | |
| try: | |
| content, fetched_type = await service.fetch_url_content(file_url) | |
| except GCSError as exc: | |
| raise HTTPException(status_code=exc.status_code, detail=exc.message) | |
| if len(content) > _MAX_UPLOAD_BYTES: | |
| raise HTTPException(status_code=413, detail=f"Content from file_url exceeds the {_MAX_UPLOAD_BYTES} byte upload limit") | |
| if fetched_type and (content_type == "application/octet-stream" or content_type == mimetypes.guess_type(cleaned_name)[0]): | |
| content_type = fetched_type | |
| else: | |
| content = await _read_file_limited(file, what="file") | |
| metadata: Dict[str, Any] = {} | |
| if metadata_json: | |
| try: | |
| metadata = json.loads(metadata_json) | |
| if not isinstance(metadata, dict): | |
| raise ValueError("metadata_json must be a JSON object") | |
| except (json.JSONDecodeError, ValueError) as exc: | |
| raise HTTPException(status_code=400, detail=f"Invalid metadata_json: {exc}") | |
| try: | |
| result = await service.upload_object( | |
| creds, bucket, cleaned_name, content, | |
| content_type=content_type, metadata=metadata, | |
| if_generation_match=if_generation_match, | |
| ) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def initiate_resumable_upload( | |
| bucket: str, | |
| body: GCSResumableUploadRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.initiate_resumable_upload( | |
| creds, bucket, body.name, | |
| content_type=body.content_type or "application/octet-stream", | |
| object_metadata=body.object_metadata, | |
| ) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data={"session_uri": result.get("session_uri")}) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def watch_all_objects( | |
| bucket: str, | |
| body: GCSWatchRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.watch_all_objects(creds, bucket, body.channel) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def download_object_raw( | |
| bucket: str, | |
| object: str, | |
| generation: Optional[int] = Query(None, description="Specific object generation to download"), | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url"), | |
| ): | |
| """Download an object's raw bytes (returns the file directly).""" | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| try: | |
| result = await service.download_object(creds, bucket, object, generation=generation) | |
| except GCSError as exc: | |
| raise HTTPException(status_code=exc.status_code, detail=exc.message) | |
| content = result.get("data") | |
| if not isinstance(content, bytes): | |
| content = content or b"" | |
| content_type = result.get("content_type") or "application/octet-stream" | |
| return Response( | |
| content=content, | |
| media_type=content_type, | |
| headers={"Content-Disposition": f'attachment; filename="{object.split("/")[-1]}"'}, | |
| ) | |
| async def download_object_base64( | |
| bucket: str, | |
| object: str, | |
| generation: Optional[int] = Query(None, description="Specific object generation to download"), | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url"), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| try: | |
| result = await service.download_object(creds, bucket, object, generation=generation) | |
| except GCSError as exc: | |
| return _error_response(GCSDownloadResponse, start, exc) | |
| content = result.get("data") | |
| content_bytes = content if isinstance(content, bytes) else (content or b"") | |
| return GCSDownloadResponse( | |
| success=True, time_ms=_elapsed_ms(start), | |
| content_base64=base64.b64encode(content_bytes).decode("ascii"), | |
| content_type=result.get("content_type"), | |
| size_bytes=len(content_bytes), | |
| ) | |
| async def public_url( | |
| bucket: str, | |
| object: str, | |
| ): | |
| start = time.perf_counter() | |
| url = GCSService.public_url(bucket, object) | |
| return GCSURLResponse(success=True, time_ms=_elapsed_ms(start), url=url, method="GET") | |
| # --------------------------------------------------------------------------- | |
| # Objects — IAM | |
| # --------------------------------------------------------------------------- | |
| async def get_object_iam( | |
| bucket: str, | |
| object: str, | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url"), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| try: | |
| result = await service.get_object_iam(creds, bucket, object) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def set_object_iam( | |
| bucket: str, | |
| object: str, | |
| body: GCSBucketIAMRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.set_object_iam(creds, bucket, object, body.policy) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| # --------------------------------------------------------------------------- | |
| # Objects — ACLs | |
| # --------------------------------------------------------------------------- | |
| async def list_object_acl( | |
| bucket: str, | |
| object: str, | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url"), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| try: | |
| result = await service.list_object_acl(creds, bucket, object) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def insert_object_acl( | |
| bucket: str, | |
| object: str, | |
| body: GCSACLRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| acl = {"entity": body.entity, "role": body.role} | |
| try: | |
| result = await service.insert_object_acl(creds, bucket, object, acl) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def get_object_acl( | |
| bucket: str, | |
| object: str, | |
| entity: str, | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url"), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| try: | |
| result = await service.get_object_acl(creds, bucket, object, entity) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def patch_object_acl( | |
| bucket: str, | |
| object: str, | |
| entity: str, | |
| body: GCSACLRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| acl = {"entity": body.entity, "role": body.role} | |
| try: | |
| result = await service.patch_object_acl(creds, bucket, object, entity, acl) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def update_object_acl( | |
| bucket: str, | |
| object: str, | |
| entity: str, | |
| body: GCSACLRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| acl = {"entity": body.entity, "role": body.role} | |
| try: | |
| result = await service.update_object_acl(creds, bucket, object, entity, acl) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def delete_object_acl( | |
| bucket: str, | |
| object: str, | |
| entity: str, | |
| body: Optional[Dict[str, Any]] = None, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| payload = (body or {}).get("credentials_json") or (body or {}).get("credentials") | |
| url = (body or {}).get("credentials_url") | |
| creds = await _resolve(service, payload=payload, url=url) | |
| try: | |
| await service.delete_object_acl(creds, bucket, object, entity) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data={"deleted": entity}) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def get_object( | |
| bucket: str, | |
| object: str, | |
| generation: Optional[int] = Query(None, description="Specific object generation"), | |
| service: GCSService = Depends(get_gcs_service), | |
| x_credentials: Optional[str] = Header(None, alias="X-Goog-Credentials"), | |
| x_credentials_url: Optional[str] = Header(None, alias="X-Goog-Credentials-Url"), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, payload=x_credentials, url=x_credentials_url) | |
| try: | |
| result = await service.get_object(creds, bucket, object, generation=generation) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def patch_object( | |
| bucket: str, | |
| object: str, | |
| body: GCSObjectUpdateRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.patch_object(creds, bucket, object, body.object, | |
| if_generation_match=body.if_generation_match) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def update_object( | |
| bucket: str, | |
| object: str, | |
| body: GCSObjectUpdateRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.update_object(creds, bucket, object, body.object, | |
| if_generation_match=body.if_generation_match) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def delete_object( | |
| bucket: str, | |
| object: str, | |
| body: Optional[Dict[str, Any]] = None, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| payload = (body or {}).get("credentials_json") or (body or {}).get("credentials") | |
| url = (body or {}).get("credentials_url") | |
| creds = await _resolve(service, payload=payload, url=url) | |
| generation = (body or {}).get("generation") | |
| if_generation_match = (body or {}).get("if_generation_match") | |
| try: | |
| await service.delete_object(creds, bucket, object, | |
| generation=generation, if_generation_match=if_generation_match) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data={"deleted": object}) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| # --------------------------------------------------------------------------- | |
| # Objects — copy / move / rewrite / compose | |
| # --------------------------------------------------------------------------- | |
| async def copy_object( | |
| bucket: str, | |
| object: str, | |
| body: GCSObjectTransferRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.copy_object( | |
| creds, bucket, object, body.destination_bucket, body.destination_name, | |
| source_generation=body.source_generation, body=body.object, | |
| ) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def move_object( | |
| bucket: str, | |
| object: str, | |
| body: GCSObjectTransferRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.move_object( | |
| creds, bucket, object, body.destination_bucket, body.destination_name, | |
| source_generation=body.source_generation, if_generation_match=body.if_generation_match, | |
| body=body.object, | |
| ) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def rewrite_object( | |
| bucket: str, | |
| object: str, | |
| body: GCSObjectTransferRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.rewrite_object( | |
| creds, bucket, object, body.destination_bucket, body.destination_name, | |
| rewrite_token=body.rewrite_token, | |
| max_bytes_rewritten_per_call=body.max_bytes_rewritten_per_call, | |
| body=body.object, | |
| ) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def set_object_hold( | |
| bucket: str, | |
| object: str, | |
| body: GCSBaseCredentialsRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.set_object_hold(creds, bucket, object) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def release_object_hold( | |
| bucket: str, | |
| object: str, | |
| body: GCSBaseCredentialsRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.release_object_hold(creds, bucket, object) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def lock_object_retention( | |
| bucket: str, | |
| object: str, | |
| body: GCSBucketRestoreRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.lock_object_retention( | |
| creds, bucket, object, if_metageneration_match=body.if_metageneration_match | |
| ) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| async def compose_object( | |
| bucket: str, | |
| object: str, | |
| body: GCSComposeRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| result = await service.compose_object(creds, bucket, object, body.source_objects, body=body.destination) | |
| return GCSGenericResponse(success=True, time_ms=_elapsed_ms(start), data=result.get("data")) | |
| except GCSError as exc: | |
| return _error_response(GCSGenericResponse, start, exc) | |
| # --------------------------------------------------------------------------- | |
| # Public & signed URLs | |
| # --------------------------------------------------------------------------- | |
| async def signed_download_url( | |
| bucket: str, | |
| object: str, | |
| body: GCSURLRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| url, expires = await service.download_signed_url( | |
| creds, bucket, object, | |
| expires_in_seconds=body.expires_in_seconds, | |
| response_content_type=body.response_content_type, | |
| response_disposition=body.response_disposition, | |
| ) | |
| return GCSURLResponse( | |
| success=True, time_ms=_elapsed_ms(start), url=url, | |
| method="GET", expires_in_seconds=expires, | |
| ) | |
| except GCSError as exc: | |
| return _error_response(GCSURLResponse, start, exc) | |
| async def signed_upload_url( | |
| bucket: str, | |
| object: str, | |
| body: GCSURLRequest, | |
| service: GCSService = Depends(get_gcs_service), | |
| ): | |
| start = time.perf_counter() | |
| creds = await _resolve(service, **_creds_fields(body)) | |
| try: | |
| url, expires = await service.upload_signed_url( | |
| creds, bucket, object, | |
| expires_in_seconds=body.expires_in_seconds, | |
| content_type=body.content_type, | |
| ) | |
| return GCSURLResponse( | |
| success=True, time_ms=_elapsed_ms(start), url=url, | |
| method="PUT", expires_in_seconds=expires, | |
| ) | |
| except GCSError as exc: | |
| return _error_response(GCSURLResponse, start, exc) | |