git.lucas.co / cce-system-interface
system settings
git clone https://git.lucas.co/cce-system-interface.git

src/watchers.rs (6.6K)

  1 use std::sync::Arc;
  2 use std::sync::mpsc::{channel, Receiver, Sender};
  3 use crate::pages::{Page, audio, bluetooth, browser, default_apps, network, processes, services, system_info, storage, packages, accounts, notifications, timers, power};
  4 
  5 /// The page on screen, published by the UI (`watch::Sender::send_replace`)
  6 /// to every page worker.
  7 pub type PageWatch = tokio::sync::watch::Receiver<u8>;
  8 
  9 /// Wakes the UI loop after a worker sends: the runner cannot see the std
 10 /// channels below, and until 2026-10-05 the app made it poll them every
 11 /// 250 ms (`idle_poll_interval`) whatever page was open.
 12 pub type Wake = Arc<dyn Fn() + Send + Sync>;
 13 
 14 pub struct Watchers {
 15     pub rx_audio: Receiver<audio::AudioState>,
 16     pub rx_network: Receiver<network::NetworkState>,
 17     pub rx_bluetooth: Receiver<bluetooth::BluetoothState>,
 18     pub rx_power: Receiver<power::PowerFacts>,
 19     pub rx_processes: Receiver<processes::ProcessesState>,
 20     pub rx_system: Receiver<system_info::SystemInfo>,
 21     pub rx_storage: Receiver<storage::StorageInfo>,
 22     pub rx_notifications: Receiver<notifications::NotificationsConfig>,
 23     pub rx_browser: Receiver<browser::BrowserConfig>,
 24     pub rx_services: Receiver<Vec<services::ServiceInfo>>,
 25     pub rx_default_apps: Receiver<default_apps::DefaultAppsInfo>,
 26     pub rx_timers: Receiver<Vec<timers::TimerInfo>>,
 27     pub rx_accounts: Receiver<accounts::AccountsSnapshot>,
 28     pub rx_packages: Receiver<packages::PackagesState>,
 29 }
 30 
 31 /// One page's worker: while `target_page` is on screen, fetch every
 32 /// `period_secs` (and at once on arriving, if the last fetch is that old),
 33 /// send the result, wake the UI. While another page is up it sleeps until the
 34 /// page changes — until 2026-10-05 each of the fourteen workers woke every
 35 /// 250 ms to look, whatever was open. `None` from `f` sends nothing.
 36 fn spawn_page_worker<T, F, Fut>(
 37     mut page: PageWatch,
 38     target_page: Page,
 39     period_secs: u64,
 40     wake: Wake,
 41     f: F,
 42 ) -> Receiver<T>
 43 where
 44     T: Send + 'static,
 45     F: Fn() -> Fut + Send + 'static,
 46     Fut: std::future::Future<Output = Option<T>> + Send + 'static,
 47 {
 48     let target = target_page.index() as u8;
 49     let period = std::time::Duration::from_secs(period_secs);
 50     let (tx, rx) = channel::<T>();
 51     tokio::spawn(async move {
 52         let mut last_fetch: Option<std::time::Instant> = None;
 53         loop {
 54             if *page.borrow_and_update() != target {
 55                 if page.changed().await.is_err() {
 56                     break;
 57                 }
 58                 continue;
 59             }
 60             if last_fetch.is_none_or(|t| t.elapsed() >= period) {
 61                 if let Some(val) = f().await {
 62                     if tx.send(val).is_err() {
 63                         break;
 64                     }
 65                     wake();
 66                 }
 67                 last_fetch = Some(std::time::Instant::now());
 68             }
 69             let due = period.saturating_sub(last_fetch.map_or(period, |t| t.elapsed()));
 70             tokio::select! {
 71                 _ = tokio::time::sleep(due) => {}
 72                 changed = page.changed() => {
 73                     if changed.is_err() {
 74                         break;
 75                     }
 76                 }
 77             }
 78         }
 79     });
 80     rx
 81 }
 82 
 83 /// [`spawn_page_worker`] over an async fetch that always answers.
 84 fn spawn_bg_active<T, F>(
 85     page: PageWatch,
 86     target_page: Page,
 87     period_secs: u64,
 88     wake: Wake,
 89     f: fn() -> F,
 90 ) -> Receiver<T>
 91 where
 92     T: Send + 'static,
 93     F: std::future::Future<Output = T> + Send + 'static,
 94 {
 95     spawn_page_worker(page, target_page, period_secs, wake, move || {
 96         let fut = f();
 97         async move { Some(fut.await) }
 98     })
 99 }
