Skip to content

fix: prevent document lifecycle races and add async cleanup - #106

Open
ltrinthewind wants to merge 3 commits into
nageoffer:mainfrom
ltrinthewind:fix/document-lifecycle-race
Open

fix: prevent document lifecycle races and add async cleanup#106
ltrinthewind wants to merge 3 commits into
nageoffer:mainfrom
ltrinthewind:fix/document-lifecycle-race

Conversation

@ltrinthewind

@ltrinthewind ltrinthewind commented Aug 6, 2026

Copy link
Copy Markdown

Summary

  1. 使用 document_version 作为每次文档操作的版本及写入 fencing token,解决文档相关操作之间的竞态问题以及 ABA 问题。
  2. 参考现有 KnowledgeBase 清理方案,使用 RocketMQ 事务消息异步清理文档外部资源。
  3. 处理 RocketMQ 本地事务回滚后,接口仍可能返回成功的问题。

Changes

版本隔离与并发控制

  • 新增 document_version,每次分块、删除、恢复及稳定态修改通过 CAS 推进版本。
  • 分块事件携带本次任务版本,延迟消息和事务回查同时校验 RUNNING 状态与版本。
  • 首个 Sink 写入前锁定并校验当前任务版本,chunk、向量写入与 SUCCESS 状态在同一事务中提交。
  • 成功、失败状态仅允许持有当前版本的任务写回;超时恢复会推进版本,使恢复执行的旧任务自动失效。
  • 单 Chunk 修改、文档配置修改和启停操作同样推进并校验版本,避免与分块、删除交错产生孤儿数据。

异步资源清理

  • 文档删除通过事务消息将本地软删除与清理事件绑定,提交后异步清理外部资源。
  • 清理消费者分别处理 Milvus、ES、LightRAG 和对象文件,删除操作保持幂等。
  • 各清理阶段独立尝试;存在失败项时抛出异常触发 MQ 重试,已成功阶段可安全重复执行。

事务消息失败传播

  • 事务消息仅在 RocketMQ 返回 SEND_OK 且本地事务状态为 COMMIT_MESSAGE 时视为成功。
  • 非成功发送状态、ROLLBACK_MESSAGEUNKNOW 会清理本地事务回调并抛出 ServiceException,避免调用方收到假成功。
  • RocketMQ 发送过程直接抛出的底层异常保持原样传播,同时清理回调。

数据库升级

  • document_version 增量脚本放入 upgrades/v2.0.0/

Tests

  • 文档版本并发、Sink 原子写入和分块事务消息回查的状态与版本校验测试。
  • 文档删除、清理事务回查和外部资源分阶段清理测试。
  • 事务消息成功、非 OK、回滚、未知状态、底层异常及回调清理测试。

Related to #42
Related to #44

@ltrinthewind
ltrinthewind force-pushed the fix/document-lifecycle-race branch from 820d345 to 98d49ba Compare August 6, 2026 18:26
@magestacks

Copy link
Copy Markdown
Collaborator

供参考:

