Files
company-ai-platform/app/tasks.py
JiuContinent 0cda45238a ```
feat: 添加AI记忆模块和事件调度系统

- 新增AI记忆模块,支持本地记忆召回和自动写入功能
- 实现事件调度系统,支持批量处理待定事件和重试机制
- 集成心跳监控机制,跟踪API、调度器和工作节点状态
- 扩展仪表板数据统计,包含AI记忆条目和心跳概要
- 添加企业运营分析报告功能,提供财务、采购等多维度分析
- 更新配置设置,增加事件调度和AI记忆相关参数
- 优化任务队列,添加事件分发任务类型
- 扩展审计日志,记录AI记忆操作和事件调度行为
- 实现领域事件模型,支持事件持久化和状态管理
- 添加观察性服务,监控系统组件健康状况
```
2026-07-09 17:26:19 +08:00

150 lines
4.2 KiB
Python

from typing import Any
from socket import gethostname
from celery import Celery
from app.core.config import get_settings
from app.core.constants import ActorValue
from app.core.database import SessionLocal
from app.modules.feishu.constants import FeishuReceiveIdType
settings = get_settings()
celery_app = Celery(
"company_ai_platform",
broker=settings.redis_url,
backend=settings.celery_result_backend_url or settings.redis_url,
)
celery_app.conf.task_always_eager = settings.task_queue_always_eager
@celery_app.task(name="reports.push_daily_brief")
def push_daily_brief(
receive_id: str | None = None,
receive_id_type: str = FeishuReceiveIdType.CHAT_ID,
actor: str = ActorValue.SCHEDULER,
push_run_code: str | None = None,
) -> dict[str, Any]:
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,
push_run_code=push_run_code,
)
finally:
db.close()
@celery_app.task(name="events.dispatch_pending")
def dispatch_pending_events(
limit: int | None = None,
actor: str = ActorValue.WORKER,
) -> list[dict[str, Any]]:
from app.modules.events.service import EventService
from app.modules.observability.constants import HeartbeatComponent
from app.modules.observability.service import ObservabilityService
db = SessionLocal()
try:
ObservabilityService(db).record_heartbeat(
component=HeartbeatComponent.WORKER,
instance_id=gethostname(),
actor=actor,
)
return EventService(db).dispatch_pending(
limit=limit or settings.event_dispatch_batch_size,
worker_id=f"{actor}:{gethostname()}",
)
finally:
db.close()
@celery_app.task(name="reports.push_project_weekly")
def push_project_weekly(
receive_id: str | None = None,
receive_id_type: str = FeishuReceiveIdType.CHAT_ID,
actor: str = ActorValue.SCHEDULER,
push_run_code: str | None = None,
) -> dict[str, Any]:
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,
push_run_code=push_run_code,
)
finally:
db.close()
@celery_app.task(name="risks.generate_events")
def generate_risk_events(actor: str = ActorValue.SCHEDULER) -> dict[str, Any]:
from app.modules.risk.service import RiskService
db = SessionLocal()
try:
return RiskService(db).generate_events(actor=actor)
finally:
db.close()
@celery_app.task(name="legacy.sync_projects")
def sync_legacy_projects(
source_query: str | None = None,
source_query_name: str | None = None,
field_map: dict[str, str] | None = None,
limit: int = 100,
dry_run: bool = False,
actor: str = ActorValue.SCHEDULER,
) -> dict[str, Any]:
from app.modules.legacy_mysql.service import LegacyMySQLService
db = SessionLocal()
try:
return LegacyMySQLService(db).sync_projects(
source_query=source_query,
source_query_name=source_query_name,
field_map=field_map or {},
limit=limit,
dry_run=dry_run,
actor=actor,
)
finally:
db.close()
@celery_app.task(name="legacy.sync_tasks")
def sync_legacy_tasks(
source_query: str | None = None,
source_query_name: str | None = None,
field_map: dict[str, str] | None = None,
limit: int = 100,
dry_run: bool = False,
actor: str = ActorValue.SCHEDULER,
) -> dict[str, Any]:
from app.modules.legacy_mysql.service import LegacyMySQLService
db = SessionLocal()
try:
return LegacyMySQLService(db).sync_tasks(
source_query=source_query,
source_query_name=source_query_name,
field_map=field_map or {},
limit=limit,
dry_run=dry_run,
actor=actor,
)
finally:
db.close()