5个APScheduler定时作业的独立会话隔离机制 #201

Open
opened 2026-08-21 06:14:20 +00:00 by KumaAgent · 0 comments
Member

Parent

#151(ADR-0013 ①层单一事务边界落地——请求级事务,issue #108 Q1)。事务边界方向已确认为分段
事务只包住连续无外部 I/O 打断的一段纯 DB 读写,下一步要发起任何外部调用前当前段必须先提交/关闭
(见 #151 comments"拆片前的事务边界设计方向已与用户确认")。

本片负责 5 个 APScheduler 定时作业_callback_scanentry_callback.py:35)、_drive_tick
entry_drive_tick.py:265)、_compress_pending_delta/_run_reflection/_run_unanswered_escalation
(均在 entry_offline_jobs.py,分别在 30/96/145 行附近)。这 5 处已核实确实全部是
@scheduler.scheduled_job 直接装饰的顶层协程。

What to build

design-verify 阶段 [HIGH] 发现:get_scoped_session()(第三方库 nonebot_plugin_orm 提供,本片
起草时核实其为外部依赖、非本仓代码,其内部 scope key 的精确实现细节以原始发现为准,未重新逐行核实
该库源码)的缓存 key 依赖 nonebot 的 current_event/current_matcher 上下文变量,对这 5 个跑在
事件/matcher 上下文之外的顶层协程,会全部坍缩成同一个 key,即它们实际共享同一个 Python
Session 对象。现状因为每次调用都是"开-用-commit-close"的极短窗口所以风险很小;但"分段"设计
要求把持有窗口拉长到跨越真实 LLM 调用的单个 item 处理过程(秒级)。APScheduler 默认配置不阻止
5 个不同 job 之间并发重叠执行——一旦重叠,会在同一个不支持并发使用的 AsyncSession 对象上互相
踩踏,存在真实的会话状态损坏风险。

这 5 个作业需要一套完全独立于 nonebot scoped-session 机制的会话构造方式,不能直接复用反应式
入口那一套依赖 nonebot 上下文变量的机制。

Acceptance criteria

