Skip to content

Commit d98f41c

Browse files
committed
feat: offload event deletion to asynchronous celery task & update docs
- Refactored `DeleteEventTableUseCase` to trigger a background Celery task (`delete_event_data_task`) rather than blocking the HTTP response, preventing frontend serverless timeouts. - Added `event_id` parameter to the `DELETE /api/events/delete-event-table/{event_code}` endpoint to allow the backend to clean up associated zip directories. - Implemented `delete_folder` in the MinIO storage adapter using `list_objects_v2` pagination to recursively wipe S3 directories. - Updated `ITaskQueueService` and `IStorageService` ports to reflect new deletion capabilities. - Wired the new queue dependency to the deletion use case in `di_container.py`. - docs: Completely overhauled `README.md`, `main_api/README.md`, `inference_api/README.md`, and `workflows.md` to reflect the GCP VM + PostgreSQL deployment (removing outdated mentions of Hugging Face, Redis, MinIO, and dynamic tables). - docs: Added missing API schemas to `main_api/README.md` (Generate ZIP, Check ZIP, Delete Event Table) and fixed request payload typos.
1 parent c4a3745 commit d98f41c

12 files changed

Lines changed: 144 additions & 42 deletions

File tree

‎README.md‎

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
<div align="center">
22
<h1>📸 Eventsnap API</h1>
33
<i>A High-Performance Asynchronous Facial Recognition Pipeline</i><br>
4-
<i>Powered by FastAPI, Celery, PostgreSQL pgvector, and Hugging Face</i>
4+
<i>Powered by FastAPI, Celery, PostgreSQL pgvector, and ONNX Runtime</i>
55
</div>
66

77
---
@@ -13,7 +13,7 @@ Eventsnap is a distributed, horizontally scalable microservice architecture desi
1313
It is split into two main components (each with their own dedicated README files):
1414

1515
### 1. `main_api` (The Orchestrator)
16-
A FastAPI server that acts as the entry point. It accepts requests, authenticates them, dynamically creates Postgres tables, and dumps background encoding tasks into RabbitMQ for Celery workers to pick up.
16+
A FastAPI server that acts as the entry point. It accepts requests, authenticates them, saves face embeddings to Postgres, and dumps background encoding tasks into RabbitMQ for Celery workers to pick up.
1717

