llm-ready-data / app /api /v1 /gcs.py
validops-east-1's picture
yes
efa9f90
Raw
History Blame Contribute Delete
44.7 kB
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
# ---------------------------------------------------------------------------
@router.post("/buckets", response_model=GCSGenericResponse,
summary="Create a Cloud Storage bucket")
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)
@router.get("/buckets/{bucket}", response_model=GCSGenericResponse,
summary="Get bucket metadata")
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)
@router.get("/buckets", response_model=GCSListResponse,
summary="List buckets in a project")
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)
@router.patch("/buckets/{bucket}", response_model=GCSGenericResponse,
summary="Patch bucket metadata (partial update)")
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)
@router.put("/buckets/{bucket}", response_model=GCSGenericResponse,
summary="Replace bucket metadata (full update)")
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)
@router.delete("/buckets/{bucket}", response_model=GCSGenericResponse,
summary="Delete a bucket")
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
# ---------------------------------------------------------------------------
@router.post("/buckets/{bucket}/restore", response_model=GCSGenericResponse,
summary="Restore a soft-deleted bucket")
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)
@router.get("/buckets/{bucket}/iam", response_model=GCSGenericResponse,
summary="Get bucket IAM policy")
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)
@router.put("/buckets/{bucket}/iam", response_model=GCSGenericResponse,
summary="Set bucket IAM policy")
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)
@router.get("/buckets/{bucket}/iam/testPermissions", response_model=GCSGenericResponse,
summary="Test permissions on a bucket")
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
# ---------------------------------------------------------------------------
@router.get("/buckets/{bucket}/defaultAcl", response_model=GCSGenericResponse,
summary="List 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)
@router.post("/buckets/{bucket}/defaultAcl", response_model=GCSGenericResponse,
summary="Add a default object ACL entry")
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)
@router.get("/buckets/{bucket}/defaultAcl/{entity}", response_model=GCSGenericResponse,
summary="Get a default object ACL entry")
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)
@router.patch("/buckets/{bucket}/defaultAcl/{entity}", response_model=GCSGenericResponse,
summary="Patch a default object ACL entry")
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)
@router.put("/buckets/{bucket}/defaultAcl/{entity}", response_model=GCSGenericResponse,
summary="Update a default object ACL entry")
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)
@router.delete("/buckets/{bucket}/defaultAcl/{entity}", response_model=GCSGenericResponse,
summary="Delete a default object ACL entry")
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
# ---------------------------------------------------------------------------
@router.get("/buckets/{bucket}/objects", response_model=GCSListResponse,
summary="List objects in a bucket")
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)
@router.post("/buckets/{bucket}/objects/upload", response_model=GCSGenericResponse,
summary="Upload an object from a file or a URL (up to the server upload limit)")
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)
@router.post("/buckets/{bucket}/objects/upload/resumable", response_model=GCSGenericResponse,
summary="Start a resumable upload session and return its upload URI")
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)
@router.post("/buckets/{bucket}/objects/watch", response_model=GCSGenericResponse,
summary="Register an object-change notification channel for a bucket")
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)
@router.get("/buckets/{bucket}/objects/{object:path}/download", include_in_schema=False)
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]}"'},
)
@router.get("/buckets/{bucket}/objects/{object:path}/download/base64", response_model=GCSDownloadResponse,
summary="Download an object as base64 (JSON response)")
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),
)
@router.get("/buckets/{bucket}/objects/{object:path}/public-url", response_model=GCSURLResponse,
summary="Generate the public URL for an object")
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
# ---------------------------------------------------------------------------
@router.get("/buckets/{bucket}/objects/{object:path}/iam", response_model=GCSGenericResponse,
summary="Get object IAM policy")
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)
@router.put("/buckets/{bucket}/objects/{object:path}/iam", response_model=GCSGenericResponse,
summary="Set object IAM policy")
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
# ---------------------------------------------------------------------------
@router.get("/buckets/{bucket}/objects/{object:path}/acl", response_model=GCSGenericResponse,
summary="List object 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)
@router.post("/buckets/{bucket}/objects/{object:path}/acl", response_model=GCSGenericResponse,
summary="Add an object ACL entry")
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)
@router.get("/buckets/{bucket}/objects/{object:path}/acl/{entity}", response_model=GCSGenericResponse,
summary="Get an object ACL entry")
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)
@router.patch("/buckets/{bucket}/objects/{object:path}/acl/{entity}", response_model=GCSGenericResponse,
summary="Patch an object ACL entry")
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)
@router.put("/buckets/{bucket}/objects/{object:path}/acl/{entity}", response_model=GCSGenericResponse,
summary="Update an object ACL entry")
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)
@router.delete("/buckets/{bucket}/objects/{object:path}/acl/{entity}", response_model=GCSGenericResponse,
summary="Delete an object ACL entry")
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)
@router.get("/buckets/{bucket}/objects/{object:path}", response_model=GCSGenericResponse,
summary="Get object metadata")
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)
@router.patch("/buckets/{bucket}/objects/{object:path}", response_model=GCSGenericResponse,
summary="Patch object metadata (partial update)")
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)
@router.put("/buckets/{bucket}/objects/{object:path}", response_model=GCSGenericResponse,
summary="Replace object metadata (full update)")
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)
@router.delete("/buckets/{bucket}/objects/{object:path}", response_model=GCSGenericResponse,
summary="Delete an object")
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
# ---------------------------------------------------------------------------
@router.post("/buckets/{bucket}/objects/{object:path}/copy", response_model=GCSGenericResponse,
summary="Copy an object to a destination")
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)
@router.post("/buckets/{bucket}/objects/{object:path}/move", response_model=GCSGenericResponse,
summary="Move an object to a destination")
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)
@router.post("/buckets/{bucket}/objects/{object:path}/rewrite", response_model=GCSGenericResponse,
summary="Rewrite (async copy) an object to a destination")
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)
@router.post("/buckets/{bucket}/objects/{object:path}/hold", response_model=GCSGenericResponse,
summary="Set an event-based hold on an object")
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)
@router.post("/buckets/{bucket}/objects/{object:path}/releaseHold", response_model=GCSGenericResponse,
summary="Release an event-based hold on an object")
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)
@router.post("/buckets/{bucket}/objects/{object:path}/lockRetentionPolicy", response_model=GCSGenericResponse,
summary="Lock a retention policy on an object (permanent)")
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)
@router.post("/buckets/{bucket}/objects/{object:path}/compose", response_model=GCSGenericResponse,
summary="Compose multiple source objects into this destination object")
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
# ---------------------------------------------------------------------------
@router.post("/buckets/{bucket}/objects/{object:path}/signed-url", response_model=GCSURLResponse,
summary="Generate a V4 signed download URL for an object")
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)
@router.post("/buckets/{bucket}/objects/{object:path}/signed-upload-url", response_model=GCSURLResponse,
summary="Generate a V4 signed upload URL for an object (PUT)")
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)