diff --git a/app/core/config.py b/app/core/config.py index c9270f03..9d01ffa0 100644 --- a/app/core/config.py +++ b/app/core/config.py @@ -35,6 +35,10 @@ class Settings(BaseSettings): PHOTO_APPROVAL_TIMEOUT_DAYS: int = 7 EVENT_LIFECYCLE_POLL_INTERVAL_SECONDS: int = 60 + DIRECT_UPLOAD_PRESIGN_EXPIRES_SECONDS: int = 1800 + DIRECT_UPLOAD_STALE_PENDING_MINUTES: int = 45 + DIRECT_UPLOAD_RECONCILE_POLL_INTERVAL_SECONDS: int = 300 + DIRECT_UPLOAD_MAX_BATCH_SIZE: int = 200 # Mobile auth/session defaults MOBILE_SESSION_LIMIT: int = 3 @@ -79,9 +83,20 @@ class Settings(BaseSettings): GOOGLE_CLIENT_ID: str = "" GOOGLE_CLIENT_SECRET: str = "" GOOGLE_REDIRECT_URI: str = "" + # drive.readonly alone can't write; drive.file alone can only see files + # the app itself created, which would break browsing/importing existing + # Drive folders. Both scopes together preserve the existing read/import + # flow and add write access for syncing approved direct uploads back to + # Drive. Existing staff connections keep their old readonly-only grant + # until they disconnect and reconnect through the consent screen. GOOGLE_OAUTH_SCOPES: str = ( - "https://www.googleapis.com/auth/drive.readonly openid email profile" + "https://www.googleapis.com/auth/drive.readonly " + "https://www.googleapis.com/auth/drive.file openid email profile" ) + # Folder ID (from the Drive URL) that approved direct-upload photos get + # synced into. Empty means uploads land in the connected account's Drive + # root instead of a specific folder. + GOOGLE_CLUB_DRIVE_FOLDER_ID: str = "" FACE_ENCRYPTION_KEY: str FIREBASE_CREDENTIALS_PATH: str diff --git a/app/infra/google_drive.py b/app/infra/google_drive.py index f25c4ef6..05633bf7 100644 --- a/app/infra/google_drive.py +++ b/app/infra/google_drive.py @@ -189,6 +189,67 @@ async def get_file_metadata( size_bytes=size_bytes, ) + @staticmethod + async def upload_file( + *, + access_token: str, + file_name: str, + content_type: str, + data: bytes, + folder_id: str | None, + ) -> GoogleDriveFileMetadata: + boundary = "multai-drive-upload-boundary" + metadata: dict[str, object] = {"name": file_name} + if folder_id: + metadata["parents"] = [folder_id] + + body = ( + f"--{boundary}\r\n" + "Content-Type: application/json; charset=UTF-8\r\n\r\n" + f"{json.dumps(metadata)}\r\n" + f"--{boundary}\r\n" + f"Content-Type: {content_type}\r\n\r\n" + ).encode("utf-8") + data + f"\r\n--{boundary}--".encode("utf-8") + + def _request() -> dict[str, object]: + url = ( + "https://www.googleapis.com/upload/drive/v3/files" + "?uploadType=multipart&supportsAllDrives=true&fields=id,name,mimeType,size" + ) + request = urllib.request.Request( + url, + data=body, + headers={ + "Authorization": f"Bearer {access_token}", + "Content-Type": f"multipart/related; boundary={boundary}", + }, + method="POST", + ) + try: + with urllib.request.urlopen(request, timeout=60) as response: + return json.loads(response.read().decode("utf-8")) + except urllib.error.HTTPError as exc: + details = exc.read().decode("utf-8", errors="ignore") + raise AppException.bad_request( + f"Google Drive file upload failed: {details or exc.reason}" + ) from exc + except urllib.error.URLError as exc: + raise AppException.internal_error("Unable to reach Google APIs") from exc + + result = await asyncio.to_thread(_request) + size_raw = result.get("size", "0") + try: + size_bytes = int(size_raw) if isinstance(size_raw, (str, int)) else len(data) + except (TypeError, ValueError): + size_bytes = len(data) + + return GoogleDriveFileMetadata( + id=GoogleDriveClient._require_str(result, "id"), + name=GoogleDriveClient._require_str(result, "name"), + mime_type=GoogleDriveClient._require_str(result, "mimeType"), + size_bytes=size_bytes, + ) + @staticmethod async def download_file( *, diff --git a/app/infra/minio.py b/app/infra/minio.py index e6249dae..9a666c5c 100644 --- a/app/infra/minio.py +++ b/app/infra/minio.py @@ -2,6 +2,8 @@ import random import string import uuid +from dataclasses import dataclass +from datetime import timedelta from fastapi import UploadFile from miniopy_async.commonconfig import CopySource from miniopy_async.error import S3Error @@ -36,6 +38,12 @@ async def init_minio_client( if not await Bucket.client.bucket_exists(bucket_name): await Bucket.client.make_bucket(bucket_name) +@dataclass(frozen=True) +class ObjectStat: + size: int + content_type: str + + class Bucket: bucket_name: str file_prefix: str @@ -130,6 +138,25 @@ async def copy(self, *, source_object_name: str, target_object_name: str) -> str ) return target_object_name + async def presigned_put_url(self, object_name: str, *, expires_seconds: int) -> str: + return await self.client.presigned_put_object( + bucket_name=self.bucket_name, + object_name=self._object_path(object_name), + expires=timedelta(seconds=expires_seconds), + ) + + async def stat(self, object_name: str) -> ObjectStat | None: + try: + result = await self.client.stat_object( + bucket_name=self.bucket_name, + object_name=self._object_path(object_name), + ) + except S3Error as e: + if e.code == "NoSuchKey": + return None + raise + return ObjectStat(size=result.size or 0, content_type=result.content_type or DEFAULT_CONTENT_TYPE) + image_ext_content_type_map = { "apng": ["image/apng"], "avif": ["image/avif"], diff --git a/app/infra/nats.py b/app/infra/nats.py index e17bd500..8b07e831 100644 --- a/app/infra/nats.py +++ b/app/infra/nats.py @@ -35,6 +35,7 @@ class NatsSubjects(Enum): STAFF_UPLOAD_REQUEST_APPROVED = "staff.upload_request.approved" STAFF_UPLOAD_REQUEST_REJECTED = "staff.upload_request.rejected" PHOTO_PROCESS = "photo.process" + PHOTO_DRIVE_SYNC_REQUESTED = "photo.drive_sync.requested" class NatsClient: diff --git a/app/router/staff/__init__.py b/app/router/staff/__init__.py index 35cbe0ec..11023d7b 100644 --- a/app/router/staff/__init__.py +++ b/app/router/staff/__init__.py @@ -3,8 +3,10 @@ from app.router.staff.drive import router as staff_drive_router from app.router.staff.notifications import router as staff_notifications_router from app.router.staff.uploads import router as staff_uploads_router +from app.router.staff.uploads_direct import router as staff_uploads_direct_router router = APIRouter(prefix="/staff", tags=["staff"]) router.include_router(staff_drive_router) router.include_router(staff_notifications_router) router.include_router(staff_uploads_router) +router.include_router(staff_uploads_direct_router) diff --git a/app/router/staff/uploads_direct.py b/app/router/staff/uploads_direct.py new file mode 100644 index 00000000..267cd9bd --- /dev/null +++ b/app/router/staff/uploads_direct.py @@ -0,0 +1,119 @@ +from uuid import UUID + +from fastapi import APIRouter, Depends + +from app.container import Container, get_container +from app.deps.cookie_auth import get_current_staff_user +from app.schema.internal.uploads import DirectFileInput +from app.schema.request.staff.uploads_direct import ( + CreateDirectGroupRequest, + RegisterDirectBatchRequest, +) +from app.schema.response.staff.upload_groups import UploadRequestGroupSchema +from app.schema.response.staff.uploads_direct import ( + DirectUploadFileResponse, + RegisterDirectBatchResponse, + ResumeDirectGroupResponse, +) +from db.generated.models import StaffUser + +router = APIRouter(prefix="/uploads/direct") +# this endpoint are for staff to upload images directly to the system and very large files and they can resume and restart and retry . + +@router.post("/groups", response_model=UploadRequestGroupSchema) +async def create_direct_group( + req: CreateDirectGroupRequest, + current_staff_user: StaffUser = Depends(get_current_staff_user), + container: Container = Depends(get_container), +) -> UploadRequestGroupSchema: + group = await container.upload_requests_service.create_direct_group( + event_id=req.event_id, requested_by=current_staff_user, + ) + details = await container.upload_requests_service.get_group_details( + group_id=group.id, current_staff_user=current_staff_user, + ) + return UploadRequestGroupSchema.from_details(details) + + +@router.post("/groups/{group_id}/batches", response_model=RegisterDirectBatchResponse) +async def register_direct_batch( + group_id: UUID, + req: RegisterDirectBatchRequest, + current_staff_user: StaffUser = Depends(get_current_staff_user), + container: Container = Depends(get_container), +) -> RegisterDirectBatchResponse: + results = await container.upload_requests_service.register_direct_batch( + group_id=group_id, + files=[ + DirectFileInput( + file_name=f.file_name, + mime_type=f.mime_type, + size_bytes=f.size_bytes, + taken_at=f.taken_at, + day_number=f.day_number, + visibility=f.visibility, + ) + for f in req.files + ], + requested_by=current_staff_user, + ) + return RegisterDirectBatchResponse( + group_id=group_id, + items=[ + DirectUploadFileResponse(photo_id=photo.id, file_name=photo.file_name, upload_url=url) + for photo, url in results + ], + ) + + +@router.post("/photos/{photo_id}/confirm", response_model=DirectUploadFileResponse) +async def confirm_direct_upload( + photo_id: UUID, + current_staff_user: StaffUser = Depends(get_current_staff_user), + container: Container = Depends(get_container), +) -> DirectUploadFileResponse: + photo = await container.upload_requests_service.confirm_direct_upload( + photo_id=photo_id, requested_by=current_staff_user, + ) + return DirectUploadFileResponse(photo_id=photo.id, file_name=photo.file_name, upload_url="") + + +@router.post("/photos/{photo_id}/fail", response_model=DirectUploadFileResponse) +async def fail_direct_upload( + photo_id: UUID, + current_staff_user: StaffUser = Depends(get_current_staff_user), + container: Container = Depends(get_container), +) -> DirectUploadFileResponse: + photo = await container.upload_requests_service.fail_direct_upload( + photo_id=photo_id, requested_by=current_staff_user, + ) + return DirectUploadFileResponse(photo_id=photo.id, file_name=photo.file_name, upload_url="") + + +@router.post("/groups/{group_id}/resume", response_model=ResumeDirectGroupResponse) +async def resume_direct_group( + group_id: UUID, + current_staff_user: StaffUser = Depends(get_current_staff_user), + container: Container = Depends(get_container), +) -> ResumeDirectGroupResponse: + results = await container.upload_requests_service.resume_direct_group( + group_id=group_id, requested_by=current_staff_user, + ) + return ResumeDirectGroupResponse( + items=[ + DirectUploadFileResponse(photo_id=photo.id, file_name=photo.file_name, upload_url=url) + for photo, url in results + ] + ) + + +@router.get("/groups/{group_id}", response_model=UploadRequestGroupSchema) +async def get_direct_group_status( + group_id: UUID, + current_staff_user: StaffUser = Depends(get_current_staff_user), + container: Container = Depends(get_container), +) -> UploadRequestGroupSchema: + details = await container.upload_requests_service.get_group_details( + group_id=group_id, current_staff_user=current_staff_user, + ) + return UploadRequestGroupSchema.from_details(details) diff --git a/app/schema/internal/uploads.py b/app/schema/internal/uploads.py index c8b91da5..58f82491 100644 --- a/app/schema/internal/uploads.py +++ b/app/schema/internal/uploads.py @@ -8,3 +8,13 @@ class UploadPhotoInput: taken_at: datetime | None day_number: int | None visibility: str + + +@dataclass(frozen=True) +class DirectFileInput: + file_name: str + mime_type: str + size_bytes: int + taken_at: datetime | None + day_number: int | None + visibility: str diff --git a/app/schema/request/staff/uploads_direct.py b/app/schema/request/staff/uploads_direct.py new file mode 100644 index 00000000..068c1b69 --- /dev/null +++ b/app/schema/request/staff/uploads_direct.py @@ -0,0 +1,22 @@ +from datetime import datetime +from typing import Optional +from uuid import UUID + +from pydantic import BaseModel + + +class DirectFileInputRequest(BaseModel): + file_name: str + mime_type: str + size_bytes: int + taken_at: Optional[datetime] = None + day_number: Optional[int] = None + visibility: str = "private" + + +class CreateDirectGroupRequest(BaseModel): + event_id: UUID + + +class RegisterDirectBatchRequest(BaseModel): + files: list[DirectFileInputRequest] diff --git a/app/schema/response/staff/upload_groups.py b/app/schema/response/staff/upload_groups.py index 7a264c2d..f691cc9a 100644 --- a/app/schema/response/staff/upload_groups.py +++ b/app/schema/response/staff/upload_groups.py @@ -18,10 +18,11 @@ class UploadRequestGroupSchema(BaseModel): id: UUID event_id: UUID - folder_id: str + folder_id: str | None requested_by: UUID approved_by: UUID | None status: str + source: str processing_status: str total_photo_count: int batch_count: int @@ -72,8 +73,9 @@ class UploadRequestGroupSummarySchema(BaseModel): model_config = ConfigDict(from_attributes=True) id: UUID event_id: UUID - folder_id: str + folder_id: str | None status: str + source: str processing_status: str total_photo_count: int batch_count: int diff --git a/app/schema/response/staff/uploads.py b/app/schema/response/staff/uploads.py index 863414cd..60a13544 100644 --- a/app/schema/response/staff/uploads.py +++ b/app/schema/response/staff/uploads.py @@ -11,7 +11,7 @@ class UploadRequestPhotoSchema(BaseModel): model_config = ConfigDict(from_attributes=True) id: UUID - drive_file_id: str + drive_file_id: str | None file_name: str mime_type: str size_bytes: int @@ -19,6 +19,8 @@ class UploadRequestPhotoSchema(BaseModel): day_number: int | None visibility: str status: str + source: str + transfer_status: str created_at: datetime @@ -32,6 +34,7 @@ class UploadRequestSchema(BaseModel): requested_by: UUID approved_by: UUID | None status: str + source: str photo_count: int created_at: datetime approved_at: datetime | None diff --git a/app/schema/response/staff/uploads_direct.py b/app/schema/response/staff/uploads_direct.py new file mode 100644 index 00000000..fba3d4c4 --- /dev/null +++ b/app/schema/response/staff/uploads_direct.py @@ -0,0 +1,18 @@ +from uuid import UUID + +from pydantic import BaseModel + + +class DirectUploadFileResponse(BaseModel): + photo_id: UUID + file_name: str + upload_url: str + + +class RegisterDirectBatchResponse(BaseModel): + group_id: UUID + items: list[DirectUploadFileResponse] + + +class ResumeDirectGroupResponse(BaseModel): + items: list[DirectUploadFileResponse] diff --git a/app/service/staff_drive.py b/app/service/staff_drive.py index f3b2ba00..b80a98be 100644 --- a/app/service/staff_drive.py +++ b/app/service/staff_drive.py @@ -209,6 +209,25 @@ async def get_system_access_token(self) -> str: connection = await self._refresh_connection_access_token(connection) return self.decrypt(connection.access_token) + async def upload_to_system_drive( + self, + *, + file_name: str, + content_type: str, + data: bytes, + ) -> str: + """Upload bytes to the system/club Drive using the most recently + connected active staff Drive connection. Returns the Drive file id.""" + access_token = await self.get_system_access_token() + metadata = await GoogleDriveClient.upload_file( + access_token=access_token, + file_name=file_name, + content_type=content_type, + data=data, + folder_id=settings.GOOGLE_CLUB_DRIVE_FOLDER_ID or None, + ) + return metadata.id + async def disconnect(self, staff_user_id: uuid.UUID) -> None: connection = await self.get_status(staff_user_id) if connection is None: diff --git a/app/service/staged_upload_storage.py b/app/service/staged_upload_storage.py index 813fa2bf..e8f9d97a 100644 --- a/app/service/staged_upload_storage.py +++ b/app/service/staged_upload_storage.py @@ -5,7 +5,7 @@ import uuid from app.core.exceptions import AppException -from app.infra.minio import Bucket, IMAGES_BUCKET_NAME +from app.infra.minio import Bucket, IMAGES_BUCKET_NAME, ObjectStat @dataclass(frozen=True) @@ -100,3 +100,22 @@ async def delete_storage_key(self, storage_key: str) -> None: async def get_preview(self, storage_key: str) -> PreviewObject: data, file_name, content_type = await self.bucket.get(storage_key) return PreviewObject(data=data, file_name=file_name, content_type=content_type) + + async def create_presigned_staging_upload( + self, + *, + upload_request_id: uuid.UUID, + photo_id: uuid.UUID, + file_name: str, + expires_seconds: int, + ) -> tuple[str, str]: + storage_key = self.build_staging_key( + upload_request_id=upload_request_id, + photo_id=photo_id, + file_name=file_name, + ) + url = await self.bucket.presigned_put_url(storage_key, expires_seconds=expires_seconds) + return storage_key, url + + async def stat_staging_object(self, storage_key: str) -> ObjectStat | None: + return await self.bucket.stat(storage_key) diff --git a/app/service/upload_requests.py b/app/service/upload_requests.py index 0ab404d6..5af3cdd4 100644 --- a/app/service/upload_requests.py +++ b/app/service/upload_requests.py @@ -8,6 +8,7 @@ from sqlalchemy.exc import IntegrityError +from app.core.config import settings from app.core.constant import AuditEventType from app.core.exceptions import AppException from app.core.logger import logger @@ -17,7 +18,7 @@ GoogleDriveFileMetadata, ) from app.infra.nats import NatsClient, NatsSubjects -from app.schema.internal.uploads import UploadPhotoInput +from app.schema.internal.uploads import DirectFileInput, UploadPhotoInput from app.service.audit import AuditService from app.service.staged_upload_storage import PreviewObject, StagedUploadStorageService from app.service.staff_drive import StaffDriveService @@ -282,6 +283,7 @@ async def _create_request_with_access_token( drive_file_id=None, requested_by=requested_by.id, photo_count=len(photos), + source="drive", ) ) except IntegrityError as exc: @@ -336,6 +338,16 @@ async def _approve_request_without_side_effects( if not staged_photos: raise AppException.bad_request("No staged photos found for this upload request") + not_transferred = [ + p for p in staged_photos + if getattr(p, "transfer_status", "uploaded") != "uploaded" + ] + if not_transferred: + raise AppException.bad_request( + f"{len(not_transferred)} photo(s) have not finished uploading. " + "Resume or remove them before approving." + ) + finalized_storage_keys: list[str] = [] created_photos: list[Photo] = [] try: @@ -492,6 +504,34 @@ async def _publish_photo_process_events(self, photos: list[Photo]) -> None: if photos: logger.info("Published %d photo process events", len(photos)) + async def _publish_drive_sync_events( + self, + staged_photos: list[UploadRequestPhoto], + created_photos: list[Photo], + ) -> None: + # staged_photos and created_photos are built 1:1 in the same order by + # _approve_request_without_side_effects (and accumulated in lockstep + # across multiple calls for a group approval), so zipping them here + # is safe. Only direct-uploaded photos get synced — Drive-imported + # photos already live in Drive and syncing them back would be a + # pointless round trip. + count = 0 + for staged_photo, created_photo in zip(staged_photos, created_photos): + if getattr(staged_photo, "source", "drive") != "direct": + continue + await self._publish_event( + subject=NatsSubjects.PHOTO_DRIVE_SYNC_REQUESTED, + payload={ + "photo_id": str(created_photo.id), + "storage_key": created_photo.storage_key, + "file_name": staged_photo.file_name, + "mime_type": staged_photo.mime_type, + }, + ) + count += 1 + if count: + logger.info("Published %d Drive sync events", count) + async def _mark_group_import_failed( self, *, @@ -572,6 +612,8 @@ async def create_group_from_folder( requested_by=requested_by.id, total_photo_count=0, batch_count=0, + source="drive", + processing_status="pending", ) ) except IntegrityError as exc: @@ -593,6 +635,177 @@ async def create_group_from_folder( ) return UploadRequestGroupDetails(group=upload_group, requests=[]) + async def create_direct_group( + self, + *, + event_id: uuid.UUID, + requested_by: StaffUser, + ) -> UploadRequestGroup: + try: + upload_group = await self.upload_request_group_querier.create_upload_request_group( + upload_request_group_queries.CreateUploadRequestGroupParams( + event_id=event_id, + folder_id=None, + requested_by=requested_by.id, + total_photo_count=0, + batch_count=0, + source="direct", + processing_status="completed", + ) + ) + except IntegrityError as exc: + self._raise_integrity_error(exc) + if upload_group is None: + raise AppException.internal_error("Failed to create upload group") + return upload_group + + async def register_direct_batch( + self, + *, + group_id: uuid.UUID, + files: Sequence[DirectFileInput], + requested_by: StaffUser, + ) -> list[tuple[UploadRequestPhoto, str]]: + if not files: + raise AppException.bad_request("At least one file is required") + if len(files) > settings.DIRECT_UPLOAD_MAX_BATCH_SIZE: + raise AppException.bad_request( + f"A batch can contain at most {settings.DIRECT_UPLOAD_MAX_BATCH_SIZE} files" + ) + for file in files: + if file.mime_type not in self._allowed_mime_types: + raise AppException.image_format_error(f"Unsupported image format: {file.mime_type}") + if file.size_bytes <= 0 or file.size_bytes > self._max_photo_size_bytes: + raise AppException.bad_request(f"{file.file_name} exceeds maximum allowed size") + + group = await self.upload_request_group_querier.get_upload_request_group_by_id(id=group_id) + if group is None: + raise AppException.not_found("Upload group not found") + self._ensure_group_access(current_staff_user=requested_by, upload_group=group) + + upload_request = await self.upload_request_querier.create_upload_request( + upload_request_queries.CreateUploadRequestParams( + event_id=group.event_id, + group_id=group_id, + drive_file_id=None, + requested_by=requested_by.id, + photo_count=len(files), + source="direct", + ) + ) + if upload_request is None: + raise AppException.internal_error("Failed to create upload request") + + results: list[tuple[UploadRequestPhoto, str]] = [] + for file in files: + photo_id = uuid.uuid4() + storage_key, presigned_url = await self.staged_upload_storage.create_presigned_staging_upload( + upload_request_id=upload_request.id, + photo_id=photo_id, + file_name=file.file_name, + expires_seconds=settings.DIRECT_UPLOAD_PRESIGN_EXPIRES_SECONDS, + ) + created_photo = await self.upload_request_photo_querier.create_direct_upload_request_photo( + upload_request_photo_queries.CreateDirectUploadRequestPhotoParams( + upload_request_id=upload_request.id, + file_name=file.file_name, + mime_type=file.mime_type, + size_bytes=file.size_bytes, + staging_storage_key=storage_key, + taken_at=file.taken_at, + day_number=file.day_number, + visibility=file.visibility, + ) + ) + if created_photo is None: + raise AppException.internal_error("Failed to register upload photo") + results.append((created_photo, presigned_url)) + + await self.upload_request_group_querier.increment_upload_request_group_counts( + id=group_id, total_photo_count=len(files), + ) + + return results + + async def confirm_direct_upload( + self, + *, + photo_id: uuid.UUID, + requested_by: StaffUser, + ) -> UploadRequestPhoto: + photo = await self.upload_request_photo_querier.get_upload_request_photo_by_id(id=photo_id) + if photo is None: + raise AppException.not_found("Upload photo not found") + + stat = await self.staged_upload_storage.stat_staging_object(photo.staging_storage_key) + if stat is None: + failed = await self.upload_request_photo_querier.fail_upload_request_photo_transfer(id=photo_id) + if failed is None: + raise AppException.internal_error("Failed to mark upload as failed") + raise AppException.bad_request( + "Upload did not complete — file not found in storage. Retry the upload." + ) + + confirmed = await self.upload_request_photo_querier.confirm_upload_request_photo_transfer( + id=photo_id, + size_bytes=stat.size, + mime_type=stat.content_type, + ) + if confirmed is None: + raise AppException.internal_error("Failed to confirm upload") + return confirmed + + async def fail_direct_upload( + self, + *, + photo_id: uuid.UUID, + requested_by: StaffUser, + ) -> UploadRequestPhoto: + photo = await self.upload_request_photo_querier.get_upload_request_photo_by_id(id=photo_id) + if photo is None: + raise AppException.not_found("Upload photo not found") + + failed = await self.upload_request_photo_querier.fail_upload_request_photo_transfer(id=photo_id) + if failed is None: + raise AppException.internal_error("Failed to mark upload as failed") + return failed + + async def resume_direct_group( + self, + *, + group_id: uuid.UUID, + requested_by: StaffUser, + ) -> list[tuple[UploadRequestPhoto, str]]: + group = await self.upload_request_group_querier.get_upload_request_group_by_id(id=group_id) + if group is None: + raise AppException.not_found("Upload group not found") + self._ensure_group_access(current_staff_user=requested_by, upload_group=group) + + request_ids: list[uuid.UUID] = [] + async for req in self.upload_request_querier.list_upload_requests_by_group_id(group_id=group_id): + request_ids.append(req.id) + + results: list[tuple[UploadRequestPhoto, str]] = [] + async for photo in self.upload_request_photo_querier.list_upload_request_photos_by_upload_request_ids( + dollar_1=request_ids + ): + if getattr(photo, "transfer_status", "uploaded") not in ("pending_upload", "failed"): + continue + storage_key, presigned_url = await self.staged_upload_storage.create_presigned_staging_upload( + upload_request_id=photo.upload_request_id, + photo_id=photo.id, + file_name=photo.file_name, + expires_seconds=settings.DIRECT_UPLOAD_PRESIGN_EXPIRES_SECONDS, + ) + reset_photo = await self.upload_request_photo_querier.reset_upload_request_photo_transfer_to_pending( + id=photo.id + ) + if reset_photo is None: + continue + results.append((reset_photo, presigned_url)) + + return results + async def process_group_import( self, *, @@ -968,6 +1181,7 @@ async def approve_request( }, ) await self._publish_photo_process_events(created_photos) + await self._publish_drive_sync_events(staged_photos, created_photos) await self._audit( AuditEventType.UPLOAD_REQUEST_APPROVED, request_id=upload_request.id, @@ -1100,6 +1314,7 @@ async def approve_group( }, ) await self._publish_photo_process_events(all_created_photos) + await self._publish_drive_sync_events(all_staged_photos, all_created_photos) await self._audit( AuditEventType.UPLOAD_REQUEST_APPROVED, group_id=upload_group.id, diff --git a/app/worker/drive_sync/__init__.py b/app/worker/drive_sync/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/app/worker/drive_sync/main.py b/app/worker/drive_sync/main.py new file mode 100644 index 00000000..f642fa1b --- /dev/null +++ b/app/worker/drive_sync/main.py @@ -0,0 +1,103 @@ +import asyncio +import json +import uuid + +from pydantic import BaseModel, ValidationError + +from app.core.config import settings +from app.core.logger import logger +from app.infra.database import engine +from app.infra.minio import Bucket, IMAGES_BUCKET_NAME, init_minio_client +from app.infra.nats import NatsClient, NatsSubjects +from app.infra.redis import RedisClient +from app.service.staff_drive import StaffDriveService +from db.generated import photos as photo_queries +from db.generated import staff_drive_connections as drive_queries +from db.generated import staff_user as staff_queries + + +class PhotoDriveSyncEvent(BaseModel): + photo_id: uuid.UUID + storage_key: str + file_name: str + mime_type: str + + +def _parse_payload(raw_data: bytes) -> PhotoDriveSyncEvent | None: + try: + parsed = json.loads(raw_data.decode("utf-8")) + except (UnicodeDecodeError, json.JSONDecodeError) as exc: + logger.error("drive_sync: cannot parse payload: %s", exc) + return None + if not isinstance(parsed, dict): + return None + try: + return PhotoDriveSyncEvent.model_validate(parsed) + except ValidationError as exc: + logger.warning("drive_sync: payload validation failed: %s", exc) + return None + + +async def _handle_event(raw_data: bytes) -> None: + event = _parse_payload(raw_data) + if event is None: + return + + bucket = Bucket(IMAGES_BUCKET_NAME, "") + try: + data, _, content_type = await bucket.get(event.storage_key) + except Exception as exc: + logger.warning("drive_sync: failed to read photo %s from storage: %s", event.photo_id, exc) + return + + async with engine.begin() as conn: + staff_drive_service = StaffDriveService( + staff_user_querier=staff_queries.AsyncQuerier(conn), + drive_connection_querier=drive_queries.AsyncQuerier(conn), + redis=RedisClient.get_instance(), + ) + photo_querier = photo_queries.AsyncQuerier(conn) + + try: + drive_file_id = await staff_drive_service.upload_to_system_drive( + file_name=event.file_name, + content_type=event.mime_type or content_type, + data=data, + ) + except Exception as exc: + logger.warning("drive_sync: upload failed for photo %s: %s", event.photo_id, exc) + return + + synced = await photo_querier.mark_photo_drive_synced( + id=event.photo_id, drive_file_id=drive_file_id, + ) + if synced is None: + logger.warning("drive_sync: photo %s not found when recording sync", event.photo_id) + return + + logger.info("drive_sync: synced photo %s to Drive as %s", event.photo_id, drive_file_id) + + +async def main() -> None: + logger.info("Drive sync worker starting") + await init_minio_client( + minio_host=settings.MINIO_HOST, + minio_port=settings.MINIO_API_PORT, + minio_root_user=settings.MINIO_ROOT_USER, + minio_root_password=settings.MINIO_ROOT_PASSWORD, + ) + RedisClient.init( + host=settings.REDIS_HOST, + port=settings.REDIS_PORT, + password=settings.REDIS_PASSWORD, + ) + await NatsClient.connect() + try: + await NatsClient.subscribe(NatsSubjects.PHOTO_DRIVE_SYNC_REQUESTED, _handle_event) + await asyncio.Event().wait() + finally: + await NatsClient.close() + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/app/worker/upload_reconciler/__init__.py b/app/worker/upload_reconciler/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/app/worker/upload_reconciler/main.py b/app/worker/upload_reconciler/main.py new file mode 100644 index 00000000..915d7911 --- /dev/null +++ b/app/worker/upload_reconciler/main.py @@ -0,0 +1,60 @@ +import asyncio + +from app.core.config import settings +from app.core.logger import logger +from app.infra.database import engine +from app.service.staged_upload_storage import StagedUploadStorageService +from db.generated import upload_request_photos as upload_request_photo_queries + +storage_service = StagedUploadStorageService() + + +async def run_reconcile_pass() -> None: + async with engine.begin() as conn: + querier = upload_request_photo_queries.AsyncQuerier(conn) + + stale_photos = [ + photo + async for photo in querier.list_stale_pending_transfer_photos( + dollar_1=str(settings.DIRECT_UPLOAD_STALE_PENDING_MINUTES) + ) + ] + + if not stale_photos: + return + + confirmed = 0 + failed = 0 + for photo in stale_photos: + stat = await storage_service.stat_staging_object(photo.staging_storage_key) + if stat is not None: + await querier.confirm_upload_request_photo_transfer( + id=photo.id, size_bytes=stat.size, mime_type=stat.content_type, + ) + confirmed += 1 + else: + await querier.fail_upload_request_photo_transfer(id=photo.id) + failed += 1 + + logger.info( + "upload_reconciler: reconciled %d stale photo(s) — %d confirmed, %d failed", + len(stale_photos), confirmed, failed, + ) + + +async def main() -> None: + logger.info( + "Upload reconciler worker starting, poll_interval=%ds, stale_after=%dmin", + settings.DIRECT_UPLOAD_RECONCILE_POLL_INTERVAL_SECONDS, + settings.DIRECT_UPLOAD_STALE_PENDING_MINUTES, + ) + while True: + try: + await run_reconcile_pass() + except Exception: + logger.exception("upload_reconciler: pass failed") + await asyncio.sleep(settings.DIRECT_UPLOAD_RECONCILE_POLL_INTERVAL_SECONDS) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/db/generated/models.py b/db/generated/models.py index 86617af5..7d935995 100644 --- a/db/generated/models.py +++ b/db/generated/models.py @@ -115,6 +115,8 @@ class Photo: visibility: str status: Any created_at: datetime.datetime + drive_file_id: Optional[str] + drive_synced_at: Optional[datetime.datetime] @dataclasses.dataclass() @@ -208,13 +210,14 @@ class UploadRequest: photo_count: int rejection_reason: Optional[str] group_id: Optional[uuid.UUID] + source: str @dataclasses.dataclass() class UploadRequestGroup: id: uuid.UUID event_id: uuid.UUID - folder_id: str + folder_id: Optional[str] requested_by: uuid.UUID approved_by: Optional[uuid.UUID] status: Any @@ -227,13 +230,14 @@ class UploadRequestGroup: processed_photo_count: int failed_photo_count: int error_message: Optional[str] + source: str @dataclasses.dataclass() class UploadRequestPhoto: id: uuid.UUID upload_request_id: uuid.UUID - drive_file_id: str + drive_file_id: Optional[str] file_name: str mime_type: str size_bytes: int @@ -244,6 +248,8 @@ class UploadRequestPhoto: visibility: str status: str created_at: datetime.datetime + source: str + transfer_status: str @dataclasses.dataclass() diff --git a/db/generated/photos.py b/db/generated/photos.py index 599a46d1..a978143c 100644 --- a/db/generated/photos.py +++ b/db/generated/photos.py @@ -35,7 +35,7 @@ ) VALUES ( :p1, :p2, :p3, :p4, :p5 ) -RETURNING id, event_id, uploaded_by, storage_key, taken_at, day_number, visibility, status, created_at +RETURNING id, event_id, uploaded_by, storage_key, taken_at, day_number, visibility, status, created_at, drive_file_id, drive_synced_at """ @@ -57,12 +57,12 @@ class CreatePhotoParams: GET_PHOTO_BY_ID = """-- name: get_photo_by_id \\:one -SELECT id, event_id, uploaded_by, storage_key, taken_at, day_number, visibility, status, created_at FROM photos WHERE id = :p1 +SELECT id, event_id, uploaded_by, storage_key, taken_at, day_number, visibility, status, created_at, drive_file_id, drive_synced_at FROM photos WHERE id = :p1 """ LIST_EVENT_PHOTOS_FOR_USER = """-- name: list_event_photos_for_user \\:many -SELECT p.id, p.event_id, p.uploaded_by, p.storage_key, p.taken_at, p.day_number, p.visibility, p.status, p.created_at, +SELECT p.id, p.event_id, p.uploaded_by, p.storage_key, p.taken_at, p.day_number, p.visibility, p.status, p.created_at, p.drive_file_id, p.drive_synced_at, (SELECT COUNT(*) FROM photo_faces pf2 WHERE pf2.photo_id = p.id)\\:\\:int AS face_count FROM photos p WHERE p.event_id = :p2 @@ -106,11 +106,13 @@ class ListEventPhotosForUserRow: visibility: str status: Any created_at: datetime.datetime + drive_file_id: Optional[str] + drive_synced_at: Optional[datetime.datetime] face_count: int LIST_USER_PHOTOS = """-- name: list_user_photos \\:many -SELECT p.id, p.event_id, p.uploaded_by, p.storage_key, p.taken_at, p.day_number, p.visibility, p.status, p.created_at, +SELECT p.id, p.event_id, p.uploaded_by, p.storage_key, p.taken_at, p.day_number, p.visibility, p.status, p.created_at, p.drive_file_id, p.drive_synced_at, (SELECT COUNT(*) FROM photo_faces pf2 WHERE pf2.photo_id = p.id)\\:\\:int AS face_count FROM photos p WHERE ( @@ -152,14 +154,25 @@ class ListUserPhotosRow: visibility: str status: Any created_at: datetime.datetime + drive_file_id: Optional[str] + drive_synced_at: Optional[datetime.datetime] face_count: int +MARK_PHOTO_DRIVE_SYNCED = """-- name: mark_photo_drive_synced \\:one +UPDATE photos +SET drive_file_id = :p2, + drive_synced_at = NOW() +WHERE id = :p1 +RETURNING id, event_id, uploaded_by, storage_key, taken_at, day_number, visibility, status, created_at, drive_file_id, drive_synced_at +""" + + UPDATE_PHOTO_STATUS = """-- name: update_photo_status \\:one UPDATE photos SET status = :p2 WHERE id = :p1 -RETURNING id, event_id, uploaded_by, storage_key, taken_at, day_number, visibility, status, created_at +RETURNING id, event_id, uploaded_by, storage_key, taken_at, day_number, visibility, status, created_at, drive_file_id, drive_synced_at """ @@ -167,7 +180,7 @@ class ListUserPhotosRow: UPDATE photos SET visibility = :p2 WHERE id = :p1 -RETURNING id, event_id, uploaded_by, storage_key, taken_at, day_number, visibility, status, created_at +RETURNING id, event_id, uploaded_by, storage_key, taken_at, day_number, visibility, status, created_at, drive_file_id, drive_synced_at """ @@ -201,9 +214,11 @@ async def create_photo(self, arg: CreatePhotoParams) -> Optional[models.Photo]: visibility=row[6], status=row[7], created_at=row[8], + drive_file_id=row[9], + drive_synced_at=row[10], ) - async def get_drive_file_id_for_photo(self, *, final_storage_key: Optional[str]) -> Optional[str]: + async def get_drive_file_id_for_photo(self, *, final_storage_key: Optional[str]) -> Optional[Optional[str]]: row = (await self._conn.execute(sqlalchemy.text(GET_DRIVE_FILE_ID_FOR_PHOTO), {"p1": final_storage_key})).first() if row is None: return None @@ -223,6 +238,8 @@ async def get_photo_by_id(self, *, id: uuid.UUID) -> Optional[models.Photo]: visibility=row[6], status=row[7], created_at=row[8], + drive_file_id=row[9], + drive_synced_at=row[10], ) async def list_event_photos_for_user(self, arg: ListEventPhotosForUserParams) -> AsyncIterator[ListEventPhotosForUserRow]: @@ -244,7 +261,9 @@ async def list_event_photos_for_user(self, arg: ListEventPhotosForUserParams) -> visibility=row[6], status=row[7], created_at=row[8], - face_count=row[9], + drive_file_id=row[9], + drive_synced_at=row[10], + face_count=row[11], ) async def list_user_photos(self, arg: ListUserPhotosParams) -> AsyncIterator[ListUserPhotosRow]: @@ -266,9 +285,29 @@ async def list_user_photos(self, arg: ListUserPhotosParams) -> AsyncIterator[Lis visibility=row[6], status=row[7], created_at=row[8], - face_count=row[9], + drive_file_id=row[9], + drive_synced_at=row[10], + face_count=row[11], ) + async def mark_photo_drive_synced(self, *, id: uuid.UUID, drive_file_id: Optional[str]) -> Optional[models.Photo]: + row = (await self._conn.execute(sqlalchemy.text(MARK_PHOTO_DRIVE_SYNCED), {"p1": id, "p2": drive_file_id})).first() + if row is None: + return None + return models.Photo( + id=row[0], + event_id=row[1], + uploaded_by=row[2], + storage_key=row[3], + taken_at=row[4], + day_number=row[5], + visibility=row[6], + status=row[7], + created_at=row[8], + drive_file_id=row[9], + drive_synced_at=row[10], + ) + async def update_photo_status(self, *, id: uuid.UUID, status: Any) -> Optional[models.Photo]: row = (await self._conn.execute(sqlalchemy.text(UPDATE_PHOTO_STATUS), {"p1": id, "p2": status})).first() if row is None: @@ -283,6 +322,8 @@ async def update_photo_status(self, *, id: uuid.UUID, status: Any) -> Optional[m visibility=row[6], status=row[7], created_at=row[8], + drive_file_id=row[9], + drive_synced_at=row[10], ) async def update_photo_visibility(self, *, id: uuid.UUID, visibility: str) -> Optional[models.Photo]: @@ -299,4 +340,6 @@ async def update_photo_visibility(self, *, id: uuid.UUID, visibility: str) -> Op visibility=row[6], status=row[7], created_at=row[8], + drive_file_id=row[9], + drive_synced_at=row[10], ) diff --git a/db/generated/upload_request_groups.py b/db/generated/upload_request_groups.py index 039b1f05..ceb91cf5 100644 --- a/db/generated/upload_request_groups.py +++ b/db/generated/upload_request_groups.py @@ -20,7 +20,7 @@ rejection_reason = NULL WHERE id = :p1 AND status = 'pending' -RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message +RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source """ @@ -33,7 +33,7 @@ failed_photo_count = :p5, error_message = NULL WHERE id = :p1 -RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message +RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source """ @@ -52,21 +52,25 @@ class CompleteUploadRequestGroupProcessingParams: folder_id, requested_by, total_photo_count, - batch_count + batch_count, + source, + processing_status ) VALUES ( - :p1, :p2, :p3, :p4, :p5 + :p1, :p2, :p3, :p4, :p5, :p6, :p7 ) -RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message +RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source """ @dataclasses.dataclass() class CreateUploadRequestGroupParams: event_id: uuid.UUID - folder_id: str + folder_id: Optional[str] requested_by: uuid.UUID total_photo_count: int batch_count: int + source: str + processing_status: str DELETE_UPLOAD_REQUEST_GROUP = """-- name: delete_upload_request_group \\:exec @@ -84,7 +88,7 @@ class CreateUploadRequestGroupParams: failed_photo_count = :p5, error_message = :p6 WHERE id = :p1 -RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message +RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source """ @@ -99,21 +103,30 @@ class FailUploadRequestGroupProcessingParams: GET_UPLOAD_REQUEST_GROUP_BY_ID = """-- name: get_upload_request_group_by_id \\:one -SELECT id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message +SELECT id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source FROM upload_request_groups WHERE id = :p1 """ +INCREMENT_UPLOAD_REQUEST_GROUP_COUNTS = """-- name: increment_upload_request_group_counts \\:one +UPDATE upload_request_groups +SET total_photo_count = total_photo_count + :p2, + batch_count = batch_count + 1 +WHERE id = :p1 +RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source +""" + + LIST_UPLOAD_REQUEST_GROUPS = """-- name: list_upload_request_groups \\:many -SELECT id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message +SELECT id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source FROM upload_request_groups ORDER BY created_at DESC """ LIST_UPLOAD_REQUEST_GROUPS_BY_REQUESTER = """-- name: list_upload_request_groups_by_requester \\:many -SELECT id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message +SELECT id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source FROM upload_request_groups WHERE requested_by = :p1 ORDER BY created_at DESC @@ -121,7 +134,7 @@ class FailUploadRequestGroupProcessingParams: LIST_UPLOAD_REQUEST_GROUPS_BY_REQUESTER_AND_STATUS = """-- name: list_upload_request_groups_by_requester_and_status \\:many -SELECT id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message +SELECT id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source FROM upload_request_groups WHERE requested_by = :p1 AND status = :p2 @@ -130,7 +143,7 @@ class FailUploadRequestGroupProcessingParams: LIST_UPLOAD_REQUEST_GROUPS_BY_STATUS = """-- name: list_upload_request_groups_by_status \\:many -SELECT id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message +SELECT id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source FROM upload_request_groups WHERE status = :p1 ORDER BY created_at DESC @@ -145,7 +158,7 @@ class FailUploadRequestGroupProcessingParams: rejection_reason = :p3 WHERE id = :p1 AND status = 'pending' -RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message +RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source """ @@ -155,7 +168,7 @@ class FailUploadRequestGroupProcessingParams: error_message = NULL WHERE id = :p1 AND processing_status = 'pending' -RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message +RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source """ @@ -166,7 +179,7 @@ class FailUploadRequestGroupProcessingParams: processed_photo_count = :p4, failed_photo_count = :p5 WHERE id = :p1 -RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message +RETURNING id, event_id, folder_id, requested_by, approved_by, status, total_photo_count, batch_count, created_at, approved_at, rejection_reason, processing_status, processed_photo_count, failed_photo_count, error_message, source """ @@ -203,6 +216,7 @@ async def approve_upload_request_group(self, *, id: uuid.UUID, approved_by: Opti processed_photo_count=row[12], failed_photo_count=row[13], error_message=row[14], + source=row[15], ) async def complete_upload_request_group_processing(self, arg: CompleteUploadRequestGroupProcessingParams) -> Optional[models.UploadRequestGroup]: @@ -231,6 +245,7 @@ async def complete_upload_request_group_processing(self, arg: CompleteUploadRequ processed_photo_count=row[12], failed_photo_count=row[13], error_message=row[14], + source=row[15], ) async def create_upload_request_group(self, arg: CreateUploadRequestGroupParams) -> Optional[models.UploadRequestGroup]: @@ -240,6 +255,8 @@ async def create_upload_request_group(self, arg: CreateUploadRequestGroupParams) "p3": arg.requested_by, "p4": arg.total_photo_count, "p5": arg.batch_count, + "p6": arg.source, + "p7": arg.processing_status, })).first() if row is None: return None @@ -259,6 +276,7 @@ async def create_upload_request_group(self, arg: CreateUploadRequestGroupParams) processed_photo_count=row[12], failed_photo_count=row[13], error_message=row[14], + source=row[15], ) async def delete_upload_request_group(self, *, id: uuid.UUID) -> None: @@ -291,6 +309,7 @@ async def fail_upload_request_group_processing(self, arg: FailUploadRequestGroup processed_photo_count=row[12], failed_photo_count=row[13], error_message=row[14], + source=row[15], ) async def get_upload_request_group_by_id(self, *, id: uuid.UUID) -> Optional[models.UploadRequestGroup]: @@ -313,6 +332,30 @@ async def get_upload_request_group_by_id(self, *, id: uuid.UUID) -> Optional[mod processed_photo_count=row[12], failed_photo_count=row[13], error_message=row[14], + source=row[15], + ) + + async def increment_upload_request_group_counts(self, *, id: uuid.UUID, total_photo_count: int) -> Optional[models.UploadRequestGroup]: + row = (await self._conn.execute(sqlalchemy.text(INCREMENT_UPLOAD_REQUEST_GROUP_COUNTS), {"p1": id, "p2": total_photo_count})).first() + if row is None: + return None + return models.UploadRequestGroup( + id=row[0], + event_id=row[1], + folder_id=row[2], + requested_by=row[3], + approved_by=row[4], + status=row[5], + total_photo_count=row[6], + batch_count=row[7], + created_at=row[8], + approved_at=row[9], + rejection_reason=row[10], + processing_status=row[11], + processed_photo_count=row[12], + failed_photo_count=row[13], + error_message=row[14], + source=row[15], ) async def list_upload_request_groups(self) -> AsyncIterator[models.UploadRequestGroup]: @@ -334,6 +377,7 @@ async def list_upload_request_groups(self) -> AsyncIterator[models.UploadRequest processed_photo_count=row[12], failed_photo_count=row[13], error_message=row[14], + source=row[15], ) async def list_upload_request_groups_by_requester(self, *, requested_by: uuid.UUID) -> AsyncIterator[models.UploadRequestGroup]: @@ -355,6 +399,7 @@ async def list_upload_request_groups_by_requester(self, *, requested_by: uuid.UU processed_photo_count=row[12], failed_photo_count=row[13], error_message=row[14], + source=row[15], ) async def list_upload_request_groups_by_requester_and_status(self, *, requested_by: uuid.UUID, status: Any) -> AsyncIterator[models.UploadRequestGroup]: @@ -376,6 +421,7 @@ async def list_upload_request_groups_by_requester_and_status(self, *, requested_ processed_photo_count=row[12], failed_photo_count=row[13], error_message=row[14], + source=row[15], ) async def list_upload_request_groups_by_status(self, *, status: Any) -> AsyncIterator[models.UploadRequestGroup]: @@ -397,6 +443,7 @@ async def list_upload_request_groups_by_status(self, *, status: Any) -> AsyncIte processed_photo_count=row[12], failed_photo_count=row[13], error_message=row[14], + source=row[15], ) async def reject_upload_request_group(self, *, id: uuid.UUID, approved_by: Optional[uuid.UUID], rejection_reason: Optional[str]) -> Optional[models.UploadRequestGroup]: @@ -419,6 +466,7 @@ async def reject_upload_request_group(self, *, id: uuid.UUID, approved_by: Optio processed_photo_count=row[12], failed_photo_count=row[13], error_message=row[14], + source=row[15], ) async def start_upload_request_group_processing(self, *, id: uuid.UUID) -> Optional[models.UploadRequestGroup]: @@ -441,6 +489,7 @@ async def start_upload_request_group_processing(self, *, id: uuid.UUID) -> Optio processed_photo_count=row[12], failed_photo_count=row[13], error_message=row[14], + source=row[15], ) async def update_upload_request_group_import_progress(self, arg: UpdateUploadRequestGroupImportProgressParams) -> Optional[models.UploadRequestGroup]: @@ -469,4 +518,5 @@ async def update_upload_request_group_import_progress(self, arg: UpdateUploadReq processed_photo_count=row[12], failed_photo_count=row[13], error_message=row[14], + source=row[15], ) diff --git a/db/generated/upload_request_photos.py b/db/generated/upload_request_photos.py index 1cd3ebb4..ba33dee1 100644 --- a/db/generated/upload_request_photos.py +++ b/db/generated/upload_request_photos.py @@ -13,6 +13,50 @@ from db.generated import models +CONFIRM_UPLOAD_REQUEST_PHOTO_TRANSFER = """-- name: confirm_upload_request_photo_transfer \\:one +UPDATE upload_request_photos +SET transfer_status = 'uploaded', + size_bytes = :p2, + mime_type = :p3 +WHERE id = :p1 + AND transfer_status IN ('pending_upload', 'failed') +RETURNING id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at, source, transfer_status +""" + + +CREATE_DIRECT_UPLOAD_REQUEST_PHOTO = """-- name: create_direct_upload_request_photo \\:one +INSERT INTO upload_request_photos ( + upload_request_id, + drive_file_id, + file_name, + mime_type, + size_bytes, + staging_storage_key, + taken_at, + day_number, + visibility, + status, + source, + transfer_status +) VALUES ( + :p1, NULL, :p2, :p3, :p4, :p5, :p6, :p7, :p8, 'staged', 'direct', 'pending_upload' +) +RETURNING id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at, source, transfer_status +""" + + +@dataclasses.dataclass() +class CreateDirectUploadRequestPhotoParams: + upload_request_id: uuid.UUID + file_name: str + mime_type: str + size_bytes: int + staging_storage_key: str + taken_at: Optional[datetime.datetime] + day_number: Optional[int] + visibility: str + + CREATE_UPLOAD_REQUEST_PHOTO = """-- name: create_upload_request_photo \\:one INSERT INTO upload_request_photos ( upload_request_id, @@ -28,14 +72,14 @@ ) VALUES ( :p1, :p2, :p3, :p4, :p5, :p6, :p7, :p8, :p9, :p10 ) -RETURNING id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at +RETURNING id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at, source, transfer_status """ @dataclasses.dataclass() class CreateUploadRequestPhotoParams: upload_request_id: uuid.UUID - drive_file_id: str + drive_file_id: Optional[str] file_name: str mime_type: str size_bytes: int @@ -52,15 +96,35 @@ class CreateUploadRequestPhotoParams: """ +FAIL_UPLOAD_REQUEST_PHOTO_TRANSFER = """-- name: fail_upload_request_photo_transfer \\:one +UPDATE upload_request_photos +SET transfer_status = 'failed' +WHERE id = :p1 + AND transfer_status IN ('pending_upload', 'failed') +RETURNING id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at, source, transfer_status +""" + + GET_UPLOAD_REQUEST_PHOTO_BY_ID = """-- name: get_upload_request_photo_by_id \\:one -SELECT id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at +SELECT id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at, source, transfer_status FROM upload_request_photos WHERE id = :p1 """ +LIST_STALE_PENDING_TRANSFER_PHOTOS = """-- name: list_stale_pending_transfer_photos \\:many +SELECT id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at, source, transfer_status +FROM upload_request_photos +WHERE source = 'direct' + AND transfer_status = 'pending_upload' + AND created_at <= NOW() - (:p1 || ' minutes')\\:\\:interval +ORDER BY created_at ASC +LIMIT 500 +""" + + LIST_UPLOAD_REQUEST_PHOTOS_BY_UPLOAD_REQUEST_ID = """-- name: list_upload_request_photos_by_upload_request_id \\:many -SELECT id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at +SELECT id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at, source, transfer_status FROM upload_request_photos WHERE upload_request_id = :p1 ORDER BY created_at ASC @@ -68,19 +132,28 @@ class CreateUploadRequestPhotoParams: LIST_UPLOAD_REQUEST_PHOTOS_BY_UPLOAD_REQUEST_IDS = """-- name: list_upload_request_photos_by_upload_request_ids \\:many -SELECT id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at +SELECT id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at, source, transfer_status FROM upload_request_photos WHERE upload_request_id = ANY(:p1\\:\\:uuid[]) ORDER BY created_at ASC """ +RESET_UPLOAD_REQUEST_PHOTO_TRANSFER_TO_PENDING = """-- name: reset_upload_request_photo_transfer_to_pending \\:one +UPDATE upload_request_photos +SET transfer_status = 'pending_upload' +WHERE id = :p1 + AND transfer_status IN ('pending_upload', 'failed') +RETURNING id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at, source, transfer_status +""" + + UPDATE_UPLOAD_REQUEST_PHOTO_APPROVAL = """-- name: update_upload_request_photo_approval \\:one UPDATE upload_request_photos SET status = :p2, final_storage_key = :p3 WHERE id = :p1 -RETURNING id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at +RETURNING id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at, source, transfer_status """ @@ -88,7 +161,7 @@ class CreateUploadRequestPhotoParams: UPDATE upload_request_photos SET status = :p2 WHERE upload_request_id = :p1 -RETURNING id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at +RETURNING id, upload_request_id, drive_file_id, file_name, mime_type, size_bytes, staging_storage_key, final_storage_key, taken_at, day_number, visibility, status, created_at, source, transfer_status """ @@ -96,6 +169,59 @@ class AsyncQuerier: def __init__(self, conn: sqlalchemy.ext.asyncio.AsyncConnection): self._conn = conn + async def confirm_upload_request_photo_transfer(self, *, id: uuid.UUID, size_bytes: int, mime_type: str) -> Optional[models.UploadRequestPhoto]: + row = (await self._conn.execute(sqlalchemy.text(CONFIRM_UPLOAD_REQUEST_PHOTO_TRANSFER), {"p1": id, "p2": size_bytes, "p3": mime_type})).first() + if row is None: + return None + return models.UploadRequestPhoto( + id=row[0], + upload_request_id=row[1], + drive_file_id=row[2], + file_name=row[3], + mime_type=row[4], + size_bytes=row[5], + staging_storage_key=row[6], + final_storage_key=row[7], + taken_at=row[8], + day_number=row[9], + visibility=row[10], + status=row[11], + created_at=row[12], + source=row[13], + transfer_status=row[14], + ) + + async def create_direct_upload_request_photo(self, arg: CreateDirectUploadRequestPhotoParams) -> Optional[models.UploadRequestPhoto]: + row = (await self._conn.execute(sqlalchemy.text(CREATE_DIRECT_UPLOAD_REQUEST_PHOTO), { + "p1": arg.upload_request_id, + "p2": arg.file_name, + "p3": arg.mime_type, + "p4": arg.size_bytes, + "p5": arg.staging_storage_key, + "p6": arg.taken_at, + "p7": arg.day_number, + "p8": arg.visibility, + })).first() + if row is None: + return None + return models.UploadRequestPhoto( + id=row[0], + upload_request_id=row[1], + drive_file_id=row[2], + file_name=row[3], + mime_type=row[4], + size_bytes=row[5], + staging_storage_key=row[6], + final_storage_key=row[7], + taken_at=row[8], + day_number=row[9], + visibility=row[10], + status=row[11], + created_at=row[12], + source=row[13], + transfer_status=row[14], + ) + async def create_upload_request_photo(self, arg: CreateUploadRequestPhotoParams) -> Optional[models.UploadRequestPhoto]: row = (await self._conn.execute(sqlalchemy.text(CREATE_UPLOAD_REQUEST_PHOTO), { "p1": arg.upload_request_id, @@ -125,11 +251,35 @@ async def create_upload_request_photo(self, arg: CreateUploadRequestPhotoParams) visibility=row[10], status=row[11], created_at=row[12], + source=row[13], + transfer_status=row[14], ) async def delete_upload_request_photos_by_upload_request_id(self, *, upload_request_id: uuid.UUID) -> None: await self._conn.execute(sqlalchemy.text(DELETE_UPLOAD_REQUEST_PHOTOS_BY_UPLOAD_REQUEST_ID), {"p1": upload_request_id}) + async def fail_upload_request_photo_transfer(self, *, id: uuid.UUID) -> Optional[models.UploadRequestPhoto]: + row = (await self._conn.execute(sqlalchemy.text(FAIL_UPLOAD_REQUEST_PHOTO_TRANSFER), {"p1": id})).first() + if row is None: + return None + return models.UploadRequestPhoto( + id=row[0], + upload_request_id=row[1], + drive_file_id=row[2], + file_name=row[3], + mime_type=row[4], + size_bytes=row[5], + staging_storage_key=row[6], + final_storage_key=row[7], + taken_at=row[8], + day_number=row[9], + visibility=row[10], + status=row[11], + created_at=row[12], + source=row[13], + transfer_status=row[14], + ) + async def get_upload_request_photo_by_id(self, *, id: uuid.UUID) -> Optional[models.UploadRequestPhoto]: row = (await self._conn.execute(sqlalchemy.text(GET_UPLOAD_REQUEST_PHOTO_BY_ID), {"p1": id})).first() if row is None: @@ -148,8 +298,31 @@ async def get_upload_request_photo_by_id(self, *, id: uuid.UUID) -> Optional[mod visibility=row[10], status=row[11], created_at=row[12], + source=row[13], + transfer_status=row[14], ) + async def list_stale_pending_transfer_photos(self, *, dollar_1: Optional[str]) -> AsyncIterator[models.UploadRequestPhoto]: + result = await self._conn.stream(sqlalchemy.text(LIST_STALE_PENDING_TRANSFER_PHOTOS), {"p1": dollar_1}) + async for row in result: + yield models.UploadRequestPhoto( + id=row[0], + upload_request_id=row[1], + drive_file_id=row[2], + file_name=row[3], + mime_type=row[4], + size_bytes=row[5], + staging_storage_key=row[6], + final_storage_key=row[7], + taken_at=row[8], + day_number=row[9], + visibility=row[10], + status=row[11], + created_at=row[12], + source=row[13], + transfer_status=row[14], + ) + async def list_upload_request_photos_by_upload_request_id(self, *, upload_request_id: uuid.UUID) -> AsyncIterator[models.UploadRequestPhoto]: result = await self._conn.stream(sqlalchemy.text(LIST_UPLOAD_REQUEST_PHOTOS_BY_UPLOAD_REQUEST_ID), {"p1": upload_request_id}) async for row in result: @@ -167,6 +340,8 @@ async def list_upload_request_photos_by_upload_request_id(self, *, upload_reques visibility=row[10], status=row[11], created_at=row[12], + source=row[13], + transfer_status=row[14], ) async def list_upload_request_photos_by_upload_request_ids(self, *, dollar_1: List[uuid.UUID]) -> AsyncIterator[models.UploadRequestPhoto]: @@ -186,8 +361,32 @@ async def list_upload_request_photos_by_upload_request_ids(self, *, dollar_1: Li visibility=row[10], status=row[11], created_at=row[12], + source=row[13], + transfer_status=row[14], ) + async def reset_upload_request_photo_transfer_to_pending(self, *, id: uuid.UUID) -> Optional[models.UploadRequestPhoto]: + row = (await self._conn.execute(sqlalchemy.text(RESET_UPLOAD_REQUEST_PHOTO_TRANSFER_TO_PENDING), {"p1": id})).first() + if row is None: + return None + return models.UploadRequestPhoto( + id=row[0], + upload_request_id=row[1], + drive_file_id=row[2], + file_name=row[3], + mime_type=row[4], + size_bytes=row[5], + staging_storage_key=row[6], + final_storage_key=row[7], + taken_at=row[8], + day_number=row[9], + visibility=row[10], + status=row[11], + created_at=row[12], + source=row[13], + transfer_status=row[14], + ) + async def update_upload_request_photo_approval(self, *, id: uuid.UUID, status: str, final_storage_key: Optional[str]) -> Optional[models.UploadRequestPhoto]: row = (await self._conn.execute(sqlalchemy.text(UPDATE_UPLOAD_REQUEST_PHOTO_APPROVAL), {"p1": id, "p2": status, "p3": final_storage_key})).first() if row is None: @@ -206,6 +405,8 @@ async def update_upload_request_photo_approval(self, *, id: uuid.UUID, status: s visibility=row[10], status=row[11], created_at=row[12], + source=row[13], + transfer_status=row[14], ) async def update_upload_request_photo_status_by_upload_request_id(self, *, upload_request_id: uuid.UUID, status: str) -> AsyncIterator[models.UploadRequestPhoto]: @@ -225,4 +426,6 @@ async def update_upload_request_photo_status_by_upload_request_id(self, *, uploa visibility=row[10], status=row[11], created_at=row[12], + source=row[13], + transfer_status=row[14], ) diff --git a/db/generated/upload_requests.py b/db/generated/upload_requests.py index b0da8bb0..de8af04e 100644 --- a/db/generated/upload_requests.py +++ b/db/generated/upload_requests.py @@ -20,7 +20,7 @@ rejection_reason = NULL WHERE id = :p1 AND status = 'pending' -RETURNING id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id +RETURNING id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id, source """ @@ -30,11 +30,12 @@ group_id, drive_file_id, requested_by, - photo_count + photo_count, + source ) VALUES ( - :p1, :p2, :p3, :p4, :p5 + :p1, :p2, :p3, :p4, :p5, :p6 ) -RETURNING id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id +RETURNING id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id, source """ @@ -45,6 +46,7 @@ class CreateUploadRequestParams: drive_file_id: Optional[str] requested_by: uuid.UUID photo_count: int + source: str DELETE_UPLOAD_REQUEST = """-- name: delete_upload_request \\:exec @@ -54,21 +56,21 @@ class CreateUploadRequestParams: GET_UPLOAD_REQUEST_BY_ID = """-- name: get_upload_request_by_id \\:one -SELECT id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id +SELECT id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id, source FROM upload_requests WHERE id = :p1 """ LIST_UPLOAD_REQUESTS = """-- name: list_upload_requests \\:many -SELECT id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id +SELECT id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id, source FROM upload_requests ORDER BY created_at DESC """ LIST_UPLOAD_REQUESTS_BY_GROUP_ID = """-- name: list_upload_requests_by_group_id \\:many -SELECT id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id +SELECT id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id, source FROM upload_requests WHERE group_id = :p1 ORDER BY created_at ASC @@ -76,7 +78,7 @@ class CreateUploadRequestParams: LIST_UPLOAD_REQUESTS_BY_REQUESTER = """-- name: list_upload_requests_by_requester \\:many -SELECT id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id +SELECT id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id, source FROM upload_requests WHERE requested_by = :p1 ORDER BY created_at DESC @@ -84,7 +86,7 @@ class CreateUploadRequestParams: LIST_UPLOAD_REQUESTS_BY_REQUESTER_AND_STATUS = """-- name: list_upload_requests_by_requester_and_status \\:many -SELECT id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id +SELECT id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id, source FROM upload_requests WHERE requested_by = :p1 AND status = :p2 @@ -93,7 +95,7 @@ class CreateUploadRequestParams: LIST_UPLOAD_REQUESTS_BY_STATUS = """-- name: list_upload_requests_by_status \\:many -SELECT id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id +SELECT id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id, source FROM upload_requests WHERE status = :p1 ORDER BY created_at DESC @@ -108,7 +110,7 @@ class CreateUploadRequestParams: rejection_reason = :p3 WHERE id = :p1 AND status = 'pending' -RETURNING id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id +RETURNING id, event_id, drive_file_id, requested_by, approved_by, status, created_at, approved_at, photo_count, rejection_reason, group_id, source """ @@ -132,6 +134,7 @@ async def approve_upload_request(self, *, id: uuid.UUID, approved_by: Optional[u photo_count=row[8], rejection_reason=row[9], group_id=row[10], + source=row[11], ) async def create_upload_request(self, arg: CreateUploadRequestParams) -> Optional[models.UploadRequest]: @@ -141,6 +144,7 @@ async def create_upload_request(self, arg: CreateUploadRequestParams) -> Optiona "p3": arg.drive_file_id, "p4": arg.requested_by, "p5": arg.photo_count, + "p6": arg.source, })).first() if row is None: return None @@ -156,6 +160,7 @@ async def create_upload_request(self, arg: CreateUploadRequestParams) -> Optiona photo_count=row[8], rejection_reason=row[9], group_id=row[10], + source=row[11], ) async def delete_upload_request(self, *, id: uuid.UUID) -> None: @@ -177,6 +182,7 @@ async def get_upload_request_by_id(self, *, id: uuid.UUID) -> Optional[models.Up photo_count=row[8], rejection_reason=row[9], group_id=row[10], + source=row[11], ) async def list_upload_requests(self) -> AsyncIterator[models.UploadRequest]: @@ -194,6 +200,7 @@ async def list_upload_requests(self) -> AsyncIterator[models.UploadRequest]: photo_count=row[8], rejection_reason=row[9], group_id=row[10], + source=row[11], ) async def list_upload_requests_by_group_id(self, *, group_id: Optional[uuid.UUID]) -> AsyncIterator[models.UploadRequest]: @@ -211,6 +218,7 @@ async def list_upload_requests_by_group_id(self, *, group_id: Optional[uuid.UUID photo_count=row[8], rejection_reason=row[9], group_id=row[10], + source=row[11], ) async def list_upload_requests_by_requester(self, *, requested_by: uuid.UUID) -> AsyncIterator[models.UploadRequest]: @@ -228,6 +236,7 @@ async def list_upload_requests_by_requester(self, *, requested_by: uuid.UUID) -> photo_count=row[8], rejection_reason=row[9], group_id=row[10], + source=row[11], ) async def list_upload_requests_by_requester_and_status(self, *, requested_by: uuid.UUID, status: Any) -> AsyncIterator[models.UploadRequest]: @@ -245,6 +254,7 @@ async def list_upload_requests_by_requester_and_status(self, *, requested_by: uu photo_count=row[8], rejection_reason=row[9], group_id=row[10], + source=row[11], ) async def list_upload_requests_by_status(self, *, status: Any) -> AsyncIterator[models.UploadRequest]: @@ -262,6 +272,7 @@ async def list_upload_requests_by_status(self, *, status: Any) -> AsyncIterator[ photo_count=row[8], rejection_reason=row[9], group_id=row[10], + source=row[11], ) async def reject_upload_request(self, *, id: uuid.UUID, approved_by: Optional[uuid.UUID], rejection_reason: Optional[str]) -> Optional[models.UploadRequest]: @@ -280,4 +291,5 @@ async def reject_upload_request(self, *, id: uuid.UUID, approved_by: Optional[uu photo_count=row[8], rejection_reason=row[9], group_id=row[10], + source=row[11], ) diff --git a/db/queries/photos.sql b/db/queries/photos.sql index 3838e3c2..9904ad8a 100644 --- a/db/queries/photos.sql +++ b/db/queries/photos.sql @@ -84,3 +84,10 @@ SELECT urp.drive_file_id FROM upload_request_photos urp WHERE urp.final_storage_key = $1 LIMIT 1; + +-- name: MarkPhotoDriveSynced :one +UPDATE photos +SET drive_file_id = $2, + drive_synced_at = NOW() +WHERE id = $1 +RETURNING *; diff --git a/db/queries/upload_request_groups.sql b/db/queries/upload_request_groups.sql index 7c800f14..1dfe97ba 100644 --- a/db/queries/upload_request_groups.sql +++ b/db/queries/upload_request_groups.sql @@ -4,12 +4,21 @@ INSERT INTO upload_request_groups ( folder_id, requested_by, total_photo_count, - batch_count + batch_count, + source, + processing_status ) VALUES ( - $1, $2, $3, $4, $5 + $1, $2, $3, $4, $5, $6, $7 ) RETURNING *; +-- name: IncrementUploadRequestGroupCounts :one +UPDATE upload_request_groups +SET total_photo_count = total_photo_count + $2, + batch_count = batch_count + 1 +WHERE id = $1 +RETURNING *; + -- name: GetUploadRequestGroupById :one SELECT * FROM upload_request_groups diff --git a/db/queries/upload_request_photos.sql b/db/queries/upload_request_photos.sql index f78ab850..880f1afe 100644 --- a/db/queries/upload_request_photos.sql +++ b/db/queries/upload_request_photos.sql @@ -48,3 +48,54 @@ RETURNING *; -- name: DeleteUploadRequestPhotosByUploadRequestId :exec DELETE FROM upload_request_photos WHERE upload_request_id = $1; + +-- name: CreateDirectUploadRequestPhoto :one +INSERT INTO upload_request_photos ( + upload_request_id, + drive_file_id, + file_name, + mime_type, + size_bytes, + staging_storage_key, + taken_at, + day_number, + visibility, + status, + source, + transfer_status +) VALUES ( + $1, NULL, $2, $3, $4, $5, $6, $7, $8, 'staged', 'direct', 'pending_upload' +) +RETURNING *; + +-- name: ConfirmUploadRequestPhotoTransfer :one +UPDATE upload_request_photos +SET transfer_status = 'uploaded', + size_bytes = $2, + mime_type = $3 +WHERE id = $1 + AND transfer_status IN ('pending_upload', 'failed') +RETURNING *; + +-- name: FailUploadRequestPhotoTransfer :one +UPDATE upload_request_photos +SET transfer_status = 'failed' +WHERE id = $1 + AND transfer_status IN ('pending_upload', 'failed') +RETURNING *; + +-- name: ResetUploadRequestPhotoTransferToPending :one +UPDATE upload_request_photos +SET transfer_status = 'pending_upload' +WHERE id = $1 + AND transfer_status IN ('pending_upload', 'failed') +RETURNING *; + +-- name: ListStalePendingTransferPhotos :many +SELECT * +FROM upload_request_photos +WHERE source = 'direct' + AND transfer_status = 'pending_upload' + AND created_at <= NOW() - ($1 || ' minutes')::interval +ORDER BY created_at ASC +LIMIT 500; diff --git a/db/queries/upload_requests.sql b/db/queries/upload_requests.sql index 31eb3739..2fd571bc 100644 --- a/db/queries/upload_requests.sql +++ b/db/queries/upload_requests.sql @@ -4,9 +4,10 @@ INSERT INTO upload_requests ( group_id, drive_file_id, requested_by, - photo_count + photo_count, + source ) VALUES ( - $1, $2, $3, $4, $5 + $1, $2, $3, $4, $5, $6 ) RETURNING *; diff --git a/makefile b/makefile index 36690ed8..7e19b776 100644 --- a/makefile +++ b/makefile @@ -66,6 +66,8 @@ run-workers: uv run python -m app.worker.storage_cleaner.main & \ uv run python -m app.worker.email_worker.main & \ uv run python -m app.worker.event_lifecycle.main & \ + uv run python -m app.worker.upload_reconciler.main & \ + uv run python -m app.worker.drive_sync.main & \ wait lint: diff --git a/migrations/sql/down/add_direct_upload_support.sql b/migrations/sql/down/add_direct_upload_support.sql new file mode 100644 index 00000000..d9837913 --- /dev/null +++ b/migrations/sql/down/add_direct_upload_support.sql @@ -0,0 +1,11 @@ +ALTER TABLE upload_request_photos + ALTER COLUMN drive_file_id SET NOT NULL, + DROP COLUMN transfer_status, + DROP COLUMN source; + +ALTER TABLE upload_requests + DROP COLUMN source; + +ALTER TABLE upload_request_groups + ALTER COLUMN folder_id SET NOT NULL, + DROP COLUMN source; diff --git a/migrations/sql/down/add_drive_sync_fields_to_photos.sql b/migrations/sql/down/add_drive_sync_fields_to_photos.sql new file mode 100644 index 00000000..d6488697 --- /dev/null +++ b/migrations/sql/down/add_drive_sync_fields_to_photos.sql @@ -0,0 +1,3 @@ +ALTER TABLE photos + DROP COLUMN drive_file_id, + DROP COLUMN drive_synced_at; diff --git a/migrations/sql/up/add_direct_upload_support.sql b/migrations/sql/up/add_direct_upload_support.sql new file mode 100644 index 00000000..0319c12f --- /dev/null +++ b/migrations/sql/up/add_direct_upload_support.sql @@ -0,0 +1,11 @@ +ALTER TABLE upload_request_groups + ADD COLUMN source character varying(16) DEFAULT 'drive'::character varying NOT NULL, + ALTER COLUMN folder_id DROP NOT NULL; + +ALTER TABLE upload_requests + ADD COLUMN source character varying(16) DEFAULT 'drive'::character varying NOT NULL; + +ALTER TABLE upload_request_photos + ADD COLUMN source character varying(16) DEFAULT 'drive'::character varying NOT NULL, + ADD COLUMN transfer_status character varying(16) DEFAULT 'uploaded'::character varying NOT NULL, + ALTER COLUMN drive_file_id DROP NOT NULL; diff --git a/migrations/sql/up/add_drive_sync_fields_to_photos.sql b/migrations/sql/up/add_drive_sync_fields_to_photos.sql new file mode 100644 index 00000000..2455b45d --- /dev/null +++ b/migrations/sql/up/add_drive_sync_fields_to_photos.sql @@ -0,0 +1,3 @@ +ALTER TABLE photos + ADD COLUMN drive_file_id text, + ADD COLUMN drive_synced_at timestamp with time zone; diff --git a/migrations/versions/5425a051d68c_add_direct_upload_support.py b/migrations/versions/5425a051d68c_add_direct_upload_support.py new file mode 100644 index 00000000..0e8807ce --- /dev/null +++ b/migrations/versions/5425a051d68c_add_direct_upload_support.py @@ -0,0 +1,25 @@ +"""add_direct_upload_support + +Revision ID: 5425a051d68c +Revises: afb0a93aa21e +Create Date: 2026-08-25 02:35:19.168255 + +""" +from typing import Sequence, Union + +from migrations.helper import run_sql_up, run_sql_down + + +# revision identifiers, used by Alembic. +revision: str = '5425a051d68c' +down_revision: Union[str, Sequence[str], None] = 'afb0a93aa21e' +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + run_sql_up("add_direct_upload_support") + + +def downgrade() -> None: + run_sql_down("add_direct_upload_support") diff --git a/migrations/versions/af58506b53a9_add_drive_sync_fields_to_photos.py b/migrations/versions/af58506b53a9_add_drive_sync_fields_to_photos.py new file mode 100644 index 00000000..bb6cac94 --- /dev/null +++ b/migrations/versions/af58506b53a9_add_drive_sync_fields_to_photos.py @@ -0,0 +1,25 @@ +"""add_drive_sync_fields_to_photos + +Revision ID: af58506b53a9 +Revises: 5425a051d68c +Create Date: 2026-08-25 02:58:36.347923 + +""" +from typing import Sequence, Union + +from migrations.helper import run_sql_up, run_sql_down + + +# revision identifiers, used by Alembic. +revision: str = 'af58506b53a9' +down_revision: Union[str, Sequence[str], None] = '5425a051d68c' +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + run_sql_up("add_drive_sync_fields_to_photos") + + +def downgrade() -> None: + run_sql_down("add_drive_sync_fields_to_photos") diff --git a/tests/unit/test_direct_uploads.py b/tests/unit/test_direct_uploads.py new file mode 100644 index 00000000..b1435af3 --- /dev/null +++ b/tests/unit/test_direct_uploads.py @@ -0,0 +1,404 @@ +import uuid +from datetime import datetime, timezone +from unittest.mock import AsyncMock, patch + +import pytest + +from app.schema.internal.uploads import DirectFileInput +from app.service.staged_upload_storage import StoredObject +from app.service.upload_requests import UploadRequestsService +from db.generated.models import ( + StaffUser, + UploadRequest, + UploadRequestGroup, + UploadRequestPhoto, +) + + +@pytest.fixture +def mock_upload_request_group_querier(): + return AsyncMock() + +@pytest.fixture +def mock_upload_request_querier(): + return AsyncMock() + +@pytest.fixture +def mock_upload_request_photo_querier(): + return AsyncMock() + +@pytest.fixture +def mock_photo_querier(): + return AsyncMock() + +@pytest.fixture +def mock_staged_upload_storage(): + mock = AsyncMock() + mock.store_staging_object.return_value = StoredObject(storage_key="test_storage_key", content_type="image/jpeg", file_name="photo.jpg") + mock.create_presigned_staging_upload.return_value = ("staging/upload-requests/req1/photo1.jpg", "https://minio.local/signed-url") + return mock + +@pytest.fixture +def mock_staff_drive_service(): + mock = AsyncMock() + mock.get_access_token_for_staff_user.return_value = "fake_access_token" + mock.staff_user_querier = AsyncMock() + return mock + +@pytest.fixture +def mock_staff_notifications_service(): + return AsyncMock() + +@pytest.fixture +def mock_audit_service(): + return AsyncMock() + +@pytest.fixture +def upload_requests_service( + mock_upload_request_group_querier, + mock_upload_request_querier, + mock_upload_request_photo_querier, + mock_photo_querier, + mock_staged_upload_storage, + mock_staff_drive_service, + mock_staff_notifications_service, + mock_audit_service, +): + return UploadRequestsService( + upload_request_group_querier=mock_upload_request_group_querier, + upload_request_querier=mock_upload_request_querier, + upload_request_photo_querier=mock_upload_request_photo_querier, + photo_querier=mock_photo_querier, + staged_upload_storage=mock_staged_upload_storage, + staff_drive_service=mock_staff_drive_service, + staff_notifications_service=mock_staff_notifications_service, + audit_service=mock_audit_service, + ) + +@pytest.fixture +def mock_staff_user(): + return StaffUser( + id=uuid.uuid4(), + email="test@multai.com", + password="hash", + role="multi", + created_at=datetime.now(timezone.utc), + updated_at=datetime.now(timezone.utc), + ) + + +def _make_group(group_id, event_id, requested_by_id, **overrides): + defaults = dict( + id=group_id, event_id=event_id, folder_id=None, requested_by=requested_by_id, + approved_by=None, status="pending", total_photo_count=0, batch_count=0, + created_at=datetime.now(timezone.utc), approved_at=None, rejection_reason=None, + processing_status="completed", processed_photo_count=0, failed_photo_count=0, + error_message=None, source="direct", + ) + defaults.update(overrides) + return UploadRequestGroup(**defaults) + + +def _make_request(request_id, event_id, requested_by_id, group_id, **overrides): + defaults = dict( + id=request_id, event_id=event_id, drive_file_id=None, requested_by=requested_by_id, + approved_by=None, status="pending", created_at=datetime.now(timezone.utc), + approved_at=None, photo_count=1, rejection_reason=None, group_id=group_id, + source="direct", + ) + defaults.update(overrides) + return UploadRequest(**defaults) + + +def _make_photo(photo_id, request_id, **overrides): + defaults = dict( + id=photo_id, upload_request_id=request_id, drive_file_id=None, file_name="a.jpg", + mime_type="image/jpeg", size_bytes=1000, staging_storage_key="staging/x.jpg", + final_storage_key=None, taken_at=None, day_number=None, visibility="private", + status="staged", created_at=datetime.now(timezone.utc), + source="direct", transfer_status="pending_upload", + ) + defaults.update(overrides) + return UploadRequestPhoto(**defaults) + + +@pytest.mark.asyncio +async def test_create_direct_group_sets_source_and_completed_processing( + upload_requests_service, + mock_upload_request_group_querier, + mock_staff_user, +): + event_id = uuid.uuid4() + group_id = uuid.uuid4() + mock_upload_request_group_querier.create_upload_request_group.return_value = _make_group( + group_id, event_id, mock_staff_user.id, + ) + + group = await upload_requests_service.create_direct_group( + event_id=event_id, requested_by=mock_staff_user, + ) + + assert group.id == group_id + call_args = mock_upload_request_group_querier.create_upload_request_group.call_args + params = call_args.args[0] if call_args.args else call_args.kwargs["arg"] + assert params.folder_id is None + assert params.source == "direct" + assert params.processing_status == "completed" + + +@pytest.mark.asyncio +async def test_register_direct_batch_creates_pending_photos_and_returns_urls( + upload_requests_service, + mock_upload_request_querier, + mock_upload_request_photo_querier, + mock_upload_request_group_querier, + mock_staged_upload_storage, + mock_staff_user, +): + group_id = uuid.uuid4() + event_id = uuid.uuid4() + request_id = uuid.uuid4() + photo_id = uuid.uuid4() + + mock_upload_request_group_querier.get_upload_request_group_by_id.return_value = _make_group( + group_id, event_id, mock_staff_user.id, + ) + mock_upload_request_querier.create_upload_request.return_value = _make_request( + request_id, event_id, mock_staff_user.id, group_id, + ) + mock_upload_request_photo_querier.create_direct_upload_request_photo.return_value = _make_photo( + photo_id, request_id, staging_storage_key="staging/upload-requests/req1/photo1.jpg", + ) + + results = await upload_requests_service.register_direct_batch( + group_id=group_id, + files=[DirectFileInput( + file_name="a.jpg", mime_type="image/jpeg", size_bytes=1000, + taken_at=None, day_number=None, visibility="private", + )], + requested_by=mock_staff_user, + ) + + assert len(results) == 1 + photo, url = results[0] + assert photo.id == photo_id + assert url == "https://minio.local/signed-url" + mock_staged_upload_storage.create_presigned_staging_upload.assert_awaited() + mock_upload_request_group_querier.increment_upload_request_group_counts.assert_awaited_once_with( + id=group_id, total_photo_count=1, + ) + + +@pytest.mark.asyncio +async def test_register_direct_batch_rejects_oversized_batch( + upload_requests_service, + mock_upload_request_group_querier, + mock_staff_user, +): + group_id = uuid.uuid4() + from app.core.exceptions import AppException + + files = [ + DirectFileInput(file_name=f"{i}.jpg", mime_type="image/jpeg", size_bytes=1000, taken_at=None, day_number=None, visibility="private") + for i in range(201) + ] + + with pytest.raises(Exception): + await upload_requests_service.register_direct_batch( + group_id=group_id, files=files, requested_by=mock_staff_user, + ) + mock_upload_request_group_querier.get_upload_request_group_by_id.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_confirm_direct_upload_success_updates_transfer_status( + upload_requests_service, + mock_upload_request_photo_querier, + mock_staged_upload_storage, + mock_staff_user, +): + from app.infra.minio import ObjectStat + + photo_id = uuid.uuid4() + request_id = uuid.uuid4() + existing_photo = _make_photo(photo_id, request_id, transfer_status="pending_upload") + mock_upload_request_photo_querier.get_upload_request_photo_by_id.return_value = existing_photo + mock_staged_upload_storage.stat_staging_object.return_value = ObjectStat(size=1000, content_type="image/jpeg") + mock_upload_request_photo_querier.confirm_upload_request_photo_transfer.return_value = _make_photo( + photo_id, request_id, transfer_status="uploaded", + ) + + result = await upload_requests_service.confirm_direct_upload( + photo_id=photo_id, requested_by=mock_staff_user, + ) + + assert result.id == photo_id + mock_upload_request_photo_querier.confirm_upload_request_photo_transfer.assert_awaited_once_with( + id=photo_id, size_bytes=1000, mime_type="image/jpeg", + ) + + +@pytest.mark.asyncio +async def test_confirm_direct_upload_marks_failed_when_object_missing( + upload_requests_service, + mock_upload_request_photo_querier, + mock_staged_upload_storage, + mock_staff_user, +): + photo_id = uuid.uuid4() + request_id = uuid.uuid4() + existing_photo = _make_photo(photo_id, request_id, transfer_status="pending_upload") + mock_upload_request_photo_querier.get_upload_request_photo_by_id.return_value = existing_photo + mock_staged_upload_storage.stat_staging_object.return_value = None + mock_upload_request_photo_querier.fail_upload_request_photo_transfer.return_value = _make_photo( + photo_id, request_id, transfer_status="failed", + ) + + with pytest.raises(Exception): + await upload_requests_service.confirm_direct_upload( + photo_id=photo_id, requested_by=mock_staff_user, + ) + + mock_upload_request_photo_querier.fail_upload_request_photo_transfer.assert_awaited_once_with(id=photo_id) + + +@pytest.mark.asyncio +async def test_approve_request_blocked_when_photo_not_fully_uploaded( + upload_requests_service, + mock_upload_request_querier, + mock_upload_request_photo_querier, + mock_staff_user, +): + request_id = uuid.uuid4() + photo_id = uuid.uuid4() + + mock_upload_request_querier.get_upload_request_by_id.return_value = _make_request( + request_id, uuid.uuid4(), mock_staff_user.id, None, + ) + not_uploaded_photo = _make_photo(photo_id, request_id, transfer_status="pending_upload") + + async def _photos_iter(upload_request_id): + yield not_uploaded_photo + + mock_upload_request_photo_querier.list_upload_request_photos_by_upload_request_id = _photos_iter + + with pytest.raises(Exception) as exc_info: + await upload_requests_service.approve_request( + request_id=request_id, approved_by=mock_staff_user, + ) + assert "have not finished uploading" in str(exc_info.value) + + +@pytest.mark.asyncio +async def test_resume_direct_group_reissues_urls_for_pending_and_failed_only( + upload_requests_service, + mock_upload_request_group_querier, + mock_upload_request_querier, + mock_upload_request_photo_querier, + mock_staged_upload_storage, + mock_staff_user, +): + group_id = uuid.uuid4() + event_id = uuid.uuid4() + request_id = uuid.uuid4() + failed_photo_id = uuid.uuid4() + uploaded_photo_id = uuid.uuid4() + + mock_upload_request_group_querier.get_upload_request_group_by_id.return_value = _make_group( + group_id, event_id, mock_staff_user.id, total_photo_count=2, batch_count=1, failed_photo_count=1, + ) + + async def _requests_iter(group_id): + yield _make_request(request_id, event_id, mock_staff_user.id, group_id, photo_count=2) + + mock_upload_request_querier.list_upload_requests_by_group_id = _requests_iter + + failed_photo = _make_photo(failed_photo_id, request_id, file_name="fail.jpg", staging_storage_key="staging/fail.jpg", transfer_status="failed") + uploaded_photo = _make_photo(uploaded_photo_id, request_id, file_name="ok.jpg", staging_storage_key="staging/ok.jpg", transfer_status="uploaded") + + async def _photos_iter(dollar_1): + for p in [failed_photo, uploaded_photo]: + yield p + + mock_upload_request_photo_querier.list_upload_request_photos_by_upload_request_ids = _photos_iter + mock_staged_upload_storage.create_presigned_staging_upload.return_value = ("staging/fail.jpg", "https://minio.local/resumed") + mock_upload_request_photo_querier.reset_upload_request_photo_transfer_to_pending.return_value = _make_photo( + failed_photo_id, request_id, file_name="fail.jpg", transfer_status="pending_upload", + ) + + results = await upload_requests_service.resume_direct_group( + group_id=group_id, requested_by=mock_staff_user, + ) + + assert len(results) == 1 + photo, url = results[0] + assert photo.id == failed_photo_id + assert url == "https://minio.local/resumed" + + +@pytest.mark.asyncio +async def test_approve_request_publishes_drive_sync_event_for_direct_photo_only( + upload_requests_service, + mock_upload_request_querier, + mock_upload_request_photo_querier, + mock_photo_querier, + mock_staged_upload_storage, + mock_staff_user, +): + from db.generated.models import Photo + + request_id = uuid.uuid4() + event_id = uuid.uuid4() + photo_id = uuid.uuid4() + + mock_upload_request_querier.get_upload_request_by_id.return_value = _make_request( + request_id, event_id, mock_staff_user.id, None, + ) + + async def _photos_iter(upload_request_id): + yield _make_photo(photo_id, request_id, transfer_status="uploaded", source="direct") + + mock_upload_request_photo_querier.list_upload_request_photos_by_upload_request_id = _photos_iter + mock_staged_upload_storage.promote_to_final.return_value = "events/e1/p1.jpg" + mock_photo_querier.create_photo.return_value = Photo( + id=photo_id, event_id=event_id, uploaded_by=None, storage_key="events/e1/p1.jpg", + taken_at=None, day_number=None, visibility="private", status="pending", + created_at=datetime.now(timezone.utc), drive_file_id=None, drive_synced_at=None, + ) + mock_upload_request_photo_querier.update_upload_request_photo_approval.return_value = _make_photo( + photo_id, request_id, source="direct", transfer_status="uploaded", + ) + mock_upload_request_querier.approve_upload_request.return_value = _make_request( + request_id, event_id, mock_staff_user.id, None, + ) + + with patch("app.service.upload_requests.NatsClient.publish") as mock_publish: + await upload_requests_service.approve_request( + request_id=request_id, approved_by=mock_staff_user, + ) + + published_subjects = [call.args[0] for call in mock_publish.call_args_list] + from app.infra.nats import NatsSubjects + assert NatsSubjects.PHOTO_DRIVE_SYNC_REQUESTED in published_subjects + + +@pytest.mark.asyncio +async def test_fail_direct_upload_marks_transfer_failed( + upload_requests_service, + mock_upload_request_photo_querier, + mock_staff_user, +): + photo_id = uuid.uuid4() + request_id = uuid.uuid4() + existing_photo = _make_photo(photo_id, request_id, transfer_status="pending_upload") + mock_upload_request_photo_querier.get_upload_request_photo_by_id.return_value = existing_photo + mock_upload_request_photo_querier.fail_upload_request_photo_transfer.return_value = _make_photo( + photo_id, request_id, transfer_status="failed", + ) + + result = await upload_requests_service.fail_direct_upload( + photo_id=photo_id, requested_by=mock_staff_user, + ) + + assert result.id == photo_id + mock_upload_request_photo_querier.fail_upload_request_photo_transfer.assert_awaited_once_with(id=photo_id) diff --git a/tests/unit/test_minio.py b/tests/unit/test_minio.py index 66115611..1587c27f 100644 --- a/tests/unit/test_minio.py +++ b/tests/unit/test_minio.py @@ -8,6 +8,7 @@ from app.infra.minio import ( Bucket, ImageBucket, + ObjectStat, WaSimBucket, init_minio_client, ) @@ -200,3 +201,45 @@ async def test_wa_sim_bucket_auto_name(mock_minio_client, mock_upload_file): # WaSimBucket generates 16 digit string assert len(object_name) == 16 assert object_name.isdigit() + + +@pytest.mark.asyncio +async def test_presigned_put_url_calls_client_with_expiry(mock_minio_client): + Bucket.client = mock_minio_client + mock_minio_client.presigned_put_object = AsyncMock(return_value="https://minio.local/signed") + bucket = Bucket("test_bucket", "") + + url = await bucket.presigned_put_url("staging/foo.jpg", expires_seconds=1800) + + assert url == "https://minio.local/signed" + mock_minio_client.presigned_put_object.assert_awaited_once() + kwargs = mock_minio_client.presigned_put_object.call_args[1] + assert kwargs["object_name"] == "staging/foo.jpg" + assert kwargs["expires"].total_seconds() == 1800 + + +@pytest.mark.asyncio +async def test_stat_returns_object_stat_when_present(mock_minio_client): + Bucket.client = mock_minio_client + stat_result = MagicMock(size=12345, content_type="image/jpeg") + mock_minio_client.stat_object = AsyncMock(return_value=stat_result) + bucket = Bucket("test_bucket", "") + + result = await bucket.stat("staging/foo.jpg") + + assert result == ObjectStat(size=12345, content_type="image/jpeg") + + +@pytest.mark.asyncio +async def test_stat_returns_none_when_object_missing(mock_minio_client): + Bucket.client = mock_minio_client + error = S3Error( + code="NoSuchKey", message="not found", resource="", request_id="", + host_id="", response=MagicMock(), + ) + mock_minio_client.stat_object = AsyncMock(side_effect=error) + bucket = Bucket("test_bucket", "") + + result = await bucket.stat("staging/missing.jpg") + + assert result is None diff --git a/tests/unit/test_upload_requests.py b/tests/unit/test_upload_requests.py index 691c30c8..dd31c391 100644 --- a/tests/unit/test_upload_requests.py +++ b/tests/unit/test_upload_requests.py @@ -122,6 +122,7 @@ async def test_create_request_success( rejection_reason=None, created_at=datetime.now(timezone.utc), approved_at=None, + source="drive", ) mock_upload_request_photo_querier.create_upload_request_photo.return_value = UploadRequestPhoto( @@ -138,6 +139,8 @@ async def test_create_request_success( visibility="public", status="staged", created_at=datetime.now(timezone.utc), + source="drive", + transfer_status="uploaded", ) photos = [ @@ -191,7 +194,7 @@ async def test_create_request_duplicate_conflict( request_id = uuid.uuid4() mock_upload_request_querier.create_upload_request.return_value = UploadRequest( - id=request_id, event_id=event_id, group_id=None, drive_file_id=None, requested_by=mock_staff_user.id, photo_count=1, status="pending", approved_by=None, rejection_reason=None, created_at=datetime.now(timezone.utc), approved_at=None + id=request_id, event_id=event_id, group_id=None, drive_file_id=None, requested_by=mock_staff_user.id, photo_count=1, status="pending", approved_by=None, rejection_reason=None, created_at=datetime.now(timezone.utc), approved_at=None, source="drive" ) # Simulate DB Conflict (Duplicate) on photo insert @@ -225,7 +228,7 @@ async def test_create_group_from_folder( group_id = uuid.uuid4() mock_upload_request_group_querier.create_upload_request_group.return_value = UploadRequestGroup( - id=group_id, event_id=event_id, folder_id="folder_123", requested_by=mock_staff_user.id, total_photo_count=0, batch_count=0, processed_photo_count=0, failed_photo_count=0, processing_status="pending", error_message=None, created_at=datetime.now(timezone.utc), status="pending", approved_by=None, approved_at=None, rejection_reason=None + id=group_id, event_id=event_id, folder_id="folder_123", requested_by=mock_staff_user.id, total_photo_count=0, batch_count=0, processed_photo_count=0, failed_photo_count=0, processing_status="pending", error_message=None, created_at=datetime.now(timezone.utc), status="pending", approved_by=None, approved_at=None, rejection_reason=None, source="drive" ) with patch("app.service.upload_requests.NatsClient.publish") as mock_publish: @@ -249,7 +252,7 @@ async def test_process_group_import_no_images( ): group_id = uuid.uuid4() mock_upload_request_group_querier.start_upload_request_group_processing.return_value = UploadRequestGroup( - id=group_id, event_id=uuid.uuid4(), folder_id="folder_123", requested_by=mock_staff_user.id, total_photo_count=0, batch_count=0, processed_photo_count=0, failed_photo_count=0, processing_status="processing", error_message=None, created_at=datetime.now(timezone.utc), status="pending", approved_by=None, approved_at=None, rejection_reason=None + id=group_id, event_id=uuid.uuid4(), folder_id="folder_123", requested_by=mock_staff_user.id, total_photo_count=0, batch_count=0, processed_photo_count=0, failed_photo_count=0, processing_status="processing", error_message=None, created_at=datetime.now(timezone.utc), status="pending", approved_by=None, approved_at=None, rejection_reason=None, source="drive" ) mock_staff_drive_service.staff_user_querier.get_staff_user_by_id.return_value = mock_staff_user @@ -257,7 +260,7 @@ async def test_process_group_import_no_images( with patch("app.service.upload_requests.GoogleDriveClient.list_folder_files", return_value=[]): async def mock_get_group(*args, **kwargs): yield UploadRequestGroup( - id=group_id, event_id=uuid.uuid4(), folder_id="folder_123", requested_by=mock_staff_user.id, total_photo_count=0, batch_count=0, processed_photo_count=0, failed_photo_count=0, processing_status="processing", error_message=None, created_at=datetime.now(timezone.utc), status="pending", approved_by=None, approved_at=None, rejection_reason=None + id=group_id, event_id=uuid.uuid4(), folder_id="folder_123", requested_by=mock_staff_user.id, total_photo_count=0, batch_count=0, processed_photo_count=0, failed_photo_count=0, processing_status="processing", error_message=None, created_at=datetime.now(timezone.utc), status="pending", approved_by=None, approved_at=None, rejection_reason=None, source="drive" ) mock_upload_request_querier.list_upload_requests_by_group_id = mock_get_group @@ -267,7 +270,7 @@ async def mock_list_photos_by_ids(*args, **kwargs): upload_requests_service.upload_request_photo_querier.list_upload_request_photos_by_upload_request_ids = mock_list_photos_by_ids mock_upload_request_group_querier.get_upload_request_group_by_id.return_value = UploadRequestGroup( - id=group_id, event_id=uuid.uuid4(), folder_id="folder_123", requested_by=mock_staff_user.id, total_photo_count=0, batch_count=0, processed_photo_count=0, failed_photo_count=0, processing_status="processing", error_message=None, created_at=datetime.now(timezone.utc), status="pending", approved_by=None, approved_at=None, rejection_reason=None + id=group_id, event_id=uuid.uuid4(), folder_id="folder_123", requested_by=mock_staff_user.id, total_photo_count=0, batch_count=0, processed_photo_count=0, failed_photo_count=0, processing_status="processing", error_message=None, created_at=datetime.now(timezone.utc), status="pending", approved_by=None, approved_at=None, rejection_reason=None, source="drive" ) await upload_requests_service.process_group_import( @@ -294,12 +297,12 @@ async def test_approve_request_without_side_effects( event_id = uuid.uuid4() mock_upload_request_querier.get_upload_request_by_id.return_value = UploadRequest( - id=request_id, event_id=event_id, group_id=None, drive_file_id=None, requested_by=uuid.uuid4(), photo_count=1, status="pending", approved_by=None, rejection_reason=None, created_at=datetime.now(timezone.utc), approved_at=None + id=request_id, event_id=event_id, group_id=None, drive_file_id=None, requested_by=uuid.uuid4(), photo_count=1, status="pending", approved_by=None, rejection_reason=None, created_at=datetime.now(timezone.utc), approved_at=None, source="drive" ) async def mock_list_photos(*args, **kwargs): yield UploadRequestPhoto( - id=photo_id, upload_request_id=request_id, drive_file_id="drive_1", file_name="p.jpg", mime_type="image/jpeg", size_bytes=100, staging_storage_key="stage_key", final_storage_key=None, taken_at=None, day_number=None, visibility="public", status="staged", created_at=datetime.now(timezone.utc) + id=photo_id, upload_request_id=request_id, drive_file_id="drive_1", file_name="p.jpg", mime_type="image/jpeg", size_bytes=100, staging_storage_key="stage_key", final_storage_key=None, taken_at=None, day_number=None, visibility="public", status="staged", created_at=datetime.now(timezone.utc), source="drive", transfer_status="uploaded" ) mock_upload_request_photo_querier.list_upload_request_photos_by_upload_request_id = mock_list_photos