use std::collections::BTreeMap;
use super::workflow::WorkflowMeta;
pub const DEFAULT_PORT: u16 = 7274;
pub const EPHEMERAL_IDLE_SECS: u64 = 300;
pub const EPHEMERAL_COUNTDOWN_AFTER_SECS: u64 = 5;
pub const SESSION_START_TIMEOUT_SECS: u64 = 30;
pub const SESSION_IDLE_TIMEOUT_SECS: u64 = crate::config::DEFAULT_INACTIVITY_TIMEOUT_SECS;
pub const SESSION_NO_WORK_TIMEOUT_SECS: u64 = SESSION_IDLE_TIMEOUT_SECS * 2;
pub const MAX_STORED_SESSIONS: usize = 200;
pub const MAX_ARCHIVED_SESSIONS: usize = 2000;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DaemonMode {
Persistent,
Ephemeral,
}
impl DaemonMode {
pub fn as_str(self) -> &'static str {
match self {
DaemonMode::Persistent => "persistent",
DaemonMode::Ephemeral => "ephemeral",
}
}
pub fn parse(s: &str) -> Option<Self> {
match s {
"persistent" => Some(DaemonMode::Persistent),
"ephemeral" => Some(DaemonMode::Ephemeral),
_ => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SessionLifecycle {
Running,
Completed,
Failed,
Cancelled,
}
impl SessionLifecycle {
pub fn label(self) -> &'static str {
match self {
SessionLifecycle::Running => "running",
SessionLifecycle::Completed => "completed",
SessionLifecycle::Failed => "failed",
SessionLifecycle::Cancelled => "cancelled",
}
}
pub fn css_class(self) -> &'static str {
match self {
SessionLifecycle::Running => "running",
SessionLifecycle::Completed => "completed",
SessionLifecycle::Failed => "failed",
SessionLifecycle::Cancelled => "cancelled",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ProcStatus {
Waiting,
Running,
Ok,
Graceful,
Fail,
Skipped,
}
impl ProcStatus {
pub fn as_str(self) -> &'static str {
match self {
ProcStatus::Waiting => "waiting",
ProcStatus::Running => "running",
ProcStatus::Ok => "ok",
ProcStatus::Graceful => "graceful",
ProcStatus::Fail => "fail",
ProcStatus::Skipped => "skipped",
}
}
pub fn parse(s: &str) -> Option<Self> {
match s {
"waiting" => Some(ProcStatus::Waiting),
"running" => Some(ProcStatus::Running),
"ok" => Some(ProcStatus::Ok),
"graceful" => Some(ProcStatus::Graceful),
"fail" => Some(ProcStatus::Fail),
"skipped" => Some(ProcStatus::Skipped),
_ => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ProcKind {
Build,
Skill,
Annotate,
}
impl ProcKind {
pub fn as_str(self) -> &'static str {
match self {
ProcKind::Build => "build",
ProcKind::Skill => "skill",
ProcKind::Annotate => "annotate",
}
}
pub fn parse(s: &str) -> Option<Self> {
match s {
"build" => Some(ProcKind::Build),
"skill" => Some(ProcKind::Skill),
"annotate" => Some(ProcKind::Annotate),
_ => None,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct OutputLine {
pub at: f64,
pub text: String,
}
#[derive(Debug, Clone, PartialEq)]
pub struct ProcRecord {
pub index: usize,
pub previous_attempt: Option<usize>,
pub label: String,
pub kind: ProcKind,
pub status: ProcStatus,
pub skill_name: Option<String>,
pub harness: Option<String>,
pub model: Option<String>,
pub started_at: Option<u64>,
pub note: Option<String>,
pub detail: Option<String>,
pub fail_reason: Option<String>,
pub elapsed: Option<f64>,
pub lines: Vec<OutputLine>,
pub container_name: Option<String>,
pub container_runtime: Option<String>,
pub cast_path: Option<String>,
pub diff_path: Option<String>,
pub skill_source: Option<String>,
pub route: Option<String>,
pub result_path: Option<String>,
pub annotate_target: Option<String>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct SkillMeta {
pub name: String,
pub harness: String,
}
pub const DEFAULT_JOB_RETRIES: u32 = 25;
pub const JOB_FAIL_STREAK_CAP: u32 = 3;
#[derive(Debug, Clone, PartialEq, Default)]
pub struct SupervisorState {
pub job_attempt: u32,
pub retries: u32,
pub next_retry_at: Option<u64>,
pub fail_signature: Option<String>,
pub fail_streak: u32,
pub gave_up: Option<String>,
pub restarted_as: Option<String>,
}
impl SupervisorState {
pub fn fresh(retries: u32) -> SupervisorState {
SupervisorState { job_attempt: 1, retries, ..Default::default() }
}
pub fn supervised(&self) -> bool {
self.retries > 0
}
pub fn attempt(&self) -> u32 {
self.job_attempt.max(1)
}
pub fn inherited(&self) -> SupervisorState {
SupervisorState {
job_attempt: self.attempt() + 1,
retries: self.retries,
next_retry_at: None,
fail_signature: self.fail_signature.clone(),
fail_streak: self.fail_streak,
gave_up: None,
restarted_as: None,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct Session {
pub id: String,
pub started_at: u64,
pub ended_at: Option<u64>,
pub profile: Option<String>,
pub kind: Option<String>,
pub repo: String,
pub branch: String,
pub skills: Vec<SkillMeta>,
pub procs: Vec<ProcRecord>,
pub last_seen_at: u64,
pub client_connected: bool,
pub run_pid: Option<u32>,
pub workflow: Option<WorkflowMeta>,
pub parent_session: Option<String>,
pub supervisor: SupervisorState,
}
#[derive(Debug, Clone, PartialEq)]
pub struct OpenRepo {
pub path: String,
pub opened_at: u64,
pub clean: bool,
}
#[derive(Debug, Clone, PartialEq)]
pub struct Store {
pub mode: DaemonMode,
pub port: u16,
pub started_at: u64,
pub active_clients: u32,
pub last_activity: u64,
pub no_alive_since: Option<u64>,
pub sessions: BTreeMap<String, Session>,
pub open_repos: BTreeMap<String, OpenRepo>,
}
impl Store {
pub fn new(mode: DaemonMode, port: u16, now: u64) -> Store {
Store {
mode,
port,
started_at: now,
active_clients: 0,
last_activity: now,
no_alive_since: Some(now),
sessions: BTreeMap::new(),
open_repos: BTreeMap::new(),
}
}
pub fn touch(&mut self, now: u64) {
self.last_activity = now;
}
pub fn open_repo(&mut self, repo: OpenRepo) {
self.open_repos.insert(repo.path.clone(), repo);
}
pub fn job_running_in(&self, repo: &str, now: u64) -> bool {
self.sessions.values().any(|s| s.repo == repo && s.lifecycle_status(now) == SessionLifecycle::Running)
}
pub fn alive_clients(&self, now: u64) -> u32 {
self
.sessions
.values()
.filter(|s| s.client_connected && s.lifecycle_status(now) == SessionLifecycle::Running)
.count() as u32
}
pub fn reconcile(&mut self, now: u64) {
for session in self.sessions.values_mut() {
if session.client_connected && session.lifecycle_status(now) != SessionLifecycle::Running {
session.client_connected = false;
}
}
self.active_clients = self.sessions.values().filter(|s| s.client_connected).count() as u32;
if self.alive_clients(now) > 0 {
self.no_alive_since = None;
} else if self.no_alive_since.is_none() {
self.no_alive_since = Some(now);
}
}
pub fn ephemeral_shutdown_in_secs(&self, now: u64) -> Option<u64> {
if self.mode != DaemonMode::Ephemeral {
return None;
}
let since = self.no_alive_since?;
let idle = now.saturating_sub(since);
if idle < EPHEMERAL_COUNTDOWN_AFTER_SECS {
return None;
}
Some(EPHEMERAL_IDLE_SECS.saturating_sub(idle))
}
pub fn should_shutdown_ephemeral(&self, now: u64) -> bool {
self.mode == DaemonMode::Ephemeral
&& self.alive_clients(now) == 0
&& self.no_alive_since.is_some_and(|since| now.saturating_sub(since) >= EPHEMERAL_IDLE_SECS)
}
pub fn session_mut(&mut self, id: &str) -> Option<&mut Session> {
self.sessions.get_mut(id)
}
pub fn proc_mut(&mut self, session_id: &str, proc_index: usize) -> Option<&mut ProcRecord> {
self.session_mut(session_id).and_then(|s| s.procs.iter_mut().find(|p| p.index == proc_index))
}
pub fn insert_session(&mut self, id: String, session: Session) {
self.sessions.insert(id, session);
trim_sessions_to_cap(&mut self.sessions, crate::now_secs());
}
}
pub fn trim_sessions_to_cap(sessions: &mut std::collections::BTreeMap<String, Session>, now: u64) {
trim_sessions_to(sessions, now, MAX_STORED_SESSIONS)
}
fn trim_sessions_to(sessions: &mut std::collections::BTreeMap<String, Session>, now: u64, cap: usize) {
while sessions.len() > cap {
let Some(old_id) = sessions
.iter()
.filter(|(_, s)| s.lifecycle_status(now) != SessionLifecycle::Running)
.min_by_key(|(_, s)| {
let tier = match s.parent_session.as_deref() {
Some(parent) if !sessions.contains_key(parent) => 0u8,
Some(_) => 1,
None => 2,
};
(tier, s.started_at)
})
.map(|(id, _)| id.clone())
else {
break;
};
sessions.remove(&old_id);
let orphaned: Vec<String> = sessions
.iter()
.filter(|(_, s)| s.parent_session.as_deref() == Some(old_id.as_str()))
.filter(|(_, s)| s.lifecycle_status(now) != SessionLifecycle::Running)
.map(|(id, _)| id.clone())
.collect();
for id in orphaned {
sessions.remove(&id);
}
}
}
pub fn sessions_for_index(sessions: &BTreeMap<String, Session>, now: u64) -> Vec<&Session> {
let mut list: Vec<&Session> = sessions.values().filter(|session| session.parent_session.is_none()).collect();
list.sort_by(|a, b| {
let a_live = a.lifecycle_status(now) == SessionLifecycle::Running;
let b_live = b.lifecycle_status(now) == SessionLifecycle::Running;
match (a_live, b_live) {
(true, false) => std::cmp::Ordering::Less,
(false, true) => std::cmp::Ordering::Greater,
_ => b.started_at.cmp(&a.started_at),
}
});
list
}
impl Session {
pub fn has_incomplete_procs(&self) -> bool {
self.procs.iter().any(|p| p.status == ProcStatus::Running || p.status == ProcStatus::Waiting)
}
pub fn has_started_work(&self) -> bool {
self.procs.iter().any(|p| p.started_at.is_some() || !matches!(p.status, ProcStatus::Waiting))
}
pub(crate) fn has_live_proc(&self) -> bool {
self.procs.iter().any(|p| matches!(p.status, ProcStatus::Running | ProcStatus::Waiting))
}
pub(crate) fn last_work_end(&self) -> Option<u64> {
self
.procs
.iter()
.filter_map(|p| match (p.started_at, p.elapsed) {
(Some(start), Some(elapsed)) => Some(start.saturating_add(elapsed as u64)),
_ => None,
})
.max()
}
pub(crate) fn liveness_deadline(&self) -> u64 {
if self.has_started_work() {
self.last_seen_at.saturating_add(SESSION_IDLE_TIMEOUT_SECS)
} else {
self.started_at.saturating_add(SESSION_START_TIMEOUT_SECS)
}
}
pub(crate) fn proc_is_superseded(&self, proc: &ProcRecord) -> bool {
self.proc_next_attempt(proc).is_some()
}
pub(crate) fn proc_next_attempt(&self, proc: &ProcRecord) -> Option<&ProcRecord> {
if let Some(next) = self.procs.iter().find(|candidate| candidate.previous_attempt == Some(proc.index)) {
return Some(next);
}
if proc.previous_attempt.is_some() {
return None;
}
let name = proc.skill_name.as_deref().filter(|n| !n.is_empty())?;
self
.procs
.iter()
.filter(|later| {
later.previous_attempt.is_none()
&& later.index > proc.index
&& later.kind == proc.kind
&& later.skill_name.as_deref() == Some(name)
})
.min_by_key(|later| later.index)
}
pub(crate) fn proc_previous_attempt(&self, proc: &ProcRecord) -> Option<&ProcRecord> {
if let Some(index) = proc.previous_attempt {
return self.procs.iter().find(|candidate| candidate.index == index);
}
if self.procs.iter().any(|candidate| candidate.previous_attempt == Some(proc.index)) {
return None;
}
let name = proc.skill_name.as_deref().filter(|name| !name.is_empty())?;
self
.procs
.iter()
.filter(|earlier| {
earlier.index < proc.index
&& earlier.kind == proc.kind
&& earlier.skill_name.as_deref() == Some(name)
&& proc.previous_attempt.is_none()
})
.max_by_key(|earlier| earlier.index)
}
pub(crate) fn proc_first_attempt<'a>(&'a self, proc: &'a ProcRecord) -> &'a ProcRecord {
let mut first = proc;
let mut seen = std::collections::BTreeSet::from([first.index]);
while let Some(previous) = self.proc_previous_attempt(first) {
if !seen.insert(previous.index) {
break;
}
first = previous;
}
first
}
pub(crate) fn proc_attempt(&self, proc: &ProcRecord) -> (usize, usize) {
if proc.previous_attempt.is_some()
|| self.procs.iter().any(|candidate| candidate.previous_attempt == Some(proc.index))
{
let mut root = proc;
let mut seen = std::collections::BTreeSet::new();
seen.insert(root.index);
while let Some(previous) =
root.previous_attempt.and_then(|index| self.procs.iter().find(|candidate| candidate.index == index))
{
if !seen.insert(previous.index) {
break;
}
root = previous;
}
let mut ordinal = 1;
let mut total = 1;
let mut current = root;
let mut forward_seen = std::collections::BTreeSet::from([root.index]);
while let Some(next) = self.procs.iter().find(|candidate| candidate.previous_attempt == Some(current.index)) {
if !forward_seen.insert(next.index) {
break;
}
total += 1;
if next.index <= proc.index {
ordinal += 1;
}
current = next;
}
return (ordinal, total);
}
let Some(name) = proc.skill_name.as_deref().filter(|n| !n.is_empty()) else {
return (1, 1);
};
let mut ordinal = 0;
let mut total = 0;
for p in &self.procs {
if p.kind == proc.kind && p.skill_name.as_deref() == Some(name) {
total += 1;
if p.index <= proc.index {
ordinal += 1;
}
}
}
(ordinal.max(1), total.max(1))
}
pub fn lifecycle_status(&self, now: u64) -> SessionLifecycle {
if self.ended_at.is_some() {
if self.has_incomplete_procs() {
return SessionLifecycle::Cancelled;
}
let failed: Vec<&ProcRecord> =
self.procs.iter().filter(|p| p.status == ProcStatus::Fail && !self.proc_is_superseded(p)).collect();
let interrupted = !failed.is_empty()
&& failed.iter().all(|p| {
matches!(
p.fail_reason.as_deref(),
Some(
crate::failure::reason::FORCE_STOPPED
| crate::failure::reason::FORCE_RESTARTED
| crate::failure::reason::SESSION_END_INCOMPLETE
)
)
});
if interrupted {
return SessionLifecycle::Cancelled;
}
if !failed.is_empty() {
return SessionLifecycle::Failed;
}
return SessionLifecycle::Completed;
}
if now > self.liveness_deadline() {
return SessionLifecycle::Failed;
}
SessionLifecycle::Running
}
pub fn duration_secs(&self, now: u64) -> Option<u64> {
if let Some(end) = self.ended_at {
return Some(end.saturating_sub(self.started_at));
}
let lifecycle = self.lifecycle_status(now);
if lifecycle == SessionLifecycle::Running {
return Some(now.saturating_sub(self.started_at));
}
if lifecycle == SessionLifecycle::Failed {
return Some(self.liveness_deadline().saturating_sub(self.started_at));
}
None
}
}
#[cfg(test)]
mod tests {
use super::*;
fn test_proc(status: ProcStatus) -> ProcRecord {
ProcRecord {
index: 0,
previous_attempt: None,
label: "skill".into(),
kind: ProcKind::Skill,
status,
skill_name: None,
harness: None,
model: None,
started_at: None,
note: None,
detail: None,
fail_reason: None,
elapsed: None,
lines: Vec::new(),
container_name: None,
container_runtime: None,
cast_path: None,
diff_path: None,
skill_source: None,
route: None,
result_path: None,
annotate_target: None,
}
}
fn stored_session(id: &str, started_at: u64, ended_at: Option<u64>, last_seen_at: u64) -> Session {
Session {
id: id.into(),
started_at,
ended_at,
profile: None,
kind: None,
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: Vec::new(),
last_seen_at,
client_connected: false,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: Default::default(),
}
}
#[test]
fn distinct_skill_names_are_parallel_tasks_not_attempts() {
let named = |index: usize, name: &str| {
let mut p = test_proc(ProcStatus::Ok);
p.index = index;
p.skill_name = Some(name.to_string());
p
};
let mut quota = stored_session("quota", 1, Some(2), 2);
quota.procs = vec![named(0, "quota-claude"), named(1, "quota-codex"), named(2, "quota-grok")];
for p in "a.procs {
assert_eq!(quota.proc_attempt(p), (1, 1), "{:?} must not read as a retry", p.skill_name);
assert!(quota.proc_next_attempt(p).is_none());
}
let mut retried = stored_session("retried", 1, Some(2), 2);
retried.procs = vec![named(0, "quota-claude"), named(1, "quota-claude")];
assert_eq!(retried.proc_attempt(&retried.procs[0]), (1, 2));
assert_eq!(retried.proc_attempt(&retried.procs[1]), (2, 2));
}
#[test]
fn store_trim_never_evicts_a_running_session() {
let now = 10_000;
let mut sessions = std::collections::BTreeMap::new();
let mut live = stored_session("live-old", 1, None, now);
live.procs = vec![test_proc(ProcStatus::Running)];
sessions.insert("live-old".to_string(), live);
sessions.insert("wedged".to_string(), stored_session("wedged", 2, None, 2));
for i in 0..MAX_STORED_SESSIONS {
let id = format!("done-{i:04}");
sessions.insert(id.clone(), stored_session(&id, 100 + i as u64, Some(200 + i as u64), 200 + i as u64));
}
trim_sessions_to_cap(&mut sessions, now);
assert_eq!(sessions.len(), MAX_STORED_SESSIONS);
assert!(sessions.contains_key("live-old"), "a RUNNING session survives the cap whatever its age");
assert!(!sessions.contains_key("wedged"), "a wedged session past its deadline is evictable");
assert!(!sessions.contains_key("done-0000"), "the oldest finished session went next");
assert!(sessions.contains_key(&format!("done-{:04}", MAX_STORED_SESSIONS - 1)));
}
#[test]
fn lifecycle_uses_the_terminal_proc_status_for_ended_jobs() {
let session = Session {
id: "done".into(),
started_at: 100,
ended_at: Some(200),
profile: None,
kind: None,
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: vec![ProcRecord {
index: 0,
previous_attempt: None,
label: "skill".into(),
kind: ProcKind::Skill,
status: ProcStatus::Ok,
skill_name: None,
harness: None,
model: None,
started_at: Some(100),
note: None,
detail: None,
fail_reason: None,
elapsed: Some(5.0),
lines: Vec::new(),
container_name: None,
container_runtime: None,
cast_path: None,
diff_path: None,
skill_source: None,
route: None,
result_path: None,
annotate_target: None,
}],
last_seen_at: 200,
client_connected: false,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: Default::default(),
};
assert_eq!(session.lifecycle_status(200), SessionLifecycle::Completed);
assert_eq!(session.duration_secs(200), Some(100));
let mut invalid_result = session.clone();
invalid_result.procs[0].status = ProcStatus::Fail;
invalid_result.procs[0].fail_reason = Some(crate::failure::reason::RESULT_INVALID.into());
assert_eq!(invalid_result.lifecycle_status(200), SessionLifecycle::Failed);
}
#[test]
fn a_recovered_retry_supersedes_its_failed_attempt() {
let mut first = test_proc(ProcStatus::Fail);
first.skill_name = Some("conventions-reviewer-claude".into());
first.fail_reason = Some(crate::failure::reason::CONTAINER_TIMEOUT.into());
let mut retry = test_proc(ProcStatus::Ok);
retry.index = 1;
retry.previous_attempt = Some(0);
retry.skill_name = Some("conventions-reviewer-claude".into());
let mut session = Session {
id: "retry".into(),
started_at: 100,
ended_at: Some(200),
profile: None,
kind: None,
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: vec![first, retry],
last_seen_at: 200,
client_connected: false,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: Default::default(),
};
assert_eq!(session.lifecycle_status(200), SessionLifecycle::Completed);
session.procs[1].status = ProcStatus::Fail;
session.procs[1].fail_reason = Some(crate::failure::reason::CONTAINER_TIMEOUT.into());
assert_eq!(session.lifecycle_status(200), SessionLifecycle::Failed);
session.procs[1].status = ProcStatus::Ok;
session.procs[1].fail_reason = None;
session.procs[1].previous_attempt = None;
session.procs[0].skill_name = None;
assert_eq!(session.lifecycle_status(200), SessionLifecycle::Failed);
}
#[test]
fn explicit_attempt_lineage_reaches_the_original_from_a_third_attempt() {
let mut first = test_proc(ProcStatus::Fail);
first.skill_name = Some("review".into());
let mut second = test_proc(ProcStatus::Fail);
second.index = 1;
second.previous_attempt = Some(0);
second.skill_name = Some("review".into());
let mut third = test_proc(ProcStatus::Running);
third.index = 2;
third.previous_attempt = Some(1);
third.skill_name = Some("review".into());
let session = Session {
id: "third".into(),
started_at: 100,
ended_at: None,
profile: None,
kind: None,
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: vec![first, second, third],
last_seen_at: 200,
client_connected: true,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: Default::default(),
};
assert_eq!(session.proc_attempt(&session.procs[0]), (1, 3));
assert_eq!(session.proc_attempt(&session.procs[1]), (2, 3));
assert_eq!(session.proc_attempt(&session.procs[2]), (3, 3));
assert_eq!(session.proc_first_attempt(&session.procs[2]).index, 0);
}
#[test]
fn lifecycle_fails_start_after_thirty_seconds_without_work() {
let session = Session {
id: "stale".into(),
started_at: 100,
ended_at: None,
profile: None,
kind: None,
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: Vec::new(),
last_seen_at: 100,
client_connected: true,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: Default::default(),
};
assert_eq!(session.lifecycle_status(100 + SESSION_START_TIMEOUT_SECS), SessionLifecycle::Running);
assert_eq!(session.lifecycle_status(100 + SESSION_START_TIMEOUT_SECS + 1), SessionLifecycle::Failed);
assert_eq!(session.duration_secs(100 + SESSION_START_TIMEOUT_SECS + 1), Some(SESSION_START_TIMEOUT_SECS));
}
#[test]
fn lifecycle_allows_thirty_minutes_idle_after_work_starts() {
let mut session = Session {
id: "idle".into(),
started_at: 100,
ended_at: None,
profile: None,
kind: None,
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: Vec::new(),
last_seen_at: 150,
client_connected: true,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: Default::default(),
};
let mut proc = test_proc(ProcStatus::Running);
proc.started_at = Some(110);
session.procs.push(proc);
assert_eq!(session.lifecycle_status(150 + SESSION_IDLE_TIMEOUT_SECS), SessionLifecycle::Running);
assert_eq!(session.lifecycle_status(150 + SESSION_IDLE_TIMEOUT_SECS + 1), SessionLifecycle::Failed);
assert_eq!(session.duration_secs(150 + SESSION_IDLE_TIMEOUT_SECS + 1), Some(50 + SESSION_IDLE_TIMEOUT_SECS));
}
#[test]
fn lifecycle_cancelled_when_ended_with_incomplete_procs() {
let session = Session {
id: "cancel".into(),
started_at: 1,
ended_at: Some(50),
profile: None,
kind: None,
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: vec![ProcRecord {
index: 0,
previous_attempt: None,
label: "skill".into(),
kind: ProcKind::Skill,
status: ProcStatus::Running,
skill_name: None,
harness: None,
model: None,
started_at: Some(1),
note: None,
detail: None,
fail_reason: None,
elapsed: None,
lines: Vec::new(),
container_name: None,
container_runtime: None,
cast_path: None,
diff_path: None,
skill_source: None,
route: None,
result_path: None,
annotate_target: None,
}],
last_seen_at: 50,
client_connected: false,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: Default::default(),
};
assert_eq!(session.lifecycle_status(50), SessionLifecycle::Cancelled);
}
#[test]
fn lifecycle_running_while_incomplete_procs_and_recent() {
let session = Session {
id: "test".into(),
started_at: 1,
ended_at: None,
profile: None,
kind: None,
repo: "/repo".into(),
branch: "main".into(),
skills: Vec::new(),
procs: vec![
ProcRecord {
index: 0,
previous_attempt: None,
label: "done".into(),
kind: ProcKind::Skill,
status: ProcStatus::Ok,
skill_name: None,
harness: None,
model: None,
started_at: None,
note: None,
detail: None,
fail_reason: None,
elapsed: None,
lines: Vec::new(),
container_name: None,
container_runtime: None,
cast_path: None,
diff_path: None,
skill_source: None,
route: None,
result_path: None,
annotate_target: None,
},
ProcRecord {
index: 1,
previous_attempt: None,
label: "still going".into(),
kind: ProcKind::Skill,
status: ProcStatus::Waiting,
skill_name: None,
harness: None,
model: None,
started_at: None,
note: None,
detail: None,
fail_reason: None,
elapsed: None,
lines: Vec::new(),
container_name: None,
container_runtime: None,
cast_path: None,
diff_path: None,
skill_source: None,
route: None,
result_path: None,
annotate_target: None,
},
],
last_seen_at: 1,
client_connected: true,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: Default::default(),
};
assert!(session.has_incomplete_procs());
assert_eq!(session.lifecycle_status(2), SessionLifecycle::Running);
}
#[test]
fn sessions_for_index_puts_running_first_then_recent() {
let mut running = Session {
id: "run".into(),
started_at: 10,
ended_at: None,
profile: None,
kind: None,
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: Vec::new(),
last_seen_at: 100,
client_connected: true,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: Default::default(),
};
let mut running_proc = test_proc(ProcStatus::Running);
running_proc.started_at = Some(10);
running.procs.push(running_proc);
let done = Session {
id: "done".into(),
started_at: 200,
ended_at: Some(250),
profile: None,
kind: None,
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: Vec::new(),
last_seen_at: 250,
client_connected: false,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: Default::default(),
};
let mut sessions = BTreeMap::new();
sessions.insert(done.id.clone(), done);
sessions.insert(running.id.clone(), running);
let mut annotation = sessions["done"].clone();
annotation.id = "annotate".into();
annotation.started_at = 300;
annotation.parent_session = Some("done".into());
sessions.insert(annotation.id.clone(), annotation);
let ordered = sessions_for_index(&sessions, 100);
assert_eq!(ordered.len(), 2);
assert_eq!(ordered[0].id, "run");
assert_eq!(ordered[1].id, "done");
}
#[test]
fn insert_session_evicts_oldest_when_over_cap() {
let mut store = Store::new(DaemonMode::Persistent, DEFAULT_PORT, 0);
for i in 0..=MAX_STORED_SESSIONS {
store.insert_session(
format!("{i:06}"),
Session {
id: format!("{i:06}"),
started_at: i as u64,
ended_at: None,
profile: None,
kind: None,
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: Vec::new(),
last_seen_at: i as u64,
client_connected: false,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: Default::default(),
},
);
}
assert_eq!(store.sessions.len(), MAX_STORED_SESSIONS);
assert!(!store.sessions.contains_key("000000"));
assert!(store.sessions.contains_key(&format!("{MAX_STORED_SESSIONS:06}")));
}
#[test]
fn annotations_are_evicted_before_the_jobs_they_belong_to() {
let session = |id: &str, at: u64, parent: Option<&str>| Session {
id: id.into(),
started_at: at,
ended_at: Some(at + 1),
profile: parent.map(|_| "annotate".to_string()),
kind: None,
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: Vec::new(),
last_seen_at: at,
client_connected: false,
run_pid: None,
workflow: None,
parent_session: parent.map(str::to_string),
supervisor: Default::default(),
};
let mut store = Store::new(DaemonMode::Persistent, DEFAULT_PORT, 0);
store.insert_session("oldjob".into(), session("oldjob", 1, None));
store.insert_session("livejob".into(), session("livejob", 2, None));
for i in 0..(MAX_STORED_SESSIONS - 2) {
let id = format!("ann{i:04}");
store.insert_session(id.clone(), session(&id, 1000 + i as u64, Some("livejob")));
}
assert_eq!(store.sessions.len(), MAX_STORED_SESSIONS);
store.insert_session("annnew".into(), session("annnew", 9000, Some("livejob")));
assert_eq!(store.sessions.len(), MAX_STORED_SESSIONS);
assert!(store.sessions.contains_key("oldjob"), "the oldest JOB outlives newer annotations");
assert!(store.sessions.contains_key("livejob"));
assert!(!store.sessions.contains_key("ann0000"), "the oldest annotation went instead");
let mut small: std::collections::BTreeMap<String, Session> = Default::default();
small.insert("parent".into(), session("parent", 1, None));
small.insert("kid1".into(), session("kid1", 2, Some("parent")));
small.insert("kid2".into(), session("kid2", 3, Some("parent")));
small.insert("other".into(), session("other", 4, None));
trim_sessions_to(&mut small, 10, 2);
assert_eq!(
small.keys().collect::<Vec<_>>(),
vec!["other", "parent"],
"the two jobs survive; their annotations are what the cap reclaims"
);
let mut orphans: std::collections::BTreeMap<String, Session> = Default::default();
orphans.insert("job".into(), session("job", 1, None));
orphans.insert("lost".into(), session("lost", 5, Some("evicted-long-ago")));
trim_sessions_to(&mut orphans, 10, 1);
assert_eq!(orphans.keys().collect::<Vec<_>>(), vec!["job"], "the orphan goes, the older job stays");
}
#[test]
fn session_id_is_six_lowercase_letters() {
let id = crate::runtime::random_nonce_6();
assert_eq!(id.len(), 6);
assert!(id.chars().all(|c| c.is_ascii_lowercase()));
}
#[test]
fn ephemeral_shutdown_after_idle() {
let now = 1_000_000;
let mut store = Store::new(DaemonMode::Ephemeral, DEFAULT_PORT, now);
store.reconcile(now + 100);
assert!(!store.should_shutdown_ephemeral(now + 100));
assert!(!store.should_shutdown_ephemeral(now + EPHEMERAL_IDLE_SECS - 1));
assert!(store.should_shutdown_ephemeral(now + EPHEMERAL_IDLE_SECS));
}
#[test]
fn startup_timed_out_client_not_counted_alive() {
let now = 100;
let mut store = Store::new(DaemonMode::Ephemeral, DEFAULT_PORT, now);
store.insert_session(
"stale".into(),
Session {
id: "stale".into(),
started_at: now,
ended_at: None,
profile: None,
kind: None,
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: Vec::new(),
last_seen_at: now,
client_connected: true,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: Default::default(),
},
);
assert_eq!(store.alive_clients(now + SESSION_START_TIMEOUT_SECS), 1);
assert_eq!(store.alive_clients(now + SESSION_START_TIMEOUT_SECS + 1), 0);
store.reconcile(now + SESSION_START_TIMEOUT_SECS + 1);
assert_eq!(store.active_clients, 0);
assert!(store.no_alive_since.is_some());
}
#[test]
fn ephemeral_countdown_after_no_alive_grace() {
let now = 0;
let store = Store::new(DaemonMode::Ephemeral, DEFAULT_PORT, now);
assert!(store.ephemeral_shutdown_in_secs(now + EPHEMERAL_COUNTDOWN_AFTER_SECS - 1).is_none());
assert_eq!(
store.ephemeral_shutdown_in_secs(now + EPHEMERAL_COUNTDOWN_AFTER_SECS),
Some(EPHEMERAL_IDLE_SECS - EPHEMERAL_COUNTDOWN_AFTER_SECS)
);
assert_eq!(store.ephemeral_shutdown_in_secs(now + EPHEMERAL_IDLE_SECS), Some(0));
}
#[test]
fn proc_kind_annotate_round_trips() {
assert_eq!(ProcKind::Annotate.as_str(), "annotate");
assert_eq!(ProcKind::parse("annotate"), Some(ProcKind::Annotate));
assert_eq!(ProcKind::parse("build"), Some(ProcKind::Build));
assert_eq!(ProcKind::parse("skill"), Some(ProcKind::Skill));
assert_eq!(ProcKind::parse("other"), None);
}
#[test]
fn persistent_never_auto_shutdown() {
let store = Store::new(DaemonMode::Persistent, DEFAULT_PORT, 0);
assert!(!store.should_shutdown_ephemeral(u64::MAX));
}
}