1use crate::daemon::Daemon;
2use crate::daemon_id::DaemonId;
3use crate::error::FileError;
4use crate::{Result, env};
5use once_cell::sync::Lazy;
6use std::collections::{BTreeMap, BTreeSet};
7use std::fmt::Debug;
8use std::path::{Path, PathBuf};
9use std::sync::Mutex;
10use std::sync::atomic::{AtomicBool, Ordering};
11
12#[derive(Debug, serde::Serialize, serde::Deserialize)]
13pub struct StateFile {
14 #[serde(default)]
15 pub daemons: BTreeMap<DaemonId, Daemon>,
16 #[serde(default)]
17 pub disabled: BTreeSet<DaemonId>,
18 #[serde(default)]
19 pub shell_dirs: BTreeMap<String, PathBuf>,
20 #[serde(default)]
24 pub project_sessions: BTreeMap<String, BTreeMap<PathBuf, ProjectSession>>,
25 #[serde(skip)]
26 pub(crate) path: PathBuf,
27 #[serde(skip)]
28 pub(crate) dirty: AtomicBool,
29 #[serde(skip)]
34 pub(crate) last_content: Mutex<Option<String>>,
35}
36
37#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)]
41pub struct ProjectSession {
42 #[serde(skip_serializing_if = "Option::is_none", default)]
43 pub liveness_title: Option<String>,
44}
45
46impl StateFile {
47 pub fn new(path: PathBuf) -> Self {
48 Self {
49 daemons: Default::default(),
50 disabled: Default::default(),
51 shell_dirs: Default::default(),
52 project_sessions: Default::default(),
53 path,
54 dirty: AtomicBool::new(false),
55 last_content: Mutex::new(None),
56 }
57 }
58
59 pub fn get() -> &'static Self {
60 static STATE_FILE: Lazy<StateFile> = Lazy::new(|| {
61 let path = &*env::PITCHFORK_STATE_FILE;
62 StateFile::read(path).unwrap_or_else(|e| {
63 error!(
64 "failed to read state file {}: {}. Falling back to in-memory empty state",
65 path.display(),
66 e
67 );
68 StateFile::new(path.to_path_buf())
69 })
70 });
71 &STATE_FILE
72 }
73
74 pub fn read<P: AsRef<Path>>(path: P) -> Result<Self> {
75 let path = path.as_ref();
76 if !path.exists() {
77 return Ok(Self::new(path.to_path_buf()));
78 }
79 let canonical_path = normalized_lock_path(path);
80 let _lock = xx::fslock::get(&canonical_path, false)?;
81 let raw = xx::file::read_to_string(path).unwrap_or_else(|e| {
82 warn!("Error reading state file {path:?}: {e}");
83 String::new()
84 });
85
86 match toml::from_str::<Self>(&raw) {
88 Ok(mut state_file) => {
89 state_file.path = path.to_path_buf();
90 state_file.dirty = AtomicBool::new(false);
91 for (id, daemon) in state_file.daemons.iter_mut() {
92 daemon.id = id.clone();
93 }
94 state_file.last_content = Mutex::new(Some(raw));
97 Ok(state_file)
98 }
99 Err(parse_err) => {
100 if Self::looks_like_old_format(&raw) {
101 debug!(
103 "State file at {} appears to be in old format, attempting silent migration",
104 path.display()
105 );
106 match Self::migrate_old_format(&raw) {
107 Ok(migrated) => {
108 let mut state_file = migrated;
109 state_file.path = path.to_path_buf();
110 if let Err(e) = state_file.write_unlocked() {
112 warn!("State file migration write failed: {e}");
113 }
114 debug!("State file migrated successfully");
115 return Ok(state_file);
116 }
117 Err(e) => {
118 error!(
119 "State file migration failed: {e}. \
120 Raw content preserved at {}. Starting with empty state.",
121 path.display()
122 );
123 return Err(miette::miette!(
124 "Failed to migrate state file {}: {e}",
125 path.display()
126 ));
127 }
128 }
129 }
130 Err(miette::miette!(
132 "Failed to parse state file {}: {parse_err}",
133 path.display()
134 ))
135 }
136 }
137 }
138
139 fn looks_like_old_format(raw: &str) -> bool {
144 use toml::Value;
145 let Ok(Value::Table(doc)) = toml::from_str::<Value>(raw) else {
146 return false;
147 };
148 let Some(Value::Table(daemons)) = doc.get("daemons") else {
149 return false;
150 };
151 !daemons.is_empty() && daemons.keys().any(|k| !k.contains('/'))
153 }
154
155 fn migrate_old_format(raw: &str) -> Result<Self> {
158 use toml::Value;
159
160 const LEGACY_NAMESPACE: &str = "legacy";
161
162 let mut doc: toml::map::Map<String, Value> = toml::from_str(raw)
164 .map_err(|e| miette::miette!("failed to parse old state file: {e}"))?;
165
166 if let Some(Value::Table(daemons)) = doc.get_mut("daemons") {
168 let old_keys: Vec<String> = daemons.keys().cloned().collect();
169 for key in old_keys {
170 if !key.contains('/')
171 && let Some(val) = daemons.remove(&key)
172 {
173 let mut new_key = format!("{LEGACY_NAMESPACE}/{key}");
174 if daemons.contains_key(&new_key) {
176 let base = format!("{key}-legacy");
177 let mut candidate = format!("{LEGACY_NAMESPACE}/{base}");
178 let mut n: u32 = 2;
179 while daemons.contains_key(&candidate) {
180 candidate = format!("{LEGACY_NAMESPACE}/{base}-{n}");
181 n += 1;
182 }
183 warn!(
184 "Legacy daemon key '{}' collides with '{}'; migrating as '{}'",
185 key,
186 format_args!("{LEGACY_NAMESPACE}/{key}"),
187 candidate
188 );
189 new_key = candidate;
190 }
191 let val = if let Value::Table(mut tbl) = val {
193 tbl.insert("id".to_string(), Value::String(new_key.clone()));
194 Value::Table(tbl)
195 } else {
196 val
197 };
198 daemons.insert(new_key, val);
199 }
200 }
201 }
202
203 if let Some(Value::Array(disabled)) = doc.get_mut("disabled") {
205 for entry in disabled.iter_mut() {
206 if let Value::String(s) = entry
207 && !s.contains('/')
208 {
209 *s = format!("{LEGACY_NAMESPACE}/{s}");
210 }
211 }
212 }
213
214 let new_raw =
215 toml::to_string(&Value::Table(doc)).map_err(|e| FileError::SerializeError {
216 path: PathBuf::new(),
217 source: e,
218 })?;
219
220 let mut state_file: Self = toml::from_str(&new_raw)
221 .map_err(|e| miette::miette!("failed to parse migrated state file: {e}"))?;
222 for (id, daemon) in state_file.daemons.iter_mut() {
224 daemon.id = id.clone();
225 }
226 Ok(state_file)
227 }
228
229 fn mark_dirty(&self) {
232 self.dirty.store(true, Ordering::Relaxed);
233 }
234
235 pub fn is_dirty(&self) -> bool {
237 self.dirty.load(Ordering::Relaxed)
238 }
239
240 pub fn insert_daemon(&mut self, id: &DaemonId, daemon: Daemon) {
242 self.daemons.insert(id.clone(), daemon);
243 self.mark_dirty();
244 }
245
246 pub fn remove_daemon(&mut self, id: &DaemonId) {
248 if self.daemons.remove(id).is_some() {
249 self.mark_dirty();
250 }
251 }
252
253 pub fn disable_daemon(&mut self, id: &DaemonId) -> bool {
256 let inserted = self.disabled.insert(id.clone());
257 if inserted {
258 self.mark_dirty();
259 }
260 inserted
261 }
262
263 pub fn enable_daemon(&mut self, id: &DaemonId) -> bool {
266 let removed = self.disabled.remove(id);
267 if removed {
268 self.mark_dirty();
269 }
270 removed
271 }
272
273 pub fn set_active_port(&mut self, id: &DaemonId, port: u16) -> bool {
276 if let Some(d) = self.daemons.get_mut(id) {
277 d.active_port = Some(port);
278 self.mark_dirty();
279 true
280 } else {
281 false
282 }
283 }
284
285 pub fn clear_active_port(&mut self, id: &DaemonId) -> bool {
288 if let Some(d) = self.daemons.get_mut(id) {
289 d.active_port = None;
290 self.mark_dirty();
291 true
292 } else {
293 false
294 }
295 }
296
297 pub fn set_last_cron_triggered(
300 &mut self,
301 id: &DaemonId,
302 time: chrono::DateTime<chrono::Local>,
303 ) -> bool {
304 if let Some(d) = self.daemons.get_mut(id) {
305 d.last_cron_triggered = Some(time);
306 self.mark_dirty();
307 true
308 } else {
309 false
310 }
311 }
312
313 pub fn set_shell_dir(&mut self, shell_pid: u32, dir: PathBuf) {
315 self.shell_dirs.insert(shell_pid.to_string(), dir);
316 self.mark_dirty();
317 }
318
319 pub fn remove_shell_dir(&mut self, shell_pid: u32) -> bool {
322 let removed = self.shell_dirs.remove(&shell_pid.to_string()).is_some();
323 if removed {
324 self.mark_dirty();
325 }
326 removed
327 }
328
329 pub fn set_project_session(
332 &mut self,
333 pid: u32,
334 dir: PathBuf,
335 session: ProjectSession,
336 ) -> Option<ProjectSession> {
337 let inner = self.project_sessions.entry(pid.to_string()).or_default();
338 let old = inner.insert(dir, session);
339 self.mark_dirty();
340 old
341 }
342
343 pub fn remove_project_session(&mut self, pid: u32, dir: &Path) -> Option<ProjectSession> {
347 let pid_str = pid.to_string();
348 if let std::collections::btree_map::Entry::Occupied(mut entry) =
349 self.project_sessions.entry(pid_str)
350 {
351 let removed = entry.get_mut().remove(dir);
352 if removed.is_some() {
353 if entry.get().is_empty() {
354 entry.remove();
355 }
356 self.mark_dirty();
357 }
358 removed
359 } else {
360 None
361 }
362 }
363
364 pub fn get_project_session(&self, pid: u32, dir: &Path) -> Option<&ProjectSession> {
366 self.project_sessions
367 .get(&pid.to_string())
368 .and_then(|inner| inner.get(dir))
369 }
370
371 pub fn iter_project_sessions(&self) -> Vec<(&str, &PathBuf, &ProjectSession)> {
374 let mut out: Vec<(&str, &PathBuf, &ProjectSession)> = Vec::new();
375 for (pid_str, inner) in &self.project_sessions {
376 for (dir, session) in inner {
377 out.push((pid_str.as_str(), dir, session));
378 }
379 }
380 out
381 }
382
383 pub fn retain_daemons<F>(&mut self, mut f: F)
386 where
387 F: FnMut(&DaemonId, &Daemon) -> bool,
388 {
389 let before = self.daemons.len();
390 self.daemons.retain(|id, daemon| f(id, daemon));
391 if self.daemons.len() != before {
392 self.mark_dirty();
393 }
394 }
395
396 pub fn write(&self) -> Result<()> {
401 let canonical_path = normalized_lock_path(&self.path);
402 let _lock = xx::fslock::get(&canonical_path, false)?;
403 let raw = toml::to_string(self).map_err(|e| FileError::SerializeError {
404 path: self.path.clone(),
405 source: e,
406 })?;
407 if self
408 .last_content
409 .lock()
410 .unwrap()
411 .as_ref()
412 .is_some_and(|last| last == &raw)
413 {
414 self.dirty.store(false, Ordering::Relaxed);
416 return Ok(());
417 }
418 Self::write_raw(&self.path, &raw)?;
419 *self.last_content.lock().unwrap() = Some(raw);
420 self.dirty.store(false, Ordering::Relaxed);
421 Ok(())
422 }
423
424 fn write_unlocked(&self) -> Result<()> {
427 let raw = toml::to_string(self).map_err(|e| FileError::SerializeError {
428 path: self.path.clone(),
429 source: e,
430 })?;
431 Self::write_raw(&self.path, &raw)?;
432 *self.last_content.lock().unwrap() = Some(raw);
433 self.dirty.store(false, Ordering::Relaxed);
434 Ok(())
435 }
436
437 pub(crate) fn write_raw(path: &Path, raw: &str) -> Result<()> {
441 if let Some(parent) = path.parent() {
442 std::fs::create_dir_all(parent).map_err(|e| FileError::WriteError {
443 path: parent.to_path_buf(),
444 details: Some(format!("failed to create state file directory: {e}")),
445 })?;
446 }
447 let temp_path = path.with_extension("toml.tmp");
448 xx::file::write(&temp_path, raw).map_err(|e| FileError::WriteError {
449 path: temp_path.clone(),
450 details: Some(e.to_string()),
451 })?;
452 std::fs::rename(&temp_path, path).map_err(|e| FileError::WriteError {
453 path: path.to_path_buf(),
454 details: Some(format!("failed to rename temp file: {e}")),
455 })?;
456 Ok(())
457 }
458}
459
460fn normalized_lock_path(path: &Path) -> PathBuf {
461 if let Ok(canonical) = path.canonicalize() {
462 return canonical;
463 }
464
465 if let Some(parent) = path.parent()
466 && let Ok(canonical_parent) = parent.canonicalize()
467 && let Some(file_name) = path.file_name()
468 {
469 return canonical_parent.join(file_name);
470 }
471
472 path.to_path_buf()
473}
474
475#[cfg(test)]
476mod tests {
477 use super::*;
478 use crate::daemon_status::DaemonStatus;
479
480 #[test]
481 fn test_state_file_toml_roundtrip_stopped() {
482 let mut state = StateFile::new(PathBuf::from("/tmp/test.toml"));
483 let daemon_id = DaemonId::new("project", "test");
484 state.daemons.insert(
485 daemon_id.clone(),
486 Daemon {
487 id: daemon_id,
488 status: DaemonStatus::Stopped,
489 last_exit_success: Some(true),
490 user: Some("postgres".to_string()),
491 ..Daemon::default()
492 },
493 );
494
495 let toml_str = toml::to_string(&state).unwrap();
496 println!("Serialized TOML:\n{toml_str}");
497
498 let parsed: StateFile = toml::from_str(&toml_str).expect("Failed to parse TOML");
499 println!("Parsed: {parsed:?}");
500
501 assert!(
502 parsed
503 .daemons
504 .contains_key(&DaemonId::new("project", "test"))
505 );
506 let daemon = parsed
507 .daemons
508 .get(&DaemonId::new("project", "test"))
509 .unwrap();
510 assert_eq!(daemon.user.as_deref(), Some("postgres"));
511 }
512
513 #[test]
514 fn test_looks_like_old_format_bare_names() {
515 let old = r#"
516[daemons.api]
517id = "api"
518autostop = false
519retry = 0
520retry_count = 0
521status = "stopped"
522"#;
523 assert!(StateFile::looks_like_old_format(old));
524 }
525
526 #[test]
527 fn test_looks_like_old_format_new_format() {
528 let new = r#"
529 disabled = []
530
531 [daemons."legacy/api"]
532 id = "legacy/api"
533autostop = false
534retry = 0
535retry_count = 0
536status = "stopped"
537"#;
538 assert!(!StateFile::looks_like_old_format(new));
539 }
540
541 #[test]
542 fn test_looks_like_old_format_empty() {
543 assert!(!StateFile::looks_like_old_format(""));
544 assert!(!StateFile::looks_like_old_format("[shell_dirs]"));
545 }
546
547 #[test]
548 fn test_migrate_old_format_basic() {
549 let old = r#"
550[daemons.api]
551id = "api"
552autostop = false
553retry = 0
554retry_count = 0
555status = "stopped"
556
557[daemons.worker]
558id = "worker"
559autostop = false
560retry = 0
561retry_count = 0
562status = "stopped"
563last_exit_success = true
564"#;
565 let migrated = StateFile::migrate_old_format(old).expect("migration should succeed");
566 assert!(
567 migrated
568 .daemons
569 .contains_key(&DaemonId::new("legacy", "api")),
570 "api should be migrated to legacy/api"
571 );
572 assert!(
573 migrated
574 .daemons
575 .contains_key(&DaemonId::new("legacy", "worker")),
576 "worker should be migrated to legacy/worker"
577 );
578 assert_eq!(migrated.daemons.len(), 2);
579 }
580
581 #[test]
582 fn test_migrate_old_format_preserves_disabled() {
583 let old = r#"
584disabled = ["api", "worker"]
585
586[daemons.api]
587id = "api"
588autostop = false
589retry = 0
590retry_count = 0
591status = "stopped"
592"#;
593 let migrated = StateFile::migrate_old_format(old).expect("migration should succeed");
594 assert!(
595 migrated.disabled.contains(&DaemonId::new("legacy", "api")),
596 "disabled 'api' should become 'legacy/api'"
597 );
598 assert!(
599 migrated
600 .disabled
601 .contains(&DaemonId::new("legacy", "worker")),
602 "disabled 'worker' should become 'legacy/worker'"
603 );
604 }
605
606 #[test]
607 fn test_migrate_old_format_already_qualified_unchanged() {
608 let mixed = r#"
610[daemons.bare]
611id = "bare"
612autostop = false
613retry = 0
614retry_count = 0
615status = "stopped"
616"#;
617 let migrated = StateFile::migrate_old_format(mixed).expect("migration should succeed");
618 assert!(
620 migrated
621 .daemons
622 .contains_key(&DaemonId::new("legacy", "bare")),
623 "bare key should become legacy/bare"
624 );
625 assert_eq!(migrated.daemons.len(), 1);
627 }
628
629 #[test]
630 fn test_migrate_old_format_does_not_overwrite_existing_qualified_entry() {
631 let mixed = r#"
632[daemons.api]
633id = "api"
634cmd = ["echo", "old"]
635autostop = false
636retry = 0
637retry_count = 0
638status = "stopped"
639
640[daemons."legacy/api"]
641id = "legacy/api"
642cmd = ["echo", "new"]
643autostop = false
644retry = 0
645retry_count = 0
646status = "stopped"
647"#;
648
649 let migrated = StateFile::migrate_old_format(mixed).expect("migration should succeed");
650 let key = DaemonId::new("legacy", "api");
651 let daemon = migrated.daemons.get(&key).expect("legacy/api should exist");
652
653 let cmd = daemon.cmd.as_ref().expect("cmd should exist");
654 assert_eq!(cmd, &vec!["echo".to_string(), "new".to_string()]);
655
656 let preserved = DaemonId::new("legacy", "api-legacy");
658 let preserved_daemon = migrated
659 .daemons
660 .get(&preserved)
661 .expect("colliding bare key should be preserved as legacy/api-legacy");
662 let preserved_cmd = preserved_daemon
663 .cmd
664 .as_ref()
665 .expect("preserved cmd should exist");
666 assert_eq!(preserved_cmd, &vec!["echo".to_string(), "old".to_string()]);
667 assert_eq!(migrated.daemons.len(), 2);
668 }
669
670 #[test]
671 fn test_project_sessions_nested_map_roundtrip() {
672 let mut state = StateFile::new(PathBuf::from("/tmp/test.toml"));
673 state.set_project_session(
674 1234,
675 PathBuf::from("/projects/a"),
676 ProjectSession {
677 liveness_title: Some("sleep".to_string()),
678 },
679 );
680 state.set_project_session(
681 1234,
682 PathBuf::from("/projects/b"),
683 ProjectSession {
684 liveness_title: None,
685 },
686 );
687 state.set_project_session(
688 5678,
689 PathBuf::from("/projects/a"),
690 ProjectSession {
691 liveness_title: Some("code".to_string()),
692 },
693 );
694
695 let toml_str = toml::to_string(&state).unwrap();
696 let parsed: StateFile = toml::from_str(&toml_str).expect("roundtrip parse");
697
698 assert_eq!(parsed.iter_project_sessions().len(), 3);
700 assert!(
701 parsed
702 .get_project_session(1234, &PathBuf::from("/projects/a"))
703 .is_some()
704 );
705 assert!(
706 parsed
707 .get_project_session(1234, &PathBuf::from("/projects/b"))
708 .is_some()
709 );
710 assert!(
711 parsed
712 .get_project_session(5678, &PathBuf::from("/projects/a"))
713 .is_some()
714 );
715
716 let mut state = parsed;
718 state.remove_project_session(5678, &PathBuf::from("/projects/a"));
719 assert!(
720 state
721 .get_project_session(5678, &PathBuf::from("/projects/a"))
722 .is_none()
723 );
724 assert!(!state.project_sessions.contains_key("5678"));
725 assert_eq!(state.iter_project_sessions().len(), 2);
726 }
727}