pub mod run;
pub mod store;
#[cfg(test)]
pub(crate) mod test_support;
mod types;
use std::fs;
use std::path::Path;
use std::time::{SystemTime, UNIX_EPOCH};
use anyhow::{Context, Result};
use tracing::warn;
use crate::agent_identity::classify_agent_kind;
use crate::multiplexer::{AgentStatus, Multiplexer};
pub use store::StateStore;
pub use types::{AgentState, LastDoneCycleState, PaneKey, RuntimeState};
pub(crate) fn write_atomic(path: &Path, content: &[u8]) -> Result<()> {
let tmp = path.with_extension("json.tmp");
fs::write(&tmp, content).context("Failed to write temp file")?;
fs::rename(&tmp, path).context("Failed to rename temp file")?;
Ok(())
}
pub fn persist_agent_update(
mux: &dyn Multiplexer,
pane_id: &str,
status: Option<AgentStatus>,
title_override: Option<String>,
agent_session_id: Option<String>,
) {
persist_agent_snapshot(mux, pane_id, status, title_override, agent_session_id, true);
}
pub fn persist_agent_registration(mux: &dyn Multiplexer, pane_id: &str) {
persist_agent_snapshot(mux, pane_id, None, None, None, false);
}
fn persist_agent_snapshot(
mux: &dyn Multiplexer,
pane_id: &str,
status: Option<AgentStatus>,
title_override: Option<String>,
agent_session_id: Option<String>,
preserve_existing: bool,
) {
let pane_key = PaneKey {
backend: mux.name().to_string(),
instance: mux.instance_id(),
pane_id: pane_id.to_string(),
};
let live_info = match mux.get_live_pane_info(pane_id) {
Ok(Some(info)) => info,
Ok(None) => {
warn!(%pane_id, "pane not found, skipping state persist");
return;
}
Err(e) => {
warn!(error = %e, "failed to get live pane info, skipping state persist");
return;
}
};
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let (live_pid, boot_id) = if mux.name() == "tmux" {
let Some(live_pid) = live_info.pid else {
warn!(%pane_id, "tmux pane PID unavailable, skipping state persist");
return;
};
let boot_id = match mux.server_boot_id() {
Ok(Some(boot_id)) => Some(boot_id),
Ok(None) => {
warn!(%pane_id, "tmux server identity unavailable, skipping state persist");
return;
}
Err(error) => {
warn!(%pane_id, %error, "failed to get tmux server identity, skipping state persist");
return;
}
};
(live_pid, boot_id)
} else {
(
live_info.pid.unwrap_or(0),
mux.server_boot_id().unwrap_or(None),
)
};
let existing = preserve_existing
.then(|| {
StateStore::new()
.ok()
.and_then(|store| store.get_agent(&pane_key).ok().flatten())
.filter(|state| {
mux.name() != "tmux" || (state.pane_pid == live_pid && state.boot_id == boot_id)
})
})
.flatten();
let final_status = status.or(existing.as_ref().and_then(|e| e.status));
let previous_status = existing.as_ref().and_then(|state| state.status);
let status_ts = match final_status {
None => None,
Some(_) if final_status == previous_status => {
Some(existing.as_ref().and_then(|e| e.status_ts).unwrap_or(now))
}
Some(_) => Some(now),
};
let previous_activity_ts = existing.as_ref().and_then(AgentState::activity_ts);
let activity_ts = resolve_activity_ts(
preserve_existing,
status,
previous_status,
previous_activity_ts,
now,
);
let existing_agent_kind = existing.as_ref().and_then(|e| e.agent_kind.clone());
let agent_session_id = agent_session_id.or_else(|| {
existing
.as_ref()
.and_then(|state| state.agent_session_id.clone())
});
let live_title_for_classify = live_info.title.clone();
let pane_title = title_override
.or_else(|| existing.as_ref().and_then(|state| state.pane_title.clone()))
.or(live_info.title);
let agent_kind = merge_agent_kind(
classify_agent_kind(
live_info.current_command.as_deref(),
live_title_for_classify.as_deref(),
),
existing_agent_kind,
);
let state = AgentState {
pane_key,
workdir: live_info.working_dir,
status: final_status,
status_ts,
activity_ts,
pane_title,
pane_pid: live_pid,
command: live_info.current_command.unwrap_or_default(),
updated_ts: now,
window_name: live_info.window,
session_name: live_info.session,
boot_id,
agent_kind,
agent_session_id,
};
if let Ok(store) = StateStore::new()
&& let Err(e) = store.upsert_agent(&state)
{
warn!(error = %e, "failed to persist agent state");
}
}
fn resolve_activity_ts(
preserve_existing: bool,
explicit_status: Option<AgentStatus>,
previous_status: Option<AgentStatus>,
previous_activity_ts: Option<u64>,
now: u64,
) -> Option<u64> {
if !preserve_existing
|| previous_activity_ts.is_none()
|| (explicit_status.is_some() && explicit_status != previous_status)
{
Some(now)
} else {
previous_activity_ts
}
}
fn merge_agent_kind(new: Option<String>, existing: Option<String>) -> Option<String> {
existing.or(new)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn registration_starts_activity_without_status() {
assert_eq!(
resolve_activity_ts(false, None, Some(AgentStatus::Done), Some(10), 20),
Some(20)
);
}
#[test]
fn first_status_starts_new_activity() {
assert_eq!(
resolve_activity_ts(true, Some(AgentStatus::Working), None, Some(10), 20),
Some(20)
);
}
#[test]
fn status_transition_starts_new_activity() {
assert_eq!(
resolve_activity_ts(
true,
Some(AgentStatus::Done),
Some(AgentStatus::Working),
Some(10),
20,
),
Some(20)
);
}
#[test]
fn repeated_status_preserves_activity() {
assert_eq!(
resolve_activity_ts(
true,
Some(AgentStatus::Working),
Some(AgentStatus::Working),
Some(10),
20,
),
Some(10)
);
}
#[test]
fn metadata_update_preserves_activity() {
assert_eq!(
resolve_activity_ts(true, None, None, Some(10), 20),
Some(10)
);
}
#[test]
fn merge_keeps_existing_when_new_is_none() {
let merged = merge_agent_kind(None, Some("claude".into()));
assert_eq!(merged, Some("claude".into()));
}
#[test]
fn merge_preserves_existing_against_drift() {
let merged = merge_agent_kind(Some("vibe".into()), Some("claude".into()));
assert_eq!(merged, Some("claude".into()));
}
#[test]
fn merge_returns_none_when_both_none() {
assert_eq!(merge_agent_kind(None, None), None);
}
#[test]
fn merge_classifies_when_existing_is_none() {
let merged = merge_agent_kind(Some("claude".into()), None);
assert_eq!(merged, Some("claude".into()));
}
}