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 #[cfg(unix)]
322 pub fn remove_shell_dir(&mut self, shell_pid: u32) -> bool {
323 let removed = self.shell_dirs.remove(&shell_pid.to_string()).is_some();
324 if removed {
325 self.mark_dirty();
326 }
327 removed
328 }
329
330 pub fn set_project_session(
333 &mut self,
334 pid: u32,
335 dir: PathBuf,
336 session: ProjectSession,
337 ) -> Option<ProjectSession> {
338 let inner = self.project_sessions.entry(pid.to_string()).or_default();
339 let old = inner.insert(dir, session);
340 self.mark_dirty();
341 old
342 }
343
344 pub fn remove_project_session(&mut self, pid: u32, dir: &Path) -> Option<ProjectSession> {
348 let pid_str = pid.to_string();
349 if let std::collections::btree_map::Entry::Occupied(mut entry) =
350 self.project_sessions.entry(pid_str)
351 {
352 let removed = entry.get_mut().remove(dir);
353 if removed.is_some() {
354 if entry.get().is_empty() {
355 entry.remove();
356 }
357 self.mark_dirty();
358 }
359 removed
360 } else {
361 None
362 }
363 }
364
365 #[cfg(any(unix, test))]
367 pub fn get_project_session(&self, pid: u32, dir: &Path) -> Option<&ProjectSession> {
368 self.project_sessions
369 .get(&pid.to_string())
370 .and_then(|inner| inner.get(dir))
371 }
372
373 pub fn iter_project_sessions(&self) -> Vec<(&str, &PathBuf, &ProjectSession)> {
376 let mut out: Vec<(&str, &PathBuf, &ProjectSession)> = Vec::new();
377 for (pid_str, inner) in &self.project_sessions {
378 for (dir, session) in inner {
379 out.push((pid_str.as_str(), dir, session));
380 }
381 }
382 out
383 }
384
385 pub fn retain_daemons<F>(&mut self, mut f: F)
388 where
389 F: FnMut(&DaemonId, &Daemon) -> bool,
390 {
391 let before = self.daemons.len();
392 self.daemons.retain(|id, daemon| f(id, daemon));
393 if self.daemons.len() != before {
394 self.mark_dirty();
395 }
396 }
397
398 pub fn write(&self) -> Result<()> {
403 let canonical_path = normalized_lock_path(&self.path);
404 let _lock = xx::fslock::get(&canonical_path, false)?;
405 let raw = toml::to_string(self).map_err(|e| FileError::SerializeError {
406 path: self.path.clone(),
407 source: e,
408 })?;
409 if self
410 .last_content
411 .lock()
412 .unwrap()
413 .as_ref()
414 .is_some_and(|last| last == &raw)
415 {
416 self.dirty.store(false, Ordering::Relaxed);
418 return Ok(());
419 }
420 Self::write_raw(&self.path, &raw)?;
421 *self.last_content.lock().unwrap() = Some(raw);
422 self.dirty.store(false, Ordering::Relaxed);
423 Ok(())
424 }
425
426 fn write_unlocked(&self) -> Result<()> {
429 let raw = toml::to_string(self).map_err(|e| FileError::SerializeError {
430 path: self.path.clone(),
431 source: e,
432 })?;
433 Self::write_raw(&self.path, &raw)?;
434 *self.last_content.lock().unwrap() = Some(raw);
435 self.dirty.store(false, Ordering::Relaxed);
436 Ok(())
437 }
438
439 pub(crate) fn write_raw(path: &Path, raw: &str) -> Result<()> {
443 if let Some(parent) = path.parent() {
444 std::fs::create_dir_all(parent).map_err(|e| FileError::WriteError {
445 path: parent.to_path_buf(),
446 details: Some(format!("failed to create state file directory: {e}")),
447 })?;
448 }
449 let temp_path = path.with_extension("toml.tmp");
450 xx::file::write(&temp_path, raw).map_err(|e| FileError::WriteError {
451 path: temp_path.clone(),
452 details: Some(e.to_string()),
453 })?;
454 std::fs::rename(&temp_path, path).map_err(|e| FileError::WriteError {
455 path: path.to_path_buf(),
456 details: Some(format!("failed to rename temp file: {e}")),
457 })?;
458 Ok(())
459 }
460}
461
462fn normalized_lock_path(path: &Path) -> PathBuf {
463 if let Ok(canonical) = path.canonicalize() {
464 return canonical;
465 }
466
467 if let Some(parent) = path.parent()
468 && let Ok(canonical_parent) = parent.canonicalize()
469 && let Some(file_name) = path.file_name()
470 {
471 return canonical_parent.join(file_name);
472 }
473
474 path.to_path_buf()
475}
476
477#[cfg(test)]
478mod tests {
479 use super::*;
480 use crate::daemon_status::DaemonStatus;
481
482 #[test]
483 fn test_state_file_toml_roundtrip_stopped() {
484 let mut state = StateFile::new(PathBuf::from("/tmp/test.toml"));
485 let daemon_id = DaemonId::new("project", "test");
486 state.daemons.insert(
487 daemon_id.clone(),
488 Daemon {
489 id: daemon_id,
490 status: DaemonStatus::Stopped,
491 last_exit_success: Some(true),
492 user: Some("postgres".to_string()),
493 ..Daemon::default()
494 },
495 );
496
497 let toml_str = toml::to_string(&state).unwrap();
498 println!("Serialized TOML:\n{toml_str}");
499
500 let parsed: StateFile = toml::from_str(&toml_str).expect("Failed to parse TOML");
501 println!("Parsed: {parsed:?}");
502
503 assert!(
504 parsed
505 .daemons
506 .contains_key(&DaemonId::new("project", "test"))
507 );
508 let daemon = parsed
509 .daemons
510 .get(&DaemonId::new("project", "test"))
511 .unwrap();
512 assert_eq!(daemon.user.as_deref(), Some("postgres"));
513 }
514
515 #[test]
516 fn test_looks_like_old_format_bare_names() {
517 let old = r#"
518[daemons.api]
519id = "api"
520autostop = false
521retry = 0
522retry_count = 0
523status = "stopped"
524"#;
525 assert!(StateFile::looks_like_old_format(old));
526 }
527
528 #[test]
529 fn test_looks_like_old_format_new_format() {
530 let new = r#"
531 disabled = []
532
533 [daemons."legacy/api"]
534 id = "legacy/api"
535autostop = false
536retry = 0
537retry_count = 0
538status = "stopped"
539"#;
540 assert!(!StateFile::looks_like_old_format(new));
541 }
542
543 #[test]
544 fn test_looks_like_old_format_empty() {
545 assert!(!StateFile::looks_like_old_format(""));
546 assert!(!StateFile::looks_like_old_format("[shell_dirs]"));
547 }
548
549 #[test]
550 fn test_migrate_old_format_basic() {
551 let old = r#"
552[daemons.api]
553id = "api"
554autostop = false
555retry = 0
556retry_count = 0
557status = "stopped"
558
559[daemons.worker]
560id = "worker"
561autostop = false
562retry = 0
563retry_count = 0
564status = "stopped"
565last_exit_success = true
566"#;
567 let migrated = StateFile::migrate_old_format(old).expect("migration should succeed");
568 assert!(
569 migrated
570 .daemons
571 .contains_key(&DaemonId::new("legacy", "api")),
572 "api should be migrated to legacy/api"
573 );
574 assert!(
575 migrated
576 .daemons
577 .contains_key(&DaemonId::new("legacy", "worker")),
578 "worker should be migrated to legacy/worker"
579 );
580 assert_eq!(migrated.daemons.len(), 2);
581 }
582
583 #[test]
584 fn test_migrate_old_format_preserves_disabled() {
585 let old = r#"
586disabled = ["api", "worker"]
587
588[daemons.api]
589id = "api"
590autostop = false
591retry = 0
592retry_count = 0
593status = "stopped"
594"#;
595 let migrated = StateFile::migrate_old_format(old).expect("migration should succeed");
596 assert!(
597 migrated.disabled.contains(&DaemonId::new("legacy", "api")),
598 "disabled 'api' should become 'legacy/api'"
599 );
600 assert!(
601 migrated
602 .disabled
603 .contains(&DaemonId::new("legacy", "worker")),
604 "disabled 'worker' should become 'legacy/worker'"
605 );
606 }
607
608 #[test]
609 fn test_migrate_old_format_already_qualified_unchanged() {
610 let mixed = r#"
612[daemons.bare]
613id = "bare"
614autostop = false
615retry = 0
616retry_count = 0
617status = "stopped"
618"#;
619 let migrated = StateFile::migrate_old_format(mixed).expect("migration should succeed");
620 assert!(
622 migrated
623 .daemons
624 .contains_key(&DaemonId::new("legacy", "bare")),
625 "bare key should become legacy/bare"
626 );
627 assert_eq!(migrated.daemons.len(), 1);
629 }
630
631 #[test]
632 fn test_migrate_old_format_does_not_overwrite_existing_qualified_entry() {
633 let mixed = r#"
634[daemons.api]
635id = "api"
636cmd = ["echo", "old"]
637autostop = false
638retry = 0
639retry_count = 0
640status = "stopped"
641
642[daemons."legacy/api"]
643id = "legacy/api"
644cmd = ["echo", "new"]
645autostop = false
646retry = 0
647retry_count = 0
648status = "stopped"
649"#;
650
651 let migrated = StateFile::migrate_old_format(mixed).expect("migration should succeed");
652 let key = DaemonId::new("legacy", "api");
653 let daemon = migrated.daemons.get(&key).expect("legacy/api should exist");
654
655 let cmd = daemon.cmd.as_ref().expect("cmd should exist");
656 assert_eq!(cmd, &vec!["echo".to_string(), "new".to_string()]);
657
658 let preserved = DaemonId::new("legacy", "api-legacy");
660 let preserved_daemon = migrated
661 .daemons
662 .get(&preserved)
663 .expect("colliding bare key should be preserved as legacy/api-legacy");
664 let preserved_cmd = preserved_daemon
665 .cmd
666 .as_ref()
667 .expect("preserved cmd should exist");
668 assert_eq!(preserved_cmd, &vec!["echo".to_string(), "old".to_string()]);
669 assert_eq!(migrated.daemons.len(), 2);
670 }
671
672 #[test]
673 fn test_project_sessions_nested_map_roundtrip() {
674 let mut state = StateFile::new(PathBuf::from("/tmp/test.toml"));
675 state.set_project_session(
676 1234,
677 PathBuf::from("/projects/a"),
678 ProjectSession {
679 liveness_title: Some("sleep".to_string()),
680 },
681 );
682 state.set_project_session(
683 1234,
684 PathBuf::from("/projects/b"),
685 ProjectSession {
686 liveness_title: None,
687 },
688 );
689 state.set_project_session(
690 5678,
691 PathBuf::from("/projects/a"),
692 ProjectSession {
693 liveness_title: Some("code".to_string()),
694 },
695 );
696
697 let toml_str = toml::to_string(&state).unwrap();
698 let parsed: StateFile = toml::from_str(&toml_str).expect("roundtrip parse");
699
700 assert_eq!(parsed.iter_project_sessions().len(), 3);
702 assert!(
703 parsed
704 .get_project_session(1234, &PathBuf::from("/projects/a"))
705 .is_some()
706 );
707 assert!(
708 parsed
709 .get_project_session(1234, &PathBuf::from("/projects/b"))
710 .is_some()
711 );
712 assert!(
713 parsed
714 .get_project_session(5678, &PathBuf::from("/projects/a"))
715 .is_some()
716 );
717
718 let mut state = parsed;
720 state.remove_project_session(5678, &PathBuf::from("/projects/a"));
721 assert!(
722 state
723 .get_project_session(5678, &PathBuf::from("/projects/a"))
724 .is_none()
725 );
726 assert!(!state.project_sessions.contains_key("5678"));
727 assert_eq!(state.iter_project_sessions().len(), 2);
728 }
729}