1use serde::{Deserialize, Serialize};
12use std::path::PathBuf;
13
14#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
16#[serde(rename_all = "snake_case")]
17pub enum Phase {
18 Running,
20 Done,
22 Error,
24 Stopped,
26}
27
28#[derive(Debug, Clone, Serialize, Deserialize)]
30pub struct SessionState {
31 #[serde(default, skip_serializing_if = "Option::is_none")]
32 pub handoff_reason: Option<String>,
33 #[serde(default, skip_serializing_if = "Option::is_none")]
34 pub workflow: Option<WorkflowProgress>,
35 #[serde(default)]
36 pub version: u32,
37 #[serde(default)]
38 pub call_id: String,
39 #[serde(default)]
40 pub seq: u64,
41 #[serde(default)]
42 pub started_ms: u64,
43 #[serde(default)]
44 pub finished_ms: Option<u64>,
45 #[serde(default)]
46 pub action: Option<crate::activity::ActionState>,
47 #[serde(default)]
48 pub recent_actions: Vec<crate::activity::ActionState>,
49 #[serde(default)]
50 pub error_code: Option<String>,
51 pub pid: u32,
52 pub command: String,
53 pub app: Option<String>,
55 pub step: Option<String>,
57 pub phase: Phase,
58 pub ts_ms: u64,
60}
61
62impl SessionState {
63 pub fn now(
64 pid: u32,
65 command: &str,
66 app: Option<&str>,
67 step: Option<String>,
68 phase: Phase,
69 ) -> Self {
70 Self {
71 handoff_reason: None,
72 workflow: None,
73 version: 1,
74 call_id: String::new(),
75 seq: 0,
76 started_ms: unix_ms(),
77 finished_ms: None,
78 action: None,
79 recent_actions: Vec::new(),
80 error_code: None,
81 pid,
82 command: command.to_string(),
83 app: app.map(str::to_string),
84 step,
85 phase,
86 ts_ms: unix_ms(),
87 }
88 }
89}
90
91#[derive(Debug, Clone, Serialize, Deserialize)]
92pub struct WorkflowProgress {
93 pub run_id: String,
94 pub title: String,
95 pub status: String,
96 #[serde(default, skip_serializing_if = "Option::is_none")]
97 pub reason: Option<String>,
98 #[serde(default, skip_serializing_if = "Option::is_none")]
99 pub step_id: Option<String>,
100 #[serde(default, skip_serializing_if = "Option::is_none")]
101 pub error_code: Option<crate::ErrorCode>,
102 #[serde(default, skip_serializing_if = "Vec::is_empty")]
103 pub actions: Vec<String>,
104}
105
106pub fn unix_ms() -> u64 {
107 std::time::SystemTime::now()
108 .duration_since(std::time::UNIX_EPOCH)
109 .map(|d| d.as_millis() as u64)
110 .unwrap_or(0)
111}
112
113#[derive(Debug, Clone)]
115pub struct SignalPaths {
116 pub dir: PathBuf,
117}
118
119impl Default for SignalPaths {
120 fn default() -> Self {
121 let dir = std::env::var_os("LOCALAPPDATA")
122 .map(PathBuf::from)
123 .unwrap_or_else(std::env::temp_dir)
124 .join("actl");
125 Self { dir }
126 }
127}
128
129impl SignalPaths {
130 pub fn at(dir: impl Into<PathBuf>) -> Self {
131 Self { dir: dir.into() }
132 }
133
134 pub fn state_file(&self) -> PathBuf {
135 self.dir.join("state.json")
136 }
137
138 pub fn stop_file(&self) -> PathBuf {
139 self.dir.join("stop-requested")
140 }
141
142 pub fn write_state(&self, st: &SessionState) {
144 if let Err(error) = self.try_write_state(st) {
145 crate::log_health::failure(
146 &self.dir,
147 &st.call_id,
148 st.workflow.as_ref().map(|w| w.run_id.as_str()),
149 "state_write",
150 0,
151 &error,
152 );
153 eprintln!("[actl-state] INTERNAL: state publication failed: {error}");
154 }
155 }
156 fn try_write_state(&self, st: &SessionState) -> std::io::Result<()> {
157 std::fs::create_dir_all(&self.dir)?;
158 if st.version >= 2 && !st.call_id.is_empty() {
159 let dir = self.dir.join("calls");
160 std::fs::create_dir_all(&dir)?;
161 let file = dir.join(format!("{}.json", st.call_id));
163 Self::atomic_write(&file, st)?;
164 }
165 Self::atomic_write(&self.state_file(), st)
166 }
167
168 fn atomic_write(file: &std::path::Path, st: &SessionState) -> std::io::Result<()> {
169 let tmp = file.with_extension(format!("{}.{}.tmp", st.pid, st.call_id));
170 let result = (|| {
171 std::fs::write(&tmp, serde_json::to_vec(st)?)?;
172 std::fs::rename(&tmp, file)
173 })();
174 let _ = std::fs::remove_file(tmp);
175 result
176 }
177
178 pub fn read_calls(&self) -> Vec<SessionState> {
180 let now = unix_ms();
181 let mut states = Vec::new();
182 if let Ok(entries) = std::fs::read_dir(self.dir.join("calls")) {
183 for entry in entries.flatten() {
184 let path = entry.path();
185 if path.extension().is_none_or(|e| e != "json") {
186 continue;
187 }
188 let Some(st) = std::fs::read(&path)
189 .ok()
190 .and_then(|b| serde_json::from_slice::<SessionState>(&b).ok())
191 else {
192 continue;
193 };
194 if now.saturating_sub(st.ts_ms) > 60_000 {
195 let _ = std::fs::remove_file(path);
196 } else {
197 states.push(st);
198 }
199 }
200 }
201 if states.is_empty()
202 && let Some(st) = self.read_state()
203 {
204 states.push(st);
205 }
206 states.sort_by_key(|s| s.ts_ms);
207 let terminal_count = states.iter().filter(|s| s.phase != Phase::Running).count();
208 let mut excess = terminal_count.saturating_sub(128);
209 states.retain(|s| {
210 if excess > 0 && s.phase != Phase::Running && !s.call_id.is_empty() {
211 if s.call_id
213 .bytes()
214 .all(|c| c.is_ascii_alphanumeric() || c == b'-')
215 {
216 let _ = std::fs::remove_file(
217 self.dir.join("calls").join(format!("{}.json", s.call_id)),
218 );
219 }
220 excess -= 1;
221 false
222 } else {
223 true
224 }
225 });
226 states
227 }
228
229 pub fn read_state(&self) -> Option<SessionState> {
230 let text = std::fs::read_to_string(self.state_file()).ok()?;
231 serde_json::from_str(&text).ok()
232 }
233}
234
235#[cfg(test)]
236mod tests {
237 use super::*;
238
239 fn tmp() -> SignalPaths {
240 static N: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0);
241 let n = N.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
242 let dir = std::env::temp_dir().join(format!("actl-state-test-{}-{n}", std::process::id()));
243 let _ = std::fs::remove_dir_all(&dir);
244 SignalPaths::at(&dir)
245 }
246
247 #[test]
248 fn state_roundtrip_and_fields() {
249 let p = tmp();
250 let st = SessionState::now(
251 42,
252 "press",
253 Some("记事本"),
254 Some("3/8".into()),
255 Phase::Running,
256 );
257 p.write_state(&st);
258 let back = p.read_state().expect("read back");
259 assert_eq!(back.command, "press");
260 assert_eq!(back.app.as_deref(), Some("记事本"));
261 assert_eq!(back.step.as_deref(), Some("3/8"));
262 assert_eq!(back.phase, Phase::Running);
263 assert!(back.ts_ms > 0);
264 let _ = std::fs::remove_dir_all(&p.dir);
265 }
266
267 #[test]
268 fn stop_flag_is_sticky_until_cleared() {
269 let p = tmp();
270 assert!(!p.stop_requested());
271 p.request_stop().unwrap();
272 assert!(p.stop_requested());
273 p.recover_execution().unwrap();
274 assert!(!p.stop_requested());
275 let _ = std::fs::remove_dir_all(&p.dir);
276 }
277
278 #[test]
279 fn phase_serializes_snake_case() {
280 assert_eq!(
281 serde_json::to_string(&Phase::Stopped).unwrap(),
282 r#""stopped""#
283 );
284 }
285}