pitchfork_cli/supervisor/
state.rs1use super::Supervisor;
6use crate::Result;
7use crate::daemon::Daemon;
8use crate::daemon::RunOptions;
9use crate::daemon_id::DaemonId;
10use crate::daemon_status::DaemonStatus;
11use crate::pitchfork_toml::CpuLimit;
12use crate::pitchfork_toml::CronRetrigger;
13use crate::pitchfork_toml::MemoryLimit;
14use crate::pitchfork_toml::PitchforkToml;
15use crate::pitchfork_toml::PortConfig;
16use crate::pitchfork_toml::ReadyCmd;
17use crate::pitchfork_toml::ReadyHttp;
18use crate::pitchfork_toml::ReadyOutput;
19use crate::pitchfork_toml::ReadyPort;
20use crate::pitchfork_toml::Retry;
21use crate::pitchfork_toml::StopConfig;
22use crate::pitchfork_toml::WatchMode;
23use crate::procs::PROCS;
24use indexmap::IndexMap;
25use std::collections::{HashMap, HashSet};
26use std::path::PathBuf;
27
28#[derive(Debug, Default)]
32pub(crate) struct UpsertDaemonOpts {
33 pub id: DaemonId,
34 pub pid: Option<u32>,
35 pub status: DaemonStatus,
36 pub shell_pid: Option<u32>,
37 pub dir: Option<PathBuf>,
38 pub cmd: Option<Vec<String>>,
39 pub run: Option<String>,
40 pub autostop: bool,
41 pub cron_schedule: Option<String>,
42 pub cron_retrigger: Option<CronRetrigger>,
43 pub cron_immediate: Option<bool>,
44 pub last_exit_success: Option<bool>,
45 pub retry: Option<Retry>,
46 pub retry_count: Option<u32>,
47 pub ready_delay: Option<u64>,
48 pub ready_output: Option<ReadyOutput>,
49 pub ready_http: Option<ReadyHttp>,
50 pub ready_port: Option<ReadyPort>,
51 pub ready_cmd: Option<ReadyCmd>,
52 pub port: Option<PortConfig>,
54 pub resolved_port: Vec<u16>,
56 pub active_port: Option<u16>,
58 pub slug: Option<String>,
60 pub proxy: Option<bool>,
62 pub depends: Option<Vec<DaemonId>>,
63 pub env: Option<IndexMap<String, String>>,
64 pub watch: Option<Vec<String>>,
65 pub watch_mode: Option<WatchMode>,
66 pub watch_base_dir: Option<PathBuf>,
67 pub mise: Option<bool>,
68 pub user: Option<String>,
70 pub memory_limit: Option<MemoryLimit>,
72 pub cpu_limit: Option<CpuLimit>,
74 pub stop_signal: Option<StopConfig>,
76 pub archive_hook: Option<String>,
78 pub log_format: Option<String>,
80 pub pty: Option<bool>,
82 pub config_registered: bool,
84}
85
86#[derive(Debug)]
98pub(crate) struct UpsertDaemonOptsBuilder {
99 pub opts: UpsertDaemonOpts,
100}
101
102impl UpsertDaemonOpts {
103 pub fn builder(id: DaemonId) -> UpsertDaemonOptsBuilder {
105 UpsertDaemonOptsBuilder {
106 opts: UpsertDaemonOpts {
107 id,
108 ..Default::default()
109 },
110 }
111 }
112
113 pub(crate) fn from_run_options(
117 opts: &RunOptions,
118 status: DaemonStatus,
119 ) -> UpsertDaemonOptsBuilder {
120 UpsertDaemonOpts::builder(opts.id.clone()).set(|o| {
121 o.status = status;
122 o.shell_pid = opts.shell_pid;
123 o.dir = Some(opts.dir.0.clone());
124 o.cmd = Some(opts.cmd.clone());
125 o.run = opts.run.clone();
126 o.autostop = opts.autostop;
127 o.cron_schedule = opts.cron_schedule.clone();
128 o.cron_retrigger = opts.cron_retrigger;
129 o.cron_immediate = opts.cron_immediate;
130 o.retry = Some(opts.retry);
131 o.retry_count = Some(opts.retry_count);
132 o.ready_delay = opts.ready_delay;
133 o.ready_output = opts.ready_output.clone();
134 o.ready_http = opts.ready_http.clone();
135 o.ready_port = opts.ready_port.clone();
136 o.ready_cmd = opts.ready_cmd.clone();
137 o.port = opts.port.clone();
138 o.depends = Some(opts.depends.clone());
139 o.env = opts.env.clone();
140 o.watch = Some(opts.watch.clone());
141 o.watch_mode = Some(opts.watch_mode);
142 o.watch_base_dir = opts.watch_base_dir.clone();
143 o.mise = opts.mise;
144 o.user = opts.user.clone();
145 o.memory_limit = opts.memory_limit;
146 o.cpu_limit = opts.cpu_limit;
147 o.stop_signal = opts.stop_signal;
148 o.pty = opts.pty;
149 o.archive_hook = opts.archive_hook.clone();
150 o.log_format = opts.log_format.clone();
151 })
152 }
153}
154
155impl UpsertDaemonOptsBuilder {
156 pub fn set<F: FnOnce(&mut UpsertDaemonOpts)>(mut self, f: F) -> Self {
158 f(&mut self.opts);
159 self
160 }
161
162 pub fn build(self) -> UpsertDaemonOpts {
164 self.opts
165 }
166}
167
168impl Supervisor {
169 pub(crate) async fn upsert_daemon(&self, opts: UpsertDaemonOpts) -> Result<Daemon> {
171 info!(
172 "upserting daemon: {} pid: {} status: {}",
173 opts.id,
174 opts.pid.unwrap_or(0),
175 opts.status
176 );
177 let mut state_file = self.state_file.lock().await;
178 let existing = state_file.daemons.get(&opts.id);
179 let daemon = Daemon {
180 id: opts.id.clone(),
181 title: opts.pid.and_then(|pid| {
187 PROCS.title(pid).or_else(|| {
188 existing
189 .filter(|d| d.pid == Some(pid))
190 .and_then(|d| d.title.clone())
191 })
192 }),
193 start_time: opts.pid.and_then(|pid| {
194 PROCS.start_time(pid).or_else(|| {
195 existing
196 .filter(|d| d.pid == Some(pid))
197 .and_then(|d| d.start_time)
198 })
199 }),
200 boot_time: opts.pid.map(|_| PROCS.boot_time()),
204 pid: opts.pid,
205 status: opts.status,
206 shell_pid: opts.shell_pid,
207 autostop: opts.autostop || existing.is_some_and(|d| d.autostop),
208 dir: opts.dir.or(existing.and_then(|d| d.dir.clone())),
209 cmd: opts.cmd.or(existing.and_then(|d| d.cmd.clone())),
210 run: opts.run.or(existing.and_then(|d| d.run.clone())),
211 cron_schedule: opts
212 .cron_schedule
213 .or(existing.and_then(|d| d.cron_schedule.clone())),
214 cron_retrigger: opts
215 .cron_retrigger
216 .or(existing.and_then(|d| d.cron_retrigger)),
217 cron_immediate: opts
218 .cron_immediate
219 .or(existing.and_then(|d| d.cron_immediate)),
220 last_cron_triggered: existing.and_then(|d| d.last_cron_triggered),
221 last_exit_success: opts
222 .last_exit_success
223 .or(existing.and_then(|d| d.last_exit_success)),
224 retry: opts
225 .retry
226 .unwrap_or_else(|| existing.map(|d| d.retry).unwrap_or_default()),
227 retry_count: opts
228 .retry_count
229 .unwrap_or(existing.map(|d| d.retry_count).unwrap_or(0)),
230 ready_delay: opts.ready_delay.or(existing.and_then(|d| d.ready_delay)),
231 ready_output: opts
232 .ready_output
233 .or(existing.and_then(|d| d.ready_output.clone())),
234 ready_http: opts
235 .ready_http
236 .or(existing.and_then(|d| d.ready_http.clone())),
237 ready_port: opts
238 .ready_port
239 .or(existing.and_then(|d| d.ready_port.clone())),
240 ready_cmd: opts
241 .ready_cmd
242 .or(existing.and_then(|d| d.ready_cmd.clone())),
243 port: opts.port.or_else(|| existing.and_then(|d| d.port.clone())),
244 resolved_port: if opts.resolved_port.is_empty() {
245 existing
246 .map(|d| d.resolved_port.clone())
247 .unwrap_or_default()
248 } else {
249 opts.resolved_port
250 },
251 depends: opts
252 .depends
253 .unwrap_or_else(|| existing.map(|d| d.depends.clone()).unwrap_or_default()),
254 env: opts.env.or(existing.and_then(|d| d.env.clone())),
255 watch: opts
256 .watch
257 .unwrap_or_else(|| existing.map(|d| d.watch.clone()).unwrap_or_default()),
258 watch_mode: opts
259 .watch_mode
260 .unwrap_or_else(|| existing.map(|d| d.watch_mode).unwrap_or_default()),
261 watch_base_dir: opts
262 .watch_base_dir
263 .or(existing.and_then(|d| d.watch_base_dir.clone())),
264 mise: opts.mise.or(existing.and_then(|d| d.mise)),
265 user: opts.user.or(existing.and_then(|d| d.user.clone())),
266 proxy: opts.proxy.or(existing.and_then(|d| d.proxy)),
267 active_port: opts.active_port,
273 slug: opts.slug.or(existing.and_then(|d| d.slug.clone())),
274 memory_limit: opts.memory_limit.or(existing.and_then(|d| d.memory_limit)),
275 cpu_limit: opts.cpu_limit.or(existing.and_then(|d| d.cpu_limit)),
276 stop_signal: opts.stop_signal.or(existing.and_then(|d| d.stop_signal)),
277 archive_hook: opts
278 .archive_hook
279 .or(existing.and_then(|d| d.archive_hook.clone())),
280 log_format: opts
281 .log_format
282 .or(existing.and_then(|d| d.log_format.clone())),
283 pty: opts.pty.or(existing.and_then(|d| d.pty)),
284 config_registered: opts.config_registered,
285 };
286 state_file.insert_daemon(&opts.id, daemon.clone());
287 Ok(daemon)
288 }
289
290 pub async fn enable(&self, id: &DaemonId) -> Result<bool> {
292 info!("enabling daemon: {id}");
293 let config = PitchforkToml::all_merged_all_namespaces()?;
294 let mut state_file = self.state_file.lock().await;
295 let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
296 if !exists {
297 return Err(miette::miette!("daemon '{}' not found", id));
298 }
299 let result = state_file.enable_daemon(id);
300 Ok(result)
301 }
302
303 pub async fn disable(&self, id: &DaemonId) -> Result<bool> {
305 info!("disabling daemon: {id}");
306 let config = PitchforkToml::all_merged_all_namespaces()?;
307 let mut state_file = self.state_file.lock().await;
308 let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
309 if !exists {
310 return Err(miette::miette!("daemon '{}' not found", id));
311 }
312 let result = state_file.disable_daemon(id);
313 Ok(result)
314 }
315
316 pub(crate) async fn get_daemon(&self, id: &DaemonId) -> Option<Daemon> {
318 self.state_file.lock().await.daemons.get(id).cloned()
319 }
320
321 pub(crate) async fn active_daemons(&self) -> Vec<Daemon> {
323 let pitchfork_id = DaemonId::pitchfork();
324 self.state_file
325 .lock()
326 .await
327 .daemons
328 .values()
329 .filter(|d| d.pid.is_some() && d.id != pitchfork_id)
330 .cloned()
331 .collect()
332 }
333
334 pub(crate) async fn remove_daemon(&self, id: &DaemonId) -> Result<()> {
336 let mut state_file = self.state_file.lock().await;
337 state_file.remove_daemon(id);
338 Ok(())
339 }
340
341 pub(crate) async fn set_shell_dir(&self, shell_pid: u32, dir: PathBuf) -> Result<()> {
343 let mut state_file = self.state_file.lock().await;
344 state_file.set_shell_dir(shell_pid, dir);
345 Ok(())
346 }
347
348 pub(crate) async fn get_shell_dir(&self, shell_pid: u32) -> Option<PathBuf> {
350 self.state_file
351 .lock()
352 .await
353 .shell_dirs
354 .get(&shell_pid.to_string())
355 .cloned()
356 }
357
358 pub(crate) async fn remove_shell_pid(&self, shell_pid: u32) -> Result<()> {
360 let mut state_file = self.state_file.lock().await;
361 state_file.remove_shell_dir(shell_pid);
362 Ok(())
363 }
364
365 pub(crate) async fn get_dirs_with_shell_pids(&self) -> HashMap<PathBuf, Vec<u32>> {
367 self.state_file.lock().await.shell_dirs.iter().fold(
368 HashMap::new(),
369 |mut acc, (pid, dir)| {
370 if let Ok(pid) = pid.parse() {
371 acc.entry(dir.clone()).or_default().push(pid);
372 }
373 acc
374 },
375 )
376 }
377
378 pub(crate) async fn get_notifications(&self) -> Vec<(log::LevelFilter, String)> {
380 self.pending_notifications.lock().await.drain(..).collect()
381 }
382
383 pub(crate) async fn clean(&self) -> Result<()> {
385 let mut state_file = self.state_file.lock().await;
386 state_file.retain_daemons(|_id, d| d.pid.is_some());
387 Ok(())
388 }
389
390 pub(crate) async fn get_active_directories(&self) -> Vec<PathBuf> {
394 let state = self.state_file.lock().await;
395 let mut dirs: HashSet<PathBuf> = state.shell_dirs.values().cloned().collect();
396 for (_, dir, _) in state.iter_project_sessions() {
397 dirs.insert(dir.clone());
398 }
399 dirs.into_iter().collect()
400 }
401
402 pub(crate) async fn get_liveness_sessions(&self) -> Vec<(u32, PathBuf, Option<String>)> {
406 self.state_file
407 .lock()
408 .await
409 .iter_project_sessions()
410 .into_iter()
411 .filter_map(|(pid_str, dir, session)| {
412 pid_str
413 .parse()
414 .ok()
415 .map(|pid| (pid, dir.clone(), session.liveness_title.clone()))
416 })
417 .collect()
418 }
419
420 pub(crate) async fn get_project_sessions_info(&self) -> Vec<crate::ipc::ProjectSessionInfo> {
424 let sessions: Vec<(u32, PathBuf, Option<String>)> = self
425 .state_file
426 .lock()
427 .await
428 .iter_project_sessions()
429 .into_iter()
430 .filter_map(|(pid_str, dir, session)| {
431 pid_str
432 .parse()
433 .ok()
434 .map(|pid| (pid, dir.clone(), session.liveness_title.clone()))
435 })
436 .collect();
437 let pids: Vec<u32> = sessions.iter().map(|(pid, _, _)| *pid).collect();
438 if !pids.is_empty() {
439 PROCS.refresh_pids(&pids);
440 }
441 sessions
442 .into_iter()
443 .map(
444 |(pid, directory, liveness_title)| crate::ipc::ProjectSessionInfo {
445 pid,
446 directory,
447 liveness_title,
448 alive: PROCS.is_running(pid),
449 current_title: PROCS.title(pid),
450 },
451 )
452 .collect()
453 }
454
455 pub(crate) async fn enter_project_session(
459 &self,
460 pid: u32,
461 dir: PathBuf,
462 ) -> Result<Option<crate::state_file::ProjectSession>> {
463 if pid == 0 {
464 return Err(miette::miette!("invalid host PID 0"));
465 }
466 if pid > i32::MAX as u32 {
467 return Err(miette::miette!("host PID {pid} exceeds i32::MAX"));
468 }
469 PROCS.refresh_pids(&[pid]);
470 let liveness_title = PROCS.title(pid);
471 #[cfg(unix)]
477 if !PROCS.is_running(pid) {
478 return Err(miette::miette!("host PID {pid} is not running"));
479 }
480 let mut state_file = self.state_file.lock().await;
481 let previous = state_file.set_project_session(
482 pid,
483 dir,
484 crate::state_file::ProjectSession { liveness_title },
485 );
486 Ok(previous)
487 }
488
489 pub(crate) async fn leave_project_session(
493 &self,
494 pid: u32,
495 dir: &std::path::Path,
496 ) -> Result<Option<PathBuf>> {
497 let mut state_file = self.state_file.lock().await;
498 if state_file.remove_project_session(pid, dir).is_some() {
499 Ok(Some(dir.to_path_buf()))
500 } else {
501 Ok(None)
502 }
503 }
504}