一、独立会话构造机制

  • 设计并落地一套不依赖 current_event/current_matcher 的 session/segment 构造方式——建议
    方向:每次 job 内部按处理单元(_callback_scan 每条 callback、_drive_tick 每个 chat、
    _compress_pending_delta 每个 chat、_run_reflection 每条 attempt、_run_unanswered_escalation
    每条候选)各自独立构造一个新的 DB 连接/session,不复用 get_scoped_session() 的缓存机制,具体
    API 形状与「地基片」协调(begin_segment() 是否需要一个不依赖 nonebot 上下文的变体,或
    PersistentStorage 内部提供另一条不经过 get_scoped_session() 的构造路径——见地基片"五、
    APScheduler 5 个作业的独立会话构造方式不在本票范围"那条留下的接口空间)
  • 补一条测试直接验证并发场景:模拟两个 job(或同一 job 的两次调用)在真实存在时间重叠的情况下
    各自开启并使用自己的 session,断言互不干扰(不是简单断言"代码跑通",要能证明两次调用确实使用了
    不同的底层连接/session 对象,或者构造一个如果共享会暴露出来的具体错误场景,如两者交替写入不同
    记录并断言写入结果符合各自预期而非互相覆盖)

原文(design-verify HIGH 发现,issue #151 comments):

get_scoped_session() 的 scope key (id(current_event.get(None)), current_matcher.get(None))
对这 5 个作业全部坍缩成同一个 (id(None), None)……APScheduler 默认配置不会阻止 5 个不同 job
之间的并发重叠执行——一旦重叠,会在同一个不支持并发使用的 AsyncSession 对象上互相踩踏,存在
真实的会话状态损坏风险。

二、5 个作业各自按"两次外部 I/O 之间"分段(复用"一"节交付的独立会话机制)

  • _callback_scan:逐条 callback 调用 _deliver_callbackentry_callback.py:59)——
    _deliver_callback 函数体内部的分段边界已由「四个门控触发入口」片的"六"节处理(它只负责
    函数体内部三步怎么分段),本片负责 job 循环层面:每条 callback 的处理是否共享同一个
    session/segment 构造,还是每条各自独立构造——建议每条独立构造(避免一条 callback 的处理异常
    影响循环里其它 callback 已经在用的连接状态),与"一"节的独立会话机制配合
  • _drive_tick:逐 chat 循环调用 _evaluate_drive_tick_chat(「四个门控触发入口」片"三"节
    已处理该函数体内部分段);循环尾部 clear_recall_timing_overrideentry_drive_tick.py:329
    算哪一段——本片核实并确定:这次清理紧跟在 _evaluate_drive_tick_chat 返回之后,若两者之间无
    外部 I/O,可以合并进它的 post-LLM 段,若中间有其它 chat 的并发处理插入则需要独立成段,实现期
    按实际执行顺序核实
  • _compress_pending_delta:调用 run_delta_compression_cycle(逐 chat 驱动
    _compress_one_chat,「Delta压缩周期」片负责该函数体内部分段)——本片负责 job 循环层面的 session
    供应
  • _run_reflection:调用 run_reflection_cycle(逐 attempt 驱动 _apply_outcome/
    _commit_lesson,「反思闭环家族」片负责函数体内部分段)——本片负责 job 循环层面的 session 供应
  • _run_unanswered_escalation:核实其内部结构后(「反思闭环家族」片"四"节已要求先核实
    unanswered_escalation.py),确定 job 循环层面的 session 供应方式

三、并发重叠的显式处理策略

  • 每个 job 内部核实是否已有防重入机制(如成本池熔断 _pool_is_broken 只管 LLM 调用预算,
    不管 DB session 并发);若没有,评估是否需要新增(例如用一个简单的进程内锁/标志位防止同一个
    job 的上一次调用还没跑完、下一次调度又触发),或者确认"一"节的独立会话机制本身已经足以安全
    处理并发(每次调用各自独立连接,不共享可变状态),不需要额外加锁——两种方案选哪个、为什么,
    写进本片的实现说明,不要含糊带过(这是原始 AC 就点名要求的:"APScheduler 定时作业……需要显式
    判断:是否也该在各自的作业执行范围内共享一个事务,还是维持独立提交,写进本票的实现说明,不要
    含糊带过")

四、测试

  • "一"节的并发场景测试
  • 5 个作业各自补一条基本的分段行为测试(复用姊妹片已经交付的函数体内部分段,本片只需验证
    job 循环层面正确供应了 session)
  • 全量测试通过

Not in scope

  • _deliver_callback/_evaluate_drive_tick_chat/_compress_one_chat/_apply_outcome/
    _commit_lesson 等函数体内部的分段边界——分属「四个门控触发入口」/「Delta压缩周期」/
    「反思闭环家族」三个姊妹片,本片只负责这些函数被 job 循环调用时,session/segment 怎么从 job
    层面往下供应
  • 5 个作业的调度频率/触发条件本身——不改变现有的 APScheduler 配置,只改事务/会话边界

Blocked by

  • #197(存储层分段事务基建:session 传递契约与 SAVEPOINT 修复)——需要它交付的 session/事务传递
    API,以及本片"一"节需要的"不依赖 nonebot 上下文"这条扩展点
## Parent #151(ADR-0013 ①层单一事务边界落地——请求级事务,issue #108 Q1)。事务边界方向已确认为**分段**: 事务只包住连续无外部 I/O 打断的一段纯 DB 读写,下一步要发起任何外部调用前当前段必须先提交/关闭 (见 #151 comments"拆片前的事务边界设计方向已与用户确认")。 本片负责 **5 个 APScheduler 定时作业**:`_callback_scan`(`entry_callback.py:35`)、`_drive_tick` (`entry_drive_tick.py:265`)、`_compress_pending_delta`/`_run_reflection`/`_run_unanswered_escalation` (均在 `entry_offline_jobs.py`,分别在 30/96/145 行附近)。这 5 处已核实确实全部是 `@scheduler.scheduled_job` 直接装饰的顶层协程。 ## What to build design-verify 阶段 [HIGH] 发现:`get_scoped_session()`(第三方库 `nonebot_plugin_orm` 提供,本片 起草时核实其为外部依赖、非本仓代码,其内部 scope key 的精确实现细节以原始发现为准,未重新逐行核实 该库源码)的缓存 key 依赖 nonebot 的 `current_event`/`current_matcher` 上下文变量,对这 5 个跑在 事件/matcher 上下文**之外**的顶层协程,会全部坍缩成同一个 key,即它们实际共享同一个 Python `Session` 对象。现状因为每次调用都是"开-用-commit-close"的极短窗口所以风险很小;但"分段"设计 要求把持有窗口拉长到跨越真实 LLM 调用的单个 item 处理过程(秒级)。APScheduler 默认配置不阻止 5 个不同 job 之间并发重叠执行——一旦重叠,会在同一个不支持并发使用的 `AsyncSession` 对象上互相 踩踏,存在真实的会话状态损坏风险。 **这 5 个作业需要一套完全独立于 nonebot scoped-session 机制的会话构造方式**,不能直接复用反应式 入口那一套依赖 nonebot 上下文变量的机制。 ## Acceptance criteria **一、独立会话构造机制** - [ ] 设计并落地一套不依赖 `current_event`/`current_matcher` 的 session/segment 构造方式——建议 方向:每次 job 内部按处理单元(`_callback_scan` 每条 callback、`_drive_tick` 每个 chat、 `_compress_pending_delta` 每个 chat、`_run_reflection` 每条 attempt、`_run_unanswered_escalation` 每条候选)各自独立构造一个新的 DB 连接/session,不复用 `get_scoped_session()` 的缓存机制,具体 API 形状与「地基片」协调(`begin_segment()` 是否需要一个不依赖 nonebot 上下文的变体,或 `PersistentStorage` 内部提供另一条不经过 `get_scoped_session()` 的构造路径——见地基片"五、 APScheduler 5 个作业的独立会话构造方式不在本票范围"那条留下的接口空间) - [ ] 补一条测试直接验证并发场景:模拟两个 job(或同一 job 的两次调用)在真实存在时间重叠的情况下 各自开启并使用自己的 session,断言互不干扰(不是简单断言"代码跑通",要能证明两次调用确实使用了 不同的底层连接/session 对象,或者构造一个如果共享会暴露出来的具体错误场景,如两者交替写入不同 记录并断言写入结果符合各自预期而非互相覆盖) 原文(design-verify HIGH 发现,issue #151 comments): > `get_scoped_session()` 的 scope key `(id(current_event.get(None)), current_matcher.get(None))` > 对这 5 个作业全部坍缩成同一个 `(id(None), None)`……APScheduler 默认配置不会阻止 5 个不同 job > 之间的并发重叠执行——一旦重叠,会在同一个不支持并发使用的 AsyncSession 对象上互相踩踏,存在 > 真实的会话状态损坏风险。 **二、5 个作业各自按"两次外部 I/O 之间"分段(复用"一"节交付的独立会话机制)** - [ ] `_callback_scan`:逐条 callback 调用 `_deliver_callback`(`entry_callback.py:59`)—— `_deliver_callback` 函数体**内部**的分段边界已由「四个门控触发入口」片的"六"节处理(它只负责 函数体内部三步怎么分段),**本片负责 job 循环层面**:每条 callback 的处理是否共享同一个 session/segment 构造,还是每条各自独立构造——建议每条独立构造(避免一条 callback 的处理异常 影响循环里其它 callback 已经在用的连接状态),与"一"节的独立会话机制配合 - [ ] `_drive_tick`:逐 chat 循环调用 `_evaluate_drive_tick_chat`(「四个门控触发入口」片"三"节 已处理该函数体内部分段);循环尾部 `clear_recall_timing_override`(`entry_drive_tick.py:329`) 算哪一段——本片核实并确定:这次清理紧跟在 `_evaluate_drive_tick_chat` 返回之后,若两者之间无 外部 I/O,可以合并进它的 post-LLM 段,若中间有其它 chat 的并发处理插入则需要独立成段,实现期 按实际执行顺序核实 - [ ] `_compress_pending_delta`:调用 `run_delta_compression_cycle`(逐 chat 驱动 `_compress_one_chat`,「Delta压缩周期」片负责该函数体内部分段)——本片负责 job 循环层面的 session 供应 - [ ] `_run_reflection`:调用 `run_reflection_cycle`(逐 attempt 驱动 `_apply_outcome`/ `_commit_lesson`,「反思闭环家族」片负责函数体内部分段)——本片负责 job 循环层面的 session 供应 - [ ] `_run_unanswered_escalation`:核实其内部结构后(「反思闭环家族」片"四"节已要求先核实 `unanswered_escalation.py`),确定 job 循环层面的 session 供应方式 **三、并发重叠的显式处理策略** - [ ] 每个 job 内部核实是否已有防重入机制(如成本池熔断 `_pool_is_broken` 只管 LLM 调用预算, 不管 DB session 并发);若没有,评估是否需要新增(例如用一个简单的进程内锁/标志位防止同一个 job 的上一次调用还没跑完、下一次调度又触发),或者确认"一"节的独立会话机制本身已经足以安全 处理并发(每次调用各自独立连接,不共享可变状态),不需要额外加锁——两种方案选哪个、为什么, 写进本片的实现说明,不要含糊带过(这是原始 AC 就点名要求的:"APScheduler 定时作业……需要显式 判断:是否也该在各自的作业执行范围内共享一个事务,还是维持独立提交,写进本票的实现说明,不要 含糊带过") **四、测试** - [ ] "一"节的并发场景测试 - [ ] 5 个作业各自补一条基本的分段行为测试(复用姊妹片已经交付的函数体内部分段,本片只需验证 job 循环层面正确供应了 session) - [ ] 全量测试通过 ## Not in scope - `_deliver_callback`/`_evaluate_drive_tick_chat`/`_compress_one_chat`/`_apply_outcome`/ `_commit_lesson` 等函数体**内部**的分段边界——分属「四个门控触发入口」/「Delta压缩周期」/ 「反思闭环家族」三个姊妹片,本片只负责这些函数被 job 循环调用时,session/segment 怎么从 job 层面往下供应 - 5 个作业的调度频率/触发条件本身——不改变现有的 APScheduler 配置,只改事务/会话边界 ## Blocked by - #197(存储层分段事务基建:session 传递契约与 SAVEPOINT 修复)——需要它交付的 session/事务传递 API,以及本片"一"节需要的"不依赖 nonebot 上下文"这条扩展点
Sign in to join this conversation.
No milestone
No project
No assignees
1 participant
Notifications
Due date
The due date is invalid or out of range. Please use the format "yyyy-mm-dd".

No due date set.

Dependencies

No dependencies set.

Reference
ProjectKuma/arise#201
No description provided.