100 
101 /// [`spawn_page_worker`] over a blocking read, run off the runtime's workers.
102 fn spawn_bg_blocking<T>(
103     page: PageWatch,
104     target_page: Page,
105     period_secs: u64,
106     wake: Wake,
107     f: fn() -> T,
108 ) -> Receiver<T>
109 where
110     T: Send + 'static,
111 {
112     spawn_page_worker(page, target_page, period_secs, wake, move || async move {
113         tokio::task::spawn_blocking(f).await.ok()
114     })
115 }
116 
117 pub fn spawn_all(
118     page: PageWatch,
119     wake: Wake,
120 ) -> (
121     Watchers,
122     Sender<storage::StorageMessage>,
123     Receiver<storage::StorageMessage>,
124     Sender<packages::PackagesMessage>,
125     Receiver<packages::PackagesMessage>,
126 ) {
127     let w = || (page.clone(), wake.clone());
128     let (p, k) = w(); let rx_audio = spawn_bg_active(p, Page::Audio, 3, k, || audio::fetch_audio_state());
129     let (p, k) = w(); let rx_network = spawn_bg_active(p, Page::Network, 5, k, || network::fetch_network_state());
130     let (p, k) = w(); let rx_bluetooth = spawn_bg_active(p, Page::Bluetooth, 5, k, || bluetooth::fetch_bluetooth_page_state());
131     let (p, k) = w(); let rx_power = spawn_bg_active(p, Page::Power, 5, k, || power::fetch_power_state());
132 
133     let (p, k) = w(); let rx_system = spawn_bg_active(p, Page::System, 5, k, || system_info::fetch_system_state());
134     let (p, k) = w(); let rx_processes = spawn_bg_active(p, Page::Processes, 3, k, || processes::fetch_processes_state());
135     let (p, k) = w(); let rx_storage = spawn_bg_active(p, Page::Storage, 10, k, || storage::fetch_storage_state());
136 
137     let (p, k) = w(); let rx_notifications = spawn_bg_blocking(p, Page::Notifications, 30, k, notifications::read_notifications_config);
138     // Config-file poll while the Browser page is open: catches edits made
139     // outside this app (the browser itself, cce-data-editor).
140     let (p, k) = w(); let rx_browser = spawn_bg_blocking(p, Page::Browser, 5, k, browser::read_browser_config);
141 
142     let (p, k) = w(); let rx_services = spawn_bg_active(p, Page::Services, 3, k, || services::fetch_services());
143     let (p, k) = w(); let rx_default_apps = spawn_bg_active(p, Page::DefaultApps, 10, k, || default_apps::fetch_default_apps());
144     let (p, k) = w(); let rx_timers = spawn_bg_active(p, Page::Timers, 5, k, || timers::fetch_timers());
145     let (p, k) = w(); let rx_accounts = spawn_bg_active(p, Page::Accounts, 3, k, || accounts::fetch_accounts());
146 
147     let (tx_backup, rx_backup) = channel();
148     let (p, k) = w(); let rx_packages = spawn_bg_active(p, Page::Packages, 30, k, || packages::fetch_packages_state());
149     let (tx_update, rx_update) = channel();
150 
151     (
152         Watchers {
153             rx_audio,
154             rx_network,
155             rx_bluetooth,
156             rx_power,
157             rx_processes,
158             rx_system,
159             rx_storage,
160             rx_notifications,
161             rx_browser,
162             rx_services,
163             rx_default_apps,
164             rx_timers,
165             rx_accounts,
166             rx_packages,
167         },
168         tx_backup,
169         rx_backup,
170         tx_update,
171         rx_update,
172     )
173 }