Skip to main content

cranpose_services/
memory_pressure.rs

1//! Memory pressure the platform reports for this process.
2//!
3//! Android delivers it through `onTrimMemory`; other hosts publish their own
4//! signal. There is no backlog: pressure describes a moment, so an observer
5//! that registers later waits for the next report. Applications collect the
6//! stream and give back what they can rebuild — caches, warm model sessions,
7//! pools.
8
9use std::sync::{
10    Arc, Mutex, OnceLock,
11    atomic::{AtomicU64, Ordering},
12};
13
14use cranpose_core::{EventStream, rememberEventStream};
15
16/// How hard the platform asks for memory back.
17#[derive(Clone, Copy, Debug, PartialEq, Eq)]
18pub enum MemoryPressure {
19    /// The UI left the screen. Anything held only to draw the next frame fast
20    /// can go.
21    UiHidden,
22    /// The process should give back what it can rebuild.
23    Low,
24    /// The system reclaims by force next. Free everything that can go.
25    Critical,
26}
27
28impl MemoryPressure {
29    /// Maps an Android `ComponentCallbacks2` trim level.
30    pub fn from_android_trim_level(level: i32) -> Self {
31        match level {
32            20 => Self::UiHidden,
33            level if level >= 60 || level == 15 => Self::Critical,
34            _ => Self::Low,
35        }
36    }
37}
38
39type Observer = Arc<dyn Fn(MemoryPressure) + Send + Sync>;
40
41struct Registry {
42    observers: Vec<(u64, Observer)>,
43}
44
45impl Registry {
46    fn new() -> Self {
47        Self {
48            observers: Vec::new(),
49        }
50    }
51
52    fn observe(&mut self, id: u64, observer: Observer) {
53        self.observers.push((id, observer));
54    }
55
56    /// The observers a report should reach. Handed back rather than run here
57    /// so the caller can leave the lock before running application code.
58    fn publish(&self) -> Vec<Observer> {
59        self.observers
60            .iter()
61            .map(|(_, observer)| Arc::clone(observer))
62            .collect()
63    }
64
65    fn remove_observer(&mut self, id: u64) {
66        self.observers.retain(|(existing, _)| *existing != id);
67    }
68}
69
70fn registry() -> &'static Mutex<Registry> {
71    static REGISTRY: OnceLock<Mutex<Registry>> = OnceLock::new();
72    REGISTRY.get_or_init(|| Mutex::new(Registry::new()))
73}
74
75static NEXT_ID: AtomicU64 = AtomicU64::new(1);
76
77/// Keeps an observer registered until it is dropped.
78pub struct MemoryPressureObserver {
79    id: u64,
80}
81
82impl Drop for MemoryPressureObserver {
83    fn drop(&mut self) {
84        if let Ok(mut registry) = registry().lock() {
85            registry.remove_observer(self.id);
86        }
87    }
88}
89
90/// Registers `observer` for pressure reports.
91///
92/// Applications collect the stream from [`rememberMemoryPressure`] instead of
93/// calling this.
94pub fn observe_memory_pressure(
95    observer: impl Fn(MemoryPressure) + Send + Sync + 'static,
96) -> MemoryPressureObserver {
97    let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
98    if let Ok(mut registry) = registry().lock() {
99        registry.observe(id, Arc::new(observer));
100    }
101    MemoryPressureObserver { id }
102}
103
104/// Publishes a pressure report. Callable from any thread; the framework moves
105/// each report onto the UI thread before a composition sees it.
106pub fn publish_memory_pressure(pressure: MemoryPressure) {
107    let observers = {
108        let Ok(registry) = registry().lock() else {
109            return;
110        };
111        registry.publish()
112    };
113    for observer in observers {
114        observer(pressure);
115    }
116}
117
118/// Collects pressure reports for as long as this call stays in the
119/// composition.
120///
121/// ```rust,no_run
122/// use cranpose_macros::composable;
123/// use cranpose_services::rememberMemoryPressure;
124///
125/// #[composable]
126/// fn Caches() {
127///     let pressure = rememberMemoryPressure();
128///     cranpose_core::CollectEvents(pressure, (), |report| {
129///         log::info!("memory pressure: {report:?}");
130///     });
131/// }
132/// ```
133#[allow(non_snake_case)]
134#[track_caller]
135pub fn rememberMemoryPressure() -> EventStream<MemoryPressure> {
136    rememberEventStream((), |sender| {
137        observe_memory_pressure(move |pressure| sender.send(pressure))
138    })
139}
140
141#[cfg(test)]
142mod tests {
143    use super::*;
144
145    fn recording_observer() -> (Observer, Arc<Mutex<Vec<MemoryPressure>>>) {
146        let seen = Arc::new(Mutex::new(Vec::new()));
147        let recorder = Arc::clone(&seen);
148        let observer: Observer = Arc::new(move |pressure| {
149            recorder
150                .lock()
151                .unwrap_or_else(|error| error.into_inner())
152                .push(pressure)
153        });
154        (observer, seen)
155    }
156
157    /// These exercise `Registry` directly rather than the process-global one,
158    /// for the same reason the incoming-share tests do: the global is shared
159    /// by every test in this binary and the harness runs them on parallel
160    /// threads.
161    #[test]
162    fn trim_levels_map_to_the_three_kinds() {
163        assert_eq!(
164            MemoryPressure::from_android_trim_level(20),
165            MemoryPressure::UiHidden
166        );
167        for level in [5, 10, 40] {
168            assert_eq!(
169                MemoryPressure::from_android_trim_level(level),
170                MemoryPressure::Low
171            );
172        }
173        for level in [15, 60, 80] {
174            assert_eq!(
175                MemoryPressure::from_android_trim_level(level),
176                MemoryPressure::Critical
177            );
178        }
179    }
180
181    #[test]
182    fn publish_reaches_every_observer() {
183        let mut registry = Registry::new();
184        let (first, first_seen) = recording_observer();
185        let (second, second_seen) = recording_observer();
186        registry.observe(1, first);
187        registry.observe(2, second);
188
189        for observer in registry.publish() {
190            observer(MemoryPressure::Critical);
191        }
192
193        for seen in [first_seen, second_seen] {
194            assert_eq!(
195                seen.lock().unwrap_or_else(|e| e.into_inner()).as_slice(),
196                [MemoryPressure::Critical]
197            );
198        }
199    }
200
201    #[test]
202    fn a_removed_observer_stops_seeing_reports() {
203        let mut registry = Registry::new();
204        let (observer, seen) = recording_observer();
205        registry.observe(7, observer);
206
207        for observer in registry.publish() {
208            observer(MemoryPressure::Low);
209        }
210        registry.remove_observer(7);
211        for observer in registry.publish() {
212            observer(MemoryPressure::Low);
213        }
214
215        assert_eq!(
216            seen.lock().unwrap_or_else(|e| e.into_inner()).as_slice(),
217            [MemoryPressure::Low]
218        );
219    }
220
221    #[test]
222    fn a_report_with_no_observers_goes_nowhere() {
223        let registry = Registry::new();
224        assert!(
225            registry.publish().is_empty(),
226            "pressure describes a moment; nothing is kept for late observers"
227        );
228    }
229}