use std::path::{Path, PathBuf};
use std::sync::atomic::Ordering;
use std::sync::{Arc, Mutex};
use rpi_agent::AgentTool;
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::{BranchBounds, EntryQuery, SessionMetadata};
use rpi_harness::session::Session;
use rpi_harness::system_prompt::compose_system_prompt;
use rpi_harness::types::{
AgentHarnessOptions, AgentHarnessResources, AgentHarnessStreamOptions, CompactionSettings,
DrivingMode, HarnessTool, HarnessToolExecution, RetryPolicy, ToolReplay,
};
use rpi_tools::{
create_bash_tool, create_edit_tool, create_read_tool, create_write_tool, ExecutionToolContext,
MutationQueueRegistry, OsExecutionEnv,
};
use crate::args::Args;
use crate::docs_tool::create_docs_tool;
use crate::extension_api::ExtensionBackend;
use crate::provider::ResolvedModel;
use crate::resource_dirs::{
discover_append_system_prompt_file_with_packages, discover_system_prompt_file_with_packages,
extension_dirs, global_extension_dirs, global_prompt_template_dirs, global_skill_dirs,
load_prompt_templates_with_precedence, load_skills_with_precedence,
project_prompt_template_dirs, project_skill_dirs, prompt_template_dirs, skill_dirs,
};
use rpi_extensions::{
emit_resources_discover, ExtensionEmitter, ExtensionSession, NullDiagnostics,
PluginDiagnostics, PluginToolAdapter, TeeEmitter,
};
pub const BUILTIN_TOOL_NAMES: &[&str] = &["read", "bash", "edit", "write", "docs"];
pub(crate) fn should_load_js_packages(args: &Args) -> bool {
args.enable_pi_packages && !args.no_extensions
}
pub(crate) fn package_resources_for(
args: &Args,
cwd: &Path,
project_trusted: bool,
) -> crate::packages::PackageResources {
if should_load_js_packages(args) {
if crate::args::offline_mode_enabled(args.offline) {
if project_trusted {
crate::packages::resolve_offline_from_settings(cwd)
} else {
crate::packages::resolve_offline_from_global_settings(cwd)
}
} else if project_trusted {
crate::packages::resolve_from_settings(cwd)
} else {
crate::packages::resolve_from_global_settings(cwd)
}
} else {
crate::packages::PackageResources::default()
}
}
pub(crate) fn package_resources_for_update_check(
args: &Args,
cwd: &Path,
project_trusted: bool,
) -> crate::packages::PackageResources {
if args.dev_local_only {
crate::packages::PackageResources::default()
} else if project_trusted {
crate::packages::discover_from_settings(cwd)
} else {
crate::packages::discover_from_global_settings(cwd)
}
}
pub fn default_system_prompt(cwd: &str) -> String {
format!(
"You are an expert coding assistant operating inside rpi, 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
- docs — Look up rpi usage, extension, package, and compatibility documentation
Guidelines:
- Be concise in your responses
- Show file paths clearly when working with files
- Prefer the smallest change that solves the problem
- When unsure about rpi commands, extensions, Pi package compatibility, or .rpi configuration, consult the project documentation before guessing
Current working directory: {cwd}"
)
}
#[derive(Debug, Clone)]
pub enum SessionSelection {
Ephemeral,
New { dir: PathBuf, name: Option<String> },
Latest,
ById { id: String },
ByExactId { id: String },
Fork { source: 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.fork {
return SessionSelection::Fork { source: s.clone() };
}
if let Some(s) = &args.session_id {
return SessionSelection::ByExactId { id: s.clone() };
}
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 {
let preferred = cwd.join(".rpi").join("sessions");
let legacy = cwd.join(".pi").join("sessions");
if preferred.exists() || !legacy.exists() {
preferred
} else {
legacy
}
}
pub async fn build(
resolved: &ResolvedModel,
args: &Args,
cwd: &Path,
project_trusted: bool,
) -> Result<
(
AgentHarness,
tokio::sync::broadcast::Receiver<rpi_agent::AgentEvent>,
ReloadContext,
),
BuildError,
> {
let cwd_str = cwd.to_string_lossy().to_string();
if !project_trusted && args.verbose {
eprintln!(
"warning: current-project settings, resources, and discovered extensions are explicitly disabled (use --approve or /trust yes to re-enable)"
);
}
let package_resources = if args.dev_local_only {
crate::packages::PackageResources::default()
} else {
package_resources_for(args, cwd, project_trusted)
};
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(),
args.unknown_flags.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, project_trusted, Some(Arc::clone(&action_bridge)))
};
let js_extension_session = if !should_load_js_packages(args) {
None
} else {
let paths = js_extension_paths(args, cwd, project_trusted, &package_resources);
let js_context = serde_json::json!({
"cwd": cwd_str,
"theme": resolved.theme.clone(),
"currentModel": resolved.model.clone(),
"models": catalog.clone(),
"thinkingLevel": resolved.thinking_level,
});
match crate::js_extensions::JsExtensionSession::load_with_context(
&paths,
args.verbose,
js_context,
) {
Ok(session) => session,
Err(error) => {
eprintln!("warning: JS/TS extensions were not loaded: {error}");
None
}
}
};
if js_extension_session.is_some() {
eprintln!(
"warning: enabled Pi JS/TS extensions execute with the current user's permissions"
);
}
if let Some(session) = &js_extension_session {
if let Err(error) =
session.enable_provider_runtime(resolved.provider.clone(), runtime.clone())
{
if args.verbose {
eprintln!("warning: JS provider runtime was not enabled: {error}");
}
}
}
if args.verbose {
if let Some(session) = &js_extension_session {
let info = session.backend_info();
eprintln!(
"JS extension backend: {} v{} ({})",
info.name,
info.api_version,
info.capability_names().join(", ")
);
}
if let Some(s) = extension_session.summary() {
eprintln!("extensions: {s}");
}
report_deferred_renderers(&extension_session);
}
merge_extension_tools(&mut tools, &extension_session, args);
if let Some(session) = &js_extension_session {
merge_js_extension_tools(&mut tools, session, args);
if args.verbose && !session.commands.is_empty() {
eprintln!("JS extension commands: {}", session.commands.join(", "));
}
}
let mut active = active_tool_names(&tools, args);
if let Some(session) = &js_extension_session {
let js_names = session.tool_names();
if let Some(js_active) = session.active_tools() {
active.retain(|name| !js_names.iter().any(|js| js == name));
active.extend(js_active.into_iter().filter(|name| {
js_names.iter().any(|js| js == name) && tool_name_allowed(name, args)
}));
}
}
active = filter_active_tool_names(active, args);
let selection = select_session(args, cwd);
let session = build_session(&selection, &cwd_str).await?;
if let Some(js) = &js_extension_session {
let session_id = session
.get_metadata()
.await
.ok()
.map(|metadata| metadata.id);
let leaf_id = session.get_leaf_id().await.ok().flatten();
if let Some(session_id) = session_id {
let branch = session
.find_entries_on_branch(&EntryQuery::default(), &BranchBounds::default())
.await
.ok()
.unwrap_or_default();
let branch_json =
serde_json::to_value(&branch).unwrap_or_else(|_| serde_json::json!([]));
let runtime_context = serde_json::json!({
"session": {
"id": session_id,
"leafId": leaf_id,
"branch": branch_json,
"entries": branch_json.clone(),
},
});
if let Err(error) = js.set_runtime_context(runtime_context) {
if args.verbose {
eprintln!("warning: could not sync JS session context: {error}");
}
}
}
}
let base_prompt = match args.system_prompt.as_deref() {
Some(explicit) => explicit.to_string(),
None if project_trusted => {
match discover_system_prompt_file_with_packages(cwd, &package_resources) {
Some(path) => std::fs::read_to_string(&path)
.unwrap_or_else(|_| default_system_prompt(&cwd_str)),
None => 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) = project_trusted
.then(|| discover_append_system_prompt_file_with_packages(cwd, &package_resources))
.flatten()
{
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 = if args.dev_local_only {
project_skill_dirs(cwd)
} else if project_trusted {
skill_dirs(cwd)
} else {
global_skill_dirs()
};
dirs.extend(args.skill.iter().cloned());
dirs.extend(discovered.skill_paths.iter().map(PathBuf::from));
if let Some(session) = &js_extension_session {
dirs.extend(session.resources.skill_paths.iter().cloned());
}
dirs.extend(package_resources.skill_dirs());
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 = if args.dev_local_only {
project_prompt_template_dirs(cwd)
} else if project_trusted {
prompt_template_dirs(cwd)
} else {
global_prompt_template_dirs()
};
dirs.extend(args.prompt_template.iter().cloned());
dirs.extend(discovered.prompt_paths.iter().map(PathBuf::from));
if let Some(session) = &js_extension_session {
dirs.extend(session.resources.prompt_paths.iter().cloned());
}
dirs.extend(package_resources.prompt_dirs());
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 || !project_trusted {
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 &package_resources.diagnostics {
eprintln!("warning: package {}: {}", d.spec, d.message);
}
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_with_packages(cwd, &package_resources).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) ---",
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 { .. }
| SessionSelection::ByExactId { .. }
| SessionSelection::Fork { .. }
),
stream_options: AgentHarnessStreamOptions {
timeout: args.timeout,
..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)),
js_extension_session: js_extension_session.clone(),
package_resources: Arc::new(package_resources.clone()),
action_bridge: Arc::new(Mutex::new(Some(Arc::clone(&action_bridge)))),
catalog,
gateway: resolved.provider.clone(),
runtime: runtime.clone(),
cwd: cwd.to_path_buf(),
project_trusted,
args: args.clone(),
resolved_model: resolved.model.clone(),
broadcast: broadcast_for_context,
mailbox: reload_mailbox,
dev_extension: None,
};
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 js_extension_session: Option<crate::js_extensions::JsExtensionSession>,
pub package_resources: Arc<crate::packages::PackageResources>,
pub action_bridge: ActionBridgeCell,
pub catalog: Vec<rpi_ai::Model>,
pub gateway: Arc<dyn Provider>,
pub runtime: tokio::runtime::Handle,
pub cwd: PathBuf,
pub project_trusted: bool,
pub args: Args,
pub resolved_model: rpi_ai::Model,
pub broadcast: Arc<dyn rpi_agent::AgentEmitter>,
pub mailbox: rpi_extensions::ReloadMailbox,
pub dev_extension: Option<Arc<crate::dev_extension::DevExtension>>,
}
pub struct ReloadOutcome {
pub summary: String,
pub had_warnings: bool,
}
struct PreparedReloadInputs {
package_resources: crate::packages::PackageResources,
extension_dirs: Vec<PathBuf>,
skill_base_dirs: Vec<PathBuf>,
prompt_base_dirs: Vec<PathBuf>,
}
fn append_reload_resource_paths(
mut paths: Vec<PathBuf>,
explicit: &[PathBuf],
discovered: &[String],
js_paths: &[PathBuf],
package_paths: &[PathBuf],
) -> Vec<PathBuf> {
paths.extend(explicit.iter().cloned());
paths.extend(discovered.iter().map(PathBuf::from));
paths.extend(js_paths.iter().cloned());
paths.extend(package_paths.iter().cloned());
paths
}
pub async fn reload_extension_resources(
harness: &AgentHarness,
ctx: &ReloadContext,
) -> ReloadOutcome {
reload_extension_resources_inner(harness, ctx, || {}).await
}
async fn reload_extension_resources_inner<F>(
harness: &AgentHarness,
ctx: &ReloadContext,
after_prepare: F,
) -> ReloadOutcome
where
F: FnOnce() + Send,
{
let mut effective_args = ctx.args.clone();
if let Some(dev) = &ctx.dev_extension {
if let Err(error) = dev
.rebuild()
.and_then(|_| dev.apply_to_args(&mut effective_args))
{
return ReloadOutcome {
summary: format!(
"Extension build failed for {}: {error}. Keeping the currently loaded version.",
dev.package_name()
),
had_warnings: true,
};
}
}
let cwd_str = ctx.cwd.to_string_lossy().to_string();
let project_trusted = resolve_project_trust(&effective_args, &ctx.cwd);
let prepared = match prepare_reload_inputs(&effective_args, &ctx.cwd, project_trusted) {
Ok(prepared) => prepared,
Err(error) => {
return ReloadOutcome {
summary: format!(
"Settings reload failed: {error}. Keeping the currently loaded resources."
),
had_warnings: true,
};
}
};
after_prepare();
let PreparedReloadInputs {
package_resources,
extension_dirs,
skill_base_dirs,
prompt_base_dirs,
} = prepared;
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(),
ctx.args.unknown_flags.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 effective_args.no_extensions {
rpi_extensions::ExtensionSession::none()
} else {
load_extensions_from_dirs(
&effective_args,
&extension_dirs,
Some(Arc::clone(&fresh_bridge)),
)
};
if extension_session.is_empty() && !effective_args.no_extensions {
}
if effective_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 !effective_args.no_skills {
let js_paths: &[PathBuf] = ctx
.js_extension_session
.as_ref()
.map(|js| js.resources.skill_paths.as_slice())
.unwrap_or(&[]);
let package_paths = package_resources.skill_dirs();
let dirs = append_reload_resource_paths(
skill_base_dirs,
&effective_args.skill,
&discovered.skill_paths,
js_paths,
&package_paths,
);
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 !effective_args.no_prompt_templates {
let js_paths: &[PathBuf] = ctx
.js_extension_session
.as_ref()
.map(|js| js.resources.prompt_paths.as_slice())
.unwrap_or(&[]);
let package_paths = package_resources.prompt_dirs();
let dirs = append_reload_resource_paths(
prompt_base_dirs,
&effective_args.prompt_template,
&discovered.prompt_paths,
js_paths,
&package_paths,
);
let result = load_prompt_templates_with_precedence(&env_dyn, &dirs).await;
prompt_templates = result.prompt_templates;
prompt_diags = result.diagnostics;
}
let context_block = if effective_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()
|| !package_resources.diagnostics.is_empty()
{
warnings = true;
if effective_args.verbose {
for d in &package_resources.diagnostics {
eprintln!("warning: package {}: {}", d.spec, d.message);
}
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 effective_args.system_prompt.as_deref() {
Some(explicit) => explicit.to_string(),
None => match discover_system_prompt_file_with_packages(&ctx.cwd, &package_resources) {
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 &effective_args.append_system_prompt {
let text = read_append_target(extra).unwrap_or_else(|| extra.clone());
append_texts.push(text);
}
if effective_args.append_system_prompt.is_empty() {
if let Some(path) =
discover_append_system_prompt_file_with_packages(&ctx.cwd, &package_resources)
{
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, &effective_args);
merge_extension_tools(&mut tools, &extension_session, &effective_args);
if let Some(js) = &ctx.js_extension_session {
merge_js_extension_tools(&mut tools, js, &effective_args);
}
let mut active = active_tool_names(&tools, &effective_args);
if let Some(js) = &ctx.js_extension_session {
let js_names = js.tool_names();
if let Some(js_active) = js.active_tools() {
active.retain(|name| {
tool_name_allowed(name, &effective_args)
&& !js_names.iter().any(|js_name| js_name == name)
});
active.extend(js_active.into_iter().filter(|name| {
js_names.iter().any(|js_name| js_name == name)
&& tool_name_allowed(name, &effective_args)
}));
}
}
active = filter_active_tool_names(active, &effective_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,
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
struct ReloadSettingsFields {
packages: Option<Vec<crate::settings::PackageSetting>>,
npm_command: Option<Vec<String>>,
skill_dirs: Option<Vec<String>>,
prompt_dirs: Option<Vec<String>>,
extension_dirs: Option<Vec<String>>,
}
impl From<crate::settings::Settings> for ReloadSettingsFields {
fn from(settings: crate::settings::Settings) -> Self {
Self {
packages: settings.packages,
npm_command: settings.npm_command,
skill_dirs: settings.skill_dirs,
prompt_dirs: settings.prompt_dirs,
extension_dirs: settings.extension_dirs,
}
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
struct ReloadSettingsSnapshot {
global: ReloadSettingsFields,
project: Option<ReloadSettingsFields>,
}
fn reload_reads_settings(args: &Args) -> bool {
!args.dev_local_only
&& (should_load_js_packages(args)
|| !args.no_extensions
|| !args.no_skills
|| !args.no_prompt_templates)
}
fn load_reload_settings_snapshot(
args: &Args,
cwd: &Path,
project_trusted: bool,
) -> Result<Option<ReloadSettingsSnapshot>, String> {
if !reload_reads_settings(args) {
return Ok(None);
}
let global = crate::settings::load_settings()
.map_err(|error| format!("could not load global settings: {error}"))?
.into();
let project = if project_trusted {
crate::settings::load_active_project_settings(cwd)
.map_err(|error| format!("could not load project settings: {error}"))?
.map(|(_, settings)| settings.into())
} else {
None
};
Ok(Some(ReloadSettingsSnapshot { global, project }))
}
fn validate_settings_for_reload(
args: &Args,
cwd: &Path,
project_trusted: bool,
) -> Result<(), String> {
load_reload_settings_snapshot(args, cwd, project_trusted).map(|_| ())
}
fn prepare_reload_inputs(
args: &Args,
cwd: &Path,
project_trusted: bool,
) -> Result<PreparedReloadInputs, String> {
prepare_reload_inputs_inner(args, cwd, project_trusted, || {})
}
fn prepare_reload_inputs_inner<F>(
args: &Args,
cwd: &Path,
project_trusted: bool,
before_verify: F,
) -> Result<PreparedReloadInputs, String>
where
F: FnOnce(),
{
let settings_before = load_reload_settings_snapshot(args, cwd, project_trusted)?;
let package_resources = if args.dev_local_only {
crate::packages::PackageResources::default()
} else {
package_resources_for(args, cwd, project_trusted)
};
let mut extension_dirs = if args.no_extensions || args.dev_local_only {
Vec::new()
} else if project_trusted {
extension_dirs(cwd)
} else {
global_extension_dirs()
};
if !args.no_extensions {
extension_dirs.extend(args.extensions_dir.iter().cloned());
}
let skill_base_dirs = if args.no_skills {
Vec::new()
} else if args.dev_local_only {
project_skill_dirs(cwd)
} else if project_trusted {
skill_dirs(cwd)
} else {
global_skill_dirs()
};
let prompt_base_dirs = if args.no_prompt_templates {
Vec::new()
} else if args.dev_local_only {
project_prompt_template_dirs(cwd)
} else if project_trusted {
prompt_template_dirs(cwd)
} else {
global_prompt_template_dirs()
};
before_verify();
let settings_after = load_reload_settings_snapshot(args, cwd, project_trusted)?;
if settings_before != settings_after {
return Err(
"settings changed while reload inputs were being prepared; retry /reload".to_string(),
);
}
Ok(PreparedReloadInputs {
package_resources,
extension_dirs,
skill_base_dirs,
prompt_base_dirs,
})
}
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, {} message-render, {} entry-render (active)",
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
}
pub(crate) fn resolve_project_trust(args: &Args, cwd: &Path) -> bool {
if let Some(override_value) = args.trust_override {
return override_value;
}
crate::config::project_trust_decision(cwd)
.ok()
.flatten()
.unwrap_or(true)
}
fn load_extensions(
args: &Args,
cwd: &Path,
project_trusted: bool,
action_bridge: Option<Arc<rpi_extensions::ActionBridge>>,
) -> ExtensionSession {
let mut dirs = if args.dev_local_only {
Vec::new()
} else if project_trusted {
extension_dirs(cwd)
} else {
global_extension_dirs()
};
dirs.extend(args.extensions_dir.iter().cloned());
load_extensions_from_dirs(args, &dirs, action_bridge)
}
fn load_extensions_from_dirs(
args: &Args,
dirs: &[PathBuf],
action_bridge: Option<Arc<rpi_extensions::ActionBridge>>,
) -> ExtensionSession {
let diagnostics: Arc<dyn PluginDiagnostics> = Arc::new(NullDiagnostics);
rpi_extensions::load_session_mixed(dirs, &args.extension, diagnostics, action_bridge)
}
fn js_extension_paths(
args: &Args,
cwd: &Path,
project_trusted: bool,
packages: &crate::packages::PackageResources,
) -> Vec<PathBuf> {
if args.dev_local_only {
return Vec::new();
}
let mut paths = packages.extension_paths();
let discovered_dirs = if project_trusted {
extension_dirs(cwd)
} else {
global_extension_dirs()
};
for dir in discovered_dirs {
if let Ok(entries) = std::fs::read_dir(dir) {
paths.extend(entries.flatten().map(|entry| entry.path()).filter(|path| {
matches!(
path.extension()
.and_then(|ext| ext.to_str())
.map(|ext| ext.to_ascii_lowercase())
.as_deref(),
Some("js" | "mjs" | "cjs" | "ts" | "tsx")
)
}));
}
}
paths.extend(
args.extension
.iter()
.filter(|path| {
matches!(
path.extension()
.and_then(|ext| ext.to_str())
.map(|ext| ext.to_ascii_lowercase())
.as_deref(),
Some("js" | "mjs" | "cjs" | "ts" | "tsx")
)
})
.cloned(),
);
paths
}
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 !tool_name_allowed(name, args) {
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 merge_js_extension_tools(
tools: &mut Vec<HarnessTool>,
session: &crate::js_extensions::JsExtensionSession,
args: &Args,
) {
for adapter in session.tools() {
let name = adapter.schema().name.clone();
if !tool_name_allowed(&name, args) {
continue;
}
let harness_tool = HarnessTool::new(Arc::new(adapter));
match tools
.iter_mut()
.find(|tool| tool.tool.schema().name == name)
{
Some(slot) => *slot = harness_tool,
None => tools.push(harness_tool),
}
}
}
pub fn bash_options() -> rpi_tools::tools::bash::BashToolOptions {
use rpi_tools::tools::bash::BashToolOptions;
let default = std::env::var("RPI_BASH_TIMEOUT")
.ok()
.and_then(|v| v.parse::<f64>().ok())
.unwrap_or(120.0);
BashToolOptions {
command_prefix: None,
default_timeout: Some(default),
}
}
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, Some(bash_options()))),
),
("edit", HarnessTool::new(create_edit_tool(ctx))),
("write", HarnessTool::new(create_write_tool(ctx))),
("docs", HarnessTool::new(create_docs_tool())),
];
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> {
filter_active_tool_names(
tools.iter().map(|tool| tool.tool.schema().name.clone()),
args,
)
}
pub(crate) fn tool_name_allowed(name: &str, args: &Args) -> bool {
if args.no_tools {
return false;
}
if args
.tools
.as_ref()
.is_some_and(|allow| !allow.iter().any(|value| value == name))
{
return false;
}
if args
.exclude_tools
.as_ref()
.is_some_and(|deny| deny.iter().any(|value| value == name))
{
return false;
}
true
}
pub(crate) fn filter_active_tool_names<I>(names: I, args: &Args) -> Vec<String>
where
I: IntoIterator<Item = String>,
{
names
.into_iter()
.filter(|name| tool_name_allowed(name, args))
.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 { .. }
| SessionSelection::ByExactId { .. } => restore_session(selection, cwd).await,
SessionSelection::Fork { source } => fork_session_at_launch(source, 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),
}),
SessionSelection::ByExactId { id } => {
let metas = list_session_metadata(cwd).await?;
if let Some(meta) = metas.iter().find(|m| m.id == *id) {
return open_session(meta, cwd).await;
}
let dir = default_session_dir(Path::new(cwd));
std::fs::create_dir_all(&dir)
.map_err(|e| BuildError::SessionDir(format!("{}: {e}", dir.display())))?;
create_jsonl_session_with_id(&dir, cwd, Some(id.clone()))
.await
.map_err(|e| BuildError::SessionDir(format!("{}: {e}", dir.display())))
}
_ => unreachable!("restore_session only called for Latest/ById/ByExactId"),
}
}
async fn fork_session_at_launch(source: &str, cwd: &str) -> Result<Session, BuildError> {
use rpi_harness::session::jsonl::{
JsonlSessionCreateOptions, JsonlSessionRepo, JsonlSessionRepoOptions,
};
use rpi_harness::session::types::{ForkOptions, SessionStorage};
use rpi_tools::FileSystem;
let dir = default_session_dir(Path::new(cwd));
std::fs::create_dir_all(&dir)
.map_err(|e| BuildError::SessionDir(format!("{}: {e}", dir.display())))?;
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 metas = repo
.list_typed(&rpi_harness::session::jsonl::JsonlSessionListOptions::default())
.await
.map_err(|e| BuildError::SessionDir(format!("list sessions: {e}")))?;
let source_meta = metas
.iter()
.find(|m| m.id == *source || m.path.contains(source) || source.contains(&m.id))
.ok_or_else(|| BuildError::SessionNotFound {
requested: format!("--fork {source}"),
dir: dir.display().to_string(),
})?;
let fork_storage = repo
.fork_typed(
source_meta,
&JsonlSessionCreateOptions {
id: None,
parent_session_id: Some(source_meta.id.clone()),
cwd: cwd.to_string(),
metadata: None,
},
&ForkOptions::default(),
)
.await
.map_err(|e| BuildError::SessionDir(format!("fork {}: {e}", source_meta.path)))?;
let storage_arc: Arc<dyn SessionStorage> = Arc::new(fork_storage);
Ok(Session::new(storage_arc, None))
}
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> {
create_jsonl_session_with_id(dir, cwd, None).await
}
pub(crate) async fn create_jsonl_session_with_id(
dir: &Path,
cwd: &str,
id: Option<String>,
) -> 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, 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("docs"));
assert!(!p.contains("- grep"));
assert!(!p.contains("- find"));
assert!(!p.contains("- ls"));
assert!(!p.contains("powershell"));
}
#[test]
fn default_tools_include_docs_lookup() {
let env = Arc::new(OsExecutionEnv::with_cwd(PathBuf::from(".")));
let env_dyn: Arc<dyn rpi_tools::ExecutionEnv> = env.clone();
let mut_env: Arc<dyn rpi_tools::MutatingEnv> = env.clone();
let context = ExecutionToolContext::new(env_dyn, Some(mut_env));
let names: Vec<String> = build_tools(&context, &Args::default())
.iter()
.map(|tool| tool.tool.schema().name.clone())
.collect();
assert_eq!(names, vec!["read", "bash", "edit", "write", "docs"]);
}
#[test]
fn pi_package_loading_is_opt_in_and_respects_no_extensions() {
let args = Args::default();
assert!(!should_load_js_packages(&args));
let resources = package_resources_for(&args, Path::new("."), false);
assert!(resources.packages.is_empty());
let args = Args {
enable_pi_packages: true,
..Args::default()
};
assert!(should_load_js_packages(&args));
let args = Args {
enable_pi_packages: true,
no_extensions: true,
..Args::default()
};
assert!(!should_load_js_packages(&args));
assert!(package_resources_for(&args, Path::new("."), false)
.packages
.is_empty());
}
#[test]
fn pi_offline_env_disables_startup_package_remediation() {
struct RestoreEnv {
name: &'static str,
value: Option<std::ffi::OsString>,
}
impl Drop for RestoreEnv {
fn drop(&mut self) {
match self.value.take() {
Some(value) => std::env::set_var(self.name, value),
None => std::env::remove_var(self.name),
}
}
}
let _guard = crate::config::test_support::env_lock().lock().unwrap();
let _restore_config = RestoreEnv {
name: crate::config::CONFIG_DIR_ENV,
value: std::env::var_os(crate::config::CONFIG_DIR_ENV),
};
let _restore_offline = RestoreEnv {
name: crate::args::PI_OFFLINE_ENV,
value: std::env::var_os(crate::args::PI_OFFLINE_ENV),
};
let tmp = tempfile::tempdir().unwrap();
let agent = tmp.path().join("agent");
let cwd = tmp.path().join("project");
let package = agent.join("npm/node_modules/demo");
std::fs::create_dir_all(package.join("extensions")).unwrap();
std::fs::create_dir_all(&cwd).unwrap();
std::fs::write(
package.join("package.json"),
r#"{"name":"demo","version":"1.0.0"}"#,
)
.unwrap();
std::fs::write(
package.join("extensions/index.js"),
"export default () => {};",
)
.unwrap();
std::fs::write(
agent.join("settings.json"),
r#"{"npmCommand":[""],"packages":["npm:demo@2.0.0"]}"#,
)
.unwrap();
std::env::set_var(crate::config::CONFIG_DIR_ENV, &agent);
std::env::set_var(crate::args::PI_OFFLINE_ENV, "TrUe");
let args = Args {
enable_pi_packages: true,
offline: false,
..Args::default()
};
let resources = package_resources_for(&args, &cwd, false);
assert!(resources.packages.is_empty());
assert_eq!(resources.diagnostics.len(), 1);
assert!(resources.diagnostics[0].message.contains("offline"));
assert_eq!(
std::fs::read(package.join("package.json")).unwrap(),
br#"{"name":"demo","version":"1.0.0"}"#
);
}
#[test]
fn package_discovery_uses_the_callers_trust_snapshot() {
let _guard = crate::config::test_support::env_lock().lock().unwrap();
let previous = std::env::var_os(crate::config::CONFIG_DIR_ENV);
let tmp = tempfile::tempdir().unwrap();
let agent = tmp.path().join("agent");
let cwd = tmp.path().join("project");
let package = cwd.join("package");
std::fs::create_dir_all(&agent).unwrap();
std::fs::create_dir_all(cwd.join(".rpi")).unwrap();
std::fs::create_dir_all(&package).unwrap();
std::fs::write(agent.join("settings.json"), "{}").unwrap();
std::fs::write(
package.join("package.json"),
r#"{"name":"snapshot-package","version":"1.0.0"}"#,
)
.unwrap();
std::fs::write(
cwd.join(".rpi/settings.json"),
serde_json::json!({"packages": [package]}).to_string(),
)
.unwrap();
std::env::set_var(crate::config::CONFIG_DIR_ENV, &agent);
let args = Args {
enable_pi_packages: true,
..Args::default()
};
assert_eq!(package_resources_for(&args, &cwd, true).packages.len(), 1);
assert!(package_resources_for(&args, &cwd, false)
.packages
.is_empty());
assert_eq!(
package_resources_for_update_check(&args, &cwd, true)
.packages
.len(),
1
);
assert!(package_resources_for_update_check(&args, &cwd, false)
.packages
.is_empty());
match previous {
Some(value) => std::env::set_var(crate::config::CONFIG_DIR_ENV, value),
None => std::env::remove_var(crate::config::CONFIG_DIR_ENV),
}
}
#[test]
fn reload_settings_preflight_is_independent_of_packages_and_respects_trust() {
struct RestoreConfigDir(Option<std::ffi::OsString>);
impl Drop for RestoreConfigDir {
fn drop(&mut self) {
match self.0.take() {
Some(value) => std::env::set_var(crate::config::CONFIG_DIR_ENV, value),
None => std::env::remove_var(crate::config::CONFIG_DIR_ENV),
}
}
}
let _guard = crate::config::test_support::env_lock().lock().unwrap();
let _restore = RestoreConfigDir(std::env::var_os(crate::config::CONFIG_DIR_ENV));
let tmp = tempfile::tempdir().unwrap();
let agent = tmp.path().join("agent");
let cwd = tmp.path().join("project");
std::fs::create_dir_all(&agent).unwrap();
std::fs::create_dir_all(cwd.join(".rpi")).unwrap();
std::fs::write(agent.join("settings.json"), "{}").unwrap();
std::fs::write(cwd.join(".rpi/settings.json"), "{ malformed").unwrap();
std::env::set_var(crate::config::CONFIG_DIR_ENV, &agent);
let packages_disabled = Args::default();
let packages_enabled = Args {
enable_pi_packages: true,
..Args::default()
};
for args in [&packages_disabled, &packages_enabled] {
let trusted = validate_settings_for_reload(args, &cwd, true);
let untrusted = validate_settings_for_reload(args, &cwd, false);
assert!(trusted
.unwrap_err()
.contains("could not load project settings"));
assert!(untrusted.is_ok());
}
std::fs::write(agent.join("settings.json"), "{ malformed").unwrap();
for args in [&packages_disabled, &packages_enabled] {
assert!(validate_settings_for_reload(args, &cwd, false)
.unwrap_err()
.contains("could not load global settings"));
}
let local_only = Args {
dev_local_only: true,
..Args::default()
};
assert!(validate_settings_for_reload(&local_only, &cwd, true).is_ok());
}
#[tokio::test(flavor = "current_thread")]
async fn reload_with_packages_disabled_preserves_live_resources_when_settings_break() {
struct RestoreConfigDir(Option<std::ffi::OsString>);
impl Drop for RestoreConfigDir {
fn drop(&mut self) {
match self.0.take() {
Some(value) => std::env::set_var(crate::config::CONFIG_DIR_ENV, value),
None => std::env::remove_var(crate::config::CONFIG_DIR_ENV),
}
}
}
let _guard = crate::config::test_support::env_lock().lock().unwrap();
let _restore = RestoreConfigDir(std::env::var_os(crate::config::CONFIG_DIR_ENV));
let tmp = tempfile::tempdir().unwrap();
let agent = tmp.path().join("agent");
let cwd = tmp.path().join("project");
let skill_dir = tmp.path().join("configured-skills");
std::fs::create_dir_all(&agent).unwrap();
std::fs::create_dir_all(&cwd).unwrap();
std::fs::create_dir_all(&skill_dir).unwrap();
std::fs::write(
skill_dir.join("SKILL.md"),
"---\nname: keep-me\ndescription: Reload sentinel\n---\nKeep this skill loaded.",
)
.unwrap();
std::fs::write(
agent.join("settings.json"),
serde_json::json!({"skillDirs": [skill_dir]}).to_string(),
)
.unwrap();
std::env::set_var(crate::config::CONFIG_DIR_ENV, &agent);
let resolved = crate::provider::resolve(
Some("anthropic"),
Some(crate::provider::DEFAULT_MODEL_ID),
None,
Some("test-key"),
None,
)
.unwrap();
let args = Args {
trust_override: Some(false),
no_session: true,
no_extensions: true,
no_prompt_templates: true,
no_context_files: true,
system_prompt: Some("stable system prompt".into()),
..Args::default()
};
assert!(!should_load_js_packages(&args));
let (harness, _events, context) = build(&resolved, &args, &cwd, false).await.unwrap();
let before_resources = harness.get_resources().await.unwrap();
assert_eq!(before_resources.skills.as_ref().unwrap().len(), 1);
assert_eq!(before_resources.skills.as_ref().unwrap()[0].name, "keep-me");
let before_prompt = harness.get_system_prompt().await.unwrap();
let before_tools: Vec<String> = harness
.get_tools()
.await
.unwrap()
.iter()
.map(|tool| tool.tool.schema().name.clone())
.collect();
let before_bridge = context
.action_bridge
.lock()
.unwrap()
.as_ref()
.unwrap()
.clone();
std::fs::write(agent.join("settings.json"), "{ malformed").unwrap();
let outcome = reload_extension_resources(&harness, &context).await;
assert!(outcome.had_warnings);
assert!(outcome.summary.contains("Settings reload failed"));
assert!(outcome.summary.contains("could not load global settings"));
assert_eq!(harness.get_resources().await.unwrap(), before_resources);
assert_eq!(harness.get_system_prompt().await.unwrap(), before_prompt);
let after_tools: Vec<String> = harness
.get_tools()
.await
.unwrap()
.iter()
.map(|tool| tool.tool.schema().name.clone())
.collect();
assert_eq!(after_tools, before_tools);
let after_bridge = context
.action_bridge
.lock()
.unwrap()
.as_ref()
.unwrap()
.clone();
assert!(Arc::ptr_eq(&after_bridge, &before_bridge));
}
#[test]
fn reload_preparation_rejects_settings_changed_during_derivation() {
struct RestoreConfigDir(Option<std::ffi::OsString>);
impl Drop for RestoreConfigDir {
fn drop(&mut self) {
match self.0.take() {
Some(value) => std::env::set_var(crate::config::CONFIG_DIR_ENV, value),
None => std::env::remove_var(crate::config::CONFIG_DIR_ENV),
}
}
}
let _guard = crate::config::test_support::env_lock().lock().unwrap();
let _restore = RestoreConfigDir(std::env::var_os(crate::config::CONFIG_DIR_ENV));
let tmp = tempfile::tempdir().unwrap();
let agent = tmp.path().join("agent");
let cwd = tmp.path().join("project");
let first_skill_dir = tmp.path().join("first-skills");
let second_skill_dir = tmp.path().join("second-skills");
std::fs::create_dir_all(&agent).unwrap();
std::fs::create_dir_all(&cwd).unwrap();
std::fs::write(
agent.join("settings.json"),
serde_json::json!({"skillDirs": [first_skill_dir]}).to_string(),
)
.unwrap();
std::env::set_var(crate::config::CONFIG_DIR_ENV, &agent);
let args = Args {
trust_override: Some(false),
no_extensions: true,
no_prompt_templates: true,
..Args::default()
};
let result = prepare_reload_inputs_inner(&args, &cwd, false, || {
std::fs::write(
agent.join("settings.json"),
serde_json::json!({"skillDirs": [second_skill_dir]}).to_string(),
)
.unwrap();
});
let error = result
.err()
.expect("settings mutation must fail preparation");
assert!(error.contains("settings changed while reload inputs were being prepared"));
}
#[tokio::test(flavor = "current_thread")]
async fn reload_uses_frozen_settings_inputs_after_preparation() {
struct RestoreConfigDir(Option<std::ffi::OsString>);
impl Drop for RestoreConfigDir {
fn drop(&mut self) {
match self.0.take() {
Some(value) => std::env::set_var(crate::config::CONFIG_DIR_ENV, value),
None => std::env::remove_var(crate::config::CONFIG_DIR_ENV),
}
}
}
let _guard = crate::config::test_support::env_lock().lock().unwrap();
let _restore = RestoreConfigDir(std::env::var_os(crate::config::CONFIG_DIR_ENV));
let tmp = tempfile::tempdir().unwrap();
let agent = tmp.path().join("agent");
let cwd = tmp.path().join("project");
let skill_dir = tmp.path().join("configured-skills");
std::fs::create_dir_all(&agent).unwrap();
std::fs::create_dir_all(&cwd).unwrap();
std::fs::create_dir_all(&skill_dir).unwrap();
std::fs::write(
skill_dir.join("SKILL.md"),
"---\nname: frozen-skill\ndescription: Reload sentinel\n---\nFrozen input.",
)
.unwrap();
std::fs::write(
agent.join("settings.json"),
serde_json::json!({"skillDirs": [skill_dir]}).to_string(),
)
.unwrap();
std::env::set_var(crate::config::CONFIG_DIR_ENV, &agent);
let resolved = crate::provider::resolve(
Some("anthropic"),
Some(crate::provider::DEFAULT_MODEL_ID),
None,
Some("test-key"),
None,
)
.unwrap();
let args = Args {
trust_override: Some(false),
no_session: true,
no_extensions: true,
no_prompt_templates: true,
no_context_files: true,
system_prompt: Some("stable system prompt".into()),
..Args::default()
};
let (harness, _events, context) = build(&resolved, &args, &cwd, false).await.unwrap();
let outcome = reload_extension_resources_inner(&harness, &context, || {
std::fs::write(agent.join("settings.json"), "{ malformed").unwrap();
})
.await;
assert!(!outcome.summary.contains("Settings reload failed"));
assert_eq!(
outcome.summary,
"Reloaded 0 plugin(s), 1 skill(s), 0 prompt(s)."
);
let resources = harness.get_resources().await.unwrap();
let skills = resources.skills.as_ref().unwrap();
assert_eq!(skills.len(), 1);
assert_eq!(skills[0].name, "frozen-skill");
}
#[tokio::test]
async fn reload_resource_paths_include_explicit_skill_and_prompt_files() {
let tmp = tempfile::tempdir().unwrap();
let skill_path = tmp.path().join("explicit-skill.md");
std::fs::write(
&skill_path,
"---\nname: explicit-skill\ndescription: Explicit skill\n---\nSkill body",
)
.unwrap();
let prompt_path = tmp.path().join("explicit-prompt.md");
std::fs::write(
&prompt_path,
"---\ndescription: Explicit prompt\n---\nPrompt body",
)
.unwrap();
let args = Args {
skill: vec![skill_path.clone()],
prompt_template: vec![prompt_path.clone()],
..Args::default()
};
let skill_paths = append_reload_resource_paths(Vec::new(), &args.skill, &[], &[], &[]);
let prompt_paths =
append_reload_resource_paths(Vec::new(), &args.prompt_template, &[], &[], &[]);
let env = Arc::new(OsExecutionEnv::with_cwd(tmp.path().to_path_buf()));
let env_dyn: Arc<dyn rpi_tools::ExecutionEnv> = env;
let skills = load_skills_with_precedence(&env_dyn, &skill_paths).await;
assert_eq!(skills.skills.len(), 1, "{:?}", skills.diagnostics);
assert_eq!(skills.skills[0].name, "explicit-skill");
let prompts = load_prompt_templates_with_precedence(&env_dyn, &prompt_paths).await;
assert_eq!(
prompts.prompt_templates.len(),
1,
"{:?}",
prompts.diagnostics
);
assert_eq!(prompts.prompt_templates[0].name, "explicit-prompt");
}
#[test]
fn tool_policy_applies_to_rust_and_js_active_names() {
let names = vec![
"read".to_string(),
"ask_user_question".to_string(),
"write".to_string(),
];
let args = Args {
tools: Some(vec!["read".into(), "ask_user_question".into()]),
..Args::default()
};
assert_eq!(
filter_active_tool_names(names.clone(), &args),
vec!["read", "ask_user_question"]
);
let args = Args {
exclude_tools: Some(vec!["ask_user_question".into()]),
..Args::default()
};
assert_eq!(
filter_active_tool_names(names.clone(), &args),
vec!["read", "write"]
);
let args = Args {
no_tools: true,
..Args::default()
};
assert!(filter_active_tool_names(names, &args).is_empty());
}
#[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 project_resources_load_by_default_and_allow_explicit_opt_out() {
let defaults = Args::default();
assert!(resolve_project_trust(
&defaults,
Path::new("C:/definitely-not-a-project")
));
let denied = Args {
trust_override: Some(false),
..Args::default()
};
assert!(!resolve_project_trust(
&denied,
Path::new("C:/definitely-not-a-project")
));
}
#[tokio::test(flavor = "current_thread")]
async fn build_preserves_the_supplied_startup_project_snapshot() {
struct RestoreConfigDir(Option<std::ffi::OsString>);
impl Drop for RestoreConfigDir {
fn drop(&mut self) {
match self.0.take() {
Some(value) => std::env::set_var(crate::config::CONFIG_DIR_ENV, value),
None => std::env::remove_var(crate::config::CONFIG_DIR_ENV),
}
}
}
let _guard = crate::config::test_support::env_lock().lock().unwrap();
let _restore = RestoreConfigDir(std::env::var_os(crate::config::CONFIG_DIR_ENV));
let tmp = tempfile::tempdir().unwrap();
let agent = tmp.path().join("agent");
let cwd = tmp.path().join("project");
std::fs::create_dir_all(&agent).unwrap();
std::fs::create_dir_all(&cwd).unwrap();
std::env::set_var(crate::config::CONFIG_DIR_ENV, &agent);
let resolved = crate::provider::resolve(
Some("anthropic"),
Some(crate::provider::DEFAULT_MODEL_ID),
None,
Some("test-key"),
None,
)
.unwrap();
let args = Args {
trust_override: Some(false),
no_session: true,
no_tools: true,
no_extensions: true,
no_skills: true,
no_prompt_templates: true,
no_context_files: true,
system_prompt: Some("test prompt".into()),
..Args::default()
};
let (_, _, context) = build(&resolved, &args, &cwd, true).await.unwrap();
assert_eq!(context.cwd, cwd);
assert!(context.project_trusted);
}
#[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/.rpi/sessions"));
}
other => panic!("expected New, got {other:?}"),
}
}
#[test]
fn default_session_dir_prefers_rpi_but_reads_legacy_pi() {
let tmp = tempfile::tempdir().unwrap();
let cwd = tmp.path();
std::fs::create_dir_all(cwd.join(".pi/sessions")).unwrap();
assert_eq!(default_session_dir(cwd), cwd.join(".pi/sessions"));
std::fs::create_dir_all(cwd.join(".rpi/sessions")).unwrap();
assert_eq!(default_session_dir(cwd), cwd.join(".rpi/sessions"));
}
#[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());
}
}