use std::path::{Path, PathBuf};
use std::sync::atomic::Ordering;
use std::sync::{Arc, Mutex};
use rpi_ai::Provider;
use rpi_harness::agent_harness::AgentHarness;
use rpi_harness::context_files::{format_project_context, load_project_context_files};
use rpi_harness::session::memory::{InMemorySessionStorage, SystemClock};
use rpi_harness::session::session::DefaultIdGenerator;
use rpi_harness::session::types::SessionMetadata;
use rpi_harness::session::Session;
use rpi_harness::system_prompt::compose_system_prompt;
use rpi_harness::types::{
AgentHarnessOptions, AgentHarnessResources, CompactionSettings, DrivingMode,
HarnessToolExecution, HarnessTool, RetryPolicy, ToolReplay,
};
use rpi_tools::{
create_bash_tool, create_edit_tool, create_find_tool, create_grep_tool, create_ls_tool,
create_read_tool, create_write_tool, ExecutionToolContext, MutationQueueRegistry,
OsExecutionEnv,
};
use crate::args::Args;
use crate::provider::ResolvedModel;
use crate::resource_dirs::{
discover_append_system_prompt_file, discover_system_prompt_file, global_dir,
load_prompt_templates_with_precedence, load_skills_with_precedence, project_dir,
prompt_template_dirs, skill_dirs,
};
use rpi_extensions::{
ExtensionEmitter, ExtensionSession, NullDiagnostics, PluginDiagnostics, PluginToolAdapter,
TeeEmitter, emit_resources_discover, load_session,
};
const EXTENSIONS_SUBDIR: &str = "extensions";
pub const BUILTIN_TOOL_NAMES: &[&str] = &["read", "bash", "edit", "write", "grep", "find", "ls"];
pub fn default_system_prompt(cwd: &str) -> String {
format!(
"You are an expert coding assistant operating inside pi, a coding agent harness. \
You help users by reading files, executing commands, editing code, and writing new files.
Available tools:
- read — Read file contents
- bash — Execute shell commands
- edit — Find/replace edits to existing files
- write — Create or overwrite files
- grep — Search file contents for a pattern
- find — Search for files by glob pattern
- ls — List directory contents
Guidelines:
- Be concise in your responses
- Show file paths clearly when working with files
- Prefer the smallest change that solves the problem
Current working directory: {cwd}"
)
}
#[derive(Debug, Clone)]
pub enum SessionSelection {
Ephemeral,
New { dir: PathBuf, name: Option<String> },
Latest,
ById { id: String },
}
pub fn select_session(args: &Args, cwd: &Path) -> SessionSelection {
if args.no_session {
return SessionSelection::Ephemeral;
}
if args.continue_session || args.resume {
return SessionSelection::Latest;
}
if let Some(s) = &args.session {
return SessionSelection::ById { id: s.clone() };
}
let dir = args
.session_dir
.clone()
.unwrap_or_else(|| default_session_dir(cwd));
SessionSelection::New { dir, name: args.name.clone() }
}
pub fn default_session_dir(cwd: &Path) -> PathBuf {
cwd.join(".pi").join("sessions")
}
pub async fn build(
resolved: &ResolvedModel,
args: &Args,
cwd: &Path,
) -> Result<
(
AgentHarness,
tokio::sync::broadcast::Receiver<rpi_agent::AgentEvent>,
ReloadContext,
),
BuildError,
> {
let cwd_str = cwd.to_string_lossy().to_string();
let runtime = tokio::runtime::Handle::try_current()
.map_err(|e| BuildError::HarnessCreate(format!("no tokio runtime for action bridge: {e}")))?;
let catalog = crate::provider::available_catalog(resolved);
let (action_host, harness_cell) = crate::extensions_actions::HarnessActionHost::new_empty(
catalog.clone(),
cwd.to_path_buf(),
runtime.clone(),
);
let host_arc: Arc<dyn rpi_extensions::RuntimeActionHost> = Arc::new(action_host);
let reload_mailbox = rpi_extensions::ReloadMailbox::new();
let action_bridge = rpi_extensions::ActionBridge::with_reload(
runtime.clone(),
host_arc,
rpi_extensions::reload_callback_from_mailbox(reload_mailbox.clone()),
);
let env = Arc::new(OsExecutionEnv::with_cwd(cwd.to_path_buf()));
let env_dyn: Arc<dyn rpi_tools::ExecutionEnv> = env.clone();
let mut_env: Arc<dyn rpi_tools::MutatingEnv> = env.clone();
let _registry = Arc::new(MutationQueueRegistry::new());
let ctx = ExecutionToolContext::new(env_dyn.clone(), Some(mut_env));
let tools = build_tools(&ctx, args);
let mut tools = tools;
let extension_session = if args.no_extensions {
ExtensionSession::none()
} else {
load_extensions(args, cwd, Some(Arc::clone(&action_bridge)))
};
if args.verbose {
if let Some(s) = extension_session.summary() {
eprintln!("extensions: {s}");
}
report_deferred_renderers(&extension_session);
}
merge_extension_tools(&mut tools, &extension_session, args);
let active = active_tool_names(&tools, args);
let selection = select_session(args, cwd);
let session = build_session(&selection, &cwd_str).await?;
let base_prompt = match args.system_prompt.as_deref() {
Some(explicit) => explicit.to_string(),
None => match discover_system_prompt_file(cwd) {
Some(path) => std::fs::read_to_string(&path).unwrap_or_else(|_| {
default_system_prompt(&cwd_str)
}),
None => default_system_prompt(&cwd_str),
},
};
let mut append_texts: Vec<String> = Vec::new();
for extra in &args.append_system_prompt {
let text = read_append_target(extra).unwrap_or_else(|| extra.clone());
append_texts.push(text);
}
if args.append_system_prompt.is_empty() {
if let Some(path) = discover_append_system_prompt_file(cwd) {
if let Ok(text) = std::fs::read_to_string(&path) {
append_texts.push(text);
}
}
}
let append_join = if append_texts.is_empty() {
None
} else {
Some(append_texts.join("\n\n"))
};
let agent_dir = crate::config::agent_dir().ok();
let discovered = extension_session
.snapshot_arc()
.map(|snap| emit_resources_discover(&cwd_str, "startup", &snap))
.unwrap_or_default();
let mut skills: Vec<rpi_harness::types::Skill> = Vec::new();
let mut skill_diags: Vec<rpi_harness::skills::SkillDiagnostic> = Vec::new();
if !args.no_skills {
let mut dirs = skill_dirs(cwd);
dirs.extend(discovered.skill_paths.iter().map(PathBuf::from));
let result = load_skills_with_precedence(&env_dyn, &dirs).await;
skills = result.skills;
skill_diags = result.diagnostics;
}
let mut prompt_templates: Vec<rpi_harness::types::PromptTemplate> = Vec::new();
let mut prompt_diags: Vec<rpi_harness::prompt_templates::PromptTemplateDiagnostic> = Vec::new();
if !args.no_prompt_templates {
let mut dirs = prompt_template_dirs(cwd);
dirs.extend(discovered.prompt_paths.iter().map(PathBuf::from));
let result = load_prompt_templates_with_precedence(&env_dyn, &dirs).await;
prompt_templates = result.prompt_templates;
prompt_diags = result.diagnostics;
}
let context_block = if args.no_context_files {
String::new()
} else {
let agent_dir_path = agent_dir.clone().unwrap_or_else(|| cwd.to_path_buf());
let files = load_project_context_files(&env_dyn, cwd, &agent_dir_path).await;
format_project_context(&files)
};
if args.verbose {
for d in &skill_diags {
eprintln!("warning: skill {} ({}): {}", d.path, d.code.as_str(), d.message);
}
for d in &prompt_diags {
eprintln!(
"warning: prompt template {} ({}): {}",
d.path,
d.code.as_str(),
d.message
);
}
}
let system_prompt = compose_system_prompt(
Some(&base_prompt),
&[], if context_block.is_empty() { None } else { Some(&context_block) },
append_join.as_deref(),
);
if args.debug_system_prompt {
eprintln!("=== --debug-system-prompt ===");
let base_src = if args.system_prompt.is_some() {
"--system-prompt"
} else if discover_system_prompt_file(cwd).is_some() {
"SYSTEM.md"
} else {
"default"
};
eprintln!("[base source: {base_src}]");
eprintln!("--- base ---\n{base_prompt}");
if let Some(append) = append_join.as_deref() {
eprintln!("--- append ---\n{append}");
} else {
eprintln!("--- append: (none) ---");
}
if context_block.is_empty() {
eprintln!("--- context: (none) ---");
} else {
eprintln!("--- context ---{context_block}");
}
let visible_skills = skills
.iter()
.filter(|s| s.disable_model_invocation != Some(true))
.count();
eprintln!(
"--- skills: {} loaded ({} model-visible, {} hidden) ---",
skills.len(),
visible_skills,
skills.len() - visible_skills
);
for s in &skills {
let hidden = if s.disable_model_invocation == Some(true) { " [hidden]" } else { "" };
eprintln!(" {}{hidden} — {}", s.name, s.description);
}
eprintln!("--- prompt templates: {} ---", prompt_templates.len());
for t in &prompt_templates {
eprintln!(" /{}", t.name);
}
eprintln!(
"--- discovered via resources_discover: {} skill(s), {} prompt(s), {} theme(s) (ignored) ---",
discovered.skill_paths.len(),
discovered.prompt_paths.len(),
discovered.theme_paths.len(),
);
for p in &discovered.skill_paths {
eprintln!(" skill: {p}");
}
for p in &discovered.prompt_paths {
eprintln!(" prompt: {p}");
}
eprintln!(
"--- final composed base+append+context (skills listing added by harness) ---\n{system_prompt}"
);
eprintln!("=== end --debug-system-prompt ===");
}
let (broadcast, event_rx) = rpi_agent::events::BroadcastEmitter::new(256);
let broadcast_emitter: Arc<dyn rpi_agent::AgentEmitter> = Arc::new(broadcast);
let broadcast_for_context: Arc<dyn rpi_agent::AgentEmitter> = Arc::clone(&broadcast_emitter);
let emitter: Arc<dyn rpi_agent::AgentEmitter> =
match extension_session.snapshot_arc() {
Some(snapshot) => {
let ext = ExtensionEmitter::new(snapshot, extension_session.keepalive());
Arc::new(TeeEmitter::new(vec![
broadcast_emitter,
Arc::new(ext),
]))
}
None => broadcast_emitter,
};
let options = AgentHarnessOptions {
model: resolved.model.clone(),
thinking_level: resolved.thinking_level,
active_tool_names: active,
tools,
system_prompt: Some(system_prompt),
resources: AgentHarnessResources {
skills: if skills.is_empty() { None } else { Some(skills) },
prompt_templates: if prompt_templates.is_empty() {
None
} else {
Some(prompt_templates)
},
},
allow_existing_session: matches!(
selection,
SessionSelection::Latest | SessionSelection::ById { .. }
),
stream_options: Default::default(),
retry: RetryPolicy::default(),
compaction: CompactionSettings::default(),
steering_mode: Default::default(),
follow_up_mode: Default::default(),
tool_execution: HarnessToolExecution::default(),
drive: DrivingMode::default(),
session,
models: build_models_with_extensions(resolved, &extension_session, runtime.clone()),
to_provider_messages: None,
entry_projectors: Default::default(),
agent_emitter: Some(emitter),
before_tool_call: None,
after_tool_call: None,
transform_context: None,
entry_transforms: Vec::new(),
provider_hooks: rpi_extensions::ExtensionProviderHooks::from_session(
&extension_session,
)
.map(|h| Arc::new(h) as Arc<dyn rpi_ai::ProviderHooks>),
};
let harness = match AgentHarness::create(options).await {
Ok(h) => {
crate::extensions_actions::HarnessActionHost::set_harness(
&harness_cell,
Arc::new(h.clone()),
);
h
}
Err(e) => return Err(BuildError::HarnessCreate(e.to_string())),
};
let reload_context = ReloadContext {
extension_session: Arc::new(Mutex::new(extension_session)),
action_bridge: Arc::new(Mutex::new(Some(Arc::clone(&action_bridge)))),
catalog,
gateway: resolved.provider.clone(),
runtime: runtime.clone(),
cwd: cwd.to_path_buf(),
args: args.clone(),
resolved_model: resolved.model.clone(),
broadcast: broadcast_for_context,
mailbox: reload_mailbox,
};
Ok((harness, event_rx, reload_context))
}
pub type ExtensionSessionCell = Arc<Mutex<ExtensionSession>>;
pub type ActionBridgeCell = Arc<Mutex<Option<Arc<rpi_extensions::ActionBridge>>>>;
#[derive(Clone)]
pub struct ReloadContext {
pub extension_session: ExtensionSessionCell,
pub action_bridge: ActionBridgeCell,
pub catalog: Vec<rpi_ai::Model>,
pub gateway: Arc<dyn Provider>,
pub runtime: tokio::runtime::Handle,
pub cwd: PathBuf,
pub args: Args,
pub resolved_model: rpi_ai::Model,
pub broadcast: Arc<dyn rpi_agent::AgentEmitter>,
pub mailbox: rpi_extensions::ReloadMailbox,
}
pub struct ReloadOutcome {
pub summary: String,
pub had_warnings: bool,
}
pub async fn reload_extension_resources(
harness: &AgentHarness,
ctx: &ReloadContext,
) -> ReloadOutcome {
let cwd_str = ctx.cwd.to_string_lossy().to_string();
let mut warnings = false;
let old_bridge = ctx.action_bridge.lock().unwrap().clone();
let host: Arc<dyn rpi_extensions::RuntimeActionHost> = match &old_bridge {
Some(b) => b.clone_host(),
None => {
let (action_host, _cell) =
crate::extensions_actions::HarnessActionHost::new_empty(
ctx.catalog.clone(),
ctx.cwd.clone(),
ctx.runtime.clone(),
);
crate::extensions_actions::HarnessActionHost::set_harness(
&_cell,
Arc::new(harness.clone()),
);
Arc::new(action_host)
}
};
let reload_cb = rpi_extensions::reload_callback_from_mailbox(ctx.mailbox.clone());
let fresh_bridge =
rpi_extensions::ActionBridge::with_reload(ctx.runtime.clone(), host, reload_cb);
let extension_session = if ctx.args.no_extensions {
rpi_extensions::ExtensionSession::none()
} else {
load_extensions(&ctx.args, &ctx.cwd, Some(Arc::clone(&fresh_bridge)))
};
if extension_session.is_empty() && !ctx.args.no_extensions {
}
if ctx.args.verbose {
if let Some(s) = extension_session.summary() {
eprintln!("reload: {s}");
}
report_deferred_renderers(&extension_session);
}
{
let mut session_guard = ctx.extension_session.lock().unwrap();
let old_session =
std::mem::replace(&mut *session_guard, rpi_extensions::ExtensionSession::none());
if let Some(old_snap) = old_session.snapshot_arc() {
old_snap.active_flag().store(false, Ordering::SeqCst);
}
}
if let Some(old_b) = old_bridge {
old_b.invalidate();
}
*ctx.action_bridge.lock().unwrap() = Some(Arc::clone(&fresh_bridge));
*ctx.extension_session.lock().unwrap() = extension_session.clone();
let discovered = extension_session
.snapshot_arc()
.map(|snap| rpi_extensions::emit_resources_discover(&cwd_str, "reload", &snap))
.unwrap_or_default();
let env = Arc::new(rpi_tools::OsExecutionEnv::with_cwd(ctx.cwd.clone()));
let env_dyn: Arc<dyn rpi_tools::ExecutionEnv> = env.clone();
let mut skills: Vec<rpi_harness::types::Skill> = Vec::new();
let mut skill_diags: Vec<rpi_harness::skills::SkillDiagnostic> = Vec::new();
if !ctx.args.no_skills {
let mut dirs = skill_dirs(&ctx.cwd);
dirs.extend(discovered.skill_paths.iter().map(PathBuf::from));
let result = load_skills_with_precedence(&env_dyn, &dirs).await;
skills = result.skills;
skill_diags = result.diagnostics;
}
let mut prompt_templates: Vec<rpi_harness::types::PromptTemplate> = Vec::new();
let mut prompt_diags: Vec<rpi_harness::prompt_templates::PromptTemplateDiagnostic> = Vec::new();
if !ctx.args.no_prompt_templates {
let mut dirs = prompt_template_dirs(&ctx.cwd);
dirs.extend(discovered.prompt_paths.iter().map(PathBuf::from));
let result = load_prompt_templates_with_precedence(&env_dyn, &dirs).await;
prompt_templates = result.prompt_templates;
prompt_diags = result.diagnostics;
}
let context_block = if ctx.args.no_context_files {
String::new()
} else {
let agent_dir = crate::config::agent_dir().ok();
let agent_dir_path = agent_dir.unwrap_or_else(|| ctx.cwd.clone());
let files = load_project_context_files(&env_dyn, &ctx.cwd, &agent_dir_path).await;
format_project_context(&files)
};
if !skill_diags.is_empty() || !prompt_diags.is_empty() {
warnings = true;
if ctx.args.verbose {
for d in &skill_diags {
eprintln!("warning: skill {} ({}): {}", d.path, d.code.as_str(), d.message);
}
for d in &prompt_diags {
eprintln!(
"warning: prompt template {} ({}): {}",
d.path,
d.code.as_str(),
d.message
);
}
}
}
let base_prompt = match ctx.args.system_prompt.as_deref() {
Some(explicit) => explicit.to_string(),
None => match discover_system_prompt_file(&ctx.cwd) {
Some(path) => std::fs::read_to_string(&path)
.unwrap_or_else(|_| default_system_prompt(&cwd_str)),
None => default_system_prompt(&cwd_str),
},
};
let mut append_texts: Vec<String> = Vec::new();
for extra in &ctx.args.append_system_prompt {
let text = read_append_target(extra).unwrap_or_else(|| extra.clone());
append_texts.push(text);
}
if ctx.args.append_system_prompt.is_empty() {
if let Some(path) = discover_append_system_prompt_file(&ctx.cwd) {
if let Ok(text) = std::fs::read_to_string(&path) {
append_texts.push(text);
}
}
}
let append_join = if append_texts.is_empty() {
None
} else {
Some(append_texts.join("\n\n"))
};
let system_prompt = compose_system_prompt(
Some(&base_prompt),
&[],
if context_block.is_empty() { None } else { Some(&context_block) },
append_join.as_deref(),
);
let emitter: Arc<dyn rpi_agent::AgentEmitter> =
match extension_session.snapshot_arc() {
Some(snapshot) => {
let ext = ExtensionEmitter::new(snapshot, extension_session.keepalive());
Arc::new(TeeEmitter::new(vec![
ctx.broadcast.clone(),
Arc::new(ext),
]))
}
None => ctx.broadcast.clone(),
};
let resources = AgentHarnessResources {
skills: if skills.is_empty() { None } else { Some(skills.clone()) },
prompt_templates: if prompt_templates.is_empty() {
None
} else {
Some(prompt_templates.clone())
},
};
let _ = harness.set_system_prompt(Some(system_prompt)).await;
let _ = harness.set_resources(resources).await;
let _ = harness.set_agent_emitter(Some(emitter)).await;
let _ = harness
.set_models(build_models_with_extensions_for_reload(
&ctx.gateway,
&extension_session,
ctx.runtime.clone(),
))
.await;
let _ = harness
.set_provider_hooks(
rpi_extensions::ExtensionProviderHooks::from_session(&extension_session)
.map(|h| Arc::new(h) as Arc<dyn rpi_ai::ProviderHooks>),
)
.await;
let mut_env: Arc<dyn rpi_tools::MutatingEnv> = env.clone();
let tool_ctx = rpi_tools::ExecutionToolContext::new(env_dyn.clone(), Some(mut_env));
let mut tools = build_tools(&tool_ctx, &ctx.args);
merge_extension_tools(&mut tools, &extension_session, &ctx.args);
let active = active_tool_names(&tools, &ctx.args);
let _ = harness.set_tools(tools, Some(active)).await;
let summary = format!(
"Reloaded {} plugin(s), {} skill(s), {} prompt(s).",
extension_session.loaded_paths().len(),
skills.len(),
prompt_templates.len(),
);
ReloadOutcome { summary, had_warnings: warnings }
}
fn build_models_with_extensions_for_reload(
gateway: &Arc<dyn Provider>,
extension_session: &ExtensionSession,
runtime: tokio::runtime::Handle,
) -> Vec<Arc<dyn Provider>> {
let mut models: Vec<Arc<dyn Provider>> = vec![gateway.clone()];
let pluggable = rpi_extensions::PluggableProvider::from_session(extension_session, runtime);
models.extend(pluggable);
models
}
fn report_deferred_renderers(session: &ExtensionSession) {
let Some(snap) = session.snapshot_arc() else {
return;
};
let all = snap.renderers();
let markdown = all
.iter()
.filter(|r| r.kind == rpi_extensions::RegisteredRendererKind::Markdown)
.count();
let message = all
.iter()
.filter(|r| r.kind == rpi_extensions::RegisteredRendererKind::Message)
.count();
let entry = all
.iter()
.filter(|r| r.kind == rpi_extensions::RegisteredRendererKind::Entry)
.count();
if markdown + message + entry == 0 {
return;
}
eprintln!(
"renderers: {} markdown-transform (active), {} message-render (deferred), {} entry-render (deferred)",
markdown, message, entry
);
}
#[derive(Debug, thiserror::Error)]
pub enum BuildError {
#[error("Could not create the session directory: {0}")]
SessionDir(String),
#[error("No session found for {requested} in {dir}. Start a fresh session instead (drop --continue/--resume/--session).")]
SessionNotFound { requested: String, dir: String },
#[error("Could not build the harness: {0}")]
HarnessCreate(String),
}
fn build_models_with_extensions(
resolved: &ResolvedModel,
extension_session: &ExtensionSession,
runtime: tokio::runtime::Handle,
) -> Vec<Arc<dyn Provider>> {
let mut models: Vec<Arc<dyn Provider>> = vec![resolved.provider.clone() as Arc<dyn Provider>];
let pluggable = rpi_extensions::PluggableProvider::from_session(extension_session, runtime);
models.extend(pluggable);
models
}
fn load_extensions(
args: &Args,
cwd: &Path,
action_bridge: Option<Arc<rpi_extensions::ActionBridge>>,
) -> ExtensionSession {
let mut dirs = vec![project_dir(cwd, EXTENSIONS_SUBDIR)];
if let Some(g) = global_dir(EXTENSIONS_SUBDIR) {
dirs.push(g);
}
dirs.extend(args.extensions_dir.iter().cloned());
let diagnostics: Arc<dyn PluginDiagnostics> = Arc::new(NullDiagnostics);
load_session(&dirs, diagnostics, action_bridge)
}
fn merge_extension_tools(tools: &mut Vec<HarnessTool>, session: &ExtensionSession, args: &Args) {
let Some(snapshot) = session.snapshot() else { return };
for et in snapshot.tools() {
let name = &et.tool.name;
if let Some(allow) = &args.tools {
if !allow.iter().any(|a| a == name) {
continue;
}
}
if let Some(deny) = &args.exclude_tools {
if deny.iter().any(|d| d == name) {
continue;
}
}
let adapter = PluginToolAdapter::new(et.tool.clone(), et.handle(), session.keepalive());
let harness_tool = HarnessTool::new(Arc::new(adapter));
match tools.iter_mut().find(|t| t.tool.schema().name == *name) {
Some(slot) => *slot = harness_tool,
None => tools.push(harness_tool),
}
}
}
fn build_tools(ctx: &ExecutionToolContext, args: &Args) -> Vec<HarnessTool> {
if args.no_tools {
return Vec::new();
}
let mut all: Vec<(&'static str, HarnessTool)> = vec![
("read", HarnessTool::new(create_read_tool(ctx, None))),
("bash", HarnessTool::new(create_bash_tool(ctx, None))),
("edit", HarnessTool::new(create_edit_tool(ctx))),
("write", HarnessTool::new(create_write_tool(ctx))),
("grep", HarnessTool::new(create_grep_tool(ctx, None))),
("find", HarnessTool::new(create_find_tool(ctx, None))),
("ls", HarnessTool::new(create_ls_tool(ctx, None))),
];
if args.no_builtin_tools {
all.clear();
}
if let Some(allow) = &args.tools {
all.retain(|(name, _)| allow.iter().any(|a| a == name));
}
if let Some(deny) = &args.exclude_tools {
all.retain(|(name, _)| !deny.iter().any(|d| d == name));
}
all.into_iter().map(|(_, t)| t.with_replay(ToolReplay::Safe)).collect()
}
fn active_tool_names(tools: &[HarnessTool], args: &Args) -> Vec<String> {
if args.no_tools {
return Vec::new();
}
if let Some(allow) = &args.tools {
let names: Vec<String> = tools.iter().map(|t| t.tool.schema().name.clone()).collect();
return allow.iter().filter(|a| names.iter().any(|n| n == *a)).cloned().collect();
}
tools.iter().map(|t| t.tool.schema().name.clone()).collect()
}
async fn build_session(selection: &SessionSelection, cwd: &str) -> Result<Session, BuildError> {
match selection {
SessionSelection::Ephemeral => Ok(ephemeral_session()),
SessionSelection::New { dir, .. } => {
std::fs::create_dir_all(dir)
.map_err(|e| BuildError::SessionDir(format!("{}: {e}", dir.display())))?;
let session = create_jsonl_session(dir, cwd)
.await
.map_err(|e| BuildError::SessionDir(format!("{}: {e}", dir.display())))?;
Ok(session)
}
SessionSelection::Latest | SessionSelection::ById { .. } => {
restore_session(selection, cwd).await
}
}
}
async fn restore_session(selection: &SessionSelection, cwd: &str) -> Result<Session, BuildError> {
match selection {
SessionSelection::Latest => {
let metas = list_session_metadata(cwd).await?;
let Some(meta) = metas.first() else {
return Err(BuildError::SessionNotFound {
requested: "the most recent session".to_string(),
dir: default_session_dir(Path::new(cwd)).display().to_string(),
});
};
open_session(meta, cwd).await
}
SessionSelection::ById { id } => {
open_session_by_id(id, cwd).await.map_err(|e| match e {
OpenError::NotFound { requested } => BuildError::SessionNotFound {
requested,
dir: default_session_dir(Path::new(cwd)).display().to_string(),
},
OpenError::Other(msg) => BuildError::SessionDir(msg),
})
}
_ => unreachable!("restore_session only called for Latest/ById"),
}
}
pub enum OpenError {
NotFound { requested: String },
Other(String),
}
impl std::fmt::Display for OpenError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
OpenError::NotFound { requested } => write!(f, "no session matches {requested}"),
OpenError::Other(msg) => write!(f, "{msg}"),
}
}
}
pub async fn list_session_metadata(cwd: &str) -> Result<Vec<rpi_harness::session::jsonl::JsonlSessionMetadata>, BuildError> {
use rpi_harness::session::jsonl::{JsonlSessionListOptions, JsonlSessionRepo, JsonlSessionRepoOptions};
use rpi_tools::FileSystem;
let dir = default_session_dir(Path::new(cwd));
let env = Arc::new(OsExecutionEnv::with_cwd(PathBuf::from(cwd)));
let fs: Arc<dyn FileSystem> = env.clone();
let repo = JsonlSessionRepo::with_env_cwd(JsonlSessionRepoOptions {
fs: fs.clone(),
sessions_root: dir.to_string_lossy().into_owned(),
clock: Arc::new(SystemClock),
ids: Arc::new(DefaultIdGenerator::new()),
});
repo.list_typed(&JsonlSessionListOptions::default())
.await
.map_err(|e| BuildError::SessionDir(format!("list sessions: {e}")))
}
pub async fn open_session_by_id(id: &str, cwd: &str) -> Result<Session, OpenError> {
let metas = list_session_metadata(cwd)
.await
.map_err(|e| OpenError::Other(e.to_string()))?;
let Some(meta) = metas
.iter()
.find(|m| m.id == id || m.path.contains(id) || id.contains(&m.id))
else {
return Err(OpenError::NotFound { requested: format!("session {id}") });
};
open_session(meta, cwd)
.await
.map_err(|e| OpenError::Other(e.to_string()))
}
pub(crate) async fn fork_session_storage(
harness: &AgentHarness,
cwd: &str,
) -> Result<Session, String> {
use rpi_harness::session::jsonl::{JsonlSessionRepo, JsonlSessionRepoOptions};
use rpi_tools::FileSystem;
let dir = default_session_dir(Path::new(cwd));
let env = Arc::new(OsExecutionEnv::with_cwd(PathBuf::from(cwd)));
let fs: Arc<dyn FileSystem> = env.clone();
let repo = JsonlSessionRepo::with_env_cwd(JsonlSessionRepoOptions {
fs,
sessions_root: dir.to_string_lossy().into_owned(),
clock: Arc::new(SystemClock),
ids: Arc::new(DefaultIdGenerator::new()),
});
let id = harness.session().storage().metadata().id.clone();
let metas = list_session_metadata(cwd).await.map_err(|e| e.to_string())?;
let Some(source) = metas.iter().find(|m| m.id == id) else {
return Err(format!("current session {id} not found on disk"));
};
let fork_storage = repo
.fork_typed(
source,
&rpi_harness::session::jsonl::JsonlSessionCreateOptions {
id: None,
parent_session_id: Some(source.id.clone()),
cwd: cwd.to_string(),
metadata: None,
},
&rpi_harness::session::types::ForkOptions::default(),
)
.await
.map_err(|e| e.to_string())?;
Ok(Session::new(Arc::new(fork_storage), None))
}
async fn open_session(
meta: &rpi_harness::session::jsonl::JsonlSessionMetadata,
cwd: &str,
) -> Result<Session, BuildError> {
use rpi_harness::session::jsonl::{JsonlSessionRepo, JsonlSessionRepoOptions};
use rpi_harness::session::types::SessionStorage;
use rpi_tools::FileSystem;
let dir = default_session_dir(Path::new(cwd));
let env = Arc::new(OsExecutionEnv::with_cwd(PathBuf::from(cwd)));
let fs: Arc<dyn FileSystem> = env.clone();
let repo = JsonlSessionRepo::with_env_cwd(JsonlSessionRepoOptions {
fs: fs.clone(),
sessions_root: dir.to_string_lossy().into_owned(),
clock: Arc::new(SystemClock),
ids: Arc::new(DefaultIdGenerator::new()),
});
let storage = repo
.open_by_jsonl_metadata(meta)
.await
.map_err(|e| BuildError::SessionDir(format!("open {}: {e}", meta.path)))?;
let storage_arc: Arc<dyn SessionStorage> = Arc::new(storage);
Ok(Session::new(storage_arc, None))
}
fn ephemeral_session() -> Session {
let storage = Arc::new(InMemorySessionStorage::new(
SessionMetadata {
id: "ephemeral".into(),
created_at: 0,
parent_session_id: None,
},
Arc::new(SystemClock),
Arc::new(DefaultIdGenerator::new()),
));
Session::new(storage, None)
}
pub(crate) async fn create_jsonl_session(dir: &Path, cwd: &str) -> Result<Session, String> {
use rpi_harness::session::jsonl::{
JsonlSessionCreateOptions, JsonlSessionRepo, JsonlSessionRepoOptions,
};
use rpi_tools::FileSystem;
let env = Arc::new(OsExecutionEnv::with_cwd(PathBuf::from(cwd)));
let fs: Arc<dyn FileSystem> = env.clone();
let repo = JsonlSessionRepo::with_env_cwd(JsonlSessionRepoOptions {
fs: fs.clone(),
sessions_root: dir.to_string_lossy().into_owned(),
clock: Arc::new(SystemClock),
ids: Arc::new(DefaultIdGenerator::new()),
});
let opts = JsonlSessionCreateOptions {
id: None, parent_session_id: None,
cwd: cwd.to_string(),
metadata: None,
};
let storage = repo
.create_typed(&opts)
.await
.map_err(|e| format!("create session: {e}"))?;
let storage_arc: Arc<dyn rpi_harness::session::types::SessionStorage> = Arc::new(storage);
Ok(Session::new(storage_arc, None))
}
fn read_append_target(target: &str) -> Option<String> {
let path = Path::new(target);
if path.is_file() {
std::fs::read_to_string(path).ok()
} else {
None
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::args::Args;
#[test]
fn default_prompt_mentions_cwd_and_tools() {
let p = default_system_prompt("/tmp/proj");
assert!(p.contains("/tmp/proj"));
assert!(p.contains("read"));
assert!(p.contains("bash"));
assert!(p.contains("edit"));
assert!(p.contains("write"));
assert!(p.contains("grep"));
assert!(p.contains("find"));
assert!(p.contains("ls"));
}
#[test]
fn select_ephemeral_when_no_session() {
let args = Args { no_session: true, ..Args::default() };
let cwd = Path::new("/tmp");
assert!(matches!(select_session(&args, cwd), SessionSelection::Ephemeral));
}
#[test]
fn select_latest_for_continue_and_resume() {
let args = Args { continue_session: true, ..Args::default() };
let cwd = Path::new("/tmp");
assert!(matches!(select_session(&args, cwd), SessionSelection::Latest));
let args = Args { resume: true, ..Args::default() };
assert!(matches!(select_session(&args, cwd), SessionSelection::Latest));
}
#[test]
fn select_by_id_for_session_flag() {
let args = Args {
session: Some("01a02ece".into()),
..Args::default()
};
let cwd = Path::new("/tmp");
assert!(matches!(
select_session(&args, cwd),
SessionSelection::ById { id } if id == "01a02ece"
));
}
#[test]
fn select_new_with_custom_dir() {
let args = Args {
session_dir: Some(PathBuf::from("/tmp/sess")),
..Args::default()
};
let cwd = Path::new("/tmp");
match select_session(&args, cwd) {
SessionSelection::New { dir, .. } => assert_eq!(dir, PathBuf::from("/tmp/sess")),
other => panic!("expected New, got {other:?}"),
}
}
#[test]
fn select_new_default_dir() {
let args = Args::default();
let cwd = Path::new("/proj");
match select_session(&args, cwd) {
SessionSelection::New { dir, .. } => {
assert_eq!(dir, Path::new("/proj/.pi/sessions"));
}
other => panic!("expected New, got {other:?}"),
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn ephemeral_session_builds_roundtrips() {
let s = ephemeral_session();
let leaf = s.get_leaf_id().await;
assert!(leaf.is_ok());
}
}