```
feat: 添加仪表板路由和响应数据脱敏功能 - 添加了仪表板模块路由并集成到主路由器中 - 实现了敏感数据响应脱敏配置和功能 - 增加了 Feishu 审批卡片操作处理功能 - 支持通过任务队列异步推送日常简报和项目周报 - 添加了风险事件生成的任务队列支持 - 在 smoke 测试中增加了相关功能验证 refactor: 格式化模型注册模块导入列表 - 将单行导入列表改为多行格式以提高可读性 ```
This commit is contained in:
107
app/core/task_queue.py
Normal file
107
app/core/task_queue.py
Normal file
@@ -0,0 +1,107 @@
|
||||
from collections.abc import Callable
|
||||
from typing import Any
|
||||
|
||||
from app.core.config import get_settings
|
||||
from app.modules.feishu.constants import FeishuReceiveIdType
|
||||
|
||||
TASK_PUSH_DAILY_BRIEF = "reports.push_daily_brief"
|
||||
TASK_PUSH_PROJECT_WEEKLY = "reports.push_project_weekly"
|
||||
TASK_GENERATE_RISK_EVENTS = "risks.generate_events"
|
||||
|
||||
|
||||
def dispatch_task(
|
||||
task_name: str,
|
||||
kwargs: dict[str, Any],
|
||||
inline: Callable[[], Any],
|
||||
) -> dict[str, Any]:
|
||||
"""Dispatch a task through Celery when enabled, otherwise run inline."""
|
||||
|
||||
settings = get_settings()
|
||||
if settings.task_queue_enabled:
|
||||
from app.tasks import celery_app
|
||||
|
||||
async_result = celery_app.signature(task_name, kwargs=kwargs).apply_async()
|
||||
return {
|
||||
"queued": True,
|
||||
"mode": "celery",
|
||||
"task_name": task_name,
|
||||
"task_id": async_result.id,
|
||||
}
|
||||
return {
|
||||
"queued": False,
|
||||
"mode": "inline",
|
||||
"task_name": task_name,
|
||||
"result": inline(),
|
||||
}
|
||||
|
||||
|
||||
def enqueue_daily_brief_push(
|
||||
receive_id: str | None = None,
|
||||
receive_id_type: str = FeishuReceiveIdType.CHAT_ID,
|
||||
actor: str = "scheduler",
|
||||
) -> dict[str, Any]:
|
||||
def inline() -> Any:
|
||||
from app.core.database import SessionLocal
|
||||
from app.modules.reports.service import ReportService
|
||||
|
||||
db = SessionLocal()
|
||||
try:
|
||||
report = ReportService(db).daily_brief()
|
||||
return ReportService(db).push_report(report, receive_id, receive_id_type, actor)
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
return dispatch_task(
|
||||
TASK_PUSH_DAILY_BRIEF,
|
||||
{
|
||||
"receive_id": receive_id,
|
||||
"receive_id_type": receive_id_type,
|
||||
"actor": actor,
|
||||
},
|
||||
inline,
|
||||
)
|
||||
|
||||
|
||||
def enqueue_project_weekly_push(
|
||||
receive_id: str | None = None,
|
||||
receive_id_type: str = FeishuReceiveIdType.CHAT_ID,
|
||||
actor: str = "scheduler",
|
||||
) -> dict[str, Any]:
|
||||
def inline() -> Any:
|
||||
from app.core.database import SessionLocal
|
||||
from app.modules.reports.service import ReportService
|
||||
|
||||
db = SessionLocal()
|
||||
try:
|
||||
report = ReportService(db).project_weekly()
|
||||
return ReportService(db).push_report(report, receive_id, receive_id_type, actor)
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
return dispatch_task(
|
||||
TASK_PUSH_PROJECT_WEEKLY,
|
||||
{
|
||||
"receive_id": receive_id,
|
||||
"receive_id_type": receive_id_type,
|
||||
"actor": actor,
|
||||
},
|
||||
inline,
|
||||
)
|
||||
|
||||
|
||||
def enqueue_risk_event_generation(actor: str = "scheduler") -> dict[str, Any]:
|
||||
def inline() -> Any:
|
||||
from app.core.database import SessionLocal
|
||||
from app.modules.risk.service import RiskService
|
||||
|
||||
db = SessionLocal()
|
||||
try:
|
||||
return RiskService(db).generate_events(actor=actor)
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
return dispatch_task(
|
||||
TASK_GENERATE_RISK_EVENTS,
|
||||
{"actor": actor},
|
||||
inline,
|
||||
)
|
||||
Reference in New Issue
Block a user