use std::path::{Path, PathBuf};
use anyhow::{Context, Result};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ThreadState {
Idle,
Running,
AwaitingInput,
Cancelled,
Staged,
Done,
Failed,
Orphaned,
}
impl ThreadState {
pub const ALL: [ThreadState; 8] = [
ThreadState::Idle,
ThreadState::Running,
ThreadState::AwaitingInput,
ThreadState::Cancelled,
ThreadState::Staged,
ThreadState::Done,
ThreadState::Failed,
ThreadState::Orphaned,
];
pub fn as_str(self) -> &'static str {
match self {
ThreadState::Idle => "idle",
ThreadState::Running => "running",
ThreadState::AwaitingInput => "awaiting_input",
ThreadState::Cancelled => "cancelled",
ThreadState::Staged => "staged",
ThreadState::Done => "done",
ThreadState::Failed => "failed",
ThreadState::Orphaned => "orphaned",
}
}
pub fn describe(self) -> &'static str {
match self {
ThreadState::Idle => "bound to a session and a workspace; nothing running",
ThreadState::Running => "a run is in flight",
ThreadState::AwaitingInput => "the run is blocked on an approval or a question",
ThreadState::Cancelled => "stopped at a safe point; the partial turn was kept",
ThreadState::Staged => "finished, and it left drafts in the outbox",
ThreadState::Done => "finished, nothing pending",
ThreadState::Failed => "the run errored",
ThreadState::Orphaned => "the connector restarted while this run was in flight",
}
}
pub fn resolved_by(self) -> &'static str {
match self {
ThreadState::Idle => "an owner message starts a run",
ThreadState::Running => "the run ends, or the owner presses Stop",
ThreadState::AwaitingInput => {
"the owner answers, or the timeout fires and the call is refused"
}
ThreadState::Cancelled => "an owner message starts a new run on the same conversation",
ThreadState::Staged => "release or reject the drafts, from here or any outbox surface",
ThreadState::Done => "an owner message starts a run",
ThreadState::Failed => {
"an owner message starts a run; the error was posted, not just logged"
}
ThreadState::Orphaned => {
"the restart sweep announces it in the thread and resets to idle"
}
}
}
pub fn is_active(self) -> bool {
matches!(self, ThreadState::Running | ThreadState::AwaitingInput)
}
pub fn accepts_new_prompt(self) -> bool {
matches!(
self,
ThreadState::Idle
| ThreadState::Cancelled
| ThreadState::Staged
| ThreadState::Done
| ThreadState::Failed
)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Event {
OwnerSpoke,
AskedForInput,
InputSettled,
Finished { staged: bool },
Errored,
StopPressed,
ConnectorRestarted,
OrphanAnnounced,
}
pub fn next(state: ThreadState, event: Event) -> Option<ThreadState> {
use Event::*;
use ThreadState::*;
match (state, event) {
(s, OwnerSpoke) if s.accepts_new_prompt() => Some(Running),
(Running, OwnerSpoke) => Some(Running),
(AwaitingInput, OwnerSpoke) => Some(AwaitingInput),
(Orphaned, OwnerSpoke) => None,
(Running, AskedForInput) => Some(AwaitingInput),
(AwaitingInput, InputSettled) => Some(Running),
(s, Finished { staged }) if s.is_active() => Some(if staged { Staged } else { Done }),
(s, Errored) if s.is_active() => Some(Failed),
(s, StopPressed) if s.is_active() => Some(Cancelled),
(s, ConnectorRestarted) if s.is_active() => Some(Orphaned),
(s, ConnectorRestarted) => Some(s),
(Orphaned, OrphanAnnounced) => Some(Idle),
_ => None,
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RunMarker {
pub pid: u32,
pub started_at: DateTime<Utc>,
}
impl RunMarker {
pub fn here() -> Self {
Self {
pid: std::process::id(),
started_at: Utc::now(),
}
}
pub fn is_live(&self) -> bool {
mecha_core::process_alive(self.pid)
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ThreadRecord {
pub key: String,
pub channel_id: String,
pub thread_ts: String,
pub state: ThreadState,
#[serde(default)]
pub session_id: Option<String>,
#[serde(default)]
pub workspace: Option<PathBuf>,
pub mode: String,
#[serde(default)]
pub last_seen_ts: Option<String>,
#[serde(default)]
pub run: Option<RunMarker>,
#[serde(default)]
pub stream_ts: Option<String>,
#[serde(default)]
pub controls_ts: Option<String>,
pub updated_at: DateTime<Utc>,
}
impl ThreadRecord {
fn new(channel_id: &str, thread_ts: &str, mode: &str) -> Self {
Self {
key: key_for(channel_id, thread_ts),
channel_id: channel_id.to_string(),
thread_ts: thread_ts.to_string(),
state: ThreadState::Idle,
session_id: None,
workspace: None,
mode: mode.to_string(),
last_seen_ts: None,
run: None,
stream_ts: None,
controls_ts: None,
updated_at: Utc::now(),
}
}
}
pub fn key_for(channel_id: &str, thread_ts: &str) -> String {
fn safe(s: &str) -> String {
s.chars()
.map(|c| if c.is_ascii_alphanumeric() { c } else { '-' })
.collect()
}
format!("{}-{}", safe(channel_id), safe(thread_ts))
}
pub struct ThreadStore {
root: PathBuf,
}
impl ThreadStore {
pub fn open(root: impl Into<PathBuf>) -> Result<Self> {
let root = root.into();
mecha_slack::store::create_private_dir(&root)
.with_context(|| format!("creating {}", root.display()))?;
Ok(Self { root })
}
pub fn root(&self) -> &Path {
&self.root
}
fn path(&self, key: &str) -> PathBuf {
self.root.join(format!("{key}.json"))
}
pub fn get(&self, key: &str) -> Result<Option<ThreadRecord>> {
mecha_slack::store::read_json(&self.path(key))
.with_context(|| format!("reading thread {key}"))
}
pub fn put(&self, record: &ThreadRecord) -> Result<()> {
mecha_slack::store::write_private_json(&self.path(&record.key), record)
.with_context(|| format!("writing thread {}", record.key))
}
pub fn all(&self) -> Result<Vec<ThreadRecord>> {
let mut out = Vec::new();
let entries = match std::fs::read_dir(&self.root) {
Ok(e) => e,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(out),
Err(e) => return Err(e).context("listing threads"),
};
for entry in entries.filter_map(std::result::Result::ok) {
let path = entry.path();
if path.extension().and_then(|e| e.to_str()) != Some("json") {
continue;
}
match mecha_slack::store::read_json::<ThreadRecord>(&path) {
Ok(Some(r)) => out.push(r),
Ok(None) => {}
Err(e) => {
tracing::warn!("unreadable thread record {}: {e}", path.display());
}
}
}
out.sort_by_key(|r| std::cmp::Reverse(r.updated_at));
Ok(out)
}
pub fn ensure(
&self,
channel_id: &str,
thread_ts: &str,
default_mode: &str,
) -> Result<ThreadRecord> {
let key = key_for(channel_id, thread_ts);
if let Some(existing) = self.get(&key)? {
return Ok(existing);
}
let record = ThreadRecord::new(channel_id, thread_ts, default_mode);
self.put(&record)?;
Ok(record)
}
pub fn apply(&self, key: &str, event: Event) -> Result<Option<ThreadRecord>> {
let Some(mut record) = self.get(key)? else {
return Ok(None);
};
let Some(state) = next(record.state, event) else {
return Ok(None);
};
record.state = state;
record.updated_at = Utc::now();
if !state.is_active() {
record.run = None;
}
self.put(&record)?;
Ok(Some(record))
}
pub fn sweep(&self) -> Result<Vec<ThreadRecord>> {
let mut orphaned = Vec::new();
for record in self.all()? {
if !record.state.is_active() {
continue;
}
if record.run.as_ref().is_some_and(RunMarker::is_live) {
continue;
}
if let Some(updated) = self.apply(&record.key, Event::ConnectorRestarted)? {
orphaned.push(updated);
}
}
Ok(orphaned)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn scratch(name: &str) -> PathBuf {
let dir = std::env::temp_dir().join(format!(
"mecha-slack-threads-{name}-{}-{}",
std::process::id(),
Utc::now().timestamp_nanos_opt().unwrap_or_default()
));
let _ = std::fs::remove_dir_all(&dir);
dir
}
#[test]
fn the_typed_name_and_the_stored_name_are_the_same_name() {
for state in ThreadState::ALL {
let stored = serde_json::to_string(&state).unwrap();
assert_eq!(
stored.trim_matches('"'),
state.as_str(),
"{state:?} disagrees with itself"
);
}
}
#[test]
fn every_state_says_what_it_means_and_how_to_leave_it() {
for state in ThreadState::ALL {
assert!(!state.describe().is_empty(), "{state:?} has no meaning");
assert!(
!state.resolved_by().is_empty(),
"{state:?} has no way out — a thread in it cannot be rescued"
);
}
}
#[test]
fn every_state_is_reachable_and_leaveable() {
let events = [
Event::OwnerSpoke,
Event::AskedForInput,
Event::InputSettled,
Event::Finished { staged: true },
Event::Finished { staged: false },
Event::Errored,
Event::StopPressed,
Event::ConnectorRestarted,
Event::OrphanAnnounced,
];
for target in ThreadState::ALL {
let reachable = ThreadState::ALL.iter().any(|&from| {
from != target && events.iter().any(|&e| next(from, e) == Some(target))
});
assert!(
reachable || target == ThreadState::Idle,
"{target:?} unreachable"
);
let leaveable = events
.iter()
.any(|&e| matches!(next(target, e), Some(s) if s != target));
assert!(leaveable, "{target:?} cannot be left");
}
}
#[test]
fn a_message_starts_a_run_when_idle_and_steers_one_when_running() {
assert_eq!(
next(ThreadState::Idle, Event::OwnerSpoke),
Some(ThreadState::Running)
);
assert_eq!(
next(ThreadState::Running, Event::OwnerSpoke),
Some(ThreadState::Running),
"steering does not restart the run"
);
assert_eq!(
next(ThreadState::AwaitingInput, Event::OwnerSpoke),
Some(ThreadState::AwaitingInput),
"a message while a question is pending does not answer it"
);
}
#[test]
fn an_orphan_must_be_announced_before_the_thread_is_reused() {
assert_eq!(next(ThreadState::Orphaned, Event::OwnerSpoke), None);
assert_eq!(
next(ThreadState::Orphaned, Event::OrphanAnnounced),
Some(ThreadState::Idle)
);
}
#[test]
fn finishing_with_drafts_is_a_different_state_from_finishing_without() {
assert_eq!(
next(ThreadState::Running, Event::Finished { staged: true }),
Some(ThreadState::Staged)
);
assert_eq!(
next(ThreadState::Running, Event::Finished { staged: false }),
Some(ThreadState::Done)
);
}
#[test]
fn events_that_do_not_apply_are_refused_rather_than_ignored() {
assert_eq!(next(ThreadState::Idle, Event::AskedForInput), None);
assert_eq!(next(ThreadState::Done, Event::StopPressed), None);
assert_eq!(next(ThreadState::Idle, Event::InputSettled), None);
}
#[test]
fn a_restart_is_only_news_for_a_thread_that_was_running() {
assert_eq!(
next(ThreadState::Running, Event::ConnectorRestarted),
Some(ThreadState::Orphaned)
);
assert_eq!(
next(ThreadState::AwaitingInput, Event::ConnectorRestarted),
Some(ThreadState::Orphaned)
);
assert_eq!(
next(ThreadState::Done, Event::ConnectorRestarted),
Some(ThreadState::Done)
);
}
#[test]
fn a_key_cannot_escape_the_directory() {
let key = key_for("../../etc", "passwd/../..");
assert!(!key.contains('/'), "{key}");
assert!(!key.contains(".."), "{key}");
}
#[test]
fn the_store_round_trips_and_ensure_is_idempotent() {
let dir = scratch("roundtrip");
let store = ThreadStore::open(&dir).unwrap();
let first = store.ensure("D1", "1000.1", "ask").unwrap();
assert_eq!(first.state, ThreadState::Idle);
assert_eq!(first.mode, "ask");
let again = store.ensure("D1", "1000.1", "allow").unwrap();
assert_eq!(
again.mode, "ask",
"ensure must not rewrite a thread's mode from a default"
);
assert_eq!(store.all().unwrap().len(), 1);
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn applying_an_event_persists_it_and_clears_a_finished_runs_marker() {
let dir = scratch("apply");
let store = ThreadStore::open(&dir).unwrap();
let mut record = store.ensure("D1", "1000.1", "ask").unwrap();
record.run = Some(RunMarker::here());
record.state = ThreadState::Running;
store.put(&record).unwrap();
let done = store
.apply(&record.key, Event::Finished { staged: false })
.unwrap()
.expect("the event applies");
assert_eq!(done.state, ThreadState::Done);
assert!(
done.run.is_none(),
"a marker outliving its run is what a later sweep would reason about"
);
let reread = store.get(&record.key).unwrap().unwrap();
assert_eq!(reread.state, ThreadState::Done);
assert!(
store
.apply(&record.key, Event::StopPressed)
.unwrap()
.is_none(),
"an event that does not apply changes nothing"
);
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn the_sweep_orphans_a_dead_run_and_leaves_a_live_one_alone() {
let dir = scratch("sweep");
let store = ThreadStore::open(&dir).unwrap();
let mut dead = store.ensure("D1", "1.0", "ask").unwrap();
dead.state = ThreadState::Running;
dead.run = Some(RunMarker {
pid: i32::MAX as u32, started_at: Utc::now(),
});
store.put(&dead).unwrap();
let mut live = store.ensure("D2", "2.0", "ask").unwrap();
live.state = ThreadState::Running;
live.run = Some(RunMarker::here());
store.put(&live).unwrap();
let mut idle = store.ensure("D3", "3.0", "ask").unwrap();
idle.state = ThreadState::Done;
store.put(&idle).unwrap();
let orphaned = store.sweep().unwrap();
assert_eq!(orphaned.len(), 1, "only the dead one");
assert_eq!(orphaned[0].key, dead.key);
assert_eq!(orphaned[0].state, ThreadState::Orphaned);
assert_eq!(
store.get(&live.key).unwrap().unwrap().state,
ThreadState::Running,
"a live run must not be swept out from under itself"
);
assert_eq!(
store.get(&idle.key).unwrap().unwrap().state,
ThreadState::Done
);
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn a_thread_with_no_marker_at_all_is_orphaned_too() {
let dir = scratch("nomarker");
let store = ThreadStore::open(&dir).unwrap();
let mut record = store.ensure("D1", "1.0", "ask").unwrap();
record.state = ThreadState::Running;
record.run = None;
store.put(&record).unwrap();
let orphaned = store.sweep().unwrap();
assert_eq!(orphaned.len(), 1);
std::fs::remove_dir_all(&dir).ok();
}
}