"""定时调度(工作台侧,dev 默认关):每日组合再平衡(遍历 signed 客户)。 架构 §7 已选型 APScheduler;仅当 config.advisor.scheduler_enabled=true 时启动。 dev 常以 uvicorn --reload 运行(重载会重复调度),故默认关闭,生产单进程开启。 边界:调度器无需前端/无需请求上下文,需自造投顾 JWT(service.auth.create_token)与 trace_id(new_request_id)调用 Agent;大额申赎扫描/回访到期/风评到期等本地定时任务 P1 后续补(依赖 fin_transaction 模型等)。 """ from __future__ import annotations import logging from common_const import CRON_PORTFOLIO_REBALANCE, CUSTOMER_REL_STATUS_SIGNED from config.settings import settings from repositories.customer_relation import CustomerRelationRepo logger = logging.getLogger("service.advisor.scheduler") class AdvisorScheduler: def __init__(self): self._scheduler = None def start(self) -> None: if not settings.advisor.scheduler_enabled: return # 懒加载:未安装 apscheduler 时不阻塞应用启动(默认关闭本调度器) from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.cron import CronTrigger self._scheduler = AsyncIOScheduler() self._scheduler.add_job( self.daily_rebalance, CronTrigger.from_crontab(CRON_PORTFOLIO_REBALANCE), id="advisor_daily_rebalance", name="每日组合再平衡(遍历 signed 客户)", ) self._scheduler.start() logger.info("advisor scheduler started (cron=%s)", CRON_PORTFOLIO_REBALANCE) def shutdown(self) -> None: if self._scheduler is not None: self._scheduler.shutdown(wait=False) self._scheduler = None async def daily_rebalance(self) -> None: """每日遍历已签约客户,按投顾逐个调 Agent rebalance/run(异步受理)。 结果靠事件 event:rebalance_draft_created 回执 → 工作台消费生成待办。 """ from config.database.mysql import get_session_factory from service.advisor.agent_client import get_agent_client from service.auth import create_token from utils.request_id import new_request_id client = get_agent_client() if not client.configured: logger.warning("投顾Agent 未配置,跳过每日 rebalance") return async with get_session_factory()() as session: relations = await CustomerRelationRepo(session).list_by_status( CUSTOMER_REL_STATUS_SIGNED ) by_advisor: dict[int, list[int]] = {} for rel in relations: by_advisor.setdefault(rel.advisor_id, []).append(rel.customer_id) for advisor_id, customer_ids in by_advisor.items(): # 后台任务无请求上下文,自造投顾 JWT + trace_id token = create_token(int(advisor_id)) auth = f"Bearer {token}" for customer_id in customer_ids: try: await client.rebalance_run( customer_id, auth_header=auth, trace_id=new_request_id() ) except Exception: logger.exception( "rebalance run failed advisor=%s customer=%s", advisor_id, customer_id )