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

src/watchers.rs (6.8K)

  1 use std::sync::Arc;
  2 use std::sync::atomic::{AtomicU8, Ordering};
  3 use std::sync::mpsc::{channel, Receiver, Sender};
  4 use crate::pages::{Page, audio, bluetooth, browser, default_apps, network, processes, services, system_info, storage, packages, accounts, notifications, timers, power};
  5 
  6 pub struct Watchers {
  7     pub rx_audio: Receiver<audio::AudioState>,
  8     pub rx_network: Receiver<network::NetworkState>,
  9     pub rx_bluetooth: Receiver<bluetooth::BluetoothState>,
 10     pub rx_power: Receiver<power::PowerFacts>,
 11     pub rx_processes: Receiver<processes::ProcessesState>,
 12     pub rx_system: Receiver<system_info::SystemInfo>,
 13     pub rx_storage: Receiver<storage::StorageInfo>,
 14     pub rx_notifications: Receiver<notifications::NotificationsConfig>,
 15     pub rx_browser: Receiver<browser::BrowserConfig>,
 16     pub rx_services: Receiver<Vec<services::ServiceInfo>>,
 17     pub rx_default_apps: Receiver<default_apps::DefaultAppsInfo>,
 18     pub rx_timers: Receiver<Vec<timers::TimerInfo>>,
 19     pub rx_accounts: Receiver<accounts::AccountsSnapshot>,
 20     pub rx_packages: Receiver<packages::PackagesState>,
 21 }
 22 
 23 fn spawn_bg_active<T, F>(
 24     current_page_shared: Arc<AtomicU8>,
 25     target_page_idx: u8,
 26     period_secs: u64,
 27     f: fn() -> F,
 28 ) -> Receiver<T>
 29 where
 30     T: Send + 'static,
 31     F: std::future::Future<Output = T> + Send + 'static,
 32 {
 33     let (tx, rx) = channel::<T>();
 34     tokio::spawn(async move {
 35         let mut last_fetch: Option<std::time::Instant> = None;
 36         loop {
 37             let current_page = current_page_shared.load(Ordering::SeqCst);
 38             if current_page == target_page_idx {
 39                 let should_fetch = match last_fetch {
 40                     None => true,
 41                     Some(t) => t.elapsed() >= std::time::Duration::from_secs(period_secs),
 42                 };
 43                 if should_fetch {
 44                     let val = f().await;
 45                     if tx.send(val).is_err() { break; }
 46                     last_fetch = Some(std::time::Instant::now());
 47                 }
 48             }
 49             tokio::time::sleep(std::time::Duration::from_millis(250)).await;
 50         }
 51     });
 52     rx
 53 }
 54 
 55 pub fn spawn_all(
 56     current_page_shared: Arc<AtomicU8>,
 57 ) -> (
 58     Watchers,
 59     Sender<storage::StorageMessage>,
 60     Receiver<storage::StorageMessage>,
 61     Sender<packages::PackagesMessage>,
 62     Receiver<packages::PackagesMessage>,
 63 ) {
 64     let rx_audio = spawn_bg_active(current_page_shared.clone(), Page::Audio.index() as u8, 3, || audio::fetch_audio_state());
 65     let rx_network = spawn_bg_active(current_page_shared.clone(), Page::Network.index() as u8, 5, || network::fetch_network_state());
 66     let rx_bluetooth = spawn_bg_active(current_page_shared.clone(), Page::Bluetooth.index() as u8, 5, || bluetooth::fetch_bluetooth_page_state());
 67     let rx_power = spawn_bg_active(current_page_shared.clone(), Page::Power.index() as u8, 5, || power::fetch_power_state());
 68 
 69     let rx_system = spawn_bg_active(current_page_shared.clone(), Page::System.index() as u8, 5, || system_info::fetch_system_state());
 70     let rx_processes = spawn_bg_active(current_page_shared.clone(), Page::Processes.index() as u8, 3, || processes::fetch_processes_state());
 71     let rx_storage = spawn_bg_active(current_page_shared.clone(), Page::Storage.index() as u8, 10, || storage::fetch_storage_state());
 72 
 73     let rx_notifications = {
 74         let (tx, rx) = channel::<notifications::NotificationsConfig>();
 75         let current_page_shared = current_page_shared.clone();
 76         tokio::spawn(async move {
 77             let mut last_fetch: Option<std::time::Instant> = None;
 78             loop {
 79                 let current_page = current_page_shared.load(Ordering::SeqCst);
 80                 if current_page == Page::Notifications.index() as u8 {
 81                     let should_fetch = match last_fetch {
 82                         None => true,
 83                         Some(t) => t.elapsed() >= std::time::Duration::from_secs(30),
 84                     };
 85                     if should_fetch {
 86                         let val = tokio::task::spawn_blocking(|| notifications::read_notifications_config()).await;
 87                         if let Ok(val) = val {
 88                             if tx.send(val).is_err() { break; }
 89                         }
 90                         last_fetch = Some(std::time::Instant::now());
 91                     }
 92                 }
 93                 tokio::time::sleep(std::time::Duration::from_millis(250)).await;
 94             }
 95         });
 96         rx
 97     };
 98 
 99     // Config-file poll while the Browser page is open: catches edits made
100     // outside this app (the browser itself, cce-data-editor).
101     let rx_browser = {
102         let (tx, rx) = channel::<browser::BrowserConfig>();
103         let current_page_shared = current_page_shared.clone();
104         tokio::spawn(async move {
105             let mut last_fetch: Option<std::time::Instant> = None;
106             loop {
107                 let current_page = current_page_shared.load(Ordering::SeqCst);
108                 if current_page == Page::Browser.index() as u8 {
109                     let should_fetch = match last_fetch {
110                         None => true,
111                         Some(t) => t.elapsed() >= std::time::Duration::from_secs(5),
112                     };
113                     if should_fetch {
114                         let val = tokio::task::spawn_blocking(browser::read_browser_config).await;
115                         if let Ok(val) = val {
116                             if tx.send(val).is_err() { break; }
117                         }
118                         last_fetch = Some(std::time::Instant::now());
119                     }
120                 }
121                 tokio::time::sleep(std::time::Duration::from_millis(250)).await;
122             }
123         });
124         rx
125     };
126 
127     let rx_services = spawn_bg_active(current_page_shared.clone(), Page::Services.index() as u8, 3, || services::fetch_services());
128     let rx_default_apps = spawn_bg_active(current_page_shared.clone(), Page::DefaultApps.index() as u8, 10, || default_apps::fetch_default_apps());
129     let rx_timers = spawn_bg_active(current_page_shared.clone(), Page::Timers.index() as u8, 5, || timers::fetch_timers());
130     let rx_accounts = spawn_bg_active(current_page_shared.clone(), Page::Accounts.index() as u8, 3, || accounts::fetch_accounts());
131 
132     let (tx_backup, rx_backup) = channel();
133     let rx_packages = spawn_bg_active(current_page_shared.clone(), Page::Packages.index() as u8, 30, || packages::fetch_packages_state());
134     let (tx_update, rx_update) = channel();
135 
136     (
137         Watchers {
138             rx_audio,
139             rx_network,
140             rx_bluetooth,
141             rx_power,
142             rx_processes,
143             rx_system,
144             rx_storage,
145             rx_notifications,
146             rx_browser,
147             rx_services,
148             rx_default_apps,
149             rx_timers,
150             rx_accounts,
151             rx_packages,
152         },
153         tx_backup,
154         rx_backup,
155         tx_update,
156         rx_update,
157     )
158 }