subetha_cxc/
shared_async_pointer.rs1use std::path::Path;
37use std::sync::Arc;
38use std::thread;
39
40use crate::shared_once_cell::{SharedOnceCell, SharedOnceError};
41
42pub mod strategy {
44 pub const RESOLVED: u32 = 0;
45 pub const LAZY: u32 = 1;
46 pub const SPECULATIVE: u32 = 2;
47}
48
49#[derive(Debug, Clone, Copy, PartialEq, Eq)]
50pub enum SharedAsyncError {
51 Once(SharedOnceError),
52 AllWorkersDied,
53}
54
55impl From<SharedOnceError> for SharedAsyncError {
56 fn from(e: SharedOnceError) -> Self { Self::Once(e) }
57}
58
59pub struct SharedAsyncPointer<T: Copy + Send + Sync + 'static> {
60 cell: Arc<SharedOnceCell<T>>,
61 header_sidecar: subetha_core::HandshakeHeader,
62 ring_sidecar: Box<subetha_core::ObservationRing>,
63}
64
65impl<T: Copy + Send + Sync + 'static>
66 subetha_sidecar::AdaptiveInstance for SharedAsyncPointer<T>
67{
68 fn header(&self) -> &subetha_core::HandshakeHeader { &self.header_sidecar }
69 fn ring(&self) -> &subetha_core::ObservationRing { &self.ring_sidecar }
70 fn make_policy(&self) -> Box<dyn subetha_sidecar::Policy> {
71 Box::new(subetha_sidecar::NoMigrationPolicy)
72 }
73}
74
75impl<T: Copy + Send + Sync + 'static> SharedAsyncPointer<T> {
76 pub const SIGNATURE: subetha_core::AxisMask = subetha_core::AxisMask::from_axes(
80 &[subetha_core::Axis::Async],
81 );
82
83 pub fn create(path: impl AsRef<Path>) -> Result<Self, SharedAsyncError> {
86 let cell = Arc::new(SharedOnceCell::create(path)?);
87 Ok(Self {
88 cell,
89 header_sidecar: subetha_core::HandshakeHeader::new(),
90 ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
91 })
92 }
93
94 pub fn open(path: impl AsRef<Path>) -> Result<Self, SharedAsyncError> {
96 let cell = Arc::new(SharedOnceCell::open(path)?);
97 Ok(Self {
98 cell,
99 header_sidecar: subetha_core::HandshakeHeader::new(),
100 ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
101 })
102 }
103
104 pub fn is_resolved(&self) -> bool {
106 self.cell.is_initialized()
107 }
108
109 pub fn try_get(&self) -> Option<T> {
111 let r = self.cell.get();
112 self.ring_sidecar.push_op(
113 crate::sidecar_ops::async_pointer::OP_TRY_GET,
114 if r.is_none() { 2 } else { 0 },
115 );
116 r
117 }
118
119 pub fn set_resolved(&self, value: T) -> bool {
123 self.cell.set(value)
124 }
125
126 pub fn get_or_lazy<F>(&self, f: F) -> T
132 where F: FnOnce() -> T,
133 {
134 let was_resolved = self.cell.is_initialized();
135 if let Some(v) = self.cell.get() {
136 self.ring_sidecar
137 .push_op(crate::sidecar_ops::async_pointer::OP_GET_OR_FETCH, 0);
138 return v;
139 }
140 let computed = f();
141 if !self.cell.set(computed) {
142 while !self.cell.is_initialized() {
146 std::hint::spin_loop();
147 }
148 }
149 let v = self.cell.get().expect("INITIALIZED after set or read-back");
150 self.ring_sidecar.push_op(
151 crate::sidecar_ops::async_pointer::OP_GET_OR_FETCH,
152 if was_resolved { 0 } else { 1 }, );
154 v
155 }
156
157 pub fn get_or_speculative<F>(&self, n: usize, f: F) -> T
166 where F: Fn() -> T + Send + Sync + 'static + Clone,
167 {
168 if let Some(v) = self.cell.get() { return v; }
169 assert!(n >= 1, "speculative race needs at least 1 worker");
170 let mut handles = Vec::with_capacity(n);
171 for _ in 0..n {
172 let cell = self.cell.clone();
173 let f = f.clone();
174 handles.push(thread::spawn(move || {
175 if cell.is_initialized() { return; }
178 let v = f();
179 cell.set(v);
182 }));
183 }
184 for h in handles { h.join().ok(); }
185 self.cell.get().expect("at least one worker should publish")
186 }
187
188 pub fn get_or_speculative_with<I, F>(&self, fs: I) -> T
191 where I: IntoIterator<Item = F>, F: FnOnce() -> T + Send + 'static,
192 {
193 if let Some(v) = self.cell.get() { return v; }
194 let mut handles = vec![];
195 for f in fs {
196 let cell = self.cell.clone();
197 handles.push(thread::spawn(move || {
198 if cell.is_initialized() { return; }
199 let v = f();
200 cell.set(v);
202 }));
203 }
204 assert!(!handles.is_empty(), "speculative race needs at least 1 closure");
205 for h in handles { h.join().ok(); }
206 self.cell.get().expect("at least one worker should publish")
207 }
208
209 pub fn get_or_speculative_resilient<F>(&self, n: usize, f: F) -> Result<T, SharedAsyncError>
214 where F: Fn() -> T + Send + Sync + 'static + Clone,
215 {
216 if let Some(v) = self.cell.get() { return Ok(v); }
217 assert!(n >= 1);
218 let mut handles = Vec::with_capacity(n);
219 for _ in 0..n {
220 let cell = self.cell.clone();
221 let f = f.clone();
222 handles.push(thread::spawn(move || {
223 if cell.is_initialized() { return; }
224 let v = f();
225 cell.set(v);
227 }));
228 }
229 for h in handles { h.join().ok(); }
231 self.cell.get().ok_or(SharedAsyncError::AllWorkersDied)
232 }
233
234 pub fn flush(&self) -> Result<(), SharedAsyncError> {
236 Ok(self.cell.flush()?)
237 }
238
239 pub fn flush_async(&self) -> Result<(), SharedAsyncError> {
243 Ok(self.cell.flush_async()?)
244 }
245}
246
247#[cfg(test)]
248mod tests {
249 use super::*;
250 use std::sync::atomic::{AtomicU32, Ordering};
251 use std::time::Duration;
252
253 fn tmp(name: &str) -> std::path::PathBuf {
254 let mut p = std::env::temp_dir();
255 let pid = std::process::id();
256 p.push(format!("subetha-async-{name}-{pid}.bin"));
257 p
258 }
259
260 #[test]
261 fn resolved_strategy_returns_immediately() {
262 let p = tmp("resolved");
263 let sap: SharedAsyncPointer<u64> = SharedAsyncPointer::create(&p).unwrap();
264 assert!(!sap.is_resolved());
265 sap.set_resolved(42);
266 assert!(sap.is_resolved());
267 assert_eq!(sap.try_get(), Some(42));
268 std::fs::remove_file(&p).ok();
269 }
270
271 #[test]
272 fn lazy_runs_closure_exactly_once_across_threads() {
273 let p = tmp("lazy-once");
274 let sap: Arc<SharedAsyncPointer<u64>>
275 = Arc::new(SharedAsyncPointer::create(&p).unwrap());
276 let counter = Arc::new(AtomicU32::new(0));
277 let mut handles = vec![];
278 for _ in 0..8 {
279 let sap = sap.clone();
280 let counter = counter.clone();
281 handles.push(thread::spawn(move || {
282 sap.get_or_lazy(|| {
283 counter.fetch_add(1, Ordering::AcqRel);
284 thread::sleep(Duration::from_millis(2));
285 777
286 })
287 }));
288 }
289 let results: Vec<u64> = handles.into_iter().map(|h| h.join().unwrap()).collect();
290 assert!(results.iter().all(|v| *v == 777));
292 let runs = counter.load(Ordering::Acquire);
294 assert!((1..=8).contains(&runs),
295 "closure runs should be between 1 and 8 (one per non-fast-pathed thread); got {runs}");
296 std::fs::remove_file(&p).ok();
303 }
304
305 #[test]
306 fn speculative_race_returns_one_published_result() {
307 let p = tmp("speculative");
308 let sap: SharedAsyncPointer<u64> = SharedAsyncPointer::create(&p).unwrap();
309 let runs = Arc::new(AtomicU32::new(0));
310 let runs_clone = runs.clone();
311 let result = sap.get_or_speculative(4, move || {
312 runs_clone.fetch_add(1, Ordering::AcqRel);
313 thread::sleep(Duration::from_millis(5));
315 123u64
316 });
317 assert_eq!(result, 123);
318 assert!(runs.load(Ordering::Acquire) >= 1,
319 "at least one worker must run");
320 assert!(sap.is_resolved());
322 std::fs::remove_file(&p).ok();
323 }
324
325 #[test]
326 fn speculative_first_finisher_wins() {
327 let p = tmp("first-wins");
328 let sap: SharedAsyncPointer<u64> = SharedAsyncPointer::create(&p).unwrap();
329 let result = sap.get_or_speculative_with([
332 Box::new(|| {
333 thread::sleep(Duration::from_millis(100));
334 999u64
335 }) as Box<dyn FnOnce() -> u64 + Send>,
336 Box::new(|| {
337 thread::sleep(Duration::from_millis(2));
338 100u64
339 }) as Box<dyn FnOnce() -> u64 + Send>,
340 ]);
341 assert_eq!(result, 100, "fast closure should win");
342 std::fs::remove_file(&p).ok();
343 }
344
345 #[test]
346 fn speculative_resilient_tolerates_panicking_workers() {
347 let p = tmp("resilient");
348 let sap: SharedAsyncPointer<u64> = SharedAsyncPointer::create(&p).unwrap();
349 let attempt = Arc::new(AtomicU32::new(0));
350 let attempt_clone = attempt.clone();
351 let result = sap.get_or_speculative_resilient(8, move || {
353 let n = attempt_clone.fetch_add(1, Ordering::AcqRel);
354 if n.is_multiple_of(2) {
355 panic!("simulated worker death on attempt {n}");
356 }
357 42u64
358 }).expect("at least one survivor publishes");
359 assert_eq!(result, 42);
360 std::fs::remove_file(&p).ok();
361 }
362
363 #[test]
364 fn second_call_after_resolution_returns_cached_value() {
365 let p = tmp("cached");
366 let sap: SharedAsyncPointer<u64> = SharedAsyncPointer::create(&p).unwrap();
367 let _r1 = sap.get_or_lazy(|| 5);
368 let r2 = sap.get_or_lazy(|| panic!("must not run after resolution"));
371 assert_eq!(r2, 5);
372 std::fs::remove_file(&p).ok();
373 }
374
375 #[test]
376 fn cross_handle_speculative_race_shares_one_winner() {
377 let p = tmp("cross-handle-spec");
378 let sap_a = SharedAsyncPointer::<u64>::create(&p).unwrap();
379 let sap_b = SharedAsyncPointer::<u64>::open(&p).unwrap();
380 let r_a = sap_a.get_or_speculative(2, || 9999u64);
382 assert_eq!(r_a, 9999);
383 let r_b = sap_b.get_or_lazy(|| panic!("must not run on already-resolved cell"));
385 assert_eq!(r_b, 9999);
386 std::fs::remove_file(&p).ok();
387 }
388
389 #[test]
390 fn try_get_does_not_force_resolution() {
391 let p = tmp("try-get");
392 let sap: SharedAsyncPointer<u64> = SharedAsyncPointer::create(&p).unwrap();
393 assert_eq!(sap.try_get(), None);
394 assert!(!sap.is_resolved());
395 sap.set_resolved(100);
396 assert_eq!(sap.try_get(), Some(100));
397 std::fs::remove_file(&p).ok();
398 }
399
400 #[test]
401 fn speculative_with_one_worker_is_just_lazy() {
402 let p = tmp("spec-1");
403 let sap: SharedAsyncPointer<u64> = SharedAsyncPointer::create(&p).unwrap();
404 let runs = Arc::new(AtomicU32::new(0));
405 let runs_clone = runs.clone();
406 let r = sap.get_or_speculative(1, move || {
407 runs_clone.fetch_add(1, Ordering::AcqRel);
408 17u64
409 });
410 assert_eq!(r, 17);
411 assert_eq!(runs.load(Ordering::Acquire), 1);
412 std::fs::remove_file(&p).ok();
413 }
414
415 #[test]
419 fn signature_engages_the_async_axis() {
420 let sig = SharedAsyncPointer::<u8>::SIGNATURE;
421 assert_ne!(sig, subetha_core::AxisMask::EMPTY);
422 assert!(sig.contains(subetha_core::Axis::Async));
423 }
424}