1use std::{
4 path::{Path, PathBuf},
5 sync::{
6 Arc, Mutex, OnceLock,
7 atomic::{AtomicU8, AtomicU64, Ordering},
8 },
9};
10
11#[derive(Clone, Debug, Eq, PartialEq)]
13pub struct PlatformDirectories {
14 pub data: PathBuf,
15 pub config: PathBuf,
16 pub cache: PathBuf,
17 pub documents: Option<PathBuf>,
18 pub temporary: PathBuf,
19 pub shared: Option<PathBuf>,
20}
21
22impl PlatformDirectories {
23 fn scoped(&self, application_id: &str) -> Self {
24 Self {
25 data: self.data.join(application_id),
26 config: self.config.join(application_id),
27 cache: self.cache.join(application_id),
28 documents: self
29 .documents
30 .as_ref()
31 .map(|path| path.join(application_id)),
32 temporary: self.temporary.join(application_id),
33 shared: self.shared.as_ref().map(|path| path.join(application_id)),
34 }
35 }
36}
37
38#[derive(Clone, Debug, Eq, PartialEq, thiserror::Error)]
39pub enum PlatformDirectoryError {
40 #[error("application id must be one non-empty path component")]
41 InvalidApplicationId,
42 #[error("no application id has been registered")]
43 NoApplicationId,
44 #[error("platform directories are unavailable")]
45 Unavailable,
46}
47
48#[derive(Clone, Copy, Debug, Eq, PartialEq)]
50pub enum LifecycleState {
51 Created,
52 Started,
53 Resumed,
54 Paused,
55 Stopped,
56 Destroyed,
57}
58
59#[derive(Clone, Copy, Debug, Eq, PartialEq)]
61pub struct LifecycleEvent {
62 pub from: LifecycleState,
63 pub to: LifecycleState,
64}
65
66pub trait HostController: Send + Sync {
68 fn set_keep_screen_on(&self, enabled: bool);
70 fn platform_directories(&self) -> Option<PlatformDirectories>;
72 fn exit(&self);
74 fn background(&self);
76 fn durable_save_deadline(&self) -> std::time::Duration {
81 DEFAULT_DURABLE_SAVE_DEADLINE
82 }
83}
84
85pub const DEFAULT_DURABLE_SAVE_DEADLINE: std::time::Duration = std::time::Duration::from_secs(2);
87
88pub fn durable_save_deadline() -> std::time::Duration {
90 host_controller().map_or(DEFAULT_DURABLE_SAVE_DEADLINE, |host| {
91 host.durable_save_deadline()
92 })
93}
94
95pub type HostControllerRef = Arc<dyn HostController>;
96
97fn controller() -> &'static Mutex<Option<HostControllerRef>> {
98 static SLOT: OnceLock<Mutex<Option<HostControllerRef>>> = OnceLock::new();
99 SLOT.get_or_init(|| Mutex::new(None))
100}
101
102pub fn set_host_controller(value: HostControllerRef) {
104 if let Ok(mut slot) = controller().lock() {
105 *slot = Some(value);
106 }
107 KEEP_SCREEN_ON.store(0, Ordering::Release);
108}
109pub fn clear_host_controller() {
111 if let Ok(mut slot) = controller().lock() {
112 *slot = None;
113 }
114 KEEP_SCREEN_ON.store(0, Ordering::Release);
115}
116pub fn host_controller() -> Option<HostControllerRef> {
118 controller().lock().ok().and_then(|slot| slot.clone())
119}
120pub fn set_keep_screen_on(enabled: bool) {
122 let value = if enabled { 2 } else { 1 };
123 if KEEP_SCREEN_ON.swap(value, Ordering::AcqRel) == value {
124 return;
125 }
126 if let Some(host) = host_controller() {
127 host.set_keep_screen_on(enabled);
128 }
129}
130fn valid_application_id(application_id: &str) -> bool {
131 !application_id.is_empty()
132 && Path::new(application_id).components().count() == 1
133 && application_id != "."
134 && application_id != ".."
135}
136
137fn desktop_platform_directories() -> Option<PlatformDirectories> {
138 let base = directories::BaseDirs::new()?;
139 let documents = directories::UserDirs::new()
140 .and_then(|directories| directories.document_dir().map(Path::to_path_buf));
141 Some(PlatformDirectories {
142 data: base.data_dir().to_path_buf(),
143 config: base.config_dir().to_path_buf(),
144 cache: base.cache_dir().to_path_buf(),
145 documents,
146 temporary: std::env::temp_dir(),
147 shared: None,
148 })
149}
150
151fn application_id_slot() -> &'static Mutex<Option<String>> {
152 static SLOT: OnceLock<Mutex<Option<String>>> = OnceLock::new();
153 SLOT.get_or_init(|| Mutex::new(None))
154}
155
156pub fn set_application_id(application_id: &str) -> Result<(), PlatformDirectoryError> {
163 if !valid_application_id(application_id) {
164 return Err(PlatformDirectoryError::InvalidApplicationId);
165 }
166 if let Ok(mut slot) = application_id_slot().lock() {
167 *slot = Some(application_id.to_string());
168 }
169 Ok(())
170}
171
172pub fn clear_application_id() {
174 if let Ok(mut slot) = application_id_slot().lock() {
175 *slot = None;
176 }
177}
178
179pub fn application_id() -> Option<String> {
181 application_id_slot()
182 .lock()
183 .ok()
184 .and_then(|slot| slot.clone())
185}
186
187pub fn application_directories() -> Result<PlatformDirectories, PlatformDirectoryError> {
189 let application_id = application_id().ok_or(PlatformDirectoryError::NoApplicationId)?;
190 let roots = host_controller()
191 .and_then(|host| host.platform_directories())
192 .or_else(desktop_platform_directories)
193 .ok_or(PlatformDirectoryError::Unavailable)?;
194 Ok(roots.scoped(&application_id))
195}
196pub fn exit_app() {
198 if let Some(host) = host_controller() {
199 host.exit();
200 }
201}
202pub fn background_app() {
204 if let Some(host) = host_controller() {
205 host.background();
206 }
207}
208
209#[cfg(not(target_arch = "wasm32"))]
210type Observer = Arc<dyn Fn(LifecycleEvent) + Send + Sync>;
211#[cfg(target_arch = "wasm32")]
212type Observer = std::rc::Rc<dyn Fn(LifecycleEvent)>;
213
214#[cfg(not(target_arch = "wasm32"))]
215fn observers() -> &'static Mutex<Vec<(u64, Observer)>> {
216 static SLOT: OnceLock<Mutex<Vec<(u64, Observer)>>> = OnceLock::new();
217 SLOT.get_or_init(|| Mutex::new(Vec::new()))
218}
219#[cfg(target_arch = "wasm32")]
220thread_local! {
221 static OBSERVERS: std::cell::RefCell<Vec<(u64, Observer)>> = const { std::cell::RefCell::new(Vec::new()) };
222}
223static NEXT_ID: AtomicU64 = AtomicU64::new(1);
224static LIFECYCLE_STATE: AtomicU8 = AtomicU8::new(LifecycleState::Created as u8);
225static KEEP_SCREEN_ON: AtomicU8 = AtomicU8::new(0);
226
227pub struct LifecycleObserver {
229 id: u64,
230}
231impl Drop for LifecycleObserver {
232 fn drop(&mut self) {
233 #[cfg(not(target_arch = "wasm32"))]
234 if let Ok(mut list) = observers().lock() {
235 list.retain(|(id, _)| *id != self.id);
236 }
237 #[cfg(target_arch = "wasm32")]
238 OBSERVERS.with(|list| list.borrow_mut().retain(|(id, _)| *id != self.id));
239 }
240}
241#[cfg(not(target_arch = "wasm32"))]
243pub fn observe_lifecycle(
244 observer: impl Fn(LifecycleEvent) + Send + Sync + 'static,
245) -> LifecycleObserver {
246 let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
247 if let Ok(mut list) = observers().lock() {
248 list.push((id, Arc::new(observer)));
249 }
250 LifecycleObserver { id }
251}
252#[cfg(target_arch = "wasm32")]
254pub fn observe_lifecycle(observer: impl Fn(LifecycleEvent) + 'static) -> LifecycleObserver {
255 let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
256 OBSERVERS.with(|list| list.borrow_mut().push((id, std::rc::Rc::new(observer))));
257 LifecycleObserver { id }
258}
259pub fn dispatch_lifecycle(event: LifecycleEvent) {
261 LIFECYCLE_STATE.store(event.to as u8, Ordering::Release);
262 crate::media::on_lifecycle(event);
263 #[cfg(not(target_arch = "wasm32"))]
264 let callbacks = observers()
265 .lock()
266 .map(|list| {
267 list.iter()
268 .map(|(_, cb)| Arc::clone(cb))
269 .collect::<Vec<_>>()
270 })
271 .unwrap_or_default();
272 #[cfg(target_arch = "wasm32")]
273 let callbacks = OBSERVERS.with(|list| {
274 list.borrow()
275 .iter()
276 .map(|(_, callback)| std::rc::Rc::clone(callback))
277 .collect::<Vec<_>>()
278 });
279 for callback in callbacks {
280 callback(event);
281 }
282}
283
284pub fn current_lifecycle_state() -> LifecycleState {
286 match LIFECYCLE_STATE.load(Ordering::Acquire) {
287 0 => LifecycleState::Created,
288 1 => LifecycleState::Started,
289 2 => LifecycleState::Resumed,
290 3 => LifecycleState::Paused,
291 4 => LifecycleState::Stopped,
292 _ => LifecycleState::Destroyed,
293 }
294}
295
296pub fn dispatch_lifecycle_state(to: LifecycleState) {
303 let from = current_lifecycle_state();
304 if from == to {
305 return;
306 }
307 #[cfg(not(target_arch = "wasm32"))]
308 if matches!(to, LifecycleState::Paused) {
309 let outcome = run_durable_saves(durable_save_deadline());
310 if outcome == DurableSaveOutcome::TimedOut {
311 log::warn!("cranpose: durable saves overran the host deadline; they keep running");
312 }
313 }
314 dispatch_lifecycle(LifecycleEvent { from, to });
315}
316
317pub fn local_lifecycle_state() -> cranpose_core::CompositionLocal<LifecycleState> {
323 thread_local! {
324 static LOCAL: std::cell::RefCell<Option<cranpose_core::CompositionLocal<LifecycleState>>> =
325 const { std::cell::RefCell::new(None) };
326 }
327 LOCAL.with(|cell| {
328 cell.borrow_mut()
329 .get_or_insert_with(|| cranpose_core::compositionLocalOf(current_lifecycle_state))
330 .clone()
331 })
332}
333
334#[allow(non_snake_case)]
336#[track_caller]
337pub fn rememberLifecycleState() -> cranpose_core::State<LifecycleState> {
338 let transitions = rememberLifecycleEvents();
339 let state = cranpose_core::collectAsState(
340 transitions,
341 (),
342 LifecycleEvent {
343 from: current_lifecycle_state(),
344 to: current_lifecycle_state(),
345 },
346 );
347 cranpose_core::derivedStateOf(move || state.get().to)
348}
349
350#[allow(non_snake_case)]
352#[track_caller]
353pub fn rememberLifecycleEvents() -> cranpose_core::EventStream<LifecycleEvent> {
354 cranpose_core::rememberEventStream((), |sender| {
355 observe_lifecycle(move |event| sender.send(event))
356 })
357}
358
359#[cranpose_macros::composable]
364pub fn ProvideLifecycle(content: impl FnOnce()) {
365 let state = rememberLifecycleState();
366 let local = local_lifecycle_state();
367 cranpose_core::CompositionLocalProvider(vec![local.provides(state.get())], move || {
368 content();
369 });
370}
371
372#[derive(Clone, Copy, Debug, Eq, PartialEq)]
374pub enum DurableSaveOutcome {
375 Nothing,
377 Completed,
379 TimedOut,
382}
383
384type SaveWork = Arc<dyn Fn() + Send + Sync>;
385
386fn durable_saves() -> &'static Mutex<Vec<(u64, SaveWork)>> {
387 static SLOT: OnceLock<Mutex<Vec<(u64, SaveWork)>>> = OnceLock::new();
388 SLOT.get_or_init(|| Mutex::new(Vec::new()))
389}
390
391pub struct DurableSaveRegistration {
393 id: u64,
394}
395
396impl Drop for DurableSaveRegistration {
397 fn drop(&mut self) {
398 if let Ok(mut saves) = durable_saves().lock() {
399 saves.retain(|(id, _)| *id != self.id);
400 }
401 }
402}
403
404pub fn register_durable_save(save: impl Fn() + Send + Sync + 'static) -> DurableSaveRegistration {
409 let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
410 if let Ok(mut saves) = durable_saves().lock() {
411 saves.push((id, Arc::new(save)));
412 }
413 DurableSaveRegistration { id }
414}
415
416#[allow(non_snake_case)]
422#[track_caller]
423pub fn DurableSaveEffect<K: PartialEq + 'static>(keys: K, save: impl Fn() + Send + Sync + 'static) {
424 cranpose_core::__disposable_effect_impl(
425 cranpose_core::caller_location_key()
426 ^ cranpose_core::location_key(file!(), line!(), column!()),
427 keys,
428 move |scope| {
429 let registration = register_durable_save(save);
430 scope.on_dispose(move || drop(registration))
431 },
432 );
433}
434
435#[cfg(not(target_arch = "wasm32"))]
443pub fn run_durable_saves(deadline: std::time::Duration) -> DurableSaveOutcome {
444 let saves: Vec<SaveWork> = durable_saves()
445 .lock()
446 .map(|saves| saves.iter().map(|(_, save)| Arc::clone(save)).collect())
447 .unwrap_or_default();
448 if saves.is_empty() {
449 return DurableSaveOutcome::Nothing;
450 }
451
452 let outstanding = Arc::new(std::sync::atomic::AtomicUsize::new(saves.len()));
453 let finished = Arc::new((Mutex::new(false), std::sync::Condvar::new()));
454 for save in saves {
455 let worker_outstanding = Arc::clone(&outstanding);
456 let worker_finished = Arc::clone(&finished);
457 let lease = crate::background::acquire_background_work();
458 let spawned = std::thread::Builder::new()
459 .name("cranpose-durable-save".to_string())
460 .spawn(move || {
461 save();
462 drop(lease);
463 if worker_outstanding.fetch_sub(1, Ordering::AcqRel) == 1 {
464 let (done, wake) = &*worker_finished;
465 if let Ok(mut done) = done.lock() {
466 *done = true;
467 }
468 wake.notify_all();
469 }
470 });
471 if spawned.is_err() {
472 log::warn!("cranpose: a durable save could not be started");
473 if outstanding.fetch_sub(1, Ordering::AcqRel) == 1 {
474 let (done, wake) = &*finished;
475 if let Ok(mut done) = done.lock() {
476 *done = true;
477 }
478 wake.notify_all();
479 }
480 }
481 }
482
483 let (done, wake) = &*finished;
484 let Ok(mut guard) = done.lock() else {
485 return DurableSaveOutcome::TimedOut;
486 };
487 let mut remaining = deadline;
488 let started = web_time::Instant::now();
489 while !*guard {
490 let Ok((next, timeout)) = wake.wait_timeout(guard, remaining) else {
491 return DurableSaveOutcome::TimedOut;
492 };
493 guard = next;
494 if timeout.timed_out() {
495 break;
496 }
497 remaining = deadline.saturating_sub(started.elapsed());
498 if remaining.is_zero() {
499 break;
500 }
501 }
502 if *guard {
503 DurableSaveOutcome::Completed
504 } else {
505 DurableSaveOutcome::TimedOut
506 }
507}
508
509#[cfg(test)]
510mod tests {
511 use std::sync::PoisonError;
512
513 use super::*;
514
515 fn test_lock() -> std::sync::MutexGuard<'static, ()> {
516 static LOCK: Mutex<()> = Mutex::new(());
517 LOCK.lock().unwrap_or_else(PoisonError::into_inner)
518 }
519
520 #[test]
521 fn observer_is_removed_on_drop() {
522 let _guard = test_lock();
523 let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
524 let seen = Arc::clone(&calls);
525 let handle = observe_lifecycle(move |_| {
526 seen.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
527 });
528 dispatch_lifecycle(LifecycleEvent {
529 from: LifecycleState::Created,
530 to: LifecycleState::Started,
531 });
532 assert_eq!(calls.load(std::sync::atomic::Ordering::Relaxed), 1);
533 drop(handle);
534 dispatch_lifecycle(LifecycleEvent {
535 from: LifecycleState::Started,
536 to: LifecycleState::Resumed,
537 });
538 assert_eq!(calls.load(std::sync::atomic::Ordering::Relaxed), 1);
539 }
540
541 #[test]
542 fn state_dispatch_derives_the_previous_state() {
543 let _guard = test_lock();
544 dispatch_lifecycle_state(LifecycleState::Paused);
545 assert_eq!(current_lifecycle_state(), LifecycleState::Paused);
546 dispatch_lifecycle_state(LifecycleState::Stopped);
547 assert_eq!(current_lifecycle_state(), LifecycleState::Stopped);
548 }
549
550 #[test]
551 fn repeated_keep_screen_value_reaches_the_host_once() {
552 let _guard = test_lock();
553 struct RecordingHost(Arc<std::sync::atomic::AtomicUsize>);
554 impl HostController for RecordingHost {
555 fn set_keep_screen_on(&self, _enabled: bool) {
556 self.0.fetch_add(1, Ordering::Relaxed);
557 }
558 fn platform_directories(&self) -> Option<PlatformDirectories> {
559 Some(PlatformDirectories {
560 data: PathBuf::from("data"),
561 config: PathBuf::from("config"),
562 cache: PathBuf::from("cache"),
563 documents: Some(PathBuf::from("documents")),
564 temporary: PathBuf::from("temporary"),
565 shared: Some(PathBuf::from("shared")),
566 })
567 }
568 fn exit(&self) {}
569 fn background(&self) {}
570 }
571 let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
572 set_host_controller(Arc::new(RecordingHost(Arc::clone(&calls))));
573 set_keep_screen_on(true);
574 set_keep_screen_on(true);
575 assert_eq!(calls.load(Ordering::Relaxed), 1);
576 set_application_id("sample").expect("a plain id is valid");
577 assert_eq!(
578 application_directories().unwrap().data,
579 PathBuf::from("data/sample")
580 );
581 clear_application_id();
582 clear_host_controller();
583 }
584
585 #[test]
586 fn durable_saves_run_and_report_completion() {
587 let _services = crate::registry::test_service_guard();
588 let _guard = test_lock();
589 let ran = Arc::new(std::sync::atomic::AtomicUsize::new(0));
590 let first = Arc::clone(&ran);
591 let second = Arc::clone(&ran);
592 let a = register_durable_save(move || {
593 first.fetch_add(1, Ordering::Relaxed);
594 });
595 let b = register_durable_save(move || {
596 second.fetch_add(1, Ordering::Relaxed);
597 });
598 assert_eq!(
599 run_durable_saves(std::time::Duration::from_secs(5)),
600 DurableSaveOutcome::Completed
601 );
602 assert_eq!(ran.load(Ordering::Relaxed), 2);
603 drop((a, b));
604 assert_eq!(
605 run_durable_saves(std::time::Duration::from_secs(1)),
606 DurableSaveOutcome::Nothing
607 );
608 }
609
610 #[test]
611 fn a_save_that_overruns_the_deadline_reports_a_timeout() {
612 let _services = crate::registry::test_service_guard();
613 let _guard = test_lock();
614 let registration = register_durable_save(|| {
615 std::thread::sleep(std::time::Duration::from_millis(400));
616 });
617 assert_eq!(
618 run_durable_saves(std::time::Duration::from_millis(30)),
619 DurableSaveOutcome::TimedOut
620 );
621 drop(registration);
622 }
623
624 #[test]
625 fn a_dropped_registration_is_no_longer_saved() {
626 let _services = crate::registry::test_service_guard();
627 let _guard = test_lock();
628 let ran = Arc::new(std::sync::atomic::AtomicUsize::new(0));
629 let counted = Arc::clone(&ran);
630 let registration = register_durable_save(move || {
631 counted.fetch_add(1, Ordering::Relaxed);
632 });
633 drop(registration);
634 assert_eq!(
635 run_durable_saves(std::time::Duration::from_secs(1)),
636 DurableSaveOutcome::Nothing
637 );
638 assert_eq!(ran.load(Ordering::Relaxed), 0);
639 }
640
641 #[test]
642 fn application_id_must_be_one_component() {
643 let _guard = test_lock();
644 assert_eq!(
645 set_application_id("../sample"),
646 Err(PlatformDirectoryError::InvalidApplicationId)
647 );
648 assert_eq!(
649 set_application_id(""),
650 Err(PlatformDirectoryError::InvalidApplicationId)
651 );
652 clear_application_id();
653 assert_eq!(
654 application_directories(),
655 Err(PlatformDirectoryError::NoApplicationId)
656 );
657 }
658
659 #[test]
660 fn a_surviving_durable_save_keeps_its_registration_when_a_leader_leaves() {
661 let _guard = test_lock();
662 durable_saves()
663 .lock()
664 .unwrap_or_else(PoisonError::into_inner)
665 .clear();
666 let ran: Arc<Mutex<Vec<&'static str>>> = Arc::new(Mutex::new(Vec::new()));
667 let show_first = std::rc::Rc::new(std::cell::Cell::new(true));
668
669 fn saves(show_first: bool, ran: &Arc<Mutex<Vec<&'static str>>>) {
670 if show_first {
671 let ran = Arc::clone(ran);
672 DurableSaveEffect((), move || {
673 ran.lock()
674 .unwrap_or_else(PoisonError::into_inner)
675 .push("first");
676 });
677 }
678 let ran = Arc::clone(ran);
679 DurableSaveEffect((), move || {
680 ran.lock()
681 .unwrap_or_else(PoisonError::into_inner)
682 .push("tail");
683 });
684 }
685
686 let mut composition = cranpose_core::Composition::new(cranpose_core::MemoryApplier::new());
687 let root_key = cranpose_core::location_key(file!(), line!(), column!());
688 let mut pass = {
689 let ran = Arc::clone(&ran);
690 let show_first = std::rc::Rc::clone(&show_first);
691 move || saves(show_first.get(), &ran)
692 };
693
694 composition
695 .render(root_key, &mut pass)
696 .expect("initial composition");
697 show_first.set(false);
698 composition
699 .render(root_key, &mut pass)
700 .expect("drop the leading save");
701
702 let outcome = run_durable_saves(std::time::Duration::from_secs(5));
703 assert_eq!(outcome, DurableSaveOutcome::Completed);
704 assert_eq!(
705 ran.lock()
706 .unwrap_or_else(PoisonError::into_inner)
707 .as_slice(),
708 ["tail"],
709 "the surviving effect must keep its own registration; adopting the \
710 departed leader's group keeps the wrong save alive"
711 );
712 }
713}