subetha_cxc/
shared_condvar.rs1use std::fs::OpenOptions;
49use std::io;
50use std::path::{Path, PathBuf};
51use std::sync::Arc;
52use std::sync::atomic::{AtomicU64, Ordering};
53use std::time::{Duration, Instant};
54
55use memmap2::{MmapMut, MmapOptions};
56
57use crate::cross_process_waker::{
58 CrossProcessWaker, MAX_WAITERS_DEFAULT, WakerError,
59};
60
61const CONDVAR_GEN_MAGIC: u64 = 0x434F_4E44_5641_5230; const GEN_REGION_SIZE: usize = 64; const GEN_OFFSET: usize = 8;
66
67#[derive(Debug, Clone, Copy, PartialEq, Eq)]
69pub enum CondvarError {
70 WakerFull,
73 Timeout,
76 LayoutMismatch,
79 Io(io::ErrorKind),
81}
82
83impl From<WakerError> for CondvarError {
84 fn from(e: WakerError) -> Self {
85 match e {
86 WakerError::Full => Self::WakerFull,
87 WakerError::Timeout => Self::Timeout,
88 WakerError::LayoutMismatch => Self::LayoutMismatch,
89 WakerError::IoError(k) => Self::Io(k),
90 }
91 }
92}
93
94impl From<io::Error> for CondvarError {
95 fn from(e: io::Error) -> Self { Self::Io(e.kind()) }
96}
97
98#[allow(dead_code)]
107enum GenBacking {
108 Anon(MmapMut),
109 File(std::fs::File, MmapMut),
110}
111
112struct GenAtom {
113 #[allow(dead_code)]
115 backing: GenBacking,
116 ptr: *const AtomicU64,
117}
118
119unsafe impl Send for GenAtom {}
122unsafe impl Sync for GenAtom {}
123
124impl GenAtom {
125 fn create_anon() -> Result<Self, CondvarError> {
126 let mut mmap = MmapOptions::new().len(GEN_REGION_SIZE).map_anon()?;
127 let base = mmap.as_mut_ptr();
128 unsafe {
129 (base as *mut u64).write(CONDVAR_GEN_MAGIC);
130 (base.add(GEN_OFFSET) as *mut AtomicU64).write(AtomicU64::new(0));
131 }
132 let ptr = unsafe { base.add(GEN_OFFSET) as *const AtomicU64 };
133 Ok(Self { backing: GenBacking::Anon(mmap), ptr })
134 }
135
136 fn create_file(path: &Path) -> Result<Self, CondvarError> {
140 let (file, mut mmap) = crate::mmf_attach::create_or_attach(
141 path,
142 GEN_REGION_SIZE,
143 |base| unsafe {
144 (base.add(GEN_OFFSET) as *mut AtomicU64).write(AtomicU64::new(0));
145 std::ptr::write_volatile(base as *mut u64, CONDVAR_GEN_MAGIC);
146 },
147 |base| unsafe { (base as *const u64).read() == CONDVAR_GEN_MAGIC },
148 )?;
149 let base = mmap.as_mut_ptr();
150 let ptr = unsafe { base.add(GEN_OFFSET) as *const AtomicU64 };
151 Ok(Self { backing: GenBacking::File(file, mmap), ptr })
152 }
153
154 fn open_file(path: &Path) -> Result<Self, CondvarError> {
155 let file = OpenOptions::new().read(true).write(true).open(path)?;
156 let meta = file.metadata()?;
157 if (meta.len() as usize) < GEN_REGION_SIZE {
158 return Err(CondvarError::LayoutMismatch);
159 }
160 let mut mmap = unsafe { MmapOptions::new().len(GEN_REGION_SIZE).map_mut(&file)? };
161 let base = mmap.as_mut_ptr();
162 let magic = unsafe { (base as *const u64).read() };
163 if magic != CONDVAR_GEN_MAGIC {
164 return Err(CondvarError::LayoutMismatch);
165 }
166 let ptr = unsafe { base.add(GEN_OFFSET) as *const AtomicU64 };
167 Ok(Self { backing: GenBacking::File(file, mmap), ptr })
168 }
169
170 #[inline]
171 fn atom(&self) -> &AtomicU64 {
172 unsafe { &*self.ptr }
176 }
177}
178
179pub struct SharedCondvar {
183 waker: Arc<CrossProcessWaker>,
184 gen_atom: Arc<GenAtom>,
185}
186
187impl SharedCondvar {
188 pub fn create_anon() -> Result<Self, CondvarError> {
191 Self::create_anon_with_capacity(MAX_WAITERS_DEFAULT)
192 }
193
194 pub fn create_anon_with_capacity(max_waiters: usize) -> Result<Self, CondvarError> {
196 let waker = Arc::new(CrossProcessWaker::create_anon(max_waiters)?);
197 let gen_atom = Arc::new(GenAtom::create_anon()?);
198 Ok(Self { waker, gen_atom })
199 }
200
201 pub fn create(base_path: impl AsRef<Path>) -> Result<Self, CondvarError> {
205 Self::create_with_capacity(base_path, MAX_WAITERS_DEFAULT)
206 }
207
208 pub fn create_with_capacity(
209 base_path: impl AsRef<Path>,
210 max_waiters: usize,
211 ) -> Result<Self, CondvarError> {
212 let (waker_path, gen_path) = side_paths(base_path.as_ref());
213 let waker = Arc::new(CrossProcessWaker::create(waker_path, max_waiters)?);
214 let gen_atom = Arc::new(GenAtom::create_file(&gen_path)?);
215 Ok(Self { waker, gen_atom })
216 }
217
218 pub fn open(base_path: impl AsRef<Path>) -> Result<Self, CondvarError> {
222 Self::open_with_capacity(base_path, MAX_WAITERS_DEFAULT)
223 }
224
225 pub fn open_with_capacity(
226 base_path: impl AsRef<Path>,
227 expected_max_waiters: usize,
228 ) -> Result<Self, CondvarError> {
229 let (waker_path, gen_path) = side_paths(base_path.as_ref());
230 let waker = Arc::new(CrossProcessWaker::open(waker_path, expected_max_waiters)?);
231 let gen_atom = Arc::new(GenAtom::open_file(&gen_path)?);
232 Ok(Self { waker, gen_atom })
233 }
234
235 pub fn wait<F: FnMut() -> bool>(&self, mut predicate: F) -> Result<(), CondvarError> {
239 loop {
240 if predicate() {
241 return Ok(());
242 }
243 let snapshot = self.gen_atom.atom().load(Ordering::Acquire);
249 let token = self.waker.try_park(snapshot + 1)?;
250 if predicate() {
252 self.waker.release(token);
253 return Ok(());
254 }
255 self.waker.wait(token, None)?;
256 }
257 }
258
259 pub fn wait_timeout<F: FnMut() -> bool>(
263 &self,
264 mut predicate: F,
265 timeout: Duration,
266 ) -> Result<(), CondvarError> {
267 let deadline = Instant::now() + timeout;
268 loop {
269 if predicate() {
270 return Ok(());
271 }
272 let snapshot = self.gen_atom.atom().load(Ordering::Acquire);
273 let token = self.waker.try_park(snapshot + 1)?;
274 if predicate() {
275 self.waker.release(token);
276 return Ok(());
277 }
278 let now = Instant::now();
279 if now >= deadline {
280 self.waker.release(token);
281 return Err(CondvarError::Timeout);
282 }
283 let remaining = deadline - now;
284 match self.waker.wait(token, Some(remaining)) {
285 Ok(()) => continue,
286 Err(WakerError::Timeout) => {
287 if predicate() {
288 return Ok(());
289 }
290 return Err(CondvarError::Timeout);
291 }
292 Err(e) => return Err(CondvarError::from(e)),
293 }
294 }
295 }
296
297 pub fn notify_one(&self) -> usize {
301 let new_gen = self.gen_atom.atom().fetch_add(1, Ordering::Release) + 1;
302 self.waker.wake_one_up_to(new_gen)
303 }
304
305 pub fn notify_all(&self) -> usize {
307 let new_gen = self.gen_atom.atom().fetch_add(1, Ordering::Release) + 1;
308 self.waker.wake_up_to(new_gen)
309 }
310
311 pub fn generation(&self) -> u64 {
314 self.gen_atom.atom().load(Ordering::Acquire)
315 }
316
317 pub fn waker(&self) -> &Arc<CrossProcessWaker> { &self.waker }
320}
321
322fn side_paths(base: &Path) -> (PathBuf, PathBuf) {
323 let mut w = base.as_os_str().to_owned();
324 w.push(".waker.bin");
325 let mut g = base.as_os_str().to_owned();
326 g.push(".gen.bin");
327 (PathBuf::from(w), PathBuf::from(g))
328}
329
330#[cfg(test)]
331mod tests {
332 use super::*;
333 use std::sync::atomic::AtomicBool;
334 use std::thread;
335
336 #[test]
337 fn notify_one_wakes_exactly_one_waiter() {
338 let cv = Arc::new(SharedCondvar::create_anon().expect("create"));
339 let pred = Arc::new(AtomicBool::new(false));
340 let waiters: Vec<_> = (0..3)
341 .map(|_| {
342 let cv2 = Arc::clone(&cv);
343 let pred2 = Arc::clone(&pred);
344 thread::spawn(move || {
345 cv2.wait(|| pred2.load(Ordering::Acquire)).unwrap();
346 })
347 })
348 .collect();
349 thread::sleep(Duration::from_millis(50));
350
351 assert_eq!(cv.notify_one(), 1);
355 thread::sleep(Duration::from_millis(20));
357 pred.store(true, Ordering::Release);
358 cv.notify_all();
359 for h in waiters {
360 h.join().unwrap();
361 }
362 }
363
364 #[test]
365 fn notify_all_wakes_every_waiter() {
366 let cv = Arc::new(SharedCondvar::create_anon().expect("create"));
367 let pred = Arc::new(AtomicBool::new(false));
368 let waiters: Vec<_> = (0..4)
369 .map(|_| {
370 let cv2 = Arc::clone(&cv);
371 let pred2 = Arc::clone(&pred);
372 thread::spawn(move || {
373 cv2.wait(|| pred2.load(Ordering::Acquire)).unwrap();
374 })
375 })
376 .collect();
377 thread::sleep(Duration::from_millis(30));
378 pred.store(true, Ordering::Release);
379 let woken = cv.notify_all();
380 assert!(woken >= 1, "at least one waiter woken (got {woken})");
381 for h in waiters {
382 h.join().unwrap();
383 }
384 }
385
386 #[test]
387 fn wait_timeout_returns_timeout() {
388 let cv = SharedCondvar::create_anon().expect("create");
389 let t0 = Instant::now();
390 let err = cv.wait_timeout(|| false, Duration::from_millis(60));
391 assert_eq!(err, Err(CondvarError::Timeout));
392 assert!(t0.elapsed() >= Duration::from_millis(50));
393 }
394
395 #[test]
396 fn wait_returns_immediately_if_predicate_already_true() {
397 let cv = SharedCondvar::create_anon().expect("create");
398 let pred = AtomicBool::new(true);
399 let t0 = Instant::now();
400 cv.wait(|| pred.load(Ordering::Acquire)).unwrap();
401 assert!(t0.elapsed() < Duration::from_millis(10));
402 }
403
404 #[test]
417 fn file_backed_create_then_arc_clone_round_trip() {
418 let dir = std::env::temp_dir();
419 let path = dir.join(format!("subetha_condvar_test_{}", std::process::id()));
420 for suffix in [".waker.bin", ".gen.bin"] {
422 let mut p = path.as_os_str().to_owned();
423 p.push(suffix);
424 drop(std::fs::remove_file(PathBuf::from(p)));
425 }
426 let cv = Arc::new(SharedCondvar::create(&path).expect("create"));
427 let pred = Arc::new(AtomicBool::new(false));
428 let cv2 = Arc::clone(&cv);
429 let pred2 = Arc::clone(&pred);
430 let waiter = thread::spawn(move || {
431 cv2.wait(|| pred2.load(Ordering::Acquire)).unwrap();
432 });
433 thread::sleep(Duration::from_millis(30));
434 pred.store(true, Ordering::Release);
435 cv.notify_all();
436 waiter.join().unwrap();
437
438 for suffix in [".waker.bin", ".gen.bin"] {
440 let mut p = path.as_os_str().to_owned();
441 p.push(suffix);
442 drop(std::fs::remove_file(PathBuf::from(p)));
443 }
444 }
445}