From 1f55ec1287875f045f9dea73c60df4dedad0caaf Mon Sep 17 00:00:00 2001 From: fylorn <249551762+fylorn@users.noreply.github.com> Date: Fri, 25 Sep 2026 00:46:31 +0800 Subject: [PATCH] test(watch): make the "a burst is one signal" tests independent of disk stalls `a_burst_of_writes_collapses_into_one_signal` failed now and then on Linux CI with a second signal. The watcher was right; the test's premise was not. It writes ten times 10 ms apart and expects one signal, which holds only while no two writes are more than the 200 ms debounce window apart. On a CI disk shared with the other tests running in parallel, one `fs::write` (truncate, write, close) regularly stalls for 200-350 ms in writeback and dirty-page throttling; that splits the burst into two, and two signals is the correct answer. Reproduced in Docker with four `dd ... conv=fsync` loops as disk load: 21 of 40 runs failed, and timing each write showed every failing run had one write gap of 200 ms or more and every passing run had none. The watch tests in tw-watch and tw-config now use a directory on tmpfs (/dev/shm) on Linux, where a write never waits for the disk, so the premise holds. Under the same load: 0 of 40 failures, for both crates, and 0 of 10 for the whole tw-config suite. macOS keeps the system temporary directory. The debounce loop is also pulled out into `debounce_loop` and tested with no filesystem and no timing in the way: ten events queued before it starts must give exactly one signal (read from a channel large enough to see a second one, unlike the capacity-1 channel `watch` uses, which would swallow it), a later event another, and a closed source ends it. Making the loop send per event fails that test. Co-Authored-By: Claude Opus 5.5 --- crates/tw-config/src/watch.rs | 26 +++++++--- crates/tw-watch/src/lib.rs | 92 ++++++++++++++++++++++++++++------- 2 files changed, 94 insertions(+), 24 deletions(-) diff --git a/crates/tw-config/src/watch.rs b/crates/tw-config/src/watch.rs index ea10526..e46ee15 100644 --- a/crates/tw-config/src/watch.rs +++ b/crates/tw-config/src/watch.rs @@ -54,6 +54,20 @@ mod tests { Closed, } + /// 测试用的目录,Linux 上在内存盘上。 + /// + /// 这里好几条断言是「只报一次」,而它的前提是写入之间的间隔不超过去抖窗口。 + /// 在 CI 的磁盘上,一次 `fs::write` 卡两三百毫秒是常事(别的测试在同时写盘), + /// 那时报两次是监听对、测试错 —— 理由和复现办法见 `tw_watch` 的同名函数。 + fn scratch() -> tempfile::TempDir { + let shm = Path::new("/dev/shm"); + if cfg!(target_os = "linux") && shm.is_dir() { + tempfile::tempdir_in(shm).unwrap() + } else { + tempfile::tempdir().unwrap() + } + } + async fn recv(rx: &mut tokio::sync::mpsc::Receiver<()>, within: Duration) -> Signal { match tokio::time::timeout(within, rx.recv()).await { Ok(Some(())) => Signal::Got, @@ -64,7 +78,7 @@ mod tests { #[tokio::test] async fn an_edit_produces_exactly_one_signal() { - let d = tempfile::tempdir().unwrap(); + let d = scratch(); let p = d.path().join("config.yaml"); std::fs::write(&p, "a: 1\n").unwrap(); let (_w, mut rx) = watch(&p).unwrap(); @@ -87,7 +101,7 @@ mod tests { async fn a_save_that_writes_a_temp_file_and_renames_still_works() { // **编辑器就是这么保存的。**盯着文件而不是目录的话,这里之后 // 监听会永久失效,而且不报任何错。 - let d = tempfile::tempdir().unwrap(); + let d = scratch(); let p = d.path().join("config.yaml"); std::fs::write(&p, "a: 1\n").unwrap(); let (_w, mut rx) = watch(&p).unwrap(); @@ -106,7 +120,7 @@ mod tests { #[tokio::test] async fn a_burst_of_writes_collapses_into_one_signal() { - let d = tempfile::tempdir().unwrap(); + let d = scratch(); let p = d.path().join("config.yaml"); std::fs::write(&p, "a: 0\n").unwrap(); let (_w, mut rx) = watch(&p).unwrap(); @@ -126,7 +140,7 @@ mod tests { #[tokio::test] async fn a_sibling_file_does_not_wake_us_up() { // 历史目录就在配置旁边,每存一版都会动那个目录。 - let d = tempfile::tempdir().unwrap(); + let d = scratch(); let p = d.path().join("config.yaml"); // **故意不建 config.yaml。**过滤器是按文件名判的,所以这个测试 // 里唯一可能穿过它的事件就是 config.yaml 自己的 —— 而它不存在, @@ -154,7 +168,7 @@ mod tests { async fn deleting_the_file_is_reported_too() { // 删掉配置文件是一件必须知道的事 —— 不然网关会拿着一份内存里的 // 幽灵配置继续跑,而用户以为自己已经清空了它。 - let d = tempfile::tempdir().unwrap(); + let d = scratch(); let p = d.path().join("config.yaml"); std::fs::write(&p, "a: 1\n").unwrap(); let (_w, mut rx) = watch(&p).unwrap(); @@ -169,7 +183,7 @@ mod tests { #[tokio::test] async fn watching_a_file_that_does_not_exist_yet_still_works() { // 首次运行时配置还没生成,而监听可能先起来。 - let d = tempfile::tempdir().unwrap(); + let d = scratch(); let p = d.path().join("config.yaml"); let (_w, mut rx) = watch(&p).unwrap(); std::fs::write(&p, "a: 1\n").unwrap(); diff --git a/crates/tw-watch/src/lib.rs b/crates/tw-watch/src/lib.rs index 6f474b7..b9a6da2 100644 --- a/crates/tw-watch/src/lib.rs +++ b/crates/tw-watch/src/lib.rs @@ -79,24 +79,36 @@ pub fn watch( } // notify 的回调跑在它自己的线程上,这里把它接到 tokio 上并去抖。 - std::thread::spawn(move || { - while raw_rx.recv().is_ok() { - // 收到一个之后,把窗口内后续的全部吞掉 —— 一次保存的那一串事件因此只 - // 产生一个信号。 - loop { - match raw_rx.recv_timeout(debounce) { - Ok(()) => continue, - Err(std::sync::mpsc::RecvTimeoutError::Timeout) => break, - Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => return, - } - } - let _ = tx.try_send(()); - } - }); + std::thread::spawn(move || debounce_loop(&raw_rx, debounce, &tx)); Ok((Watch { _inner: w }, rx)) } +/// 去抖:收到一个事件之后,把间隔不到 `window` 的后续事件全部并进来,静下来 +/// `window` 之后发**一个**信号。事件源断开就返回。 +/// +/// **「一阵」是按事件之间的间隔定义的,不是按谁发起的。**同一个人连写十次, +/// 只要中间有一次写入本身卡了超过 `window`(磁盘忙时 `write` 会被限流),那就是 +/// 两阵、两个信号 —— 这是对的:第一阵之后读到的文件是那一刻的真实内容。 +fn debounce_loop( + raw_rx: &std::sync::mpsc::Receiver<()>, + window: Duration, + tx: &tokio::sync::mpsc::Sender<()>, +) { + while raw_rx.recv().is_ok() { + // 收到一个之后,把窗口内后续的全部吞掉 —— 一次保存的那一串事件因此只 + // 产生一个信号。 + loop { + match raw_rx.recv_timeout(window) { + Ok(()) => continue, + Err(std::sync::mpsc::RecvTimeoutError::Timeout) => break, + Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => return, + } + } + let _ = tx.try_send(()); + } +} + #[cfg(test)] mod tests { use super::*; @@ -124,9 +136,53 @@ mod tests { p.extension().is_some_and(|e| e == "md") } + /// 测试用的目录。**Linux 上放在内存盘(`/dev/shm`)上。** + /// + /// 「连写十次是一阵」的前提是任意两次写之间不超过去抖窗口。放在磁盘上时 + /// 这个前提不归测试管:CI 的机器上别的测试同时在写盘,一次 `fs::write` 被 + /// 截断之后的刷盘和脏页限流卡住两三百毫秒是常事(在 Docker 里加上 `dd` + /// 压盘就能复现,失败的那几次正好都有一次写超过了 200ms)。那时监听报两次 + /// 是对的,错的是测试。内存盘上的写不碰磁盘,前提才真的成立。 + /// + /// macOS 没有 `/dev/shm`,用系统的临时目录 —— 那边从来没有因此失败过。 + fn scratch() -> tempfile::TempDir { + let shm = Path::new("/dev/shm"); + if cfg!(target_os = "linux") && shm.is_dir() { + tempfile::tempdir_in(shm).unwrap() + } else { + tempfile::tempdir().unwrap() + } + } + + /// 去抖本身,不经过文件系统:事件之间有没有空隙完全由测试决定。 + #[tokio::test] + async fn events_already_waiting_are_one_signal_and_a_later_one_is_another() { + let (raw_tx, raw_rx) = std::sync::mpsc::channel(); + // 通道放得下每一个信号:`watch` 里那个容量 1 的通道会把多发的信号吞掉, + // 在这里用它就看不出去抖有没有并起来 + let (tx, mut rx) = tokio::sync::mpsc::channel(16); + // 十个事件在去抖开始之前就排好了:它们之间没有任何间隔 + for _ in 0..10 { + raw_tx.send(()).unwrap(); + } + std::thread::spawn(move || debounce_loop(&raw_rx, WINDOW, &tx)); + assert_eq!(recv(&mut rx, Duration::from_secs(3)).await, Signal::Got); + assert_eq!( + recv(&mut rx, WINDOW * 3).await, + Signal::Quiet, + "排在一起的十个事件报了不止一次" + ); + // 静下来之后的下一个事件是新的一阵 + raw_tx.send(()).unwrap(); + assert_eq!(recv(&mut rx, Duration::from_secs(3)).await, Signal::Got); + // 发端没了,去抖跟着结束 + drop(raw_tx); + assert_eq!(recv(&mut rx, Duration::from_secs(3)).await, Signal::Closed); + } + #[tokio::test] async fn a_burst_of_writes_is_one_signal() { - let d = tempfile::tempdir().unwrap(); + let d = scratch(); let (_w, mut rx) = watch(&[d.path().to_path_buf()], WINDOW, only_md).unwrap(); for i in 0..10 { std::fs::write(d.path().join("a.md"), format!("{i}\n")).unwrap(); @@ -142,7 +198,7 @@ mod tests { #[tokio::test] async fn what_the_filter_turns_down_stays_quiet() { - let d = tempfile::tempdir().unwrap(); + let d = scratch(); let (_w, mut rx) = watch(&[d.path().to_path_buf()], WINDOW, only_md).unwrap(); std::fs::write(d.path().join("x.log"), "x\n").unwrap(); assert_eq!( @@ -155,7 +211,7 @@ mod tests { #[tokio::test] async fn a_subdirectory_is_not_watched() { - let d = tempfile::tempdir().unwrap(); + let d = scratch(); std::fs::create_dir_all(d.path().join("history")).unwrap(); let (_w, mut rx) = watch(&[d.path().to_path_buf()], WINDOW, only_md).unwrap(); std::fs::write(d.path().join("history/1.md"), "x\n").unwrap(); @@ -168,7 +224,7 @@ mod tests { #[tokio::test] async fn several_directories_are_watched_at_once() { - let d = tempfile::tempdir().unwrap(); + let d = scratch(); let (a, b) = (d.path().join("a"), d.path().join("b")); std::fs::create_dir_all(&a).unwrap(); std::fs::create_dir_all(&b).unwrap();