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();