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#[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#[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 pub fn new() -> Self {
99 Self::default()
100 }
101
102 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 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 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 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}