Skip to main content

glycin_core/
pool.rs

1static DEFAULT_POOL: LazyLock<Arc<Pool>> = LazyLock::new(|| Arc::new(Pool::default()));
2
3use std::collections::BTreeMap;
4use std::path::PathBuf;
5use std::sync::atomic::Ordering;
6use std::sync::{Arc, LazyLock, Mutex};
7use std::time::{Duration, Instant};
8
9use gio::glib;
10use gio::prelude::*;
11
12#[cfg(feature = "external")]
13use crate::DBusProxy;
14use crate::config::{ConfigEntry, ConfigEntryHash};
15use crate::error::ErrorKind;
16use crate::util::{AsyncMutex, TimerHandle, spawn_timeout};
17use crate::{Error, SandboxMechanism, config, dbus};
18
19#[derive(Debug)]
20pub struct PooledProcess<P: DBusProxy> {
21    last_use: Mutex<Instant>,
22    _timeout: Arc<Mutex<Option<TimerHandle>>>,
23    process: Arc<dbus::RemoteProcess<P>>,
24    useage_tracker: Mutex<std::sync::Weak<UsageTracker>>,
25}
26
27#[derive(Debug)]
28pub struct UsageTracker {
29    pool: Arc<Pool>,
30    timeout: Arc<Mutex<Option<TimerHandle>>>,
31}
32
33impl UsageTracker {
34    pub fn new(pool: Arc<Pool>, timeout: Arc<Mutex<Option<TimerHandle>>>) -> Self {
35        Self { pool, timeout }
36    }
37}
38
39impl Drop for UsageTracker {
40    fn drop(&mut self) {
41        tracing::trace!("One process occupation dropped");
42        let pool = self.pool.clone();
43
44        *self.timeout.lock().unwrap() = Some(spawn_timeout(
45            self.pool.config.loader_retention_time,
46            async {
47                pool.clean_loaders().await;
48            },
49        ));
50    }
51}
52
53impl<P: DBusProxy> PooledProcess<P> {
54    pub fn use_(&self) -> Arc<dbus::RemoteProcess<P>> {
55        tracing::trace!("Using pooled process");
56        *self.last_use.lock().unwrap() = Instant::now();
57        self.process.clone()
58    }
59
60    pub fn n_users(&self) -> usize {
61        self.useage_tracker.lock().unwrap().strong_count()
62    }
63}
64
65/// Configuration and set of processor processes
66///
67/// Pools store a configuration based on which processor processes are spawned.
68/// Pool can retain spawned processed for later use based on the configuration.
69#[derive(Debug, Default)]
70pub struct Pool {
71    loaders: AsyncMutex<
72        BTreeMap<config::ConfigEntryHash, Vec<Arc<PooledProcess<dbus::LoaderProxy<'static>>>>>,
73    >,
74    editors: AsyncMutex<
75        BTreeMap<config::ConfigEntryHash, Vec<Arc<PooledProcess<dbus::EditorProxy<'static>>>>>,
76    >,
77    config: PoolConfig,
78}
79
80/// [Pool](Pool) configuration
81#[derive(Debug)]
82pub struct PoolConfig {
83    loader_retention_time: Duration,
84    max_parallel_operations: usize,
85}
86
87impl Default for PoolConfig {
88    fn default() -> Self {
89        Self {
90            loader_retention_time: Duration::from_secs(30),
91            max_parallel_operations: usize::MAX,
92        }
93    }
94}
95
96impl PoolConfig {
97    /// Default pool configuration. See below for default values.
98    pub fn new() -> Self {
99        Self::default()
100    }
101
102    /// Maximum of operations one process will be tasked with. The default value
103    /// is [`usize::MAX`].
104    pub fn max_parallel_operations(mut self, max_parallel_operations: usize) -> Self {
105        if max_parallel_operations == 0 {
106            self.max_parallel_operations = usize::MAX;
107        } else {
108            self.max_parallel_operations = max_parallel_operations;
109        }
110        self
111    }
112
113    /// Time after last use after which a processor process will be terminated.
114    /// The default value is 30 seconds.
115    pub fn retention_time(mut self, retention_time: Duration) -> Self {
116        self.loader_retention_time = retention_time;
117        self
118    }
119}
120
121impl Pool {
122    pub fn new(config: PoolConfig) -> Arc<Self> {
123        Arc::new(Self {
124            config,
125            ..Default::default()
126        })
127    }
128
129    pub fn global() -> Arc<Self> {
130        DEFAULT_POOL.clone()
131    }
132
133    pub(crate) async fn get_loader(
134        self: Arc<Self>,
135        loader_config: config::LoaderConfig,
136        sandbox_mechanism: SandboxMechanism,
137        base_dir: Option<PathBuf>,
138        cancellable: &gio::Cancellable,
139    ) -> Result<
140        (
141            Arc<PooledProcess<dbus::LoaderProxy<'static>>>,
142            Arc<UsageTracker>,
143        ),
144        Error,
145    > {
146        let pooled_loaders = &self.loaders;
147
148        let pp = self
149            .clone()
150            .get_process(
151                pooled_loaders,
152                ConfigEntry::Loader(loader_config.clone()),
153                sandbox_mechanism,
154                base_dir,
155                cancellable,
156            )
157            .await?;
158
159        Ok(pp)
160    }
161
162    /// Spawns loader if not available yet
163    pub(crate) async fn get_editor(
164        self: Arc<Self>,
165        editor_config: config::EditorConfig,
166        sandbox_mechanism: SandboxMechanism,
167        base_dir: Option<PathBuf>,
168        cancellable: &gio::Cancellable,
169    ) -> Result<
170        (
171            Arc<PooledProcess<dbus::EditorProxy<'static>>>,
172            Arc<UsageTracker>,
173        ),
174        Error,
175    > {
176        let pooled_editors = &self.editors;
177
178        let pp = self
179            .clone()
180            .get_process(
181                pooled_editors,
182                ConfigEntry::Editor(editor_config.clone()),
183                sandbox_mechanism,
184                base_dir,
185                cancellable,
186            )
187            .await?;
188
189        Ok(pp)
190    }
191
192    /// Spawns process if not available yet
193    pub(crate) async fn get_process<P: DBusProxy>(
194        self: Arc<Self>,
195        pooled_processes: &AsyncMutex<BTreeMap<ConfigEntryHash, Vec<Arc<PooledProcess<P>>>>>,
196        config: config::ConfigEntry,
197        sandbox_mechanism: SandboxMechanism,
198        base_dir: Option<PathBuf>,
199        cancellable: &gio::Cancellable,
200    ) -> Result<(Arc<PooledProcess<P>>, Arc<UsageTracker>), Error> {
201        let config_hash = config.hash_value(base_dir.clone(), sandbox_mechanism);
202        let mut pooled_processes = pooled_processes.lock().await;
203        let pooled_processes = pooled_processes.entry(config_hash).or_default();
204
205        for process in pooled_processes.iter() {
206            if process.process.process_disconnected.load(Ordering::Relaxed) {
207                tracing::debug!("Existing loader/editor in pool is disconnected. Trying next.");
208            } else if process.n_users() >= self.config.max_parallel_operations {
209                tracing::debug!(
210                    "Existing loader/editor in pool is at 'max_parallel_operations'. Trying next."
211                );
212            } else {
213                tracing::debug!("Using existing loader from pool.");
214                let mut current_usage_tracker = process.useage_tracker.lock().unwrap();
215                let usage_tracker = current_usage_tracker.upgrade().unwrap_or_else(|| {
216                    Arc::new(UsageTracker::new(self.clone(), process._timeout.clone()))
217                });
218                *current_usage_tracker = Arc::downgrade(&usage_tracker);
219                return Ok((process.clone(), usage_tracker));
220            }
221        }
222
223        tracing::debug!("No existing loader/editor in pool. Spawning new one.");
224
225        let process_cancellable = gio::Cancellable::new();
226        let Some(process_cancellable_tie) = cancellable.connect_cancelled(glib::clone!(
227            #[weak]
228            process_cancellable,
229            move |_| process_cancellable.cancel()
230        )) else {
231            return Err(ErrorKind::Canceled(None).err());
232        };
233
234        let process = Arc::new(
235            dbus::RemoteProcess::new(
236                config.clone(),
237                sandbox_mechanism,
238                base_dir,
239                &process_cancellable,
240            )
241            .await?,
242        );
243
244        cancellable.disconnect_cancelled(process_cancellable_tie);
245
246        let _timeout = Arc::new(Mutex::new(None));
247
248        let usage_tracker = Arc::new(UsageTracker::new(self.clone(), _timeout.clone()));
249
250        let pp = Arc::new(PooledProcess {
251            last_use: Mutex::new(Instant::now()),
252            _timeout,
253            process: process.clone(),
254            useage_tracker: Mutex::new(Arc::downgrade(&usage_tracker)),
255        });
256
257        pooled_processes.push(pp.clone());
258
259        Ok((pp, usage_tracker))
260    }
261
262    pub(crate) async fn clean_loaders(self: Arc<Self>) {
263        tracing::debug!("Cleaning up loaders");
264        let mut loader_map = self.loaders.lock().await;
265
266        for (cfg, loaders) in loader_map.iter_mut() {
267            loaders.retain(|loader| {
268                let n_users = loader.n_users();
269                let idle = loader.last_use.lock().unwrap().elapsed();
270                let drop = n_users == 0 && idle > self.config.loader_retention_time;
271
272                tracing::debug!(
273                    "Loader {:?}: drop {drop} users {n_users} (max {}), idle {idle:?} (max {:?})",
274                    cfg.exec(),
275                    self.config.max_parallel_operations,
276                    self.config.loader_retention_time
277                );
278
279                if drop {
280                    tracing::debug!(
281                        "Dropping loader {:?} {}",
282                        cfg.exec(),
283                        Arc::strong_count(&loader.process)
284                    )
285                }
286                !drop
287            });
288        }
289    }
290}