|
1 | 1 | import logging |
| 2 | +from contextlib import contextmanager |
2 | 3 | from typing import Annotated, Any |
3 | 4 |
|
| 5 | +from celery_library import task |
| 6 | +from celery_library.errors import TaskNotFoundError |
4 | 7 | from common_library.error_codes import create_error_code |
5 | 8 | from common_library.logging.logging_errors import create_troubleshooting_log_kwargs |
6 | 9 | from fastapi import APIRouter, Depends, FastAPI, HTTPException, status |
|
18 | 21 | from servicelib.celery.models import TaskState, TaskUUID |
19 | 22 | from servicelib.fastapi.dependencies import get_app |
20 | 23 |
|
| 24 | +from ...exceptions.backend_errors import CeleryTaskNotFoundError |
21 | 25 | from ...models.domain.celery_models import ( |
22 | 26 | ApiServerOwnerMetadata, |
23 | 27 | ) |
|
42 | 46 | } |
43 | 47 |
|
44 | 48 |
|
| 49 | +@contextmanager |
| 50 | +def _exception_mapper(task_uuid: TaskUUID): |
| 51 | + try: |
| 52 | + yield |
| 53 | + except TaskNotFoundError as exc: |
| 54 | + raise CeleryTaskNotFoundError(task_uuid=task_uuid) from exc |
| 55 | + |
| 56 | + |
45 | 57 | @router.get( |
46 | 58 | "", |
47 | 59 | response_model=ApiServerEnvelope[list[TaskGet]], |
@@ -109,10 +121,11 @@ async def get_task_status( |
109 | 121 | user_id=user_id, |
110 | 122 | product_name=product_name, |
111 | 123 | ) |
112 | | - task_status = await task_manager.get_task_status( |
113 | | - owner_metadata=owner_metadata, |
114 | | - task_uuid=TaskUUID(f"{task_uuid}"), |
115 | | - ) |
| 124 | + with _exception_mapper(task_uuid=task_uuid): |
| 125 | + task_status = await task_manager.get_task_status( |
| 126 | + owner_metadata=owner_metadata, |
| 127 | + task_uuid=TaskUUID(f"{task_uuid}"), |
| 128 | + ) |
116 | 129 |
|
117 | 130 | return TaskStatus( |
118 | 131 | task_progress=TaskProgress( |
@@ -147,10 +160,11 @@ async def cancel_task( |
147 | 160 | user_id=user_id, |
148 | 161 | product_name=product_name, |
149 | 162 | ) |
150 | | - await task_manager.cancel_task( |
151 | | - owner_metadata=owner_metadata, |
152 | | - task_uuid=TaskUUID(f"{task_uuid}"), |
153 | | - ) |
| 163 | + with _exception_mapper(task_uuid=task_uuid): |
| 164 | + await task_manager.cancel_task( |
| 165 | + owner_metadata=owner_metadata, |
| 166 | + task_uuid=TaskUUID(f"{task_uuid}"), |
| 167 | + ) |
154 | 168 |
|
155 | 169 |
|
156 | 170 | @router.get( |
@@ -183,37 +197,38 @@ async def get_task_result( |
183 | 197 | product_name=product_name, |
184 | 198 | ) |
185 | 199 |
|
186 | | - task_status = await task_manager.get_task_status( |
187 | | - owner_metadata=owner_metadata, |
188 | | - task_uuid=TaskUUID(f"{task_uuid}"), |
189 | | - ) |
190 | | - |
191 | | - if not task_status.is_done: |
192 | | - raise HTTPException( |
193 | | - status_code=status.HTTP_404_NOT_FOUND, |
194 | | - detail="Task result not available yet", |
| 200 | + with _exception_mapper(task_uuid=task_uuid): |
| 201 | + task_status = await task_manager.get_task_status( |
| 202 | + owner_metadata=owner_metadata, |
| 203 | + task_uuid=TaskUUID(f"{task_uuid}"), |
195 | 204 | ) |
196 | 205 |
|
197 | | - task_result = await task_manager.get_task_result( |
198 | | - owner_metadata=owner_metadata, |
199 | | - task_uuid=TaskUUID(f"{task_uuid}"), |
200 | | - ) |
201 | | - |
202 | | - if task_status.task_state == TaskState.FAILURE: |
203 | | - assert isinstance(task_result, Exception) |
204 | | - user_error_msg = f"The execution of task {task_uuid} failed" |
205 | | - support_id = create_error_code(task_result) |
206 | | - _logger.exception( |
207 | | - **create_troubleshooting_log_kwargs( |
208 | | - user_error_msg, |
209 | | - error=task_result, |
210 | | - error_code=support_id, |
211 | | - tip="Unexpected error in Celery", |
| 206 | + if not task_status.is_done: |
| 207 | + raise HTTPException( |
| 208 | + status_code=status.HTTP_404_NOT_FOUND, |
| 209 | + detail="Task result not available yet", |
212 | 210 | ) |
| 211 | + |
| 212 | + task_result = await task_manager.get_task_result( |
| 213 | + owner_metadata=owner_metadata, |
| 214 | + task_uuid=TaskUUID(f"{task_uuid}"), |
213 | 215 | ) |
214 | | - raise HTTPException( |
215 | | - status_code=status.HTTP_503_SERVICE_UNAVAILABLE, |
216 | | - detail=user_error_msg, |
217 | | - ) |
| 216 | + |
| 217 | + if task_status.task_state == TaskState.FAILURE: |
| 218 | + assert isinstance(task_result, Exception) |
| 219 | + user_error_msg = f"The execution of task {task_uuid} failed" |
| 220 | + support_id = create_error_code(task_result) |
| 221 | + _logger.exception( |
| 222 | + **create_troubleshooting_log_kwargs( |
| 223 | + user_error_msg, |
| 224 | + error=task_result, |
| 225 | + error_code=support_id, |
| 226 | + tip="Unexpected error in Celery", |
| 227 | + ) |
| 228 | + ) |
| 229 | + raise HTTPException( |
| 230 | + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, |
| 231 | + detail=user_error_msg, |
| 232 | + ) |
218 | 233 |
|
219 | 234 | return TaskResult(result=task_result, error=None) |
0 commit comments