Skip to main content

actl_core/
state.rs

1//! 显示信号协议(docs/11 §10 / docs/13 §4.1):actl 侧状态源。
2//!
3//! 原语四件套:①`state.json` 命令级状态快照(开始/结束/batch 步间写入);
4//! ②`stop-requested` 当前停止记录 + `stop-generation` 调用取消代次;
5//! ③输入占用锁(actl-uia 的命名互斥体,"正在操控"的权威信号,本模块只约定
6//! 不实现);④时间戳心跳(消费端判失联)。
7//! 消费端(actl-signal 等)全部拉取、零订阅——"事件当触发器,拉取当真相源"
8//! (docs/13 的事故价目表:常驻事件订阅路线不可取)。
9//! 路径可注入(`SignalPaths::at`)以便测试;`default()` 用 %LOCALAPPDATA%。
10
11use serde::{Deserialize, Serialize};
12use std::path::PathBuf;
13
14/// 命令生命周期相位。
15#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
16#[serde(rename_all = "snake_case")]
17pub enum Phase {
18    /// 命令开始(batch:第 i/N 步)
19    Running,
20    /// 命令成功结束
21    Done,
22    /// 命令失败结束
23    Error,
24    /// 用户停止中止
25    Stopped,
26}
27
28/// 一次命令执行的状态快照(单写者:当前 actl 进程;多读者)。
29#[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    /// 目标窗口摘要(--app 原样;无则 None)
54    pub app: Option<String>,
55    /// batch 步序 "i/N";单命令 None
56    pub step: Option<String>,
57    pub phase: Phase,
58    /// unix 毫秒时间戳(心跳)
59    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/// 信号目录/文件的解析。默认 %LOCALAPPDATA%\actl\。
114#[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    /// 状态快照落盘;诊断故障不改变业务结果,但写入独立健康记录。
143    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            // Each invocation owns its file: concurrent commands cannot overwrite it.
162            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    /// Recent call snapshots, bounded retention; active calls are never evicted.
179    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                // File names are enumerated/generated locally; never trust a serialized path.
212                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}