结论:建议 Request changes,暂不直接合并 [PR #106](https://github.com/nageoffer/ragent/pull/106)。方案修复了最基础的 TOCTOU,但“状态等于所有权”这一前提不成立,仍存在明确竞态。

主要问题

  1. [P1] 缺少 execution token,旧任务会篡改新任务

当前成功、失败和写入校验都只检查 status=RUNNING,没有验证“这个 RUNNING 属于谁”:

确定可复现的时序:

A: RUNNING
恢复任务: RUNNING → FAILED
B: FAILED → RUNNING
A: 捕获异常后执行 RUNNING → FAILED

最后一步会把 B 的运行权错误改成 FAILED。另一种情况是 A 在 B 启动后看到 RUNNING,直接提交 A 的旧分块结果。

这不是普通“已知优化项”,而是所有权模型不成立。应增加 execution_id 或单调递增 version,并让 MQ、调度、行锁检查、成功/失败 CAS 全部携带:

WHERE id = ?
  AND status = 'running'
  AND execution_id = ?
  1. [P1] 删除状态机没有覆盖 enable,删除后仍可能重新写入孤儿向量

[enable 方法](https://github.com/ltrinthewind/ragent/blob/98d49ba31d9ca9055420507ddf68386aa160a1e4/bootstrap/src/main/java/com/nageoffer/ai/ragent/knowledge/service/impl/KnowledgeDocumentServiceImpl.java#L678-L723) 会在事务外计算向量,随后忽略 documentMapper.updateById() 的返回值,并继续写向量。

enable: 读取文档并计算 vectorChunks
delete: CAS → DELETING,删除文档、chunk、向量后提交
enable: updateById 返回 0,但仍 indexDocumentChunks

结果是文档已删除,但向量被重新创建。至少需要检查更新行数;更完整的方式是让 enable 也领取带版本的操作权,外部索引任务携带 document generation。

  1. [P2] 定时刷新过早标记 SUCCESS,留下删除窗口

分块事务内已经写回 SUCCESS,但定时刷新直到随后才切换新文件元数据:

如果删除在两者之间抢到 SUCCESS → DELETING,文件元数据更新会返回 0,新上传文件最终被明确“保留待处理”,成为无文档引用的孤儿文件。

新文件元数据、chunk 数量、SUCCESS 应在同一 execution token 校验和同一事务中提交。

  1. [P2] Outbox 缺失仍会造成跨系统不一致

删除在数据库事务提交前同步删除索引和文件;文件失败还会被静默吞掉:

需要说明:补充材料中的“主库 MySQL、PgVector 独立 PostgreSQL”与当前仓库不完全一致。当前默认 PgVector 使用同一 JdbcTemplate,可加入本地数据库事务,见 PgVectorStoreService.java。但文件存储、Milvus、ES、LightRAG 依然不可回滚,仍需要幂等 Outbox。

同时,行锁目前覆盖整个 sinks.forEach。启用外部 Sink 后,数据库连接和行锁会跨越远程 I/O;建议只在短数据库事务内校验 token、写本地数据并登记 Outbox。

  1. [P2] 并发测试没有真正覆盖生产 Mapper 调用链

[DocumentLifecycleConcurrencyTest](https://github.com/ltrinthewind/ragent/blob/98d49ba31d9ca9055420507ddf68386aa160a1e4/bootstrap/src/test/java/com/nageoffer/ai/ragent/knowledge/dao/mapper/DocumentLifecycleConcurrencyTest.java#L42-L119) 手写了另一份 SQL,没有调用实际 Mapper/Service;staleWriterObservesDeleting... 也只是读取状态,没有执行真实 Sink。

建议补充实际 Spring/MyBatis 集成测试,至少覆盖:

  • A 超时、B 重启、A 失败回写。
  • A 超时、B 重启、A 旧结果写入。
  • 定时刷新完成分块与文件切换之间发生删除。
  • enable 计算向量期间发生删除。
  • 外部删除成功、数据库提交失败。

总体评价

CAS、DELETING 和提交前行锁的方向是正确的,能够修复“没有超时恢复、没有重新执行、只有分块与删除两个参与者”的基础竞态。但第一性原理要求“运行权必须属于某次具体执行”,而不能只属于 RUNNING 这个状态。

建议改造顺序是:execution token → 定时刷新原子收尾 → 外部清理 Outbox → 补真实并发集成测试

验证方面:PR 已在最新 main 上应用并编译通过,14 个非容器定向测试通过;3 个 PostgreSQL Testcontainers 测试因本机 Docker 不可用而被自动跳过。审查过程未修改工作区文件。

@magestacks

Copy link
Copy Markdown
Collaborator

其次,删除代码中的AI味。比如注释句尾的中文句号。

- 新增 document_version 和 DELETING 状态,统一文档操作所有权
- 使用版本 CAS 保护分块领取、超时恢复、删除、启禁用及 Chunk CRUD
- MQ、事务回查和定时刷新携带文档版本,跳过陈旧或重复任务
- Sink 写入前锁定当前运行版本,将 Chunk、向量、文件元数据和成功状态统一提交
- 修复旧任务覆盖新执行、删除后重建向量及定时刷新文件切换竞态
- 新增本地 PostgreSQL 并发、事务回滚和写入顺序测试
@ltrinthewind
ltrinthewind force-pushed the fix/document-lifecycle-race branch from 98d49ba to 0c56729 Compare August 9, 2026 21:29
- 使用 RocketMQ 事务消息投递文档清理事件
- 在本地事务中删除文档、Chunk、调度、日志及 PgVector 数据
- 事务提交后异步清理 Milvus、ES、LightRAG 和对象文件
- 补充事务回查、幂等消费、失败重试及存储分阶段测试
@ltrinthewind
ltrinthewind force-pushed the fix/document-lifecycle-race branch from 0c56729 to 514d907 Compare August 9, 2026 21:31
@ltrinthewind

Copy link
Copy Markdown
Author

感谢详细审查,已通过两个提交完成调整:

  1. 新增 document_version 作为版本号,结合行锁和 CAS 解决并发问题。
  2. 参考现有 KnowledgeBase 清理方案,使用 RocketMQ 事务消息异步清理文档外部资源,未引入 Outbox。

@magestacks

Copy link
Copy Markdown
Collaborator

一个代码问题,一个迁移问题,1.1.0早上我发布了,SQL位置需要改动下:

阻断问题

  1. [P1] RocketMQ 本地事务回滚后,接口仍可能返回成功

[RocketMQProducerAdapter#sendInTransaction](https://github.com/ltrinthewind/ragent/blob/514d907e75e278d344dc8b81466fdc145cd1155e/framework/src/main/java/com/nageoffer/ai/ragent/framework/mq/producer/RocketMQProducerAdapter.java#L81-L95) 只检查 sendStatus 来清理回调,却不会抛异常,也没有检查 localTransactionState

[DelegatingTransactionListener](https://github.com/ltrinthewind/ragent/blob/514d907e75e278d344dc8b81466fdc145cd1155e/framework/src/main/java/com/nageoffer/ai/ragent/framework/mq/producer/DelegatingTransactionListener.java#L69-L83) 会捕获本地事务异常并返回 ROLLBACK

因此可能出现:

删除 CAS 失败
→ 本地数据库事务回滚
→ RocketMQ 半消息回滚
→ sendInTransaction 正常返回
→ delete() 对外返回成功并记录删除日志
→ 实际文档仍然存在

FLUSH_DISK_TIMEOUT 等非 SEND_OK 状态也存在同样的“假成功”。当前测试甚至明确允许这种行为:[RocketMQProducerAdapterTest](https://github.com/ltrinthewind/ragent/blob/514d907e75e278d344dc8b81466fdc145cd1155e/framework/src/test/java/com/nageoffer/ai/ragent/framework/mq/producer/RocketMQProducerAdapterTest.java#L90-L109)。

建议:只有 SEND_OK + COMMIT_MESSAGE 才正常返回,其余状态向调用方抛出异常,并增加本地事务回滚测试。

  1. [P1] 必需的数据库迁移被追加到了已经发布的 v1.1.0

PR 把 document_version 迁移放在 [upgrades/v1.1.0/260810_document_version.sql](https://github.com/ltrinthewind/ragent/blob/514d907e75e278d344dc8b81466fdc145cd1155e/resources/database/upgrades/v1.1.0/260810_document_version.sql),但 [1.1.0 已经正式发布](https://github.com/nageoffer/ragent/releases/tag/1.1.0),发布说明明确列出了当时的 9 个迁移脚本。

已经升级过 1.1.0 的用户通常不会重新扫描该目录。新代码又强依赖 document_version NOT NULL,会导致存量部署的文档查询或更新直接报字段不存在。

下一个版本升级是2.0,resources/database/upgrades 下创建 v2.0.0 目录,同步更新数据库升级说明。

  1. 这个注释可以适当精简下,中英文不要空行:当前文档操作版本及写入 fencing token

修复上述两个 P1,并更新已经过时的 PR 描述后,可以再进入合并评审。考虑到改动较多,会进行多次CheckReview,见谅。

- 增加事务消息对本地事务执行结果的检查
- 将 document_version 升级脚本移至 v2.0.0
- 更新数据库升级说明和文档版本字段注释
@ltrinthewind

Copy link
Copy Markdown
Author

感谢 review,相关问题已在 24906fe 中修复:

  1. 事务消息增加本地事务执行结果检查,仅 SEND_OK + COMMIT_MESSAGE 视为成功;本地事务回滚或状态未知时会清理回调并抛出异常。
  2. document_version 迁移脚本已从 v1.1.0 移至 v2.0.0,并同步更新数据库升级说明。
  3. KnowledgeDocumentDO.documentVersion 注释已按建议精简。
  4. PR 描述已更新。

@ltrinthewind ltrinthewind changed the title fix: prevent races between document chunking and deletion fix: prevent document lifecycle races and add async cleanup Aug 11, 2026
@magestacks

Copy link
Copy Markdown
Collaborator

这个代码改动不对吧,怎么这么多文件提交?建议这个PR关掉,重新提一个干净的PR,不需要那么多单元测试,不写也可以。

其次考虑下这个问题:

[P1] 定时刷新可能删除数据库正在引用的新文件

在刷新写入成功后,事务已经把文档状态和 fileUrl 切换为新文件;但调度器随后重新读取的是全局 status,没有确认这个 SUCCESS 是否属于自己的 documentVersion。并发交错如下:

  1. 刷新 V 上传新文件并成功提交,数据库指向新文件。
  2. 手工切片 W 随即把状态改成 RUNNING
  3. 刷新 V 回读到 RUNNING,误判自己失败。
  4. finally 因本地阶段仍是 CHUNK_STARTED,删除新文件。
  5. 数据库最终指向一个已经被删除的对象。

反向交错也有问题:V 失败后,W 成功,V 可能把 W 的 SUCCESS 当成自己的成功并删除仍在使用的旧文件。

根因在于 ScheduleRefreshProcessor 用全局状态推断本次执行结果,随后按本地 phase [决定删除哪个文件](https://github.com/ltrinthewind/ragent/blob/24906fe1f39b6776097b586c3cc1943820783d08/bootstrap/src/main/java/com/nageoffer/ai/ragent/knowledge/schedule/ScheduleRefreshProcessor.java#L259-L270);而底层 runChunkTask 捕获异常且不返回归属明确的结果

修复应让切片接口返回“本 owner/version 是否已经提交文件切换”的明确结果,不能靠事后读取全局状态推断。

@ltrinthewind

ltrinthewind commented Aug 11, 2026

Copy link
Copy Markdown
Author

你好,提交文件多的原因有以下:

  1. 原始 issue: 文档删除与分块的并发竞态问题 #42 就包括文档删除与分块之间的竞态问题以及文档删除时,外部资源异步删除问题需要处理.
  2. 后续发现文档启用禁用,分块的 crud 这些操作之间,也同样存在竞态问题.
  3. 后续又解决了 Review 中又指出了RocketMQ 事务消息在本地事务回滚后仍可能返回成功的历史遗留问题.
  4. 单元测试较多.

想确认一下,更期望采用哪种方式重新提交:

  1. 保留一个 PR,整理 commit 并精简非必要测试.
  2. 拆成“文档操作并发控制”和“文档外部资源异步清理”两个 PR.
  3. 在 2 基础上,再将 RocketMQ 事务消息假成功问题单独拆成第三个 PR.

关于最新提出的定时刷新问题,我理解是:当前 chunkDocument() 没有返回本次分块任务的执行结果,调度器只能在任务结束后重新查询全局 status。这个状态可能已经被其他 documentVersion 对应的任务修改,导致调度器误判本次刷新是否成功,进而删除错误的文件。初步修改思路如下:

  1. KnowledgeDocumentServiceImpl#runChunkTask() 增加 boolean 返回值,表示本次 documentVersion 的分块任务是否执行成功.
  2. ingestionKernel.run() 正常完成后返回 true;核心流程发生异常时,处理失败状态并返回 false.
  3. KnowledgeDocumentServiceImpl#runChunkTask() 中的 updateChunkLog 日志更新单独捕获异常处理,不改变分块任务的执行结果.
  4. KnowledgeDocumentServiceImpl#chunkDocument() 返回 KnowledgeDocumentServiceImpl#runChunkTask() 的执行结果.
  5. 定时刷新根据 KnowledgeDocumentServiceImpl#chunkDocument() 返回值决定清理旧文件还是本次上传的新文件,不再重新查询全局文档状态.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants