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> {
137 let file = OpenOptions::new()
138 .read(true).write(true).create(true).truncate(true)
139 .open(path)?;
140 file.set_len(GEN_REGION_SIZE as u64)?;
141 let mut mmap = unsafe { MmapOptions::new().len(GEN_REGION_SIZE).map_mut(&file)? };
142 let base = mmap.as_mut_ptr();
143 unsafe {
144 (base as *mut u64).write(CONDVAR_GEN_MAGIC);
145 (base.add(GEN_OFFSET) as *mut AtomicU64).write(AtomicU64::new(0));
146 }
147 let ptr = unsafe { base.add(GEN_OFFSET) as *const AtomicU64 };
148 Ok(Self { backing: GenBacking::File(file, mmap), ptr })
149 }
150
151 fn open_file(path: &Path) -> Result<Self, CondvarError> {
152 let file = OpenOptions::new().read(true).write(true).open(path)?;
153 let meta = file.metadata()?;
154 if (meta.len() as usize) < GEN_REGION_SIZE {
155 return Err(CondvarError::LayoutMismatch);
156 }
157 let mut mmap = unsafe { MmapOptions::new().len(GEN_REGION_SIZE).map_mut(&file)? };
158 let base = mmap.as_mut_ptr();
159 let magic = unsafe { (base as *const u64).read() };
160 if magic != CONDVAR_GEN_MAGIC {
161 return Err(CondvarError::LayoutMismatch);
162 }
163 let ptr = unsafe { base.add(GEN_OFFSET) as *const AtomicU64 };
164 Ok(Self { backing: GenBacking::File(file, mmap), ptr })
165 }
166
167 #[inline]
168 fn atom(&self) -> &AtomicU64 {
169 unsafe { &*self.ptr }
173 }
174}
175
176pub struct SharedCondvar {
180 waker: Arc<CrossProcessWaker>,
181 gen_atom: Arc<GenAtom>,
182}
183
184impl SharedCondvar {
185 pub fn create_anon() -> Result<Self, CondvarError> {
188 Self::create_anon_with_capacity(MAX_WAITERS_DEFAULT)
189 }
190
191 pub fn create_anon_with_capacity(max_waiters: usize) -> Result<Self, CondvarError> {
193 let waker = Arc::new(CrossProcessWaker::create_anon(max_waiters)?);
194 let gen_atom = Arc::new(GenAtom::create_anon()?);
195 Ok(Self { waker, gen_atom })
196 }
197
198 pub fn create(base_path: impl AsRef<Path>) -> Result<Self, CondvarError> {
202 Self::create_with_capacity(base_path, MAX_WAITERS_DEFAULT)
203 }
204
205 pub fn create_with_capacity(
206 base_path: impl AsRef<Path>,
207 max_waiters: usize,
208 ) -> Result<Self, CondvarError> {
209 let (waker_path, gen_path) = side_paths(base_path.as_ref());
210 let waker = Arc::new(CrossProcessWaker::create(waker_path, max_waiters)?);
211 let gen_atom = Arc::new(GenAtom::create_file(&gen_path)?);
212 Ok(Self { waker, gen_atom })
213 }
214
215 pub fn open(base_path: impl AsRef<Path>) -> Result<Self, CondvarError> {
219 Self::open_with_capacity(base_path, MAX_WAITERS_DEFAULT)
220 }
221
222 pub fn open_with_capacity(
223 base_path: impl AsRef<Path>,
224 expected_max_waiters: usize,
225 ) -> Result<Self, CondvarError> {
226 let (waker_path, gen_path) = side_paths(base_path.as_ref());
227 let waker = Arc::new(CrossProcessWaker::open(waker_path, expected_max_waiters)?);
228 let gen_atom = Arc::new(GenAtom::open_file(&gen_path)?);
229 Ok(Self { waker, gen_atom })
230 }
231
232 pub fn wait<F: FnMut() -> bool>(&self, mut predicate: F) -> Result<(), CondvarError> {
236 loop {
237 if predicate() {
238 return Ok(());
239 }
240 let snapshot = self.gen_atom.atom().load(Ordering::Acquire);
246 let token = self.waker.try_park(snapshot + 1)?;
247 if predicate() {
249 self.waker.release(token);
250 return Ok(());
251 }
252 self.waker.wait(token, None)?;
253 }
254 }
255
256 pub fn wait_timeout<F: FnMut() -> bool>(
260 &self,
261 mut predicate: F,
262 timeout: Duration,
263 ) -> Result<(), CondvarError> {
264 let deadline = Instant::now() + timeout;
265 loop {
266 if predicate() {
267 return Ok(());
268 }
269 let snapshot = self.gen_atom.atom().load(Ordering::Acquire);
270 let token = self.waker.try_park(snapshot + 1)?;
271 if predicate() {
272 self.waker.release(token);
273 return Ok(());
274 }
275 let now = Instant::now();
276 if now >= deadline {
277 self.waker.release(token);
278 return Err(CondvarError::Timeout);
279 }
280 let remaining = deadline - now;
281 match self.waker.wait(token, Some(remaining)) {
282 Ok(()) => continue,
283 Err(WakerError::Timeout) => {
284 if predicate() {
285 return Ok(());
286 }
287 return Err(CondvarError::Timeout);
288 }
289 Err(e) => return Err(CondvarError::from(e)),
290 }
291 }
292 }
293
294 pub fn notify_one(&self) -> usize {
298 let new_gen = self.gen_atom.atom().fetch_add(1, Ordering::Release) + 1;
299 self.waker.wake_one_up_to(new_gen)
300 }
301
302 pub fn notify_all(&self) -> usize {
304 let new_gen = self.gen_atom.atom().fetch_add(1, Ordering::Release) + 1;
305 self.waker.wake_up_to(new_gen)
306 }
307
308 pub fn generation(&self) -> u64 {
311 self.gen_atom.atom().load(Ordering::Acquire)
312 }
313
314 pub fn waker(&self) -> &Arc<CrossProcessWaker> { &self.waker }
317}
318
319fn side_paths(base: &Path) -> (PathBuf, PathBuf) {
320 let mut w = base.as_os_str().to_owned();
321 w.push(".waker.bin");
322 let mut g = base.as_os_str().to_owned();
323 g.push(".gen.bin");
324 (PathBuf::from(w), PathBuf::from(g))
325}
326
327#[cfg(test)]
328mod tests {
329 use super::*;
330 use std::sync::atomic::AtomicBool;
331 use std::thread;
332
333 #[test]
334 fn notify_one_wakes_exactly_one_waiter() {
335 let cv = Arc::new(SharedCondvar::create_anon().expect("create"));
336 let pred = Arc::new(AtomicBool::new(false));
337 let waiters: Vec<_> = (0..3)
338 .map(|_| {
339 let cv2 = Arc::clone(&cv);
340 let pred2 = Arc::clone(&pred);
341 thread::spawn(move || {
342 cv2.wait(|| pred2.load(Ordering::Acquire)).unwrap();
343 })
344 })
345 .collect();
346 thread::sleep(Duration::from_millis(50));
347
348 assert_eq!(cv.notify_one(), 1);
352 thread::sleep(Duration::from_millis(20));
354 pred.store(true, Ordering::Release);
355 cv.notify_all();
356 for h in waiters {
357 h.join().unwrap();
358 }
359 }
360
361 #[test]
362 fn notify_all_wakes_every_waiter() {
363 let cv = Arc::new(SharedCondvar::create_anon().expect("create"));
364 let pred = Arc::new(AtomicBool::new(false));
365 let waiters: Vec<_> = (0..4)
366 .map(|_| {
367 let cv2 = Arc::clone(&cv);
368 let pred2 = Arc::clone(&pred);
369 thread::spawn(move || {
370 cv2.wait(|| pred2.load(Ordering::Acquire)).unwrap();
371 })
372 })
373 .collect();
374 thread::sleep(Duration::from_millis(30));
375 pred.store(true, Ordering::Release);
376 let woken = cv.notify_all();
377 assert!(woken >= 1, "at least one waiter woken (got {woken})");
378 for h in waiters {
379 h.join().unwrap();
380 }
381 }
382
383 #[test]
384 fn wait_timeout_returns_timeout() {
385 let cv = SharedCondvar::create_anon().expect("create");
386 let t0 = Instant::now();
387 let err = cv.wait_timeout(|| false, Duration::from_millis(60));
388 assert_eq!(err, Err(CondvarError::Timeout));
389 assert!(t0.elapsed() >= Duration::from_millis(50));
390 }
391
392 #[test]
393 fn wait_returns_immediately_if_predicate_already_true() {
394 let cv = SharedCondvar::create_anon().expect("create");
395 let pred = AtomicBool::new(true);
396 let t0 = Instant::now();
397 cv.wait(|| pred.load(Ordering::Acquire)).unwrap();
398 assert!(t0.elapsed() < Duration::from_millis(10));
399 }
400
401 #[test]
414 fn file_backed_create_then_arc_clone_round_trip() {
415 let dir = std::env::temp_dir();
416 let path = dir.join(format!("subetha_condvar_test_{}", std::process::id()));
417 for suffix in [".waker.bin", ".gen.bin"] {
419 let mut p = path.as_os_str().to_owned();
420 p.push(suffix);
421 drop(std::fs::remove_file(PathBuf::from(p)));
422 }
423 let cv = Arc::new(SharedCondvar::create(&path).expect("create"));
424 let pred = Arc::new(AtomicBool::new(false));
425 let cv2 = Arc::clone(&cv);
426 let pred2 = Arc::clone(&pred);
427 let waiter = thread::spawn(move || {
428 cv2.wait(|| pred2.load(Ordering::Acquire)).unwrap();
429 });
430 thread::sleep(Duration::from_millis(30));
431 pred.store(true, Ordering::Release);
432 cv.notify_all();
433 waiter.join().unwrap();
434
435 for suffix in [".waker.bin", ".gen.bin"] {
437 let mut p = path.as_os_str().to_owned();
438 p.push(suffix);
439 drop(std::fs::remove_file(PathBuf::from(p)));
440 }
441 }
442}