use std::path::Path;
use anyhow::Result;
use serde_json::Value;
use tracing::{info, warn};
use super::{
HarnessAdapter, HarnessEvent, NeedsInputReason, SpawnContext, SpawnPlan, resolve_pulpo_bin,
};
const EXTENSION_TEMPLATE: &str = include_str!("pulpo.ts.tmpl");
const PULPO_BIN_PLACEHOLDER: &str = "__PULPO_BIN_PATH__";
const SESSION_SELECTION_FLAGS: &[&str] = &[
"--session",
"--continue",
"-c",
"--resume",
"-r",
"--fork",
"--no-session",
];
const RESUME_CONFLICT_FLAGS: &[&str] =
&["--session", "--continue", "-c", "--resume", "-r", "--fork"];
const KNOWN_SUBCOMMANDS: &[&str] = &[
"install",
"remove",
"uninstall",
"update",
"list",
"config",
"auth",
];
pub struct PiAdapter;
impl HarnessAdapter for PiAdapter {
fn id(&self) -> &'static str {
"pi"
}
fn matches(&self, argv0: &str) -> bool {
argv0 == "pi"
}
fn prepare_spawn(&self, ctx: &SpawnContext) -> Result<SpawnPlan> {
let Ok(tokens) = shell_words::split(ctx.command) else {
warn!(
session = %ctx.session_name,
"pi adapter: command did not parse as shell words, spawning unchanged"
);
return Ok(SpawnPlan::unchanged(ctx.command));
};
let Some(pi_idx) = pi_token_index(&tokens) else {
return Ok(SpawnPlan::unchanged(ctx.command));
};
if tokens
.get(pi_idx + 1)
.is_some_and(|token| KNOWN_SUBCOMMANDS.contains(&token.as_str()))
{
info!(
session = %ctx.session_name,
"pi adapter: command invokes a pi subcommand, spawning unchanged"
);
return Ok(SpawnPlan::unchanged(ctx.command));
}
if has_flag(&tokens, SESSION_SELECTION_FLAGS) {
info!(
session = %ctx.session_name,
"pi adapter: command already pins session selection via another flag, spawning unchanged"
);
return Ok(SpawnPlan::unchanged(ctx.command));
}
match rewrite_spawn(ctx, tokens, pi_idx) {
Ok(plan) => Ok(plan),
Err(error) => {
warn!(
session = %ctx.session_name,
%error,
"pi adapter: failed to prepare spawn, spawning unchanged"
);
Ok(SpawnPlan::unchanged(ctx.command))
}
}
}
fn resume_command(&self, original_command: &str, harness_session_id: &str) -> Option<String> {
let mut tokens = shell_words::split(original_command).ok()?;
let pi_idx = pi_token_index(&tokens)?;
if has_flag(&tokens, RESUME_CONFLICT_FLAGS) {
return None;
}
strip_flag_with_value(&mut tokens, "--session-id");
tokens.splice(
(pi_idx + 1)..=pi_idx,
["--session-id".to_owned(), harness_session_id.to_owned()],
);
Some(shell_words::join(&tokens))
}
fn parse_event(&self, raw: &Value) -> Result<Option<HarnessEvent>> {
let event_name = raw.get("event").and_then(Value::as_str).unwrap_or("");
Ok(match event_name {
"session_start" => {
let harness_session_id = raw
.get("session_id")
.and_then(Value::as_str)
.map(str::to_owned);
let resumed = raw.get("reason").and_then(Value::as_str) != Some("new");
Some(HarnessEvent::SessionStarted {
harness_session_id,
resumed,
})
}
"agent_start" | "ui_prompt_end" => Some(HarnessEvent::Working),
"agent_settled" => Some(raw.get("error").and_then(Value::as_str).map_or_else(
|| {
HarnessEvent::TurnFinished {
summary: raw
.get("last_assistant_message")
.and_then(Value::as_str)
.map(str::to_owned),
}
},
|error| HarnessEvent::Failed {
rate_limited: is_rate_limited(error),
error: error.to_owned(),
},
)),
"ui_prompt_start" => {
let reason = if raw.get("kind").and_then(Value::as_str) == Some("confirm") {
NeedsInputReason::Permission
} else {
NeedsInputReason::Question
};
Some(HarnessEvent::NeedsInput { reason })
}
"session_shutdown" => {
(raw.get("reason").and_then(Value::as_str) == Some("quit")).then(|| {
HarnessEvent::SessionEnded {
reason: Some("quit".to_owned()),
}
})
}
_ => None,
})
}
fn emits_events(&self) -> bool {
true
}
}
fn rewrite_spawn(ctx: &SpawnContext, mut tokens: Vec<String>, pi_idx: usize) -> Result<SpawnPlan> {
let harness_dir = ctx.data_dir.join("harness").join(ctx.session_id);
std::fs::create_dir_all(&harness_dir)?;
let ext_path = harness_dir.join("pulpo.ts");
let pulpo_bin = resolve_pulpo_bin();
std::fs::write(&ext_path, render_extension(&pulpo_bin))?;
let ext_arg = ext_path.to_string_lossy().into_owned();
let harness_session_id = if let Some(idx) = flag_scan_region(&tokens)
.iter()
.position(|t| t == "--session-id")
{
if let Some(existing) = tokens.get(idx + 1).cloned() {
tokens.splice(idx + 2..idx + 2, ["-e".to_owned(), ext_arg]);
Some(existing)
} else {
let sid = ctx.session_id.to_owned();
tokens.splice((idx + 1)..=idx, [sid.clone(), "-e".to_owned(), ext_arg]);
Some(sid)
}
} else {
let sid = ctx.session_id.to_owned();
tokens.splice(
(pi_idx + 1)..=pi_idx,
[
"--session-id".to_owned(),
sid.clone(),
"-e".to_owned(),
ext_arg,
],
);
Some(sid)
};
Ok(SpawnPlan {
command: shell_words::join(&tokens),
env: Vec::new(),
files: vec![ext_path],
harness_session_id,
})
}
fn render_extension(pulpo_bin: &str) -> String {
EXTENSION_TEMPLATE.replace(PULPO_BIN_PLACEHOLDER, &escape_ts_string(pulpo_bin))
}
fn escape_ts_string(value: &str) -> String {
value.replace('\\', "\\\\").replace('"', "\\\"")
}
fn is_rate_limited(error: &str) -> bool {
let lower = error.to_lowercase();
contains_standalone_429(&lower)
|| lower.contains("too many requests")
|| ["rate limit", "ratelimit", "rate-limit", "rate_limit"]
.iter()
.any(|needle| lower.contains(needle))
}
fn contains_standalone_429(text: &str) -> bool {
let bytes = text.as_bytes();
let mut start = 0;
while let Some(rel) = text[start..].find("429") {
let idx = start + rel;
let before_is_digit = idx > 0 && bytes[idx - 1].is_ascii_digit();
let after_idx = idx + 3;
let after_is_digit = bytes.get(after_idx).is_some_and(u8::is_ascii_digit);
if !before_is_digit && !after_is_digit {
return true;
}
start = idx + 1;
}
false
}
fn flag_scan_region(tokens: &[String]) -> &[String] {
tokens
.iter()
.position(|t| t == "--")
.map_or(tokens, |idx| &tokens[..idx])
}
fn has_flag(tokens: &[String], flags: &[&str]) -> bool {
flag_scan_region(tokens)
.iter()
.any(|token| flags.contains(&token.as_str()))
}
fn pi_token_index(tokens: &[String]) -> Option<usize> {
let mut idx = 0;
if Path::new(tokens.first()?)
.file_name()
.and_then(|f| f.to_str())
== Some("env")
{
idx = 1;
while idx < tokens.len() {
let token = &tokens[idx];
if token.starts_with('-') || token.contains('=') {
idx += 1;
continue;
}
break;
}
}
let token = tokens.get(idx)?;
(Path::new(token).file_name().and_then(|f| f.to_str()) == Some("pi")).then_some(idx)
}
fn strip_flag_with_value(tokens: &mut Vec<String>, flag: &str) {
let mut i = 0;
while i < tokens.len() {
if tokens[i] == flag {
tokens.remove(i);
if i < tokens.len() {
tokens.remove(i);
}
} else {
i += 1;
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn ctx<'a>(data_dir: &'a Path, command: &'a str) -> SpawnContext<'a> {
SpawnContext {
session_id: "11111111-1111-1111-1111-111111111111",
session_name: "sess",
workdir: "/tmp/repo",
command,
data_dir,
}
}
#[test]
fn test_matches_pi_only() {
assert!(PiAdapter.matches("pi"));
assert!(!PiAdapter.matches("claude"));
assert!(!PiAdapter.matches(""));
}
#[test]
fn test_matches_does_not_match_npx_pi() {
assert!(!PiAdapter.matches("npx"));
}
#[test]
fn test_id_and_emits_events() {
assert_eq!(PiAdapter.id(), "pi");
assert!(PiAdapter.emits_events());
}
#[test]
fn test_prepare_spawn_rewrites_plain_pi() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter.prepare_spawn(&ctx(tmp.path(), "pi")).unwrap();
let sid = "11111111-1111-1111-1111-111111111111";
assert_eq!(plan.harness_session_id.as_deref(), Some(sid));
let tokens = shell_words::split(&plan.command).unwrap();
assert_eq!(tokens[0], "pi");
assert_eq!(tokens[1], "--session-id");
assert_eq!(tokens[2], sid);
assert_eq!(tokens[3], "-e");
assert_eq!(plan.files.len(), 1);
assert!(plan.files[0].exists());
assert_eq!(tokens[4], plan.files[0].to_string_lossy());
}
#[test]
fn test_prepare_spawn_preserves_original_args() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter
.prepare_spawn(&ctx(tmp.path(), "pi --model x"))
.unwrap();
let tokens = shell_words::split(&plan.command).unwrap();
assert_eq!(&tokens[5..], ["--model", "x"]);
}
#[test]
fn test_prepare_spawn_preserves_print_prompt_arg() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter
.prepare_spawn(&ctx(tmp.path(), "pi -p 'prompt'"))
.unwrap();
assert!(plan.command.ends_with("-p prompt"));
}
#[test]
fn test_prepare_spawn_writes_extension_under_harness_dir() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter.prepare_spawn(&ctx(tmp.path(), "pi")).unwrap();
let expected_dir = tmp
.path()
.join("harness")
.join("11111111-1111-1111-1111-111111111111");
assert_eq!(plan.files[0].parent().unwrap(), expected_dir);
assert_eq!(plan.files[0].file_name().unwrap(), "pulpo.ts");
}
#[test]
fn test_prepare_spawn_extension_contains_templated_bin_and_every_handler() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter.prepare_spawn(&ctx(tmp.path(), "pi")).unwrap();
let written = std::fs::read_to_string(&plan.files[0]).unwrap();
let pulpo_bin = resolve_pulpo_bin();
assert!(written.contains(&format!("const PULPO_BIN = \"{pulpo_bin}\";")));
assert!(!written.contains(PULPO_BIN_PLACEHOLDER));
for event in [
"session_start",
"agent_start",
"agent_end",
"agent_settled",
"ui_prompt_start",
"ui_prompt_end",
"session_shutdown",
] {
assert!(
written.contains(&format!("pi.on(\"{event}\",")),
"missing pi.on(\"{event}\", ...) registration"
);
}
}
#[test]
fn test_prepare_spawn_handles_absolute_path_argv0() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter
.prepare_spawn(&ctx(tmp.path(), "/usr/local/bin/pi -p hi"))
.unwrap();
assert!(plan.harness_session_id.is_some());
assert!(plan.command.starts_with("/usr/local/bin/pi --session-id"));
}
#[test]
fn test_prepare_spawn_handles_env_prefix() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter
.prepare_spawn(&ctx(tmp.path(), "env FOO=bar pi -p hi"))
.unwrap();
let tokens = shell_words::split(&plan.command).unwrap();
assert_eq!(tokens[0], "env");
assert_eq!(tokens[1], "FOO=bar");
assert_eq!(tokens[2], "pi");
assert_eq!(tokens[3], "--session-id");
}
#[test]
fn test_prepare_spawn_keeps_existing_session_id_and_adds_extension_only() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter
.prepare_spawn(&ctx(tmp.path(), "pi --session-id my-sid --model x"))
.unwrap();
assert_eq!(plan.harness_session_id.as_deref(), Some("my-sid"));
let tokens = shell_words::split(&plan.command).unwrap();
assert_eq!(tokens[0], "pi");
assert_eq!(tokens[1], "--session-id");
assert_eq!(tokens[2], "my-sid");
assert_eq!(tokens[3], "-e");
assert_eq!(&tokens[5..], ["--model", "x"]);
assert_eq!(plan.files.len(), 1);
assert!(plan.files[0].exists());
}
#[test]
fn test_prepare_spawn_trailing_valueless_session_id_mints_id() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter
.prepare_spawn(&ctx(tmp.path(), "pi --session-id"))
.unwrap();
let sid = "11111111-1111-1111-1111-111111111111";
assert_eq!(plan.harness_session_id.as_deref(), Some(sid));
let tokens = shell_words::split(&plan.command).unwrap();
assert_eq!(tokens[0], "pi");
assert_eq!(tokens[1], "--session-id");
assert_eq!(tokens[2], sid);
assert_eq!(tokens[3], "-e");
assert_eq!(tokens[4], plan.files[0].to_string_lossy());
}
#[test]
fn test_prepare_spawn_fresh_and_resume_produce_identical_shape() {
let tmp = tempfile::tempdir().unwrap();
let fresh = PiAdapter.prepare_spawn(&ctx(tmp.path(), "pi")).unwrap();
let sid = fresh.harness_session_id.clone().unwrap();
let resumed_base = PiAdapter.resume_command("pi", &sid).unwrap();
let resumed = PiAdapter
.prepare_spawn(&ctx(tmp.path(), &resumed_base))
.unwrap();
assert_eq!(fresh.command, resumed.command);
}
#[test]
fn test_prepare_spawn_skips_for_each_session_selection_flag() {
let tmp = tempfile::tempdir().unwrap();
for command in [
"pi --session /tmp/x.jsonl",
"pi --continue",
"pi -c",
"pi --resume",
"pi -r",
"pi --fork abc",
"pi --no-session",
] {
let plan = PiAdapter.prepare_spawn(&ctx(tmp.path(), command)).unwrap();
assert_eq!(plan.command, command, "expected no-op for {command:?}");
assert!(plan.harness_session_id.is_none());
assert!(plan.files.is_empty());
}
}
#[test]
fn test_prepare_spawn_skips_for_each_known_subcommand() {
let tmp = tempfile::tempdir().unwrap();
for command in [
"pi install some-extension",
"pi remove some-extension",
"pi uninstall some-extension",
"pi update",
"pi list",
"pi config",
"pi auth print-api-key",
] {
let plan = PiAdapter.prepare_spawn(&ctx(tmp.path(), command)).unwrap();
assert_eq!(plan.command, command, "expected no-op for {command:?}");
assert!(plan.harness_session_id.is_none());
assert!(plan.files.is_empty());
}
}
#[test]
fn test_prepare_spawn_known_subcommand_with_env_prefix_is_skipped() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter
.prepare_spawn(&ctx(tmp.path(), "env FOO=bar pi update"))
.unwrap();
assert_eq!(plan.command, "env FOO=bar pi update");
assert!(plan.harness_session_id.is_none());
}
#[test]
fn test_prepare_spawn_subcommand_named_word_after_flags_is_not_a_subcommand() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter
.prepare_spawn(&ctx(tmp.path(), "pi -p update"))
.unwrap();
assert!(plan.harness_session_id.is_some());
}
#[test]
fn test_prepare_spawn_ignores_session_selection_flags_after_double_dash() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter
.prepare_spawn(&ctx(tmp.path(), "pi -- --resume"))
.unwrap();
assert!(plan.harness_session_id.is_some());
let tokens = shell_words::split(&plan.command).unwrap();
assert_eq!(tokens[0], "pi");
assert_eq!(tokens[1], "--session-id");
assert_eq!(&tokens[5..], ["--", "--resume"]);
}
#[test]
fn test_prepare_spawn_ignores_session_id_after_double_dash() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter
.prepare_spawn(&ctx(tmp.path(), "pi -- --session-id foo"))
.unwrap();
let sid = "11111111-1111-1111-1111-111111111111";
assert_eq!(plan.harness_session_id.as_deref(), Some(sid));
let tokens = shell_words::split(&plan.command).unwrap();
assert_eq!(tokens[0], "pi");
assert_eq!(tokens[1], "--session-id");
assert_eq!(tokens[2], sid);
assert_eq!(tokens[3], "-e");
assert_eq!(&tokens[5..], ["--", "--session-id", "foo"]);
}
#[test]
fn test_prepare_spawn_falls_back_on_unparseable_command() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter
.prepare_spawn(&ctx(tmp.path(), "pi \"unterminated"))
.unwrap();
assert_eq!(plan.command, "pi \"unterminated");
assert!(plan.harness_session_id.is_none());
}
#[test]
fn test_prepare_spawn_not_pi_returns_unchanged() {
let tmp = tempfile::tempdir().unwrap();
let plan = PiAdapter.prepare_spawn(&ctx(tmp.path(), "bash")).unwrap();
assert_eq!(plan.command, "bash");
assert!(plan.harness_session_id.is_none());
}
#[test]
fn test_prepare_spawn_falls_back_when_data_dir_cannot_be_created() {
let tmp = tempfile::tempdir().unwrap();
let blocker = tmp.path().join("blocker");
std::fs::write(&blocker, b"not a directory").unwrap();
let plan = PiAdapter.prepare_spawn(&ctx(&blocker, "pi")).unwrap();
assert_eq!(plan.command, "pi");
assert!(plan.harness_session_id.is_none());
assert!(plan.files.is_empty());
}
#[test]
fn test_resume_command_plain() {
let cmd = PiAdapter.resume_command("pi", "sid-1").unwrap();
assert_eq!(cmd, "pi --session-id sid-1");
}
#[test]
fn test_resume_command_preserves_other_args() {
let cmd = PiAdapter.resume_command("pi --model x", "sid-1").unwrap();
assert_eq!(cmd, "pi --session-id sid-1 --model x");
}
#[test]
fn test_resume_command_preserves_print_prompt_arg() {
let cmd = PiAdapter.resume_command("pi -p 'prompt'", "sid-1").unwrap();
assert_eq!(cmd, "pi --session-id sid-1 -p prompt");
}
#[test]
fn test_resume_command_strips_existing_session_id() {
let cmd = PiAdapter
.resume_command("pi --session-id old-sid --model x", "sid-new")
.unwrap();
assert_eq!(cmd, "pi --session-id sid-new --model x");
}
#[test]
fn test_resume_command_preserves_env_prefix() {
let cmd = PiAdapter
.resume_command("env FOO=bar pi -p hi", "sid-1")
.unwrap();
let tokens = shell_words::split(&cmd).unwrap();
assert_eq!(
tokens,
["env", "FOO=bar", "pi", "--session-id", "sid-1", "-p", "hi"]
);
}
#[test]
fn test_resume_command_none_when_session_flag_present() {
assert!(
PiAdapter
.resume_command("pi --session /tmp/x.jsonl", "sid-1")
.is_none()
);
}
#[test]
fn test_resume_command_none_when_continue_flag_present() {
assert!(PiAdapter.resume_command("pi --continue", "sid-1").is_none());
assert!(PiAdapter.resume_command("pi -c", "sid-1").is_none());
}
#[test]
fn test_resume_command_none_when_resume_flag_present() {
assert!(PiAdapter.resume_command("pi --resume", "sid-1").is_none());
assert!(PiAdapter.resume_command("pi -r", "sid-1").is_none());
}
#[test]
fn test_resume_command_none_when_not_pi() {
assert!(PiAdapter.resume_command("bash", "sid-1").is_none());
}
#[test]
fn test_resume_command_none_when_unparseable() {
assert!(
PiAdapter
.resume_command("pi \"unterminated", "sid-1")
.is_none()
);
}
#[test]
fn test_resume_command_none_when_fork_flag_present() {
assert!(PiAdapter.resume_command("pi --fork abc", "sid-1").is_none());
}
#[test]
fn test_resume_command_no_session_is_not_a_conflict() {
assert!(
PiAdapter
.resume_command("pi --no-session", "sid-1")
.is_some()
);
}
#[test]
fn test_parse_event_session_start_new_is_not_resumed() {
let raw = serde_json::json!({
"event": "session_start",
"session_id": "sid-1",
"reason": "new",
});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::SessionStarted {
harness_session_id: Some("sid-1".into()),
resumed: false,
}
);
}
#[test]
fn test_parse_event_session_start_other_reasons_are_resumed() {
for reason in ["startup", "reload", "resume", "fork"] {
let raw = serde_json::json!({
"event": "session_start",
"session_id": "sid-1",
"reason": reason,
});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::SessionStarted {
harness_session_id: Some("sid-1".into()),
resumed: true,
},
"reason {reason:?} should be treated as resumed"
);
}
}
#[test]
fn test_parse_event_agent_start_is_working() {
let raw = serde_json::json!({"event": "agent_start"});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::Working
);
}
#[test]
fn test_parse_event_agent_settled_turn_finished_with_summary() {
let raw = serde_json::json!({
"event": "agent_settled",
"last_assistant_message": "Fixed the bug",
"stop_reason": "stop",
"error": null,
});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::TurnFinished {
summary: Some("Fixed the bug".into())
}
);
}
#[test]
fn test_parse_event_agent_settled_no_summary() {
let raw = serde_json::json!({"event": "agent_settled"});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::TurnFinished { summary: None }
);
}
#[test]
fn test_parse_event_agent_settled_failed_rate_limited() {
for error in [
"rate limit exceeded",
"ratelimit hit",
"rate-limit",
"rate_limit_error",
"429",
"HTTP 429 received",
"Too Many Requests",
] {
let raw = serde_json::json!({"event": "agent_settled", "error": error});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::Failed {
error: error.into(),
rate_limited: true,
},
"error {error:?} should be classified as rate-limited"
);
}
}
#[test]
fn test_parse_event_agent_settled_failed_not_rate_limited() {
let raw = serde_json::json!({"event": "agent_settled", "error": "connection reset"});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::Failed {
error: "connection reset".into(),
rate_limited: false,
}
);
}
#[test]
fn test_parse_event_ui_prompt_start_confirm_is_permission() {
let raw = serde_json::json!({"event": "ui_prompt_start", "kind": "confirm"});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::NeedsInput {
reason: NeedsInputReason::Permission
}
);
}
#[test]
fn test_parse_event_ui_prompt_start_other_kinds_are_question() {
for kind in ["select", "input", "editor", "custom"] {
let raw = serde_json::json!({"event": "ui_prompt_start", "kind": kind});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::NeedsInput {
reason: NeedsInputReason::Question
},
"kind {kind:?} should map to Question"
);
}
}
#[test]
fn test_parse_event_ui_prompt_start_missing_kind_is_question() {
let raw = serde_json::json!({"event": "ui_prompt_start"});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::NeedsInput {
reason: NeedsInputReason::Question
}
);
}
#[test]
fn test_parse_event_ui_prompt_end_is_working() {
let raw = serde_json::json!({"event": "ui_prompt_end"});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::Working
);
}
#[test]
fn test_parse_event_session_shutdown_quit_is_session_ended() {
let raw = serde_json::json!({"event": "session_shutdown", "reason": "quit"});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::SessionEnded {
reason: Some("quit".into())
}
);
}
#[test]
fn test_parse_event_session_shutdown_non_quit_reasons_are_ignored() {
for reason in ["reload", "new", "resume", "fork"] {
let raw = serde_json::json!({"event": "session_shutdown", "reason": reason});
assert!(
PiAdapter.parse_event(&raw).unwrap().is_none(),
"reason {reason:?} should be ignored"
);
}
}
#[test]
fn test_parse_event_session_shutdown_missing_reason_is_ignored() {
let raw = serde_json::json!({"event": "session_shutdown"});
assert!(PiAdapter.parse_event(&raw).unwrap().is_none());
}
#[test]
fn test_parse_event_unknown_event_is_none() {
let raw = serde_json::json!({"event": "turn_start"});
assert!(PiAdapter.parse_event(&raw).unwrap().is_none());
}
#[test]
fn test_parse_event_missing_event_field_is_none() {
let raw = serde_json::json!({});
assert!(PiAdapter.parse_event(&raw).unwrap().is_none());
}
#[test]
fn test_parse_event_agent_settled_429_is_not_rate_limited_when_embedded_in_a_longer_number() {
let raw = serde_json::json!({"event": "agent_settled", "error": "request 14290 done"});
assert_eq!(
PiAdapter.parse_event(&raw).unwrap().unwrap(),
HarnessEvent::Failed {
error: "request 14290 done".into(),
rate_limited: false,
}
);
}
#[test]
fn test_parse_event_round_trips_template_payload_field_names() {
let session_start = serde_json::json!({
"event": "session_start",
"session_id": "sid-1",
"session_file": "/home/u/.pi/agent/sessions/--repo--/x.jsonl",
"cwd": "/repo",
"reason": "new",
"previous_session_file": null,
});
assert_eq!(
PiAdapter.parse_event(&session_start).unwrap().unwrap(),
HarnessEvent::SessionStarted {
harness_session_id: Some("sid-1".into()),
resumed: false,
}
);
let agent_start = serde_json::json!({
"event": "agent_start",
"session_id": "sid-1",
"session_file": "/home/u/.pi/agent/sessions/--repo--/x.jsonl",
"cwd": "/repo",
});
assert_eq!(
PiAdapter.parse_event(&agent_start).unwrap().unwrap(),
HarnessEvent::Working
);
let agent_settled_ok = serde_json::json!({
"event": "agent_settled",
"session_id": "sid-1",
"session_file": "/home/u/.pi/agent/sessions/--repo--/x.jsonl",
"cwd": "/repo",
"last_assistant_message": "Fixed the bug",
"stop_reason": "stop",
"error": null,
});
assert_eq!(
PiAdapter.parse_event(&agent_settled_ok).unwrap().unwrap(),
HarnessEvent::TurnFinished {
summary: Some("Fixed the bug".into())
}
);
let agent_settled_err = serde_json::json!({
"event": "agent_settled",
"session_id": "sid-1",
"session_file": "/home/u/.pi/agent/sessions/--repo--/x.jsonl",
"cwd": "/repo",
"last_assistant_message": null,
"stop_reason": "error",
"error": "rate limited: 429",
});
assert_eq!(
PiAdapter.parse_event(&agent_settled_err).unwrap().unwrap(),
HarnessEvent::Failed {
error: "rate limited: 429".into(),
rate_limited: true,
}
);
let ui_prompt_start = serde_json::json!({
"event": "ui_prompt_start",
"session_id": "sid-1",
"session_file": "/home/u/.pi/agent/sessions/--repo--/x.jsonl",
"cwd": "/repo",
"kind": "confirm",
"title": "Allow this action?",
});
assert_eq!(
PiAdapter.parse_event(&ui_prompt_start).unwrap().unwrap(),
HarnessEvent::NeedsInput {
reason: NeedsInputReason::Permission
}
);
let ui_prompt_end = serde_json::json!({
"event": "ui_prompt_end",
"session_id": "sid-1",
"session_file": "/home/u/.pi/agent/sessions/--repo--/x.jsonl",
"cwd": "/repo",
"kind": "confirm",
"title": "Allow this action?",
});
assert_eq!(
PiAdapter.parse_event(&ui_prompt_end).unwrap().unwrap(),
HarnessEvent::Working
);
let session_shutdown = serde_json::json!({
"event": "session_shutdown",
"session_id": "sid-1",
"session_file": "/home/u/.pi/agent/sessions/--repo--/x.jsonl",
"cwd": "/repo",
"reason": "quit",
});
assert_eq!(
PiAdapter.parse_event(&session_shutdown).unwrap().unwrap(),
HarnessEvent::SessionEnded {
reason: Some("quit".into())
}
);
}
#[test]
fn test_escape_ts_string_escapes_backslash_and_quote() {
assert_eq!(escape_ts_string(r#"C:\weird"path"#), r#"C:\\weird\"path"#);
assert_eq!(
escape_ts_string("/usr/local/bin/pulpo"),
"/usr/local/bin/pulpo"
);
}
#[test]
fn test_has_flag_false_for_clean_command() {
let tokens = vec!["pi".to_owned(), "-p".to_owned(), "hi".to_owned()];
assert!(!has_flag(&tokens, SESSION_SELECTION_FLAGS));
}
#[test]
fn test_has_flag_ignores_tokens_after_double_dash() {
let tokens = vec!["pi".to_owned(), "--".to_owned(), "--resume".to_owned()];
assert!(!has_flag(&tokens, SESSION_SELECTION_FLAGS));
}
#[test]
fn test_has_flag_true_before_double_dash() {
let tokens = vec![
"pi".to_owned(),
"--resume".to_owned(),
"--".to_owned(),
"hello".to_owned(),
];
assert!(has_flag(&tokens, SESSION_SELECTION_FLAGS));
}
#[test]
fn test_flag_scan_region_no_double_dash_returns_everything() {
let tokens = vec!["pi".to_owned(), "-p".to_owned(), "hi".to_owned()];
assert_eq!(flag_scan_region(&tokens), tokens.as_slice());
}
#[test]
fn test_flag_scan_region_stops_before_double_dash() {
let tokens = vec![
"pi".to_owned(),
"--model".to_owned(),
"x".to_owned(),
"--".to_owned(),
"--resume".to_owned(),
];
assert_eq!(flag_scan_region(&tokens), &tokens[..3]);
}
#[test]
fn test_contains_standalone_429_bare_number() {
assert!(contains_standalone_429("429"));
assert!(contains_standalone_429("http 429 received"));
}
#[test]
fn test_contains_standalone_429_not_embedded_in_longer_number() {
assert!(!contains_standalone_429("request 14290 done"));
assert!(!contains_standalone_429("code 94290"));
assert!(!contains_standalone_429("id 4297"));
}
#[test]
fn test_pi_token_index_not_found() {
let tokens = vec!["bash".to_owned()];
assert!(pi_token_index(&tokens).is_none());
}
#[test]
fn test_pi_token_index_env_with_no_command_after() {
let tokens = vec!["env".to_owned(), "FOO=bar".to_owned()];
assert!(pi_token_index(&tokens).is_none());
}
#[test]
fn test_strip_flag_with_value_removes_flag_and_value() {
let mut tokens = vec![
"pi".to_owned(),
"--session-id".to_owned(),
"old".to_owned(),
"--model".to_owned(),
"x".to_owned(),
];
strip_flag_with_value(&mut tokens, "--session-id");
assert_eq!(tokens, vec!["pi", "--model", "x"]);
}
#[test]
fn test_strip_flag_with_value_noop_when_absent() {
let mut tokens = vec!["pi".to_owned(), "--model".to_owned(), "x".to_owned()];
strip_flag_with_value(&mut tokens, "--session-id");
assert_eq!(tokens, vec!["pi", "--model", "x"]);
}
}