Skip to main content

actl_core/
stop.rs

1//! Persistent cancellation generations. Recovery never revives an old invocation.
2use crate::{
3    CtlError, ErrorCode,
4    state::{SignalPaths, unix_ms},
5};
6use serde::{Deserialize, Serialize};
7use std::{
8    cell::RefCell,
9    fs,
10    io::{ErrorKind, Write},
11};
12
13#[derive(Debug, Serialize, Deserialize)]
14#[serde(deny_unknown_fields)]
15pub struct StopRecord {
16    pub version: u32,
17    pub generation: String,
18    pub requested_ms: u64,
19    pub source: String,
20    pub pid: u32,
21    pub display_epoch: Option<String>,
22}
23fn invalid(detail: impl ToString) -> CtlError {
24    CtlError::with_evidence(
25        ErrorCode::NotActionable,
26        "execution state requires repair",
27        serde_json::json!({
28            "reason":"execution_state_invalid", "detail":detail.to_string(),
29            "recovery":"execution-recover", "retry_input":false
30        }),
31    )
32}
33fn cancelled() -> CtlError {
34    CtlError::with_evidence(
35        ErrorCode::Aborted,
36        "this invocation was stopped; start a new invocation and verify unfinished work",
37        serde_json::json!({
38            "reason":"execution_generation_changed", "retry_input":false
39        }),
40    )
41}
42impl SignalPaths {
43    pub fn execution_status(&self) -> serde_json::Value {
44        let snapshot = (|| -> Result<_, CtlError> {
45            let _lock = self.control_lock()?;
46            let generation = self.read_generation()?;
47            let stop = self.read_stop(&generation)?;
48            Ok((generation, stop))
49        })();
50        match snapshot {
51            Ok((generation, stop)) => {
52                serde_json::json!({"state":if stop.is_some(){"stopped"}else{"ready"},
53                "generation":generation,"stop":stop,"new_task_requires_consent":true,"old_calls_resume":false})
54            }
55            Err(e) => {
56                serde_json::json!({"state":"repair_required","error":e.evidence,"new_task_requires_consent":true,"old_calls_resume":false})
57            }
58        }
59    }
60    fn generation_file(&self) -> std::path::PathBuf {
61        self.dir.join("stop-generation")
62    }
63    fn control_lock(&self) -> Result<fs::File, CtlError> {
64        fs::create_dir_all(&self.dir).map_err(invalid)?;
65        let file = fs::OpenOptions::new()
66            .create(true)
67            .truncate(false)
68            .read(true)
69            .write(true)
70            .open(self.dir.join("stop-control.lock"))
71            .map_err(invalid)?;
72        file.lock().map_err(invalid)?;
73        Ok(file)
74    }
75    fn read_generation(&self) -> Result<String, CtlError> {
76        match fs::read_to_string(self.generation_file()) {
77            Ok(s)
78                if !s.is_empty()
79                    && s.len() <= 128
80                    && s.bytes().all(|c| c.is_ascii_alphanumeric() || c == b'-') =>
81            {
82                Ok(s)
83            }
84            Ok(_) => Err(invalid("invalid generation")),
85            Err(e) if e.kind() == ErrorKind::NotFound => Ok("initial".into()),
86            Err(e) => Err(invalid(e)),
87        }
88    }
89    fn read_stop(&self, generation: &str) -> Result<Option<StopRecord>, CtlError> {
90        let data = match fs::read(self.stop_file()) {
91            Ok(data) => data,
92            Err(e) if e.kind() == ErrorKind::NotFound => return Ok(None),
93            Err(e) => return Err(invalid(e)),
94        };
95        let r: StopRecord = serde_json::from_slice(&data).map_err(invalid)?;
96        if r.version != 2 || r.generation != generation {
97            return Err(invalid("unsupported stop record or generation mismatch"));
98        }
99        Ok(Some(r))
100    }
101    pub fn execution_generation(&self) -> Result<String, CtlError> {
102        let _lock = self.control_lock()?;
103        let generation = self.read_generation()?;
104        self.read_stop(&generation)?;
105        Ok(generation)
106    }
107    pub fn stop_requested(&self) -> bool {
108        match fs::symlink_metadata(self.stop_file()) {
109            Ok(_) => true,
110            Err(e) => e.kind() != ErrorKind::NotFound,
111        }
112    }
113    fn write_control(&self, path: &std::path::Path, data: &[u8]) -> Result<(), CtlError> {
114        let tmp = path.with_extension(crate::snapshot::new_snapshot_id());
115        let result = (|| {
116            let mut f = fs::OpenOptions::new()
117                .write(true)
118                .create_new(true)
119                .open(&tmp)
120                .map_err(invalid)?;
121            f.write_all(data)
122                .and_then(|_| f.sync_all())
123                .map_err(invalid)?;
124            fs::rename(&tmp, path).map_err(invalid)
125        })();
126        if result.is_err() {
127            let _ = fs::remove_file(tmp);
128        }
129        result
130    }
131    pub fn request_stop(&self) -> Result<(), CtlError> {
132        self.request_stop_from("explicit_request", None)
133    }
134    pub fn request_stop_from(&self, source: &str, epoch: Option<&str>) -> Result<(), CtlError> {
135        let _lock = self.control_lock()?;
136        let generation = crate::snapshot::new_snapshot_id();
137        self.write_control(&self.generation_file(), generation.as_bytes())?;
138        let record = StopRecord {
139            version: 2,
140            generation,
141            requested_ms: unix_ms(),
142            source: source.into(),
143            pid: std::process::id(),
144            display_epoch: epoch.map(str::to_owned),
145        };
146        self.write_control(
147            &self.stop_file(),
148            &serde_json::to_vec(&record).map_err(invalid)?,
149        )
150    }
151    /// Only an explicit user's Start for this generation may acknowledge the stop.
152    pub fn acknowledge_stop(&self, expected: &str) -> Result<(), CtlError> {
153        let _lock = self.control_lock()?;
154        let generation = self.read_generation()?;
155        if generation != expected {
156            return Err(cancelled());
157        }
158        self.read_stop(&generation)?;
159        if self.stop_requested() {
160            let archive = self.dir.join("stop-history").join(&generation);
161            fs::create_dir_all(&archive).map_err(invalid)?;
162            fs::copy(self.stop_file(), archive.join("stop-requested")).map_err(invalid)?;
163        }
164        match fs::remove_file(self.stop_file()) {
165            Ok(()) => Ok(()),
166            Err(e) if e.kind() == ErrorKind::NotFound => Ok(()),
167            Err(e) => Err(invalid(e)),
168        }
169    }
170    /// Explicit recovery: invalidate callers and archive raw state; no consent or resume.
171    pub fn recover_execution(&self) -> Result<serde_json::Value, CtlError> {
172        let _lock = self.control_lock()?;
173        let generation = crate::snapshot::new_snapshot_id();
174        let archive = self.dir.join("stop-history").join(&generation);
175        let mut archived = Vec::new();
176        for path in [self.stop_file(), self.generation_file()] {
177            match fs::symlink_metadata(&path) {
178                Ok(meta) if meta.is_file() && !meta.file_type().is_symlink() => {
179                    fs::create_dir_all(&archive).map_err(invalid)?;
180                    let name = path
181                        .file_name()
182                        .ok_or_else(|| invalid("missing filename"))?;
183                    let dest = archive.join(name);
184                    fs::copy(&path, &dest).map_err(invalid)?;
185                    archived.push(dest.to_string_lossy().into_owned());
186                }
187                Ok(_) => {
188                    return Err(invalid(
189                        "state path is not a regular file; no automatic removal",
190                    ));
191                }
192                Err(e) if e.kind() == ErrorKind::NotFound => {}
193                Err(e) => return Err(invalid(e)),
194            }
195        }
196        self.write_control(&self.generation_file(), generation.as_bytes())?;
197        match fs::remove_file(self.stop_file()) {
198            Ok(()) => {}
199            Err(e) if e.kind() == ErrorKind::NotFound => {}
200            Err(e) => return Err(invalid(e)),
201        }
202        Ok(
203            serde_json::json!({"recovered":true,"generation":generation,"permission_granted":false,"resumed":false,"archived":archived}),
204        )
205    }
206    pub fn stop_record(&self) -> Option<StopRecord> {
207        let _lock = self.control_lock().ok()?;
208        self.read_stop(&self.read_generation().ok()?).ok().flatten()
209    }
210    pub fn stop_timestamp(&self) -> Option<u64> {
211        self.stop_requested()
212            .then(|| self.stop_record().map_or(u64::MAX, |r| r.requested_ms))
213    }
214    pub fn stop_details(&self) -> String {
215        if !self.stop_requested() {
216            return String::new();
217        }
218        match self.stop_record() {
219            Some(r) => format!("\n已停止上一轮执行(来源:{})。新任务可在交接提示中点击“开始”;旧调用不会恢复。", r.source),
220            None => "\n运行状态损坏或版本不支持。请使用“修复运行状态”或 execution-recover;原始记录会归档,旧任务不会恢复。".into()
221        }
222    }
223}
224#[derive(Clone)]
225pub struct Invocation(Result<String, (ErrorCode, String, Option<serde_json::Value>)>);
226impl Invocation {
227    pub fn capture(paths: &SignalPaths) -> Self {
228        Self(
229            paths
230                .execution_generation()
231                .map_err(|e| (e.code, e.message, e.evidence)),
232        )
233    }
234    pub fn generation(&self) -> Result<String, CtlError> {
235        self.0
236            .clone()
237            .map_err(|(code, message, evidence)| CtlError {
238                code,
239                message,
240                evidence,
241            })
242    }
243    pub fn check(&self, paths: &SignalPaths) -> Result<(), CtlError> {
244        if self.generation()? != paths.execution_generation()? {
245            return Err(cancelled());
246        }
247        Ok(())
248    }
249}
250thread_local! { static CURRENT: RefCell<Option<Invocation>> = const { RefCell::new(None) }; }
251pub struct Scope(Option<Invocation>);
252impl Scope {
253    pub fn enter() -> Self {
254        Self(CURRENT.with(|c| c.replace(Some(Invocation::capture(&SignalPaths::default())))))
255    }
256}
257impl Drop for Scope {
258    fn drop(&mut self) {
259        CURRENT.with(|c| c.replace(self.0.take()));
260    }
261}
262pub fn generation() -> Result<String, CtlError> {
263    CURRENT.with(|c| {
264        c.borrow_mut()
265            .get_or_insert_with(|| Invocation::capture(&SignalPaths::default()))
266            .generation()
267    })
268}
269
270/// Bind a queued child to the generation selected by its parent UI action.
271pub fn bind_generation(expected: String) {
272    CURRENT.with(|c| *c.borrow_mut() = Some(Invocation(Ok(expected))));
273}
274pub fn check_current() -> Result<(), CtlError> {
275    CURRENT.with(|c| {
276        c.borrow_mut()
277            .get_or_insert_with(|| Invocation::capture(&SignalPaths::default()))
278            .check(&SignalPaths::default())
279    })
280}