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, status_reason};
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 fallback_resume_command(&self, _original_command: &str) -> Option<String> {
None
}
fn fallback_resume_is_cwd_scoped(&self) -> bool {
true
}
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 status_reason: Option<String>,
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) -> StateUpdate {
match event {
HarnessEvent::SessionStarted {
harness_session_id,
resumed: _,
} => StateUpdate {
status: Some(SessionStatus::Working),
metadata_clear: vec![meta::NEEDS_INPUT],
harness_session_id: harness_session_id.clone(),
..Default::default()
},
HarnessEvent::Working => StateUpdate {
status: Some(SessionStatus::Working),
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::Waiting),
status_reason: Some(status_reason::IDLE.to_owned()),
metadata_set,
metadata_clear: vec![meta::NEEDS_INPUT],
set_idle_since: true,
..Default::default()
}
}
HarnessEvent::NeedsInput { reason } => StateUpdate {
status: Some(SessionStatus::Waiting),
status_reason: Some(status_reason::needs_input(&truncate_chars(
&reason.to_string(),
500,
))),
metadata_clear: vec![meta::NEEDS_INPUT],
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::Waiting),
status_reason: Some(status_reason::IDLE.to_owned()),
metadata_set,
metadata_clear: vec![meta::NEEDS_INPUT],
notify: true,
..Default::default()
}
}
HarnessEvent::SessionEnded { reason: _ } => StateUpdate::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_fallback_resume_command_default_is_none() {
assert!(
generic::GenericAdapter
.fallback_resume_command("bash")
.is_none()
);
}
#[test]
fn test_fallback_resume_is_cwd_scoped_default_is_true() {
assert!(generic::GenericAdapter.fallback_resume_is_cwd_scoped());
assert!(claude::ClaudeAdapter.fallback_resume_is_cwd_scoped());
assert!(pi::PiAdapter.fallback_resume_is_cwd_scoped());
}
#[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_working_and_id() {
let update = transition_for_event(&HarnessEvent::SessionStarted {
harness_session_id: Some("sid-1".into()),
resumed: false,
});
assert_eq!(update.status, Some(SessionStatus::Working));
assert_eq!(update.status_reason, None);
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_working() {
let update = transition_for_event(&HarnessEvent::SessionStarted {
harness_session_id: None,
resumed: true,
});
assert_eq!(update.status, Some(SessionStatus::Working));
assert!(update.harness_session_id.is_none());
}
#[test]
fn test_transition_working_clears_needs_input_and_error() {
let update = transition_for_event(&HarnessEvent::Working);
assert_eq!(update.status, Some(SessionStatus::Working));
assert_eq!(update.status_reason, None);
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_waiting_idle_and_summary() {
let update = transition_for_event(&HarnessEvent::TurnFinished {
summary: Some("Fixed the bug".into()),
});
assert_eq!(update.status, Some(SessionStatus::Waiting));
assert_eq!(update.status_reason.as_deref(), Some(status_reason::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),
});
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 });
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 });
assert_eq!(update.metadata_clear, vec![meta::NEEDS_INPUT]);
}
#[test]
fn test_transition_needs_input_sets_waiting_and_reason_and_notifies() {
let update = transition_for_event(&HarnessEvent::NeedsInput {
reason: NeedsInputReason::Permission,
});
assert_eq!(update.status, Some(SessionStatus::Waiting));
assert_eq!(
update.status_reason.as_deref(),
Some("needs_input:permission")
);
assert!(update.metadata_set.is_empty());
assert!(update.notify);
}
#[test]
fn test_transition_needs_input_idle_reason_is_distinct_from_plain_idle() {
let update = transition_for_event(&HarnessEvent::NeedsInput {
reason: NeedsInputReason::Idle,
});
assert_eq!(update.status_reason.as_deref(), Some("needs_input:idle"));
assert_ne!(update.status_reason.as_deref(), Some(status_reason::IDLE));
}
#[test]
fn test_transition_failed_sets_waiting_idle_error_and_notifies() {
let update = transition_for_event(&HarnessEvent::Failed {
error: "API error".into(),
rate_limited: false,
});
assert_eq!(update.status, Some(SessionStatus::Waiting));
assert_eq!(update.status_reason.as_deref(), Some(status_reason::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,
});
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,
});
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,
});
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),
});
assert_eq!(
update.status_reason.as_deref().unwrap().len()
- status_reason::NEEDS_INPUT_PREFIX.len(),
500
);
}
#[test]
fn test_transition_session_ended_leaves_status_alone() {
let update = transition_for_event(&HarnessEvent::SessionEnded { reason: None });
assert_eq!(update.status, None);
assert_eq!(update.status_reason, None);
let update = transition_for_event(&HarnessEvent::SessionEnded {
reason: Some("exit".into()),
});
assert_eq!(update.status, None);
}
}