use serde::{Deserialize, Serialize};
use std::path::PathBuf;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Phase {
Running,
Done,
Error,
Stopped,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SessionState {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub handoff_reason: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub workflow: Option<WorkflowProgress>,
#[serde(default)]
pub version: u32,
#[serde(default)]
pub call_id: String,
#[serde(default)]
pub seq: u64,
#[serde(default)]
pub started_ms: u64,
#[serde(default)]
pub finished_ms: Option<u64>,
#[serde(default)]
pub action: Option<crate::activity::ActionState>,
#[serde(default)]
pub recent_actions: Vec<crate::activity::ActionState>,
#[serde(default)]
pub error_code: Option<String>,
pub pid: u32,
pub command: String,
pub app: Option<String>,
pub step: Option<String>,
pub phase: Phase,
pub ts_ms: u64,
}
impl SessionState {
pub fn now(
pid: u32,
command: &str,
app: Option<&str>,
step: Option<String>,
phase: Phase,
) -> Self {
Self {
handoff_reason: None,
workflow: None,
version: 1,
call_id: String::new(),
seq: 0,
started_ms: unix_ms(),
finished_ms: None,
action: None,
recent_actions: Vec::new(),
error_code: None,
pid,
command: command.to_string(),
app: app.map(str::to_string),
step,
phase,
ts_ms: unix_ms(),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkflowProgress {
pub run_id: String,
pub title: String,
pub status: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reason: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub step_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error_code: Option<crate::ErrorCode>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub actions: Vec<String>,
}
pub fn unix_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
#[derive(Debug, Clone)]
pub struct SignalPaths {
pub dir: PathBuf,
}
impl Default for SignalPaths {
fn default() -> Self {
let dir = std::env::var_os("LOCALAPPDATA")
.map(PathBuf::from)
.unwrap_or_else(std::env::temp_dir)
.join("actl");
Self { dir }
}
}
impl SignalPaths {
pub fn at(dir: impl Into<PathBuf>) -> Self {
Self { dir: dir.into() }
}
pub fn state_file(&self) -> PathBuf {
self.dir.join("state.json")
}
pub fn stop_file(&self) -> PathBuf {
self.dir.join("stop-requested")
}
pub fn write_state(&self, st: &SessionState) {
if let Err(error) = self.try_write_state(st) {
crate::log_health::failure(
&self.dir,
&st.call_id,
st.workflow.as_ref().map(|w| w.run_id.as_str()),
"state_write",
0,
&error,
);
eprintln!("[actl-state] INTERNAL: state publication failed: {error}");
}
}
fn try_write_state(&self, st: &SessionState) -> std::io::Result<()> {
std::fs::create_dir_all(&self.dir)?;
if st.version >= 2 && !st.call_id.is_empty() {
let dir = self.dir.join("calls");
std::fs::create_dir_all(&dir)?;
let file = dir.join(format!("{}.json", st.call_id));
Self::atomic_write(&file, st)?;
}
Self::atomic_write(&self.state_file(), st)
}
fn atomic_write(file: &std::path::Path, st: &SessionState) -> std::io::Result<()> {
let tmp = file.with_extension(format!("{}.{}.tmp", st.pid, st.call_id));
let result = (|| {
std::fs::write(&tmp, serde_json::to_vec(st)?)?;
std::fs::rename(&tmp, file)
})();
let _ = std::fs::remove_file(tmp);
result
}
pub fn read_calls(&self) -> Vec<SessionState> {
let now = unix_ms();
let mut states = Vec::new();
if let Ok(entries) = std::fs::read_dir(self.dir.join("calls")) {
for entry in entries.flatten() {
let path = entry.path();
if path.extension().is_none_or(|e| e != "json") {
continue;
}
let Some(st) = std::fs::read(&path)
.ok()
.and_then(|b| serde_json::from_slice::<SessionState>(&b).ok())
else {
continue;
};
if now.saturating_sub(st.ts_ms) > 60_000 {
let _ = std::fs::remove_file(path);
} else {
states.push(st);
}
}
}
if states.is_empty()
&& let Some(st) = self.read_state()
{
states.push(st);
}
states.sort_by_key(|s| s.ts_ms);
let terminal_count = states.iter().filter(|s| s.phase != Phase::Running).count();
let mut excess = terminal_count.saturating_sub(128);
states.retain(|s| {
if excess > 0 && s.phase != Phase::Running && !s.call_id.is_empty() {
if s.call_id
.bytes()
.all(|c| c.is_ascii_alphanumeric() || c == b'-')
{
let _ = std::fs::remove_file(
self.dir.join("calls").join(format!("{}.json", s.call_id)),
);
}
excess -= 1;
false
} else {
true
}
});
states
}
pub fn read_state(&self) -> Option<SessionState> {
let text = std::fs::read_to_string(self.state_file()).ok()?;
serde_json::from_str(&text).ok()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn tmp() -> SignalPaths {
static N: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0);
let n = N.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
let dir = std::env::temp_dir().join(format!("actl-state-test-{}-{n}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
SignalPaths::at(&dir)
}
#[test]
fn state_roundtrip_and_fields() {
let p = tmp();
let st = SessionState::now(
42,
"press",
Some("记事本"),
Some("3/8".into()),
Phase::Running,
);
p.write_state(&st);
let back = p.read_state().expect("read back");
assert_eq!(back.command, "press");
assert_eq!(back.app.as_deref(), Some("记事本"));
assert_eq!(back.step.as_deref(), Some("3/8"));
assert_eq!(back.phase, Phase::Running);
assert!(back.ts_ms > 0);
let _ = std::fs::remove_dir_all(&p.dir);
}
#[test]
fn stop_flag_is_sticky_until_cleared() {
let p = tmp();
assert!(!p.stop_requested());
p.request_stop().unwrap();
assert!(p.stop_requested());
p.clear_stop().unwrap();
assert!(!p.stop_requested());
let _ = std::fs::remove_dir_all(&p.dir);
}
#[test]
fn phase_serializes_snake_case() {
assert_eq!(
serde_json::to_string(&Phase::Stopped).unwrap(),
r#""stopped""#
);
}
}