1use 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 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 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
270pub 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}