Refactor task routes and add control routes (#749)

This commit is contained in:
Wu Clan
2025-08-04 13:18:07 +08:00
committed by GitHub
parent dedf4e7bae
commit 8591d4e592
6 changed files with 66 additions and 29 deletions
+3 -1
View File
@@ -2,11 +2,13 @@
# -*- coding: utf-8 -*-
from fastapi import APIRouter
from backend.app.task.api.v1.control import router as task_control_router
from backend.app.task.api.v1.result import router as task_result_router
from backend.app.task.api.v1.scheduler import router as task_scheduler_router
from backend.core.conf import settings
v1 = APIRouter(prefix=f'{settings.FASTAPI_API_V1_PATH}/task', tags=['任务'])
v1 = APIRouter(prefix=f'{settings.FASTAPI_API_V1_PATH}/tasks', tags=['任务'])
v1.include_router(task_control_router)
v1.include_router(task_result_router, prefix='/results')
v1.include_router(task_scheduler_router, prefix='/schedulers')
+51
View File
@@ -0,0 +1,51 @@
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
from typing import Annotated
from fastapi import APIRouter, Depends, Path
from starlette.concurrency import run_in_threadpool
from backend.app.task import celery_app
from backend.app.task.schema.control import TaskRegisteredDetail
from backend.common.exception import errors
from backend.common.response.response_schema import ResponseModel, ResponseSchemaModel, response_base
from backend.common.security.jwt import DependsJwtAuth
from backend.common.security.permission import RequestPermission
from backend.common.security.rbac import DependsRBAC
router = APIRouter()
@router.get('/registered', summary='获取已注册的任务', dependencies=[DependsJwtAuth])
async def get_task_registered() -> ResponseSchemaModel[list[TaskRegisteredDetail]]:
inspector = celery_app.control.inspect(timeout=0.5)
registered = await run_in_threadpool(inspector.registered)
if not registered:
raise errors.ServerError(msg='Celery Worker 暂不可用,请稍后重试')
task_registered = []
celery_app_tasks = celery_app.tasks
for _, tasks in registered.items():
for task in tasks:
task_ins = celery_app_tasks.get(task)
if task_ins:
task_doc = task_ins.__doc__
task_registered.append({'name': task_doc or task_ins, 'task': task_ins})
else:
task_registered.append({'name': task, 'task': task})
return response_base.success(data=task_registered)
@router.delete(
'/{task_id}/cancel',
summary='撤销任务',
dependencies=[
Depends(RequestPermission('sys:task:revoke')),
DependsRBAC,
],
)
async def revoke_task(task_id: Annotated[str, Path(description='任务 UUID')]) -> ResponseModel:
workers = await run_in_threadpool(celery_app.control.ping, timeout=0.5)
if not workers:
raise errors.ServerError(msg='Celery Worker 暂不可用,请稍后重试')
celery_app.control.revoke(task_id)
return response_base.success()
-13
View File
@@ -119,16 +119,3 @@ async def delete_task_scheduler(pk: Annotated[int, Path(description='任务调
async def execute_task(pk: Annotated[int, Path(description='任务调度 ID')]) -> ResponseModel:
await task_scheduler_service.execute(pk=pk)
return response_base.success()
@router.delete(
'/{task_id}/cancel',
summary='撤销任务',
dependencies=[
Depends(RequestPermission('sys:task:revoke')),
DependsRBAC,
],
)
async def revoke_task(task_id: Annotated[str, Path(description='任务 UUID')]) -> ResponseModel:
await task_scheduler_service.revoke(task_id=task_id)
return response_base.success()
+4 -2
View File
@@ -13,7 +13,8 @@ from backend.core.path_conf import BASE_PATH
def find_task_packages():
packages = []
for root, dirs, files in os.walk(os.path.join(BASE_PATH, 'app', 'task', 'tasks')):
task_dir = os.path.join(BASE_PATH, 'app', 'task', 'tasks')
for root, dirs, files in os.walk(task_dir):
if 'tasks.py' in files:
package = root.replace(str(BASE_PATH.parent) + os.path.sep, '').replace(os.path.sep, '.')
packages.append(package)
@@ -54,7 +55,8 @@ def init_celery() -> celery.Celery:
)
# 自动发现任务
app.autodiscover_tasks(find_task_packages())
packages = find_task_packages()
app.autodiscover_tasks(packages)
return app
+8
View File
@@ -0,0 +1,8 @@
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
from backend.common.schema import SchemaBase
class TaskRegisteredDetail(SchemaBase):
name: str
task: str
@@ -142,18 +142,5 @@ class TaskSchedulerService:
else:
celery_app.send_task(name=task_scheduler.task, args=args, kwargs=kwargs)
@staticmethod
async def revoke(*, task_id: str) -> None:
"""
撤销指定的任务
:param task_id: 任务 UUID
:return:
"""
workers = await run_in_threadpool(celery_app.control.ping, timeout=0.5)
if not workers:
raise errors.ServerError(msg='Celery Worker 暂不可用,请稍后重试')
celery_app.control.revoke(task_id)
task_scheduler_service: TaskSchedulerService = TaskSchedulerService()