1818
### 2. `inference_api` (The GPU Worker)
1919
A strictly mathematical, stateless ONNX Runtime container. It receives Base64 encoded photos, runs the powerful `insightface` SCRFD and ArcFace models on the NVIDIA GPU, and returns precise bounding boxes and 512-dimension `glintr100` embeddings.
@@ -25,9 +25,9 @@ A strictly mathematical, stateless ONNX Runtime container. It receives Base64 en
2525
* **API Framework:** FastAPI (Python 3.14, Native AsyncIO)
2626
* **Architecture Pattern:** Hexagonal Architecture (Ports and Adapters) with Dependency Injection
2727
* **Package Manager:** uv (Ultra-fast Python package installer)
28-
* **Background Tasks:** Celery + RabbitMQ (Broker) + Redis (Result Backend)
28+
* **Background Tasks:** Celery + RabbitMQ (Broker) + PostgreSQL (Result Backend)
2929
* **Database:** PostgreSQL + [`pgvector`](https://github.com/pgvector/pgvector) extension (Cosine Similarity matching)
30-
* **Object Storage:** MinIO (S3 Compatible)
30+
* **Object Storage:** Storage Bucket (S3 Compatible)
3131
* **Machine Learning:** ONNX runtime (CUDA 11.8), InsightFace
3232
* **Containerization:** Docker & Docker Compose
3333

@@ -49,7 +49,7 @@ docker build -t main_api:dev ./main_api
4949
```
5050

5151
### 2. Spin Up the Stack
52-
Bring up all the containers (Postgres DB, RabbitMQ, MinIO, Inference API, Main API, and Celery Worker). The orchestrated services will automatically wait for their database dependencies to become healthy before starting.
52+
Bring up all the containers (Postgres DB, RabbitMQ, Storage Bucket, Inference API, Main API, and Celery Worker). The orchestrated services will automatically wait for their database dependencies to become healthy before starting.
5353

5454
```bash
5555
docker compose up -d
@@ -65,8 +65,8 @@ docker compose up -d
6565

6666
*(For detailed sequence diagrams of the complete asynchronous system, see [workflows.md](./workflows.md))*
6767

68-
1. A user uploads a ZIP of an event directly via the Next.js frontend, which extracts and pushes the images into **MinIO**.
68+
1. A user uploads a ZIP of an event directly via the Next.js frontend, which extracts and pushes the images into **Storage Bucket**.
6969
2. The frontend hits the **Main API** `/encode-event/` endpoint, passing the `event_code` in the JSON payload.
7070
3. The **Main API** creates a Celery Task and immediately returns a `task_id` so the user isn't stuck waiting.
71-
4. The background **Celery Worker** picks up the task, pre-fetches images from **MinIO** using an aggressive 64-connection pool, beams them (Base64) to the **Inference API**, and bulk-inserts the generated 512D vectors directly into **PostgreSQL**.
71+
4. The background **Celery Worker** picks up the task, pre-fetches images from **Storage Bucket** using an aggressive 64-connection pool, beams them (Base64) to the **Inference API**, and bulk-inserts the generated 512D vectors directly into **PostgreSQL**.
7272
5. An attendee hits the **Main API** `/sort-attendee/` endpoint with their selfies and the `event_code`. The orchestrator gets the embeddings for those selfies, averages them, and executes a sub-millisecond `<=>` cosine similarity search in `pgvector` to find all photos they appear in!

‎inference_api/README.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22

33
The `inference_api` is an ultra-fast, stateless FastAPI server designed purely for mathematical facial processing. It uses the `insightface` library powered by `onnxruntime-gpu` (CUDA 11) to extract 512-dimension vector embeddings from raw image data.
44

5-
It is designed to be horizontally scaled or deployed as a Serverless Endpoint on Hugging Face (e.g., Nvidia T4 hardware).
5+
It is designed to be horizontally scaled or deployed on a GCP VM (e.g., Nvidia T4 hardware).
66

77
## Core Responsibilities
88
* **Vectorization**: Converts human faces into mathematical matrices (ResNet100 model).

‎main_api/README.md‎

Lines changed: 58 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,7 @@ This API strictly follows the **Ports and Adapters (Hexagonal) Architecture** us
1515

1616
* **`domain/`**: Enterprise logic, entities, and domain exceptions. Completely independent of any external frameworks.
1717
* **`application/`**: Business use cases (e.g., `events.py`, `attendees.py`, `background_tasks.py`), DTOs, and Ports (interfaces).
18-
* **`infrastructure/`**: External adapters (Celery, PostgreSQL/pgvector, MinIO, HuggingFace Inference API, Image Augmentation).
18+
* **`infrastructure/`**: External adapters (Celery, PostgreSQL/pgvector, Storage Bucket, Internal Inference API, Image Augmentation).
1919
* **`presentation/`**: FastAPI routers, schemas, and exception handlers.
2020

2121
---
@@ -28,7 +28,7 @@ Content-Type: application/json
2828
```
2929

3030
### 1. Encode Event (Asynchronous)
31-
Triggers the background Celery worker to download a folder from MinIO, process every image through the GPU inference API, and save the 512-dimension vector embeddings to Postgres.
31+
Triggers the background Celery worker to download a folder from Storage Bucket, process every image through the GPU inference API, and save the 512-dimension vector embeddings to Postgres.
3232

3333
**Example `curl` Request:**
3434
```bash
@@ -38,8 +38,8 @@ curl -X 'POST' \
3838
-d '{
3939
"event_code": "DCAYTI",
4040
"max_faces": 0,
41-
"det_conf": 0.5,
42-
"nms_thresh": 0.4
41+
"detection_conf": 0.5,
42+
"nms_threshold": 0.4
4343
}'
4444
```
4545

@@ -50,8 +50,8 @@ import axios from 'axios';
5050
const response = await axios.post('http://localhost:8000/api/events/encode-event/', {
5151
event_code: 'DCAYTI',
5252
max_faces: 0,
53-
det_conf: 0.5,
54-
nms_thresh: 0.4
53+
detection_conf: 0.5,
54+
nms_threshold: 0.4
5555
}, {
5656
headers: {
5757
'Content-Type': 'application/json'
@@ -200,6 +200,58 @@ console.log(response.data);
200200
}
201201
```
202202

