pub mod claude;
pub mod codex;
pub mod generic;
pub mod pi;
pub mod registry;
use std::path::{Path, PathBuf};
use anyhow::Result;
use pulpo_common::session::{SessionStatus, meta};
use serde::{Deserialize, Serialize};
pub use registry::HarnessRegistry;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum NeedsInputReason {
Permission,
Question,
Idle,
Other(String),
}
impl std::fmt::Display for NeedsInputReason {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Permission => write!(f, "permission"),
Self::Question => write!(f, "question"),
Self::Idle => write!(f, "idle"),
Self::Other(label) => write!(f, "{label}"),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum HarnessEvent {
SessionStarted {
harness_session_id: Option<String>,
resumed: bool,
},
Working,
TurnFinished { summary: Option<String> },
NeedsInput { reason: NeedsInputReason },
Failed { error: String, rate_limited: bool },
SessionEnded { reason: Option<String> },
}
pub struct SpawnContext<'a> {
pub session_id: &'a str,
pub session_name: &'a str,
pub workdir: &'a str,
pub command: &'a str,
pub data_dir: &'a Path,
}
pub struct SpawnPlan {
pub command: String,
pub env: Vec<(String, String)>,
pub files: Vec<PathBuf>,
pub harness_session_id: Option<String>,
}
impl SpawnPlan {
#[must_use]
pub fn unchanged(command: &str) -> Self {
Self {
command: command.to_owned(),
env: Vec::new(),
files: Vec::new(),
harness_session_id: None,
}
}
}
pub(crate) fn resolve_pulpo_bin() -> String {
resolve_pulpo_bin_from(std::env::current_exe().ok().as_deref())
}
pub(crate) fn resolve_pulpo_bin_from(current_exe: Option<&Path>) -> String {
current_exe
.and_then(Path::parent)
.map(|dir| dir.join("pulpo"))
.filter(|candidate| candidate.is_file())
.map_or_else(
|| "pulpo".to_owned(),
|candidate| candidate.to_string_lossy().into_owned(),
)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct HarnessSignals {
pub lifecycle: bool,
pub rate_limit: bool,
pub error: bool,
}
impl HarnessSignals {
#[must_use]
pub const fn all() -> Self {
Self {
lifecycle: true,
rate_limit: true,
error: true,
}
}
#[must_use]
pub const fn lifecycle_only() -> Self {
Self {
lifecycle: true,
rate_limit: false,
error: false,
}
}
#[must_use]
pub const fn none() -> Self {
Self {
lifecycle: false,
rate_limit: false,
error: false,
}
}
}
pub trait HarnessAdapter: Send + Sync {
fn id(&self) -> &'static str;
fn matches(&self, argv0: &str) -> bool;
fn prepare_spawn(&self, ctx: &SpawnContext) -> Result<SpawnPlan>;
fn resume_command(&self, original_command: &str, harness_session_id: &str) -> Option<String>;
fn parse_event(&self, raw: &serde_json::Value) -> Result<Option<HarnessEvent>>;
fn emits_events(&self) -> bool;
fn owned_signals(&self) -> HarnessSignals {
HarnessSignals::all()
}
}
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct StateUpdate {
pub status: Option<SessionStatus>,
pub metadata_set: Vec<(&'static str, String)>,
pub metadata_clear: Vec<&'static str>,
pub harness_session_id: Option<String>,
pub set_idle_since: bool,
pub notify: bool,
}
fn truncate_chars(s: &str, max_chars: usize) -> String {
s.chars().take(max_chars).collect()
}
#[must_use]
pub fn transition_for_event(event: &HarnessEvent, backend_alive: bool) -> StateUpdate {
match event {
HarnessEvent::SessionStarted {
harness_session_id,
resumed: _,
} => StateUpdate {
status: Some(SessionStatus::Active),
metadata_clear: vec![meta::NEEDS_INPUT],
harness_session_id: harness_session_id.clone(),
..Default::default()
},
HarnessEvent::Working => StateUpdate {
status: Some(SessionStatus::Active),
metadata_clear: vec![meta::NEEDS_INPUT, meta::ERROR_STATUS, meta::ERROR_STATUS_AT],
..Default::default()
},
HarnessEvent::TurnFinished { summary } => {
let metadata_set = summary
.as_ref()
.map(|s| (meta::LAST_SUMMARY, truncate_chars(s, 200)))
.into_iter()
.collect();
StateUpdate {
status: Some(SessionStatus::Idle),
metadata_set,
metadata_clear: vec![meta::NEEDS_INPUT],
set_idle_since: true,
..Default::default()
}
}
HarnessEvent::NeedsInput { reason } => StateUpdate {
status: Some(SessionStatus::Idle),
metadata_set: vec![(meta::NEEDS_INPUT, truncate_chars(&reason.to_string(), 500))],
notify: true,
..Default::default()
},
HarnessEvent::Failed {
error,
rate_limited,
} => {
let now = chrono::Utc::now().to_rfc3339();
let error = truncate_chars(error, 500);
let mut metadata_set = vec![
(meta::ERROR_STATUS, error.clone()),
(meta::ERROR_STATUS_AT, now.clone()),
];
if *rate_limited {
metadata_set.push((meta::RATE_LIMIT, error));
metadata_set.push((meta::RATE_LIMIT_AT, now));
}
StateUpdate {
status: Some(SessionStatus::Idle),
metadata_set,
metadata_clear: vec![meta::NEEDS_INPUT],
notify: true,
..Default::default()
}
}
HarnessEvent::SessionEnded { reason: _ } => StateUpdate {
status: Some(if backend_alive {
SessionStatus::Ready
} else {
SessionStatus::Stopped
}),
..Default::default()
},
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_harness_signals_all() {
let signals = HarnessSignals::all();
assert!(signals.lifecycle);
assert!(signals.rate_limit);
assert!(signals.error);
}
#[test]
fn test_harness_signals_lifecycle_only() {
let signals = HarnessSignals::lifecycle_only();
assert!(signals.lifecycle);
assert!(!signals.rate_limit);
assert!(!signals.error);
}
#[test]
fn test_harness_signals_none() {
let signals = HarnessSignals::none();
assert!(!signals.lifecycle);
assert!(!signals.rate_limit);
assert!(!signals.error);
}
#[test]
fn test_owned_signals_default_is_all() {
assert_eq!(
generic::GenericAdapter.owned_signals(),
HarnessSignals::all()
);
assert_eq!(claude::ClaudeAdapter.owned_signals(), HarnessSignals::all());
}
#[test]
fn test_needs_input_reason_display() {
assert_eq!(NeedsInputReason::Permission.to_string(), "permission");
assert_eq!(NeedsInputReason::Question.to_string(), "question");
assert_eq!(NeedsInputReason::Idle.to_string(), "idle");
assert_eq!(
NeedsInputReason::Other("custom_thing".into()).to_string(),
"custom_thing"
);
}
#[test]
fn test_needs_input_reason_serde_roundtrip() {
let reason = NeedsInputReason::Other("weird".into());
let json = serde_json::to_string(&reason).unwrap();
let back: NeedsInputReason = serde_json::from_str(&json).unwrap();
assert_eq!(back, reason);
}
#[test]
fn test_spawn_plan_unchanged() {
let plan = SpawnPlan::unchanged("claude -p hi");
assert_eq!(plan.command, "claude -p hi");
assert!(plan.env.is_empty());
assert!(plan.files.is_empty());
assert!(plan.harness_session_id.is_none());
}
#[test]
fn test_transition_session_started_sets_active_and_id() {
let update = transition_for_event(
&HarnessEvent::SessionStarted {
harness_session_id: Some("sid-1".into()),
resumed: false,
},
true,
);
assert_eq!(update.status, Some(SessionStatus::Active));
assert_eq!(update.harness_session_id.as_deref(), Some("sid-1"));
assert_eq!(update.metadata_clear, vec![meta::NEEDS_INPUT]);
assert!(!update.notify);
}
#[test]
fn test_transition_session_started_resumed_still_active() {
let update = transition_for_event(
&HarnessEvent::SessionStarted {
harness_session_id: None,
resumed: true,
},
true,
);
assert_eq!(update.status, Some(SessionStatus::Active));
assert!(update.harness_session_id.is_none());
}
#[test]
fn test_transition_working_clears_needs_input_and_error() {
let update = transition_for_event(&HarnessEvent::Working, true);
assert_eq!(update.status, Some(SessionStatus::Active));
assert_eq!(
update.metadata_clear,
vec![meta::NEEDS_INPUT, meta::ERROR_STATUS, meta::ERROR_STATUS_AT]
);
assert!(!update.notify);
}
#[test]
fn test_transition_turn_finished_sets_idle_and_summary() {
let update = transition_for_event(
&HarnessEvent::TurnFinished {
summary: Some("Fixed the bug".into()),
},
true,
);
assert_eq!(update.status, Some(SessionStatus::Idle));
assert!(update.set_idle_since);
assert_eq!(
update.metadata_set,
vec![(meta::LAST_SUMMARY, "Fixed the bug".to_owned())]
);
assert!(!update.notify);
}
#[test]
fn test_transition_turn_finished_truncates_summary_to_200_chars() {
let long = "x".repeat(500);
let update = transition_for_event(
&HarnessEvent::TurnFinished {
summary: Some(long),
},
true,
);
assert_eq!(update.metadata_set[0].1.chars().count(), 200);
}
#[test]
fn test_transition_turn_finished_no_summary_sets_no_metadata() {
let update = transition_for_event(&HarnessEvent::TurnFinished { summary: None }, true);
assert!(update.metadata_set.is_empty());
assert!(update.set_idle_since);
}
#[test]
fn test_transition_turn_finished_clears_needs_input() {
let update = transition_for_event(&HarnessEvent::TurnFinished { summary: None }, true);
assert_eq!(update.metadata_clear, vec![meta::NEEDS_INPUT]);
}
#[test]
fn test_transition_needs_input_sets_idle_and_reason_and_notifies() {
let update = transition_for_event(
&HarnessEvent::NeedsInput {
reason: NeedsInputReason::Permission,
},
true,
);
assert_eq!(update.status, Some(SessionStatus::Idle));
assert_eq!(
update.metadata_set,
vec![(meta::NEEDS_INPUT, "permission".to_owned())]
);
assert!(update.notify);
}
#[test]
fn test_transition_failed_sets_idle_error_and_notifies() {
let update = transition_for_event(
&HarnessEvent::Failed {
error: "API error".into(),
rate_limited: false,
},
true,
);
assert_eq!(update.status, Some(SessionStatus::Idle));
assert!(
update
.metadata_set
.contains(&(meta::ERROR_STATUS, "API error".to_owned()))
);
assert!(
!update
.metadata_set
.iter()
.any(|(k, _)| *k == meta::RATE_LIMIT)
);
assert!(update.notify);
}
#[test]
fn test_transition_failed_rate_limited_sets_rate_limit_metadata() {
let update = transition_for_event(
&HarnessEvent::Failed {
error: "429 rate limited".into(),
rate_limited: true,
},
true,
);
assert!(
update
.metadata_set
.iter()
.any(|(k, v)| *k == meta::RATE_LIMIT && v == "429 rate limited")
);
assert!(
update
.metadata_set
.iter()
.any(|(k, _)| *k == meta::RATE_LIMIT_AT)
);
}
#[test]
fn test_transition_failed_clears_needs_input() {
let update = transition_for_event(
&HarnessEvent::Failed {
error: "API error".into(),
rate_limited: false,
},
true,
);
assert_eq!(update.metadata_clear, vec![meta::NEEDS_INPUT]);
}
#[test]
fn test_transition_failed_truncates_error_to_500_chars() {
let long = "x".repeat(2000);
let update = transition_for_event(
&HarnessEvent::Failed {
error: long,
rate_limited: true,
},
true,
);
for (key, value) in &update.metadata_set {
if *key == meta::ERROR_STATUS || *key == meta::RATE_LIMIT {
assert_eq!(value.chars().count(), 500, "key {key} not truncated");
}
}
}
#[test]
fn test_transition_needs_input_truncates_other_reason_to_500_chars() {
let long = "y".repeat(2000);
let update = transition_for_event(
&HarnessEvent::NeedsInput {
reason: NeedsInputReason::Other(long),
},
true,
);
assert_eq!(update.metadata_set[0].1.chars().count(), 500);
}
#[test]
fn test_transition_session_ended_ready_when_backend_alive() {
let update = transition_for_event(&HarnessEvent::SessionEnded { reason: None }, true);
assert_eq!(update.status, Some(SessionStatus::Ready));
}
#[test]
fn test_transition_session_ended_stopped_when_backend_gone() {
let update = transition_for_event(
&HarnessEvent::SessionEnded {
reason: Some("exit".into()),
},
false,
);
assert_eq!(update.status, Some(SessionStatus::Stopped));
}
}