par_term/profile/dynamic/
manager.rs1use std::collections::HashMap;
8use std::sync::Arc;
9use std::time::{Duration, SystemTime};
10use tokio::sync::mpsc;
11
12use par_term_config::{ConflictResolution, DynamicProfileSource};
13
14use crate::profile::dynamic::fetch::fetch_profiles;
15
16#[derive(Debug, Clone)]
18pub struct DynamicProfileUpdate {
19 pub url: String,
21 pub profiles: Vec<par_term_config::Profile>,
23 pub conflict_resolution: ConflictResolution,
25 pub error: Option<String>,
27}
28
29#[derive(Debug, Clone)]
31pub struct SourceStatus {
32 pub url: String,
34 pub enabled: bool,
36 pub last_fetch: Option<SystemTime>,
38 pub last_error: Option<String>,
40 pub profile_count: usize,
42 pub fetching: bool,
44}
45
46pub struct DynamicProfileManager {
51 pub update_rx: mpsc::UnboundedReceiver<DynamicProfileUpdate>,
53 update_tx: mpsc::UnboundedSender<DynamicProfileUpdate>,
55 pub statuses: HashMap<String, SourceStatus>,
57 task_handles: Vec<tokio::task::JoinHandle<()>>,
59}
60
61impl DynamicProfileManager {
62 pub fn new() -> Self {
64 let (update_tx, update_rx) = mpsc::unbounded_channel();
65 Self {
66 update_rx,
67 update_tx,
68 statuses: HashMap::new(),
69 task_handles: Vec::new(),
70 }
71 }
72
73 pub fn start(
80 &mut self,
81 sources: &[DynamicProfileSource],
82 runtime: &Arc<tokio::runtime::Runtime>,
83 ) {
84 use crate::profile::dynamic::cache::read_cache;
85
86 self.stop();
88
89 for source in sources {
90 if !source.enabled || source.url.is_empty() {
91 continue;
92 }
93
94 self.statuses.insert(
96 source.url.clone(),
97 SourceStatus {
98 url: source.url.clone(),
99 enabled: source.enabled,
100 last_fetch: None,
101 last_error: None,
102 profile_count: 0,
103 fetching: false,
104 },
105 );
106
107 if let Ok((profiles, meta)) = read_cache(&source.url) {
109 let update = DynamicProfileUpdate {
110 url: source.url.clone(),
111 profiles,
112 conflict_resolution: source.conflict_resolution.clone(),
113 error: None,
114 };
115 let _ = self.update_tx.send(update);
116
117 if let Some(status) = self.statuses.get_mut(&source.url) {
118 status.last_fetch = Some(meta.last_fetched);
119 status.profile_count = meta.profile_count;
120 }
121 }
122
123 let tx = self.update_tx.clone();
125 let source_clone = source.clone();
126 let url_for_log = source.url.clone();
127 let handle = runtime.spawn(async move {
128 let src = source_clone.clone();
130 let conflict = source_clone.conflict_resolution.clone();
131 match tokio::task::spawn_blocking(move || fetch_profiles(&src)).await {
132 Ok(result) => {
133 if tx
134 .send(DynamicProfileUpdate {
135 url: result.url.clone(),
136 profiles: result.profiles,
137 conflict_resolution: conflict,
138 error: result.error,
139 })
140 .is_err()
141 {
142 return; }
144 }
145 Err(e) => {
146 log::error!(
147 "Dynamic profile fetch task panicked for {}: {}",
148 url_for_log,
149 e
150 );
151 }
152 }
153
154 let mut interval =
156 tokio::time::interval(Duration::from_secs(source_clone.refresh_interval_secs));
157 interval.tick().await; loop {
159 interval.tick().await;
160 let src = source_clone.clone();
161 let source_clone2 = source_clone.clone();
162 let tx_clone = tx.clone();
163 match tokio::task::spawn_blocking(move || fetch_profiles(&src)).await {
164 Ok(result) => {
165 if tx_clone
166 .send(DynamicProfileUpdate {
167 url: result.url.clone(),
168 profiles: result.profiles,
169 conflict_resolution: source_clone2.conflict_resolution.clone(),
170 error: result.error,
171 })
172 .is_err()
173 {
174 break; }
176 }
177 Err(e) => {
178 log::error!(
179 "Dynamic profile fetch task panicked for {}: {}",
180 url_for_log,
181 e
182 );
183 }
184 }
185 }
186 });
187
188 self.task_handles.push(handle);
189
190 if let Some(status) = self.statuses.get_mut(&source.url) {
191 status.fetching = true;
192 }
193 }
194 }
195
196 pub fn stop(&mut self) {
198 for handle in self.task_handles.drain(..) {
199 handle.abort();
200 }
201 }
202
203 pub fn refresh_all(
205 &mut self,
206 sources: &[DynamicProfileSource],
207 runtime: &Arc<tokio::runtime::Runtime>,
208 ) {
209 for source in sources {
210 if !source.enabled || source.url.is_empty() {
211 continue;
212 }
213 self.refresh_source(source, runtime);
214 }
215 }
216
217 pub fn refresh_source(
219 &mut self,
220 source: &DynamicProfileSource,
221 runtime: &Arc<tokio::runtime::Runtime>,
222 ) {
223 let tx = self.update_tx.clone();
224 let source_clone = source.clone();
225 let url_for_log = source.url.clone();
226 runtime.spawn(async move {
227 let conflict = source_clone.conflict_resolution.clone();
228 match tokio::task::spawn_blocking(move || fetch_profiles(&source_clone)).await {
229 Ok(result) => {
230 let _ = tx.send(DynamicProfileUpdate {
231 url: result.url.clone(),
232 profiles: result.profiles,
233 conflict_resolution: conflict,
234 error: result.error,
235 });
236 }
237 Err(e) => {
238 log::error!(
239 "Dynamic profile fetch task panicked for {}: {}",
240 url_for_log,
241 e
242 );
243 }
244 }
245 });
246
247 if let Some(status) = self.statuses.get_mut(&source.url) {
248 status.fetching = true;
249 }
250 }
251
252 pub fn try_recv(&mut self) -> Option<DynamicProfileUpdate> {
254 self.update_rx.try_recv().ok()
255 }
256
257 pub fn update_status(&mut self, update: &DynamicProfileUpdate) {
259 if let Some(status) = self.statuses.get_mut(&update.url) {
260 status.fetching = false;
261 status.last_error = update.error.clone();
262 if update.error.is_none() {
263 status.last_fetch = Some(SystemTime::now());
264 status.profile_count = update.profiles.len();
265 }
266 }
267 }
268}
269
270impl Default for DynamicProfileManager {
271 fn default() -> Self {
272 Self::new()
273 }
274}
275
276impl Drop for DynamicProfileManager {
277 fn drop(&mut self) {
278 self.stop();
279 }
280}