Skip to content

【No.11】Megatron 慢卡诊断与平台上报 - #378

Draft
shanyulu wants to merge 57 commits into
redai-studio:mainfrom
shanyulu:feat/task11-straggler-profiler
Draft

shanyulu wants to merge 57 commits into
redai-studio:mainfrom
shanyulu:feat/task11-straggler-profiler

Conversation

@shanyulu

@shanyulu shanyulu commented Sep 25, 2026 •

Copy link
Copy Markdown

改动目的

为 Megatron 训练增加默认关闭的慢 rank 诊断:沿用已有 timer 记录 host/CUDA 区间,按可比 rank 的连续窗口定位阶段性偏差,并接入现有日志与 TensorBoard。设计、取舍和完整实验解释见 RFC #357。

保持 Draft。 当前产品 / PR 头为 927c5de:在冻结的 a48a23b 之上加入确认告警指标、静默摘要读取、collector 锁隔离修复与指南文档修正(4 个提交;straggler+tools 套件 483 passed / 2 skipped,含锁隔离 old-fail/new-pass ×3);CI 在新头上重跑中(fork PR 需维护者批准)。C1 尚无确认性结论,C2 参数与 overlap、C3 平台与事件链仍需补齐。

实现与审查入口

flowchart TB
  subgraph capture["各训练 rank · 采样与读回"]
    direction LR
    T["Megatron timer<br/>host 计时 / CUDA event"] --> Q["有界队列<br/>区间结束时绑定上下文"]
    Q --> R["后台读回<br/>事件就绪后读取耗时"]
  end
  subgraph aggregate["global rank 0 · 进程内汇总"]
    direction LR
    C["Collector<br/>去重 / 迟到 / 丢弃计数"] --> D["Detector<br/>可比集合 / 连续窗口"]
    D --> J["JSONL<br/>计时事实与判定"]
  end
  R -->|"本 rank 直送;其他 rank 经后台 TCP"| C
  D --> S["rollout 节奏读取摘要"]
  S --> L["既有日志 / TensorBoard<br/>确认告警指标已加入,最终版本平台归属待回归"]
  classDef blue fill:#eaf2ff,stroke:#4778b8,color:#142d4e
  classDef green fill:#eaf6ee,stroke:#488766,color:#173c29
  classDef amber fill:#fff4dd,stroke:#bc8734,color:#51360b
  style capture fill:#f3f6fa,stroke:#9eacc1,color:#203651
  style aggregate fill:#f3f6fa,stroke:#9eacc1,color:#203651
  class T,Q,R blue
  class C,D,J green
  class S,L amber
Loading

图示为 CUDA 路径且配置共享 collector 地址时的行为。汇总位于 global rank 0 进程内,不是独立服务。未配置共享地址时,各 rank 本地汇总;纯 host-only 模式可能在调用线程交付,不能声称所有处理都在后台。

建议按以下顺序审查:

入口 重点检查
model.py、actor.py timer 注入位置、开关关闭行为、workload 上下文
timer shim、observer.py 区间结束时固定上下文;后台 event 读回;host-only 降级不越序
runtime.py、collector.py 有界队列、进程内汇总、重复 / 迟到 / 丢弃、退出与 flush
detector.py 可比集合内选参照;候选连续计数与已生效告警分离
reporter.py、train_metric_utils.py 既有平台路径、确认告警归属、读取摘要的副作用
单元与回归测试 顺序、退化、窗口间隙、低负载状态机及故障边界

关键语义

  • 相对偏差以该 rank 的窗口 host 耗时中位数,对比工作量可比集合中的最快成员;不使用全 cohort 最快值替代有效参照。
  • a48a23b 已修复低负载窗口清除告警、onset 连续计数跨间隙串联等问题。不确定窗口不能被当成恢复证据。
  • 927c5de 进一步加入确认告警指标(confirmed_straggler_deviation/rank,排除 uncertain/recovered,保留原 worst_rank 语义)、静默摘要读取,并将 collector 的文件写入与外部回调移出状态锁(写阻塞/写失败/回调阻塞三条 old-fail/new-pass 回归钉住)。
  • 工作量缺失会标记降级,仍可能产生判定;没有足够可比成员与工作量缺失不是同一分支。
  • CUDA event 表示流上可观察区间。host / stream stall 是诊断提示,不是硬件根因证明。
  • 窗口由后续记录推进或显式 flush 关闭;默认 5 秒窗口不是平台可见时延上界。

当前验收

项目 证据与版本 当前结论
当前头 CI a48a23b 8/8 通过(通用 CI、GPU 单测、训练集成);927c5de 重跑中(fork PR 需维护者批准) 待新头结果;本地全套 483 passed / 2 skipped
C1:整体开销 <0.5% e961661;6 对 AB/BA,均值 +0.417%,95% CI [-1.295%, +2.023%] INCONCLUSIVE;旧版本 pilot,不是当前头确认验收
C2:训练指标 a48a23b;两对 ON/OFF,loss / 梯度范数在冻结包络内,学习率一致、token 差为 0 指标检查通过;不代表整体 C2 通过
C2:参数等价 a48a23b;checkpoint 树哈希不同,旧分析器未实际比较张量;复核判定 未证明(复核判定 INCOMPLETE:0 指标违规、2 项参数证据缺失)
C2:通算重叠 a48a23b;DP4 + overlap-grad-reduce;冻结包络超界约 6.2e-5 NOT PASS;原结论保留
C3:定位与上报 a48a23b;6 条 straggler:rank 3 五条、rank 2 一条;本轮健康对照零 straggler 部分完成;平台归属、可靠延迟、非目标告警和末窗仍待验证

证据分支现固定于 b9c01f2(含复核更正记录、严格验收门禁工具、测量臂原始日志与 16 件原始 trace 归档;原始实验仍按其 manifest 的产品 SHA 解释):C1 配对结果、C2 复核判定、overlap 判定、C3 原始判定。每条结论按该次实验的产品 SHA 解释,不跨版本继承。

C1 实验图与判定方法
%%{init: {"themeVariables": {"xyChart": {"plotColorPalette": "#2969a6,#bd6530"}}}}%%
xychart-beta
  title "C1 pilot:配对开销与验收阈值"
  x-axis ["S1 AB", "S2 BA", "S3 AB", "S4 BA", "S5 AB", "S6 BA"]
  y-axis "相对 OFF(%)" -3.5 --> 3
  line [2.177, 1.779, -2.308, 2.433, 1.393, -2.970]
  line [0.5, 0.5, 0.5, 0.5, 0.5, 0.5]
Loading

A=OFF,B=ON,每臂 48 步。折线各点为各对整段 perf/train_time 均值相对差,水平线为 0.5% 参考线;连线不作趋势拟合。预注册判定使用配对差均值的区间上界;+0.417% 的点估计不能抵消 +2.023% 的上界。bootstrap 以 pair 而非 step 为重采样单位,不把同一训练中的相关步当作独立试验。

仍需修复或补证的内容

  1. 确认告警平台回归。 293390e 已兼容加入确认告警指标(confirmed_straggler_deviation/rank),4654e5a 静默摘要读取,f120aa8 将文件写入与回调移出状态锁(old-fail/new-pass ×3)。本地测试不能替代最终产品版本上的平台归属回归,该验证待 GPU 复验。
  2. 参数与公开复算。 loss 包络不能替代参数比较;保留可比较 checkpoint 并实际检查。测量臂原始 job.log 已随复核更正入库存档(原始件 + RFC1918 脱敏公开件 + 双哈希台账 + gitleaks 扫描报告,见 f69343f8);16 件原始 chrome trace(4 臂 × 4 rank,260MB)亦已入库存档并经 gitleaks 全量扫描 PASS(16/16、零发现,无需脱敏),附双哈希台账——干净克隆可直接复算 overlap 判定(已验证:exit 1、NOT_PASS、13/13 字段一致)。
  3. Overlap 独立验证。 原实验 NOT PASS 不改写。跨会话方差须通过独立 OFF/OFF 校准,冻结新包络后再看新的 ON 数据。
  4. C3 事件链。 关联区间结束、判定、落盘、平台接收的同一事件;解释 rank 2 告警,验证尾窗口与无后续记录场景。旧 0.60/1.20 秒“延迟”已撤回:它没有事件身份关联。
  5. 支持范围。 当前证据为单机 dense DP4。PP>1 时,global rank 0 collector 与末 PP stage 导出主 rank 的归属未闭合;attention/MoE 仅 schema,不能写成已完成插桩。

四臂 cudaDeviceSynchronize 计数均为 68,只说明该 API 计数一致;不作“零新增全局同步”的扩大结论。

受影响验收声明。 产品头由 a48a23b 变更为 927c5de(行为修复),受其影响的验收——C2 参数比较、C3 事件链与平台归属、overlap 校准——须在新版本上重新验证;上表历史实验结论仍按各自版本 SHA 解释,不迁移。

转正式审查的条件

  • 公开产品缺口修复,变更后的当前头测试与 CI 通过。
  • 导师确定 C1 主指标后,冻结最终产品、样本量、顺序、分析器和停止规则,取得正式确认结论;不运行到通过为止。
  • C2 完成参数比较与 overlap 验证,或维护者明确接受替代验收。
  • C3 完成平台定位与事件链验证;实时 cadence、attention/MoE 范围及 C1 主指标取得明确裁决。

产品代码由本 PR 承载,设计以 RFC 为准,实验失败及原始记录保留在 evidence。本文只列当前可支持的结论,不以 CI 绿灯代替官方验收。

Relax disables Megatron's config.timers because Timer.start/stop call
torch.cuda.synchronize() (and a barrier when barrier_with_L1_time is set),
serialising the step they measure. The cost is that the framework has no
per-stage timing at all, so a slow rank can only be seen as "the step got
slower" with no attribution.

This adds relax/utils/straggler, off unless RELAX_STRAGGLER_ENABLE is set:

- megatron_timer_shim: drop-in config.timers object with the same call
  shapes and log-level filtering Megatron uses (61 accesses, 23 names, only
  start/stop), but start/stop only record a host timestamp and a CUDA event
  on the current stream. No device sync, no collective, no added allocation
  in the hot path; every entry point fails open and counts anomalies instead
  of raising.
- observer: fixed CUDA-event pool plus a daemon readout thread. Intervals
  whose event is not complete yet are retried; exhaustion, a full queue, a
  readout timeout and consumer failures all degrade to host timing with a
  counted reason, so profiling can never block or abort training.
- identity: defensive rank/cohort discovery (TP/PP/VPP/CP/EP position with
  DP excluded) so only equivalent ranks are ever compared.
- config + Envs entries; every knob defaults to off.

model.py assigns the timers to both the model's TransformerConfig and the
optimizer's OptimizerConfig (Megatron reads them separately), after model
construction because attention/MoE deepcopy their config while being built.
With the profiler off both assignments are None, unchanged behaviour.

Verified on 4x RTX 4090 with real CUDA events: 9/9 envelopes read back with
device times, 0 timeouts, 3 suppressed barriers. 53 CPU tests, pre-commit
clean.
Add the judgement and aggregation half of the straggler profiler, still off by
default and still free of synchronisation:

- detector compares each rank against the fastest peer in its cohort over a
  sliding window, requires the deviation to persist, and reports an incomplete
  cohort as uncertain instead of guessing;
- collector aggregates every rank's envelopes, writes JSONL evidence when an
  output directory is configured, and emits a periodic summary line;
- same-host TCP transport ships envelopes to rank 0, dropping (and counting)
  rather than blocking when the collector is away or the queue is full;
- runtime wires observer, collector and transport per role, and exposes a
  status() snapshot for the existing logging and metrics pipeline.
Assert, by parsing the backend module, that config.timers is decided in exactly
three places: the optimizer config and the training TransformerConfig take the
profiler's object, while the evaluation path keeps upstream's explicit None.
Also assert the profiler's object is never passed into a config constructor,
which would leak its events and threads into the attention/MoE configs through
copy.deepcopy.
Real four-process runs exposed three defects that unit tests alone would not
have caught, and each is fixed with the evidence that found it:

- an identical rank was reported as a straggler at +23% in its second window,
  purely from lazy CUDA event allocation, so the opening windows are excluded
  and the exclusion is counted (warmup_windows, default 2);
- window membership used an absolute clock, so ranks whose processes started
  seconds apart shared no window and every judgement degraded to 'no peer';
  windows are now relative to each rank's first observation;
- the boundary arithmetic truncated (4.1 - 0.1 floors to 3), silently merging
  two windows; indices are computed in integer microseconds.

Also stop emitting one 'uncertain' verdict per rank per window: a lone rank of a
multi-rank cohort is a coverage gap and is counted (incomplete_windows) instead of
burying real findings, and the device attribution is renamed to what CUDA events
can actually show: host_only_stall, gpu_stream_stall, attribution_unknown.
… measure

Three additions that make the observation trustworthy rather than merely
present:

- The cohort key now carries the topology epoch, the stage schema and the
  model chunk in addition to TP/PP/VPP/CP. Two ranks may be compared only when
  they run the same role for the same chunk under the same layout; a re-shard
  invalidates every comparison instead of silently pooling incomparable ranks.
  EP/ETP/EDP are carried in the schema but never used for grouping, because
  comparing expert-parallel roles is not claimed to work.
- stages.py maps the timer names the pinned Megatron actually emits onto coarse
  forward/backward/optimizer/communication/data/setup/eval groups, with the
  inventory command in the docstring. Attention and MoE are declared
  schema-only: Megatron core 0.19 exposes no such timer, and trading an honest
  'not instrumented' for an invented stage name is not acceptable. Unknown
  names stay unclassified rather than being guessed into a group.
- A test-only slow-rank injection (three env knobs, inert by default) that
  delays inside the measured interval, so the detector's sensitivity curve can
  be measured against a known +X% instead of whatever an external injection
  happens to produce. It changes nothing but the sleep and is resolved through
  the sink, so exactly one rank can be targeted.

The recipe in scripts/training/genrm is the real four-GPU observer smoke: the
proven end-to-end DAPO/GRPO/GenRM job of this machine with the observer bolted
on through environment variables only, default off, so an OFF arm and an ON arm
differ in nothing else.
… step context

Completes the observation chain from the measured interval to something a
platform can consume, without adding work to the training thread:

- protocol.py defines the wire contract -- protocol name, schema version,
  field validation and the idempotency key -- so a malformed or replayed
  packet is counted and dropped at the collector instead of reaching the
  detector or the JSONL evidence.
- context.py publishes the training loop's rollout id and optimizer-step
  index as an immutable snapshot behind a lock; train_one_step records it
  through one counted call that is a no-op while the profiler is off, which
  lets a verdict name the step it belongs to.
- reporter.py turns the runtime's counters and pending verdicts into scalar
  perf/straggler/* metrics and train_metric_utils merges them into the perf
  line only when the profiler is enabled; a value that cannot be measured is
  omitted rather than reported as zero.
- collector.py gains the validation/dedup gates and a non-flushing summary,
  runtime.py exposes summary() for the training thread while report() keeps
  the flushing behaviour for diagnostics, and the detector/observer pick up
  the streak and comparability changes pinned by the new tests.

The suite grows to 229 CPU tests covering the protocol gates, the context
snapshot and the metric surface.
…FT set

The genrm smoke is a real end-to-end job but its actor world is two data-
parallel ranks; the acceptance protocol asks for a cohort of four, so this
recipe makes all four local RTX 4090s identical SFT replicas (TP = PP =
CP = 1) and leaves the profiler opt-in and default-off exactly like the
other observer recipe, so an OFF arm and an ON arm differ only in the six
RELAX_STRAGGLER_* variables.

The OpenMathReasoning-mini parquet the 8-GPU parent recipe expects is not on
this machine, so tools/straggler/make_sft_dataset.py derives a small, fixed
JSONL from the local dapo-math-17k prompts in the string+label shape the
Relax SFT loader accepts; the 256- and 512-row outputs are committed with
the recipe so its default PROMPT_SET resolves without a network fetch.
Add the Phase-1 invariant audit as executable evidence:

A. zero training-path collective: an AST guard over relax/utils/straggler/
   (no dist.* collective, no torch.cuda.synchronize/.item/.tolist, no
   blocking requests/socket in the training path; sockets confined to the
   background transport) plus a runtime instrumentation test that wraps the
   forbidden entry points with raisers and drives the real
   timers -> observer -> sender -> socket -> receiver -> collector path,
   asserting zero forbidden calls and that socket/persistence operations
   happen only off the training thread.

B. strictly bounded: the event pool never allocates past capacity, its full
   condition drops and counts, and every documented structure cap is asserted
   directly.

C. failure cannot escape: injected event create/record/read failures, a
   failing consumer, a broken serialiser, an unreachable collector, a crashing
   detector, malformed/duplicate/late packets and shutdown with unfinished
   events are each counted and the training-side call returns normally.

tests/utils/straggler: 239 passed.
The finish path performs one forced synchronous save even with
--save-interval effectively off, which would dominate the whole-run wall time
of the short OFF/ON timing arms. SAVE=0 omits --save/--save-interval/--no-save-rng
so both arms run without a terminal checkpoint; SAVE defaults to 1, so the
proven smoke behaviour is unchanged. Requested by the Task 11 lead for Phase 3.
Symptom: a submitted arm could import another worktree's relax package and run
with no straggler module at all, yet still look green.

Root cause: ray job submit was invoked without --working-dir unless the
operator set WORKING_DIR by hand, so the Ray daemon's own cwd (another
worktree) decided which relax the driver and actors imported.

Fix: both observer recipes now default WORKING_DIR to this repository root.

Evidence: every remaining GPU arm must be relaunched, and its manifest must
record the resolved relax.__file__ for the driver and one actor; an arm that
does not resolve under task11-c2 is invalid.
Symptom: a submitted arm could import another worktree's relax package and run
with no straggler module, yet still look green; the runtime-env PYTHONPATH in
the submit log did not prove what the driver and actors actually imported.

Root cause: nothing recorded the import root at runtime, so the manifest guard
("resolved relax.__file__ for the driver and one actor") had no artefact to
read.

Fix: relax.utils.straggler logs relax_root and its own module path once per
process at import. model.py imports this package unconditionally, so the line
appears in the Ray driver and in every Megatron actor for both OFF and ON arms.

Evidence: every remaining GPU arm must be relaunched and its manifest must
record the first provenance line (driver) and a MegatronTrainRayActor
provenance line (actor); a line that does not resolve under task11-c2
invalidates the arm. tests/utils/straggler + test_train_metric_utils: 241 passed.
Symptom: the acceptance protocol reads observer/sender/collector counters
(queue high-water, dropped pool/queue, malformed/duplicate/late, coverage,
dedup and detector evictions), and analyze_run.py already reads
collector_status.json from the evidence directory, but nothing wrote either
file. Phase 3 would have had no counter artefact to report.

Root cause: StragglerRuntime.status() aggregated timers/observer/collector/
sender/receiver counters but only close() called collector.report(), which
logs a single summary line and persists nothing.

Fix: close() writes runtime_status.json (all components) and
collector_status.json into RELAX_STRAGGLER_OUTPUT_DIR, inside the existing
fail-open try/except; no output_dir means a clean no-op. Runs at shutdown,
after the readout threads stop, never on the training thread.

Evidence: test_close_persists_the_runtime_and_collector_counters and
test_close_persists_the_sender_side_counters (new); 244 passed (straggler +
test_train_metric_utils).
# 🐛 Bug Fix

## Stop the FrozenInstanceError on `config.timers`

- Symptom: the first enabled real-machine arm (Ray job `mbTivCsZz9YY23ft`)
  failed at step 0 with `dataclasses.FrozenInstanceError: cannot assign to
  field 'enabled'` on the path `MegatronTrainRayActor.train_async ->
  update_weights_fully_async -> megatron_to_hf -> remove_non_pickleables`.
- Root cause: `remove_non_pickleables(config, max_depth=3)` in
  `megatron/bridge/models/conversion/utils.py:250-263` recurses through every
  attribute of any object that has `__dict__`, doing `copy.copy(obj)` and then
  `setattr` for each attribute in `vars(obj)`. The enabled shim in
  `config.timers` exposed the frozen `StragglerConfig` as `_config`, so the walk
  assigned its `enabled` field and raised.
- Fix: `StragglerTimers` is now `__slots__`-based, so it has no instance
  `__dict__` and the walk leaves it untouched instead of recursing into the
  frozen config. `StragglerConfig` stays frozen and `remove_non_pickleables` is
  unchanged.
- Fix: `StragglerTimers.__reduce__` reconstructs an inert, sinkless shim via
  `_inert_straggler_timers`, because the walk is immediately followed by
  `broadcast_obj_from_pp_rank`, which pickles the cleaned config; a live shim
  carries a readout thread, a socket, locks and a CUDA event pool.
- Evidence to rerun: every ON arm of the §5-§13 acceptance phases, because the
  existing ON evidence is classified `INVALID-observer-bug` (report Appendix A,
  A7). This is a GPU-window item.

---

# ✅ Tests

## Add the A7 pickle-boundary regression

- `tests/utils/straggler/test_straggler_pickle_boundary.py`: asserts the shim
  has no instance `__dict__`, that pickling yields an inert but functional copy,
  and that the real upstream `remove_non_pickleables` accepts a real
  `TransformerConfig` with `get_straggler_timers()` installed and leaves the
  live shim usable afterwards.
- The end-to-end test loads the upstream module by file and runs when
  `MEGATRON_PATH`/`MEGATRON` points at a Megatron-LM checkout; otherwise it
  skips with the reason, because the CPU venv has no complete Megatron bridge
  stack.
- `247 passed`: `tests/utils/straggler` + `tests/utils/test_train_metric_utils.py`.
# 🔩 Chore

## Stop tracking the generated SFT payloads

- `git rm --cached` the two generated payloads that `6e2aca0` tracked:
  `scripts/training/sft/data/dapo-math-17k-sft-256.jsonl` (146137 bytes) and
  `scripts/training/sft/data/dapo-math-17k-sft-512.jsonl` (292510 bytes). The
  working-tree copies are kept so the DP4 recipe still runs; only the index
  entries are removed.
- Add a narrowly scoped ignore rule `scripts/training/**/data/*.jsonl` so the
  derived payloads cannot be re-added, while the generator
  `tools/straggler/make_sft_dataset.py` and both observer recipes stay tracked.
- Verified the generator reproduces both payloads byte-for-byte:
  `python tools/straggler/make_sft_dataset.py --seed 42 --num-rows 256` and
  `--num-rows 512` yield the committed SHA-256 values `44f9ddac...` and
  `4f778f6f...`.

Evidence: `git ls-files -- scripts/training | grep jsonl` is empty, and
`git check-ignore -v` resolves both paths to the new `.gitignore` rule.
Training-path isolation:
- collector: never write JSONL from the constructing/training thread; buffer
  host-only ingest under a cap and drain it off-thread (RT-01).
- collector: swap the pending buffer under a lock and write outside it, so
  concurrent reader threads can no longer duplicate or drop persisted lines
  (RT-02).

Bounded structures:
- receiver: prune finished reader threads before applying MAX_CONNECTIONS, so
  the cap counts live connections instead of lifetime ones (RT-03).
- shim: cap the timer-name table at MAX_TIMER_NAMES (RT-06).
- protocol: length-cap free-text fields so DEDUP_MAX_ENTRIES is the binding
  memory constraint (RT-05).

Correctness and attribution:
- detector: drop non-finite/negative host and device samples instead of letting
  one bad packet become the peer reference and silence a window (RT-04).
- identity: include expert-parallel roles in the cohort key so ranks holding
  different experts are never compared (RT-07).
- observer/protocol: carry workload counters over the wire with a bounded
  sanitiser so workload_delta is reachable in a real collector run (RT-08).

Regression tests live in tests/utils/straggler/test_straggler_redteam_regressions.py;
the pickle boundary now proves the pre-fix frozen-config crash with a
Megatron-free local walker.
# 🐛 Bug Fix

## Persist the status files while the run is alive

- Symptom: the first real ON arm that completed (DP4 SFT, Ray job
  `raysubmit_6WPA4CUy5EYWtKGY`) wrote `straggler_envelopes.jsonl` and
  `straggler_verdicts.jsonl` but no `runtime_status.json` /
  `collector_status.json`, so `analyze_run.py` reported
  `collector_status: missing` for a valid run.
- Root cause: only `StragglerRuntime.close()` calls `_write_status()`. Ray
  terminates the actor with SIGTERM at job end, so the registered atexit hook
  never runs and `close()` is never called.
- Fix: `start()` launches one daemon status-writer thread (only when
  `output_dir` is set) that writes the counters every
  `min(report_interval_seconds, 10 s)`; `close()` stops it. The writer stays off
  the training thread, is bounded (one thread, constant interval) and fail-open
  (a write error is counted, never raised).
- Evidence to rerun: every ON arm of the §5-§13 acceptance phases, because the
  status files were absent for the ON smoke.

---

# ✅ Tests

## Add a regression for the no-close path

- `tests/utils/straggler/test_straggler_runtime.py::
  test_status_files_are_written_without_a_graceful_close`: starts a runtime with
  `output_dir` and asserts both status files appear without calling `close()`.
- `263 passed, 1 skipped, 1 xfailed`: `tests/utils/straggler` +
  `tests/utils/test_train_metric_utils.py` at the new tip (includes the
  red-team commit `a61543c`).
A rank that started >= 3 windows after its fastest peer was closed out before
its first sample and was never compared for the rest of the run. Windows are
now indexed from a bounded per-cohort anchor (_cohort_epoch), so every live
rank writes the same wall-clock window, and warmup is applied per rank against
its own first valid observation (_epoch stays as the warmup baseline).

- detector: reject a non-finite host_start before it can poison the cohort
  anchor or the warmup baseline; a warmup sample creates its window (so the
  skip is counted) but records no sample.
- detector: report warmup_samples_skipped, cohort_epoch_evictions,
  aligned_cohorts, MAX_COHORT_ANCHORS.
- test_straggler_redteam_regressions.py: the RT-09 xfail becomes a real
  passing regression test; new test proves a late rank's own cold-start
  windows are warmup, not a false straggler.
- test_straggler_detector.py::test_staggered_rank_starts_are_aligned_by_relative_time
  adjusted. Justification: the old form fed each rank's own relative window at
  a different wall time -- exactly the non-contemporaneous samples the cohort
  anchor now separates. It now feeds the three ranks in wall-clock lockstep;
  the assertions (straggler == {2}, no uncertain, aligned_ranks == 3) are
  unchanged, so the new number is not weakened, only the input is corrected.

Fail-before: with detector.py reverted, the two RT-09 tests plus the staggered
test fail (3 failed); after the fix all 49 detector+redteam tests pass.
The collector ingests from one thread per accepted connection plus the
observer readout thread, and its dedup tracker, detector, counters and
verdict deque were mutated lock-free. Under a 1 us interpreter switch
interval that loses evidence silently, reproducibly:

- BoundedDedup._check raised 'OrderedDict mutated during iteration' in its
  expiry walk; check() swallowed it and returned 'new' without recording the
  key, so 13-15 of 1600 distinct packets per run were never remembered and a
  resend of those keys could be judged twice.
- A window close racing a concurrent _Window.add raised 'StatisticsError: no
  median for empty data'; observe() swallowed it and dropped that whole
  window's verdicts with only a log line, leaving every counter unchanged
  (windows_closed even went up).

Changes:
- collector: add a re-entrant _state_lock around ingest, flush, status,
  drain_verdicts and the counter mutations; _write_batch counts flushed_lines
  and write_errors under _write_lock (two concurrent flushes could lose a
  +=). No file I/O ever runs under _state_lock, so a training-thread ingest
  still cannot block behind persistence.
- protocol: BoundedDedup.check/stats take an internal lock, making the public
  class self-contained for callers outside the collector.
- detector: document that it is not internally synchronised because its only
  production caller (TimingCollector) now serialises every call; a second lock
  there would only duplicate that.

Regression test test_concurrent_ingest_conserves_every_tally hammers the real
ingest path from 8 writers plus a status reader and asserts exact
conservation: envelopes == judged == dedup new == dedup entries == flushed
lines == persisted lines, with zero invalid/duplicate/late-ingest errors.
Fail-before: with collector.py and protocol.py reverted it fails with the
StatisticsError; after the lock it passes on 5/5 repeated runs.
# 🐛 Bug Fix

## Keep the observer evidence after a SIGTERM

- Symptom: the clean ON smoke completed (job `raysubmit_pUjysBz9MeYsNVKh`)
  and the collector log shows `envelopes=132 judged=84 windows=4`, but
  `straggler_envelopes.jsonl` stayed at 0 bytes and `analyze_run.py` reported
  `envelopes: 0`.
- Root cause: only `StragglerRuntime.close()` flushed the collector, and Ray
  terminates the actor with SIGTERM at job end, so `close()` never ran and the
  buffered lines were lost.
- Root cause: all four ranks share one `RELAX_STRAGGLER_OUTPUT_DIR`, so they
  raced on one unsuffixed `runtime_status.json` / `collector_status.json` and
  whichever rank wrote last left its counters behind.
- Fix: the periodic status writer now flushes the collector before persisting
  (still a daemon thread, never the training thread), so a killed arm keeps
  usable data at most one interval stale.
- Fix: output is run-scoped and per-rank
  (`<base>/run_<ray-job-id>/runtime_status_<label>.json`). The unsuffixed names
  `analyze_run.py` reads are written only by rank 0, so there is exactly one
  writer; writes go through a temp file plus `os.replace`.
- Evidence to rerun: the ON arms of the §5-§13 acceptance phases, which is what
  this unblocks.

---

# ✅ Tests

## Cover the no-close path, the flush and the path uniqueness

- `test_periodic_writer_flushes_the_jsonl_without_a_close` asserts envelopes
  reach disk with no `close()`.
- `test_output_paths_are_run_scoped_and_per_rank` asserts four ranks produce
  four distinct per-rank files and one rank-0 canonical status.
- `test_status_writer_never_touches_a_file_on_the_training_thread` spies on
  `_write_status`/`TimingCollector.flush` and asserts neither runs on the main
  thread.
- `269 passed, 1 skipped`: `tests/utils/straggler` +
  `tests/utils/test_train_metric_utils.py`.
# 🐛 Bug Fix

## Bound the evidence lost to an abrupt termination

- Symptom: the collector ingested 132 envelopes / judged 84 / closed 4 windows
  on the clean ON smoke, yet the JSONL was zero bytes after the job, because
  Ray SIGTERMs the actor and `close()` never runs.
- Fix: `_write_json` writes a temporary file, flushes, then `os.replace()`s it
  into place, so a reader never sees a half-written status snapshot.
- Fix: the periodic writer (daemon thread, never the training thread) flushes
  the collector before each snapshot, so a SIGTERM loses at most one flush
  interval instead of the whole run. There is no per-record fsync, which would
  turn the collector into cost noise.
- No new lock: the collector state read here is already guarded by the RT-10
  `_state_lock` and the `BoundedDedup` lock, so this change deliberately adds
  none.

---

# ✅ Tests

## Add the deliberate termination experiment

- `test_a_sigterm_loses_at_most_one_flush_interval` forks a child that delivers
  six intervals, waits for a flush, prints `ready`, then is SIGTERMed without
  `close()`; it asserts at least two envelope lines and a collector status
  survive.
- `21 passed` in `tests/utils/straggler/test_straggler_runtime.py`.
# 🎨 Style

## Clear the docformatter debt

- `pre-commit run --all-files` failed only on `docformatter`, which reflowed
  docstrings in 8 files; this commit is that tool's output, verbatim.
- Formatting-only: no logic, assertion or token changes.
New verifier-only suite (unique filename; no production file touched):

- abrupt death: SIGTERM and SIGKILL each lose at most one flush interval, never
  the whole run, and every persisted line is complete;
- a kill during a status write leaves only valid JSON;
- four simulated ranks write four unique runtime_status files, with exactly one
  writer for the unsuffixed names analyze_run.py reads;
- a concurrent reader never observes a partial status snapshot;
- output paths are run-scoped (with the punctuation-collision limit pinned);
- training-thread purity for both the host-only and the event-pool-exhausted
  ingest branches: zero file writes on the training thread;
- concurrent ingest conserves accepted/duplicate/late/malformed/evicted/
  closed-window tallies exactly;
- RT-09 staggered start and the A7 Megatron-free walker are re-verified

Two findings are pinned as strict xfails (VERIFY-01/02, reported separately):
close() can overlap the status writer (1.0 s join timeout, then it writes the
same files itself), and _write_json shares one '<final>.tmp' path so two
concurrent writers tear the published snapshot.
# 🐛 Bug Fix

## Make the two evidence writers race-free

- VERIFY-01: `close()` joined the status writer with a 1 s timeout and then
  wrote a snapshot itself, so a slow flush let both threads write the same
  `run_<id>` names (observed `max_concurrent == 2`). The final snapshot is now
  skipped when the writer is still alive after the join; its last write is at
  most one interval old.
- VERIFY-02: `_write_json` used one fixed `<final>.tmp`, so two writers tore it
  before `os.replace` (published snapshot invalid JSON 5/5 under a forced
  race). It now uses `tempfile.mkstemp` for a unique temp path, then
  `os.replace`, with unlink-on-error.
- VERIFY-03 was taken as the documented-limit branch: `_sanitise` stays
  non-injective for punctuation-only differences, and the `_run_dir` docstring
  now claims uniqueness only for real (alphanumeric) Ray job ids instead of
  "can never interleave".

---

# ✅ Tests

## Convert the verifier's strict xfails into passing assertions

- `test_close_does_not_race_the_status_writer` and
  `test_write_json_is_not_safe_for_two_concurrent_writers` keep their names and
  now pass against the fixed writer.
- `285 passed, 1 skipped`: `tests/utils/straggler` +
  `tests/utils/test_train_metric_utils.py`.
Rename the VERIFY-02 test to state the fixed contract and assert that neither
concurrent writer lost its own temp file, so a future regression to a shared
'<final>.tmp' fails even when the torn snapshot happens to parse.

Fail-before: with the pre-1e4bb6a fixed-temp implementation the same body raises
json.JSONDecodeError (Extra data); pass-after: 15 passed.
# 🐛 Bug Fix

## Restore the session after the A7 regression test

- Symptom: this file loaded the bridge converter from the Megatron checkout
  and left both the checkout on `sys.path` and its modules in `sys.modules`.
  `importorskip("megatron")` is satisfied from `sys.modules`, so 23
  pre-existing files then ran tests they normally skip; one of them
  (`test_metric_utils_rloo.py::test_eval_logger_disables_training_rloo_diagnostics`)
  failed on its own missing `transfer_queue` dep. Repo-wide outcomes therefore
  differed from baseline by one failure.
- Fix: prepend the checkout with `monkeypatch.syspath_prepend`, and add an
  autouse fixture that removes any `megatron*` module this file imported, so
  neither `sys.path` nor `sys.modules` leaks past the test.

---

# ✅ Tests

## Pin the isolation, not just the pass

- Same-session runs in both orders
  (`test_straggler_pickle_boundary.py` + `test_metric_utils_rloo.py` and the
  reverse) both give `12 passed, 1 skipped`, so the previously extra failure no
  longer appears.
# 🐛 Bug Fix

## Stop calling host jitter a straggler

- Defect, proven on a healthy run: the DP4 SFT ON smoke (injection knobs
  absent) produced 22 `straggler` verdicts, 29/33 with an observed magnitude
  under 5 ms. On a metadata stage a peer takes 0.20 ms and a rank 2.1 ms, which
  reads as 10x but is host/launch jitter, so the output is useless for
  localization.
- Measured inventory (gpu_campaign/on-smoke-final): metadata stages run
  0.060-3.965 ms (`params-all-gather` 0.906, `optimizer-inner-step` 1.220,
  `optimizer-copy-*` 1.981/2.414, `all-grads-sync` 3.965) while the real compute
  stages run 81.214 ms (`backward-compute`) and 133.521 ms (`forward-compute`).
- Fix: `RELAX_STRAGGLER_MIN_STAGE_MS` (default 5.0 ms, declared in env.py,
  documented with the rationale). A stage whose observed or peer-fastest
  magnitude is below the floor, or whose absolute gap does not exceed it, is
  reported `uncertain` with reason `below_absolute_floor` and counted in the new
  `sub_floor_judgements`; it is never a `straggler`. The relative tolerance and
  the `recovered` transition are unchanged.
- Fix: envelopes now carry `measurement_kind` (`device`/`host_only`) stamped at
  build time; the healthy smoke had it None on all 1155 envelopes.
- Evidence impact: the verdicts in `on-smoke-valid/` and `on-smoke-final/` were
  produced by the buggy detector and are NOT valid detector evidence; those runs
  stay on disk, labelled superseded.

---

# ✅ Tests

## Pin the floor and the stamp

- Sub-millisecond 10x jitter yields no straggler, only counted uncertainty.
- A `forward-compute`-sized stage still fires at +20 ms, +5.5% and +10%, while
  exactly +5% (the strict tolerance boundary) and an equal rank do not.
- The floor is config-driven (`min_stage_ms=0` restores the old behaviour) and
  the measured inventory is asserted to fall on either side of the default.
- The envelope stamp round-trips and validates; a legacy payload without it
  derives the kind.
- `290 passed, 1 skipped`: `tests/utils/straggler` +
  `tests/utils/test_train_metric_utils.py`.
# ✨ New Feature

## A measurement instrument for work comparability

- The detector could not distinguish "this rank does more work" from "this rank
  is slow", because every envelope shipped ``workload=None``. It now publishes
  this rank's LOCAL ``(tokens, sequences, microbatches)`` per optimizer step and
  stamps them on the envelope from the context on the straggler-readout thread.
- Tokens are a pure Python ``sum`` over the existing ``total_lengths`` list and
  its ``step_local_sample_counts`` slices; ``sequences`` is the local sample
  count. Publishing a rank-invariant figure such as ``step_global_batch_size``
  would make ``workload_delta`` identically zero and leave the gate inert, so
  only local per-rank values are published.
- Comparability gate: when a cohort's own work differs by more than
  ``work_tolerance``, the window is classified ``uncertain`` with reason
  ``workload_incomparable`` and counted in the new visible
  ``workload_incomparable_windows``. It is never dropped and never a straggler:
  an unequal-work cohort is itself the finding. The streak resets, so an
  existing a straggler still recovers normally.
- Coverage: the SFT prepack path (the DP4 vehicle) publishes. The RL main path,
  the streaming path and non-prepacked SFT do not yet, and ``workload`` stays
  absent there rather than guessed; this gap is disclosed in the PR/RFC.
- Justification: this lands as MEASUREMENT (it settles whether the ranks do
  equal work, closes RT-08, and serves criterion 3), explicitly NOT as a fix for
  the single healthy-run ``forward-compute`` rank3 verdict.

## Invariant A, demonstrated rather than asserted

- No training-thread I/O or socket: the whole publication path is exercised with
  ``open`` and ``socket.socket`` spied and neither fires.
- Source-level check over ``context.py``, the only training-thread file this
  adds: no ``all_reduce``/``all_gather``/``broadcast``/``synchronize``/tensor
  conversion/``torch.``/``dist.`` tokens.
- The new ``actor.py`` block passes the same token scan (0 matches), and the
  value is published once per rollout into a bounded dict (8 rollouts,
  evictions counted).

---

# ✅ Tests

## Comparability, per-rank values, boundedness, invariant A

- Over-tolerance cohort yields no straggler and increments the counter;
  within-tolerance cohorts still judge normally.
- Local counts are per-rank and do not leak between rollouts; an unpublished
  step reads as absent.
- The two RT-08 tests keep every ``workload_delta`` assertion and now assert the
  authorised classification (``uncertain``/``workload_incomparable``) plus a
  zero straggler count, since over-tolerance work no longer judges slowness.
- `296 passed, 1 skipped`: `tests/utils/straggler` +
  `tests/utils/test_train_metric_utils.py`.
… alone

Two correctness blockers in the Task 11 detector evidence.

P0-1: the per-step workload was published inside the SFT prefetch worker from
the rank-local k partition. Training executes the DP-wide max_k partition and
repacks when local_k < max_k, so the published metadata could describe a
partition that was thrown away. The publish now happens after the DP-wide K
decision and uses the final partition, re-derived with the same pure function,
and is skipped unless it is self-consistent with the executed K and the local
sample count.

P0-2: comparability summed tokens + sequences + microbatches, mixing units, so a
2% token gap could be cancelled by a 100% sequence gap. Comparability now uses
tokens alone -- the measure that dominates compute -- and the other fields are
reported as evidence via sequences_delta / microbatches_delta, with
workload_rank_tokens, workload_peer_tokens, tokens_delta and workload_comparable
carried in the verdict facts.

The frozen estimator, the detector thresholds and the training schedule are
untouched; the profiler still fails open and never raises into train().
P0-3: commit fb4d371 added a required `num_steps_per_rollout` parameter to
`train_one_step` so the straggler profiler could derive a run-wide global step
inside it. That widened the low-level training API for a profiler. The
parameter is removed and the signature is identical to fb4d371^ again; the
profiler is now called from `train()`, where rollout_id, step_id and
num_steps_per_rollout are already locals. The global-step arithmetic itself is
NOT ours and is not changed: it mirrors the platform identity that `train()`
already derives as `accumulated_step_id` (model.py), which upstream has
computed since the initial commit.

P0-1: add the missing regression test that forces local_k < max_k. It pins that
the rank-local and DP-wide partitions really do differ, that the published
workload is self-consistent with the executed max_k partition, and that the
single publish site sits after `local_k = max_k` inside the training-thread
path, not in the prefetch worker. The "57 workload_incomparable windows" stay
unusable as final detector evidence until a healthy ON carrier is re-run.

P1-4: the detector module docstring claimed a workload difference never
suppresses a verdict, while the code turns an over-tolerance token gap into
`uncertain` / `workload_incomparable`. The docstring now describes what the
code does, including that non-token fields are evidence only.
…a runtime harness

Splits the per-step workload derivation out of the Megatron actor path into
`relax/utils/straggler/workload.py`, a megatron-free pure function
`step_workloads(sample_lengths, k_partitions, num_steps_per_rollout)`. The
arithmetic that the detector depends on is now provable in the CPU venv, where
the actor module cannot even be imported.

`actor.py` keeps only a thin call at the single site after the DP-wide K is
final, passing the EXECUTED `max_k` (never `k_local`), still behind the cached
enablement gate and inside the non-raising handler.

Tests, all green:
- `test_straggler_workload.py` (CPU venv, 6 passed): one optimizer step consumes
  the window, so the payload is the step total `(180, 6, 3)`, not the first
  group `(90, 1, 1)`; equal totals such as [97,1,1,1] and [25,25,25,25] report
  equal tokens; impossible shapes raise instead of inventing a payload.
- `test_straggler_workload_publish_runtime.py` (training venv, 5 passed): drives
  the real `_get_prefetched_sft_window` with local_k=1 < max_k=3 and asserts the
  values read back for the step are the step totals. The fake all_reduce now
  answers by tensor shape and operation instead of a call counter, which
  desynchronised when the function ran twice in one process and fed the global
  batch size into the DP-wide K reduce; the out-of-range read is asserted against
  the real `(None, None, None)` contract. Skips in the CPU venv.
- the source-order guard's `publish > final_k` assertion is replaced by a
  semantic coupling check, because it passed on revert (the removed worker site
  had a HIGHER line number than the final-K line); it now asserts the single call
  site passes `max_k` and never `k_local`.

The detector algorithm, the frozen estimator and the disabled-path cost are
unchanged.
Cover all 23 timer names Relax's core path emits, not a sample, and pin the
strict exact-match fallback: an unrecognised name, including a suffix-shaped
one, stays ``other`` instead of being guessed into a measured group.

Also pin that attention and MoE remain schema-only and can never appear as a
measured group, that every taxonomy value is a measured group, that the
observed set has no duplicate, and that the reporter logs the raw timer name
next to its coarse stage (unclassified names as ``other``, and no label
fabricated for a verdict without a stage).
Removes the F841 in the disabled-path runtime test by making it assert what it
was meant to assert: with the profiler disabled the training path makes zero
calls into the pure derivation and publishes nothing. ruff check now passes on
every file in this change.

Records at the actor call site why the optimizer-step literal is 1 -- the single
caller passes a one-element prepared_num_microbatches list and train() derives
num_steps_per_rollout = len(num_microbatches) -- and what would make it wrong: a
second caller, or a window consumed by several optimizer steps, at which point
the argument must be that real step count.
ray-job.sh takes an exclusive flock on /root/autodl-tmp/relax-ray-gpu.lock
before the cluster-wide cleanup, so two submissions from this machine cannot
interleave. Design (the acquisition/cleanup/sidecar block itself landed in
d0b4eeb after a concurrent agent swept it into a straggler commit; this commit
adds the remaining safety probe):

- Acquisition sits right after mode detection + kernel-cache setup and before
  "Cleaning up residual Relax/SGLang worker processes", so the lock covers the
  whole cleanup + submit + wait window. Contention is non-blocking by default:
  it prints "another job holds the GPU lock", dumps the <lock>.holder sidecar
  (pid, ppid, host, started, cwd, project, run_script, submit_command) and exits
  75 (EX_TEMPFAIL) before touching the cluster. RELAX_GPU_LOCK_WAIT=<seconds>
  opts into a short bounded wait instead.
- The lock lives on fd 200. `exec bash <run-script>` inherits the descriptor, so
  the submitting process tree holds it across the wait. EXIT/INT/TERM/HUP traps
  release it and remove the sidecar on every shell-controlled path; SIGINT /
  SIGTERM / SIGHUP exit 130 / 143 / 129. SIGKILL still releases it because the
  kernel closes the descriptor.
- If the lock file does not exist yet `exec 200>>...` creates it without
  truncating, and it is never unlinked, so every participant locks one inode.
- Nested/repeated invocations cannot self-deadlock: the RELAX_ENTRYPOINT_MODE
  guard returns early, RELAX_GPU_LOCK_HELD marks an ancestor holder, and an
  inherited fd 200 pointing at the lock file is reused (re-flocking an inherited
  descriptor is a no-op). Even if all three were bypassed, acquisition is
  non-blocking, so it fails loudly instead of hanging.

The hole flock cannot close: a Ray driver is spawned by the raylet, not by the
submitting shell, so a SIGKILLed submitter drops the lock while its job keeps
RUNNING and the next submission's cleanup would stop it. After acquiring a free
lock we now list RUNNING jobs whose entrypoint contains
`relax.entrypoints.train` and warn on stderr before the cleanup runs. The probe
is warn-only and tolerates a missing/failing/slow Ray CLI (`|| true`), so it can
never block a legitimate submission; exit codes, argument handling and launch
semantics are unchanged.
…ndently

Comparability was evaluated pair-wide: one over-worked peer set a flag that
withheld every rank of the pair, so an equal-work genuine straggler (hosts
100/100/100/300 ms, tokens 1000/1200/1000/1000) was reported uncertain with its
own facts saying tokens_delta=0.0. The test is now per rank: a rank whose own
tokens are within work_tolerance of its peer median is judged even when a peer
is not, and the reason, workload_comparable and
workload_delta_beyond_tolerance are all derived from the same two-sided flag so
they cannot disagree.

The per-window workload was the last-arriving sample's tokens, which made the
gate depend on packet order in a time window that has no step order. It is now
the median of the rank's per-sample token readings: order-independent and an
estimate of per-step work, where a sum would confound work with how many steps
fitted in the window. The evidence-only sequences/microbatches deltas follow
the same rule.
… counters

global_step is no longer synthesised from rollout_id * num_steps_per_rollout +
optimizer_step: the platform's accumulated_step_id is documented as
non-monotonic under dynamic batching, and the no-length case used to emit the
bare in-rollout index as if it were global. global_step is now Optional and
stays None unless a caller supplies a genuinely run-wide value; optimizer_step
is typed and documented as in-rollout only; and the profiler publishes its own
strictly increasing, collision-free step_ordinal, pinned by a four-step then
one-step rollout test.

Observability: the training-context publish counters now have a production
caller. They and the workload gate counters are emitted in the reporter metrics
(perf/straggler/workload_publish_skipped|errors, workload_missing_windows,
workload_incomparable_windows) and embedded in runtime_status/collector_status
JSON, so a human sees them without a new request, socket or collective.

The reporter's dropped key now sums every detector eviction counter, and the
judged/envelopes fallback moved to its own judged_fraction key: it was
mislabelled coverage and read 1.0 while cohort coverage was far lower.

PP>1: the metrics-exporting Megatron primary rank (tp0, pipeline-last, dp0) owns
the collector only when pp_size == 1. Instead of exporting nothing, the reporter
emits perf/straggler/collector_status_available=0 for a runtime without a
collector.

Known limits documented in the reporter: RELAX_STRAGGLER_TOPOLOGY_EPOCH is inert
and there is no rollout/topology reset.
…s are judged

A timer interval took its wire sequence number at start (acquire) time, but the
token only reached the pending queue at stop time, and the readout thread
delivers in completion order. Megatron nests timers: the level-1
forward-backward timer wraps the level-2 forward-compute/backward-compute
timers, so the enclosing interval got the lowest sequence but completed last.
The collector requires a per-rank sequence that never decreases and drops a
packet whose sequence is behind the newest seen for its rank, so every
whole-phase interval was counted late and never judged: a straggler whose extra
time sits in the phase wrapper was invisible behind a late_packets counter that
reads like transport noise.

Stamp the sequence in complete_interval instead, so the sequence follows
completion order (which is also the order tokens are appended to _pending).
This is the minimal observer-only change: IntervalToken is a mutable __slots__
class, so the sequence can be assigned at stop time and the token need not be
rebuilt; no other code reads token.seq between acquire and complete, window
attribution uses host timestamps only, and the host-only path already allocated
its sequence at completion. Duplicate detection is unchanged (same token, same
sequence) and a genuinely stale resend still lands behind the newest sequence
for its rank, so it is still classified late.

Regression test drives the real shim, observer and collector together on CPU
with a fake CUDA-event backend: a nested pair for four ranks where one rank
spends an extra 200 ms only inside forward-backward. It asserts zero late
packets, every packet judged, and a host_only_stall straggler verdict on the
enclosing interval; it fails against the start-time sequencing with 24 late
packets (one enclosing interval per rank per step) and no straggler verdict.
The readout loop re-queued a not-yet-readable interval at the tail of
_pending. Because seq is stamped at completion, _pending is in sequence
order, so a deferred interval was shipped after later, higher-sequence
intervals and the collector's deduplicator discarded it as late.

On a real 4-GPU run this cost 18 of 48 forward-backward intervals -- the
longest interval, so the one most likely to be deferred -- while the short
intervals lost none, giving late=631 of 2112 envelopes with invalid=0 and
duplicate=0.

Re-queue at the head instead. _process already bounds the wait by shipping
an unreadable interval as readout_timeout past _readout_timeout_s, so the
head cannot block the queue indefinitely. The strict late contract, and the
seven tests that pin it, are unchanged.
…the bridge regressions

Two GitHub-CI-only test defects (product code untouched; the product
freeze at cac4cb6 stands):

- test_straggler_workload_publish_runtime.py failed collection on the
  CPU CI matrix: the transfer_queue importorskip guard does not cover
  runners that ship transfer_queue but no Megatron, and the actor
  module imports Megatron at collection time. An explicit
  importorskip("megatron") now skips there (verified: full pass under
  the training venv, clean module skip under the CPU venv).

- test_straggler_pickle_boundary.py crashed on the H20 GPU runner with
  TypeError: Path(None): a distribution-installed Megatron can be a
  namespace package whose module __file__ is None. _megatron_root now
  derives candidate roots from __path__ as well, and only accepts a
  root that actually contains megatron/bridge/models/conversion/utils.py
  (the rest of the behaviour, including the skip path, is unchanged).
…tries

English + Chinese user guide for the straggler profiler: what it is
(observe-only, fail-open, no collectives), how to enable it (off by
default), the full environment-variable table with defaults and
minimums, the four-stage architecture, the verdict schema, and the known
limitations (restart contract, late-packet loss, stale workload stamps,
inert topology epoch, no-CUDA judgement thread, no rollout reset).
Sidebar entries added to the bilingual VitePress config.

Docs-only change; the product freeze at cac4cb6 stands.
The speed reference and the workload comparability gate used different
peer sets: timing was judged against the fastest of ALL ranks in the
window while comparability was gated against the peer-median workload.
Whenever workloads differed, a peer doing less work -- fast because
under-worked, not because healthy -- became the speed baseline, and every
normal rank was reported as a straggler. Four ranks at {500 tokens /
50 ms, 1000 / 100 ms, 1000 / 100 ms, 1000 / 100 ms} with identical
per-token speed produced three straggler verdicts at ratio 2.0 with
workload_comparable=True (pinned by the new regression, which fails on
the previous build).

Each rank is now judged only against peers whose workload is pairwise
comparable with its own; a rank with no comparable peer is
workload_incomparable and its timing verdict is withheld, with the
all-ranks figures retained as labelled evidence rather than a conviction
baseline. A rank whose peers report no workload at all stays DEGRADED
(judged, comparable=None), preserving the documented contract. Verdict
facts now expose reference_ranks and comparable_peers so a verdict says
which peers it was actually judged against.
…d at completion

Two delivery-correctness defects in the observer, both demonstrated by
regressions that fail on the previous build:

1. Ordering. A host-only interval (event pool exhausted while the readout
   thread is live) was delivered inline on the training thread, bypassing
   the pending queue. With a device interval still waiting for event
   readback at the head of the queue, the later host-only envelope was
   shipped first; the collector saw a sequence regression and discarded
   the device envelope as late -- the same failure shape the tail
   re-queue caused, reached through the degradation path. A finished
   host-only envelope now travels through the same queue (_ReadyEnvelope)
   so wire delivery preserves the completion-stamped sequence. The pure
   host-only deployment (no device timing, no readout thread) still
   delivers inline, as documented. As a side effect the degraded path no
   longer runs consumer work on the training thread at all, which the
   RT-01 red-team test now asserts directly.

2. Workload binding. The envelope's workload context was read at
   envelope-build time on the readout thread, so whenever the readback
   lagged the training loop the envelope married step N's timing to step
   N+1's workload -- exactly the evidence the comparability gate
   consumes. The workload is now captured on the training thread at
   interval completion and carried on the token; the readout thread only
   reads a value that was already fixed when the interval closed.

The observer-test sync assertion for pool-exhausted host-only intervals
now waits for delivery (arrival moved from synchronous to queued; the
assertion strength -- three distinct sequence numbers -- is unchanged).
…rsists, effective cohort minimum

Three boundary completions (no feature growth), each pinned by regressions
that fail on 3c1376b:

1. Degenerate workload values are degraded evidence, never poison.
   _workload_metric now accepts a count field only when it is a real number
   (booleans rejected despite bool subclassing int), finite and strictly
   positive. Zero, negative, NaN/Inf, boolean and missing values read as
   absent: the rank keeps its timing verdict with workload_comparable=None
   and the window survives. The pairwise comparability introduced in 9275809
   divides by peer token counts; a peer reporting tokens=0 raised
   ZeroDivisionError inside the judge and lost the whole window (pinned by
   test_straggler_workload_sanitization, which fails on 3c1376b).

2. A completion-time None workload stays None forever. _build_envelope no
   longer falls back to reading the context at build time: that fallback
   silently re-bound envelopes whose completion-time capture was None to
   whatever step was current on the readout thread — the cross-step
   contamination the binding fix exists to prevent, reached through the None
   path. The host-only path captures explicitly at completion; token paths
   carry token.workload. test_straggler_workload_binding adds None-stays-None
   and captured-value-survives-disappearance regressions.

3. min_cohort_size applies to the EFFECTIVE comparable class, not the raw
   reporting count. judged_slow now also requires
   len(reference_ranks[rank]) >= min_cohort_size, and a slow rank inside a
   smaller comparable class emits cohort_below_min_size uncertain instead of
   a conviction: four ranks split into two comparable pairs with
   min_cohort_size=4 no longer convict either pair, while min_cohort_size=2
   still detects within a pair. Verdict facts distinguish raw topological
   coverage (coverage_ratio) from effective reference coverage
   (comparable_class_size, comparable_coverage_ratio).

WHY_UNFREEZE (measurement-correctness, post-3c1376b): each fix changes
detector semantics or envelope evidence, so every future C1/C2/C3 run must be
produced by this SHA or later; no already-published acceptance number depended
on the removed behaviours (the cac4cb6/3c1376b arms verified delivery order
and platform wiring, not final acceptance). Suite: 370 passed / 2 skipped
under CI conditions, 376 passed with the training stack.
…ws hold the alert

Reproduced on ef6516e (minimal in-memory case, archived in the commit that
follows on the evidence branch): a window in which a flagged rank's evidence
went unusable (own workload incomparable, or effective comparable class
below the minimum) reset the streak, emitted a 'within_tolerance' RECOVERED
verdict it never earned, and cleared the active alert — sometimes in the same
window that also emitted an uncertain verdict for the same rank.

Recovery now requires positive evidence: a comparable measurement back within
tolerance. An unusable window HOLDS the streak and the active flag and emits
the explicit 'recovery_evidence_insufficient' uncertain (new counter
recovery_evidence_withheld). Onset likewise requires positive evidence, so a
held streak cannot re-fire the onset count. The floor gate no longer blocks
recovery: a rank back at par has no gap at all, and its ~0 absolute delta is
always below the floor, which would have made recovery unreachable.

Regressions (fail on ef6516e): flagged -> incomparable window -> evidence
returns still-slow (alert held) -> genuinely at par (NOW recovers); small
effective class holds the alert; not-slow-but-unusable reports unknown, not
recovered; withheld counter consistency. Suite: 374 passed / 2 skipped under
CI conditions.
@shanyulu

This comment has been minimized.

@shanyulu shanyulu changed the title 【No.11】feat(straggler): add a default-off failure-isolated rank observer 【No.11】feat(straggler): add opt-in Megatron straggler profiling Sep 27, 2026
# 🐛 Bug Fix

- Require measured relative recovery; sub-floor slow evidence cannot clear an alert.
- Count consecutive window indices and reset onset persistence across unusable evidence.
- Retain bounded onset markers across active-map eviction; count withheld recovery once.

# ✅ Tests

- Add old-fail/new-pass recovery, coverage-gap and counter regressions.
- Cover duplicate onset after active-cap eviction and explicit tail-window flushing.
- Validate the full straggler suite in the installed training stack: 392 passed.
- Run pre-commit across all tracked files.
# 🛠️ fix
- Add opt-in safe submission that bypasses global process, job, Serve and placement-group cleanup.
- Require a verified single local Ray node and complete idle inventories; refuse unknown or busy states.
- Validate 21 guard tests and the actual idle dashboard without changing cluster resources.
# 🛠️ fix
- Separate background JSONL persistence from terminal detector finalization.
- Fence already-closed event windows per cohort in the bounded alignment map; count late samples without re-judging them.
- Preserve distinct cohort anchors and document the silent-tail limitation.

# 🧪 test
- Reproduce premature window closure and false recovery before the fix.
- Verify repeated live-window writes retain three-window onset persistence.
- Training-stack regression: 416 passed including 21 launcher guard tests; full pre-commit passed.
# ✅ Tests

- Forward explicit training overrides through the real SFT recipe.
- Carry worker proxy settings and a verified GCS address into the runtime environment.
- Allow unique submission IDs for bounded, ownership-scoped supervision.
- Verify the dry-run arguments and worker environment; full pre-commit passes.
# 🐛 Bug Fix

- Add confirmed-alert metrics without changing raw worst-value semantics.
- Exclude uncertain, recovered and invalid measurements from confirmed selection.
- Include rank, window and reason in existing stage records.

# ✅ Tests

- Reproduce uncertain-rank masking before the fix.
- Pass all 398 straggler tests with the training stack available.

# 📝 Documentation

- Document raw and confirmed metrics in English and Chinese.
@shanyulu shanyulu changed the title 【No.11】feat(straggler): add opt-in Megatron straggler profiling 【No.11】Megatron 慢卡诊断与平台上报 Sep 27, 2026
# 🐛 Bug Fix

- Return collector status without synchronous summary logging.
- Remove per-verdict text output from the training metrics path.
- Preserve bounded verdict JSONL and explicit diagnostic logging.

# ✅ Tests

- Verify training ingest and metric export do not call log handlers.
- Keep explicit stage rendering and report coverage.
- Pass 399 straggler tests with the training stack.

# 📝 Documentation

- Document output-directory requirements and shared-lock limits in both languages.
ingest() held _state_lock across the whole ingest path, and the batch
flush (_append -> _flush_path -> _write_batch) plus the external
on_verdict callback both ran inside that section. A disk write that
stuck, or a callback that blocked, therefore froze every concurrent
summary()/status() read -- the training thread's per-rollout metrics
path. 4654e5a removed read-side logging but left this untouched; a
controlled write-block injection still reproduced the hang.

The locked section now only buffers: _append records dirty paths
instead of flushing, and _handle collects verdicts for deferred
callbacks. ingest()/flush() drain both after releasing the lock
(mirroring flush()'s existing swap-under-lock, write-outside pattern),
with failures counted as write_errors so ingest() keeps its
never-raises contract and no packet loses its detector observation.

Old-fail/new-pass regressions (3 failed on the pre-fix collector,
pass after): blocked batch write, persistence raise, blocked verdict
callback -- each asserts summary() completes within a bound from
another thread.
…tion/MoE scope note

The profiler guide (en/zh) described the pre-fix behavior: it claimed
the workload is looked up at envelope-build time on the readout thread
and can therefore belong to a later step. The code binds the workload
in complete_interval on the training thread, at interval completion
(observer.py: the envelope belongs to the step that closed the
interval), and a completion-time missing workload stays missing -- it
is never backfilled. Both guides now state the completion-time
binding, the no-backfill rule and their workload_missing_windows
consequence.

stages.py claimed 'the mentor already agreed attention/MoE may be
reserved behind a flag for the first phase'. No such agreement is on
record; the scope question is open as Decision B in RFC redai-studio#357. The note
now says exactly that.
shanyulu added a commit to shanyulu/Relax that referenced this pull request Sep 27, 2026
Product head moved a48a23b -> 927c5de (293390e confirmed-alert metrics,
4654e5a quiet summary reads, f120aa8 state-lock isolation, docs fixes),
pushed to PR redai-studio#378 with the affected-acceptance declaration. Per
REVIEW_CORRECTIONS items 2-4, preregister BEFORE any new data:

- C2_PARAMETER_PROTOCOL_927C5DE.md: retained checkpoints, per-tensor
  key/shape/dtype + numeric metrics, tolerance frozen from OFF/OFF
  calibration contrasts (unit = one arm-pair contrast; the withdrawn
  720-sample phrasing is explicitly superseded), lock-before-see;
- C3_EVENT_CHAIN_PROTOCOL_927C5DE.md: event-identity latency chain,
  preregistered clock-comparability check, tail-window evidence
  requirements, fixed non-target-rank classification rule, platform
  attribution regression for confirmed_straggler_* metrics;
- OVERLAP_CALIBRATION_PROTOCOL_927C5DE.md: independent calibration
  session at the new build, own frozen envelope, fixed N one-shot, old
  a48a23b NOT_PASS preserved as version-pinned history.

README gains a post-review addendum indexing the corrections, the
strict gates, the raw-log and trace archives, and these protocols.
REVIEW_CORRECTIONS_20260927.md is re-aligned by mdformat (table padding
only; content unchanged from f69343f).

Hook transparency: the gitleaks_tracked hook (always_run, whole-tree)
fails on this branch with 93 findings that are identical at afd45fb —
pre-existing archive-grade raw evidence (RFC1918 addresses in raw
job.logs and early campaign notes) deliberately published under the
two-tier scheme (publication gate = scan_public_evidence + sanitized
publics, not this hook). The four new .md files and all session .py
files pass every file-scoped hook (ruff/ruff-format/docformatter/
mdformat/gitleaks findings attributable to new files: 0).

This branch has not been deployed

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants