- 修复 scheduler store 末尾多余逗号 - 将智能体名称从“语析”改为“Kris” - 新增获取运行记录UID的权限校验接口 - 重构操作日志写入逻辑,使用独立会话避免事务污染 - 注册外部系统调度处理器 - 优化全局错误处理器,统一响应格式与序列化处理 - 新增调度任务运行日志的handler_name字段与索引 - 重构任务执行函数,新增payload覆盖与操作人审计参数 - 重写外部系统store,新增侧边栏折叠状态与持久化 - 优化会话、消息等模型的索引、约束与字段类型 - 新增多项配置项与环境变量支持 - 重构渠道相关模型,新增索引、约束与字段优化 - 新增渠道插件模型与内容审核模型的完善
182 lines
7.3 KiB
Python
182 lines
7.3 KiB
Python
"""scheduler 的 ARQ 桥接层。
|
||
|
||
本模块是 ARQ 队列协议(模块级函数 + ``ctx``)与 ``SchedulerService``(类实例方法)
|
||
之间的适配器。ARQ 只认模块级函数,无法直接调用持有依赖的 ``SchedulerService`` 实例;
|
||
本模块在 ARQ 函数内部从 ``ctx`` 取出 worker 进程级 ``HandlerRegistry`` 与 ``arq_pool``,
|
||
按请求创建 ``SchedulerService`` 实例,转调其 ``tick`` / ``execute_task`` 方法。
|
||
|
||
对外暴露:
|
||
|
||
- ``run_scheduler_tick(ctx)``:ARQ cron 入口,周期性扫描到期任务并入队执行。
|
||
由 ``WorkerSettings.cron_jobs`` 注册,默认每分钟触发(对齐
|
||
``config.scheduler_tick_interval_seconds`` 默认 60s)。
|
||
- ``execute_scheduled_task(ctx, task_id, run_id, triggered_by, scheduled_at=None,
|
||
payload_override=None, operator=None)``:
|
||
ARQ worker 入口,执行单个到期 / 手动触发的任务。函数名与
|
||
``SchedulerService._EXECUTE_TASK_FUNCTION`` 对齐,由 ``tick`` / ``trigger_task``
|
||
通过 ``pool.enqueue_job`` 入队。``operator`` 为触发人 uid(manual 时由 Admin API
|
||
传入,auto 时为 ``"system"``),用于写入 ``run_log.created_by`` / ``updated_by``
|
||
审计字段。
|
||
- ``register_builtin_handlers(registry)``:注册 scheduler 自带的维护型 handler
|
||
(``RunLogCleanupHandler`` / ``IdempotencyCleanupHandler``),由 worker 启动钩子调用。
|
||
|
||
依赖方向:依赖 ``scheduler/infrastructure/container``(factory)+
|
||
``scheduler/framework/handlers``(维护型 handler)+ ``storage/postgres/manager``
|
||
(db 会话)。不依赖 ``use_cases`` 的具体实现(仅通过 factory 间接装配)。
|
||
|
||
进程级生命周期:
|
||
- ``HandlerRegistry``:worker 启动时创建并填充,塞入 ``ctx['scheduler_handler_registry']``
|
||
- ``arq_pool``:ARQ worker 自动注入到 ``ctx['pool']``
|
||
- ``SchedulerService``:请求级,每次 ARQ 函数调用新建(含独立 db session)
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
from datetime import datetime
|
||
from typing import TYPE_CHECKING
|
||
|
||
from yuxi.config.app import config as app_config
|
||
from yuxi.scheduler.framework.handlers import (
|
||
IdempotencyCleanupHandler,
|
||
RunLogCleanupHandler,
|
||
)
|
||
from yuxi.scheduler.framework.runtime import HandlerRegistry
|
||
from yuxi.scheduler.infrastructure.container import create_worker_scheduler_service
|
||
from yuxi.storage.postgres.manager import pg_manager
|
||
from yuxi.utils.logging_config import logger
|
||
|
||
if TYPE_CHECKING:
|
||
from arq import ArqRedis
|
||
from arq.worker import WorkerContext
|
||
|
||
__all__ = [
|
||
"execute_scheduled_task",
|
||
"register_builtin_handlers",
|
||
"run_scheduler_tick",
|
||
]
|
||
|
||
# ctx 中存储 HandlerRegistry 的键名(与 _worker_startup 约定)
|
||
_HANDLER_REGISTRY_CTX_KEY = "scheduler_handler_registry"
|
||
|
||
|
||
def register_builtin_handlers(registry: HandlerRegistry) -> None:
|
||
"""注册 scheduler 自带的维护型 handler 到 ``HandlerRegistry``。
|
||
|
||
由 worker 启动钩子(``_worker_startup``)在创建 registry 后调用。
|
||
注册的 handler 通过 ``scheduled_tasks`` 表配置触发周期(如每天凌晨清理)。
|
||
|
||
Args:
|
||
registry: worker 进程级 handler 注册表。
|
||
"""
|
||
session_factory = pg_manager.get_async_session_context
|
||
registry.register(
|
||
RunLogCleanupHandler(
|
||
session_factory=session_factory,
|
||
default_retention_days=app_config.scheduler_run_log_retention_days,
|
||
)
|
||
)
|
||
registry.register(
|
||
IdempotencyCleanupHandler(
|
||
session_factory=session_factory,
|
||
default_retention_hours=app_config.scheduler_idempotency_retention_hours,
|
||
)
|
||
)
|
||
logger.info(
|
||
"scheduler_builtin_handlers_registered",
|
||
extra={"handlers": ["run_log_cleanup", "idempotency_cleanup"]},
|
||
)
|
||
|
||
|
||
async def run_scheduler_tick(ctx: WorkerContext) -> None:
|
||
"""ARQ cron 入口:周期性扫描到期任务并入队执行。
|
||
|
||
由 ``WorkerSettings.cron_jobs`` 注册,默认每分钟触发。内部转调
|
||
``SchedulerService.tick()``,扫描 ``next_run_at <= now`` 的任务,应用
|
||
stagger 抖动与 block_strategy 后入队 ``execute_scheduled_task``。
|
||
|
||
Args:
|
||
ctx: ARQ worker 上下文,含进程级 ``HandlerRegistry`` 与 ``arq_pool``。
|
||
"""
|
||
if not app_config.scheduler_enabled:
|
||
return
|
||
|
||
registry: HandlerRegistry | None = ctx.get(_HANDLER_REGISTRY_CTX_KEY)
|
||
if registry is None:
|
||
logger.error("scheduler_tick_handler_registry_missing")
|
||
return
|
||
|
||
arq_pool: ArqRedis | None = ctx.get("pool")
|
||
if arq_pool is None:
|
||
logger.error("scheduler_tick_arq_pool_missing")
|
||
return
|
||
|
||
async with pg_manager.get_async_session_context() as db:
|
||
service = create_worker_scheduler_service(db, handler_registry=registry, arq_pool=arq_pool)
|
||
await service.tick()
|
||
|
||
|
||
async def execute_scheduled_task(
|
||
ctx: WorkerContext,
|
||
task_id: str,
|
||
run_id: str,
|
||
triggered_by: str,
|
||
scheduled_at: str | None = None,
|
||
payload_override: dict | None = None,
|
||
operator: str | None = None,
|
||
) -> None:
|
||
"""ARQ worker 入口:执行单个到期 / 手动触发的任务。
|
||
|
||
由 ``tick``(auto 触发)或 ``trigger_task``(manual 触发)通过
|
||
``pool.enqueue_job("execute_scheduled_task", ...)`` 入队。函数名与
|
||
``SchedulerService._EXECUTE_TASK_FUNCTION`` 对齐。
|
||
|
||
Args:
|
||
ctx: ARQ worker 上下文,含进程级 ``HandlerRegistry`` 与 ``arq_pool``。
|
||
task_id: 任务标识。
|
||
run_id: 本次执行标识(由 tick 或 trigger_task 生成)。
|
||
triggered_by: 触发来源(``auto`` / ``manual``)。
|
||
scheduled_at: 本次计划执行时间的 ISO 字符串(tick 触发时为 acquire 前的
|
||
原 ``next_run_at``;手动触发时为 None)。
|
||
payload_override: 手动触发时传入的 payload 覆盖值,None 时使用任务定义中的
|
||
payload。仅 ``triggered_by='manual'`` 时有意义。
|
||
operator: 触发人 uid(manual 时为管理员 uid,auto 时为 ``"system"``),
|
||
用于写入 ``run_log.created_by`` / ``updated_by`` 审计字段。None 时
|
||
回退为 ``"system"``。
|
||
"""
|
||
registry: HandlerRegistry | None = ctx.get(_HANDLER_REGISTRY_CTX_KEY)
|
||
if registry is None:
|
||
logger.error(
|
||
"scheduler_execute_handler_registry_missing",
|
||
extra={"task_id": task_id, "run_id": run_id},
|
||
)
|
||
return
|
||
|
||
arq_pool: ArqRedis | None = ctx.get("pool")
|
||
if arq_pool is None:
|
||
logger.error(
|
||
"scheduler_execute_arq_pool_missing",
|
||
extra={"task_id": task_id, "run_id": run_id},
|
||
)
|
||
return
|
||
|
||
scheduled_at_dt: datetime | None = None
|
||
if scheduled_at is not None:
|
||
try:
|
||
scheduled_at_dt = datetime.fromisoformat(scheduled_at)
|
||
except ValueError:
|
||
logger.warning(
|
||
"scheduler_execute_invalid_scheduled_at",
|
||
extra={"task_id": task_id, "run_id": run_id, "scheduled_at": scheduled_at},
|
||
)
|
||
|
||
async with pg_manager.get_async_session_context() as db:
|
||
service = create_worker_scheduler_service(db, handler_registry=registry, arq_pool=arq_pool)
|
||
await service.execute_task(
|
||
task_id=task_id,
|
||
run_id=run_id,
|
||
triggered_by=triggered_by,
|
||
scheduled_at=scheduled_at_dt,
|
||
payload_override=payload_override,
|
||
operator=operator,
|
||
)
|