cranpose_services/
memory_pressure.rs1use std::sync::{
10 Arc, Mutex, OnceLock,
11 atomic::{AtomicU64, Ordering},
12};
13
14use cranpose_core::{EventStream, rememberEventStream};
15
16#[derive(Clone, Copy, Debug, PartialEq, Eq)]
18pub enum MemoryPressure {
19 UiHidden,
22 Low,
24 Critical,
26}
27
28impl MemoryPressure {
29 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 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
77pub 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
90pub 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
104pub 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#[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 #[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}