Conversation
1979eb8 to
b68b221
Compare
c2b2b5d to
d4f0f01
Compare
|
@claude review |
|
Claude: Summary本 PR 为 ProduceBatchResult impact: not affected —— RoutedExperts impact: 本 PR 核心 —— 新增 owner 标记、把释放从 Ray concurrency impact: not affected —— 未改动 Verdict: REQUEST_CHANGES Main Flowchart after this PRflowchart TD
A[Rollout producer<br/>worker / vllm / verl tool] -->|owner = rollout| S[RolloutState.routed_experts]
B[TraceStore export_training_trace] -->|owner = trace_store 借用| S
S --> C{ReplayBuffer<br/>_apply_staleness_lifecycle}
C -->|retryable 且 owner == rollout| D[release_owned_routed_experts]
C -->|retryable 且 owner == trace_store| E[仅 detach]
D --> F[reset_rollout_response<br/>纯状态重置]
E --> F
C -->|non-retryable| G[release_and_discard_rollout_groups<br/>discard_rollout_state release_refs=True]
F -->|session_id 保留,retry 重新生成| H[SessionServer.on_response]
H --> I[Trie.insert 同 key 覆盖写]
I --> J[新增:_free_ray_refs old_value]
B -.借用同一 ObjectRef.-> J
G --> K[Trie.release<br/>session 最终释放点]
style C fill:#fff3cd,stroke:#d39e00
style D fill:#fff3cd,stroke:#d39e00
style I fill:#f8d7da,stroke:#c82333
style J fill:#f8d7da,stroke:#c82333
核心原理实现与单测核心原理是“谁创建谁释放”:直接 rollout 产出的 ref 归 单测覆盖情况: 抽象与信息隐藏评估
单测建议
其他 Issues
|
| old_value = node.value | ||
| if old_value is not None and old_value is not value: | ||
| # A rerolled turn may overwrite an existing key. Release refs | ||
| # that are no longer reachable, while preserving refs shared by | ||
| # another trie value or by the replacement value itself. | ||
| retained_refs = _collect_ray_ref_keys(value) | ||
|
|
||
| def collect_other_values(current: TreeNode) -> None: | ||
| if current is not node and current.value is not None: | ||
| _collect_ray_ref_keys(current.value, retained_refs) | ||
| for child in current.children.values(): | ||
| collect_other_values(child) | ||
|
|
||
| collect_other_values(self.root) | ||
| if retained_refs: | ||
| _free_ray_refs(old_value, _exclude=retained_refs) | ||
| else: | ||
| _free_ray_refs(old_value) |
There was a problem hiding this comment.
Claude: [正确性] [复杂] 这里的 retained_refs 只统计了 trie 内其他节点与替换值中的 ref,看不到已经导出给 RolloutState 借用的同一批 ref,因此覆盖写会在 session release 之前就 ray.internal.free 掉借用中的对象,等于在第二个位置重新引入了本 PR 要修复的 premature free。
RoutedExperts impact: TraceStore ref 的释放点从“仅 session release”变成“session release + trie 覆盖写”,而后者无法感知 XTuner 侧的借用方。
可复现链路:
agent_in_sandbox_loop.py的_build_rollout_states会从同一个session_id展开出多个 segment state,每个都持有data["routed_experts"]并标记owner="trace_store"。replay_buffer.py的_apply_staleness_lifecycle按expired_mask逐条判定,只有过期条目会被 reset/detach,未过期的兄弟 segment 仍在借用这些 ref。reset_rollout_response不清除session_id,retry 会复用同一 session 重新生成,SessionServer.on_response以相同的old_prompt/new_prompt再次store.insert。- 命中本段覆盖写分支后
_free_ray_refs(old_value)释放旧 ref,而兄弟 segment 仍指向它,后续 trainer 侧ray.get会变成对象取不到或 routed-experts/seq-len 校验失败。
建议二选一:覆盖写不做释放,保持 Trie.release 作为唯一释放点;或者在 session 内维护 ref 的引用计数(导出时 +1、session release 时统一归零),只在计数归零时才 free。
| if rollout_state.routed_experts_owner != "rollout": | ||
| logger.warning( | ||
| "Skipping release of routed_experts with unknown owner " | ||
| f"{rollout_state.routed_experts_owner!r} (session_id={rollout_state.session_id!r})." | ||
| ) | ||
| return |
There was a problem hiding this comment.
Claude: [兼容性] owner is None 会直接 return 而不释放。旧 checkpoint 恢复时 _restore_nested_objectrefs 用 ray.put 生成的是真 ref 且 owner 为 None,随后被 reset_rollout_response 丢弃,相对 PR 前反而新增泄漏。建议未知 owner 时按 rollout 处理或 free 后清空。
RoutedExperts impact: 恢复 checkpoint 后过期的 state 其 ref 永久留在 Ray 对象存储。
| # Direct rollout refs belong to this state/rollout path; | ||
| # TraceStore refs remain borrowed until the session is | ||
| # released by TraceStore. | ||
| if item.routed_experts_owner == "rollout": | ||
| release_owned_routed_experts(item) | ||
| reset_rollout_response(item) |
There was a problem hiding this comment.
Claude: [设计] owner 判定规则同时写在调用方和 release_owned_routed_experts 内部,属于同一规则散落在多个调用者:helper 已经能识别 "trace_store" 并跳过,这里再加一层 == "rollout" 并没有增加安全性,反而把“未知 owner 怎么办”的策略切成了两半——调用方按“不是 rollout 就不释放”,helper 按“不是 rollout 就 warn 并保留”,于是 owner is None 的 state 两边都不负责,直接导致上面 rl_data.py 的兼容性问题。
建议调用方无条件调用 release_owned_routed_experts(item),把 owner 语义完全收敛到 rl_data.py 一处(信息隐藏),调用方只表达“这条 state 到期了,释放它自己拥有的资源”这一意图。这样以后新增 owner 取值时也只需改一个地方。
RoutedExperts impact: 当前分工使 owner 未知的 state 不被任何一方释放。
| retained_refs = _collect_ray_ref_keys(value) | ||
|
|
||
| def collect_other_values(current: TreeNode) -> None: | ||
| if current is not node and current.value is not None: | ||
| _collect_ray_ref_keys(current.value, retained_refs) | ||
| for child in current.children.values(): | ||
| collect_other_values(child) | ||
|
|
||
| collect_other_values(self.root) |
There was a problem hiding this comment.
Claude: [性能] 每次覆盖写都遍历整棵 session trie,且 _collect_ray_ref_keys 对每个 TokenizedSegment 走 model_dump(),会深拷贝全部 token_ids/labels/logprobs,复杂度约 O(session 总 token 数)。这发生在单线程 RolloutTraceStore actor 的 on_response 每轮热路径上。建议只收集 expert_key,或维护 per-session ref 计数。
| ) | ||
| stale.routed_experts_owner = "rollout" | ||
|
|
||
| with patch("xtuner.v1.rl.replay_buffer.release_owned_routed_experts") as release_refs: |
There was a problem hiding this comment.
Claude: [测试] 这两个用例 mock 了项目内的 release_owned_routed_experts,只验证了调用方那层重复的 owner 判断,helper 真实行为(哪些 ref 被 free、routed_experts 是否置空)完全没覆盖——即便 helper 误释放 trace_store ref 也照样通过。建议改为 patch 项目外边界 ray_utils.free_object_refs,断言真实释放结果。
d4f0f01 to
4f3da9d
Compare
|
Addressed the review feedback in 4f3da9d:
|
| if rollout_state.status == Status.FAILED: | ||
| error_msg = rollout_state.error_msg | ||
| status = rollout_state.status | ||
| release_owned_routed_experts(rollout_state) |
There was a problem hiding this comment.
能否只修改reset_rollout_response 和 discard_rollout_state 内部的行为,这样就不用在每个调用点新增调用 release_owned_routed_experts 了
| status: Status, | ||
| prompt_tokens: int, | ||
| completion_tokens: int, | ||
| release_input_routed_experts: bool = False, |
There was a problem hiding this comment.
为什么要增加这个参数,什么情况下 release_input_routed_experts = False
There was a problem hiding this comment.
handler 无法判断它是本进程产出的还是从 TraceStore 借来的,现在就一个调用方,也可以先简化掉。
There was a problem hiding this comment.
我建议这个简化掉,不要增加接口的复杂性,默认就是True的配置就好
|
长期设计看,ObjectRef 的生命周期管理仍然过于分散:ReplayBuffer、PartialRolloutHandler、RolloutWorker、TraceStore、TrainerWorker 和 RLTrainer 都需要理解引用所有权并参与清理,reset/discard 还需要通过 release_refs 等参数控制行为,后续很容易继续引入重复释放或泄漏问题。 建议后续将 routed experts 抽象为独立的分布式 RoutedExpertsStore:
这样可以将所有物理释放操作收敛到 NodeStore,业务组件只管理 Handle 和 lease,并进一步删除 routed_experts_owner、release_owned_routed_experts()、discard_rollout_state(release_refs=...)、TraceStore 递归遍历 ObjectRef,以及 TrainerWorker 直接调用 free_object_refs() 等分散逻辑。 |
|
@claude review |
|
Claude: Summary本 PR 为 ProduceBatchResult impact: not affected —— RoutedExperts impact: 本 PR 核心 —— 新增 owner 标记;释放从 Ray concurrency impact: not affected —— 未改动 Verdict: REQUEST_CHANGES Main Flowchart after this PRflowchart TD
A[Rollout producer<br/>worker / vllm / verl tool] -->|owner = rollout| S[RolloutState.routed_experts]
B[TraceStore export_training_trace] -->|owner = trace_store 借用| S
S --> P{enable_partial_rollout?}
P -->|是| Q[PartialRolloutHandler.postprocess<br/>release_input_routed_experts=True]
Q --> Q1[owner == rollout 才释放 history<br/>owner is None 被跳过]
S --> C{ReplayBuffer<br/>_apply_staleness_lifecycle}
C -->|retryable| D[release_owned_routed_experts<br/>trace_store 仅 detach,其余释放]
D --> F[reset_rollout_response<br/>纯状态重置]
C -->|non-retryable| G[release_and_discard_rollout_groups<br/>discard_rollout_state release_refs=True]
H[SessionServer.on_response] --> I[Trie.insert 同 key 覆盖写]
I --> J[新增:_overwritten_values.append<br/>旧值被 actor 长期持有]
J --> K[Trie.release<br/>session 最终释放点]
G --> K
style D fill:#fff3cd,stroke:#d39e00
style Q1 fill:#f8d7da,stroke:#c82333
style I fill:#f8d7da,stroke:#c82333
style J fill:#f8d7da,stroke:#c82333
核心原理实现与单测核心原理是“谁创建谁释放”:直接 rollout 产出的 ref 归 单测方面,上一轮“mock 项目内 helper”的问题已修正: 抽象与信息隐藏评估
其他 Issues
|
| old_value = node.value | ||
| if old_value is not None and old_value is not value: | ||
| # A rerolled turn may overwrite an existing key. The old value's | ||
| # refs may still be borrowed by RolloutStates, so park the value | ||
| # instead of freeing immediately; the session release frees it. | ||
| self._overwritten_values.append(old_value) | ||
| node.value = value |
There was a problem hiding this comment.
Claude: [性能] 改动前 node.value = value 丢弃句柄后,Ray 引用计数即可回收旧对象(borrower 自持句柄,不会提前失效);改为 park 整个 TokenizedSegment 后,actor 会把旧 routed-experts 及 token_ids/labels/logprobs 一直 pin 到 session release,reroll 密集的 MoE session 会持续累积对象存储与 actor 内存。建议恢复丢弃句柄,或只 park expert_key。
RoutedExperts impact: TraceStore 侧 ref 的存活期从“覆盖即可回收”延长到“整个 session 生命周期”。
| if release_input_routed_experts: | ||
| if history_routed_experts_owner == "rollout": | ||
| free_object_refs( | ||
| history_routed_experts_ref | ||
| if isinstance(history_routed_experts_ref, list) | ||
| else [history_routed_experts_ref] |
There was a problem hiding this comment.
Claude: [设计] 这里自行复刻了 owner 判定,且与公共 helper 的策略相反:release_owned_routed_experts 是“非 trace_store 即释放”(作者上一轮已明确采纳“未标记即释放”的规则),而此处是“仅 rollout 才释放”。于是 owner 为 None 的 history ref 两边都不负责:这类 ref 真实存在——旧 checkpoint 经 replay_buffer._restore_nested_objectrefs 重新 ray.put 后是活跃对象但字段默认 None,若该 state 为 ABORTED 并走 partial rollout 续跑,第 218 行的拼接结果会直接覆盖 routed_experts,旧对象再无句柄可释放。
建议此处复用同一条规则(改为 != "trace_store",或直接调用 helper 的语义),把 owner 策略收敛到 rl_data.py 一处;同时补一个 routed_experts_owner=None 的 history 用例——现有 tests/rl/test_rollout_logic.py:1726 与 :1775 只覆盖了 "rollout" 与 "trace_store"。
RoutedExperts impact: 未标记 owner 的 history ref 在 partial rollout 拼接后泄漏在 Ray 对象存储中。
Motivation
When
enable_return_routed_experts=True, an RL rollout can attach RayObjectRefs to aRolloutState. The same reference may either be created by the rollout worker or be borrowed from aTraceStoresession. Before this change,RolloutStatedid not record which component owned the reference, while several cleanup helpers unconditionally freed references. This made retryable stale samples unsafe and could invalidate references that were still held by theTraceStore.Failure chain before this PR
The problematic path is easiest to see for a retryable stale sample:
ReplayBufferand is selected for retry.reset_rollout_response()to keep the prompt and clear the generated response.reset_rollout_response()also callsfree_object_refs()for every routed-expertsObjectRefit sees. It cannot distinguish a locally created rollout reference from a borrowedTraceStorereference.TraceStoresample, theTriestill contains the same reference because the session has not been released yet. The reset therefore frees a shared object too early.Trie. The reference is already invalid, which can surface as an object-fetch failure or as routed-experts/sequence-length validation errors.There were two related versions of the same ownership bug:
TraceStorestill owned it.TraceStoresession alive.Finally, replacing a value in the
TraceStoretrie did not release routed-experts references from the overwritten value, which could leak Ray objects during reroll/overwrite workloads.What changed
RolloutState.routed_experts_ownerwith the values"rollout"and"trace_store".reset_rollout_response()a pure state reset: it clears response fields and routed-experts fields, but never callsray.free.release_owned_routed_experts()for explicit caller-owned release, and makediscard_rollout_state(..., release_refs=...)opt in to releasing resources.RolloutWorker, vLLM rollout parsing, and the VERL tool loop) and at everyTraceStoreexport boundary.release_owned_routed_experts()unconditionally:"rollout"refs are freed,"trace_store"refs are only detached (they stay valid until their session is released), and untagged refs (restored from legacy checkpoints viaray.put) are also freed so they cannot leak.RolloutWorker.generate()opts in for direct rollout inputs; the handler default is non-releasing for safe reuse by other callers.TraceStoresession release as the final release point for its references, detach released routes before generic discard, and ensure the trainer performs session release in afinallyblock.RolloutStates may still borrow those refs, so the overwrite path cannot safely decide reachability. Parked refs are freed (deduplicated against the live tree) by the same session release that frees the trie.Ownership after this change
The LMDeploy generation protocol is unchanged. Ownership is established at XTuner's producer/export boundaries, and old checkpoints remain loadable because the new field defaults to
None.Impact
This prevents premature freeing of shared
TraceStorereferences, makes direct-rollout cleanup explicit, and closes exception/overwrite cleanup gaps without introducing a global reference registry or changing stale reroll/session semantics.Tests
The following targeted checks pass locally:
ruff checkfor all changed source and test filesruff format --checkfor all changed source and test filespython -m py_compilefor all changed source and test filesThe full replay-buffer suite was also inspected; an existing async save/resume test does not complete in this shared test environment, so it was not used as a passing signal.
Related issue
Related to #2025. This PR fixes routed-experts
ObjectReflifetime leaks and premature frees on the XTuner RL path. It does not, by itself, bound learner-side materialization or change sequence-parallel transfer order; those memory-footprint issues remain separate follow-up work. The LMDeploy endpoint and generation protocol are unchanged.