203+
### 5. Generate ZIP
204+
Generates a ZIP archive containing the matched photos for an attendee.
205+
206+
**Example Request:**
207+
```json
208+
{
209+
"event_id": "uuid",
210+
"user_id": "uuid",
211+
"image_paths": [{}]
212+
}
213+
```
214+
**Response:**
215+
```json
216+
{
217+
"success": true,
218+
"task_id": "task-uuid",
219+
"message": "message"
220+
}
221+
```
222+
223+
### 6. Check ZIP Status
224+
Checks if a generated zip is available in the storage bucket.
225+
226+
**Example Request:**
227+
```bash
228+
curl -X 'GET' 'http://localhost:8000/api/attendees/check-zip/{event_id}/{user_id}'
229+
```
230+
**Response:**
231+
```json
232+
{
233+
"exists": true,
234+
"zip_path": "zip/event_id/user_id.zip",
235+
"filename": "user_id.zip"
236+
}
237+
```
238+
239+
### 7. Delete Event Table
240+
Safely clears event table from DB and ML storage in the background.
241+
242+
**Example Request:**
243+
```bash
244+
curl -X 'DELETE' 'http://localhost:8000/api/events/delete-event-table/{event_code}'
245+
```
246+
**Response:**
247+
```json
248+
{
249+
"success": true,
250+
"message": "message",
251+
"table_name": "event_encodings"
252+
}
253+
```
254+
203255
## Running the Service
204256
The service is booted automatically via `docker-compose up`, but can be tested locally using:
205257
```bash

‎main_api/application/ports/queue.py‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,3 +32,7 @@ def enqueue_create_zip(
3232
def get_task_status(self, task_id: str) -> Dict[str, Any]:
3333
"""Returns dict with state, info, result"""
3434
pass
35+
36+
def enqueue_delete_event(self, event_code: str, event_id: str | None = None) -> str:
37+
"""Returns the task ID"""
38+
pass

‎main_api/application/ports/storage.py‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,3 +19,6 @@ async def create_zip_from_images(
1919

2020
async def check_zip_exists(self, zip_key: str) -> bool:
2121
pass
22+
23+
async def delete_folder(self, prefix: str) -> None:
24+
pass

‎main_api/application/use_cases/events.py‎

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -41,14 +41,14 @@ async def execute(self, event_code: str) -> EncodedCountDTO:
4141

4242

4343
class DeleteEventTableUseCase:
44-
def __init__(self, uow: IUnitOfWork):
45-
self.uow = uow
44+
def __init__(self, queue_service: ITaskQueueService):
45+
self.queue_service = queue_service
4646

47-
async def execute(self, event_code: str) -> DeleteTableDTO:
48-
async with self.uow as uow:
49-
await uow.event_repo.delete_event_data(event_code)
50-
await uow.commit()
47+
async def execute(
48+
self, event_code: str, event_id: str | None = None
49+
) -> DeleteTableDTO:
50+
self.queue_service.enqueue_delete_event(event_code, event_id)
5151
return DeleteTableDTO(
5252
success=True,
53-
message=f"Data for event '{event_code}' deleted successfully if it existed.",
53+
message=f"Enqueued deletion task for event '{event_code}'.",
5454
)

‎main_api/infrastructure/di_container.py‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -84,7 +84,9 @@ class Container(containers.DeclarativeContainer):
8484

8585
get_encoded_count_use_case = providers.Factory(GetEncodedCountUseCase, uow=uow)
8686

87-
delete_event_table_use_case = providers.Factory(DeleteEventTableUseCase, uow=uow)
87+
delete_event_table_use_case = providers.Factory(
88+
DeleteEventTableUseCase, queue_service=queue_service
89+
)
8890

8991
encode_attendee_use_case = providers.Factory(
9092
EncodeAttendeeUseCase,

‎main_api/infrastructure/queue/celery_service.py‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,12 @@ def enqueue_create_zip(
4242
task = create_event_zip_task.delay(event_id, user_id, image_paths)
4343
return task.id
4444

45+
def enqueue_delete_event(self, event_code: str, event_id: str | None = None) -> str:
46+
from infrastructure.queue.celery_workers import delete_event_data_task
47+
48+
task = delete_event_data_task.delay(event_code, event_id)
49+
return task.id
50+
4551
def get_task_status(self, task_id: str) -> Dict[str, Any]:
4652
res = AsyncResult(task_id, app=celery_app)
4753

‎main_api/infrastructure/queue/celery_workers.py‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -88,3 +88,21 @@ def update_state_cb(state_name, meta_dict):
8888
if dataclasses.is_dataclass(result):
8989
return dataclasses.asdict(result)
9090
return result
91+
92+
93+
@shared_task(bind=True, name="delete_event_data_task", acks_late=True)
94+
def delete_event_data_task(self, event_code: str, event_id: str | None = None):
95+
container = get_container()
96+
uow = container.uow()
97+
storage = container.storage_service()
98+
99+
async def _delete():
100+
async with uow:
101+
await uow.event_repo.delete_event_data(event_code)
102+
await uow.commit()
103+
await storage.delete_folder(f"event/{event_code}/")
104+
if event_id:
105+
await storage.delete_folder(f"zip/{event_id}/")
106+
107+
asyncio.run(_delete())
108+
return {"success": True, "message": f"Deleted event {event_code}"}

‎main_api/infrastructure/storage/minio_service.py‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -142,3 +142,20 @@ async def check_zip_exists(self, zip_key: str) -> bool:
142142
return True
143143
except Exception:
144144
return False
145+
146+
async def delete_folder(self, prefix: str) -> None:
147+
session = self._get_session()
148+
async with session.client("s3", endpoint_url=self.endpoint_url) as s3:
149+
paginator = s3.get_paginator("list_objects_v2")
150+
async for page in paginator.paginate(
151+
Bucket=self.bucket_name, Prefix=prefix
152+
):
153+
if "Contents" in page:
154+
objects_to_delete = [
155+
{"Key": obj["Key"]} for obj in page["Contents"]
156+
]
157+
if objects_to_delete:
158+
await s3.delete_objects(
159+
Bucket=self.bucket_name,
160+
Delete={"Objects": objects_to_delete, "Quiet": True},
161+
)

0 commit comments

Comments
 (0)