Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 70 additions & 25 deletions src/lib-jj-helpers/src/pipeline_lock.rs
Original file line number Diff line number Diff line change
Expand Up @@ -157,6 +157,11 @@ pub fn acquire_pipeline_lock_at(
/// (remove + create_new の不在窓に親 fast-path が割り込む) も同 doc に記録。
/// SIM-NEW-pipeline_lock-L146。
///
/// reclaim 経由で確保した場合は gate marker も **takeover 完了まで保持**し、ここで
/// sentinel と一緒に除去する。marker を [`finish_reclaim`] 内で先に除去すると、同じ
/// stale content を cache した出遅れスレッドが takeover 実行中に同一 gate を再作成できる
/// (PR #342 CI で 2 `Acquired` として実測。詳細は [`finish_reclaim`] の doc)。
///
/// sentinel 自体の取得は [`acquire_takeover_sentinel`] に委譲する (age ベースの
/// 自己修復込み、SIM-NEW-pipeline_lock-L157)。
fn takeover_stale_lock(
Expand All @@ -168,14 +173,17 @@ fn takeover_stale_lock(
) -> PipelineLockResult {
let sentinel = takeover_sentinel_path(&path);
match acquire_takeover_sentinel(&sentinel, now_unix) {
SentinelGate::Acquired => {
SentinelGate::Acquired { reclaim_marker } => {
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);
}
}
if let Some(marker) = reclaim_marker {
let _ = std::fs::remove_file(&marker);
}
result
}
SentinelGate::Busy => busy_from_disk(&path, stale_threshold_secs, now_unix),
Expand All @@ -192,8 +200,12 @@ const SENTINEL_STALE_SECS: i64 = 30;

/// [`takeover_stale_lock`] 用 sentinel 取得結果。
enum SentinelGate {
/// sentinel を確保した (自分が takeover 実行権を持つ)。
Acquired,
/// sentinel を確保した (自分が takeover 実行権を持つ)。`reclaim_marker` は自己修復
/// (reclaim) 経由で確保した場合の gate marker path。**takeover 完了まで保持**し、
/// 同じ stale content を cache した出遅れスレッドが takeover 実行中に同一 gate を
/// 再作成する経路を塞ぐ (呼び出し元が takeover 後に除去する。PR #342 CI 実測の
/// 2 `Acquired` 対策の第 2 層)。
Acquired { reclaim_marker: Option<PathBuf> },
/// 別プロセスが sentinel を保持中 (fresh、または create_new 直後の書き込み待ち)。
Busy,
/// I/O エラーで判定不能。
Expand All @@ -205,7 +217,11 @@ enum SentinelGate {
/// [`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,
Ok(()) => {
return SentinelGate::Acquired {
reclaim_marker: None,
}
}
Err(e) if e.kind() != std::io::ErrorKind::AlreadyExists => {
return SentinelGate::Unavailable(e.to_string());
}
Expand Down Expand Up @@ -245,7 +261,7 @@ fn try_create_sentinel(sentinel: &Path, now_unix: i64) -> std::io::Result<()> {
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),
Ok(()) => finish_reclaim(sentinel, &reclaim, stale_content, now_unix),
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
reap_orphaned_reclaim_marker(&reclaim, now_unix);
SentinelGate::Busy
Expand All @@ -254,22 +270,50 @@ fn reclaim_stale_sentinel(sentinel: &Path, stale_content: &str, now_unix: i64) -
}
}

/// reclaim gate の勝者だけが呼ぶ: sentinel を除去して再作成する。
/// reclaim gate の勝者だけが呼ぶ: sentinel を検証してから除去・再作成する。
///
/// **除去の前に、sentinel が今も gate を正当化した stale content のままかを検証する**
/// (PR #342 CI windows-latest で実測した 2 `Acquired` の根本対策)。gate path は content
/// 由来のため、除去前に同じ stale content を読んでいた出遅れスレッドは、勝者が marker を
/// 除去した後に**同じ gate path を再作成して勝てて**しまう。旧実装はそこで sentinel を
/// 無条件除去していたため、勝者が `perform_takeover` 実行中に fresh sentinel を破壊され、
/// 2 本目の takeover 実行権が発生した (勝者の rename 前に lock を読めば双方 Stale 判定 →
/// 双方 `Acquired`)。content が期待と異なる / 消えた / 読めない場合は一切触れず `Busy` に
/// 倒す。この検証と、marker を takeover 完了まで保持する第 2 層 ([`SentinelGate`]) の
/// 2 段で経路を塞ぐ。
///
/// この関数を呼べるのは「この 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 {
/// 残余 TOCTOU (検証 → 除去の間に content が変わる窓): この窓を突くには「孤立 stale
/// marker の reap が再作成直後の fresh marker を誤回収し、同じ stale content を cache
/// した 2 スレッドが同時に検証を通過する」という 30 秒クラッシュ + 複数 µs 窓の重畳が
/// 必要で、最終的には `try_create_sentinel` の `create_new` 排他が 1 勝者に収束させる。
/// [`PipelineLock`] の Drop の残余 TOCTOU と同様、実用上安全な残余として明示する。
fn finish_reclaim(
sentinel: &Path,
reclaim: &Path,
expected_stale_content: &str,
now_unix: i64,
) -> SentinelGate {
let still_expected = std::fs::read_to_string(sentinel)
.map(|current| current == expected_stale_content)
.unwrap_or(false);
if !still_expected {
let _ = std::fs::remove_file(reclaim);
return SentinelGate::Busy;
}
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
match try_create_sentinel(sentinel, now_unix) {
Ok(()) => SentinelGate::Acquired {
reclaim_marker: Some(reclaim.to_path_buf()),
},
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
let _ = std::fs::remove_file(reclaim);
SentinelGate::Busy
}
Err(e) => {
let _ = std::fs::remove_file(reclaim);
SentinelGate::Unavailable(e.to_string())
}
}
}

/// reclaim gate 自体に `pid=` / `start_unix=` を書き込む (sentinel と同形式、
Expand All @@ -282,13 +326,14 @@ fn try_create_reclaim_marker(reclaim: &Path, now_unix: i64) -> std::io::Result<(
f.write_all(format!("pid={}\nstart_unix={}\n", std::process::id(), now_unix).as_bytes())
}

/// reclaim gate 保持者が `finish_reclaim` の 2 手 (remove + create_new) の間でクラッシュし、
/// gate 自体が孤立した場合の回収。
/// reclaim gate 保持者が takeover 完了前にクラッシュし、gate 自体が孤立した場合の回収。
///
/// この gate は特定の stale content 専用 (path が content 由来のため) で、除去後に
/// **別の正当な世代が同じ path に再作成される余地が無い**。よって sentinel と違い、
/// stale 判定さえできれば無条件 remove で安全に回収できる (複数スレッドが同時に
/// 回収を試みても、`remove_file` 自体の排他性により実害は出ない)。
/// この gate は特定の stale content 専用 (path が content 由来のため)。同じ stale content
/// を cache した出遅れスレッドが除去後に同じ path を再作成する余地はある
/// ([`finish_reclaim`] の doc の経路) が、再作成直後の gate は fresh のため本関数の
/// stale 判定 (`SENTINEL_STALE_SECS` 経過) を通らず、誤回収には至らない。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!(
Expand Down
112 changes: 112 additions & 0 deletions src/lib-jj-helpers/src/pipeline_lock/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -300,6 +300,118 @@ fn fresh_sentinel_blocks_takeover_without_being_stolen() {
let _ = std::fs::remove_file(&sentinel);
}

/// PR #342 CI (windows-latest、2 コア runner) で実測された「同時 Acquired 2 つ」の
/// 決定論再現。
///
/// 出遅れスレッドは**除去前に読んだ** stale sentinel content を握ったまま
/// `reclaim_stale_sentinel` に入る。勝者が自己修復 (`finish_reclaim`) を終えて reclaim
/// marker を除去した後でも、gate path は content 由来のため同じ path を再作成できて
/// しまう。そこで sentinel を検証なしに除去すると、**勝者の fresh sentinel を破壊**して
/// 2 本目の takeover 実行権が生まれ、勝者の rename 前に lock を読めば双方 `Acquired`
/// になる (CI の「得た数: 2」)。
///
/// 期待動作: 現在の sentinel が gate を正当化した stale content と一致しない場合は
/// `Busy` に倒し、勝者の sentinel には一切触れない。
#[test]
fn reclaim_with_cached_stale_content_leaves_winner_sentinel_intact() {
let path = temp_lock_path("reclaim-cached-stale");
let sentinel = takeover_sentinel_path(&path);
let now = 1_001_800;
let cached_stale = "pid=88888\nstart_unix=1000000\n";
let winner_fresh = "pid=77777\nstart_unix=1001799\n";
std::fs::write(&sentinel, winner_fresh).unwrap();

let gate = reclaim_stale_sentinel(&sentinel, cached_stale, now);

let sentinel_after = std::fs::read_to_string(&sentinel);
let marker = reclaim_gate_path(&sentinel, cached_stale);
let marker_leaked = marker.exists();
let _ = std::fs::remove_file(&sentinel);
let _ = std::fs::remove_file(&marker);

assert!(
matches!(gate, SentinelGate::Busy),
"cache 済み stale content で再確立済み sentinel を奪ってはならない"
);
assert_eq!(
sentinel_after.expect("勝者の sentinel は存在し続ける"),
winner_fresh,
"勝者の fresh sentinel は無傷で残る"
);
assert!(!marker_leaked, "出遅れ側が作った reclaim marker は後始末される");
}

/// 孤立 stale sentinel からの自己修復で `Acquired` した後、reclaim marker が残留しない
/// こと。marker は takeover 完了まで保持される (除去が早すぎると上記の cache 済み
/// content 再作成経路が開く) ため、除去タイミングの回帰を leak として検出する。
#[test]
fn orphaned_sentinel_selfheal_leaves_no_reclaim_marker() {
let path = temp_lock_path("orphan-selfheal-no-marker");
std::fs::write(&path, "pid=99999\nstart_unix=1000000\nlabel=crashed\n").unwrap();
let sentinel = takeover_sentinel_path(&path);
let stale_sentinel = "pid=88888\nstart_unix=1000000\n";
std::fs::write(&sentinel, stale_sentinel).unwrap();

let result = acquire_pipeline_lock_at(path.clone(), "push", 1800, 1_000_000 + 1800);

assert!(matches!(result, PipelineLockResult::Acquired(_)));
assert!(!sentinel.exists(), "takeover 完了後に sentinel は除去される");
assert!(
!reclaim_gate_path(&sentinel, stale_sentinel).exists(),
"reclaim marker は takeover 完了後に除去され残留しない"
);
drop(result);
let _ = std::fs::remove_file(&path);
}

/// 上記 regression guard の高競合版 (opt-in)。
///
/// PR #342 の CI Windows leg で `concurrent_takeover_with_orphaned_sentinel_single_winner`
/// が「得た数: 2」で落ちたが、開発機 (Windows) では 40 回連続で再現しなかった。runner の
/// コア数・負荷でしか踏まない窓があるため、スレッド数と round 数を上げて窓を広げる。
/// 失敗を即 assert せず全 round 分集計するのは、**再現率**自体が修正の効果測定になるため。
#[test]
#[ignore = "stress: 高競合の lock race 再現用 (数秒〜数十秒)。`cargo test -- --ignored --test-threads=1` で実行"]
fn concurrent_takeover_orphaned_sentinel_stress() {
const ROUNDS: usize = 400;
const THREADS: usize = 32;

let mut failures: Vec<(usize, usize)> = Vec::new();
for round in 0..ROUNDS {
let path = temp_lock_path(&format!("orphaned-sentinel-hicon-{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..THREADS)
.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();
if acquired != 1 {
failures.push((round, acquired));
}
drop(results);
let _ = std::fs::remove_file(&path);
let _ = std::fs::remove_file(&sentinel);
}

assert!(
failures.is_empty(),
"{} / {ROUNDS} round で同時 Acquired が 1 以外になった (round, 得た数): {:?}",
failures.len(),
&failures[..failures.len().min(10)]
);
}

/// SIM-NEW-pipeline_lock-L157 の regression guard (高競合版): 孤立した stale
/// sentinel が存在する状態で N スレッドが同時に取得を試みても、自己修復後も
/// **同時に** `Acquired` になるのはちょうど 1 つ (自己修復が単一 winner 性を
Expand Down
Loading