From 53fc8aa8d14e757fcd9f4b628eba8e4c86b15a8c Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 02:35:56 +0100 Subject: [PATCH 01/13] feat: add schema support for direct (non-Drive) bulk uploads --- .../sql/down/add_direct_upload_support.sql | 11 ++++++++ .../sql/up/add_direct_upload_support.sql | 11 ++++++++ .../5425a051d68c_add_direct_upload_support.py | 25 +++++++++++++++++++ 3 files changed, 47 insertions(+) create mode 100644 migrations/sql/down/add_direct_upload_support.sql create mode 100644 migrations/sql/up/add_direct_upload_support.sql create mode 100644 migrations/versions/5425a051d68c_add_direct_upload_support.py 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 0000000..d983791 --- /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/up/add_direct_upload_support.sql b/migrations/sql/up/add_direct_upload_support.sql new file mode 100644 index 0000000..0319c12 --- /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/versions/5425a051d68c_add_direct_upload_support.py b/migrations/versions/5425a051d68c_add_direct_upload_support.py new file mode 100644 index 0000000..0e8807c --- /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") From 3e881b85847514923b882acca0bee6d09b112cac Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 02:36:52 +0100 Subject: [PATCH 02/13] feat: add SQLC queries for direct upload registration and transfer tracking --- db/generated/models.py | 8 +- db/generated/upload_request_groups.py | 80 ++++++++-- db/generated/upload_request_photos.py | 217 +++++++++++++++++++++++++- db/generated/upload_requests.py | 34 ++-- db/queries/upload_request_groups.sql | 13 +- db/queries/upload_request_photos.sql | 51 ++++++ db/queries/upload_requests.sql | 5 +- 7 files changed, 369 insertions(+), 39 deletions(-) diff --git a/db/generated/models.py b/db/generated/models.py index 86617af..2d6bcf6 100644 --- a/db/generated/models.py +++ b/db/generated/models.py @@ -208,13 +208,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 +228,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 +246,8 @@ class UploadRequestPhoto: visibility: str status: str created_at: datetime.datetime + source: str + transfer_status: str @dataclasses.dataclass() diff --git a/db/generated/upload_request_groups.py b/db/generated/upload_request_groups.py index 039b1f0..ceb91cf 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 1cd3ebb..ba33dee 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 b0da8bb..de8af04 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/upload_request_groups.sql b/db/queries/upload_request_groups.sql index 7c800f1..1dfe97b 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 f78ab85..880f1af 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 31eb373..2fd571b 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 *; From 0e7a9470f202367cae5cb2b567986d024c41f12a Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 02:41:27 +0100 Subject: [PATCH 03/13] feat: add presigned PUT/stat support; fix Drive-flow params for new source/transfer_status columns --- app/infra/minio.py | 27 +++++++++++++++++ app/service/staged_upload_storage.py | 21 +++++++++++++- app/service/upload_requests.py | 3 ++ tests/unit/test_minio.py | 43 ++++++++++++++++++++++++++++ tests/unit/test_upload_requests.py | 17 ++++++----- 5 files changed, 103 insertions(+), 8 deletions(-) diff --git a/app/infra/minio.py b/app/infra/minio.py index e6249da..9a666c5 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/service/staged_upload_storage.py b/app/service/staged_upload_storage.py index 813fa2b..e8f9d97 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 0ab404d..1d19eef 100644 --- a/app/service/upload_requests.py +++ b/app/service/upload_requests.py @@ -282,6 +282,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: @@ -572,6 +573,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: diff --git a/tests/unit/test_minio.py b/tests/unit/test_minio.py index 6611561..1587c27 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 691c30c..dd31c39 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 From 3f0677916899e303deeb867335c6a43447b1c07f Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 02:41:43 +0100 Subject: [PATCH 04/13] feat: add config settings for direct upload --- app/core/config.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/app/core/config.py b/app/core/config.py index c9270f0..5a7fe12 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 From 43e54a5025523907aeabbaa93286facf6b1de01d Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 02:43:08 +0100 Subject: [PATCH 05/13] feat: add direct upload batch registration, confirm, and fail to UploadRequestsService --- app/schema/internal/uploads.py | 10 ++ app/service/upload_requests.py | 138 ++++++++++++++- tests/unit/test_direct_uploads.py | 284 ++++++++++++++++++++++++++++++ 3 files changed, 431 insertions(+), 1 deletion(-) create mode 100644 tests/unit/test_direct_uploads.py diff --git a/app/schema/internal/uploads.py b/app/schema/internal/uploads.py index c8b91da..58f8249 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/service/upload_requests.py b/app/service/upload_requests.py index 1d19eef..c857836 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 @@ -596,6 +597,141 @@ 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 process_group_import( self, *, diff --git a/tests/unit/test_direct_uploads.py b/tests/unit/test_direct_uploads.py new file mode 100644 index 0000000..5846e84 --- /dev/null +++ b/tests/unit/test_direct_uploads.py @@ -0,0 +1,284 @@ +import uuid +from datetime import datetime, timezone +from unittest.mock import AsyncMock + +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_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) From 7ac3169dd3327c4d8589ecc2dde709fff2f42a30 Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 02:44:04 +0100 Subject: [PATCH 06/13] feat: add resume and pre-approval transfer guard for direct uploads --- app/service/upload_requests.py | 46 +++++++++++++++++++ tests/unit/test_direct_uploads.py | 74 +++++++++++++++++++++++++++++++ 2 files changed, 120 insertions(+) diff --git a/app/service/upload_requests.py b/app/service/upload_requests.py index c857836..f64e346 100644 --- a/app/service/upload_requests.py +++ b/app/service/upload_requests.py @@ -338,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: @@ -732,6 +742,42 @@ async def fail_direct_upload( 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, *, diff --git a/tests/unit/test_direct_uploads.py b/tests/unit/test_direct_uploads.py index 5846e84..ca73237 100644 --- a/tests/unit/test_direct_uploads.py +++ b/tests/unit/test_direct_uploads.py @@ -262,6 +262,80 @@ async def test_confirm_direct_upload_marks_failed_when_object_missing( 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_fail_direct_upload_marks_transfer_failed( upload_requests_service, From 745fd69c8d2c3c1ace825545011b9e98056b0b7e Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 02:44:55 +0100 Subject: [PATCH 07/13] feat: add request/response schemas for direct upload endpoints --- app/schema/request/staff/uploads_direct.py | 22 +++++++++++++++++++++ app/schema/response/staff/upload_groups.py | 6 ++++-- app/schema/response/staff/uploads.py | 5 ++++- app/schema/response/staff/uploads_direct.py | 18 +++++++++++++++++ 4 files changed, 48 insertions(+), 3 deletions(-) create mode 100644 app/schema/request/staff/uploads_direct.py create mode 100644 app/schema/response/staff/uploads_direct.py diff --git a/app/schema/request/staff/uploads_direct.py b/app/schema/request/staff/uploads_direct.py new file mode 100644 index 0000000..068c1b6 --- /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 7a264c2..f691cc9 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 863414c..60a1354 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 0000000..fba3d4c --- /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] From b57af9ecbf66a3dc126c6cdad49aa7d38071f3b9 Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 02:46:19 +0100 Subject: [PATCH 08/13] feat: add staff router endpoints for direct upload --- app/router/staff/__init__.py | 2 + app/router/staff/uploads_direct.py | 119 +++++++++++++++++++++++++++++ 2 files changed, 121 insertions(+) create mode 100644 app/router/staff/uploads_direct.py diff --git a/app/router/staff/__init__.py b/app/router/staff/__init__.py index 35cbe0e..11023d7 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 0000000..0e9dbd1 --- /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, bypassing the mobile app. This is useful for bulk uploads or for users who cannot use the mobile app. + +@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) From 1291defc62b9f678ffe62048ac1b1939880cbfa5 Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 02:47:25 +0100 Subject: [PATCH 09/13] feat: add reconciliation worker for stale direct-upload transfers --- app/worker/upload_reconciler/__init__.py | 0 app/worker/upload_reconciler/main.py | 60 ++++++++++++++++++++++++ makefile | 1 + 3 files changed, 61 insertions(+) create mode 100644 app/worker/upload_reconciler/__init__.py create mode 100644 app/worker/upload_reconciler/main.py diff --git a/app/worker/upload_reconciler/__init__.py b/app/worker/upload_reconciler/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/app/worker/upload_reconciler/main.py b/app/worker/upload_reconciler/main.py new file mode 100644 index 0000000..915d791 --- /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/makefile b/makefile index 36690ed..52e37c5 100644 --- a/makefile +++ b/makefile @@ -66,6 +66,7 @@ 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 & \ wait lint: From 7154235fb5220e9859a4871b95563a93376cfcfd Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 02:47:37 +0100 Subject: [PATCH 10/13] docs: clarify direct upload router comment --- app/router/staff/uploads_direct.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/app/router/staff/uploads_direct.py b/app/router/staff/uploads_direct.py index 0e9dbd1..267cd9b 100644 --- a/app/router/staff/uploads_direct.py +++ b/app/router/staff/uploads_direct.py @@ -18,7 +18,7 @@ from db.generated.models import StaffUser router = APIRouter(prefix="/uploads/direct") -# this endpoint are for staff to upload images directly to the system, bypassing the mobile app. This is useful for bulk uploads or for users who cannot use the mobile app. +# 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( From 070d7e3c6e7ba6ed40cbcfdbd8c68de1e6a89ddc Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 03:00:16 +0100 Subject: [PATCH 11/13] feat: add Google Drive write capability for syncing approved direct uploads --- app/core/config.py | 13 +++- app/infra/google_drive.py | 61 +++++++++++++++++++ app/infra/nats.py | 1 + app/service/staff_drive.py | 19 ++++++ db/generated/models.py | 2 + db/generated/photos.py | 61 ++++++++++++++++--- db/queries/photos.sql | 7 +++ .../down/add_drive_sync_fields_to_photos.sql | 3 + .../up/add_drive_sync_fields_to_photos.sql | 3 + ...06b53a9_add_drive_sync_fields_to_photos.py | 25 ++++++++ 10 files changed, 185 insertions(+), 10 deletions(-) create mode 100644 migrations/sql/down/add_drive_sync_fields_to_photos.sql create mode 100644 migrations/sql/up/add_drive_sync_fields_to_photos.sql create mode 100644 migrations/versions/af58506b53a9_add_drive_sync_fields_to_photos.py diff --git a/app/core/config.py b/app/core/config.py index 5a7fe12..9d01ffa 100644 --- a/app/core/config.py +++ b/app/core/config.py @@ -83,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 f25c4ef..05633bf 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/nats.py b/app/infra/nats.py index e17bd50..8b07e83 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/service/staff_drive.py b/app/service/staff_drive.py index f3b2ba0..b80a98b 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/db/generated/models.py b/db/generated/models.py index 2d6bcf6..7d93599 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() diff --git a/db/generated/photos.py b/db/generated/photos.py index 599a46d..a978143 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/queries/photos.sql b/db/queries/photos.sql index 3838e3c..9904ad8 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/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 0000000..d648869 --- /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_drive_sync_fields_to_photos.sql b/migrations/sql/up/add_drive_sync_fields_to_photos.sql new file mode 100644 index 0000000..2455b45 --- /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/af58506b53a9_add_drive_sync_fields_to_photos.py b/migrations/versions/af58506b53a9_add_drive_sync_fields_to_photos.py new file mode 100644 index 0000000..bb6cac9 --- /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") From 19b2879c87729dd0fac8c0fc60413e20925d9ba8 Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 03:01:27 +0100 Subject: [PATCH 12/13] feat: publish Drive sync event for direct-uploaded photos on approval --- app/service/upload_requests.py | 30 +++++++++++++++++++ tests/unit/test_direct_uploads.py | 48 ++++++++++++++++++++++++++++++- 2 files changed, 77 insertions(+), 1 deletion(-) diff --git a/app/service/upload_requests.py b/app/service/upload_requests.py index f64e346..5af3cdd 100644 --- a/app/service/upload_requests.py +++ b/app/service/upload_requests.py @@ -504,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, *, @@ -1153,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, @@ -1285,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/tests/unit/test_direct_uploads.py b/tests/unit/test_direct_uploads.py index ca73237..b1435af 100644 --- a/tests/unit/test_direct_uploads.py +++ b/tests/unit/test_direct_uploads.py @@ -1,6 +1,6 @@ import uuid from datetime import datetime, timezone -from unittest.mock import AsyncMock +from unittest.mock import AsyncMock, patch import pytest @@ -336,6 +336,52 @@ async def _photos_iter(dollar_1): 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, From 6cdc4ee3c5cdc86ce09b6e4eb50b36017b483481 Mon Sep 17 00:00:00 2001 From: wailbentafat Date: Tue, 25 Aug 2026 03:04:09 +0100 Subject: [PATCH 13/13] feat: add drive_sync worker to sync approved direct-upload photos to Drive --- app/worker/drive_sync/__init__.py | 0 app/worker/drive_sync/main.py | 103 ++++++++++++++++++++++++++++++ makefile | 1 + 3 files changed, 104 insertions(+) create mode 100644 app/worker/drive_sync/__init__.py create mode 100644 app/worker/drive_sync/main.py diff --git a/app/worker/drive_sync/__init__.py b/app/worker/drive_sync/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/app/worker/drive_sync/main.py b/app/worker/drive_sync/main.py new file mode 100644 index 0000000..f642fa1 --- /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/makefile b/makefile index 52e37c5..7e19b77 100644 --- a/makefile +++ b/makefile @@ -67,6 +67,7 @@ run-workers: 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: