From d5f09506dbace283a4379f3140e6da36deae2899 Mon Sep 17 00:00:00 2001 From: aloekun Date: Tue, 21 Jul 2026 10:35:33 +0900 Subject: [PATCH] =?UTF-8?q?fix(pipeline-lock):=20stale=20takeover=20?= =?UTF-8?q?=E3=81=AE=E5=90=8C=E6=99=82=E5=8F=96=E5=BE=97=E3=82=92=E6=8E=92?= =?UTF-8?q?=E9=99=A4=E3=81=97=20race=20=E3=83=86=E3=82=B9=E3=83=88?= =?UTF-8?q?=E3=81=AE=E5=81=BD=E9=99=B0=E6=80=A7=E3=82=82=E4=BF=AE=E6=AD=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit pipeline lock (push/merge の二重起動防止) の stale takeover が、高競合下で 2 スレッドとも Acquired になり得た。`concurrent_stale_takeover_only_one_wins` テストが Windows で ~10% flaky に落ちていた症状の根本原因。 ## 実測で確定させた真因 (2 段階) Windows + WSL Linux の双方でシーケンス番号付きトレースを取り、憶測ではなく 観測で 2 つの独立した欠陥を特定した。 ### 欠陥 1: 本番コードの TOCTOU (2 スレッド版テストの本物の race) 旧 takeover は `read → 再比較 → remove_file → create_new`。**再比較と remove が 非アトミック**で、A が再比較を通過した直後に B が takeover を完走 (fresh lock 作成) → A の remove が B の fresh lock を破壊 → A も create_new 成功で 2 スレッド とも Acquired。旧 doc の「128bit token 偶然一致が必要=無視できる」は誤りで、 通常のスケジューリング窓だった。 修正は 2 重の防御: - **sentinel** (`.takeover` の create_new) で「takeover を試みるスレッド」を 1 つに直列化。 - sentinel 保持者は remove+create ではなく **temp へ全内容を書いてから rename で path を atomic 置換**する。path が一度も不在にならないため、他スレッドの 親 fast-path create_new が割り込めず、読み手が空/部分書き込みを観測する窓も消える。 - create_new 直後・write_all 前の空ファイルを stale と誤判定しないよう、lock content を Fresh / Held(空=書き込み中) / Stale の 3 値に分類 (姉妹 lock.rs が WP-15 で塞いだ bug class と同型)。 sentinel だけ / rename だけ では不十分だった経緯 (Linux 高競合で PARENT 経路と TAKEOVER 経路が各 1 Acquired) は関数 doc に実測ログ付きで記録。 ### 欠陥 2: 追加した stress テスト自身の偽陽性 8 スレッド stress テストの初版は `map(join).filter().count()` の遅延イテレータで 数えており、**まだ acquire 中の他スレッドの傍らで先行結果が drop** され、その `PipelineLock::drop` が lock を削除 → 走行中スレッドが正当に再取得し「2 Acquired」に 見えた (同時保持ではなく解放後の再取得 = テストアーティファクト)。全ガードを Vec に collect してから数える形に修正し、「同時点で 2 つ保持され得るか」だけを検証する。 ## 検証 - 2 スレッド版 (両ガード保持 = 本物の race を突く) は master で 3/30 失敗、本修正で Windows/Linux とも多数回 pass。 - 8 スレッド stress は旧本番コードで確実に失敗、本修正で Windows 30/30・Linux 30/30 pass。 - 公開 API (hold_pipeline_lock / acquire_pipeline_lock) は不変で呼び出し元は無改修。 - cargo test --workspace 全 pass / clippy clean を Windows + WSL Linux で確認。 --- src/lib-jj-helpers/src/pipeline_lock.rs | 540 ++++++++++-------- src/lib-jj-helpers/src/pipeline_lock/tests.rs | 337 +++++++++++ 2 files changed, 635 insertions(+), 242 deletions(-) create mode 100644 src/lib-jj-helpers/src/pipeline_lock/tests.rs diff --git a/src/lib-jj-helpers/src/pipeline_lock.rs b/src/lib-jj-helpers/src/pipeline_lock.rs index 401ff42e..7d319ef9 100644 --- a/src/lib-jj-helpers/src/pipeline_lock.rs +++ b/src/lib-jj-helpers/src/pipeline_lock.rs @@ -11,12 +11,27 @@ //! - age ベースの stale 判定 + takeover (クラッシュした pipeline の lock が永続しない) //! - RAII guard (Drop で削除) //! +//! **相互排他の要件が姉妹 lock と異なる**: `cli-pr-monitor/src/lock.rs` は「複数 takeover が +//! 同時成功しても無害」(last-write-wins) を許容するが、本 lock は pipeline の二重起動を +//! 防ぐため **stale takeover でちょうど 1 プロセスのみ `Acquired`** を要件とする。この単一 +//! winner 性は `.takeover` sentinel の `create_new` で「path を除去し得るスレッド」を +//! 1 つに限定して担保する (`takeover_stale_lock` の doc に詳細)。 +//! //! 相違点: timestamp は ISO8601 ではなく unix epoch 秒を直接記録する (parser 不要)。 //! future timestamp は stale 扱い (破損 lock が永続 fresh 化する bug class の再発防止、 //! lock.rs の PastTime と同じ invariant)。 //! //! ファイル形式は `key=value` 行 (pid / start_unix / label)。外部 config ではなく //! 内部の一時ファイルのため、依存追加 (serde/toml) を避けた最小形式とする。 +//! +//! sentinel 自身も age ベースで自己修復する (SIM-NEW-pipeline_lock-L157): sentinel +//! 保持者が `perform_takeover` 実行中にクラッシュすると sentinel が孤立し、以降 +//! 本物の stale lock があっても永久に `Busy` へ倒れていた。`classify_lock_content` +//! を流用して stale (`SENTINEL_STALE_SECS` 経過) と判定できた sentinel は、その +//! content 固有の reclaim gate (`reclaim_gate_path`) 経由で回収する。単純な +//! 「stale 判定 → 無条件 remove」は、判定から除去までの間隙で他スレッドが正当に +//! 再確立した sentinel を巻き添えで消しうる (8 スレッド高競合の実測で 2 `Acquired` +//! が再現した regression) ため不採用とした。詳細は `reclaim_stale_sentinel` の doc。 use std::fs::OpenOptions; use std::io::Write; @@ -104,17 +119,24 @@ pub fn acquire_pipeline_lock_at( PipelineLockResult::Acquired(PipelineLock { path, token }) } Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => { - let snapshot = std::fs::read_to_string(&path).ok(); - if let Some(raw) = &snapshot { - if let Some((pid, age_secs)) = is_fresh_content(raw, stale_threshold_secs, now_unix) - { - return PipelineLockResult::Busy { - holder_pid: pid, - holder_age_secs: age_secs, - }; + if let Ok(raw) = std::fs::read_to_string(&path) { + match classify_lock_content(&raw, stale_threshold_secs, now_unix) { + LockState::Fresh(pid, age_secs) => { + return PipelineLockResult::Busy { + holder_pid: pid, + holder_age_secs: age_secs, + }; + } + LockState::Held => { + return PipelineLockResult::Busy { + holder_pid: 0, + holder_age_secs: 0, + }; + } + LockState::Stale => {} } } - takeover_stale_lock(path, token, content, snapshot, stale_threshold_secs, now_unix) + takeover_stale_lock(path, token, content, stale_threshold_secs, now_unix) } Err(e) => PipelineLockResult::Unavailable { reason: e.to_string(), @@ -122,48 +144,255 @@ pub fn acquire_pipeline_lock_at( } } -/// stale と判定した lock を takeover する。`create_new` により先着 1 プロセスのみ -/// `Acquired` になることを保証する (CodeRabbit re-review Major 対応)。 +/// stale と判定した lock を takeover する。ちょうど 1 プロセスのみ `Acquired` になる。 /// -/// `remove_file` 直前に `stale_snapshot` (呼び出し元が stale と判定した時点の生 content) -/// と現在の content を再比較し、一致する場合のみ削除する。他プロセスが同じ間隙で先に -/// takeover 済み (= content が変化済み) なら削除をスキップし、後続の `create_new` が -/// 自然に `AlreadyExists` で失敗して `busy_from_disk` に落ちる。無条件 `remove_file` だと -/// 先着プロセスが作った fresh lock を後発側が検証なしに消してしまいうる -/// (2 プロセスとも `Acquired` になる実際に再現した regression)。 +/// **単一 takeover 権を sentinel の `create_new` で選出する**。`.takeover` を +/// atomic な `create_new` で作れたスレッドだけが takeover 権限を持ち、実際の +/// 「stale 判定 → 除去 → fresh 作成」を行う (`perform_takeover`)。sentinel 取得に +/// 負けたスレッドは `Busy` に倒す (権限保持者が fresh lock を設置するため)。 /// -/// 残余 TOCTOU: この再比較と `remove_file` 呼出の間隙のみ (128bit token を含む同一 -/// content が別プロセスにより偶然この一瞬で再現される確率は無視できる)。本 lock は -/// advisory (fail-open, ADR-043) であり、この残余は許容する。 +/// sentinel は「takeover を試みるスレッド」を 1 つに絞る第 1 段の直列化。実際の単一 winner +/// 性は sentinel 保持者が [`perform_takeover`] で **rename による atomic 置換**を使い、path を +/// 一度も不在にしないことで担保する (詳細はそちらの doc)。sentinel だけでは不十分だった経緯 +/// (remove + create_new の不在窓に親 fast-path が割り込む) も同 doc に記録。 +/// SIM-NEW-pipeline_lock-L146。 +/// +/// sentinel 自体の取得は [`acquire_takeover_sentinel`] に委譲する (age ベースの +/// 自己修復込み、SIM-NEW-pipeline_lock-L157)。 fn takeover_stale_lock( path: PathBuf, token: String, content: String, - stale_snapshot: Option, stale_threshold_secs: i64, now_unix: i64, ) -> PipelineLockResult { - if std::fs::read_to_string(&path).ok() == stale_snapshot { - if let Err(e) = std::fs::remove_file(&path) { - if e.kind() != std::io::ErrorKind::NotFound { - eprintln!("[pipeline-lock] takeover 時の remove 失敗 (継続): {}", e); + let sentinel = takeover_sentinel_path(&path); + match acquire_takeover_sentinel(&sentinel, now_unix) { + SentinelGate::Acquired => { + let result = + perform_takeover(&path, token, &content, stale_threshold_secs, now_unix); + if let Err(e) = std::fs::remove_file(&sentinel) { + if e.kind() != std::io::ErrorKind::NotFound { + eprintln!("[pipeline-lock] takeover sentinel の除去に失敗 (継続): {}", e); + } } + result } + SentinelGate::Busy => busy_from_disk(&path, stale_threshold_secs, now_unix), + SentinelGate::Unavailable(reason) => PipelineLockResult::Unavailable { reason }, } - match OpenOptions::new().write(true).create_new(true).open(&path) { - Ok(mut f) => { - if let Err(e) = f.write_all(content.as_bytes()) { - eprintln!("[pipeline-lock] takeover 書き込み失敗 (継続): {}", e); - } - PipelineLockResult::Acquired(PipelineLock { path, token }) +} + +/// sentinel の stale 判定 threshold。sentinel は `perform_takeover` 一回分の実行区間 +/// (read/write/rename 数回、通常ミリ秒オーダー) だけ保持される想定であり、pipeline +/// 本体の `PIPELINE_LOCK_STALE_SECS` (1800s) よりずっと短くてよい。sentinel 保持者が +/// takeover 中にクラッシュしても、この秒数が経てば次の acquire が自己修復する +/// (SIM-NEW-pipeline_lock-L157: 従来は age 判定が皆無で永久 wedge していた)。 +const SENTINEL_STALE_SECS: i64 = 30; + +/// [`takeover_stale_lock`] 用 sentinel 取得結果。 +enum SentinelGate { + /// sentinel を確保した (自分が takeover 実行権を持つ)。 + Acquired, + /// 別プロセスが sentinel を保持中 (fresh、または create_new 直後の書き込み待ち)。 + Busy, + /// I/O エラーで判定不能。 + Unavailable(String), +} + +/// `.takeover` sentinel の取得を試みる。通常は `create_new` の atomic 排他により +/// 1 プロセスのみ成功する。既存 sentinel が stale と判定できた場合に限り、 +/// [`reclaim_stale_sentinel`] へ委譲して回収を試みる (SIM-NEW-pipeline_lock-L157)。 +fn acquire_takeover_sentinel(sentinel: &Path, now_unix: i64) -> SentinelGate { + match try_create_sentinel(sentinel, now_unix) { + Ok(()) => return SentinelGate::Acquired, + Err(e) if e.kind() != std::io::ErrorKind::AlreadyExists => { + return SentinelGate::Unavailable(e.to_string()); } + Err(_) => {} + } + + let Ok(raw) = std::fs::read_to_string(sentinel) else { + return SentinelGate::Busy; + }; + match classify_lock_content(&raw, SENTINEL_STALE_SECS, now_unix) { + LockState::Held | LockState::Fresh(..) => SentinelGate::Busy, + LockState::Stale => reclaim_stale_sentinel(sentinel, &raw, now_unix), + } +} + +/// sentinel に `pid=` / `start_unix=` を書き込む。`classify_lock_content` (メイン lock と +/// 共通ロジック) で stale 判定できるよう、フィールド名をメイン lock の形式と揃える。 +fn try_create_sentinel(sentinel: &Path, now_unix: i64) -> std::io::Result<()> { + let mut f = OpenOptions::new() + .write(true) + .create_new(true) + .open(sentinel)?; + f.write_all(format!("pid={}\nstart_unix={}\n", std::process::id(), now_unix).as_bytes()) +} + +/// stale と判定した sentinel を、その content 固有の reclaim gate 経由で回収する。 +/// +/// **単純な「stale 判定 → 無条件 remove」は不採用**: 判定から除去までの間隙で別スレッドが +/// 正当に再確立した sentinel を巻き添えで消しうる (8 スレッド高競合の実測で 2 `Acquired` +/// が再現した regression、SIM-NEW-pipeline_lock-L157 修正の初版で発見)。 +/// +/// 代わりに、reclaim gate の path を **読んだ content から決定論的に導出**する +/// (`reclaim_gate_path`)。同じ stale content を読んだスレッド同士だけがその 1 つの +/// `create_new` で競い、勝者だけが実際の「除去 → 再作成」(`finish_reclaim`) を行う。 +/// 負けたスレッドは「対象はまだ同じ stale content のまま (勝者が処理中)」と分かっている +/// ため、sentinel 自体には一切触れず安全に `Busy` へ倒せる。 +fn reclaim_stale_sentinel(sentinel: &Path, stale_content: &str, now_unix: i64) -> SentinelGate { + let reclaim = reclaim_gate_path(sentinel, stale_content); + match try_create_reclaim_marker(&reclaim, now_unix) { + Ok(()) => finish_reclaim(sentinel, &reclaim, now_unix), Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => { - busy_from_disk(&path, stale_threshold_secs, now_unix) + reap_orphaned_reclaim_marker(&reclaim, now_unix); + SentinelGate::Busy + } + Err(e) => SentinelGate::Unavailable(e.to_string()), + } +} + +/// reclaim gate の勝者だけが呼ぶ: sentinel を除去して再作成する。 +/// +/// この関数を呼べるのは「この stale content 専用の reclaim gate」に勝った 1 スレッドのみ +/// (同じ content を読んだ他スレッドは gate 負けで `Busy` に倒れ、sentinel に触れない)。 +/// 唯一の例外は「この stale content をまだ読んでいない、独立した新規取得試行」が +/// 除去直後の空隙で fast path の `create_new` に成功するケースで、その場合は素直に +/// `Busy` へ倒れる (単一 winner 性は保たれる)。 +fn finish_reclaim(sentinel: &Path, reclaim: &Path, now_unix: i64) -> SentinelGate { + let _ = std::fs::remove_file(sentinel); + let result = match try_create_sentinel(sentinel, now_unix) { + Ok(()) => SentinelGate::Acquired, + Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => SentinelGate::Busy, + Err(e) => SentinelGate::Unavailable(e.to_string()), + }; + let _ = std::fs::remove_file(reclaim); + result +} + +/// reclaim gate 自体に `pid=` / `start_unix=` を書き込む (sentinel と同形式、 +/// `classify_lock_content` を再利用するため)。 +fn try_create_reclaim_marker(reclaim: &Path, now_unix: i64) -> std::io::Result<()> { + let mut f = OpenOptions::new() + .write(true) + .create_new(true) + .open(reclaim)?; + f.write_all(format!("pid={}\nstart_unix={}\n", std::process::id(), now_unix).as_bytes()) +} + +/// reclaim gate 保持者が `finish_reclaim` の 2 手 (remove + create_new) の間でクラッシュし、 +/// gate 自体が孤立した場合の回収。 +/// +/// この gate は特定の stale content 専用 (path が content 由来のため) で、除去後に +/// **別の正当な世代が同じ path に再作成される余地が無い**。よって sentinel と違い、 +/// stale 判定さえできれば無条件 remove で安全に回収できる (複数スレッドが同時に +/// 回収を試みても、`remove_file` 自体の排他性により実害は出ない)。 +fn reap_orphaned_reclaim_marker(reclaim: &Path, now_unix: i64) { + if let Ok(raw) = std::fs::read_to_string(reclaim) { + if matches!( + classify_lock_content(&raw, SENTINEL_STALE_SECS, now_unix), + LockState::Stale + ) { + let _ = std::fs::remove_file(reclaim); + } + } +} + +/// sentinel content から決定論的に reclaim gate の path を導出する。 +/// 同じ (sentinel path, content) を読んだスレッドは必ず同じ path に collide する +/// (`DefaultHasher` は固定キーで決定論的、`RandomState` と異なる)。 +fn reclaim_gate_path(sentinel: &Path, content: &str) -> PathBuf { + suffixed_path(sentinel, &format!(".reclaim.{:016x}", content_fingerprint(content))) +} + +fn content_fingerprint(content: &str) -> u64 { + use std::hash::{Hash, Hasher}; + let mut hasher = std::collections::hash_map::DefaultHasher::new(); + content.hash(&mut hasher); + hasher.finish() +} + +/// sentinel 保持者だけが呼ぶ takeover 本体。stale なら新 lock を **atomic な rename で置換**し、 +/// `Acquired` を返す。fresh / 書き込み中 (`Held`) なら奪わず `Busy`。 +/// +/// **rename で置換する (remove + create_new にしない)** のが単一 winner 性の要: +/// remove + create_new は remove と create の間で **path が一瞬不在になる窓**を作り、その窓で +/// 別スレッドの親 fast-path `create_new(path)` が成功して 2 スレッドとも `Acquired` になる +/// (Linux 高競合で実測: `TAKEOVER_REMOVE → TAKEOVER_CREATE_OK → PARENT_CREATE_OK`)。 +/// temp に全内容を書いてから `rename(temp → path)` すれば、path は旧 stale 内容から新内容へ +/// **原子的に切り替わり一度も不在にならない**ため、親の `create_new` は常に `AlreadyExists` に +/// なり fast-path で割り込めない。同時に「読み手が空 / 部分書き込みを観測する」窓も消える +/// (rename は完成済みファイルを一括で差し込む)。SIM-NEW-pipeline_lock-L146。 +fn perform_takeover( + path: &Path, + token: String, + content: &str, + stale_threshold_secs: i64, + now_unix: i64, +) -> PipelineLockResult { + if let Ok(raw) = std::fs::read_to_string(path) { + match classify_lock_content(&raw, stale_threshold_secs, now_unix) { + LockState::Fresh(pid, age_secs) => { + return PipelineLockResult::Busy { + holder_pid: pid, + holder_age_secs: age_secs, + }; + } + LockState::Held => { + return PipelineLockResult::Busy { + holder_pid: 0, + holder_age_secs: 0, + }; + } + LockState::Stale => {} } - Err(e) => PipelineLockResult::Unavailable { - reason: e.to_string(), - }, } + replace_lock_atomically(path, token, content) +} + +/// 新 lock 内容を temp ファイルへ書き、`rename` で `path` に atomic 置換する。 +/// sentinel 保持中のみ呼ばれるため temp パスの競合は起きない。 +fn replace_lock_atomically(path: &Path, token: String, content: &str) -> PipelineLockResult { + let tmp = takeover_tmp_path(path, &token); + if let Err(e) = std::fs::write(&tmp, content) { + return PipelineLockResult::Unavailable { + reason: format!("takeover temp 書き込み失敗: {}", e), + }; + } + match std::fs::rename(&tmp, path) { + Ok(()) => PipelineLockResult::Acquired(PipelineLock { + path: path.to_path_buf(), + token, + }), + Err(e) => { + let _ = std::fs::remove_file(&tmp); + PipelineLockResult::Unavailable { + reason: format!("takeover rename 失敗: {}", e), + } + } + } +} + +/// takeover 権選出用の sentinel パス (`.takeover`)。元 lock と同一ディレクトリ。 +fn takeover_sentinel_path(path: &Path) -> PathBuf { + suffixed_path(path, ".takeover") +} + +/// atomic 置換用の temp パス (`.new.`)。`rename` の atomic 性を保つため +/// 元 lock と同一ディレクトリに置く。 +fn takeover_tmp_path(path: &Path, token: &str) -> PathBuf { + suffixed_path(path, &format!(".new.{token}")) +} + +fn suffixed_path(path: &Path, suffix: &str) -> PathBuf { + let mut name = path + .file_name() + .map(|s| s.to_os_string()) + .unwrap_or_default(); + name.push(suffix); + path.with_file_name(name) } /// takeover レースに負けた際、ディスク上の現在の holder 情報から `Busy` を組み立てる。 @@ -230,8 +459,40 @@ fn read_fresh_lock(path: &Path, stale_threshold_secs: i64, now_unix: i64) -> Opt is_fresh_content(&content, stale_threshold_secs, now_unix) } +/// lock ファイル content の 3 値判定。 +/// +/// **`Empty` を `Stale` と別扱いするのが要点**: `create_new` は atomic だがその直後の +/// `write_all` までにファイルは**空**で存在する。この窓の空 content を stale と誤判定して +/// takeover (= remove) すると、`create_new` に成功して自分を holder と見なした別スレッドの +/// lock を破壊し、2 スレッドとも `Acquired` になる (8 スレッド高競合の Linux 実測で顕在化。 +/// 姉妹 `cli-pr-monitor/src/lock.rs` が WP-15 で塞いだのと同じ bug class)。空は「書き込み中の +/// 保持者あり」= `Held` に倒す。 +enum LockState { + /// 有効期限内の holder が居る (pid, age)。 + Fresh(u32, i64), + /// `create_new` 直後・`write_all` 前の空ファイル。保持者が書き込み中とみなす。 + Held, + /// 破損 / 期限切れ / 未来日付。takeover 可。 + Stale, +} + +fn classify_lock_content(content: &str, stale_threshold_secs: i64, now_unix: i64) -> LockState { + if content.trim().is_empty() { + return LockState::Held; + } + match is_fresh_content(content, stale_threshold_secs, now_unix) { + Some((pid, age_secs)) => LockState::Fresh(pid, age_secs), + None => LockState::Stale, + } +} + /// `read_fresh_lock` の判定ロジック本体。生 content を直接受け取るため、呼び出し元が /// content を再利用 (takeover 直前のスナップショット比較等) できる。 +/// +/// **空 content は `None` (= 非 fresh)** を返す点に注意: 「保持者が居るか」を厳密に問う +/// 読み取り専用チェック (`pipeline_lock_holder`) では、書き込み中の空 lock を「holder あり」と +/// 報告する意味がない (pid 不明)。空を Held として扱うのは takeover 判定側 +/// (`classify_lock_content`) の責務。 fn is_fresh_content(content: &str, stale_threshold_secs: i64, now_unix: i64) -> Option<(u32, i64)> { let pid: u32 = parse_field(content, "pid")?.parse().ok()?; let start_unix: i64 = parse_field(content, "start_unix")?.parse().ok()?; @@ -301,210 +562,5 @@ fn current_unix_secs() -> i64 { } #[cfg(test)] -mod tests { - use super::*; - - fn temp_lock_path(prefix: &str) -> PathBuf { - let nanos = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .map(|d| d.subsec_nanos()) - .unwrap_or(0); - std::env::temp_dir().join(format!( - "pipeline-lock-{}-{}-{}", - prefix, - std::process::id(), - nanos - )) - } - - #[test] - fn acquire_creates_lock_and_drop_removes_it() { - let path = temp_lock_path("acquire"); - let result = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000); - assert!(matches!(result, PipelineLockResult::Acquired(_))); - assert!(path.exists()); - drop(result); - assert!(!path.exists(), "RAII drop で lock が削除される"); - } - - #[test] - fn second_acquire_is_busy_while_fresh() { - let path = temp_lock_path("busy"); - let _guard = acquire_pipeline_lock_at(path.clone(), "merge", 1800, 1_000_000); - let second = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_100); - match second { - PipelineLockResult::Busy { - holder_pid, - holder_age_secs, - } => { - assert_eq!(holder_pid, std::process::id()); - assert_eq!(holder_age_secs, 100); - } - _ => panic!("fresh lock 保持中は Busy になるべき"), - } - } - - #[test] - fn stale_lock_is_taken_over() { - let path = temp_lock_path("stale"); - std::fs::write(&path, "pid=99999\nstart_unix=1000000\nlabel=crashed\n").unwrap(); - let result = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000 + 1800); - assert!( - matches!(result, PipelineLockResult::Acquired(_)), - "threshold 到達で takeover" - ); - } - - #[test] - fn future_dated_lock_is_treated_as_stale() { - let path = temp_lock_path("future"); - std::fs::write(&path, "pid=99999\nstart_unix=2000000\nlabel=corrupt\n").unwrap(); - let result = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000); - assert!( - matches!(result, PipelineLockResult::Acquired(_)), - "future timestamp は stale 扱い (永続 fresh 化 bug class の防止)" - ); - } - - #[test] - fn corrupt_lock_is_taken_over() { - let path = temp_lock_path("corrupt"); - std::fs::write(&path, "not a lock file").unwrap(); - let result = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000); - assert!(matches!(result, PipelineLockResult::Acquired(_))); - } - - #[test] - fn read_fresh_lock_parses_fields_and_age() { - let path = temp_lock_path("read"); - std::fs::write(&path, "pid=4321\nstart_unix=1000000\nlabel=merge\n").unwrap(); - let held = read_fresh_lock(&path, 1800, 1_000_500); - assert_eq!(held, Some((4321, 500))); - let _ = std::fs::remove_file(&path); - } - - #[test] - fn missing_lock_reads_as_not_held() { - let path = temp_lock_path("missing"); - assert_eq!(read_fresh_lock(&path, 1800, 1_000_000), None); - } - - #[test] - fn acquire_writes_a_token() { - let path = temp_lock_path("token"); - let _guard = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000); - let content = std::fs::read_to_string(&path).unwrap(); - let token = parse_field(&content, "token").expect("token が書かれる"); - assert_eq!(token.len(), 32, "128bit hex"); - assert!(token.chars().all(|c| c.is_ascii_hexdigit())); - } - - #[test] - fn generate_token_is_unique_per_call() { - assert_ne!(generate_token(), generate_token(), "取得ごとに異なる token"); - } - - /// CodeRabbit Major #271 の regression guard: stale takeover 後に旧プロセスの Drop が - /// **新プロセスの lock を消さない**。A の guard を保持したまま同じパスを B が takeover - /// (別 token を上書き) し、A を drop しても B の lock ファイルが残ることを確認する。 - #[test] - fn drop_does_not_remove_lock_after_takeover() { - let path = temp_lock_path("takeover-guard"); - let a_guard = acquire_pipeline_lock_at(path.clone(), "A", 1800, 1_000_000); - assert!(matches!(a_guard, PipelineLockResult::Acquired(_))); - - let b_takeover_content = - build_lock_content("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", 55555, 1_000_100, "B"); - std::fs::write(&path, &b_takeover_content).unwrap(); - - drop(a_guard); - - assert!(path.exists(), "A の Drop が B の lock を消してはならない"); - let after = std::fs::read_to_string(&path).unwrap(); - assert!( - after.contains("token=bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"), - "B の lock がそのまま残る" - ); - let _ = std::fs::remove_file(&path); - } - - /// 通常ケース: 自分の token が残っていれば Drop で削除される。 - #[test] - fn drop_removes_lock_when_token_matches() { - let path = temp_lock_path("self-remove"); - let guard = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000); - assert!(path.exists()); - drop(guard); - assert!(!path.exists(), "自分の token の lock は削除される"); - } - - /// CodeRabbit re-review Major の regression guard: 同じ stale lock に対する - /// takeover を 2 スレッドが同時に行っても、`Acquired` になるのは 1 つだけ。 - #[test] - fn concurrent_stale_takeover_only_one_wins() { - let path = temp_lock_path("concurrent-stale"); - std::fs::write(&path, "pid=99999\nstart_unix=1000000\nlabel=crashed\n").unwrap(); - - let path_a = path.clone(); - let path_b = path.clone(); - let a = std::thread::spawn(move || { - acquire_pipeline_lock_at(path_a, "A", 1800, 1_000_000 + 1800) - }); - let b = std::thread::spawn(move || { - acquire_pipeline_lock_at(path_b, "B", 1800, 1_000_000 + 1800) - }); - let result_a = a.join().unwrap(); - let result_b = b.join().unwrap(); - - let acquired_count = [&result_a, &result_b] - .into_iter() - .filter(|r| matches!(r, PipelineLockResult::Acquired(_))) - .count(); - assert_eq!( - acquired_count, 1, - "stale takeover のレースで Acquired になるのは 1 プロセスのみのはず" - ); - - let _ = std::fs::remove_file(&path); - } - - /// SIM-NEW-pipeline_lock-L146 の regression guard: `stale_snapshot` 取得後に別プロセスが - /// 同じ隙間で先に takeover 済み (content が変化済み) の場合、`remove_file` をスキップして - /// `Busy` に落ちる (無条件 remove による 2 プロセスとも `Acquired` の再発防止)。 - #[test] - fn takeover_stale_lock_skips_remove_when_snapshot_is_stale() { - let path = temp_lock_path("snapshot-mismatch"); - let stale_snapshot_before_gap = - Some("pid=99999\nstart_unix=1000000\nlabel=crashed\n".to_string()); - - let content_written_by_concurrent_takeover_during_gap = - build_lock_content("cccccccccccccccccccccccccccccccc", 12345, 1_000_100, "other"); - std::fs::write(&path, &content_written_by_concurrent_takeover_during_gap).unwrap(); - - let result = takeover_stale_lock( - path.clone(), - "dddddddddddddddddddddddddddddddd".to_string(), - build_lock_content( - "dddddddddddddddddddddddddddddddd", - std::process::id(), - 1_000_200, - "push", - ), - stale_snapshot_before_gap, - 1800, - 1_000_200, - ); - - assert!( - matches!(result, PipelineLockResult::Busy { .. }), - "snapshot 不一致時は remove をスキップし Busy に落ちるべき" - ); - let after = std::fs::read_to_string(&path).unwrap(); - assert_eq!( - after, content_written_by_concurrent_takeover_during_gap, - "他プロセスの lock が変更されず残る" - ); - - let _ = std::fs::remove_file(&path); - } -} +#[path = "pipeline_lock/tests.rs"] +mod tests; diff --git a/src/lib-jj-helpers/src/pipeline_lock/tests.rs b/src/lib-jj-helpers/src/pipeline_lock/tests.rs new file mode 100644 index 00000000..5d724b2e --- /dev/null +++ b/src/lib-jj-helpers/src/pipeline_lock/tests.rs @@ -0,0 +1,337 @@ +use super::*; + +fn temp_lock_path(prefix: &str) -> PathBuf { + let nanos = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.subsec_nanos()) + .unwrap_or(0); + std::env::temp_dir().join(format!( + "pipeline-lock-{}-{}-{}", + prefix, + std::process::id(), + nanos + )) +} + +#[test] +fn acquire_creates_lock_and_drop_removes_it() { + let path = temp_lock_path("acquire"); + let result = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000); + assert!(matches!(result, PipelineLockResult::Acquired(_))); + assert!(path.exists()); + drop(result); + assert!(!path.exists(), "RAII drop で lock が削除される"); +} + +#[test] +fn second_acquire_is_busy_while_fresh() { + let path = temp_lock_path("busy"); + let _guard = acquire_pipeline_lock_at(path.clone(), "merge", 1800, 1_000_000); + let second = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_100); + match second { + PipelineLockResult::Busy { + holder_pid, + holder_age_secs, + } => { + assert_eq!(holder_pid, std::process::id()); + assert_eq!(holder_age_secs, 100); + } + _ => panic!("fresh lock 保持中は Busy になるべき"), + } +} + +#[test] +fn stale_lock_is_taken_over() { + let path = temp_lock_path("stale"); + std::fs::write(&path, "pid=99999\nstart_unix=1000000\nlabel=crashed\n").unwrap(); + let result = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000 + 1800); + assert!( + matches!(result, PipelineLockResult::Acquired(_)), + "threshold 到達で takeover" + ); +} + +#[test] +fn future_dated_lock_is_treated_as_stale() { + let path = temp_lock_path("future"); + std::fs::write(&path, "pid=99999\nstart_unix=2000000\nlabel=corrupt\n").unwrap(); + let result = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000); + assert!( + matches!(result, PipelineLockResult::Acquired(_)), + "future timestamp は stale 扱い (永続 fresh 化 bug class の防止)" + ); +} + +#[test] +fn corrupt_lock_is_taken_over() { + let path = temp_lock_path("corrupt"); + std::fs::write(&path, "not a lock file").unwrap(); + let result = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000); + assert!(matches!(result, PipelineLockResult::Acquired(_))); +} + +#[test] +fn read_fresh_lock_parses_fields_and_age() { + let path = temp_lock_path("read"); + std::fs::write(&path, "pid=4321\nstart_unix=1000000\nlabel=merge\n").unwrap(); + let held = read_fresh_lock(&path, 1800, 1_000_500); + assert_eq!(held, Some((4321, 500))); + let _ = std::fs::remove_file(&path); +} + +#[test] +fn missing_lock_reads_as_not_held() { + let path = temp_lock_path("missing"); + assert_eq!(read_fresh_lock(&path, 1800, 1_000_000), None); +} + +#[test] +fn acquire_writes_a_token() { + let path = temp_lock_path("token"); + let _guard = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000); + let content = std::fs::read_to_string(&path).unwrap(); + let token = parse_field(&content, "token").expect("token が書かれる"); + assert_eq!(token.len(), 32, "128bit hex"); + assert!(token.chars().all(|c| c.is_ascii_hexdigit())); +} + +#[test] +fn generate_token_is_unique_per_call() { + assert_ne!(generate_token(), generate_token(), "取得ごとに異なる token"); +} + +/// CodeRabbit Major #271 の regression guard: stale takeover 後に旧プロセスの Drop が +/// **新プロセスの lock を消さない**。A の guard を保持したまま同じパスを B が takeover +/// (別 token を上書き) し、A を drop しても B の lock ファイルが残ることを確認する。 +#[test] +fn drop_does_not_remove_lock_after_takeover() { + let path = temp_lock_path("takeover-guard"); + let a_guard = acquire_pipeline_lock_at(path.clone(), "A", 1800, 1_000_000); + assert!(matches!(a_guard, PipelineLockResult::Acquired(_))); + + let b_takeover_content = + build_lock_content("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", 55555, 1_000_100, "B"); + std::fs::write(&path, &b_takeover_content).unwrap(); + + drop(a_guard); + + assert!(path.exists(), "A の Drop が B の lock を消してはならない"); + let after = std::fs::read_to_string(&path).unwrap(); + assert!( + after.contains("token=bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"), + "B の lock がそのまま残る" + ); + let _ = std::fs::remove_file(&path); +} + +/// 通常ケース: 自分の token が残っていれば Drop で削除される。 +#[test] +fn drop_removes_lock_when_token_matches() { + let path = temp_lock_path("self-remove"); + let guard = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000); + assert!(path.exists()); + drop(guard); + assert!(!path.exists(), "自分の token の lock は削除される"); +} + +/// CodeRabbit re-review Major の regression guard: 同じ stale lock に対する +/// takeover を 2 スレッドが同時に行っても、`Acquired` になるのは 1 つだけ。 +#[test] +fn concurrent_stale_takeover_only_one_wins() { + let path = temp_lock_path("concurrent-stale"); + std::fs::write(&path, "pid=99999\nstart_unix=1000000\nlabel=crashed\n").unwrap(); + + let path_a = path.clone(); + let path_b = path.clone(); + let a = std::thread::spawn(move || { + acquire_pipeline_lock_at(path_a, "A", 1800, 1_000_000 + 1800) + }); + let b = std::thread::spawn(move || { + acquire_pipeline_lock_at(path_b, "B", 1800, 1_000_000 + 1800) + }); + let result_a = a.join().unwrap(); + let result_b = b.join().unwrap(); + + let acquired_count = [&result_a, &result_b] + .into_iter() + .filter(|r| matches!(r, PipelineLockResult::Acquired(_))) + .count(); + assert_eq!( + acquired_count, 1, + "stale takeover のレースで Acquired になるのは 1 プロセスのみのはず" + ); + + let _ = std::fs::remove_file(&path); +} + +/// 高競合ストレス: N スレッドが同一 stale lock を同時 takeover しても、**同時に** +/// `Acquired` になるのはちょうど 1 つ。2 スレッド版 (`concurrent_stale_takeover_only_one_wins`) +/// は旧実装で ~10% しか再現しなかったが、スレッド数を増やすとほぼ確実に踏む。 +/// +/// **全ガードを Vec に保持してから数える**のが要点。遅延イテレータ +/// (`map(join).filter().count()`) で数えると、まだ acquire 中の他スレッドの傍らで +/// 先行結果が drop され、その `PipelineLock::drop` が lock を削除 → 走行中スレッドが +/// 正当に取得し「2 Acquired」に見える (= 同時保持ではなく解放後の再取得。テスト +/// アーティファクトであって lock のバグではない)。collect で全ガードを保持し、 +/// 「同時点で 2 つ保持され得るか」だけを検証する。 +#[test] +fn concurrent_stale_takeover_many_threads_single_winner() { + for round in 0..40 { + let path = temp_lock_path(&format!("stress-{round}")); + std::fs::write(&path, "pid=99999\nstart_unix=1000000\nlabel=crashed\n").unwrap(); + + let handles: Vec<_> = (0..8) + .map(|i| { + let p = path.clone(); + std::thread::spawn(move || { + acquire_pipeline_lock_at(p, &format!("T{i}"), 1800, 1_000_000 + 1800) + }) + }) + .collect(); + let results: Vec<_> = handles.into_iter().map(|h| h.join().unwrap()).collect(); + let acquired = results + .iter() + .filter(|r| matches!(r, PipelineLockResult::Acquired(_))) + .count(); + + assert_eq!( + acquired, 1, + "round {round}: stale takeover の高競合で同時 Acquired は 1 つのみのはず (得た数: {acquired})" + ); + drop(results); + let _ = std::fs::remove_file(&path); + } +} + +/// SIM-NEW-pipeline_lock-L146 の regression guard: takeover 中に対象が既に fresh lock に +/// 変わっていた場合 (別 takeover が直前に完了)、奪わずに `Busy` へ倒し、その fresh lock を +/// 破壊しない (2 プロセスとも `Acquired` の再発防止)。 +#[test] +fn takeover_preserves_fresh_lock_that_appeared_during_takeover() { + let path = temp_lock_path("snapshot-mismatch"); + + // NOTE: takeover 中に対象が別 takeover 由来の fresh lock へ変わった状況を模す (age=100 < 1800=fresh)。 + let fresh_lock_from_concurrent_takeover = + build_lock_content("cccccccccccccccccccccccccccccccc", 12345, 1_000_100, "other"); + std::fs::write(&path, &fresh_lock_from_concurrent_takeover).unwrap(); + + let result = takeover_stale_lock( + path.clone(), + "dddddddddddddddddddddddddddddddd".to_string(), + build_lock_content( + "dddddddddddddddddddddddddddddddd", + std::process::id(), + 1_000_200, + "push", + ), + 1800, + 1_000_200, + ); + + assert!( + matches!(result, PipelineLockResult::Busy { .. }), + "対象が fresh 化していたら奪わず Busy に倒すべき" + ); + let after = std::fs::read_to_string(&path).unwrap(); + assert_eq!( + after, fresh_lock_from_concurrent_takeover, + "他プロセスの fresh lock は破壊されず残る" + ); + + assert!( + !takeover_sentinel_path(&path).exists(), + "takeover sentinel は完了後に残らない" + ); + + let _ = std::fs::remove_file(&path); +} + +/// SIM-NEW-pipeline_lock-L157 の regression guard: sentinel 保持者が +/// `perform_takeover` 実行中にクラッシュして `.takeover` が孤立しても、 +/// `SENTINEL_STALE_SECS` 経過後は自己修復し、本物の stale な pipeline lock を +/// 取得できる。従来は sentinel に age 判定が皆無で、この状態になると以降の +/// 取得試行が永久に `Busy` へ倒れていた。 +#[test] +fn orphaned_stale_sentinel_self_heals_and_lock_is_acquired() { + let path = temp_lock_path("orphaned-sentinel"); + // NOTE: 本物の pipeline lock も stale (クラッシュした pipeline の残骸)。 + std::fs::write(&path, "pid=99999\nstart_unix=1000000\nlabel=crashed\n").unwrap(); + // NOTE: sentinel 保持者が takeover 途中でクラッシュし孤立した状態を模す。 + let sentinel = takeover_sentinel_path(&path); + std::fs::write(&sentinel, "pid=88888\nstart_unix=1000000\n").unwrap(); + + let result = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000 + 1800); + + assert!( + matches!(result, PipelineLockResult::Acquired(_)), + "孤立した stale sentinel は自己修復され、本物の stale lock を取得できるべき" + ); + assert!(!sentinel.exists(), "取得完了後は sentinel が残らない"); + + drop(result); + let _ = std::fs::remove_file(&sentinel); +} + +/// sentinel が fresh (直近作成 = 別スレッドが takeover 実行中) な間は、本物の lock が +/// stale であっても sentinel を奪わず `Busy` に倒し、fresh sentinel を破壊しない。 +#[test] +fn fresh_sentinel_blocks_takeover_without_being_stolen() { + let path = temp_lock_path("fresh-sentinel"); + let now = 1_000_000 + 1800; + std::fs::write(&path, "pid=99999\nstart_unix=1000000\nlabel=crashed\n").unwrap(); + let sentinel = takeover_sentinel_path(&path); + // NOTE: age = 5s < SENTINEL_STALE_SECS(30s) = fresh (取得直後を模す)。 + let fresh_sentinel_content = format!("pid={}\nstart_unix={}\n", std::process::id(), now - 5); + std::fs::write(&sentinel, &fresh_sentinel_content).unwrap(); + + let result = acquire_pipeline_lock_at(path.clone(), "push", 1800, now); + + assert!( + matches!(result, PipelineLockResult::Busy { .. }), + "fresh sentinel が既にある間は奪わず Busy に倒すべき" + ); + let after = std::fs::read_to_string(&sentinel).unwrap(); + assert_eq!( + after, fresh_sentinel_content, + "fresh sentinel は破壊されず残る" + ); + + let _ = std::fs::remove_file(&path); + let _ = std::fs::remove_file(&sentinel); +} + +/// SIM-NEW-pipeline_lock-L157 の regression guard (高競合版): 孤立した stale +/// sentinel が存在する状態で N スレッドが同時に取得を試みても、自己修復後も +/// **同時に** `Acquired` になるのはちょうど 1 つ (自己修復が単一 winner 性を +/// 壊していないことの確認)。 +#[test] +fn concurrent_takeover_with_orphaned_sentinel_single_winner() { + for round in 0..20 { + let path = temp_lock_path(&format!("orphaned-sentinel-stress-{round}")); + std::fs::write(&path, "pid=99999\nstart_unix=1000000\nlabel=crashed\n").unwrap(); + let sentinel = takeover_sentinel_path(&path); + std::fs::write(&sentinel, "pid=88888\nstart_unix=1000000\n").unwrap(); + + let handles: Vec<_> = (0..8) + .map(|i| { + let p = path.clone(); + std::thread::spawn(move || { + acquire_pipeline_lock_at(p, &format!("T{i}"), 1800, 1_000_000 + 1800) + }) + }) + .collect(); + let results: Vec<_> = handles.into_iter().map(|h| h.join().unwrap()).collect(); + let acquired = results + .iter() + .filter(|r| matches!(r, PipelineLockResult::Acquired(_))) + .count(); + + assert_eq!( + acquired, 1, + "round {round}: 孤立 sentinel の自己修復後も同時 Acquired は 1 つのみのはず (得た数: {acquired})" + ); + drop(results); + let _ = std::fs::remove_file(&path); + let _ = std::fs::remove_file(&sentinel); + } +}